mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-10 16:04:04 +00:00
146 lines
5.8 KiB
Go
146 lines
5.8 KiB
Go
package api
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"fmt"
|
|
"time"
|
|
|
|
"github.com/cosmic-clash/cosmic-clash/server/domain"
|
|
"github.com/cosmic-clash/cosmic-clash/server/store"
|
|
"github.com/cosmic-clash/cosmic-clash/server/workload"
|
|
)
|
|
|
|
// AssignmentProviderFromStore adapts the durable player-scoped assignment
|
|
// projection to the HTTP boundary. The store query filters expiry and binds
|
|
// both match and player; the API still performs its response-shape checks.
|
|
func AssignmentProviderFromStore(db *sql.DB) AssignmentProvider {
|
|
return func(ctx context.Context, playerID, matchID string, now time.Time) (AssignmentView, error) {
|
|
assignment, err := store.GetAssignment(ctx, db, playerID, matchID, now)
|
|
if err != nil {
|
|
return AssignmentView{}, err
|
|
}
|
|
return AssignmentView{
|
|
MatchID: assignment.MatchID,
|
|
ServerID: assignment.ServerID,
|
|
PlayerID: assignment.PlayerID,
|
|
Slot: assignment.Slot,
|
|
ExpiresAt: assignment.ExpiresAt,
|
|
ProtocolVersion: assignment.ProtocolVersion,
|
|
Transport: assignment.Transport,
|
|
JoinAuthorisation: assignment.JoinAuthorisation,
|
|
Endpoint: assignment.Endpoint,
|
|
Revision: assignment.Revision,
|
|
}, nil
|
|
}
|
|
}
|
|
|
|
type postgresProposalBackend struct{ db *sql.DB }
|
|
|
|
func (p postgresProposalBackend) Get(ctx context.Context, playerID, proposalID string, now time.Time) (domain.Proposal, error) {
|
|
return store.GetProposal(ctx, p.db, playerID, proposalID, now)
|
|
}
|
|
|
|
func (p postgresProposalBackend) Respond(ctx context.Context, playerID, proposalID, idempotencyKey string, accept bool, expectedRevision uint64, now time.Time) (domain.Proposal, error) {
|
|
return store.RespondToProposal(ctx, p.db, playerID, proposalID, idempotencyKey, accept, expectedRevision, now)
|
|
}
|
|
|
|
func ProposalProviderFromStore(db *sql.DB) ProposalBackend {
|
|
return postgresProposalBackend{db: db}
|
|
}
|
|
|
|
// ProposalPromoterFromStore turns a durably accepted proposal into its exact
|
|
// matcher-selected ALLOCATING match. The store chooses a deterministic match
|
|
// ID so an API retry after a transient failure cannot duplicate the match.
|
|
func ProposalPromoterFromStore(db *sql.DB) ProposalPromoter {
|
|
return ProposalPromoterFunc(func(ctx context.Context, proposal domain.Proposal, now time.Time) error {
|
|
if proposal.State != domain.Accepted {
|
|
return domain.ErrIllegalTransition
|
|
}
|
|
return store.PromoteStoredAcceptedProposal(ctx, db, proposal.ProposalID, now)
|
|
})
|
|
}
|
|
|
|
type postgresServerRegistrar struct{ db *sql.DB }
|
|
|
|
func (p postgresServerRegistrar) RegisterServer(ctx context.Context, binding domain.WorkloadBinding, protocol int, assignmentReady bool, idempotencyKey string, now time.Time) error {
|
|
return store.AdvanceServerRegistration(ctx, p.db, binding, protocol, assignmentReady, idempotencyKey, now)
|
|
}
|
|
|
|
func ServerRegistrarFromStore(db *sql.DB) ServerRegistrar {
|
|
if db == nil {
|
|
return nil
|
|
}
|
|
return postgresServerRegistrar{db: db}
|
|
}
|
|
|
|
type postgresServerShutdowner struct{ db *sql.DB }
|
|
|
|
func (p postgresServerShutdowner) ShutdownServer(ctx context.Context, binding domain.WorkloadBinding, reason, idempotencyKey string, now time.Time) error {
|
|
return store.RecordServerShutdown(ctx, p.db, binding, reason, idempotencyKey, now)
|
|
}
|
|
|
|
func ServerShutdownerFromStore(db *sql.DB) ServerShutdowner {
|
|
if db == nil {
|
|
return nil
|
|
}
|
|
return postgresServerShutdowner{db: db}
|
|
}
|
|
|
|
type postgresServerConnections struct{ db *sql.DB }
|
|
|
|
func (p postgresServerConnections) ClaimPlayerConnection(ctx context.Context, binding domain.WorkloadBinding, playerID string, expectedGeneration uint64, idempotencyKey string, now time.Time) (uint64, error) {
|
|
return store.ClaimPlayerConnection(ctx, p.db, binding, playerID, expectedGeneration, idempotencyKey, now)
|
|
}
|
|
|
|
func (p postgresServerConnections) RecordPlayerDisconnected(ctx context.Context, binding domain.WorkloadBinding, playerID string, generation uint64, idempotencyKey string, now time.Time) error {
|
|
return store.RecordPlayerDisconnected(ctx, p.db, binding, playerID, generation, idempotencyKey, now)
|
|
}
|
|
|
|
func ServerConnectionsFromStore(db *sql.DB) ServerConnectionRecorder {
|
|
if db == nil {
|
|
return nil
|
|
}
|
|
return postgresServerConnections{db: db}
|
|
}
|
|
|
|
// WorkloadVerifierFromSignedToken builds WorkloadVerify from a control-plane
|
|
// -owned signed token instead of a Kubernetes-projected JWT (see
|
|
// workload/signed_token.go for why: it needs no live cluster to verify).
|
|
// secret must be kept out of source control (env var in cmd/control-plane);
|
|
// an empty secret returns nil so a misconfigured deployment fails the same
|
|
// way an unwired verifier already does today (503, not a silent bypass).
|
|
func WorkloadVerifierFromSignedToken(secret []byte, db *sql.DB) WorkloadVerifier {
|
|
if len(secret) == 0 || db == nil {
|
|
return nil
|
|
}
|
|
return func(token string, now time.Time) (domain.WorkloadBinding, error) {
|
|
claims, err := workload.ParseSignedWorkloadToken(secret, token, now)
|
|
if err != nil {
|
|
return domain.WorkloadBinding{}, err
|
|
}
|
|
// WorkloadVerifier has no context parameter (see its type in
|
|
// service.go) so the durable lookup below cannot inherit the
|
|
// caller's request context; bound it locally instead of running
|
|
// unbounded against context.Background().
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
// The token only names allocation_id (see signed_token.go for why);
|
|
// match_id/server_id come from the durable allocator record, never
|
|
// from the caller, so a token can never claim a pairing that wasn't
|
|
// actually, durably allocated.
|
|
matchID, serverID, ok, err := store.AllocationBindingByAllocationID(ctx, db, claims.AllocationID)
|
|
if err != nil {
|
|
return domain.WorkloadBinding{}, err
|
|
}
|
|
if !ok {
|
|
return domain.WorkloadBinding{}, fmt.Errorf("signed workload token names an allocation that is no longer valid")
|
|
}
|
|
return domain.WorkloadBinding{
|
|
AllocationID: claims.AllocationID,
|
|
MatchID: matchID,
|
|
ServerID: serverID,
|
|
}, nil
|
|
}
|
|
}
|