Documentation
¶
Overview ¶
Package service содержит бизнес-логику оркестратора.
Index ¶
- type AnnounceNodeParams
- type AnnounceNodeResult
- type BuildPipeline
- func (p *BuildPipeline) DeleteBuild(ctx context.Context, gameID int64, version string) error
- func (p *BuildPipeline) GetBuild(ctx context.Context, gameID int64, version string) (*domain.ServerBuild, error)
- func (p *BuildPipeline) ListBuilds(ctx context.Context, gameID int64) ([]*domain.ServerBuild, error)
- func (p *BuildPipeline) UploadBuild(ctx context.Context, params UploadBuildParams) (*domain.ServerBuild, error)
- func (p *BuildPipeline) WithWorkDir(dir string) *BuildPipeline
- type DiscoveryService
- type EnrichedInstance
- type EnrichedNode
- type GamePolicyService
- type HeartbeatService
- type InstanceService
- func (s *InstanceService) DeleteInstance(ctx context.Context, ownerID string, gameID, instanceID int64) error
- func (s *InstanceService) GetInstance(ctx context.Context, ownerID string, gameID, instanceID int64) (*EnrichedInstance, error)
- func (s *InstanceService) GetInstanceUsage(ctx context.Context, ownerID string, gameID, instanceID int64) (*domain.ResourceUsage, error)
- func (s *InstanceService) ListInstances(ctx context.Context, ownerID string, gameID int64, ...) ([]*EnrichedInstance, error)
- func (s *InstanceService) RestartInstance(ctx context.Context, ownerID string, gameID, instanceID int64) (*domain.Instance, error)
- func (s *InstanceService) ResumeInstance(ctx context.Context, ownerID string, gameID, instanceID int64) (*domain.Instance, error)
- func (s *InstanceService) StartInstance(ctx context.Context, params StartInstanceParams) (*domain.Instance, error)
- func (s *InstanceService) StopInstance(ctx context.Context, ownerID string, gameID, instanceID int64, ...) (*domain.Instance, error)
- func (s *InstanceService) StreamInstanceLogs(ctx context.Context, ownerID string, gameID, instanceID int64, ...) (domain.LogStream, error)
- type LogStreamReader
- type NodeService
- func (s *NodeService) AnnounceNode(ctx context.Context, params AnnounceNodeParams) (*AnnounceNodeResult, error)
- func (s *NodeService) DeleteNode(ctx context.Context, ownerID string, nodeID int64) error
- func (s *NodeService) GetNode(ctx context.Context, ownerID string, nodeID int64) (*EnrichedNode, error)
- func (s *NodeService) GetNodeUsage(ctx context.Context, ownerID string, nodeID int64) (*NodeUsageResult, error)
- func (s *NodeService) ListNodeInstances(ctx context.Context, ownerID string, nodeID int64) ([]*EnrichedInstance, error)
- func (s *NodeService) ListNodes(ctx context.Context, ownerID string, status *domain.NodeStatus) ([]*EnrichedNode, error)
- func (s *NodeService) RegisterNode(ctx context.Context, params RegisterNodeParams) (*domain.Node, error)
- func (s *NodeService) SyncInstances(ctx context.Context, nodeID int64, activeContainerIDs []string) error
- type NodeUsageResult
- type QueueService
- func (s *QueueService) CleanupExpired(ctx context.Context, gameID int64, policy *domain.GamePolicy) ([]string, error)
- func (s *QueueService) Count(ctx context.Context, gameID int64) (int64, error)
- func (s *QueueService) Heartbeat(ctx context.Context, gameID int64, playerID string) (*QueueStatusResult, error)
- func (s *QueueService) Join(ctx context.Context, gameID int64, playerID, mode string) (*QueueStatusResult, error)
- func (s *QueueService) Leave(ctx context.Context, gameID int64, playerID string) error
- func (s *QueueService) ProcessQueue(ctx context.Context, gameID int64) error
- func (s *QueueService) Status(ctx context.Context, gameID int64, playerID string) (*QueueStatusResult, error)
- type QueueStatusResult
- type RegisterNodeParams
- type StartInstanceParams
- type UploadBuildParams
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type AnnounceNodeParams ¶
type AnnounceNodeParams struct {
Address string
Region string
AgentVersion string
CPUCores uint32
TotalMemoryBytes uint64
TotalDiskBytes uint64
APIKey string // NODE_API_KEY ноды
ActiveContainerIDs []string
}
AnnounceNodeParams содержит параметры анонсирования ноды.
type AnnounceNodeResult ¶
type AnnounceNodeResult struct {
NodeID int64
}
AnnounceNodeResult содержит результат анонсирования ноды.
type BuildPipeline ¶
type BuildPipeline struct {
// contains filtered or unexported fields
}
BuildPipeline определяет пайплайн преобразования архива билда в Docker-образ, загрузки его на ноду и регистрации метаданных.
func NewBuildPipeline ¶
func NewBuildPipeline( buildRepo domain.BuildStorage, buildFS domain.BuildStorageFS, nodeClient domain.NodeClient, nodeRepo domain.NodeRepo, nodeState domain.NodeStateStore, limits config.LimitsConfig, ) *BuildPipeline
NewBuildPipeline создаёт пайплайн обработки билдов.
func (*BuildPipeline) DeleteBuild ¶
DeleteBuild удаляет билд по версии. Проверяет отсутствие активных инстансов, удаляет файл из файлового хранилища и метаданные из PostgreSQL. Возвращает ErrBuildInUse при наличии запущенных инстансов, ErrNotFound при отсутствии билда.
func (*BuildPipeline) GetBuild ¶
func (p *BuildPipeline) GetBuild(ctx context.Context, gameID int64, version string) (*domain.ServerBuild, error)
GetBuild возвращает билд по версии.
func (*BuildPipeline) ListBuilds ¶
func (p *BuildPipeline) ListBuilds(ctx context.Context, gameID int64) ([]*domain.ServerBuild, error)
ListBuilds возвращает список билдов игры (последние limit штук).
func (*BuildPipeline) UploadBuild ¶
func (p *BuildPipeline) UploadBuild(ctx context.Context, params UploadBuildParams) (*domain.ServerBuild, error)
UploadBuild загружает серверный билд: распаковывает архив для валидации, отправляет архив на ноду для сборки Docker-образа, сохраняет архив в хранилище и регистрирует метаданные. Возвращает ErrNoAvailableNode при отсутствии свободных нод. Превышение размера архива или лимита билдов на игру возвращает ошибку.
func (*BuildPipeline) WithWorkDir ¶
func (p *BuildPipeline) WithWorkDir(dir string) *BuildPipeline
WithWorkDir задаёт директорию для временных файлов.
type DiscoveryService ¶
type DiscoveryService struct {
// contains filtered or unexported fields
}
DiscoveryService предоставляет данные для подключения к игровым серверам.
func NewDiscoveryService ¶
func NewDiscoveryService( instanceRepo domain.InstanceRepo, instanceState domain.InstanceStateStore, nodeRepo domain.NodeRepo, buildRepo domain.BuildStorage, policyService *GamePolicyService, instanceSvc instanceStarter, queueSvc queueProcessor, ) *DiscoveryService
NewDiscoveryService создаёт сервис обнаружения серверов.
func (*DiscoveryService) DiscoverServers ¶
func (s *DiscoveryService) DiscoverServers(ctx context.Context, gameID int64, playerID string) (*domain.DiscoveryResult, error)
DiscoverServers возвращает список доступных серверов для подключения вместе со статусом, который помогает клиенту понять текущую ситуацию.
Сценарии:
- READY — есть running-инстансы со свободными слотами.
- STARTING — нет свободных running, но инстансы стартуют (или запущен асинхронный auto-start).
- CAPACITY_REACHED — все инстансы заполнены, новых запустить нельзя (лимит или queue-режим).
- UNAVAILABLE — авто-оркестрация выключена или не настроена.
- QUEUE — все серверы заняты, scale_behavior=queue, нужно встать в очередь.
- RESERVED — для этого player_id есть активная резервация слота.
type EnrichedInstance ¶
type EnrichedInstance struct {
*domain.Instance
Status domain.InstanceStatus // из KV
PlayerCount *uint32 // из KV
}
EnrichedInstance — инстанс с данными из KV.
type EnrichedNode ¶
type EnrichedNode struct {
*domain.Node
Usage *domain.ResourceUsage
ActiveInstanceCount *uint32
}
EnrichedNode — нода с данными из KV.
type GamePolicyService ¶
type GamePolicyService struct {
// contains filtered or unexported fields
}
GamePolicyService управляет политиками оркестрации серверов по проектам.
func NewGamePolicyService ¶
func NewGamePolicyService(policyRepo domain.GamePolicyRepo) *GamePolicyService
NewGamePolicyService создаёт сервис политик.
func (*GamePolicyService) Get ¶
func (s *GamePolicyService) Get(ctx context.Context, gameID int64) (*domain.GamePolicy, error)
Get возвращает политику игры. Если не найдена — возвращает политику по умолчанию (disabled).
func (*GamePolicyService) ListAll ¶
func (s *GamePolicyService) ListAll(ctx context.Context) ([]*domain.GamePolicy, error)
ListAll возвращает все сохранённые политики.
func (*GamePolicyService) Set ¶
func (s *GamePolicyService) Set(ctx context.Context, policy *domain.GamePolicy) (*domain.GamePolicy, error)
Set создаёт или обновляет политику игры.
type HeartbeatService ¶
type HeartbeatService struct {
// contains filtered or unexported fields
}
HeartbeatService — фоновый сервис мониторинга жизнеспособности нод. Запускается как горутина, периодически опрашивает ноды через gRPC Heartbeat. Также применяет политики оркестрации: авто-рестарт crashed-инстансов.
func NewHeartbeatService ¶
func NewHeartbeatService( nodeRepo domain.NodeRepo, nodeState domain.NodeStateStore, instanceRepo domain.InstanceRepo, instanceState domain.InstanceStateStore, nodeClient domain.NodeClient, buildRepo domain.BuildStorage, policyService *GamePolicyService, instanceSvc instanceOrchestrator, queueSvc *QueueService, hb config.NodeHeartbeatCfg, log *slog.Logger, ) *HeartbeatService
NewHeartbeatService создаёт сервис мониторинга нод.
func (*HeartbeatService) EnforcePolicies ¶
func (s *HeartbeatService) EnforcePolicies(ctx context.Context)
EnforcePolicies применяет политики оркестрации. Для игр в режиме keep_alive проверяет наличие target_instances и запускает недостающие. Также обрабатывает scale_to_zero для поддержки target_instances. Может вызываться как из heartbeat-цикла, так и при старте сервера.
func (*HeartbeatService) RestoreInstanceStatuses ¶
func (s *HeartbeatService) RestoreInstanceStatuses(ctx context.Context) error
RestoreInstanceStatuses восстанавливает статусы инстансов из БД в KV при старте. Вызывается один раз при инициализации.
func (*HeartbeatService) Run ¶
func (s *HeartbeatService) Run(ctx context.Context)
Run запускает цикл мониторинга. Блокирует вызывающую горутину до отмены контекста.
type InstanceService ¶
type InstanceService struct {
// contains filtered or unexported fields
}
InstanceService управляет жизненным циклом экземпляров игровых серверов.
func NewInstanceService ¶
func NewInstanceService( instanceRepo domain.InstanceRepo, instanceState domain.InstanceStateStore, buildRepo domain.BuildStorage, nodeRepo domain.NodeRepo, nodeState domain.NodeStateStore, nodeClient domain.NodeClient, limits config.LimitsConfig, ) *InstanceService
NewInstanceService создаёт сервис управления инстансами.
func (*InstanceService) DeleteInstance ¶
func (s *InstanceService) DeleteInstance(ctx context.Context, ownerID string, gameID, instanceID int64) error
DeleteInstance удаляет экземпляр игрового сервера: останавливает контейнер на ноде, удаляет запись из PostgreSQL и состояние из KV. Уменьшает счётчик активных инстансов на ноде.
func (*InstanceService) GetInstance ¶
func (s *InstanceService) GetInstance(ctx context.Context, ownerID string, gameID, instanceID int64) (*EnrichedInstance, error)
GetInstance возвращает инстанс с обогащением из KV.
func (*InstanceService) GetInstanceUsage ¶
func (s *InstanceService) GetInstanceUsage(ctx context.Context, ownerID string, gameID, instanceID int64) (*domain.ResourceUsage, error)
GetInstanceUsage возвращает метрики потребления ресурсов инстанса. Сначала пробует KV (быстро), fallback — gRPC на ноду.
func (*InstanceService) ListInstances ¶
func (s *InstanceService) ListInstances(ctx context.Context, ownerID string, gameID int64, status *domain.InstanceStatus) ([]*EnrichedInstance, error)
ListInstances возвращает список инстансов пользователя с обогащением из KV.
func (*InstanceService) RestartInstance ¶
func (s *InstanceService) RestartInstance(ctx context.Context, ownerID string, gameID, instanceID int64) (*domain.Instance, error)
RestartInstance перезапускает инстанс через docker restart на ноде.
func (*InstanceService) ResumeInstance ¶
func (s *InstanceService) ResumeInstance(ctx context.Context, ownerID string, gameID, instanceID int64) (*domain.Instance, error)
ResumeInstance запускает остановленный инстанс (docker start на ноде).
func (*InstanceService) StartInstance ¶
func (s *InstanceService) StartInstance(ctx context.Context, params StartInstanceParams) (*domain.Instance, error)
StartInstance запускает новый экземпляр игрового сервера на доступной ноде. Проверяет лимит инстансов на игру, выбирает ноду с наименьшей загрузкой, запускает контейнер через gRPC и регистрирует метаданные. При ошибке сохраняет целостность: откатывает запуск на ноде, если не удалось сохранить в PG. Возвращает ErrNoAvailableNode при отсутствии свободных нод, ErrNotFound при отсутствии билда.
func (*InstanceService) StopInstance ¶
func (s *InstanceService) StopInstance(ctx context.Context, ownerID string, gameID, instanceID int64, timeoutSec uint32) (*domain.Instance, error)
StopInstance останавливает экземпляр игрового сервера, обновляет статус в PostgreSQL и удаляет состояние из KV. Уменьшает счётчик активных инстансов на ноде. Возвращает ошибку, если instanceID не принадлежит указанной игре или владельцу.
func (*InstanceService) StreamInstanceLogs ¶
func (s *InstanceService) StreamInstanceLogs(ctx context.Context, ownerID string, gameID, instanceID int64, req domain.StreamLogsRequest) (domain.LogStream, error)
StreamInstanceLogs возвращает поток журналов инстанса. Caller обязан закрыть stream после использования. Возвращает ошибку, если инстанс не найден или не запущен на ноде.
type LogStreamReader ¶
LogStreamReader адаптер: domain.LogStream → io.Reader (для SSE).
type NodeService ¶
type NodeService struct {
// contains filtered or unexported fields
}
NodeService управляет вычислительными нодами.
func NewNodeService ¶
func NewNodeService( log *slog.Logger, nodeRepo domain.NodeRepo, nodeState domain.NodeStateStore, instanceRepo domain.InstanceRepo, instanceState domain.InstanceStateStore, nodeClient domain.NodeClient, ) *NodeService
NewNodeService создаёт сервис управления нодами.
func (*NodeService) AnnounceNode ¶
func (s *NodeService) AnnounceNode(ctx context.Context, params AnnounceNodeParams) (*AnnounceNodeResult, error)
AnnounceNode обрабатывает анонсирование ноды от самой ноды. Нода передаёт свой NODE_API_KEY как api_key — этот ключ становится токеном авторизации. Пользователь вводит тот же NODE_API_KEY для подключения.
func (*NodeService) DeleteNode ¶
DeleteNode удаляет ноду из оркестратора. Проверяет владение.
func (*NodeService) GetNode ¶
func (s *NodeService) GetNode(ctx context.Context, ownerID string, nodeID int64) (*EnrichedNode, error)
GetNode возвращает ноду с обогащением из KV. Проверяет владение.
func (*NodeService) GetNodeUsage ¶
func (s *NodeService) GetNodeUsage(ctx context.Context, ownerID string, nodeID int64) (*NodeUsageResult, error)
GetNodeUsage возвращает метрики ноды. Проверяет владение.
func (*NodeService) ListNodeInstances ¶
func (s *NodeService) ListNodeInstances(ctx context.Context, ownerID string, nodeID int64) ([]*EnrichedInstance, error)
ListNodeInstances возвращает список инстансов на указанной ноде.
func (*NodeService) ListNodes ¶
func (s *NodeService) ListNodes(ctx context.Context, ownerID string, status *domain.NodeStatus) ([]*EnrichedNode, error)
ListNodes возвращает ноды пользователя с обогащением из KV.
func (*NodeService) RegisterNode ¶
func (s *NodeService) RegisterNode(ctx context.Context, params RegisterNodeParams) (*domain.Node, error)
RegisterNode подключает ноду к оркестратору. Токен — это NODE_API_KEY ноды, единый для обоих методов (manual и authorize).
func (*NodeService) SyncInstances ¶
func (s *NodeService) SyncInstances(ctx context.Context, nodeID int64, activeContainerIDs []string) error
SyncInstances синхронизирует статусы инстансов с активными контейнерами на ноде. Примечание: в доменной модели Instance отсутствует поле ContainerID, поэтому сопоставление выполняется по строковому представлению Instance.ID.
type NodeUsageResult ¶
type NodeUsageResult struct {
NodeID int64
Usage *domain.ResourceUsage
ActiveInstanceCount uint32
}
NodeUsageResult — результат запроса метрик ноды.
type QueueService ¶
type QueueService struct {
// contains filtered or unexported fields
}
QueueService управляет очередью игроков.
func NewQueueService ¶
func NewQueueService( store domain.QueueStore, eventRepo domain.QueueEventRepo, policySvc *GamePolicyService, instanceRepo domain.InstanceRepo, instanceState domain.InstanceStateStore, nodeRepo domain.NodeRepo, log *slog.Logger, ) *QueueService
NewQueueService создаёт сервис очереди.
func (*QueueService) CleanupExpired ¶
func (s *QueueService) CleanupExpired(ctx context.Context, gameID int64, policy *domain.GamePolicy) ([]string, error)
CleanupExpired удаляет игроков с просроченным heartbeat.
func (*QueueService) Heartbeat ¶
func (s *QueueService) Heartbeat(ctx context.Context, gameID int64, playerID string) (*QueueStatusResult, error)
Heartbeat обновляет heartbeat и возвращает текущий статус.
func (*QueueService) Join ¶
func (s *QueueService) Join(ctx context.Context, gameID int64, playerID, mode string) (*QueueStatusResult, error)
Join добавляет игрока в очередь.
func (*QueueService) ProcessQueue ¶
func (s *QueueService) ProcessQueue(ctx context.Context, gameID int64) error
ProcessQueue вызывается при освобождении слота. Резервирует слот для первого игрока в очереди.
func (*QueueService) Status ¶
func (s *QueueService) Status(ctx context.Context, gameID int64, playerID string) (*QueueStatusResult, error)
Status возвращает статус без обновления heartbeat (read-only).
type QueueStatusResult ¶
type QueueStatusResult struct {
Status domain.QueueStatus
Position int32
TotalInQueue int32
EstimatedWaitSeconds int32
ReservedEndpoint *domain.ServerEndpoint
ReservedUntil time.Time
}
QueueStatusResult — результат операции с очередью.
type RegisterNodeParams ¶
type RegisterNodeParams struct {
OwnerID string
Address string
Token string
Region string
NodeID *int64
}
RegisterNodeParams содержит параметры подключения ноды.
type StartInstanceParams ¶
type StartInstanceParams struct {
OwnerID string
GameID int64
BuildVersion string
Name string
PortAllocation domain.PortAllocation
ResourceLimits *domain.ResourceLimits
EnvVars map[string]string
Args []string
DeveloperPayload map[string]string
MaxPlayers *uint32
NodePreference string // "auto" или "node-<id>"
}
StartInstanceParams содержит параметры запуска инстанса.
type UploadBuildParams ¶
type UploadBuildParams struct {
OwnerID string
GameID int64
Version string
Protocol domain.Protocol
InternalPort uint32
MaxPlayers uint32
Archive io.Reader
ArchiveSize int64
// ArchiveData — альтернатива Archive, когда данные уже в памяти (unary gRPC).
ArchiveData []byte
}
UploadBuildParams содержит параметры загрузки билда.