Files

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
}
}