mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-11 18:53:42 +00:00
feat: gate assignment publication on allocation
This commit is contained in:
@@ -49,16 +49,18 @@ type Allocator struct {
|
||||
mu sync.Mutex
|
||||
servers map[string]ReadyServer
|
||||
allocations map[string]Allocation
|
||||
assignments map[string]Assignment
|
||||
requestHashes map[string][32]byte
|
||||
}
|
||||
|
||||
var (
|
||||
ErrNoCapacity = fmt.Errorf("no compatible ready server")
|
||||
ErrAllocationInput = fmt.Errorf("invalid allocation request")
|
||||
ErrNoCapacity = fmt.Errorf("no compatible ready server")
|
||||
ErrAllocationInput = fmt.Errorf("invalid allocation request")
|
||||
ErrAllocationNotFound = fmt.Errorf("allocation not found")
|
||||
)
|
||||
|
||||
func NewAllocator(servers []ReadyServer) (*Allocator, error) {
|
||||
a := &Allocator{servers: make(map[string]ReadyServer, len(servers)), allocations: make(map[string]Allocation), requestHashes: make(map[string][32]byte)}
|
||||
a := &Allocator{servers: make(map[string]ReadyServer, len(servers)), allocations: make(map[string]Allocation), assignments: make(map[string]Assignment), requestHashes: make(map[string][32]byte)}
|
||||
for _, server := range servers {
|
||||
if server.ServerID == "" || server.Region == "" || server.Build == "" || server.Protocol <= 0 || (server.Transport != "enet" && server.Transport != "steam_sdr") || server.State != ServerReady {
|
||||
return nil, fmt.Errorf("%w: invalid ready server", ErrAllocationInput)
|
||||
@@ -106,6 +108,31 @@ func (a *Allocator) Allocate(request AllocationRequest, now time.Time) (Allocati
|
||||
return allocation, nil
|
||||
}
|
||||
|
||||
// PublishAssignment is the allocation-to-client boundary. It holds the same
|
||||
// allocator lock as the claim and exposes no assignment until the allocated
|
||||
// server, complete compatibility tuple, endpoint, and manifest signature all
|
||||
// verify. The returned assignment is stable across an identical retry.
|
||||
func (a *Allocator) PublishAssignment(allocationID string, manifest AllocationManifest, endpoint string, signature []byte, verify func([]byte, []byte) bool) (Assignment, error) {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
allocation, ok := a.allocations[allocationID]
|
||||
if !ok {
|
||||
return Assignment{}, ErrAllocationNotFound
|
||||
}
|
||||
assignment, err := VerifyAssignment(allocation, manifest, endpoint, signature, verify)
|
||||
if err != nil {
|
||||
return Assignment{}, err
|
||||
}
|
||||
if prior, exists := a.assignments[allocationID]; exists {
|
||||
if prior != assignment {
|
||||
return Assignment{}, ErrConflict
|
||||
}
|
||||
return prior, nil
|
||||
}
|
||||
a.assignments[allocationID] = assignment
|
||||
return assignment, nil
|
||||
}
|
||||
|
||||
func validateAllocationRequest(request AllocationRequest) error {
|
||||
if request.AllocationID == "" || request.MatchID == "" || request.Region == "" || request.Build == "" || request.Protocol <= 0 || (request.Transport != "enet" && request.Transport != "steam_sdr") {
|
||||
return ErrAllocationInput
|
||||
|
||||
@@ -83,3 +83,38 @@ func TestAllocatorConcurrentClaimsCannotDoubleAllocateOneServer(t *testing.T) {
|
||||
t.Fatalf("concurrent claims succeeded %d times", wins)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAllocatorPublishesOnlyVerifiedAssignmentAndReplaysIdentically(t *testing.T) {
|
||||
a, err := NewAllocator([]ReadyServer{{ServerID: "server-1", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet", State: ServerReady}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
allocation, err := a.Allocate(AllocationRequest{AllocationID: "allocation-1", MatchID: "match-1", Region: "EU", Build: "build-1", Protocol: 1, Transport: "enet"}, time.Unix(1000, 0))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
manifest := AllocationManifest{AllocationID: allocation.AllocationID, MatchID: allocation.MatchID, ServerID: allocation.ServerID, Region: allocation.Region, Build: allocation.Build, Protocol: allocation.Protocol, Transport: allocation.Transport, RosterDigest: "roster-1"}
|
||||
digest := ManifestDigest(manifest)
|
||||
verify := func(payload, signature []byte) bool {
|
||||
return string(payload) == string(manifestBytes(manifest)) && string(signature) == string(digest[:])
|
||||
}
|
||||
if _, err := a.PublishAssignment("unknown", manifest, "127.0.0.1:30001", digest[:], verify); !errors.Is(err, ErrAllocationNotFound) {
|
||||
t.Fatalf("unknown allocation error = %v", err)
|
||||
}
|
||||
bad := manifest
|
||||
bad.Build = "build-2"
|
||||
if _, err := a.PublishAssignment(allocation.AllocationID, bad, "127.0.0.1:30001", digest[:], verify); !errors.Is(err, ErrManifestRejected) {
|
||||
t.Fatalf("tampered assignment error = %v", err)
|
||||
}
|
||||
first, err := a.PublishAssignment(allocation.AllocationID, manifest, "127.0.0.1:30001", digest[:], verify)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
replay, err := a.PublishAssignment(allocation.AllocationID, manifest, "127.0.0.1:30001", digest[:], verify)
|
||||
if err != nil || replay != first {
|
||||
t.Fatalf("assignment replay = %+v err=%v", replay, err)
|
||||
}
|
||||
if _, err := a.PublishAssignment(allocation.AllocationID, manifest, "127.0.0.1:30002", digest[:], verify); !errors.Is(err, ErrConflict) {
|
||||
t.Fatalf("endpoint mutation error = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user