package service import ( "encoding/base64" "encoding/json" "errors" "strings" "sync" "testing" "time" "browser.local/platform/domain" "browser.local/platform/repo" ) type captureDistributionBuilder struct { inputs chan domain.DistributionBuildInput release <-chan struct{} payload []byte } func (builder captureDistributionBuilder) Readiness() (bool, string) { return true, "" } func (builder captureDistributionBuilder) Build(input domain.DistributionBuildInput) ([]byte, error) { builder.inputs <- input if builder.release != nil { <-builder.release } return domain.CopyBytes(builder.payload), nil } func TestCoreServiceKeepsPlatformBuildKeyOffMachineJobChannel(t *testing.T) { svc, session, instance := newDistributionTestFixture(t) plugin, err := svc.store.GamePlugins().Get(instance.PluginID) if err != nil { t.Fatalf("get plugin fixture: %v", err) } plugin.LifecycleAssets = []domain.PluginAssetFile{ {Path: "actions/install.json", Content: `{"version":1,"action":"install","mode":"oneshot"}`, Mode: 0o600}, {Path: "bin/install-server", Content: "#!/usr/bin/env sh\n", Mode: 0o700}, } plugin.RuntimeProfiles.LogSources = append(plugin.RuntimeProfiles.LogSources, domain.RuntimeLogSource{Key: "console", Kind: "process.stdout", TargetKey: "server/process", StreamKey: "console", CursorKind: "sequence", RetentionDays: 14}, domain.RuntimeLogSource{Key: "server-events", Kind: "file.tail", TargetKey: "logs/server", StreamKey: "scum.server", CursorKind: "fingerprint", RetentionDays: 90}, ) plugin.RuntimeProfiles.TransportProfiles = append(plugin.RuntimeProfiles.TransportProfiles, domain.RuntimeTransportProfile{Key: "server-files", Kind: "file", TargetKey: "server-root", Capabilities: []string{domain.JobCapabilityRemoteRunFilesRead}}, domain.RuntimeTransportProfile{Key: "world-db", Kind: "sqlite", TargetKey: "world-db", Capabilities: []string{domain.JobCapabilityRemoteRunDBSQLiteProbe}}, ) plugin.RuntimeProfiles.DataTargets = append(plugin.RuntimeProfiles.DataTargets, domain.RuntimeDataTarget{Key: "world-db", Kind: "sqlite.snapshot", TransportKey: "world-db", SourceRootKey: "server-root", SourcePath: "world/current.db", WorkspaceKey: "databases/world-db", RefreshPolicy: "on-demand-snapshot", MaxBytes: 128 * 1024 * 1024, Platforms: []string{"linux"}}, ) if err := svc.store.GamePlugins().Update(plugin); err != nil { t.Fatalf("seed plugin lifecycle assets: %v", err) } inputs := make(chan domain.DistributionBuildInput, 1) release := make(chan struct{}) var releaseOnce sync.Once t.Cleanup(func() { releaseOnce.Do(func() { close(release) }) }) svc.ConfigureDistributionBuilder(captureDistributionBuilder{inputs: inputs, release: release, payload: []byte("captured-platform-build")}) distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "platform-secret-boundary", }) if err != nil { t.Fatalf("generate platform distribution: %v", err) } var platformInput domain.DistributionBuildInput select { case platformInput = <-inputs: case <-time.After(time.Second): t.Fatal("platform builder did not receive internal build input") } if platformInput.AuthKey == "" || platformInput.JobID != distribution.BuildJobID || platformInput.RunEndpointID != instance.RunEndpointID { t.Fatalf("platform builder received incomplete internal input: %+v", platformInput) } decodedSeed, err := base64.StdEncoding.DecodeString(platformInput.WorkspaceSeed) if err != nil { t.Fatalf("decode workspace seed: %v", err) } var seedFiles []domain.PluginAssetFile if err := json.Unmarshal(decodedSeed, &seedFiles); err != nil { t.Fatalf("unmarshal workspace seed: %v", err) } if platformInput.ProfileKey != "local" || len(seedFiles) != 3 || seedFiles[1].Path != "bin/install-server" || seedFiles[1].Content == "" { t.Fatalf("platform builder received incomplete plugin workspace seed: profile=%q seed=%+v", platformInput.ProfileKey, seedFiles) } if seedFiles[2].Path != ".platform/autonomous-lifecycle-plan.json" || seedFiles[2].Content == "" || seedFiles[2].Mode != 0o600 { t.Fatalf("workspace seed did not include autonomous lifecycle plan file: %+v", seedFiles) } plan := platformInput.AutonomousLifecycle if plan == nil || plan.SchemaVersion != "1" || plan.ServerInstanceID != instance.ID || plan.PluginID != plugin.ID || plan.ProfileKey != "local" || plan.Bootstrap == nil || plan.Bootstrap.Action != domain.ServerLifecycleActionStart || plan.Bootstrap.TargetKey != "actions/start.json" { t.Fatalf("platform builder received incomplete autonomous lifecycle plan: %+v", plan) } if len(plan.DependencyProbes) != 1 || plan.DependencyProbes[0].Key != "java-runtime" || len(plan.InstallPlans) != 1 || plan.InstallPlans[0].Key != "java-install" || len(plan.LogSources) != 3 || !hasAutonomousLogSource(plan.LogSources, "process.stdout", "console") || !hasAutonomousLogSource(plan.LogSources, "file.tail", "latest-log") || !hasAutonomousLogSource(plan.LogSources, "file.tail", "scum.server") || len(plan.DataTargets) != 1 || plan.DataTargets[0].WorkspaceKey != "databases/world-db" || plan.DataTargets[0].SourcePath != "world/current.db" || plan.RuntimeBindings["logs/latest"] != "runtime.logs.latest" { t.Fatalf("autonomous lifecycle plan lost plugin runtime declarations: %+v", plan) } var seededPlan domain.RunAutonomousLifecyclePlan if err := json.Unmarshal([]byte(seedFiles[2].Content), &seededPlan); err != nil { t.Fatalf("unmarshal seeded autonomous lifecycle plan: %v", err) } if seededPlan.ServerInstanceID != plan.ServerInstanceID || seededPlan.Bootstrap == nil || seededPlan.Bootstrap.TargetKey != plan.Bootstrap.TargetKey || len(seededPlan.DataTargets) != 1 || seededPlan.DataTargets[0].WorkspaceKey != "databases/world-db" { t.Fatalf("seeded lifecycle plan differs from build input: seed=%+v input=%+v", seededPlan, plan) } auth, err := svc.AuthenticateComponent(domain.ComponentAuthenticationRequest{ ServerInstanceID: instance.ID, ComponentKind: domain.DistributionComponentRun, Generation: platformInput.KeyGeneration, Key: platformInput.AuthKey, }) if err != nil || !auth.Allowed { t.Fatalf("platform builder did not receive the active plaintext component key: auth=%+v err=%v", auth, err) } helloRequest := validRunControlHello() helloRequest.CapabilityReport.Capabilities = append(helloRequest.CapabilityReport.Capabilities, domain.JobCapabilityDistributionBuild) hello, err := svc.RegisterRunHello(helloRequest) if err != nil { t.Fatalf("register machine endpoint: %v", err) } claim, err := svc.ClaimRunJob(domain.RunJobClaim{ RunEndpointID: instance.RunEndpointID, SessionToken: hello.SessionToken, Capabilities: []string{domain.JobCapabilityDistributionBuild}, Capacity: domain.RunCapacity{MaxJobs: 1}, }) if err != nil { t.Fatalf("claim machine jobs: %v", err) } if claim.HasJob { t.Fatalf("machine endpoint received platform-owned build job: %+v", claim) } machineInput, err := svc.GetDistributionBuildInput(domain.DistributionBuildInputRequest{ RunEndpointID: instance.RunEndpointID, SessionToken: hello.SessionToken, JobID: distribution.BuildJobID, LeaseToken: "machine-cannot-hold-platform-build-lease", Attempt: 1, }) if err == nil { t.Fatalf("machine endpoint unexpectedly read platform build input: %+v", machineInput) } if machineInput.AuthKey != "" || strings.Contains(err.Error(), platformInput.AuthKey) { t.Fatalf("machine build-input rejection leaked plaintext auth key: input=%+v err=%v", machineInput, err) } releaseOnce.Do(func() { close(release) }) completeDistributionBuild(t, svc, distribution, nil) } func hasAutonomousLogSource(sources []domain.RunAutonomousLogSource, kind string, streamKey string) bool { for _, source := range sources { if source.Kind == kind && source.StreamKey == streamKey { return true } } return false } func TestCoreServiceBuildsWithoutRegisteredDistributionWorker(t *testing.T) { svc, session, instance := newDistributionTestFixture(t) bootstrapEndpoint, err := svc.store.RunEndpoints().Get(instance.RunEndpointID) if err != nil { t.Fatalf("get bootstrap endpoint: %v", err) } capabilities := bootstrapEndpoint.Capabilities[:0] for _, capability := range bootstrapEndpoint.Capabilities { if capability != domain.JobCapabilityDistributionBuild { capabilities = append(capabilities, capability) } } bootstrapEndpoint.Capabilities = capabilities bootstrapEndpoint.Status = domain.RunEndpointStatusOffline if err := svc.store.RunEndpoints().Update(bootstrapEndpoint); err != nil { t.Fatalf("remove build worker capability: %v", err) } instance.RunEndpointID = dedicatedRunEndpointID(instance.ID) instance.DeploymentTargetID = "" if err := svc.store.ServerInstances().Update(instance); err != nil { t.Fatalf("prepare unregistered dedicated Run binding: %v", err) } if _, err := svc.store.RunEndpoints().Get(instance.RunEndpointID); !errors.Is(err, repo.ErrNotFound) { t.Fatalf("dedicated Run must be unregistered before generation, got %v", err) } distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "windows", TargetArch: "amd64", IdempotencyKey: "no-machine-distribution-worker", }) if err != nil { t.Fatalf("generate without registered distribution worker: %v", err) } distribution = completeDistributionBuild(t, svc, distribution, nil) if distribution.RunEndpointID != instance.RunEndpointID || distribution.Status != domain.DistributionStatusAvailable { t.Fatalf("unexpected platform-built distribution: %+v", distribution) } } func TestCoreServiceGeneratedRunOnlyEndpointCanGenerateAnotherRun(t *testing.T) { svc, session, instance := newDistributionTestFixture(t) instance.State = domain.ServerInstanceStateFailed if err := svc.store.ServerInstances().Update(instance); err != nil { t.Fatalf("prepare server for dedicated Run generation: %v", err) } first, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "generated-run-only-first", }) if err != nil { t.Fatalf("generate first dedicated Run: %v", err) } first = completeDistributionBuild(t, svc, first, nil) packageConfig := readGeneratedPackageConfig(t, svc, session, first.ArtifactID) instance, err = svc.GetServerInstance(instance.ID) if err != nil { t.Fatalf("get dedicated Run binding: %v", err) } bootstrap, err := svc.store.RunEndpoints().Get(instance.DeploymentTargetID) if err != nil { t.Fatalf("get former bootstrap endpoint: %v", err) } bootstrap.Status = domain.RunEndpointStatusOffline if err := svc.store.RunEndpoints().Update(bootstrap); err != nil { t.Fatalf("take former bootstrap endpoint offline: %v", err) } helloRequest := validRunControlHello() helloRequest.RunEndpointID = instance.RunEndpointID helloRequest.RegistrationToken = packageConfig.AuthKey helloRequest.ServerInstanceID = instance.ID helloRequest.PluginID = instance.PluginID helloRequest.ComponentKind = domain.DistributionComponentRun helloRequest.KeyGeneration = packageConfig.KeyGeneration capabilities := helloRequest.CapabilityReport.Capabilities[:0] for _, capability := range helloRequest.CapabilityReport.Capabilities { if capability != domain.JobCapabilityDistributionBuild { capabilities = append(capabilities, capability) } } helloRequest.CapabilityReport.Capabilities = capabilities registered, err := svc.RegisterRunHello(helloRequest) if err != nil || !registered.Accepted { t.Fatalf("register generated Run: result=%+v err=%v", registered, err) } online, err := svc.store.RunEndpoints().List(domain.RunEndpointFilter{Status: domain.RunEndpointStatusOnline}) if err != nil || len(online) != 1 || online[0].ID != instance.RunEndpointID { t.Fatalf("expected generated Run to be the only online endpoint: endpoints=%+v err=%v", online, err) } for _, capability := range online[0].Capabilities { if capability == domain.JobCapabilityDistributionBuild { t.Fatalf("generated Run must not advertise distribution build: %+v", online[0]) } } claim, err := svc.ClaimRunJob(domain.RunJobClaim{ RunEndpointID: instance.RunEndpointID, SessionToken: registered.SessionToken, Capabilities: []string{domain.JobCapabilityDistributionBuild}, Capacity: domain.RunCapacity{MaxJobs: 1}, }) if err != nil || claim.HasJob { t.Fatalf("generated Run must not receive platform build work: claim=%+v err=%v", claim, err) } second, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "generated-run-only-second", }) if err != nil { t.Fatalf("regenerate with generated Run as only endpoint: %v", err) } second = completeDistributionBuild(t, svc, second, nil) job, err := svc.GetJob(second.BuildJobID) if err != nil || job.RunEndpointID != platformDistributionBuilderEndpointID || job.State != domain.JobStateSucceeded || second.Status != domain.DistributionStatusAvailable { t.Fatalf("expected platform-built regenerated Run: job=%+v distribution=%+v err=%v", job, second, err) } } func TestCoreServiceDoesNotDuplicateInFlightPlatformBuild(t *testing.T) { svc, session, instance := newDistributionTestFixture(t) inputs := make(chan domain.DistributionBuildInput, 2) release := make(chan struct{}) var releaseOnce sync.Once t.Cleanup(func() { releaseOnce.Do(func() { close(release) }) }) svc.ConfigureDistributionBuilder(captureDistributionBuilder{inputs: inputs, release: release, payload: []byte("idempotent-platform-build")}) request := domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "same-platform-build", } first, err := svc.GenerateRunDistributionForSession(session, request) if err != nil { t.Fatalf("generate first distribution: %v", err) } select { case <-inputs: case <-time.After(time.Second): t.Fatal("first platform build did not start") } second, err := svc.GenerateRunDistributionForSession(session, request) if err != nil { t.Fatalf("repeat idempotent generation: %v", err) } if second.ID != first.ID || second.BuildJobID != first.BuildJobID { t.Fatalf("idempotent generation returned different work: first=%+v second=%+v", first, second) } select { case duplicate := <-inputs: t.Fatalf("idempotent generation started duplicate platform build: %+v", duplicate) case <-time.After(20 * time.Millisecond): } releaseOnce.Do(func() { close(release) }) completeDistributionBuild(t, svc, first, nil) } func TestCoreServiceDiscardsBuildCompletedAfterKeyReset(t *testing.T) { svc, session, instance := newDistributionTestFixture(t) inputs := make(chan domain.DistributionBuildInput, 1) release := make(chan struct{}) var releaseOnce sync.Once t.Cleanup(func() { releaseOnce.Do(func() { close(release) }) }) svc.ConfigureDistributionBuilder(captureDistributionBuilder{inputs: inputs, release: release, payload: []byte("stale-key-build")}) distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "reset-during-platform-build", }) if err != nil { t.Fatalf("generate distribution: %v", err) } select { case <-inputs: case <-time.After(time.Second): t.Fatal("platform builder did not start") } if _, err := svc.ResetComponentKeyForSession(session, domain.ComponentKeyResetRequest{ ServerInstanceID: instance.ID, ComponentKind: domain.DistributionComponentRun, }); err != nil { t.Fatalf("reset component key during build: %v", err) } releaseOnce.Do(func() { close(release) }) deadline := time.Now().Add(time.Second) for time.Now().Before(deadline) { job, jobErr := svc.GetJob(distribution.BuildJobID) updated, distributionErr := svc.store.RunDistributions().Get(distribution.ID) if jobErr != nil || distributionErr != nil { t.Fatalf("read stale build state: jobErr=%v distributionErr=%v", jobErr, distributionErr) } if job.State == domain.JobStateFailed { if updated.Status != domain.DistributionStatusRevoked || !strings.Contains(job.Progress.Message, "no longer current") { t.Fatalf("stale build did not remain revoked: job=%+v distribution=%+v", job, updated) } if artifact, artifactErr := svc.GetArtifact(distribution.ArtifactID); artifactErr == nil && artifact.State == domain.ArtifactStateAvailable { t.Fatalf("stale-key build published an available artifact: %+v", artifact) } return } time.Sleep(time.Millisecond) } t.Fatal("stale-key platform build did not terminate") } func TestCoreServicePreservesSucceededArtifactOnDuplicatePlatformResult(t *testing.T) { svc, session, instance := newDistributionTestFixture(t) distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "preserve-duplicate-platform-result", }) if err != nil { t.Fatalf("generate distribution: %v", err) } distribution = completeDistributionBuild(t, svc, distribution, nil) job, err := svc.GetJob(distribution.BuildJobID) if err != nil { t.Fatalf("get completed build job: %v", err) } if job.State != domain.JobStateSucceeded || job.ResultRef != "artifact://"+distribution.ArtifactID { t.Fatalf("expected succeeded build job with distribution artifact ref, got %+v", job) } if err := svc.succeedDistributionBuildJob(job, distribution.ArtifactID); err != nil { t.Fatalf("replay succeeded platform result: %v", err) } artifact, err := svc.GetArtifact(distribution.ArtifactID) if err != nil { t.Fatalf("get distribution artifact: %v", err) } if artifact.State != domain.ArtifactStateAvailable { t.Fatalf("duplicate platform result expired active artifact: %+v", artifact) } } func TestCoreServiceResumesPendingPlatformBuildAfterEnvelopeConfiguration(t *testing.T) { store := repo.NewMemoryStore() artifactStore := NewMemoryArtifactBodyStore() first, err := NewCoreServiceWithDurableStores(store, NewMemoryLogBodyStore(), artifactStore) if err != nil { t.Fatalf("create first durable CoreService: %v", err) } const envelopeKey = "restart-test-secret-envelope-key-must-remain-stable" if err := first.ConfigureSecretEnvelopeKey(envelopeKey); err != nil { t.Fatalf("configure first envelope key: %v", err) } plugin, endpoint := createPluginAndRunEndpoint(t, first) plugin.SupportedOS = []string{"linux"} plugin.DeclaredPermissions = append(plugin.DeclaredPermissions, "server.run.distribution") plugin.BridgeActions = append(plugin.BridgeActions, string(domain.PluginBridgeActionRunDistribution)) if err := first.store.GamePlugins().Update(plugin); err != nil { t.Fatalf("enable distribution fixture plugin: %v", err) } session := createServiceUserAndLogin(t, first, domain.User{ ID: "restart-distribution-owner", DisplayName: "Restart Distribution Owner", Email: "restart-distribution-owner@example.test", Roles: []string{"server-owner"}, PasswordHash: "secret-password", }) instance, err := first.CreateServerInstanceForSession(session, domain.ServerInstance{ ID: "server-restart-distribution", PluginID: plugin.ID, RunEndpointID: endpoint.ID, Name: "Restart Distribution Server", State: domain.ServerInstanceStateReady, }) if err != nil { t.Fatalf("create restart distribution server: %v", err) } createCompleteRuntimeBinding(t, first, instance, "local") inputs := make(chan domain.DistributionBuildInput, 1) release := make(chan struct{}) var releaseOnce sync.Once t.Cleanup(func() { releaseOnce.Do(func() { close(release) }) }) first.ConfigureDistributionBuilder(captureDistributionBuilder{inputs: inputs, release: release, payload: []byte("first-process-output")}) distribution, err := first.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "resume-after-envelope-configuration", }) if err != nil { t.Fatalf("queue durable platform build: %v", err) } select { case <-inputs: case <-time.After(time.Second): t.Fatal("first CoreService did not begin the pending build") } restarted, err := NewCoreServiceWithDurableStores(store, NewMemoryLogBodyStore(), artifactStore) if err != nil { t.Fatalf("create restarted durable CoreService: %v", err) } if err := restarted.ConfigureSecretEnvelopeKey(envelopeKey); err != nil { t.Fatalf("configure restarted envelope key: %v", err) } restarted.ConfigureDistributionBuilder(staticDistributionBuilder{payload: []byte("recovered-platform-build")}) distribution = completeDistributionBuild(t, restarted, distribution, nil) job, err := restarted.GetJob(distribution.BuildJobID) if err != nil || job.State != domain.JobStateSucceeded || job.RunEndpointID != platformDistributionBuilderEndpointID { t.Fatalf("expected recovered platform build to succeed: job=%+v err=%v", job, err) } artifact, err := restarted.GetArtifact(distribution.ArtifactID) if err != nil || artifact.State != domain.ArtifactStateAvailable || artifact.OwnerKind != domain.ArtifactOwnerKindJob || artifact.OwnerID != job.ID { t.Fatalf("expected recovered job-owned artifact: artifact=%+v err=%v", artifact, err) } } func TestCoreServiceReportsPlatformBuilderUnavailable(t *testing.T) { svc, session, instance := newDistributionTestFixture(t) svc.ConfigureDistributionBuilder(staticDistributionBuilder{err: errors.New("platform builder image is unavailable")}) actions, err := svc.GetServerRuntimeActionsForSession(session, instance.ID) if err != nil { t.Fatalf("get runtime actions: %v", err) } seen := 0 for _, action := range actions.Actions { if action.Key != "generate-run" && action.Key != "generate-client-manager" { continue } seen++ if action.Available || !strings.Contains(action.Reason, "platform builder") || strings.Contains(action.Reason, "run endpoint") { t.Fatalf("expected explicit platform builder unavailable reason, got %+v", action) } } if seen != 2 { t.Fatalf("expected both build actions, got %+v", actions.Actions) } _, err = svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "builder-unavailable", }) if err == nil || !strings.Contains(err.Error(), "platform builder") || strings.Contains(err.Error(), "run endpoint") { t.Fatalf("expected generation to report platform builder failure, got %v", err) } }