mirror of
https://github.com/jcreek/CosmicClash.git
synced 2026-09-10 16:04:04 +00:00
168 lines
5.9 KiB
Go
168 lines
5.9 KiB
Go
// Package allocator coordinates provider allocation with durable control-plane
|
|
// state. It does not expose an endpoint until both boundaries succeed.
|
|
package allocator
|
|
|
|
import (
|
|
"context"
|
|
"time"
|
|
|
|
"github.com/cosmic-clash/cosmic-clash/server/agones"
|
|
"github.com/cosmic-clash/cosmic-clash/server/domain"
|
|
)
|
|
|
|
type Provider interface {
|
|
Allocate(context.Context, domain.AllocationRequest, map[string]string, time.Time) (agones.AllocatedServer, error)
|
|
}
|
|
|
|
type ProviderRecoverer interface {
|
|
RecoverAllocation(context.Context, domain.AllocationRequest, time.Time) (agones.AllocatedServer, bool, error)
|
|
}
|
|
|
|
type Durable interface {
|
|
RecordProviderAllocation(context.Context, domain.Allocation, time.Time) (domain.Allocation, error)
|
|
}
|
|
|
|
type RosterPublisher interface {
|
|
PublishRoster(context.Context, domain.Assignment, []domain.SignedJoinAuthorisation, func([]byte, []byte) bool) error
|
|
}
|
|
|
|
type AllocationBudget interface {
|
|
Allow(region string, now time.Time) error
|
|
}
|
|
|
|
type SharedAllocationQuota interface {
|
|
Consume(context.Context, string, time.Time) error
|
|
}
|
|
|
|
type AllocationMetrics interface {
|
|
ObserveAttempt(string)
|
|
ObserveSuccess(string)
|
|
ObserveFailure(string)
|
|
ObserveDenied(string)
|
|
}
|
|
|
|
type Service struct {
|
|
Provider Provider
|
|
Durable Durable
|
|
Roster RosterPublisher
|
|
Budget AllocationBudget
|
|
Quota SharedAllocationQuota
|
|
Metrics AllocationMetrics
|
|
Now func() time.Time
|
|
}
|
|
|
|
// AllocateAcceptedProposal is the hand-off from proposal consensus to server
|
|
// allocation. Keeping this check beside the provider call prevents a caller
|
|
// from allocating capacity for an OPEN/DECLINED proposal or for a request
|
|
// whose playlist does not match the proposal that produced it.
|
|
func (s Service) AllocateAcceptedProposal(ctx context.Context, proposal domain.Proposal, request domain.AllocationRequest, playlist domain.Playlist, labels map[string]string) (agones.AllocatedServer, error) {
|
|
if proposal.State != domain.Accepted || proposal.Playlist != playlist || len(proposal.Participants) == 0 || (request.Playlist != "" && request.Playlist != playlist) || (proposal.Region != "" && request.Region != proposal.Region) || (proposal.Protocol > 0 && request.Protocol != proposal.Protocol) || request.ArenaPath != proposal.ArenaPath {
|
|
return agones.AllocatedServer{}, domain.ErrAllocationInput
|
|
}
|
|
if proposal.Playlist == domain.Ranked && len(proposal.Participants) != 6 {
|
|
return agones.AllocatedServer{}, domain.ErrAllocationInput
|
|
}
|
|
if proposal.Playlist == domain.Casual && (len(proposal.Participants) < 2 || len(proposal.Participants) > 6) {
|
|
return agones.AllocatedServer{}, domain.ErrAllocationInput
|
|
}
|
|
seen := make(map[string]struct{}, len(proposal.Participants))
|
|
for _, participant := range proposal.Participants {
|
|
if participant.PlayerID == "" || participant.Response != domain.AcceptedResponse {
|
|
return agones.AllocatedServer{}, domain.ErrAllocationInput
|
|
}
|
|
if _, exists := seen[participant.PlayerID]; exists {
|
|
return agones.AllocatedServer{}, domain.ErrAllocationInput
|
|
}
|
|
seen[participant.PlayerID] = struct{}{}
|
|
}
|
|
if request.MatchID == "" {
|
|
return agones.AllocatedServer{}, domain.ErrAllocationInput
|
|
}
|
|
return s.Allocate(ctx, request, labels)
|
|
}
|
|
|
|
func (s Service) PublishRoster(ctx context.Context, assignment domain.Assignment, roster []domain.SignedJoinAuthorisation, verify func([]byte, []byte) bool) error {
|
|
if s.Roster == nil {
|
|
return errNotConfigured
|
|
}
|
|
if assignment.Allocation.State != domain.ServerAllocated || assignment.Endpoint == "" {
|
|
return domain.ErrManifestRejected
|
|
}
|
|
return s.Roster.PublishRoster(ctx, assignment, roster, verify)
|
|
}
|
|
|
|
func (s Service) Allocate(ctx context.Context, request domain.AllocationRequest, labels map[string]string) (agones.AllocatedServer, error) {
|
|
if s.Provider == nil || s.Durable == nil || s.Now == nil {
|
|
return agones.AllocatedServer{}, errNotConfigured
|
|
}
|
|
now := s.Now()
|
|
if s.Metrics != nil {
|
|
s.Metrics.ObserveAttempt(request.Region)
|
|
}
|
|
if s.Budget != nil {
|
|
if err := s.Budget.Allow(request.Region, now); err != nil {
|
|
if s.Metrics != nil {
|
|
s.Metrics.ObserveDenied(request.Region)
|
|
}
|
|
return agones.AllocatedServer{}, err
|
|
}
|
|
}
|
|
if s.Quota != nil {
|
|
if err := s.Quota.Consume(ctx, request.Region, now); err != nil {
|
|
if s.Metrics != nil {
|
|
s.Metrics.ObserveDenied(request.Region)
|
|
}
|
|
return agones.AllocatedServer{}, err
|
|
}
|
|
}
|
|
result, err := s.Provider.Allocate(ctx, request, labels, now)
|
|
if err != nil {
|
|
if s.Metrics != nil {
|
|
s.Metrics.ObserveFailure(request.Region)
|
|
}
|
|
return agones.AllocatedServer{}, err
|
|
}
|
|
if err := validateProviderAllocation(request, result); err != nil {
|
|
if s.Metrics != nil {
|
|
s.Metrics.ObserveFailure(request.Region)
|
|
}
|
|
return agones.AllocatedServer{}, err
|
|
}
|
|
recorded, err := s.Durable.RecordProviderAllocation(ctx, result.Allocation, now)
|
|
if err != nil {
|
|
if s.Metrics != nil {
|
|
s.Metrics.ObserveFailure(request.Region)
|
|
}
|
|
return agones.AllocatedServer{}, err
|
|
}
|
|
result.Allocation = recorded
|
|
if s.Metrics != nil {
|
|
s.Metrics.ObserveSuccess(request.Region)
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func (s Service) RecordProviderAllocation(ctx context.Context, result agones.AllocatedServer, now time.Time) (domain.Allocation, error) {
|
|
if s.Durable == nil || result.Allocation.State != domain.ServerAllocated || result.Endpoint == "" {
|
|
return domain.Allocation{}, domain.ErrAllocationInput
|
|
}
|
|
// Quota is consumed by Allocate before a fresh provider request. This
|
|
// method only reconciles an already-issued provider result after an
|
|
// ambiguous write, so consuming here would charge one allocation twice.
|
|
allocation, err := s.Durable.RecordProviderAllocation(ctx, result.Allocation, now)
|
|
if s.Metrics != nil {
|
|
if err != nil {
|
|
s.Metrics.ObserveFailure(result.Allocation.Region)
|
|
} else {
|
|
s.Metrics.ObserveSuccess(result.Allocation.Region)
|
|
}
|
|
}
|
|
return allocation, err
|
|
}
|
|
|
|
var errNotConfigured = &configurationError{}
|
|
|
|
type configurationError struct{}
|
|
|
|
func (*configurationError) Error() string { return "allocator service is not configured" }
|