service

package
v0.0.0-...-0687d8a Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 15, 2026 License: MIT Imports: 19 Imported by: 0

Documentation

Overview

Package service содержит бизнес-логику оркестратора.

Index

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

func (p *BuildPipeline) DeleteBuild(ctx context.Context, gameID int64, version string) error

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

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

type LogStreamReader struct {
	Stream domain.LogStream
}

LogStreamReader адаптер: domain.LogStream → io.Reader (для SSE).

func (*LogStreamReader) Read

func (r *LogStreamReader) Read(p []byte) (int, error)

Read читает следующую журнальную запись в формате JSON.

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

func (s *NodeService) DeleteNode(ctx context.Context, ownerID string, nodeID int64) error

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) Count

func (s *QueueService) Count(ctx context.Context, gameID int64) (int64, error)

Count возвращает количество игроков в очереди.

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) Leave

func (s *QueueService) Leave(ctx context.Context, gameID int64, playerID string) error

Leave удаляет игрока из очереди.

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 содержит параметры загрузки билда.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL