Files

317 lines
14 KiB
Go

package service
import (
"errors"
"fmt"
"strings"
"sync"
"time"
"browser.local/platform/domain"
"browser.local/platform/repo"
"browser.local/platform/validator"
)
const (
sourceRCONCommandTTL = 5 * time.Minute
sourceRCONTimeoutSeconds = 30
)
type sourceRCONCommandPayload struct {
command string
expiresAt time.Time
}
// sourceRCONCommandBroker deliberately retains the raw command only in memory
// until the current leased Run attempt consumes it. It is not recoverable.
type sourceRCONCommandBroker struct {
mu sync.Mutex
now func() time.Time
payloads map[string]sourceRCONCommandPayload
}
func newSourceRCONCommandBroker(now func() time.Time) *sourceRCONCommandBroker {
return &sourceRCONCommandBroker{now: now, payloads: map[string]sourceRCONCommandPayload{}}
}
func (broker *sourceRCONCommandBroker) Put(jobID string, command string) error {
broker.mu.Lock()
defer broker.mu.Unlock()
broker.pruneLocked()
if _, exists := broker.payloads[jobID]; exists {
return validationError("source RCON command idempotency key is already pending")
}
broker.payloads[jobID] = sourceRCONCommandPayload{command: command, expiresAt: broker.now().Add(sourceRCONCommandTTL)}
return nil
}
func (broker *sourceRCONCommandBroker) Consume(jobID string) (string, error) {
broker.mu.Lock()
defer broker.mu.Unlock()
broker.pruneLocked()
payload, exists := broker.payloads[jobID]
if !exists {
return "", validationError("source RCON command input is unavailable")
}
delete(broker.payloads, jobID)
return payload.command, nil
}
func (broker *sourceRCONCommandBroker) Delete(jobID string) {
broker.mu.Lock()
defer broker.mu.Unlock()
delete(broker.payloads, jobID)
}
func (broker *sourceRCONCommandBroker) pruneLocked() {
stamp := broker.now()
for jobID, payload := range broker.payloads {
if !stamp.Before(payload.expiresAt) {
delete(broker.payloads, jobID)
}
}
}
func (svc *CoreService) DispatchSourceRCONCommandForSession(sessionID string, request domain.SourceRCONCommandRequest) (domain.SourceRCONCommandDispatch, error) {
request = domain.CopySourceRCONCommandRequest(request)
if err := validator.ValidateSourceRCONCommandRequest(request); err != nil {
return domain.SourceRCONCommandDispatch{}, err
}
instance, err := svc.GetServerInstanceForSession(sessionID, request.ServerInstanceID)
if err != nil {
return domain.SourceRCONCommandDispatch{}, err
}
if instance.State != domain.ServerInstanceStateRunning {
return domain.SourceRCONCommandDispatch{}, validationError("SCUM RCON requires a running server")
}
resolution, err := svc.resolveSourceRCONDispatch(instance)
if err != nil {
return domain.SourceRCONCommandDispatch{}, err
}
if existing, err := svc.store.Jobs().GetByIdempotency(instance.RunEndpointID, request.IdempotencyKey); err == nil {
if existing.ServerInstanceID != instance.ID || existing.Capability != domain.JobCapabilityRemoteRunRCONCommand || existing.ExecutionInput.SourceRCON == nil {
return domain.SourceRCONCommandDispatch{}, validationError("idempotencyKey is already used for a different RCON command")
}
return sourceRCONDispatchFromJob(existing), nil
} else if !errors.Is(err, repo.ErrNotFound) {
return domain.SourceRCONCommandDispatch{}, err
}
command := sourceRCONCommandText(request)
jobID := jobIDFromParts("job-source-rcon", instance.ID, request.IdempotencyKey)
if err := svc.sourceRCONCommands.Put(jobID, command); err != nil {
return domain.SourceRCONCommandDispatch{}, err
}
job := domain.Job{
ID: jobID,
ServerInstanceID: instance.ID,
RunEndpointID: instance.RunEndpointID,
Capability: domain.JobCapabilityRemoteRunRCONCommand,
TargetKey: resolution.transport.TargetKey,
InputRef: "input://source-rcon/" + jobID,
IdempotencyKey: request.IdempotencyKey,
Progress: domain.JobProgress{Percent: 0, Message: "SCUM RCON command queued"},
RetryPolicy: domain.JobRetryPolicy{MaxAttempts: 1, InitialBackoffSeconds: 1, MaxBackoffSeconds: 1},
ExecutionInput: domain.JobExecutionInput{
WorkspaceScope: resolution.binding.ProfileKey,
RemoteAdapterKey: resolution.transport.Key,
RemoteAdapterKind: string(domain.RemoteAdapterRCON),
TimeoutSeconds: sourceRCONTimeoutSeconds,
PluginID: resolution.plugin.ID,
SourceRCON: resolution.plan,
},
}
created, err := svc.CreateJob(job)
if err != nil {
svc.sourceRCONCommands.Delete(jobID)
return domain.SourceRCONCommandDispatch{}, err
}
if created.ID != jobID {
svc.sourceRCONCommands.Delete(jobID)
if created.ServerInstanceID != instance.ID || created.Capability != domain.JobCapabilityRemoteRunRCONCommand || created.ExecutionInput.SourceRCON == nil {
return domain.SourceRCONCommandDispatch{}, validationError("idempotencyKey is already used for a different RCON command")
}
}
return sourceRCONDispatchFromJob(created), nil
}
func (svc *CoreService) GetSourceRCONExecutionInput(request domain.SourceRCONExecutionInputRequest) (domain.SourceRCONExecutionInput, error) {
if err := validator.ValidateSourceRCONExecutionInputRequest(request); err != nil {
return domain.SourceRCONExecutionInput{}, err
}
job, err := svc.activeFencedInputJob(request.RunEndpointID, request.SessionToken, request.JobID, request.LeaseToken, request.Attempt)
if err != nil {
return domain.SourceRCONExecutionInput{}, err
}
if job.Capability != domain.JobCapabilityRemoteRunRCONCommand || job.ExecutionInput.SourceRCON == nil {
return domain.SourceRCONExecutionInput{}, validationError("job is not a source RCON command")
}
command, err := svc.sourceRCONCommands.Consume(job.ID)
if err != nil {
return domain.SourceRCONExecutionInput{}, err
}
return domain.CopySourceRCONExecutionInput(domain.SourceRCONExecutionInput{JobID: job.ID, ServerInstanceID: job.ServerInstanceID, RunEndpointID: job.RunEndpointID, Command: command}), nil
}
type sourceRCONDispatchResolution struct {
plugin domain.GamePlugin
binding domain.RuntimeBinding
transport domain.RuntimeTransportProfile
plan *domain.RuntimeSourceRCONPlan
}
func (svc *CoreService) resolveSourceRCONDispatch(instance domain.ServerInstance) (sourceRCONDispatchResolution, error) {
plugin, err := svc.store.GamePlugins().Get(instance.PluginID)
if err != nil {
return sourceRCONDispatchResolution{}, err
}
if plugin.Status != domain.GamePluginStatusInstalled || plugin.Version != instance.PluginVersion || !plugin.Permissions.RemoteAccess || !plugin.RemoteAccess.RCON || !containsString(plugin.RequiredRunCapabilities, domain.JobCapabilityRemoteRunRCONCommand) || !containsString(plugin.RemoteAccess.RunCapabilities, domain.JobCapabilityRemoteRunRCONCommand) {
return sourceRCONDispatchResolution{}, forbiddenError("plugin does not declare SCUM RCON command access")
}
endpoint, err := svc.GetRunEndpoint(instance.RunEndpointID)
if err != nil {
return sourceRCONDispatchResolution{}, err
}
if err := svc.validateRunnableEndpoint(endpoint, domain.JobCapabilityRemoteRunRCONCommand); err != nil {
return sourceRCONDispatchResolution{}, err
}
if !strings.EqualFold(endpoint.Platform, "windows") || !strings.EqualFold(endpoint.Architecture, "amd64") {
return sourceRCONDispatchResolution{}, validationError("unsupported_extension_platform: SCUM Source RCON requires windows/amd64")
}
binding, profile, err := svc.sourceRCONRuntimeProfile(instance, plugin, endpoint)
if err != nil {
return sourceRCONDispatchResolution{}, err
}
transport, err := sourceRCONTransport(plugin.RuntimeProfiles, profile)
if err != nil {
return sourceRCONDispatchResolution{}, err
}
extension, err := sourceRCONExtension(plugin.RuntimeProfiles, profile, endpoint)
if err != nil {
return sourceRCONDispatchResolution{}, err
}
plan := &domain.RuntimeSourceRCONPlan{
Protocol: "source-rcon",
ExtensionKey: extension.Key,
ModKey: extension.ModKey,
ConfigRef: "ue4ss/Mods/" + extension.ModKey + "/config.ini",
DeploymentStateRef: "runtime/ue4ss-dll/" + extension.TargetKey + "/release.json",
Port: extension.RCONPort,
}
return sourceRCONDispatchResolution{plugin: plugin, binding: binding, transport: transport, plan: plan}, nil
}
func (svc *CoreService) sourceRCONRuntimeProfile(instance domain.ServerInstance, plugin domain.GamePlugin, endpoint domain.RunEndpoint) (domain.RuntimeBinding, domain.RuntimeLifecycleProfile, error) {
if binding, err := svc.runtimeBindingForServer(instance.ID); err == nil {
if normalized, normalizeErr := normalizeRuntimeBinding(plugin, binding); normalizeErr == nil {
if profile, exists := runtimeLifecycleProfileForKey(plugin.RuntimeProfiles, normalized.ProfileKey); exists && sourceRCONProfileSupports(plugin.RuntimeProfiles, profile, endpoint) {
return normalized, profile, nil
}
}
} else if !errors.Is(err, repo.ErrNotFound) {
return domain.RuntimeBinding{}, domain.RuntimeLifecycleProfile{}, err
}
if key := strings.TrimSpace(instance.Deployment.ProfileKey); key != "" {
if profile, exists := runtimeLifecycleProfileForKey(plugin.RuntimeProfiles, key); exists && sourceRCONProfileSupports(plugin.RuntimeProfiles, profile, endpoint) {
return transientSourceRCONBinding(instance, plugin, profile), profile, nil
}
}
var selected domain.RuntimeLifecycleProfile
for _, profile := range plugin.RuntimeProfiles.LifecycleProfiles {
if !sourceRCONProfileSupports(plugin.RuntimeProfiles, profile, endpoint) {
continue
}
if selected.Key != "" && selected.Key != profile.Key {
return domain.RuntimeBinding{}, domain.RuntimeLifecycleProfile{}, validationError("selected runtime profile has multiple SCUM RCON candidates")
}
selected = profile
}
if selected.Key == "" {
return domain.RuntimeBinding{}, domain.RuntimeLifecycleProfile{}, validationError("selected runtime profile does not support SCUM RCON")
}
return transientSourceRCONBinding(instance, plugin, selected), selected, nil
}
func sourceRCONProfileSupports(profiles domain.GamePluginRuntimeProfiles, profile domain.RuntimeLifecycleProfile, endpoint domain.RunEndpoint) bool {
if !containsString(profile.Capabilities, domain.JobCapabilityRemoteRunRCONCommand) || !runtimePlatformsContain(profile.Platforms, "windows") {
return false
}
if _, err := sourceRCONTransport(profiles, profile); err != nil {
return false
}
if _, err := sourceRCONExtension(profiles, profile, endpoint); err != nil {
return false
}
return true
}
func transientSourceRCONBinding(instance domain.ServerInstance, plugin domain.GamePlugin, profile domain.RuntimeLifecycleProfile) domain.RuntimeBinding {
return domain.RuntimeBinding{ID: "runtime-binding-" + instance.ID, ServerInstanceID: instance.ID, PluginID: plugin.ID, PluginVersion: plugin.Version, ProfileKey: profile.Key, Mode: profile.Mode, Status: domain.RuntimeBindingStatusComplete}
}
func sourceRCONTransport(profiles domain.GamePluginRuntimeProfiles, profile domain.RuntimeLifecycleProfile) (domain.RuntimeTransportProfile, error) {
return sourceRCONTransportForCapability(profiles, profile, "", domain.JobCapabilityRemoteRunRCONCommand)
}
func sourceRCONTransportForCapability(profiles domain.GamePluginRuntimeProfiles, profile domain.RuntimeLifecycleProfile, requiredKey string, capability string) (domain.RuntimeTransportProfile, error) {
var selected domain.RuntimeTransportProfile
for _, candidate := range profiles.TransportProfiles {
if requiredKey != "" && candidate.Key != requiredKey {
continue
}
if !containsString(profile.TransportKeys, candidate.Key) || candidate.Kind != "rcon" || !containsString(candidate.Capabilities, capability) {
continue
}
if selected.Key != "" {
return domain.RuntimeTransportProfile{}, validationError("selected runtime profile has multiple SCUM RCON transports")
}
selected = candidate
}
if selected.Key == "" || strings.TrimSpace(selected.TargetKey) == "" {
return domain.RuntimeTransportProfile{}, validationError("selected runtime profile has no SCUM RCON transport")
}
return selected, nil
}
func sourceRCONExtension(profiles domain.GamePluginRuntimeProfiles, profile domain.RuntimeLifecycleProfile, endpoint domain.RunEndpoint) (domain.RuntimeDLLExtensionProfile, error) {
byKey := make(map[string]domain.RuntimeDLLExtensionProfile, len(profiles.DLLExtensions))
for _, extension := range profiles.DLLExtensions {
byKey[extension.Key] = extension
}
var selected domain.RuntimeDLLExtensionProfile
for _, key := range profile.DLLExtensionRefs {
extension, exists := byKey[key]
if !exists || extension.Kind != "ue4ss-dll" || extension.ModKey != "scum_simple_rcon" || extension.ReleaseState != "ready" {
continue
}
if !runtimeDLLExtensionSupportsTarget(extension, endpoint.Platform, endpoint.Architecture) {
return domain.RuntimeDLLExtensionProfile{}, validationError("unsupported_extension_platform: SCUM Source RCON requires windows/amd64")
}
if selected.Key != "" {
return domain.RuntimeDLLExtensionProfile{}, validationError("selected runtime profile has multiple SCUM Source RCON extensions")
}
selected = extension
}
if selected.Key == "" {
return domain.RuntimeDLLExtensionProfile{}, validationError("extension_release_unavailable: ready SCUM Source RCON DLL is not selected")
}
return selected, nil
}
func sourceRCONCommandText(request domain.SourceRCONCommandRequest) string {
if request.Kind == domain.SourceRCONCommandKindCommand {
return strings.TrimSpace(request.Command)
}
message := strings.NewReplacer("\\", "\\\\", "\"", "\\\"").Replace(request.Message)
command := fmt.Sprintf("SendChat %d \"%s\"", request.ChatType, message)
if request.TargetSteamID != "" {
command += " " + request.TargetSteamID
}
return command
}
func sourceRCONDispatchFromJob(job domain.Job) domain.SourceRCONCommandDispatch {
return domain.CopySourceRCONCommandDispatch(domain.SourceRCONCommandDispatch{JobID: job.ID, ServerInstanceID: job.ServerInstanceID, Status: string(job.State), Message: "SCUM RCON command queued"})
}