Files
browser/platform/service/distribution_build_execution_test.go
T

458 lines
18 KiB
Go

package service
import (
"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)
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: "windows",
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)
}
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 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)
}
}