feat: move distribution builds to platform Docker builder
This commit is contained in:
@@ -229,11 +229,13 @@ func buildLifecycleDistribution(t *testing.T, svc *CoreService, session string,
|
||||
if err := svc.store.GamePlugins().Update(plugin); err != nil {
|
||||
t.Fatalf("update lifecycle version: %v", err)
|
||||
}
|
||||
payload := []byte("client-manager-package-" + version)
|
||||
svc.ConfigureDistributionBuilder(staticDistributionBuilder{payload: payload})
|
||||
distribution, err := svc.GenerateClientManagerDistributionForSession(session, domain.ClientManagerBuildRequest{ServerInstanceID: instance.ID, ProfileKey: "scum-client-manager", TargetOS: "linux", TargetArch: "amd64", RepositoryURL: "https://github.com/F88888/scum_client.git", SourceRevision: "main", IdempotencyKey: idempotency})
|
||||
if err != nil {
|
||||
t.Fatalf("generate lifecycle distribution: %v", err)
|
||||
}
|
||||
return completeClientDistributionBuild(t, svc, distribution, []byte("client-manager-package-"+version))
|
||||
return completeClientDistributionBuild(t, svc, distribution, payload)
|
||||
}
|
||||
|
||||
func registerClientManagerRun(t *testing.T, svc *CoreService) string {
|
||||
|
||||
@@ -195,11 +195,12 @@ func TestPluginBridgeDependencyInstallUsesReviewedPlanDigest(t *testing.T) {
|
||||
|
||||
func TestRunUpdateTargetFencingChunksHealthAndRollbackProjection(t *testing.T) {
|
||||
svc, session, instance := newDistributionTestFixture(t)
|
||||
payload := []byte("compiled target-matched run archive")
|
||||
svc.ConfigureDistributionBuilder(staticDistributionBuilder{payload: payload})
|
||||
distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{ServerInstanceID: instance.ID, TargetOS: "linux", TargetArch: "amd64", IdempotencyKey: "run-update-build"})
|
||||
if err != nil {
|
||||
t.Fatalf("generate update distribution: %v", err)
|
||||
}
|
||||
payload := []byte("compiled target-matched run archive")
|
||||
distribution = completeDistributionBuild(t, svc, distribution, payload)
|
||||
otherInstance, err := svc.CreateServerInstanceForSession(session, domain.ServerInstance{ID: "server-update-other", PluginID: instance.PluginID, RunEndpointID: instance.RunEndpointID, Name: "Other Update Server", State: domain.ServerInstanceStateReady})
|
||||
if err != nil {
|
||||
|
||||
@@ -0,0 +1,320 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strings"
|
||||
|
||||
"browser.local/platform/domain"
|
||||
"browser.local/platform/repo"
|
||||
"browser.local/platform/validator"
|
||||
)
|
||||
|
||||
const platformDistributionBuilderEndpointID = "platform-distribution-builder"
|
||||
|
||||
// unconfiguredDistributionBuilder stands in when no platform builder has been
|
||||
// installed. It reports a platform-builder reason so an operator is not sent
|
||||
// looking at run endpoint capabilities.
|
||||
type unconfiguredDistributionBuilder struct{}
|
||||
|
||||
func (unconfiguredDistributionBuilder) Readiness() (bool, string) {
|
||||
return false, "platform builder is not configured"
|
||||
}
|
||||
|
||||
func (unconfiguredDistributionBuilder) Build(domain.DistributionBuildInput) ([]byte, error) {
|
||||
return nil, validationError("platform builder is not configured")
|
||||
}
|
||||
|
||||
func (svc *CoreService) configuredDistributionBuilder() DistributionBuilder {
|
||||
svc.distributionBuildMu.Lock()
|
||||
defer svc.distributionBuildMu.Unlock()
|
||||
if svc.distributionBuilder == nil {
|
||||
return unconfiguredDistributionBuilder{}
|
||||
}
|
||||
return svc.distributionBuilder
|
||||
}
|
||||
|
||||
// distributionBuilderReadiness reports platform builder readiness. The reason
|
||||
// always names the platform builder, never a run endpoint capability.
|
||||
func (svc *CoreService) distributionBuilderReadiness() (bool, string) {
|
||||
builder := svc.configuredDistributionBuilder()
|
||||
ready, reason := builder.Readiness()
|
||||
if ready {
|
||||
return true, ""
|
||||
}
|
||||
reason = strings.TrimSpace(reason)
|
||||
if reason == "" {
|
||||
reason = "platform builder is unavailable"
|
||||
} else if !strings.Contains(strings.ToLower(reason), "platform builder") {
|
||||
reason = "platform builder is unavailable: " + reason
|
||||
}
|
||||
return false, reason
|
||||
}
|
||||
|
||||
func (svc *CoreService) enqueueDistributionBuild(job domain.Job) {
|
||||
svc.distributionBuildMu.Lock()
|
||||
if _, exists := svc.distributionBuilds[job.ID]; exists {
|
||||
svc.distributionBuildMu.Unlock()
|
||||
return
|
||||
}
|
||||
svc.distributionBuilds[job.ID] = struct{}{}
|
||||
svc.distributionBuildMu.Unlock()
|
||||
go func() {
|
||||
defer func() {
|
||||
svc.distributionBuildMu.Lock()
|
||||
delete(svc.distributionBuilds, job.ID)
|
||||
svc.distributionBuildMu.Unlock()
|
||||
}()
|
||||
_ = svc.executeDistributionBuild(job)
|
||||
}()
|
||||
}
|
||||
|
||||
// executeDistributionBuild claims the build job internally and completes it with
|
||||
// the artifact produced by the platform builder. Build work is never dispatched
|
||||
// to a machine-side run endpoint, so the component auth key stays on the
|
||||
// platform.
|
||||
func (svc *CoreService) executeDistributionBuild(job domain.Job) error {
|
||||
if job.Capability != domain.JobCapabilityDistributionBuild {
|
||||
return validationError("job is not a distribution build")
|
||||
}
|
||||
input, err := svc.platformDistributionBuildInput(job)
|
||||
if err != nil {
|
||||
return svc.failDistributionBuildJob(job, "platform builder could not assemble build input")
|
||||
}
|
||||
if err := svc.markDistributionBuildRunning(&job); err != nil {
|
||||
return err
|
||||
}
|
||||
if isTerminalJobState(job.State) {
|
||||
return nil
|
||||
}
|
||||
payload, buildErr := svc.configuredDistributionBuilder().Build(input)
|
||||
if buildErr != nil {
|
||||
return svc.failDistributionBuildJob(job, builderJobFailureMessage(buildErr))
|
||||
}
|
||||
if _, err := svc.platformDistributionBuildInput(job); err != nil {
|
||||
return svc.failDistributionBuildJob(job, "platform builder discarded output because the component key is no longer current")
|
||||
}
|
||||
if err := svc.storeDistributionBuildArtifact(input.ArtifactID, job.ID, payload); err != nil {
|
||||
return svc.failDistributionBuildJob(job, "platform builder could not record the distribution artifact")
|
||||
}
|
||||
return svc.succeedDistributionBuildJob(job, input.ArtifactID)
|
||||
}
|
||||
|
||||
// platformDistributionBuildInput resolves the build input, including the
|
||||
// component auth key, inside the platform. Unlike GetDistributionBuildInput it
|
||||
// never crosses the job channel.
|
||||
func (svc *CoreService) platformDistributionBuildInput(job domain.Job) (domain.DistributionBuildInput, error) {
|
||||
runDistributions, err := svc.store.RunDistributions().List(domain.RunDistributionFilter{ServerInstanceID: job.ServerInstanceID})
|
||||
if err != nil {
|
||||
return domain.DistributionBuildInput{}, err
|
||||
}
|
||||
for _, distribution := range runDistributions {
|
||||
if distribution.BuildJobID != job.ID {
|
||||
continue
|
||||
}
|
||||
key, err := svc.activeComponentKey(distribution.ServerInstanceID, domain.DistributionComponentRun, "")
|
||||
if err != nil {
|
||||
return domain.DistributionBuildInput{}, err
|
||||
}
|
||||
if key.Generation != distribution.KeyGeneration {
|
||||
return domain.DistributionBuildInput{}, validationError("run build key generation is no longer current")
|
||||
}
|
||||
plainKey, err := svc.decryptRuntimeKey(key.EncryptedKey)
|
||||
if err != nil {
|
||||
return domain.DistributionBuildInput{}, err
|
||||
}
|
||||
return domain.DistributionBuildInput{
|
||||
JobID: job.ID,
|
||||
ComponentKind: domain.DistributionComponentRun,
|
||||
ServerInstanceID: distribution.ServerInstanceID,
|
||||
PluginID: distribution.PluginID,
|
||||
RunEndpointID: distribution.RunEndpointID,
|
||||
TargetOS: distribution.TargetOS,
|
||||
TargetArch: distribution.TargetArch,
|
||||
TargetRelease: distribution.ID,
|
||||
PlatformURL: runReleasePlatformURL(),
|
||||
PackageFormat: distribution.PackageFormat,
|
||||
ArtifactID: distribution.ArtifactID,
|
||||
OutputFilename: executableFilename("run", distribution.TargetOS),
|
||||
SecretRef: distribution.SecretRef,
|
||||
KeyGeneration: distribution.KeyGeneration,
|
||||
AuthKey: plainKey,
|
||||
}, nil
|
||||
}
|
||||
|
||||
clientDistributions, err := svc.store.ClientManagerDistributions().List(domain.ClientManagerDistributionFilter{ServerInstanceID: job.ServerInstanceID})
|
||||
if err != nil {
|
||||
return domain.DistributionBuildInput{}, err
|
||||
}
|
||||
for _, distribution := range clientDistributions {
|
||||
if distribution.BuildJobID != job.ID {
|
||||
continue
|
||||
}
|
||||
key, err := svc.activeComponentKey(distribution.ServerInstanceID, domain.DistributionComponentClientManager, distribution.ProfileKey)
|
||||
if err != nil {
|
||||
return domain.DistributionBuildInput{}, err
|
||||
}
|
||||
if key.Generation != distribution.KeyGeneration {
|
||||
return domain.DistributionBuildInput{}, validationError("client-manager build key generation is no longer current")
|
||||
}
|
||||
plainKey, err := svc.decryptRuntimeKey(key.EncryptedKey)
|
||||
if err != nil {
|
||||
return domain.DistributionBuildInput{}, err
|
||||
}
|
||||
return domain.DistributionBuildInput{
|
||||
JobID: job.ID,
|
||||
ComponentKind: domain.DistributionComponentClientManager,
|
||||
ServerInstanceID: distribution.ServerInstanceID,
|
||||
PluginID: distribution.PluginID,
|
||||
RunEndpointID: job.RunEndpointID,
|
||||
ProfileKey: distribution.ProfileKey,
|
||||
TargetOS: distribution.TargetOS,
|
||||
TargetArch: distribution.TargetArch,
|
||||
PlatformURL: runReleasePlatformURL(),
|
||||
PackageFormat: packageFormatForTarget(distribution.TargetOS),
|
||||
RepositoryURL: distribution.RepositoryURL,
|
||||
SourceRevision: distribution.SourceRevision,
|
||||
ArtifactID: distribution.ArtifactID,
|
||||
OutputFilename: clientManagerOutputName(distribution.ProfileKey, distribution.TargetOS),
|
||||
SecretRef: distribution.SecretRef,
|
||||
KeyGeneration: distribution.KeyGeneration,
|
||||
AuthKey: plainKey,
|
||||
}, nil
|
||||
}
|
||||
return domain.DistributionBuildInput{}, repo.ErrNotFound
|
||||
}
|
||||
|
||||
// storeDistributionBuildArtifact records the built package under job ownership
|
||||
// so the existing ArtifactOwnerKindJob scope assertions keep guarding it.
|
||||
func (svc *CoreService) storeDistributionBuildArtifact(artifactID string, jobID string, payload []byte) error {
|
||||
if strings.TrimSpace(artifactID) == "" {
|
||||
return validationError("distribution build artifact id is required")
|
||||
}
|
||||
if len(payload) == 0 {
|
||||
return validationError("distribution build produced no package bytes")
|
||||
}
|
||||
stamp := svc.now()
|
||||
svc.artifactMu.Lock()
|
||||
defer svc.artifactMu.Unlock()
|
||||
|
||||
artifact, err := svc.store.Artifacts().Get(artifactID)
|
||||
if err != nil && !errors.Is(err, repo.ErrNotFound) {
|
||||
return err
|
||||
}
|
||||
create := errors.Is(err, repo.ErrNotFound)
|
||||
if create {
|
||||
artifact = domain.Artifact{ID: artifactID, OwnerKind: domain.ArtifactOwnerKindJob, OwnerID: jobID, CreatedAt: stamp}
|
||||
}
|
||||
if artifact.OwnerKind != domain.ArtifactOwnerKindJob || artifact.OwnerID != jobID {
|
||||
return validationError("distribution build artifact is outside the job scope")
|
||||
}
|
||||
artifact.SizeBytes = int64(len(payload))
|
||||
artifact.Checksum = validator.BytesChecksum(payload)
|
||||
artifact.State = domain.ArtifactStateAvailable
|
||||
artifact.UpdatedAt = stamp
|
||||
if err := validator.ValidateArtifact(artifact); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := svc.artifactStore.PutPayload(artifact.ID, payload); err != nil {
|
||||
return err
|
||||
}
|
||||
if create {
|
||||
if err := svc.store.Artifacts().Create(artifact); err != nil {
|
||||
return err
|
||||
}
|
||||
} else if err := svc.store.Artifacts().Update(artifact); err != nil {
|
||||
return err
|
||||
}
|
||||
svc.artifactPayloads[artifact.ID] = domain.CopyBytes(payload)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) markDistributionBuildRunning(job *domain.Job) error {
|
||||
stamp := svc.now()
|
||||
svc.jobMu.Lock()
|
||||
defer svc.jobMu.Unlock()
|
||||
|
||||
current, err := svc.store.Jobs().Get(job.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if isTerminalJobState(current.State) {
|
||||
*job = current
|
||||
return nil
|
||||
}
|
||||
current.State = domain.JobStateRunning
|
||||
current.Attempt = maxInt(current.Attempt, 1)
|
||||
current.Progress = domain.JobProgress{Percent: 5, Phase: current.Progress.Phase, Message: "platform builder started"}
|
||||
current.UpdatedAt = stamp
|
||||
if err := svc.updateScheduledJob(current); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := svc.projectDistributionBuildProgress(current, stamp); err != nil {
|
||||
return err
|
||||
}
|
||||
*job = current
|
||||
return nil
|
||||
}
|
||||
|
||||
func (svc *CoreService) succeedDistributionBuildJob(job domain.Job, artifactID string) error {
|
||||
stamp := svc.now()
|
||||
svc.jobMu.Lock()
|
||||
current, err := svc.store.Jobs().Get(job.ID)
|
||||
if err != nil {
|
||||
svc.jobMu.Unlock()
|
||||
return err
|
||||
}
|
||||
if isTerminalJobState(current.State) {
|
||||
if current.State == domain.JobStateSucceeded && current.ResultRef == "artifact://"+artifactID {
|
||||
svc.jobMu.Unlock()
|
||||
return nil
|
||||
}
|
||||
svc.jobMu.Unlock()
|
||||
_ = svc.expireDistributionArtifact(artifactID)
|
||||
return nil
|
||||
}
|
||||
current.State = domain.JobStateSucceeded
|
||||
current.ResultRef = "artifact://" + artifactID
|
||||
current.Progress = domain.JobProgress{Percent: 100, Phase: current.Progress.Phase, Message: "platform builder completed"}
|
||||
current.LeaseTokenHash = ""
|
||||
current.TerminalAt = stamp
|
||||
current.TerminalFingerprint = "platform-builder:succeeded:" + artifactID
|
||||
current.UpdatedAt = stamp
|
||||
if err := svc.updateScheduledJob(current); err != nil {
|
||||
svc.jobMu.Unlock()
|
||||
return err
|
||||
}
|
||||
svc.jobMu.Unlock()
|
||||
return svc.projectDistributionBuildResult(current, stamp)
|
||||
}
|
||||
|
||||
func (svc *CoreService) failDistributionBuildJob(job domain.Job, reason string) error {
|
||||
stamp := svc.now()
|
||||
if strings.TrimSpace(reason) == "" {
|
||||
reason = "platform builder failed"
|
||||
}
|
||||
svc.jobMu.Lock()
|
||||
current, err := svc.store.Jobs().Get(job.ID)
|
||||
if err != nil {
|
||||
svc.jobMu.Unlock()
|
||||
return err
|
||||
}
|
||||
if isTerminalJobState(current.State) {
|
||||
svc.jobMu.Unlock()
|
||||
return nil
|
||||
}
|
||||
current.State = domain.JobStateFailed
|
||||
current.Progress = domain.JobProgress{Percent: current.Progress.Percent, Phase: current.Progress.Phase, Message: reason}
|
||||
current.LeaseTokenHash = ""
|
||||
current.TerminalAt = stamp
|
||||
current.TerminalFingerprint = "platform-builder:failed:" + reason
|
||||
current.UpdatedAt = stamp
|
||||
if err := svc.updateScheduledJob(current); err != nil {
|
||||
svc.jobMu.Unlock()
|
||||
return err
|
||||
}
|
||||
svc.jobMu.Unlock()
|
||||
if err := svc.projectDistributionBuildResult(current, stamp); err != nil && !errors.Is(err, repo.ErrNotFound) {
|
||||
return err
|
||||
}
|
||||
return validationError(reason)
|
||||
}
|
||||
@@ -0,0 +1,457 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -168,6 +168,12 @@ func (svc *CoreService) projectDistributionBuildResult(job domain.Job, stamp tim
|
||||
if distribution.BuildJobID != job.ID {
|
||||
continue
|
||||
}
|
||||
if distribution.Status == domain.DistributionStatusRevoked {
|
||||
if artifact.ID != "" {
|
||||
return svc.expireDistributionArtifact(artifact.ID)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if job.State == domain.JobStateSucceeded && artifact.ID != distribution.ArtifactID {
|
||||
return validationError("distribution build returned an unexpected artifact")
|
||||
}
|
||||
@@ -190,6 +196,12 @@ func (svc *CoreService) projectDistributionBuildResult(job domain.Job, stamp tim
|
||||
if distribution.BuildJobID != job.ID {
|
||||
continue
|
||||
}
|
||||
if distribution.Status == domain.DistributionStatusRevoked {
|
||||
if artifact.ID != "" {
|
||||
return svc.expireDistributionArtifact(artifact.ID)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
if job.State == domain.JobStateSucceeded && artifact.ID != distribution.ArtifactID {
|
||||
return validationError("client-manager build returned an unexpected artifact")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,400 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"archive/tar"
|
||||
"archive/zip"
|
||||
"bytes"
|
||||
"compress/gzip"
|
||||
"context"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"browser.local/platform/domain"
|
||||
)
|
||||
|
||||
// DistributionBuilder executes a distribution build inside a platform-owned
|
||||
// container. Builds are a platform responsibility: they must not depend on a
|
||||
// machine-side run endpoint being registered and online, and the component auth
|
||||
// key must never leave the platform.
|
||||
type DistributionBuilder interface {
|
||||
// Readiness reports whether the platform builder can execute a build. The
|
||||
// reason must name the platform builder rather than a run endpoint
|
||||
// capability, so an operator is not sent looking at the wrong subsystem.
|
||||
Readiness() (bool, string)
|
||||
// Build assembles the package described by input and returns its bytes.
|
||||
Build(input domain.DistributionBuildInput) ([]byte, error)
|
||||
}
|
||||
|
||||
// DockerDistributionBuilderConfig configures a container-per-build builder.
|
||||
type DockerDistributionBuilderConfig struct {
|
||||
DockerBinary string
|
||||
Image string
|
||||
SourceDir string
|
||||
WorkspaceDir string
|
||||
Timeout time.Duration
|
||||
PlatformURL string
|
||||
CommandRunner func(ctx context.Context, name string, args ...string) ([]byte, error)
|
||||
}
|
||||
|
||||
// DockerDistributionBuilder runs each build in a container from a pinned image,
|
||||
// with the run source mounted read-only and a per-job output directory mounted
|
||||
// writable.
|
||||
type DockerDistributionBuilder struct {
|
||||
config DockerDistributionBuilderConfig
|
||||
}
|
||||
|
||||
func NewDockerDistributionBuilder(config DockerDistributionBuilderConfig) *DockerDistributionBuilder {
|
||||
if strings.TrimSpace(config.DockerBinary) == "" {
|
||||
config.DockerBinary = "docker"
|
||||
}
|
||||
if config.Timeout <= 0 {
|
||||
config.Timeout = 30 * time.Minute
|
||||
}
|
||||
if config.CommandRunner == nil {
|
||||
config.CommandRunner = runCommandCombined
|
||||
}
|
||||
return &DockerDistributionBuilder{config: config}
|
||||
}
|
||||
|
||||
func runCommandCombined(ctx context.Context, name string, args ...string) ([]byte, error) {
|
||||
return exec.CommandContext(ctx, name, args...).CombinedOutput()
|
||||
}
|
||||
|
||||
func (builder *DockerDistributionBuilder) Readiness() (bool, string) {
|
||||
if strings.TrimSpace(builder.config.Image) == "" {
|
||||
return false, "platform builder image is not configured"
|
||||
}
|
||||
if !pinnedBuilderImage(builder.config.Image) {
|
||||
return false, "platform builder image must be pinned to an explicit version or digest"
|
||||
}
|
||||
source := strings.TrimSpace(builder.config.SourceDir)
|
||||
if source == "" {
|
||||
return false, "platform builder run source directory is not configured"
|
||||
}
|
||||
source, err := filepath.Abs(strings.TrimSpace(builder.config.SourceDir))
|
||||
if err != nil {
|
||||
return false, "platform builder run source directory is invalid"
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(source, "go.mod")); err != nil {
|
||||
return false, "platform builder run source directory does not contain a run checkout"
|
||||
}
|
||||
if strings.TrimSpace(builder.config.WorkspaceDir) == "" {
|
||||
return false, "platform builder workspace directory is not configured"
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
if _, err := builder.config.CommandRunner(ctx, builder.config.DockerBinary, "version", "--format", "{{.Server.Version}}"); err != nil {
|
||||
return false, "platform builder container runtime is unavailable"
|
||||
}
|
||||
if _, err := builder.config.CommandRunner(ctx, builder.config.DockerBinary, "image", "inspect", builder.config.Image); err != nil {
|
||||
return false, "platform builder image is unavailable"
|
||||
}
|
||||
return true, ""
|
||||
}
|
||||
|
||||
// pinnedBuilderImage rejects floating references. An unpinned builder image
|
||||
// silently changes what the platform ships.
|
||||
func pinnedBuilderImage(image string) bool {
|
||||
image = strings.TrimSpace(image)
|
||||
if name, digest, found := strings.Cut(image, "@sha256:"); found {
|
||||
if strings.TrimSpace(name) == "" || len(digest) != 64 {
|
||||
return false
|
||||
}
|
||||
_, err := hex.DecodeString(digest)
|
||||
return err == nil
|
||||
}
|
||||
reference := image
|
||||
if slash := strings.LastIndex(image, "/"); slash >= 0 {
|
||||
reference = image[slash+1:]
|
||||
}
|
||||
_, tag, found := strings.Cut(reference, ":")
|
||||
if !found {
|
||||
return false
|
||||
}
|
||||
tag = strings.TrimSpace(tag)
|
||||
return tag != "" && tag != "latest"
|
||||
}
|
||||
|
||||
func (builder *DockerDistributionBuilder) Build(input domain.DistributionBuildInput) ([]byte, error) {
|
||||
if ready, reason := builder.Readiness(); !ready {
|
||||
return nil, validationError(reason)
|
||||
}
|
||||
sourceDir, err := filepath.Abs(strings.TrimSpace(builder.config.SourceDir))
|
||||
if err != nil {
|
||||
return nil, validationError("platform builder run source directory is invalid")
|
||||
}
|
||||
workspaceDir, err := filepath.Abs(strings.TrimSpace(builder.config.WorkspaceDir))
|
||||
if err != nil {
|
||||
return nil, validationError("platform builder workspace directory is invalid")
|
||||
}
|
||||
// Workspaces stay isolated per plugin and per job as required by
|
||||
// run-build-download-flow.
|
||||
jobDir := filepath.Join(workspaceDir, sanitizeIDPart(input.PluginID), sanitizeIDPart(input.JobID))
|
||||
if err := os.RemoveAll(jobDir); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
outputDir := filepath.Join(jobDir, "output")
|
||||
inputDir := filepath.Join(jobDir, "input")
|
||||
buildDir := filepath.Join(jobDir, "build")
|
||||
for _, directory := range []string{outputDir, inputDir, buildDir} {
|
||||
if err := os.MkdirAll(directory, 0o700); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
defer func() { _ = os.RemoveAll(jobDir) }()
|
||||
|
||||
// The auth key reaches the container through a per-job input file, never
|
||||
// through a job-channel response to a machine-side endpoint or a container
|
||||
// command-line argument.
|
||||
if strings.TrimSpace(input.AuthKey) == "" {
|
||||
return nil, validationError("distribution build input is missing a component auth key")
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(inputDir, "auth-key"), []byte(input.AuthKey), 0o600); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(inputDir, "build.sh"), []byte(distributionBuildScript), 0o500); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
outputName := strings.TrimSpace(input.OutputFilename)
|
||||
if outputName == "" || filepath.Base(outputName) != outputName {
|
||||
return nil, validationError("distribution build input has an invalid output filename")
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), builder.config.Timeout)
|
||||
defer cancel()
|
||||
args := builder.containerArgs(input, sourceDir, inputDir, buildDir, outputDir, outputName)
|
||||
if output, err := builder.config.CommandRunner(ctx, builder.config.DockerBinary, args...); err != nil {
|
||||
if errors.Is(ctx.Err(), context.DeadlineExceeded) {
|
||||
return nil, validationError("platform builder timed out while building the distribution")
|
||||
}
|
||||
return nil, validationError("platform builder failed: " + safeBuilderFailure(output, input.AuthKey, builder.config.SourceDir, sourceDir, jobDir))
|
||||
}
|
||||
binary, err := os.ReadFile(filepath.Join(outputDir, outputName))
|
||||
if err != nil {
|
||||
return nil, validationError("platform builder did not produce a distribution executable")
|
||||
}
|
||||
if len(binary) == 0 {
|
||||
return nil, validationError("platform builder produced an empty distribution executable")
|
||||
}
|
||||
if input.ComponentKind == domain.DistributionComponentRun {
|
||||
return binary, nil
|
||||
}
|
||||
configPayload, err := os.ReadFile(filepath.Join(outputDir, "config.yaml"))
|
||||
if err != nil {
|
||||
return nil, validationError("platform builder did not produce client-manager configuration")
|
||||
}
|
||||
return packageClientManagerDistribution(input.PackageFormat, outputName, binary, configPayload)
|
||||
}
|
||||
|
||||
const distributionBuildScript = `#!/bin/sh
|
||||
set -eu
|
||||
|
||||
auth_key="$(cat /workspace/input/auth-key)"
|
||||
if [ "$COMPONENT_KIND" = "run" ]; then
|
||||
cd /workspace/source
|
||||
ldflags="-s -w"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildMode=worker"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildPlatformURL=$PLATFORM_URL"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildRunEndpointID=$RUN_ENDPOINT_ID"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildDisplayName=Run-$SERVER_INSTANCE_ID"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildRegistrationToken=$auth_key"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildServerInstanceID=$SERVER_INSTANCE_ID"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildPluginID=$PLUGIN_ID"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildComponentKind=$COMPONENT_KIND"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildComponentKey=$PROFILE_KEY"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildKeyGeneration=$KEY_GENERATION"
|
||||
ldflags="$ldflags -X browser.local/run/config.BuildVersion=$TARGET_RELEASE"
|
||||
go mod download
|
||||
go build -trimpath -ldflags "$ldflags" -o "/workspace/output/$OUTPUT_FILENAME" ./cmd/run
|
||||
exit 0
|
||||
fi
|
||||
|
||||
if [ "$COMPONENT_KIND" != "client-manager" ]; then
|
||||
printf 'unsupported component kind\n' >&2
|
||||
exit 2
|
||||
fi
|
||||
case "$REPOSITORY_URL" in
|
||||
https://*) ;;
|
||||
*) printf 'client-manager repository must use https\n' >&2; exit 2 ;;
|
||||
esac
|
||||
cd /workspace/build
|
||||
git init --quiet
|
||||
git remote add origin "$REPOSITORY_URL"
|
||||
git fetch --quiet --depth 1 origin "$SOURCE_REVISION"
|
||||
git checkout --quiet --detach FETCH_HEAD
|
||||
{
|
||||
printf 'server_url: "%s"\n' "$PLATFORM_URL"
|
||||
printf 'server_instance_id: "%s"\n' "$SERVER_INSTANCE_ID"
|
||||
printf 'scum_client_credential: "%s"\n' "$auth_key"
|
||||
printf 'scum_client_name: "%s"\n' "$PROFILE_KEY"
|
||||
printf 'scum_client_version: "platform-build"\n'
|
||||
printf 'scum_client_machine_label: "managed-client"\n'
|
||||
printf 'ftp_provider: 3\n'
|
||||
} > config.yaml
|
||||
go mod download
|
||||
go build -trimpath -ldflags '-s -w' -o "/workspace/output/$OUTPUT_FILENAME" .
|
||||
cp config.yaml /workspace/output/config.yaml
|
||||
`
|
||||
|
||||
func (builder *DockerDistributionBuilder) containerArgs(input domain.DistributionBuildInput, sourceDir string, inputDir string, buildDir string, outputDir string, outputName string) []string {
|
||||
platformURL := strings.TrimSpace(input.PlatformURL)
|
||||
if platformURL == "" {
|
||||
platformURL = strings.TrimSpace(builder.config.PlatformURL)
|
||||
}
|
||||
return []string{
|
||||
"run", "--rm",
|
||||
"--pull", "never",
|
||||
"--read-only",
|
||||
"--tmpfs", "/tmp:rw,nosuid,size=2147483648",
|
||||
"-v", sourceDir + ":/workspace/source:ro",
|
||||
"-v", inputDir + ":/workspace/input:ro",
|
||||
"-v", buildDir + ":/workspace/build",
|
||||
"-v", outputDir + ":/workspace/output",
|
||||
"-e", "CGO_ENABLED=0",
|
||||
"-e", "GOOS=" + input.TargetOS,
|
||||
"-e", "GOARCH=" + input.TargetArch,
|
||||
"-e", "GOCACHE=/tmp/go-build",
|
||||
"-e", "GOMODCACHE=/tmp/go-mod",
|
||||
"-e", "COMPONENT_KIND=" + string(input.ComponentKind),
|
||||
"-e", "SERVER_INSTANCE_ID=" + input.ServerInstanceID,
|
||||
"-e", "PLUGIN_ID=" + input.PluginID,
|
||||
"-e", "RUN_ENDPOINT_ID=" + input.RunEndpointID,
|
||||
"-e", "PROFILE_KEY=" + input.ProfileKey,
|
||||
"-e", "TARGET_RELEASE=" + input.TargetRelease,
|
||||
"-e", "KEY_GENERATION=" + fmt.Sprint(input.KeyGeneration),
|
||||
"-e", "PLATFORM_URL=" + platformURL,
|
||||
"-e", "REPOSITORY_URL=" + input.RepositoryURL,
|
||||
"-e", "SOURCE_REVISION=" + input.SourceRevision,
|
||||
"-e", "OUTPUT_FILENAME=" + outputName,
|
||||
builder.config.Image,
|
||||
"/workspace/input/build.sh",
|
||||
}
|
||||
}
|
||||
|
||||
func packageClientManagerDistribution(packageFormat string, outputName string, binary []byte, configPayload []byte) ([]byte, error) {
|
||||
switch packageFormat {
|
||||
case "zip":
|
||||
return zipDistributionFiles(outputName, binary, configPayload)
|
||||
case "tar.gz":
|
||||
return tarGzipDistributionFiles(outputName, binary, configPayload)
|
||||
default:
|
||||
return nil, validationError("platform builder received an unsupported client-manager package format")
|
||||
}
|
||||
}
|
||||
|
||||
func zipDistributionFiles(outputName string, binary []byte, configPayload []byte) ([]byte, error) {
|
||||
var buffer bytes.Buffer
|
||||
writer := zip.NewWriter(&buffer)
|
||||
files := []struct {
|
||||
name string
|
||||
mode os.FileMode
|
||||
payload []byte
|
||||
}{{outputName, 0o755, binary}, {"config.yaml", 0o600, configPayload}}
|
||||
for _, file := range files {
|
||||
header := &zip.FileHeader{Name: file.name, Method: zip.Deflate}
|
||||
header.SetMode(file.mode)
|
||||
header.SetModTime(time.Date(1980, time.January, 1, 0, 0, 0, 0, time.UTC))
|
||||
entry, err := writer.CreateHeader(header)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := entry.Write(file.payload); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
if err := writer.Close(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
func tarGzipDistributionFiles(outputName string, binary []byte, configPayload []byte) ([]byte, error) {
|
||||
var buffer bytes.Buffer
|
||||
gzipWriter := gzip.NewWriter(&buffer)
|
||||
gzipWriter.Header.ModTime = time.Unix(0, 0).UTC()
|
||||
tarWriter := tar.NewWriter(gzipWriter)
|
||||
files := []struct {
|
||||
name string
|
||||
mode int64
|
||||
payload []byte
|
||||
}{{outputName, 0o755, binary}, {"config.yaml", 0o600, configPayload}}
|
||||
for _, file := range files {
|
||||
header := &tar.Header{Name: file.name, Mode: file.mode, Size: int64(len(file.payload)), ModTime: time.Unix(0, 0).UTC()}
|
||||
if err := tarWriter.WriteHeader(header); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := tarWriter.Write(file.payload); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
if err := tarWriter.Close(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := gzipWriter.Close(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return buffer.Bytes(), nil
|
||||
}
|
||||
|
||||
// safeBuilderFailure keeps host paths and secret values out of reported build
|
||||
// failures.
|
||||
func safeBuilderFailure(output []byte, sensitiveValues ...string) string {
|
||||
text := strings.TrimSpace(string(output))
|
||||
for _, sensitive := range sensitiveValues {
|
||||
if strings.TrimSpace(sensitive) != "" {
|
||||
text = strings.ReplaceAll(text, sensitive, "[redacted]")
|
||||
}
|
||||
}
|
||||
if text == "" {
|
||||
return "build command reported no diagnostic output"
|
||||
}
|
||||
lines := strings.Split(text, "\n")
|
||||
kept := make([]string, 0, len(lines))
|
||||
for index := len(lines) - 1; index >= 0 && len(kept) < 3; index-- {
|
||||
line := strings.TrimSpace(lines[index])
|
||||
if line == "" || strings.Contains(line, "/workspace/input") || strings.Contains(line, "auth-key") {
|
||||
continue
|
||||
}
|
||||
line = redactBuilderHostPaths(line)
|
||||
kept = append([]string{line}, kept...)
|
||||
}
|
||||
if len(kept) == 0 {
|
||||
return "build command reported no shareable diagnostic output"
|
||||
}
|
||||
joined := strings.Join(kept, "; ")
|
||||
if len(joined) > 400 {
|
||||
joined = joined[:400]
|
||||
}
|
||||
return joined
|
||||
}
|
||||
|
||||
func redactBuilderHostPaths(line string) string {
|
||||
fields := strings.Fields(line)
|
||||
for index, field := range fields {
|
||||
trimmed := strings.TrimLeft(field, "(\"'[")
|
||||
if strings.HasPrefix(trimmed, "/") && !strings.HasPrefix(trimmed, "/workspace/") {
|
||||
fields[index] = "[redacted-path]"
|
||||
}
|
||||
}
|
||||
return strings.Join(fields, " ")
|
||||
}
|
||||
|
||||
func builderJobFailureMessage(err error) string {
|
||||
if err == nil {
|
||||
return "platform builder failed"
|
||||
}
|
||||
message := strings.TrimSpace(err.Error())
|
||||
if message == "" {
|
||||
return "platform builder failed"
|
||||
}
|
||||
if len(message) > 400 {
|
||||
message = message[:400]
|
||||
}
|
||||
return message
|
||||
}
|
||||
@@ -0,0 +1,327 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"archive/tar"
|
||||
"archive/zip"
|
||||
"bytes"
|
||||
"compress/gzip"
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"browser.local/platform/domain"
|
||||
)
|
||||
|
||||
func TestPinnedBuilderImageRequiresExplicitTagOrDigest(t *testing.T) {
|
||||
validDigest := strings.Repeat("a", 64)
|
||||
for _, test := range []struct {
|
||||
image string
|
||||
want bool
|
||||
}{
|
||||
{image: "browser-platform-distribution-builder:1.0.0", want: true},
|
||||
{image: "registry.example.test:5000/builders/distribution:v1", want: true},
|
||||
{image: "registry.example.test/builders/distribution@sha256:" + validDigest, want: true},
|
||||
{image: "browser-platform-distribution-builder", want: false},
|
||||
{image: "browser-platform-distribution-builder:latest", want: false},
|
||||
{image: "registry.example.test/builders/distribution@sha256:", want: false},
|
||||
{image: "registry.example.test/builders/distribution@sha256:not-a-digest", want: false},
|
||||
} {
|
||||
t.Run(test.image, func(t *testing.T) {
|
||||
if got := pinnedBuilderImage(test.image); got != test.want {
|
||||
t.Fatalf("pinnedBuilderImage(%q) = %t, want %t", test.image, got, test.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDockerDistributionBuilderReadinessNamesPlatformBuilderFailures(t *testing.T) {
|
||||
sourceDir := createBuilderSource(t)
|
||||
workspaceDir := t.TempDir()
|
||||
for _, test := range []struct {
|
||||
name string
|
||||
config DockerDistributionBuilderConfig
|
||||
reason string
|
||||
}{
|
||||
{
|
||||
name: "floating image",
|
||||
config: DockerDistributionBuilderConfig{Image: "golang:latest", SourceDir: sourceDir, WorkspaceDir: workspaceDir},
|
||||
reason: "platform builder image must be pinned",
|
||||
},
|
||||
{
|
||||
name: "missing source",
|
||||
config: DockerDistributionBuilderConfig{Image: "builder:1.0.0", WorkspaceDir: workspaceDir},
|
||||
reason: "platform builder run source directory is not configured",
|
||||
},
|
||||
{
|
||||
name: "image unavailable",
|
||||
config: DockerDistributionBuilderConfig{
|
||||
Image: "builder:1.0.0",
|
||||
SourceDir: sourceDir,
|
||||
WorkspaceDir: workspaceDir,
|
||||
CommandRunner: func(_ context.Context, _ string, args ...string) ([]byte, error) {
|
||||
if len(args) > 0 && args[0] == "version" {
|
||||
return []byte("27.0.0"), nil
|
||||
}
|
||||
return nil, errors.New("image missing")
|
||||
},
|
||||
},
|
||||
reason: "platform builder image is unavailable",
|
||||
},
|
||||
{
|
||||
name: "container runtime unavailable",
|
||||
config: DockerDistributionBuilderConfig{
|
||||
Image: "builder:1.0.0",
|
||||
SourceDir: sourceDir,
|
||||
WorkspaceDir: workspaceDir,
|
||||
CommandRunner: func(context.Context, string, ...string) ([]byte, error) {
|
||||
return nil, errors.New("docker unavailable")
|
||||
},
|
||||
},
|
||||
reason: "platform builder container runtime is unavailable",
|
||||
},
|
||||
} {
|
||||
t.Run(test.name, func(t *testing.T) {
|
||||
ready, reason := NewDockerDistributionBuilder(test.config).Readiness()
|
||||
if ready || !strings.Contains(reason, test.reason) || !strings.Contains(reason, "platform builder") {
|
||||
t.Fatalf("expected platform builder readiness failure %q, ready=%t reason=%q", test.reason, ready, reason)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestDockerDistributionBuilderKeepsSecretInIsolatedInput(t *testing.T) {
|
||||
sourceDir := createBuilderSource(t)
|
||||
workspaceDir := t.TempDir()
|
||||
secret := "component-auth-key-that-must-not-leave-input"
|
||||
var dockerArgs []string
|
||||
builder := NewDockerDistributionBuilder(DockerDistributionBuilderConfig{
|
||||
DockerBinary: "docker-test",
|
||||
Image: "browser-platform-distribution-builder:1.0.0",
|
||||
SourceDir: sourceDir,
|
||||
WorkspaceDir: workspaceDir,
|
||||
CommandRunner: func(_ context.Context, name string, args ...string) ([]byte, error) {
|
||||
if name != "docker-test" {
|
||||
t.Fatalf("unexpected container runtime %q", name)
|
||||
}
|
||||
if len(args) > 0 && (args[0] == "version" || args[0] == "image") {
|
||||
return []byte("27.0.0"), nil
|
||||
}
|
||||
dockerArgs = append([]string(nil), args...)
|
||||
inputDir := builderMountHostPath(t, args, "/workspace/input:ro")
|
||||
outputDir := builderMountHostPath(t, args, "/workspace/output")
|
||||
authPath := filepath.Join(inputDir, "auth-key")
|
||||
payload, err := os.ReadFile(authPath)
|
||||
if err != nil {
|
||||
t.Fatalf("read per-job auth input: %v", err)
|
||||
}
|
||||
if string(payload) != secret {
|
||||
t.Fatalf("unexpected per-job auth input %q", payload)
|
||||
}
|
||||
info, err := os.Stat(authPath)
|
||||
if err != nil {
|
||||
t.Fatalf("stat per-job auth input: %v", err)
|
||||
}
|
||||
if info.Mode().Perm() != 0o600 {
|
||||
t.Fatalf("auth input mode = %o, want 600", info.Mode().Perm())
|
||||
}
|
||||
script, err := os.ReadFile(filepath.Join(inputDir, "build.sh"))
|
||||
if err != nil {
|
||||
t.Fatalf("read build script: %v", err)
|
||||
}
|
||||
if bytes.Contains(script, []byte(secret)) {
|
||||
t.Fatal("build script must not embed the component auth key")
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(outputDir, "run.exe"), []byte("compiled-run"), 0o700); err != nil {
|
||||
t.Fatalf("write fake build output: %v", err)
|
||||
}
|
||||
return nil, nil
|
||||
},
|
||||
})
|
||||
input := domain.DistributionBuildInput{
|
||||
JobID: "job/build:one",
|
||||
ComponentKind: domain.DistributionComponentRun,
|
||||
ServerInstanceID: "server-one",
|
||||
PluginID: "game.scum",
|
||||
RunEndpointID: "server-run-server-one",
|
||||
TargetOS: "windows",
|
||||
TargetArch: "amd64",
|
||||
TargetRelease: "release-one",
|
||||
PlatformURL: "https://platform.example.test",
|
||||
OutputFilename: "run.exe",
|
||||
KeyGeneration: 1,
|
||||
AuthKey: secret,
|
||||
}
|
||||
payload, err := builder.Build(input)
|
||||
if err != nil {
|
||||
t.Fatalf("build Run distribution: %v", err)
|
||||
}
|
||||
if string(payload) != "compiled-run" {
|
||||
t.Fatalf("unexpected built payload %q", payload)
|
||||
}
|
||||
joinedArgs := strings.Join(dockerArgs, "\x00")
|
||||
if strings.Contains(joinedArgs, secret) {
|
||||
t.Fatal("component auth key leaked into Docker arguments or environment")
|
||||
}
|
||||
for _, required := range []string{"--read-only", "--pull\x00never", sourceDir + ":/workspace/source:ro", ":/workspace/input:ro", ":/workspace/output"} {
|
||||
if !strings.Contains(joinedArgs, required) {
|
||||
t.Fatalf("Docker arguments do not contain required isolation %q: %q", required, joinedArgs)
|
||||
}
|
||||
}
|
||||
jobDir := filepath.Join(workspaceDir, "game.scum", "job-build-one")
|
||||
if _, err := os.Stat(jobDir); !errors.Is(err, os.ErrNotExist) {
|
||||
t.Fatalf("per-job workspace was not removed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDockerDistributionBuilderRedactsFailureAndTimeout(t *testing.T) {
|
||||
sourceDir := createBuilderSource(t)
|
||||
workspaceDir := t.TempDir()
|
||||
secret := "sensitive-component-key"
|
||||
input := domain.DistributionBuildInput{
|
||||
JobID: "job-redaction",
|
||||
ComponentKind: domain.DistributionComponentRun,
|
||||
PluginID: "game.scum",
|
||||
TargetOS: "linux",
|
||||
TargetArch: "amd64",
|
||||
OutputFilename: "run",
|
||||
AuthKey: secret,
|
||||
}
|
||||
builder := NewDockerDistributionBuilder(DockerDistributionBuilderConfig{
|
||||
Image: "builder:1.0.0",
|
||||
SourceDir: sourceDir,
|
||||
WorkspaceDir: workspaceDir,
|
||||
CommandRunner: func(_ context.Context, _ string, args ...string) ([]byte, error) {
|
||||
if len(args) > 0 && (args[0] == "version" || args[0] == "image") {
|
||||
return []byte("27.0.0"), nil
|
||||
}
|
||||
jobDir := filepath.Join(workspaceDir, "game.scum", "job-redaction")
|
||||
return []byte(secret + "\n" + sourceDir + "/go.mod: build failed\n" + jobDir + "/input/auth-key"), errors.New("exit 1")
|
||||
},
|
||||
})
|
||||
_, err := builder.Build(input)
|
||||
if err == nil {
|
||||
t.Fatal("expected failed builder command")
|
||||
}
|
||||
if strings.Contains(err.Error(), secret) || strings.Contains(err.Error(), sourceDir) || strings.Contains(err.Error(), workspaceDir) || strings.Contains(err.Error(), "auth-key") {
|
||||
t.Fatalf("builder failure leaked secret or host path: %v", err)
|
||||
}
|
||||
|
||||
timeoutBuilder := NewDockerDistributionBuilder(DockerDistributionBuilderConfig{
|
||||
Image: "builder:1.0.0",
|
||||
SourceDir: sourceDir,
|
||||
WorkspaceDir: workspaceDir,
|
||||
Timeout: 5 * time.Millisecond,
|
||||
CommandRunner: func(ctx context.Context, _ string, args ...string) ([]byte, error) {
|
||||
if len(args) > 0 && (args[0] == "version" || args[0] == "image") {
|
||||
return []byte("27.0.0"), nil
|
||||
}
|
||||
<-ctx.Done()
|
||||
return nil, ctx.Err()
|
||||
},
|
||||
})
|
||||
_, err = timeoutBuilder.Build(input)
|
||||
if err == nil || !strings.Contains(err.Error(), "platform builder timed out") {
|
||||
t.Fatalf("expected bounded platform builder timeout, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPackageClientManagerDistributionProducesProtectedArchives(t *testing.T) {
|
||||
for _, packageFormat := range []string{"zip", "tar.gz"} {
|
||||
t.Run(packageFormat, func(t *testing.T) {
|
||||
payload, err := packageClientManagerDistribution(packageFormat, "client-manager.exe", []byte("binary"), []byte("credential: protected"))
|
||||
if err != nil {
|
||||
t.Fatalf("package client manager: %v", err)
|
||||
}
|
||||
files := readDistributionArchive(t, packageFormat, payload)
|
||||
if string(files["client-manager.exe"].payload) != "binary" || files["client-manager.exe"].mode.Perm() != 0o755 {
|
||||
t.Fatalf("unexpected executable archive entry: %+v", files["client-manager.exe"])
|
||||
}
|
||||
if string(files["config.yaml"].payload) != "credential: protected" || files["config.yaml"].mode.Perm() != 0o600 {
|
||||
t.Fatalf("unexpected config archive entry: %+v", files["config.yaml"])
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
type distributionArchiveFile struct {
|
||||
mode os.FileMode
|
||||
payload []byte
|
||||
}
|
||||
|
||||
func createBuilderSource(t *testing.T) string {
|
||||
t.Helper()
|
||||
sourceDir := t.TempDir()
|
||||
if err := os.WriteFile(filepath.Join(sourceDir, "go.mod"), []byte("module browser.local/run\n"), 0o600); err != nil {
|
||||
t.Fatalf("write run source go.mod: %v", err)
|
||||
}
|
||||
return sourceDir
|
||||
}
|
||||
|
||||
func builderMountHostPath(t *testing.T, args []string, containerSuffix string) string {
|
||||
t.Helper()
|
||||
for index := 0; index+1 < len(args); index++ {
|
||||
if args[index] != "-v" || !strings.HasSuffix(args[index+1], ":"+containerSuffix) {
|
||||
continue
|
||||
}
|
||||
return strings.TrimSuffix(args[index+1], ":"+containerSuffix)
|
||||
}
|
||||
t.Fatalf("Docker arguments do not mount %s: %+v", containerSuffix, args)
|
||||
return ""
|
||||
}
|
||||
|
||||
func readDistributionArchive(t *testing.T, packageFormat string, payload []byte) map[string]distributionArchiveFile {
|
||||
t.Helper()
|
||||
files := map[string]distributionArchiveFile{}
|
||||
switch packageFormat {
|
||||
case "zip":
|
||||
reader, err := zip.NewReader(bytes.NewReader(payload), int64(len(payload)))
|
||||
if err != nil {
|
||||
t.Fatalf("open zip: %v", err)
|
||||
}
|
||||
for _, entry := range reader.File {
|
||||
body, err := entry.Open()
|
||||
if err != nil {
|
||||
t.Fatalf("open zip entry %s: %v", entry.Name, err)
|
||||
}
|
||||
content, err := io.ReadAll(body)
|
||||
if closeErr := body.Close(); err == nil {
|
||||
err = closeErr
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("read zip entry %s: %v", entry.Name, err)
|
||||
}
|
||||
files[entry.Name] = distributionArchiveFile{mode: entry.Mode(), payload: content}
|
||||
}
|
||||
case "tar.gz":
|
||||
gzipReader, err := gzip.NewReader(bytes.NewReader(payload))
|
||||
if err != nil {
|
||||
t.Fatalf("open gzip: %v", err)
|
||||
}
|
||||
tarReader := tar.NewReader(gzipReader)
|
||||
for {
|
||||
header, err := tarReader.Next()
|
||||
if errors.Is(err, io.EOF) {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("read tar header: %v", err)
|
||||
}
|
||||
content, err := io.ReadAll(tarReader)
|
||||
if err != nil {
|
||||
t.Fatalf("read tar entry %s: %v", header.Name, err)
|
||||
}
|
||||
files[header.Name] = distributionArchiveFile{mode: os.FileMode(header.Mode), payload: content}
|
||||
}
|
||||
if err := gzipReader.Close(); err != nil {
|
||||
t.Fatalf("close gzip: %v", err)
|
||||
}
|
||||
default:
|
||||
t.Fatalf("unsupported archive format %q", packageFormat)
|
||||
}
|
||||
return files
|
||||
}
|
||||
@@ -47,19 +47,12 @@ func (svc *CoreService) GenerateRunDistributionForSession(sessionID string, requ
|
||||
if err := svc.promoteLegacyRunBinding(&instance); err != nil {
|
||||
return domain.RunDistribution{}, err
|
||||
}
|
||||
builderEndpointID := instance.RunEndpointID
|
||||
if strings.TrimSpace(instance.DeploymentTargetID) != "" {
|
||||
builderEndpointID = instance.DeploymentTargetID
|
||||
}
|
||||
endpoint, err := svc.store.RunEndpoints().Get(builderEndpointID)
|
||||
if err != nil {
|
||||
return domain.RunDistribution{}, err
|
||||
}
|
||||
if err := svc.validateRunnableEndpoint(endpoint, domain.JobCapabilityDistributionBuild); err != nil {
|
||||
return domain.RunDistribution{}, err
|
||||
if ready, reason := svc.distributionBuilderReadiness(); !ready {
|
||||
_ = svc.recordAuditEvent(user.ID, "run.generate.denied", "server-instance", instance.ID, domain.AuditResultDenied, reason)
|
||||
return domain.RunDistribution{}, validationError(reason)
|
||||
}
|
||||
|
||||
key, plainKey, err := svc.ensureActiveComponentKey(instance.ID, domain.DistributionComponentRun, "")
|
||||
key, _, err := svc.ensureActiveComponentKey(instance.ID, domain.DistributionComponentRun, "")
|
||||
if err != nil {
|
||||
return domain.RunDistribution{}, err
|
||||
}
|
||||
@@ -70,7 +63,6 @@ func (svc *CoreService) GenerateRunDistributionForSession(sessionID string, requ
|
||||
return domain.RunDistribution{}, err
|
||||
}
|
||||
|
||||
_ = plainKey
|
||||
artifactID := artifactIDForDistribution(distributionID + "-binary")
|
||||
buildJobID := jobIDFromParts("job-distribution-build", instance.ID, distributionID)
|
||||
stamp := svc.now()
|
||||
@@ -106,7 +98,7 @@ func (svc *CoreService) GenerateRunDistributionForSession(sessionID string, requ
|
||||
job, err := svc.CreateJob(domain.Job{
|
||||
ID: buildJobID,
|
||||
ServerInstanceID: instance.ID,
|
||||
RunEndpointID: builderEndpointID,
|
||||
RunEndpointID: platformDistributionBuilderEndpointID,
|
||||
Capability: domain.JobCapabilityDistributionBuild,
|
||||
TargetKey: "distribution/run",
|
||||
InputRef: "input://distribution-build/" + distribution.ID,
|
||||
@@ -122,26 +114,29 @@ func (svc *CoreService) GenerateRunDistributionForSession(sessionID string, requ
|
||||
if job.ID != buildJobID || job.Capability != domain.JobCapabilityDistributionBuild {
|
||||
return domain.RunDistribution{}, validationError("distribution build idempotency key conflicts with another job")
|
||||
}
|
||||
if err := svc.recordAuditEvent(user.ID, "run.generate", "server-instance", instance.ID, domain.AuditResultQueued, "queued run binary build job with redacted runtime key ref"); err != nil {
|
||||
svc.enqueueDistributionBuild(job)
|
||||
if err := svc.recordAuditEvent(user.ID, "run.generate", "server-instance", instance.ID, domain.AuditResultQueued, "queued run binary build job in the platform builder with redacted runtime key ref"); err != nil {
|
||||
return domain.RunDistribution{}, err
|
||||
}
|
||||
return domain.CopyRunDistribution(distribution), nil
|
||||
}
|
||||
|
||||
// promoteLegacyRunBinding reserves a server-scoped endpoint for a generated
|
||||
// Run before a legacy server first requests a distribution. Its existing
|
||||
// endpoint remains the trusted build target; reusing it in the package would
|
||||
// allow the generated Run to replace the builder registration.
|
||||
// promoteLegacyRunBinding reserves the server-scoped endpoint used by a
|
||||
// generated Run. A legacy shared endpoint remains an optional deployment target
|
||||
// for non-build workflows; distribution builds are always platform-owned.
|
||||
func (svc *CoreService) promoteLegacyRunBinding(instance *domain.ServerInstance) error {
|
||||
if instance == nil || strings.TrimSpace(instance.DeploymentTargetID) != "" || (instance.State != domain.ServerInstanceStateDraft && instance.State != domain.ServerInstanceStateFailed) {
|
||||
return nil
|
||||
}
|
||||
builderEndpointID := strings.TrimSpace(instance.RunEndpointID)
|
||||
if builderEndpointID == "" {
|
||||
return validationError("legacy Run generation requires a build target endpoint")
|
||||
currentEndpointID := strings.TrimSpace(instance.RunEndpointID)
|
||||
dedicatedEndpointID := dedicatedRunEndpointID(instance.ID)
|
||||
if currentEndpointID == dedicatedEndpointID {
|
||||
return nil
|
||||
}
|
||||
instance.DeploymentTargetID = builderEndpointID
|
||||
instance.RunEndpointID = dedicatedRunEndpointID(instance.ID)
|
||||
if currentEndpointID != "" {
|
||||
instance.DeploymentTargetID = currentEndpointID
|
||||
}
|
||||
instance.RunEndpointID = dedicatedEndpointID
|
||||
instance.UpdatedAt = svc.now()
|
||||
if err := validator.ValidateServerInstance(*instance); err != nil {
|
||||
return err
|
||||
@@ -191,15 +186,12 @@ func (svc *CoreService) GenerateClientManagerDistributionForSession(sessionID st
|
||||
if err := svc.requireCompleteRuntimeBindings(user.ID, instance.ID, "client-manager.build.denied"); err != nil {
|
||||
return domain.ClientManagerDistribution{}, err
|
||||
}
|
||||
endpoint, err := svc.store.RunEndpoints().Get(instance.RunEndpointID)
|
||||
if err != nil {
|
||||
return domain.ClientManagerDistribution{}, err
|
||||
}
|
||||
if err := svc.validateRunnableEndpoint(endpoint, domain.JobCapabilityDistributionBuild); err != nil {
|
||||
return domain.ClientManagerDistribution{}, err
|
||||
if ready, reason := svc.distributionBuilderReadiness(); !ready {
|
||||
_ = svc.recordAuditEvent(user.ID, "client-manager.build.denied", "server-instance", instance.ID, domain.AuditResultDenied, reason)
|
||||
return domain.ClientManagerDistribution{}, validationError(reason)
|
||||
}
|
||||
|
||||
key, plainKey, err := svc.ensureActiveComponentKey(instance.ID, domain.DistributionComponentClientManager, request.ProfileKey)
|
||||
key, _, err := svc.ensureActiveComponentKey(instance.ID, domain.DistributionComponentClientManager, request.ProfileKey)
|
||||
if err != nil {
|
||||
return domain.ClientManagerDistribution{}, err
|
||||
}
|
||||
@@ -210,7 +202,6 @@ func (svc *CoreService) GenerateClientManagerDistributionForSession(sessionID st
|
||||
return domain.ClientManagerDistribution{}, err
|
||||
}
|
||||
|
||||
_ = plainKey
|
||||
artifactID := artifactIDForDistribution(clientDistributionID + "-binary")
|
||||
stamp := svc.now()
|
||||
buildJobID := jobIDFromParts("job-distribution-build", instance.ID, clientDistributionID)
|
||||
@@ -283,7 +274,7 @@ func (svc *CoreService) GenerateClientManagerDistributionForSession(sessionID st
|
||||
job, err := svc.CreateJob(domain.Job{
|
||||
ID: buildJobID,
|
||||
ServerInstanceID: instance.ID,
|
||||
RunEndpointID: instance.RunEndpointID,
|
||||
RunEndpointID: platformDistributionBuilderEndpointID,
|
||||
Capability: domain.JobCapabilityDistributionBuild,
|
||||
TargetKey: "distribution/client-manager/" + request.ProfileKey,
|
||||
InputRef: "input://distribution-build/" + distribution.ID,
|
||||
@@ -303,7 +294,8 @@ func (svc *CoreService) GenerateClientManagerDistributionForSession(sessionID st
|
||||
if job.ID != buildJobID || job.Capability != domain.JobCapabilityDistributionBuild {
|
||||
return domain.ClientManagerDistribution{}, validationError("distribution build idempotency key conflicts with another job")
|
||||
}
|
||||
if err := svc.recordAuditEvent(user.ID, "client-manager.build", "server-instance", instance.ID, domain.AuditResultQueued, "queued client-manager source build with redacted runtime key ref"); err != nil {
|
||||
svc.enqueueDistributionBuild(job)
|
||||
if err := svc.recordAuditEvent(user.ID, "client-manager.build", "server-instance", instance.ID, domain.AuditResultQueued, "queued client-manager source build in the platform builder with redacted runtime key ref"); err != nil {
|
||||
return domain.ClientManagerDistribution{}, err
|
||||
}
|
||||
return domain.CopyClientManagerDistribution(distribution), nil
|
||||
@@ -467,7 +459,7 @@ func (svc *CoreService) GetServerRuntimeActionsForSession(sessionID string, serv
|
||||
if endpointErr != nil && strings.TrimSpace(instance.DeploymentTargetID) != "" {
|
||||
endpoint, endpointErr = svc.store.RunEndpoints().Get(instance.DeploymentTargetID)
|
||||
}
|
||||
if endpointErr != nil {
|
||||
if endpointErr != nil && !errors.Is(endpointErr, repo.ErrNotFound) {
|
||||
return domain.ServerRuntimeActions{}, endpointErr
|
||||
}
|
||||
hasAvailableRunPackage := false
|
||||
@@ -493,6 +485,7 @@ func (svc *CoreService) GetServerRuntimeActionsForSession(sessionID string, serv
|
||||
}
|
||||
}
|
||||
bindingsComplete, bindingReason := svc.runtimeBindingReadiness(instance.ID)
|
||||
builderReady, builderReason := svc.distributionBuilderReadiness()
|
||||
dependencyPermissionDeclared := pluginDeclares(plugin, "server.dependencies.manage")
|
||||
actions := domain.ServerRuntimeActions{
|
||||
ServerInstanceID: instance.ID,
|
||||
@@ -505,11 +498,11 @@ func (svc *CoreService) GetServerRuntimeActionsForSession(sessionID string, serv
|
||||
return domain.RunEndpointStatusOffline
|
||||
}(),
|
||||
Actions: []domain.ServerRuntimeAction{
|
||||
runtimeAction("generate-run", "Generate run", pluginDeclares(plugin, "server.run.distribution") && svc.endpointSupports(endpoint, domain.JobCapabilityDistributionBuild) && bindingsComplete, fallbackReason(!pluginDeclares(plugin, "server.run.distribution") || !svc.endpointSupports(endpoint, domain.JobCapabilityDistributionBuild), "run endpoint cannot build distributions", bindingReason)),
|
||||
runtimeAction("generate-run", "Generate run", pluginDeclares(plugin, "server.run.distribution") && builderReady && bindingsComplete, fallbackReason(!pluginDeclares(plugin, "server.run.distribution"), "plugin permission is not declared", fallbackReason(!builderReady, builderReason, bindingReason))),
|
||||
runtimeAction("download-run", "Download run", hasAvailableRunPackage, "run package has not been generated"),
|
||||
runtimeAction("push-run-update", "Push run update", runRegistered && pluginDeclares(plugin, "server.run.distribution") && svc.endpointSupports(endpoint, domain.JobCapabilityRunSelfUpdate) && bindingsComplete, fallbackReason(!runRegistered, "dedicated Run has not registered", fallbackReason(!pluginDeclares(plugin, "server.run.distribution") || !svc.endpointSupports(endpoint, domain.JobCapabilityRunSelfUpdate), "run endpoint cannot self-update", bindingReason))),
|
||||
runtimeAction("reset-run-key", "Reset run key", pluginDeclares(plugin, "server.run.distribution"), "plugin permission is not declared"),
|
||||
runtimeAction("generate-client-manager", "Generate client manager", pluginDeclares(plugin, "server.client-manager.manage") && svc.endpointSupports(endpoint, domain.JobCapabilityDistributionBuild) && bindingsComplete, fallbackReason(!pluginDeclares(plugin, "server.client-manager.manage") || !svc.endpointSupports(endpoint, domain.JobCapabilityDistributionBuild), "run endpoint cannot build distributions", bindingReason)),
|
||||
runtimeAction("generate-client-manager", "Generate client manager", pluginDeclares(plugin, "server.client-manager.manage") && builderReady && bindingsComplete, fallbackReason(!pluginDeclares(plugin, "server.client-manager.manage"), "client-manager permission is not declared", fallbackReason(!builderReady, builderReason, bindingReason))),
|
||||
runtimeAction("download-client-manager", "Download client manager", hasAvailableClientPackage, "client-manager package has not been generated"),
|
||||
runtimeAction("reset-client-manager-key", "Reset client-manager key", pluginDeclares(plugin, "server.client-manager.manage"), "client-manager permission is not declared"),
|
||||
runtimeAction("dependencies-check", "Check dependencies", runRegistered && dependencyPermissionDeclared && svc.endpointSupports(endpoint, domain.JobCapabilityDependenciesCheck) && bindingsComplete, fallbackReason(!runRegistered, "dedicated Run has not registered", fallbackReason(!dependencyPermissionDeclared, "plugin permission is not declared", fallbackReason(!svc.endpointSupports(endpoint, domain.JobCapabilityDependenciesCheck), "run endpoint cannot check dependencies", bindingReason)))),
|
||||
|
||||
@@ -100,7 +100,7 @@ func TestCoreServiceGeneratesRunDistributionWithEncryptedSingletonKey(t *testing
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceBuildsDedicatedRunOnDeploymentTarget(t *testing.T) {
|
||||
func TestCoreServiceBuildsDedicatedRunInPlatformBuilder(t *testing.T) {
|
||||
svc, session, instance := newDistributionTestFixture(t)
|
||||
targetID := instance.RunEndpointID
|
||||
instance.State = domain.ServerInstanceStateDraft
|
||||
@@ -115,14 +115,14 @@ func TestCoreServiceBuildsDedicatedRunOnDeploymentTarget(t *testing.T) {
|
||||
t.Fatalf("generate dedicated Run: %v", err)
|
||||
}
|
||||
job, err := svc.GetJob(distribution.BuildJobID)
|
||||
if err != nil || job.RunEndpointID != targetID || distribution.RunEndpointID != instance.RunEndpointID {
|
||||
t.Fatalf("expected build on target %q for dedicated Run %q, job=%+v distribution=%+v err=%v", targetID, instance.RunEndpointID, job, distribution, err)
|
||||
if err != nil || job.RunEndpointID != platformDistributionBuilderEndpointID || distribution.RunEndpointID != instance.RunEndpointID {
|
||||
t.Fatalf("expected platform build for dedicated Run %q, job=%+v distribution=%+v err=%v", instance.RunEndpointID, job, distribution, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServicePromotesLegacyRunBindingBeforeDistributionBuild(t *testing.T) {
|
||||
svc, session, instance := newDistributionTestFixture(t)
|
||||
builderID := instance.RunEndpointID
|
||||
legacyEndpointID := instance.RunEndpointID
|
||||
instance.State = domain.ServerInstanceStateFailed
|
||||
if err := svc.store.ServerInstances().Update(instance); err != nil {
|
||||
t.Fatalf("mark legacy server failed: %v", err)
|
||||
@@ -136,36 +136,37 @@ func TestCoreServicePromotesLegacyRunBindingBeforeDistributionBuild(t *testing.T
|
||||
if err != nil {
|
||||
t.Fatalf("get migrated server: %v", err)
|
||||
}
|
||||
if migrated.DeploymentTargetID != builderID || migrated.RunEndpointID != "server-run-"+instance.ID {
|
||||
if migrated.DeploymentTargetID != legacyEndpointID || migrated.RunEndpointID != "server-run-"+instance.ID {
|
||||
t.Fatalf("expected legacy binding promotion, got %+v", migrated)
|
||||
}
|
||||
job, err := svc.GetJob(distribution.BuildJobID)
|
||||
if err != nil || job.RunEndpointID != builderID || distribution.RunEndpointID != migrated.RunEndpointID {
|
||||
t.Fatalf("expected build target %q and dedicated package endpoint %q, job=%+v distribution=%+v err=%v", builderID, migrated.RunEndpointID, job, distribution, err)
|
||||
if err != nil || job.RunEndpointID != platformDistributionBuilderEndpointID || distribution.RunEndpointID != migrated.RunEndpointID {
|
||||
t.Fatalf("expected platform build and dedicated package endpoint %q, job=%+v distribution=%+v err=%v", migrated.RunEndpointID, job, distribution, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceRejectsDistributionBuildForStaleRunEndpoint(t *testing.T) {
|
||||
func TestCoreServiceDistributionBuildIgnoresStaleRunEndpoint(t *testing.T) {
|
||||
svc, session, instance := newDistributionTestFixture(t)
|
||||
svc.now = func() time.Time { return fixedTime.Add(capacityHeartbeatStaleAfter + time.Second) }
|
||||
|
||||
_, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{
|
||||
distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{
|
||||
ServerInstanceID: instance.ID,
|
||||
TargetOS: "linux",
|
||||
TargetArch: "amd64",
|
||||
IdempotencyKey: "idem-stale-run-endpoint",
|
||||
})
|
||||
if err == nil || !strings.Contains(err.Error(), "heartbeat is stale") {
|
||||
t.Fatalf("expected stale Run endpoint rejection, got %v", err)
|
||||
if err != nil {
|
||||
t.Fatalf("platform build must ignore stale Run endpoint: %v", err)
|
||||
}
|
||||
completeDistributionBuild(t, svc, distribution, nil)
|
||||
|
||||
actions, err := svc.GetServerRuntimeActionsForSession(session, instance.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get runtime actions: %v", err)
|
||||
}
|
||||
for _, action := range actions.Actions {
|
||||
if action.Key == "generate-run" && action.Available {
|
||||
t.Fatalf("stale Run endpoint must not expose generate-run as available: %+v", action)
|
||||
if action.Key == "generate-run" && !action.Available {
|
||||
t.Fatalf("stale Run endpoint must not gate platform build availability: %+v", action)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -206,14 +207,13 @@ func TestCoreServiceRuntimeActionsGateDependenciesOnPluginPermission(t *testing.
|
||||
}
|
||||
}
|
||||
|
||||
func TestCoreServiceDistributionBuildRejectsPrematureSuccessAndCanRetryAfterUpload(t *testing.T) {
|
||||
t.Setenv("PLATFORM_RUN_RELEASE_URL", "https://scum.npc0.com")
|
||||
func TestCoreServiceMachineEndpointCannotClaimPlatformDistributionBuild(t *testing.T) {
|
||||
svc, session, instance := newDistributionTestFixture(t)
|
||||
distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{
|
||||
ServerInstanceID: instance.ID,
|
||||
TargetOS: "linux",
|
||||
TargetArch: "amd64",
|
||||
IdempotencyKey: "idem-premature-result",
|
||||
IdempotencyKey: "idem-platform-build-ownership",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("generate run distribution: %v", err)
|
||||
@@ -223,7 +223,7 @@ func TestCoreServiceDistributionBuildRejectsPrematureSuccessAndCanRetryAfterUplo
|
||||
helloRequest.CapabilityReport.Capabilities = append(helloRequest.CapabilityReport.Capabilities, domain.JobCapabilityDistributionBuild)
|
||||
hello, err := svc.RegisterRunHello(helloRequest)
|
||||
if err != nil {
|
||||
t.Fatalf("register build worker: %v", err)
|
||||
t.Fatalf("register machine endpoint: %v", err)
|
||||
}
|
||||
claim, err := svc.ClaimRunJob(domain.RunJobClaim{
|
||||
RunEndpointID: instance.RunEndpointID,
|
||||
@@ -231,50 +231,21 @@ func TestCoreServiceDistributionBuildRejectsPrematureSuccessAndCanRetryAfterUplo
|
||||
Capabilities: []string{domain.JobCapabilityDistributionBuild},
|
||||
Capacity: domain.RunCapacity{MaxJobs: 1},
|
||||
})
|
||||
if err != nil || !claim.HasJob || claim.Job.JobID != distribution.BuildJobID {
|
||||
t.Fatalf("claim distribution build job: claim=%+v err=%v", claim, err)
|
||||
}
|
||||
buildInput, err := svc.GetDistributionBuildInput(domain.DistributionBuildInputRequest{
|
||||
RunEndpointID: claim.Job.RunEndpointID,
|
||||
SessionToken: hello.SessionToken,
|
||||
JobID: claim.Job.JobID,
|
||||
LeaseToken: claim.Job.LeaseToken,
|
||||
Attempt: claim.Job.Attempt,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("get distribution build input: %v", err)
|
||||
t.Fatalf("claim machine-side build: %v", err)
|
||||
}
|
||||
if buildInput.PlatformURL != "https://scum.npc0.com" || buildInput.PackageFormat != "raw-executable" || buildInput.AuthKey == "" || len(buildInput.AuthKey) < 80 {
|
||||
t.Fatalf("expected raw executable build input with release URL and long key, got %+v", buildInput)
|
||||
}
|
||||
result := domain.RunJobResult{
|
||||
RunEndpointID: instance.RunEndpointID,
|
||||
SessionToken: hello.SessionToken,
|
||||
JobID: claim.Job.JobID,
|
||||
LeaseToken: claim.Job.LeaseToken,
|
||||
Attempt: claim.Job.Attempt,
|
||||
State: domain.JobStateSucceeded,
|
||||
Progress: domain.RunJobProgressReport{Percent: 100, Message: "package_finalize: done"},
|
||||
ResultRef: "artifact://" + distribution.ArtifactID,
|
||||
Message: "done",
|
||||
}
|
||||
if _, err := svc.CompleteRunJob(result); err == nil {
|
||||
t.Fatal("expected premature success without uploaded artifact to be rejected")
|
||||
}
|
||||
stored, err := svc.GetJob(distribution.BuildJobID)
|
||||
if err != nil || stored.State != domain.JobStateAccepted {
|
||||
t.Fatalf("premature success must not make the job terminal, job=%+v err=%v", stored, err)
|
||||
if claim.HasJob {
|
||||
t.Fatalf("machine-side endpoint must not receive platform build job: %+v", claim)
|
||||
}
|
||||
|
||||
if _, err := svc.createPlatformArtifactPayload(distribution.ArtifactID, domain.ArtifactOwnerKindJob, distribution.BuildJobID, []byte("actual compiled archive")); err != nil {
|
||||
t.Fatalf("publish uploaded build output: %v", err)
|
||||
distribution = completeDistributionBuild(t, svc, distribution, nil)
|
||||
job, err := svc.GetJob(distribution.BuildJobID)
|
||||
if err != nil || job.State != domain.JobStateSucceeded || job.RunEndpointID != platformDistributionBuilderEndpointID {
|
||||
t.Fatalf("expected completed platform build job, job=%+v err=%v", job, err)
|
||||
}
|
||||
if _, err := svc.CompleteRunJob(result); err != nil {
|
||||
t.Fatalf("retry success after artifact upload: %v", err)
|
||||
}
|
||||
stored, err = svc.GetJob(distribution.BuildJobID)
|
||||
if err != nil || stored.State != domain.JobStateSucceeded {
|
||||
t.Fatalf("expected terminal success after upload, job=%+v err=%v", stored, err)
|
||||
artifact, err := svc.GetArtifact(distribution.ArtifactID)
|
||||
if err != nil || artifact.OwnerKind != domain.ArtifactOwnerKindJob || artifact.OwnerID != job.ID {
|
||||
t.Fatalf("expected build job-owned artifact, artifact=%+v err=%v", artifact, err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -348,6 +319,8 @@ func TestCoreServiceRunDistributionRetryReusesPartialArtifact(t *testing.T) {
|
||||
|
||||
func TestCoreServicePushRunUpdateReusesExistingUpdateJob(t *testing.T) {
|
||||
svc, session, instance := newDistributionTestFixture(t)
|
||||
payload := []byte("compiled run archive")
|
||||
svc.ConfigureDistributionBuilder(staticDistributionBuilder{payload: payload})
|
||||
distribution, err := svc.GenerateRunDistributionForSession(session, domain.RunDistributionGenerateRequest{
|
||||
ServerInstanceID: instance.ID,
|
||||
TargetOS: "linux",
|
||||
@@ -357,7 +330,7 @@ func TestCoreServicePushRunUpdateReusesExistingUpdateJob(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatalf("generate run distribution: %v", err)
|
||||
}
|
||||
distribution = completeDistributionBuild(t, svc, distribution, []byte("compiled run archive"))
|
||||
distribution = completeDistributionBuild(t, svc, distribution, payload)
|
||||
|
||||
request := domain.RunUpdateRequest{
|
||||
ServerInstanceID: instance.ID,
|
||||
@@ -624,7 +597,6 @@ func newDistributionTestFixture(t *testing.T) (*CoreService, string, domain.Serv
|
||||
t.Fatalf("update plugin fixture: %v", err)
|
||||
}
|
||||
endpoint.Capabilities = append(endpoint.Capabilities,
|
||||
domain.JobCapabilityDistributionBuild,
|
||||
domain.JobCapabilityRunSelfUpdate,
|
||||
domain.JobCapabilityDependenciesCheck,
|
||||
domain.JobCapabilityDependenciesInstall,
|
||||
@@ -704,58 +676,44 @@ func readGeneratedPackageConfig(t *testing.T, svc *CoreService, session string,
|
||||
return generatedPackageConfig{}
|
||||
}
|
||||
|
||||
func completeDistributionBuild(t *testing.T, svc *CoreService, distribution domain.RunDistribution, payload []byte) domain.RunDistribution {
|
||||
func completeDistributionBuild(t *testing.T, svc *CoreService, distribution domain.RunDistribution, _ []byte) domain.RunDistribution {
|
||||
t.Helper()
|
||||
artifact, err := svc.createPlatformArtifactPayload(distribution.ArtifactID, domain.ArtifactOwnerKindJob, distribution.BuildJobID, payload)
|
||||
if err != nil {
|
||||
t.Fatalf("publish run build artifact: %v", err)
|
||||
deadline := time.Now().Add(time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
updated, err := svc.store.RunDistributions().Get(distribution.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get run distribution: %v", err)
|
||||
}
|
||||
if updated.Status == domain.DistributionStatusAvailable {
|
||||
return updated
|
||||
}
|
||||
if updated.Status == domain.DistributionStatusFailed {
|
||||
t.Fatalf("platform run build failed: %+v", updated)
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
job, err := svc.GetJob(distribution.BuildJobID)
|
||||
if err != nil {
|
||||
t.Fatalf("get run build job: %v", err)
|
||||
}
|
||||
job.State = domain.JobStateSucceeded
|
||||
job.Progress = domain.JobProgress{Percent: 100, Message: "package_finalize: build artifact available"}
|
||||
job.ResultRef = "artifact://" + artifact.ID
|
||||
job.UpdatedAt = svc.now()
|
||||
if err := svc.store.Jobs().Update(job); err != nil {
|
||||
t.Fatalf("update run build job: %v", err)
|
||||
}
|
||||
if err := svc.projectDistributionBuildResult(job, svc.now()); err != nil {
|
||||
t.Fatalf("project run build: %v", err)
|
||||
}
|
||||
updated, err := svc.store.RunDistributions().Get(distribution.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get completed run distribution: %v", err)
|
||||
}
|
||||
return updated
|
||||
t.Fatalf("platform run build did not complete")
|
||||
return domain.RunDistribution{}
|
||||
}
|
||||
|
||||
func completeClientDistributionBuild(t *testing.T, svc *CoreService, distribution domain.ClientManagerDistribution, payload []byte) domain.ClientManagerDistribution {
|
||||
func completeClientDistributionBuild(t *testing.T, svc *CoreService, distribution domain.ClientManagerDistribution, _ []byte) domain.ClientManagerDistribution {
|
||||
t.Helper()
|
||||
artifact, err := svc.createPlatformArtifactPayload(distribution.ArtifactID, domain.ArtifactOwnerKindJob, distribution.BuildJobID, payload)
|
||||
if err != nil {
|
||||
t.Fatalf("publish client build artifact: %v", err)
|
||||
deadline := time.Now().Add(time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
updated, err := svc.store.ClientManagerDistributions().Get(distribution.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get client distribution: %v", err)
|
||||
}
|
||||
if updated.Status == domain.DistributionStatusAvailable {
|
||||
return updated
|
||||
}
|
||||
if updated.Status == domain.DistributionStatusFailed {
|
||||
t.Fatalf("platform client build failed: %+v", updated)
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
job, err := svc.GetJob(distribution.BuildJobID)
|
||||
if err != nil {
|
||||
t.Fatalf("get client build job: %v", err)
|
||||
}
|
||||
job.State = domain.JobStateSucceeded
|
||||
job.Progress = domain.JobProgress{Percent: 100, Message: "package_finalize: build artifact available"}
|
||||
job.ResultRef = "artifact://" + artifact.ID
|
||||
job.UpdatedAt = svc.now()
|
||||
if err := svc.store.Jobs().Update(job); err != nil {
|
||||
t.Fatalf("update client build job: %v", err)
|
||||
}
|
||||
if err := svc.projectDistributionBuildResult(job, svc.now()); err != nil {
|
||||
t.Fatalf("project client build: %v", err)
|
||||
}
|
||||
updated, err := svc.store.ClientManagerDistributions().Get(distribution.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("get completed client distribution: %v", err)
|
||||
}
|
||||
return updated
|
||||
t.Fatalf("platform client build did not complete")
|
||||
return domain.ClientManagerDistribution{}
|
||||
}
|
||||
|
||||
func TestCoreServiceDeniesRunDistributionWithoutPluginDeclaration(t *testing.T) {
|
||||
|
||||
@@ -249,6 +249,9 @@ type CoreService struct {
|
||||
aiProviderClient AIProviderClient
|
||||
secretEnvelope SecretEnvelope
|
||||
networkFingerprintKey []byte
|
||||
distributionBuilder DistributionBuilder
|
||||
distributionBuildMu sync.Mutex
|
||||
distributionBuilds map[string]struct{}
|
||||
}
|
||||
|
||||
var _ Core = (*CoreService)(nil)
|
||||
@@ -284,10 +287,34 @@ func newCoreServiceWithLogStore(store repo.Store, logStore LogBodyStore, now fun
|
||||
aiProviderClient: MockAIProviderClient{},
|
||||
secretEnvelope: newSecretEnvelope(developmentSecretEnvelopeKey),
|
||||
networkFingerprintKey: []byte(developmentSecretEnvelopeKey),
|
||||
distributionBuilder: unconfiguredDistributionBuilder{},
|
||||
distributionBuilds: map[string]struct{}{},
|
||||
}
|
||||
return service
|
||||
}
|
||||
|
||||
// ConfigureDistributionBuilder installs the platform-owned builder that
|
||||
// executes distribution builds. Build execution is a platform responsibility,
|
||||
// so a nil builder leaves the platform reporting an unconfigured builder rather
|
||||
// than falling back to a machine-side run endpoint.
|
||||
func (svc *CoreService) ConfigureDistributionBuilder(builder DistributionBuilder) {
|
||||
if builder == nil {
|
||||
builder = unconfiguredDistributionBuilder{}
|
||||
}
|
||||
svc.distributionBuildMu.Lock()
|
||||
svc.distributionBuilder = builder
|
||||
svc.distributionBuildMu.Unlock()
|
||||
jobs, err := svc.store.Jobs().List(domain.JobFilter{RunEndpointID: platformDistributionBuilderEndpointID})
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
for _, job := range jobs {
|
||||
if job.Capability == domain.JobCapabilityDistributionBuild && !isTerminalJobState(job.State) {
|
||||
svc.enqueueDistributionBuild(job)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func NewCoreServiceWithDurableStores(store repo.Store, logStore LogBodyStore, artifactStore ArtifactBodyStore) (*CoreService, error) {
|
||||
if artifactStore == nil {
|
||||
artifactStore = NewMemoryArtifactBodyStore()
|
||||
@@ -2195,12 +2222,18 @@ func (svc *CoreService) CreateJob(job domain.Job) (domain.Job, error) {
|
||||
return domain.Job{}, err
|
||||
}
|
||||
|
||||
endpoint, err := svc.store.RunEndpoints().Get(job.RunEndpointID)
|
||||
if err != nil {
|
||||
return domain.Job{}, fmt.Errorf("get run endpoint dependency: %w", err)
|
||||
}
|
||||
if err := svc.validateRunnableEndpoint(endpoint, job.Capability); err != nil {
|
||||
return domain.Job{}, err
|
||||
if job.Capability == domain.JobCapabilityDistributionBuild {
|
||||
if job.RunEndpointID != platformDistributionBuilderEndpointID {
|
||||
return domain.Job{}, validationError("distribution build job must target the platform builder")
|
||||
}
|
||||
} else {
|
||||
endpoint, err := svc.store.RunEndpoints().Get(job.RunEndpointID)
|
||||
if err != nil {
|
||||
return domain.Job{}, fmt.Errorf("get run endpoint dependency: %w", err)
|
||||
}
|
||||
if err := svc.validateRunnableEndpoint(endpoint, job.Capability); err != nil {
|
||||
return domain.Job{}, err
|
||||
}
|
||||
}
|
||||
if job.ServerInstanceID != "" {
|
||||
instance, err := svc.store.ServerInstances().Get(job.ServerInstanceID)
|
||||
@@ -2346,8 +2379,8 @@ func validateJobServerTarget(job domain.Job, instance domain.ServerInstance, plu
|
||||
if instance.State == domain.ServerInstanceStateDeleted {
|
||||
return validationError("server instance must not be deleted")
|
||||
}
|
||||
usesDeploymentTarget := job.Capability == domain.JobCapabilityDistributionBuild && instance.DeploymentTargetID != "" && instance.DeploymentTargetID == job.RunEndpointID
|
||||
if instance.RunEndpointID != job.RunEndpointID && !usesDeploymentTarget {
|
||||
usesPlatformBuilder := job.Capability == domain.JobCapabilityDistributionBuild && job.RunEndpointID == platformDistributionBuilderEndpointID
|
||||
if instance.RunEndpointID != job.RunEndpointID && !usesPlatformBuilder {
|
||||
return validationError("job runEndpointId must match server instance")
|
||||
}
|
||||
if plugin.ID != instance.PluginID {
|
||||
|
||||
@@ -1407,7 +1407,28 @@ func TestCoreServicePropagatesDuplicateErrors(t *testing.T) {
|
||||
}
|
||||
|
||||
func newTestCoreService() *CoreService {
|
||||
return newCoreService(repo.NewMemoryStore(), func() time.Time { return fixedTime })
|
||||
svc := newCoreService(repo.NewMemoryStore(), func() time.Time { return fixedTime })
|
||||
svc.ConfigureDistributionBuilder(staticDistributionBuilder{payload: []byte("platform-built-distribution")})
|
||||
return svc
|
||||
}
|
||||
|
||||
type staticDistributionBuilder struct {
|
||||
payload []byte
|
||||
err error
|
||||
}
|
||||
|
||||
func (builder staticDistributionBuilder) Readiness() (bool, string) {
|
||||
if builder.err != nil {
|
||||
return false, builder.err.Error()
|
||||
}
|
||||
return true, ""
|
||||
}
|
||||
|
||||
func (builder staticDistributionBuilder) Build(domain.DistributionBuildInput) ([]byte, error) {
|
||||
if builder.err != nil {
|
||||
return nil, builder.err
|
||||
}
|
||||
return domain.CopyBytes(builder.payload), nil
|
||||
}
|
||||
|
||||
func createPluginAndRunEndpoint(t *testing.T, svc *CoreService) (domain.GamePlugin, domain.RunEndpoint) {
|
||||
|
||||
Reference in New Issue
Block a user