From 0a18d881fded4d50eb72c729aeb98f476e6d63db Mon Sep 17 00:00:00 2001 From: hiroTamada <88675973+hiroTamada@users.noreply.github.com> Date: Mon, 10 Aug 2026 18:29:32 +0000 Subject: [PATCH 1/2] builds: address build-system review follow-ups - Bound POST /builds multipart reads: 512 MiB source tarball and 1 MiB per form field, returning 400 on overflow instead of unbounded reads - Fail builds hard when required secrets cannot be provided: provider errors and missing or malformed secret IDs now produce clear build or request failures, and the builder agent fails instead of proceeding without secrets on timeout - Cancel and await in-flight builds on service shutdown via queue lifecycle tracking; builds still outlive request cancellation through context.WithoutCancel so trace values propagate - Recover interrupted build attempts by removing leftover source/config volumes, deleting only this build's stale builder instance when a leftover volume is still attached - Remove the fixed pre-dial sleep in waitForResult in favor of context-aware readiness retries - Drop the unimplemented AllowedDomains build policy field - Log corrupt build metadata entries instead of silently skipping them - Correct ErrBuildInProgress documentation for the queued-to-running cancel race --- cmd/api/api/builds.go | 147 ++++------ cmd/api/api/builds_test.go | 128 +++++++++ cmd/api/main.go | 9 + lib/builds/README.md | 2 +- lib/builds/builder_agent/main.go | 94 +++++-- lib/builds/builder_agent/main_test.go | 144 ++++++++++ lib/builds/builder_disk_test.go | 8 +- lib/builds/errors.go | 11 +- lib/builds/file_secret_provider.go | 24 +- lib/builds/file_secret_provider_test.go | 52 +++- lib/builds/manager.go | 214 +++++++++++--- lib/builds/manager_test.go | 360 +++++++++++++++++++++++- lib/builds/queue.go | 102 +++++-- lib/builds/queue_test.go | 203 ++++++++++--- lib/builds/storage.go | 15 +- lib/builds/storage_test.go | 47 ++++ lib/builds/types.go | 3 - lib/builds/vsock_handler.go | 1 + lib/oapi/oapi.go | 176 ++++++------ openapi.yaml | 4 +- 20 files changed, 1405 insertions(+), 339 deletions(-) create mode 100644 cmd/api/api/builds_test.go diff --git a/cmd/api/api/builds.go b/cmd/api/api/builds.go index 110422638..2c170f0cf 100644 --- a/cmd/api/api/builds.go +++ b/cmd/api/api/builds.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "io" + "mime/multipart" "net/http" "strconv" @@ -17,6 +18,28 @@ import ( "github.com/kernel/hypeman/lib/tags" ) +var ( + // maxBuildSourceSize bounds the source tarball accepted by POST /builds. + // A var so tests can exercise the limit without large uploads. + maxBuildSourceSize int64 = 512 << 20 // 512 MiB + // maxBuildFormFieldSize bounds each small non-source multipart field + // (dockerfile, secrets, tags, and friends). + maxBuildFormFieldSize int64 = 1 << 20 // 1 MiB +) + +// readLimitedPart reads a multipart part fully, erroring when its contents +// exceed limit bytes. +func readLimitedPart(part *multipart.Part, limit int64) ([]byte, error) { + data, err := io.ReadAll(io.LimitReader(part, limit+1)) + if err != nil { + return nil, fmt.Errorf("failed to read %s field", part.FormName()) + } + if int64(len(data)) > limit { + return nil, fmt.Errorf("%s exceeds the maximum size of %d bytes", part.FormName(), limit) + } + return data, nil +} + // ListBuilds returns all builds func (s *ApiService) ListBuilds(ctx context.Context, request oapi.ListBuildsRequestObject) (oapi.ListBuildsResponseObject, error) { log := logger.FromContext(ctx) @@ -65,92 +88,48 @@ func (s *ApiService) CreateBuild(ctx context.Context, request oapi.CreateBuildRe }, nil } - switch part.FormName() { - case "source": - sourceData, err = io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_source", - Message: "failed to read source data", - }, nil + name := part.FormName() + limit := maxBuildFormFieldSize + if name == "source" { + limit = maxBuildSourceSize + } + data, err := readLimitedPart(part, limit) + part.Close() + if err != nil { + code := "invalid_request" + if name == "source" { + code = "invalid_source" } + return oapi.CreateBuild400JSONResponse{ + Code: code, + Message: err.Error(), + }, nil + } + + switch name { + case "source": + sourceData = data case "base_image_digest": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read base_image_digest field", - }, nil - } baseImageDigest = string(data) case "builder_id": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read builder_id field", - }, nil - } builderID = string(data) case "cache_scope": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read cache_scope field", - }, nil - } cacheScope = string(data) case "dockerfile": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read dockerfile field", - }, nil - } dockerfile = string(data) case "timeout_seconds": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read timeout_seconds field", - }, nil - } if v, err := strconv.Atoi(string(data)); err == nil { timeoutSeconds = v } case "memory_mb": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read memory_mb field", - }, nil - } if v, err := strconv.Atoi(string(data)); err == nil { memoryMB = v } case "cpus": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read cpus field", - }, nil - } if v, err := strconv.Atoi(string(data)); err == nil { cpus = v } case "secrets": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read secrets field", - }, nil - } if err := json.Unmarshal(data, &secrets); err != nil { return oapi.CreateBuild400JSONResponse{ Code: "invalid_request", @@ -158,40 +137,12 @@ func (s *ApiService) CreateBuild(ctx context.Context, request oapi.CreateBuildRe }, nil } case "is_admin_build": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read is_admin_build field", - }, nil - } isAdminBuild = string(data) == "true" || string(data) == "1" case "global_cache_key": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read global_cache_key field", - }, nil - } globalCacheKey = string(data) case "image_name": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read image_name field", - }, nil - } imageName = string(data) case "tags": - data, err := io.ReadAll(part) - if err != nil { - return oapi.CreateBuild400JSONResponse{ - Code: "invalid_request", - Message: "failed to read tags field", - }, nil - } parsed, err := parseTagsJSON(string(data)) if err != nil { return oapi.CreateBuild400JSONResponse{ @@ -201,7 +152,6 @@ func (s *ApiService) CreateBuild(ctx context.Context, request oapi.CreateBuildRe } resourceTags = parsed } - part.Close() } if len(sourceData) == 0 { @@ -211,6 +161,17 @@ func (s *ApiService) CreateBuild(ctx context.Context, request oapi.CreateBuildRe }, nil } + // Reject malformed secret IDs at the boundary so a bad reference fails + // the request instead of the build. + for _, secret := range secrets { + if err := builds.ValidateSecretID(secret.ID); err != nil { + return oapi.CreateBuild400JSONResponse{ + Code: "invalid_request", + Message: err.Error(), + }, nil + } + } + // Validate image_name early so the user gets a fast 400 instead of // a successful build that silently falls back to builds/{id}. if imageName != "" { diff --git a/cmd/api/api/builds_test.go b/cmd/api/api/builds_test.go new file mode 100644 index 000000000..9825ad95c --- /dev/null +++ b/cmd/api/api/builds_test.go @@ -0,0 +1,128 @@ +package api + +import ( + "bytes" + "mime/multipart" + "testing" + + "github.com/kernel/hypeman/lib/oapi" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// buildMultipartBody builds a POST /builds multipart body from small form +// fields plus the source tarball part. +func buildMultipartBody(t *testing.T, fields map[string]string, source []byte) *oapi.CreateBuildRequestObject { + t.Helper() + + var body bytes.Buffer + writer := multipart.NewWriter(&body) + for k, v := range fields { + require.NoError(t, writer.WriteField(k, v)) + } + if source != nil { + part, err := writer.CreateFormFile("source", "source.tar.gz") + require.NoError(t, err) + _, err = part.Write(source) + require.NoError(t, err) + } + require.NoError(t, writer.Close()) + + return &oapi.CreateBuildRequestObject{ + Body: multipart.NewReader(&body, writer.Boundary()), + } +} + +// overrideBuildUploadLimits shrinks the upload limits for the test and +// restores them on cleanup. Tests using it must not run in parallel. +func overrideBuildUploadLimits(t *testing.T, sourceSize, fieldSize int64) { + t.Helper() + origSource, origField := maxBuildSourceSize, maxBuildFormFieldSize + maxBuildSourceSize, maxBuildFormFieldSize = sourceSize, fieldSize + t.Cleanup(func() { + maxBuildSourceSize, maxBuildFormFieldSize = origSource, origField + }) +} + +func TestCreateBuild_SourceAtLimitAccepted(t *testing.T) { + svc := newTestService(t) + overrideBuildUploadLimits(t, 1024, 1024) + + resp, err := svc.CreateBuild(ctx(), *buildMultipartBody(t, nil, bytes.Repeat([]byte("a"), 1024))) + require.NoError(t, err) + _, ok := resp.(oapi.CreateBuild202JSONResponse) + assert.True(t, ok, "expected 202 for source exactly at the limit, got %T", resp) + + // Wait for the build goroutine to finish writing before TempDir cleanup. + require.NoError(t, svc.BuildManager.Shutdown(ctx())) +} + +func TestCreateBuild_SourceOverLimitRejected(t *testing.T) { + svc := newTestService(t) + overrideBuildUploadLimits(t, 1024, 1024) + + resp, err := svc.CreateBuild(ctx(), *buildMultipartBody(t, nil, bytes.Repeat([]byte("a"), 1025))) + require.NoError(t, err) + r, ok := resp.(oapi.CreateBuild400JSONResponse) + require.True(t, ok, "expected 400 for oversized source, got %T", resp) + assert.Equal(t, "invalid_source", r.Code) + assert.Contains(t, r.Message, "source exceeds the maximum size") +} + +func TestCreateBuild_FormFieldAtLimitAccepted(t *testing.T) { + svc := newTestService(t) + overrideBuildUploadLimits(t, 1024, 1024) + + fields := map[string]string{"dockerfile": string(bytes.Repeat([]byte("a"), 1024))} + resp, err := svc.CreateBuild(ctx(), *buildMultipartBody(t, fields, []byte("source"))) + require.NoError(t, err) + _, ok := resp.(oapi.CreateBuild202JSONResponse) + assert.True(t, ok, "expected 202 for field exactly at the limit, got %T", resp) + + // Wait for the build goroutine to finish writing before TempDir cleanup. + require.NoError(t, svc.BuildManager.Shutdown(ctx())) +} + +func TestCreateBuild_FormFieldOverLimitRejected(t *testing.T) { + svc := newTestService(t) + overrideBuildUploadLimits(t, 1024, 1024) + + fields := map[string]string{"dockerfile": string(bytes.Repeat([]byte("a"), 1025))} + resp, err := svc.CreateBuild(ctx(), *buildMultipartBody(t, fields, []byte("source"))) + require.NoError(t, err) + r, ok := resp.(oapi.CreateBuild400JSONResponse) + require.True(t, ok, "expected 400 for oversized field, got %T", resp) + assert.Equal(t, "invalid_request", r.Code) + assert.Contains(t, r.Message, "dockerfile exceeds the maximum size") +} + +func TestCreateBuild_InvalidSecretIDRejected(t *testing.T) { + svc := newTestService(t) + + for _, secrets := range []string{ + `[{"id": "../escape"}]`, + `[{"id": "a/b"}]`, + `[{"id": ""}]`, + } { + fields := map[string]string{"secrets": secrets} + resp, err := svc.CreateBuild(ctx(), *buildMultipartBody(t, fields, []byte("source"))) + require.NoError(t, err) + r, ok := resp.(oapi.CreateBuild400JSONResponse) + require.True(t, ok, "expected 400 for secrets %s, got %T", secrets, resp) + assert.Equal(t, "invalid_request", r.Code) + assert.Contains(t, r.Message, "invalid secret id") + } +} + +func TestCreateBuild_ValidSecretsAccepted(t *testing.T) { + svc := newTestService(t) + + fields := map[string]string{"secrets": `[{"id": "npm_token", "env_var": "NPM_TOKEN"}]`} + resp, err := svc.CreateBuild(ctx(), *buildMultipartBody(t, fields, []byte("source"))) + require.NoError(t, err) + _, ok := resp.(oapi.CreateBuild202JSONResponse) + assert.True(t, ok, "expected 202 for valid secrets, got %T", resp) + + // Wait for the build goroutine to finish writing before TempDir cleanup. + require.NoError(t, svc.BuildManager.Shutdown(ctx())) +} diff --git a/cmd/api/main.go b/cmd/api/main.go index 98d7be720..526d2c62d 100644 --- a/cmd/api/main.go +++ b/cmd/api/main.go @@ -695,6 +695,15 @@ func run() error { logger.Info("ingress manager shutdown complete") } + // Cancel in-flight builds and wait for their goroutines to return. + // Builds still pending stay queued on disk and recover on next start. + if err := app.BuildManager.Shutdown(shutdownCtx); err != nil { + logger.Error("failed to shutdown build manager", "error", err) + // Don't return error - continue with shutdown + } else { + logger.Info("build manager shutdown complete") + } + return errors.Join(shutdownErrs...) }) diff --git a/lib/builds/README.md b/lib/builds/README.md index f7ec2844f..3ab59ae50 100644 --- a/lib/builds/README.md +++ b/lib/builds/README.md @@ -242,7 +242,7 @@ queued → building → pushing → ready 1. **Isolation**: Each build runs in a fresh microVM (Cloud Hypervisor) 2. **Rootless**: BuildKit runs without root privileges -3. **Network Control**: `network_mode: isolated` or `egress` with optional domain allowlist +3. **Network Control**: `network_mode: isolated` or `egress` (outbound allowed) 4. **Secret Handling**: Secrets fetched via vsock, never written to disk in guest 5. **Cache Isolation**: Per-tenant cache scopes prevent cross-tenant cache poisoning 6. **Registry Auth**: Short-lived JWT tokens scoped to specific repositories (builds/{id}, cache/{scope}) diff --git a/lib/builds/builder_agent/main.go b/lib/builds/builder_agent/main.go index 6138144a9..601c066e8 100644 --- a/lib/builds/builder_agent/main.go +++ b/lib/builds/builder_agent/main.go @@ -35,6 +35,13 @@ const ( vsockPort = 5001 // Build agent port (different from exec agent) ) +var ( + // secretsDir is where fetched secrets are written for buildctl secret mounts + secretsDir = "/run/secrets" + // secretsTimeout bounds how long the build waits for the host to provide secrets + secretsTimeout = 30 * time.Second +) + // BuildConfig matches the BuildConfig type from lib/builds/types.go type BuildConfig struct { JobID string `json:"job_id"` @@ -93,6 +100,7 @@ type VsockMessage struct { Log string `json:"log,omitempty"` SecretIDs []string `json:"secret_ids,omitempty"` // For secrets request to host Secrets map[string]string `json:"secrets,omitempty"` // For secrets response from host + Error string `json:"error,omitempty"` // Set on a secrets response when the host could not provide every requested secret } // Global state for the result to send when host connects @@ -106,6 +114,8 @@ var ( buildConfigLock sync.Mutex secretsReady = make(chan struct{}) secretsOnce sync.Once + secretsErr error + secretsErrLock sync.Mutex // Encoder lock protects concurrent access to json.Encoder // (the goroutine sending build_result and the main loop handling get_status) @@ -268,8 +278,10 @@ func handleHostConnection(conn net.Conn) { // Request secrets if we have any configured if err := handleSecretsRequest(encoder, decoder); err != nil { log.Printf("Failed to fetch secrets: %v", err) + setSecretsErr(err) } - // Signal that secrets are ready (even if failed, build can proceed) + // Signal that the secrets exchange completed; the build checks + // secretsErr and fails rather than proceeding without required secrets. secretsOnce.Do(func() { close(secretsReady) }) @@ -402,13 +414,29 @@ func handleSecretsRequest(encoder *json.Encoder, decoder *json.Decoder) error { return fmt.Errorf("unexpected response type: %s", resp.Type) } - // Write secrets to /run/secrets/ - if err := os.MkdirAll("/run/secrets", 0700); err != nil { + if resp.Error != "" { + return fmt.Errorf("host could not provide secrets: %s", resp.Error) + } + + // Every requested secret must be present; building without a required + // secret fails later with a much less clear error. + var missing []string + for _, id := range secretIDs { + if _, ok := resp.Secrets[id]; !ok { + missing = append(missing, id) + } + } + if len(missing) > 0 { + return fmt.Errorf("host did not provide secrets: %s", strings.Join(missing, ", ")) + } + + // Write secrets to the secrets dir + if err := os.MkdirAll(secretsDir, 0700); err != nil { return fmt.Errorf("create secrets dir: %w", err) } for id, value := range resp.Secrets { - secretPath := fmt.Sprintf("/run/secrets/%s", id) + secretPath := filepath.Join(secretsDir, id) if err := os.WriteFile(secretPath, []byte(value), 0600); err != nil { return fmt.Errorf("write secret %s: %w", id, err) } @@ -419,6 +447,45 @@ func handleSecretsRequest(encoder *json.Encoder, decoder *json.Decoder) error { return nil } +// setSecretsErr records a secrets-exchange failure so the build fails with it. +func setSecretsErr(err error) { + secretsErrLock.Lock() + secretsErr = err + secretsErrLock.Unlock() +} + +// waitForSecrets blocks until the host completes the secrets exchange for +// the configured secrets. Any failure — host-side fetch errors, a timeout, +// or a secret that never landed on disk — is fatal to the build: proceeding +// without a required secret only produces a more confusing error later. +func waitForSecrets(ctx context.Context, secrets []SecretRef) error { + if len(secrets) == 0 { + return nil + } + + select { + case <-secretsReady: + case <-time.After(secretsTimeout): + return fmt.Errorf("timed out after %s waiting for secrets from host", secretsTimeout) + case <-ctx.Done(): + return fmt.Errorf("build timeout while waiting for secrets") + } + + secretsErrLock.Lock() + err := secretsErr + secretsErrLock.Unlock() + if err != nil { + return fmt.Errorf("fetch secrets: %w", err) + } + + for _, s := range secrets { + if _, err := os.Stat(filepath.Join(secretsDir, s.ID)); err != nil { + return fmt.Errorf("required secret %q was not provided by the host", s.ID) + } + } + return nil +} + // runBuildProcess runs the actual build and stores the result func runBuildProcess() { start := time.Now() @@ -473,27 +540,20 @@ func runBuildProcess() { defer cancel() } - // Wait for secrets if any are configured + // Wait for secrets if any are configured. Secrets are required: any + // failure here fails the build instead of proceeding without them. if len(config.Secrets) > 0 { log.Printf("Waiting for secrets from host...") - select { - case <-secretsReady: - log.Printf("Secrets ready, proceeding with build") - case <-time.After(30 * time.Second): - log.Printf("Warning: Timeout waiting for secrets, proceeding anyway") - // Signal secrets ready to avoid blocking other goroutines - secretsOnce.Do(func() { - close(secretsReady) - }) - case <-ctx.Done(): + if err := waitForSecrets(ctx, config.Secrets); err != nil { setResult(BuildResult{ Success: false, - Error: "build timeout while waiting for secrets", + Error: err.Error(), Logs: logWriter.String(), DurationMS: time.Since(start).Milliseconds(), }) return } + log.Printf("Secrets ready, proceeding with build") } // Ensure Dockerfile exists (either in source or provided via config) @@ -881,7 +941,7 @@ func runBuild(ctx context.Context, config *BuildConfig, logWriter io.Writer) (st // Add secret mounts for _, secret := range config.Secrets { - secretPath := fmt.Sprintf("/run/secrets/%s", secret.ID) + secretPath := filepath.Join(secretsDir, secret.ID) args = append(args, "--secret", fmt.Sprintf("id=%s,src=%s", secret.ID, secretPath)) } diff --git a/lib/builds/builder_agent/main_test.go b/lib/builds/builder_agent/main_test.go index 4dc3d70e6..5fdca6012 100644 --- a/lib/builds/builder_agent/main_test.go +++ b/lib/builds/builder_agent/main_test.go @@ -1,12 +1,17 @@ package main import ( + "context" + "encoding/json" "errors" "fmt" + "net" "os" "path/filepath" "strings" + "sync" "testing" + "time" "github.com/docker/go-units" "github.com/stretchr/testify/assert" @@ -253,3 +258,142 @@ func TestReconcileBuildkitVersionStamp_VersionError(t *testing.T) { require.ErrorContains(t, err, "buildkitd broken") } + +// resetSecretsState isolates the package-level secrets coordination state for +// a test. Tests using it must not run in parallel. +func resetSecretsState(t *testing.T) { + t.Helper() + + origDir, origTimeout := secretsDir, secretsTimeout + secretsDir = t.TempDir() + secretsErrLock.Lock() + secretsErr = nil + secretsErrLock.Unlock() + secretsReady = make(chan struct{}) + secretsOnce = sync.Once{} + + t.Cleanup(func() { + secretsDir, secretsTimeout = origDir, origTimeout + secretsErrLock.Lock() + secretsErr = nil + secretsErrLock.Unlock() + secretsReady = make(chan struct{}) + secretsOnce = sync.Once{} + }) +} + +// setTestBuildConfig installs a build config for handleSecretsRequest. +func setTestBuildConfig(t *testing.T, cfg *BuildConfig) { + t.Helper() + buildConfigLock.Lock() + buildConfig = cfg + buildConfigLock.Unlock() + t.Cleanup(func() { + buildConfigLock.Lock() + buildConfig = nil + buildConfigLock.Unlock() + }) +} + +// doSecretsExchange runs handleSecretsRequest against a simulated host that +// answers the get_secrets request with hostResp. +func doSecretsExchange(t *testing.T, cfg *BuildConfig, hostResp VsockMessage) error { + t.Helper() + setTestBuildConfig(t, cfg) + + agentConn, hostConn := net.Pipe() + defer agentConn.Close() + defer hostConn.Close() + + go func() { + decoder := json.NewDecoder(hostConn) + encoder := json.NewEncoder(hostConn) + var req VsockMessage + if err := decoder.Decode(&req); err != nil { + return + } + encoder.Encode(hostResp) + }() + + return handleSecretsRequest(json.NewEncoder(agentConn), json.NewDecoder(agentConn)) +} + +func TestHandleSecretsRequest_HostError(t *testing.T) { + resetSecretsState(t) + + err := doSecretsExchange(t, + &BuildConfig{Secrets: []SecretRef{{ID: "npm_token"}}}, + VsockMessage{Type: "secrets_response", Error: "secret not found: npm_token"}) + + require.Error(t, err) + assert.Contains(t, err.Error(), "secret not found: npm_token") +} + +func TestHandleSecretsRequest_MissingSecret(t *testing.T) { + resetSecretsState(t) + + err := doSecretsExchange(t, + &BuildConfig{Secrets: []SecretRef{{ID: "a"}, {ID: "b"}}}, + VsockMessage{Type: "secrets_response", Secrets: map[string]string{"a": "1"}}) + + require.Error(t, err) + assert.Contains(t, err.Error(), "did not provide") + assert.Contains(t, err.Error(), "b") +} + +func TestHandleSecretsRequest_Success(t *testing.T) { + resetSecretsState(t) + + err := doSecretsExchange(t, + &BuildConfig{Secrets: []SecretRef{{ID: "a"}, {ID: "b"}}}, + VsockMessage{Type: "secrets_response", Secrets: map[string]string{"a": "1", "b": "2"}}) + + require.NoError(t, err) + for id, want := range map[string]string{"a": "1", "b": "2"} { + got, err := os.ReadFile(filepath.Join(secretsDir, id)) + require.NoError(t, err) + assert.Equal(t, want, string(got)) + } +} + +func TestWaitForSecrets_NoneConfigured(t *testing.T) { + resetSecretsState(t) + assert.NoError(t, waitForSecrets(context.Background(), nil)) +} + +func TestWaitForSecrets_TimeoutIsFatal(t *testing.T) { + resetSecretsState(t) + secretsTimeout = 20 * time.Millisecond + + err := waitForSecrets(context.Background(), []SecretRef{{ID: "a"}}) + require.Error(t, err) + assert.Contains(t, err.Error(), "timed out") +} + +func TestWaitForSecrets_HostErrorIsFatal(t *testing.T) { + resetSecretsState(t) + + setSecretsErr(errors.New("secret not found: a")) + secretsOnce.Do(func() { close(secretsReady) }) + + err := waitForSecrets(context.Background(), []SecretRef{{ID: "a"}}) + require.Error(t, err) + assert.Contains(t, err.Error(), "secret not found: a") +} + +func TestWaitForSecrets_MissingFileIsFatal(t *testing.T) { + resetSecretsState(t) + secretsOnce.Do(func() { close(secretsReady) }) + + err := waitForSecrets(context.Background(), []SecretRef{{ID: "a"}}) + require.Error(t, err) + assert.Contains(t, err.Error(), `required secret "a"`) +} + +func TestWaitForSecrets_Success(t *testing.T) { + resetSecretsState(t) + require.NoError(t, os.WriteFile(filepath.Join(secretsDir, "a"), []byte("v"), 0600)) + secretsOnce.Do(func() { close(secretsReady) }) + + assert.NoError(t, waitForSecrets(context.Background(), []SecretRef{{ID: "a"}})) +} diff --git a/lib/builds/builder_disk_test.go b/lib/builds/builder_disk_test.go index e3ac5f7f0..026142338 100644 --- a/lib/builds/builder_disk_test.go +++ b/lib/builds/builder_disk_test.go @@ -322,11 +322,11 @@ func TestBuilderQueueSerialization(t *testing.T) { started := make(chan string, 3) req := CreateBuildRequest{BuilderID: "builder-a"} - pos1 := mgr.queue.EnqueueSerial("build-1", req, req.BuilderID, func() { + pos1 := mgr.queue.EnqueueSerial(context.Background(), "build-1", req, req.BuilderID, func(context.Context) { started <- "build-1" <-release }) - pos2 := mgr.queue.EnqueueSerial("build-2", req, req.BuilderID, func() { + pos2 := mgr.queue.EnqueueSerial(context.Background(), "build-2", req, req.BuilderID, func(context.Context) { started <- "build-2" <-release }) @@ -376,8 +376,8 @@ func TestBuilderDeleteAndPruneBlockedByQueuedBuilds(t *testing.T) { release := make(chan struct{}) defer close(release) req := CreateBuildRequest{BuilderID: builder.ID} - mgr.queue.EnqueueSerial("build-1", req, builder.ID, func() { <-release }) - mgr.queue.EnqueueSerial("build-2", req, builder.ID, func() { <-release }) + mgr.queue.EnqueueSerial(context.Background(), "build-1", req, builder.ID, func(context.Context) { <-release }) + mgr.queue.EnqueueSerial(context.Background(), "build-2", req, builder.ID, func(context.Context) { <-release }) err = mgr.builderManager.DeleteBuilder(ctx, builder.ID) assert.ErrorIs(t, err, builders.ErrInUse, "delete must be blocked by a running build") diff --git a/lib/builds/errors.go b/lib/builds/errors.go index 2e8d888d8..7f031650b 100644 --- a/lib/builds/errors.go +++ b/lib/builds/errors.go @@ -30,6 +30,15 @@ var ( // ErrBuilderNotReady is returned when the builder image is not available ErrBuilderNotReady = errors.New("builder image not ready") - // ErrBuildInProgress is returned when trying to cancel a build that's already complete + // ErrBuildInProgress is returned when cancelling a queued build that was + // already picked up and started running between the status check and the + // queue removal. ErrBuildInProgress = errors.New("build in progress") + + // ErrInvalidSecretID is returned when a secret ID is empty or contains + // path separators + ErrInvalidSecretID = errors.New("invalid secret id") + + // ErrSecretNotFound is returned when a requested secret has no value + ErrSecretNotFound = errors.New("secret not found") ) diff --git a/lib/builds/file_secret_provider.go b/lib/builds/file_secret_provider.go index 611503479..fac8a40e5 100644 --- a/lib/builds/file_secret_provider.go +++ b/lib/builds/file_secret_provider.go @@ -8,6 +8,19 @@ import ( "strings" ) +// ValidateSecretID rejects secret IDs that are empty or could escape the +// secrets directory (path traversal). It is enforced both at the API +// boundary and by the file-based secret provider. +func ValidateSecretID(id string) error { + if id == "" { + return fmt.Errorf("%w: must not be empty", ErrInvalidSecretID) + } + if strings.Contains(id, "/") || strings.Contains(id, "\\") || id == ".." || id == "." { + return fmt.Errorf("%w: %q must not contain path separators", ErrInvalidSecretID, id) + } + return nil +} + // FileSecretProvider reads secrets from files in a directory. // Each secret is stored as a file named by its ID, with the secret value as the file content. // Example: /etc/hypeman/secrets/npm_token contains the npm token value. @@ -24,15 +37,15 @@ func NewFileSecretProvider(secretsDir string) *FileSecretProvider { } // GetSecrets returns the values for the given secret IDs by reading files from the secrets directory. -// Missing secrets are silently skipped (not an error). -// Returns an error only if a secret file exists but cannot be read. +// Requested secrets are required: an invalid ID or a missing secret file is an +// error, so a build never proceeds with silently empty secrets. func (p *FileSecretProvider) GetSecrets(ctx context.Context, secretIDs []string) (map[string]string, error) { result := make(map[string]string) for _, id := range secretIDs { // Validate secret ID to prevent path traversal - if strings.Contains(id, "/") || strings.Contains(id, "\\") || id == ".." || id == "." { - continue // Skip invalid IDs + if err := ValidateSecretID(id); err != nil { + return nil, err } path := filepath.Join(p.secretsDir, id) @@ -47,8 +60,7 @@ func (p *FileSecretProvider) GetSecrets(ctx context.Context, secretIDs []string) data, err := os.ReadFile(path) if err != nil { if os.IsNotExist(err) { - // Secret doesn't exist - skip it (not an error) - continue + return nil, fmt.Errorf("%w: %s", ErrSecretNotFound, id) } return nil, fmt.Errorf("read secret %s: %w", id, err) } diff --git a/lib/builds/file_secret_provider_test.go b/lib/builds/file_secret_provider_test.go index 59e8741f2..8a6a6601f 100644 --- a/lib/builds/file_secret_provider_test.go +++ b/lib/builds/file_secret_provider_test.go @@ -32,17 +32,19 @@ func TestFileSecretProvider_GetSecrets(t *testing.T) { assert.Equal(t, "github-secret-value", secrets["github_token"]) }) - t.Run("missing secrets are skipped", func(t *testing.T) { + t.Run("missing secret is an error", func(t *testing.T) { secrets, err := provider.GetSecrets(ctx, []string{"npm_token", "nonexistent"}) - require.NoError(t, err) - assert.Len(t, secrets, 1) - assert.Equal(t, "npm-secret-value", secrets["npm_token"]) + require.Error(t, err) + assert.ErrorIs(t, err, ErrSecretNotFound) + assert.Contains(t, err.Error(), "nonexistent") + assert.Nil(t, secrets) }) - t.Run("all missing secrets returns empty map", func(t *testing.T) { + t.Run("all missing secrets is an error", func(t *testing.T) { secrets, err := provider.GetSecrets(ctx, []string{"missing1", "missing2"}) - require.NoError(t, err) - assert.Empty(t, secrets) + require.Error(t, err) + assert.ErrorIs(t, err, ErrSecretNotFound) + assert.Nil(t, secrets) }) t.Run("whitespace is trimmed", func(t *testing.T) { @@ -51,16 +53,27 @@ func TestFileSecretProvider_GetSecrets(t *testing.T) { assert.Equal(t, "trimmed", secrets["with_whitespace"]) }) - t.Run("path traversal is blocked", func(t *testing.T) { + t.Run("path traversal is an error", func(t *testing.T) { secrets, err := provider.GetSecrets(ctx, []string{"../etc/passwd", "../../root/.ssh/id_rsa"}) - require.NoError(t, err) - assert.Empty(t, secrets) + require.Error(t, err) + assert.ErrorIs(t, err, ErrInvalidSecretID) + assert.Nil(t, secrets) }) - t.Run("special characters in ID are blocked", func(t *testing.T) { - secrets, err := provider.GetSecrets(ctx, []string{"foo/bar", "baz\\qux", "..", "."}) - require.NoError(t, err) - assert.Empty(t, secrets) + t.Run("special characters in ID are an error", func(t *testing.T) { + for _, id := range []string{"foo/bar", "baz\\qux", "..", "."} { + secrets, err := provider.GetSecrets(ctx, []string{id}) + require.Error(t, err, "id %q", id) + assert.ErrorIs(t, err, ErrInvalidSecretID, "id %q", id) + assert.Nil(t, secrets, "id %q", id) + } + }) + + t.Run("empty ID is an error", func(t *testing.T) { + secrets, err := provider.GetSecrets(ctx, []string{""}) + require.Error(t, err) + assert.ErrorIs(t, err, ErrInvalidSecretID) + assert.Nil(t, secrets) }) t.Run("empty request returns empty map", func(t *testing.T) { @@ -91,6 +104,17 @@ func TestFileSecretProvider_ContextCancellation(t *testing.T) { assert.True(t, err == context.Canceled || len(secrets) <= 3) } +func TestValidateSecretID(t *testing.T) { + for _, id := range []string{"npm_token", "github-token", "token.json", "a.b_c-d"} { + assert.NoError(t, ValidateSecretID(id), "id %q", id) + } + for _, id := range []string{"", "../secret", "a/b", "a\\b", "..", "."} { + err := ValidateSecretID(id) + require.Error(t, err, "id %q", id) + assert.ErrorIs(t, err, ErrInvalidSecretID, "id %q", id) + } +} + func TestNoOpSecretProvider(t *testing.T) { provider := &NoOpSecretProvider{} ctx := context.Background() diff --git a/lib/builds/manager.go b/lib/builds/manager.go index b9157d93c..1e5229361 100644 --- a/lib/builds/manager.go +++ b/lib/builds/manager.go @@ -5,7 +5,9 @@ import ( "context" _ "embed" "encoding/json" + "errors" "fmt" + "io" "log/slog" "net" "os" @@ -35,6 +37,12 @@ const ( releaseBuildMaxAttempts = 5 releaseBuildRetryDelay = time.Second builderImageRetryInterval = 5 * time.Second + + // buildAgentDialMaxAttempts bounds how often waitForResult dials the + // builder agent while the VM boots; buildAgentDialRetryInterval spaces + // the attempts. + buildAgentDialMaxAttempts = 30 + buildAgentDialRetryInterval = 2 * time.Second ) //go:embed images/generic/Dockerfile @@ -46,6 +54,10 @@ type Manager interface { // This should be called once when the API server starts. Start(ctx context.Context) error + // Shutdown cancels in-flight builds and waits for their goroutines to + // return. Pending builds are left queued for recovery on the next start. + Shutdown(ctx context.Context) error + // CreateBuild starts a new build job CreateBuild(ctx context.Context, req CreateBuildRequest, sourceData []byte) (*Build, error) @@ -230,6 +242,12 @@ func (m *manager) ReadyForBuilds() bool { return m.builderReady.Load() } +// Shutdown cancels in-flight builds and waits for their goroutines to return. +// Pending builds stay queued on disk and are recovered on the next start. +func (m *manager) Shutdown(ctx context.Context) error { + return m.queue.Shutdown(ctx) +} + // ensureBuilderImage ensures the builder image is available in the image store. // // If BUILDER_IMAGE is set, it checks whether the image is already in the store @@ -614,8 +632,13 @@ func (m *manager) CreateBuild(ctx context.Context, req CreateBuildRequest, sourc // Enqueue the build. Builds sharing a builder are serialized by builder // ID so a waiting same-builder build stays pending instead of holding a // global concurrency slot while blocked on the builder's disk. - queuePos := m.queue.EnqueueSerial(id, req, req.BuilderID, func() { - m.runBuild(context.Background(), id, req, policy) + // + // The build must outlive the request: WithoutCancel detaches request + // cancellation while preserving context values (trace context, logger). + // The queue cancels the run only on service shutdown. + runCtx := context.WithoutCancel(ctx) + queuePos := m.queue.EnqueueSerial(runCtx, id, req, req.BuilderID, func(ctx context.Context) { + m.runBuild(ctx, id, req, policy) }) build := meta.toBuild() @@ -657,7 +680,7 @@ func (m *manager) BuilderHasBuilds(builderID string) bool { if m.pendingRecovered.Load() { return false } - pending, err := listPendingBuilds(m.paths) + pending, err := listPendingBuilds(m.paths, m.logger) if err != nil { // Fail closed: blocking a delete is recoverable, deleting a // builder out from under a pending build is not. @@ -825,13 +848,8 @@ func (m *manager) executeBuild(ctx context.Context, id string, req CreateBuildRe defer sourceFile.Close() // Create volume with source (using the volume manager's archive import) - _, err = m.volumeManager.CreateVolumeFromArchive(ctx, volumes.CreateVolumeFromArchiveRequest{ - Id: &sourceVolID, - Name: sourceVolID, - SizeGb: 10, // 10GB should be enough for most source bundles - }, sourceFile) - if err != nil { - return nil, fmt.Errorf("create source volume: %w", err) + if err := m.createBuildSourceVolume(ctx, id, sourceVolID, sourceFile); err != nil { + return nil, err } defer m.volumeManager.DeleteVolume(context.Background(), sourceVolID) @@ -844,25 +862,8 @@ func (m *manager) executeBuild(ctx context.Context, id string, req CreateBuildRe defer os.Remove(configVolPath) // Clean up the config disk file // Register the config volume with the volume manager - _, err = m.volumeManager.CreateVolume(ctx, volumes.CreateVolumeRequest{ - Id: &configVolID, - Name: configVolID, - SizeGb: 1, - }) - if err != nil { - // If volume creation fails, try to use the disk file directly - // by copying it to the expected location - volPath := m.paths.VolumeData(configVolID) - if copyErr := copyFile(configVolPath, volPath); copyErr != nil { - return nil, fmt.Errorf("setup config volume: %w", copyErr) - } - } else { - // Copy our config disk over the empty volume - volPath := m.paths.VolumeData(configVolID) - if err := copyFile(configVolPath, volPath); err != nil { - m.volumeManager.DeleteVolume(context.Background(), configVolID) - return nil, fmt.Errorf("write config to volume: %w", err) - } + if err := m.registerBuildConfigVolume(ctx, id, configVolID, configVolPath); err != nil { + return nil, err } defer m.volumeManager.DeleteVolume(context.Background(), configVolID) @@ -957,20 +958,128 @@ func (m *manager) executeBuild(ctx context.Context, id string, req CreateBuildRe return result, nil } +// missingSecretsError returns an error naming any requested secret IDs that +// the provider did not return a value for. +func missingSecretsError(secretIDs []string, secrets map[string]string) error { + var missing []string + for _, id := range secretIDs { + if _, ok := secrets[id]; !ok { + missing = append(missing, id) + } + } + if len(missing) > 0 { + return fmt.Errorf("%w: %s", ErrSecretNotFound, strings.Join(missing, ", ")) + } + return nil +} + +// createBuildSourceVolume creates the source volume for a build. Build volume +// names are deterministic (build-source-), so a re-run of the same build +// after a crash (e.g. via RecoverPendingBuilds) can hit a leftover volume from +// the interrupted attempt. Tolerate that case: remove the leftover and retry +// the create once. CreateVolumeFromArchive reports ErrAlreadyExists before +// consuming the archive reader, so the retry can reuse source. +func (m *manager) createBuildSourceVolume(ctx context.Context, buildID, volID string, source io.Reader) error { + req := volumes.CreateVolumeFromArchiveRequest{ + Id: &volID, + Name: volID, + SizeGb: 10, // 10GB should be enough for most source bundles + } + _, err := m.volumeManager.CreateVolumeFromArchive(ctx, req, source) + if errors.Is(err, volumes.ErrAlreadyExists) { + m.logger.Info("removing leftover source volume from interrupted build attempt", "build_id", buildID, "volume", volID) + if delErr := m.deleteLeftoverBuildVolume(ctx, buildID, volID); delErr != nil { + return fmt.Errorf("remove leftover source volume: %w", delErr) + } + _, err = m.volumeManager.CreateVolumeFromArchive(ctx, req, source) + } + if err != nil { + return fmt.Errorf("create source volume: %w", err) + } + return nil +} + +// registerBuildConfigVolume registers the config disk as a volume and copies +// the config data onto it. Like the source volume, the config volume has a +// deterministic name (build-config-), so a re-run of the same build after +// a crash can hit a leftover from the interrupted attempt; remove it and +// retry the create once rather than silently copying over the stale volume. +func (m *manager) registerBuildConfigVolume(ctx context.Context, buildID, volID, configDiskPath string) error { + req := volumes.CreateVolumeRequest{ + Id: &volID, + Name: volID, + SizeGb: 1, + } + _, err := m.volumeManager.CreateVolume(ctx, req) + if errors.Is(err, volumes.ErrAlreadyExists) { + m.logger.Info("removing leftover config volume from interrupted build attempt", "build_id", buildID, "volume", volID) + if delErr := m.deleteLeftoverBuildVolume(ctx, buildID, volID); delErr != nil { + return fmt.Errorf("remove leftover config volume: %w", delErr) + } + if _, err = m.volumeManager.CreateVolume(ctx, req); err != nil { + // A failed recreate must surface as-is; falling through to the + // copy-over fallback would mask it and could write the config + // onto an unmanaged disk file. + return fmt.Errorf("create config volume: %w", err) + } + } else if err != nil { + // If volume creation fails, try to use the disk file directly + // by copying it to the expected location + volPath := m.paths.VolumeData(volID) + if copyErr := copyFile(configDiskPath, volPath); copyErr != nil { + return fmt.Errorf("setup config volume: %w", copyErr) + } + return nil + } + // Copy our config disk over the empty volume + volPath := m.paths.VolumeData(volID) + if err := copyFile(configDiskPath, volPath); err != nil { + m.volumeManager.DeleteVolume(context.Background(), volID) + return fmt.Errorf("write config to volume: %w", err) + } + return nil +} + +// deleteLeftoverBuildVolume removes a build volume left behind by an +// interrupted prior attempt of the same build. If the volume is still +// attached, it is almost certainly attached to the interrupted attempt's +// stale builder instance (builder-); delete that instance (which detaches +// its volumes) and retry the volume delete. A volume attached to an unknown +// instance is never force-deleted; an error is returned instead. +func (m *manager) deleteLeftoverBuildVolume(ctx context.Context, buildID, volID string) error { + err := m.volumeManager.DeleteVolume(ctx, volID) + if !errors.Is(err, volumes.ErrInUse) { + return err + } + + builderName := fmt.Sprintf("builder-%s", buildID) + inst, getErr := m.instanceManager.GetInstance(ctx, builderName) + if getErr != nil { + if errors.Is(getErr, instances.ErrNotFound) { + return fmt.Errorf("volume %s is in use but stale builder instance %q was not found; refusing to force-delete a volume attached to an unknown instance", volID, builderName) + } + return fmt.Errorf("look up stale builder instance %q: %w", builderName, getErr) + } + + m.logger.Info("deleting stale builder instance from interrupted build attempt", "build_id", buildID, "instance", inst.Id, "volume", volID) + if delErr := m.instanceManager.DeleteInstance(ctx, inst.Id); delErr != nil { + return fmt.Errorf("delete stale builder instance %s: %w", inst.Id, delErr) + } + return m.volumeManager.DeleteVolume(ctx, volID) +} + // waitForResult waits for the build result from the builder agent via vsock func (m *manager) waitForResult(ctx context.Context, buildID string, inst *instances.Instance) (*BuildResult, error) { - // Wait a bit for the VM to start and the builder agent to listen on vsock - time.Sleep(3 * time.Second) - - // Try to connect to the builder agent with retries + // Connect to the builder agent with retries: the agent only starts + // listening once the VM has booted, so early dial attempts are expected + // to fail. The retry interval bounds the wait between attempts, so no + // fixed startup sleep is needed. var conn net.Conn var err error - for attempt := 0; attempt < 30; attempt++ { - select { - case <-ctx.Done(): - return nil, ctx.Err() - default: + for attempt := 0; attempt < buildAgentDialMaxAttempts; attempt++ { + if err := ctx.Err(); err != nil { + return nil, err } dialer, dialerErr := m.instanceManager.GetVsockDialer(ctx, inst.Id) @@ -988,7 +1097,6 @@ func (m *manager) waitForResult(ctx context.Context, buildID string, inst *insta } m.logger.Debug("waiting for builder agent", "attempt", attempt+1, "error", err) - time.Sleep(2 * time.Second) // Check if instance is still running current, checkErr := m.instanceManager.GetInstance(ctx, inst.Id) @@ -1001,6 +1109,12 @@ func (m *manager) waitForResult(ctx context.Context, buildID string, inst *insta Error: "builder instance stopped unexpectedly", }, nil } + + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-time.After(buildAgentDialRetryInterval): + } } if conn == nil { @@ -1055,11 +1169,21 @@ func (m *manager) waitForResult(ctx context.Context, buildID string, inst *insta // Agent is requesting secrets m.logger.Debug("agent requesting secrets", "instance", inst.Id, "secret_ids", dr.response.SecretIDs) - // Fetch secrets from provider + // Fetch secrets from provider. Secrets are required for the builds + // that declare them: report failures (and missing values) to the + // agent instead of letting it build with empty secrets and fail + // later with a confusing error. The agent fails the build on an + // error response and still reports a build_result. secrets, err := m.secretProvider.GetSecrets(ctx, dr.response.SecretIDs) + if err == nil { + err = missingSecretsError(dr.response.SecretIDs, secrets) + } if err != nil { m.logger.Error("failed to fetch secrets", "error", err) - secrets = make(map[string]string) + if encErr := encoder.Encode(VsockMessage{Type: "secrets_response", Error: err.Error()}); encErr != nil { + return nil, fmt.Errorf("send secrets error response: %w", encErr) + } + continue } // Send secrets response @@ -1229,7 +1353,7 @@ func (m *manager) GetBuild(ctx context.Context, id string) (*Build, error) { // ListBuilds returns all builds func (m *manager) ListBuilds(ctx context.Context) ([]*Build, error) { - metas, err := listAllBuilds(m.paths) + metas, err := listAllBuilds(m.paths, m.logger) if err != nil { return nil, err } @@ -1460,7 +1584,7 @@ func (m *manager) StreamBuildEvents(ctx context.Context, id string, follow bool) // RecoverPendingBuilds recovers builds that were interrupted on restart func (m *manager) RecoverPendingBuilds() { - pending, err := listPendingBuilds(m.paths) + pending, err := listPendingBuilds(m.paths, m.logger) if err != nil { m.logger.Error("list pending builds for recovery", "error", err) return @@ -1483,12 +1607,12 @@ func (m *manager) RecoverPendingBuilds() { continue } - m.queue.EnqueueSerial(meta.ID, *meta.Request, meta.Request.BuilderID, func() { + m.queue.EnqueueSerial(context.Background(), meta.ID, *meta.Request, meta.Request.BuilderID, func(ctx context.Context) { policy := DefaultBuildPolicy() if meta.Request.BuildPolicy != nil { policy = *meta.Request.BuildPolicy } - m.runBuild(context.Background(), meta.ID, *meta.Request, &policy) + m.runBuild(ctx, meta.ID, *meta.Request, &policy) }) } } diff --git a/lib/builds/manager_test.go b/lib/builds/manager_test.go index dfdc9fca3..d00e9c4bd 100644 --- a/lib/builds/manager_test.go +++ b/lib/builds/manager_test.go @@ -669,7 +669,7 @@ func TestCancelBuild_QueuedBuild(t *testing.T) { started := make(chan struct{}) // Add a blocking build to fill the single slot - queue.Enqueue("build-1", CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { started <- struct{}{} select {} // Block forever }) @@ -678,7 +678,7 @@ func TestCancelBuild_QueuedBuild(t *testing.T) { <-started // Add a second build - this one should be queued - queue.Enqueue("build-2", CreateBuildRequest{}, func() {}) + queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) {}) // Verify it's pending assert.Equal(t, 1, queue.PendingCount()) @@ -691,6 +691,147 @@ func TestCancelBuild_QueuedBuild(t *testing.T) { assert.Equal(t, 0, queue.PendingCount()) } +// TestManagerShutdown_InFlightBuild verifies the build run detaches from +// request cancellation while keeping request context values, and that +// manager shutdown cancels the run and waits for its goroutine. +func TestManagerShutdown_InFlightBuild(t *testing.T) { + mgr, instanceMgr, _, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + type ctxKey struct{} + reqCtx, reqCancel := context.WithCancel(context.WithValue(context.Background(), ctxKey{}, "trace-abc")) + + // The instance create hook observes the context values propagated into + // the detached build run. + values := make(chan any, 1) + instanceMgr.createFunc = func(ctx context.Context, req instances.CreateInstanceRequest) (*instances.Instance, error) { + values <- ctx.Value(ctxKey{}) + inst := &instances.Instance{ + StoredMetadata: instances.StoredMetadata{Id: "inst-" + req.Name, Name: req.Name}, + State: instances.StateRunning, + } + instanceMgr.instances[inst.Id] = inst + return inst, nil + } + dialAttempted := make(chan struct{}) + var dialOnce sync.Once + instanceMgr.vsockDialerFunc = func(ctx context.Context, instanceID string) (hypervisor.VsockDialer, error) { + dialOnce.Do(func() { close(dialAttempted) }) + return nil, fmt.Errorf("vsock dialer unavailable in test") + } + + build, err := mgr.CreateBuild(reqCtx, CreateBuildRequest{}, []byte("source")) + require.NoError(t, err) + + // Cancelling the request context must not stop the build + // (context.WithoutCancel), but its values must still propagate. + reqCancel() + select { + case v := <-values: + assert.Equal(t, "trace-abc", v) + case <-time.After(10 * time.Second): + t.Fatal("build did not reach instance creation") + } + select { + case <-dialAttempted: + case <-time.After(10 * time.Second): + t.Fatal("build did not reach the vsock dial loop; request cancellation may have leaked into the run") + } + + running, err := mgr.GetBuild(context.Background(), build.ID) + require.NoError(t, err) + assert.Equal(t, StatusBuilding, running.Status, "request cancellation must not stop the build") + + // Shutdown cancels the run and waits for the goroutine to exit. + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + start := time.Now() + require.NoError(t, mgr.Shutdown(shutdownCtx)) + assert.Less(t, time.Since(start), 5*time.Second, "shutdown should cancel the build promptly") + + finished, err := mgr.GetBuild(context.Background(), build.ID) + require.NoError(t, err) + assert.Equal(t, StatusFailed, finished.Status) +} + +// TestCancelBuild_BuildingBuild verifies cancelling a running build deletes +// its builder instance and marks it cancelled, and that the duplicate +// instance delete from executeBuild's deferred cleanup is tolerated. +func TestCancelBuild_BuildingBuild(t *testing.T) { + mgr, instanceMgr, _, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + instID := "inst-builder-build-1" + instanceMgr.instances[instID] = &instances.Instance{ + StoredMetadata: instances.StoredMetadata{Id: instID, Name: "builder-build-1"}, + State: instances.StateRunning, + } + meta := &buildMetadata{ + ID: "build-1", + Status: StatusBuilding, + CreatedAt: time.Now(), + BuilderInstance: &instID, + } + require.NoError(t, writeMetadata(mgr.paths, meta)) + + require.NoError(t, mgr.CancelBuild(context.Background(), "build-1")) + + got, err := mgr.GetBuild(context.Background(), "build-1") + require.NoError(t, err) + assert.Equal(t, StatusCancelled, got.Status) + assert.Equal(t, 1, instanceMgr.deleteCallCount) + + // executeBuild's deferred cleanup deletes the same instance once the + // build goroutine unwinds; the duplicate delete must be tolerated. + mgr.instanceManager.DeleteInstance(context.Background(), instID) + assert.Equal(t, 2, instanceMgr.deleteCallCount) +} + +// TestCancelBuild_QueuedButAlreadyPickedUp pins the ErrBuildInProgress +// semantics: a build read as queued but already picked up by the queue +// reports the race rather than a success or a confusing error. +func TestCancelBuild_QueuedButAlreadyPickedUp(t *testing.T) { + mgr, _, _, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + meta := &buildMetadata{ + ID: "build-1", + Status: StatusQueued, + CreatedAt: time.Now(), + } + require.NoError(t, writeMetadata(mgr.paths, meta)) + + // The build is not in the pending queue (it was already picked up), so + // the queue cancel misses and the sentinel reports the race. + err := mgr.CancelBuild(context.Background(), "build-1") + assert.ErrorIs(t, err, ErrBuildInProgress) +} + +// TestUpdateBuildComplete_PreservesTerminalStatus pins the cancellation +// terminal-state protection: a cancelled build whose runBuild goroutine +// later fails must keep its cancelled status. +func TestUpdateBuildComplete_PreservesTerminalStatus(t *testing.T) { + mgr, _, _, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + meta := &buildMetadata{ + ID: "build-1", + Status: StatusCancelled, + CreatedAt: time.Now(), + } + require.NoError(t, writeMetadata(mgr.paths, meta)) + + errMsg := "wait for result: context canceled" + duration := int64(42) + mgr.updateBuildComplete("build-1", StatusFailed, nil, &errMsg, nil, &duration) + + got, err := readMetadata(mgr.paths, "build-1") + require.NoError(t, err) + assert.Equal(t, StatusCancelled, got.Status) + assert.Nil(t, got.Error) + assert.Nil(t, got.CompletedAt) +} + func TestCancelBuild_NotFound(t *testing.T) { mgr, _, _, tempDir := setupTestManager(t) defer os.RemoveAll(tempDir) @@ -806,7 +947,7 @@ func TestBuildQueue_ConcurrencyLimit(t *testing.T) { // Enqueue 5 builds with blocking start functions for i := 0; i < 5; i++ { id := string(rune('A' + i)) - queue.Enqueue(id, CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), id, CreateBuildRequest{}, func(context.Context) { started <- id // Block until test completes - simulates long-running build select {} @@ -1163,6 +1304,219 @@ eventLoop: } } +// TestExecuteBuild_SourceVolumeAlreadyExists verifies that a leftover source +// volume from an interrupted prior attempt of the same build is deleted and +// recreated, allowing the re-run to proceed past source volume creation. +func TestExecuteBuild_SourceVolumeAlreadyExists(t *testing.T) { + mgr, instanceMgr, volumeMgr, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + buildID := "build-crash-src" + req := CreateBuildRequest{Dockerfile: "FROM alpine"} + prepareBuildOnDisk(t, mgr, buildID, req) + + var archiveCalls int + var deleted []string + volumeMgr.createFromArchiveFunc = func(ctx context.Context, req volumes.CreateVolumeFromArchiveRequest, archive io.Reader) (*volumes.Volume, error) { + archiveCalls++ + if archiveCalls == 1 { + // Leftover from the interrupted prior attempt + return nil, volumes.ErrAlreadyExists + } + vol := &volumes.Volume{Id: *req.Id, Name: req.Name} + volumeMgr.volumes[vol.Id] = vol + return vol, nil + } + volumeMgr.deleteFunc = func(ctx context.Context, id string) error { + deleted = append(deleted, id) + delete(volumeMgr.volumes, id) + return nil + } + // Short-circuit once volume setup is done so the test doesn't enter the + // vsock wait loop. + instanceMgr.createFunc = func(ctx context.Context, req instances.CreateInstanceRequest) (*instances.Instance, error) { + return nil, fmt.Errorf("stop after volume setup") + } + + policy := DefaultBuildPolicy() + _, err := mgr.executeBuild(context.Background(), buildID, req, &policy) + + // The build must get past source volume creation (it may fail later). + require.Error(t, err) + assert.NotContains(t, err.Error(), "create source volume") + assert.Equal(t, 2, archiveCalls, "expected delete + retry of source volume creation") + require.NotEmpty(t, deleted) + assert.Equal(t, "build-source-"+buildID, deleted[0]) +} + +// TestExecuteBuild_SourceVolumeInUse_StaleBuilder verifies that when the +// leftover source volume is still attached to the interrupted attempt's stale +// builder instance, the builder is deleted (detaching the volume) before the +// volume is deleted and recreated. +func TestExecuteBuild_SourceVolumeInUse_StaleBuilder(t *testing.T) { + mgr, instanceMgr, volumeMgr, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + buildID := "build-crash-inuse" + req := CreateBuildRequest{Dockerfile: "FROM alpine"} + prepareBuildOnDisk(t, mgr, buildID, req) + + // Seed the stale builder instance from the interrupted attempt + builderName := "builder-" + buildID + instanceMgr.instances[builderName] = &instances.Instance{ + StoredMetadata: instances.StoredMetadata{Id: builderName, Name: builderName}, + State: instances.StateRunning, + } + + var archiveCalls int + var deleteCalls int + volumeMgr.createFromArchiveFunc = func(ctx context.Context, req volumes.CreateVolumeFromArchiveRequest, archive io.Reader) (*volumes.Volume, error) { + archiveCalls++ + if archiveCalls == 1 { + return nil, volumes.ErrAlreadyExists + } + vol := &volumes.Volume{Id: *req.Id, Name: req.Name} + volumeMgr.volumes[vol.Id] = vol + return vol, nil + } + volumeMgr.deleteFunc = func(ctx context.Context, id string) error { + deleteCalls++ + if deleteCalls == 1 { + // Still attached to the stale builder + return volumes.ErrInUse + } + delete(volumeMgr.volumes, id) + return nil + } + // Short-circuit once volume setup is done so the test doesn't enter the + // vsock wait loop. + instanceMgr.createFunc = func(ctx context.Context, req instances.CreateInstanceRequest) (*instances.Instance, error) { + return nil, fmt.Errorf("stop after volume setup") + } + + policy := DefaultBuildPolicy() + _, err := mgr.executeBuild(context.Background(), buildID, req, &policy) + + require.Error(t, err) + assert.NotContains(t, err.Error(), "create source volume") + assert.Equal(t, 2, archiveCalls, "expected delete + retry of source volume creation") + assert.Equal(t, 1, instanceMgr.deleteCallCount, "expected stale builder instance to be deleted") + _, getErr := instanceMgr.GetInstance(context.Background(), builderName) + assert.ErrorIs(t, getErr, instances.ErrNotFound, "stale builder instance should be gone") +} + +// TestExecuteBuild_SourceVolumeInUse_NoStaleBuilder verifies that when the +// leftover source volume is in use but the interrupted attempt's builder +// instance is gone, the build fails with a clear error and nothing is +// force-deleted. +func TestExecuteBuild_SourceVolumeInUse_NoStaleBuilder(t *testing.T) { + mgr, instanceMgr, volumeMgr, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + buildID := "build-crash-orphan" + req := CreateBuildRequest{Dockerfile: "FROM alpine"} + prepareBuildOnDisk(t, mgr, buildID, req) + + var archiveCalls int + volumeMgr.createFromArchiveFunc = func(ctx context.Context, req volumes.CreateVolumeFromArchiveRequest, archive io.Reader) (*volumes.Volume, error) { + archiveCalls++ + return nil, volumes.ErrAlreadyExists + } + volumeMgr.deleteFunc = func(ctx context.Context, id string) error { + return volumes.ErrInUse + } + + policy := DefaultBuildPolicy() + _, err := mgr.executeBuild(context.Background(), buildID, req, &policy) + + require.Error(t, err) + assert.Contains(t, err.Error(), "refusing to force-delete") + assert.Equal(t, 1, archiveCalls, "create must not be retried when the leftover cannot be removed") + assert.Equal(t, 0, instanceMgr.deleteCallCount, "no instance should be deleted") +} + +// TestRegisterBuildConfigVolume_AlreadyExists verifies that a leftover config +// volume from an interrupted prior attempt of the same build is explicitly +// deleted and recreated rather than silently masked by the copy-over fallback. +func TestRegisterBuildConfigVolume_AlreadyExists(t *testing.T) { + mgr, _, volumeMgr, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + buildID := "build-crash-config" + configVolID := "build-config-" + buildID + + // Dummy config disk to copy over the recreated volume + configData := []byte("fake-ext4-config-disk") + configDiskPath := filepath.Join(tempDir, "config.ext4") + require.NoError(t, os.WriteFile(configDiskPath, configData, 0644)) + + var createCalls int + var deleted []string + volumeMgr.createFunc = func(ctx context.Context, req volumes.CreateVolumeRequest) (*volumes.Volume, error) { + createCalls++ + if createCalls == 1 { + // Leftover from the interrupted prior attempt + return nil, volumes.ErrAlreadyExists + } + vol := &volumes.Volume{Id: *req.Id, Name: req.Name} + volumeMgr.volumes[vol.Id] = vol + return vol, nil + } + volumeMgr.deleteFunc = func(ctx context.Context, id string) error { + deleted = append(deleted, id) + delete(volumeMgr.volumes, id) + return nil + } + + err := mgr.registerBuildConfigVolume(context.Background(), buildID, configVolID, configDiskPath) + + require.NoError(t, err) + assert.Equal(t, 2, createCalls, "expected delete + retry of config volume creation") + assert.Equal(t, []string{configVolID}, deleted) + assert.Contains(t, volumeMgr.volumes, configVolID, "config volume should be recreated") + // The config data must have been copied onto the recreated volume + copied, readErr := os.ReadFile(mgr.paths.VolumeData(configVolID)) + require.NoError(t, readErr) + assert.Equal(t, configData, copied) +} + +// TestRegisterBuildConfigVolume_RecreateFailure verifies that when recreating +// a deleted leftover config volume fails, the recreate error is surfaced +// rather than silently masked by the copy-over fallback. +func TestRegisterBuildConfigVolume_RecreateFailure(t *testing.T) { + mgr, _, volumeMgr, tempDir := setupTestManager(t) + defer os.RemoveAll(tempDir) + + buildID := "build-crash-config-fail" + configVolID := "build-config-" + buildID + + configDiskPath := filepath.Join(tempDir, "config.ext4") + require.NoError(t, os.WriteFile(configDiskPath, []byte("fake-ext4-config-disk"), 0644)) + + var createCalls int + recreateErr := fmt.Errorf("recreate failed") + volumeMgr.createFunc = func(ctx context.Context, req volumes.CreateVolumeRequest) (*volumes.Volume, error) { + createCalls++ + if createCalls == 1 { + return nil, volumes.ErrAlreadyExists + } + return nil, recreateErr + } + volumeMgr.deleteFunc = func(ctx context.Context, id string) error { + delete(volumeMgr.volumes, id) + return nil + } + + err := mgr.registerBuildConfigVolume(context.Background(), buildID, configVolID, configDiskPath) + + require.Error(t, err) + assert.ErrorIs(t, err, recreateErr) + assert.Contains(t, err.Error(), "create config volume") + assert.Equal(t, 2, createCalls, "expected delete + retry of config volume creation") + _, statErr := os.Stat(mgr.paths.VolumeData(configVolID)) + assert.ErrorIs(t, statErr, os.ErrNotExist, "config volume data must not be copied when recreate fails") +} + func TestExtractInternalBaseImageRepos(t *testing.T) { registryURL := "http://10.102.0.1:8085" diff --git a/lib/builds/queue.go b/lib/builds/queue.go index 4e0b356a9..176ad829a 100644 --- a/lib/builds/queue.go +++ b/lib/builds/queue.go @@ -1,6 +1,13 @@ package builds -import "sync" +import ( + "context" + "sync" +) + +// StartFunc executes a queued build. The context carries the enqueue-time +// values (e.g. trace context) and is cancelled when the queue shuts down. +type StartFunc func(ctx context.Context) // QueuedBuild represents a build waiting to be executed type QueuedBuild struct { @@ -11,7 +18,11 @@ type QueuedBuild struct { // once no active build holds the same key, so it never occupies a // concurrency slot while waiting on another build of its group. SerialKey string - StartFn func() + StartFn StartFunc + // ctx carries the enqueue-time context values into the run; cancel + // releases it and aborts the run on shutdown. + ctx context.Context + cancel context.CancelFunc } // BuildQueue manages concurrent builds with a configurable limit. @@ -29,8 +40,11 @@ type QueuedBuild struct { type BuildQueue struct { maxConcurrent int active map[string]bool - activeSerialKeys map[string]string // build ID -> serial key + activeSerialKeys map[string]string // build ID -> serial key + activeCancels map[string]context.CancelFunc // build ID -> run cancel pending []QueuedBuild + wg sync.WaitGroup // tracks active build goroutines + shutdown bool mu sync.Mutex } @@ -43,6 +57,7 @@ func NewBuildQueue(maxConcurrent int) *BuildQueue { maxConcurrent: maxConcurrent, active: make(map[string]bool), activeSerialKeys: make(map[string]string), + activeCancels: make(map[string]context.CancelFunc), pending: make([]QueuedBuild, 0), } } @@ -50,14 +65,19 @@ func NewBuildQueue(maxConcurrent int) *BuildQueue { // Enqueue adds a build to the queue. Returns queue position (0 if started // immediately, >0 if queued among pending builds that can currently start). // If the build is already building or queued, returns its current position without re-enqueueing. -func (q *BuildQueue) Enqueue(buildID string, req CreateBuildRequest, startFn func()) int { - return q.EnqueueSerial(buildID, req, "", startFn) +// +// The ctx supplies values (trace context, logger) to the run and bounds its +// lifetime together with Shutdown: the run is cancelled when ctx is +// cancelled or the queue shuts down. Callers that want the build to survive +// request cancellation should pass context.WithoutCancel(ctx). +func (q *BuildQueue) Enqueue(ctx context.Context, buildID string, req CreateBuildRequest, startFn StartFunc) int { + return q.EnqueueSerial(ctx, buildID, req, "", startFn) } // EnqueueSerial is Enqueue with a serial key: builds sharing a non-empty key // are serialized so a waiting build stays pending instead of occupying a // concurrency slot while blocked on another build of its group. -func (q *BuildQueue) EnqueueSerial(buildID string, req CreateBuildRequest, serialKey string, startFn func()) int { +func (q *BuildQueue) EnqueueSerial(ctx context.Context, buildID string, req CreateBuildRequest, serialKey string, startFn StartFunc) int { q.mu.Lock() defer q.mu.Unlock() @@ -73,17 +93,14 @@ func (q *BuildQueue) EnqueueSerial(buildID string, req CreateBuildRequest, seria } } - // Wrap the function to auto-complete - wrappedFn := func() { - defer q.MarkComplete(buildID) - startFn() - } - + runCtx, cancel := context.WithCancel(ctx) build := QueuedBuild{ BuildID: buildID, Request: req, SerialKey: serialKey, - StartFn: wrappedFn, + StartFn: startFn, + ctx: runCtx, + cancel: cancel, } // Start immediately if under concurrency limit and the serial key is free @@ -92,13 +109,17 @@ func (q *BuildQueue) EnqueueSerial(buildID string, req CreateBuildRequest, seria return 0 } - // Otherwise queue it + // Otherwise queue it. After Shutdown the build stays pending and never + // starts; its on-disk metadata lets startup recovery re-enqueue it. q.pending = append(q.pending, build) return *q.pendingPositionLocked(buildID) } // canStartLocked reports whether a build with the given serial key may start. func (q *BuildQueue) canStartLocked(serialKey string) bool { + if q.shutdown { + return false + } if len(q.active) >= q.maxConcurrent { return false } @@ -111,7 +132,14 @@ func (q *BuildQueue) startLocked(build QueuedBuild) { if build.SerialKey != "" { q.activeSerialKeys[build.BuildID] = build.SerialKey } - go build.StartFn() + q.activeCancels[build.BuildID] = build.cancel + q.wg.Add(1) + go func() { + defer q.wg.Done() + defer build.cancel() + defer q.MarkComplete(build.BuildID) + build.StartFn(build.ctx) + }() } // MarkComplete marks a build as complete and starts pending builds while @@ -123,6 +151,7 @@ func (q *BuildQueue) MarkComplete(buildID string) { delete(q.active, buildID) delete(q.activeSerialKeys, buildID) + delete(q.activeCancels, buildID) q.startPendingLocked() } @@ -139,7 +168,7 @@ func (q *BuildQueue) ReleaseSerialKey(buildID string) { key, held := q.activeSerialKeys[buildID] delete(q.activeSerialKeys, buildID) - if held { + if held && !q.shutdown { for i, build := range q.pending { if build.SerialKey == key { q.pending = append(q.pending[:i], q.pending[i+1:]...) @@ -154,6 +183,9 @@ func (q *BuildQueue) ReleaseSerialKey(buildID string) { // startPendingLocked starts pending builds while capacity and serial keys // allow. func (q *BuildQueue) startPendingLocked() { + if q.shutdown { + return + } for len(q.active) < q.maxConcurrent { idx := -1 for i, build := range q.pending { @@ -226,6 +258,7 @@ func (q *BuildQueue) Cancel(buildID string) bool { for i, build := range q.pending { if build.BuildID == buildID { q.pending = append(q.pending[:i], q.pending[i+1:]...) + build.cancel() return true } } @@ -233,6 +266,43 @@ func (q *BuildQueue) Cancel(buildID string) bool { return false } +// Shutdown stops scheduling new work, cancels every active build's run +// context, and waits for their goroutines to return. Pending builds are left +// unstarted; their persisted metadata lets startup recovery re-enqueue them +// on the next boot. Returns ctx.Err() if the wait outlives ctx. +func (q *BuildQueue) Shutdown(ctx context.Context) error { + q.mu.Lock() + if q.shutdown { + q.mu.Unlock() + return nil + } + q.shutdown = true + for _, build := range q.pending { + build.cancel() + } + cancels := make([]context.CancelFunc, 0, len(q.activeCancels)) + for _, cancel := range q.activeCancels { + cancels = append(cancels, cancel) + } + q.mu.Unlock() + + for _, cancel := range cancels { + cancel() + } + + done := make(chan struct{}) + go func() { + q.wg.Wait() + close(done) + }() + select { + case <-done: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + // IsActive returns true if the build is actively running func (q *BuildQueue) IsActive(buildID string) bool { q.mu.Lock() diff --git a/lib/builds/queue_test.go b/lib/builds/queue_test.go index 8c05270e5..2def0729c 100644 --- a/lib/builds/queue_test.go +++ b/lib/builds/queue_test.go @@ -1,6 +1,8 @@ package builds import ( + "context" + "fmt" "sync" "testing" "time" @@ -16,7 +18,7 @@ func TestBuildQueue_EnqueueStartsImmediately(t *testing.T) { done := make(chan struct{}) // Enqueue first build - should start immediately - pos := queue.Enqueue("build-1", CreateBuildRequest{}, func() { + pos := queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { started <- "build-1" <-done // Wait for signal }) @@ -42,7 +44,7 @@ func TestBuildQueue_QueueWhenAtCapacity(t *testing.T) { // Start first build wg.Add(1) - pos1 := queue.Enqueue("build-1", CreateBuildRequest{}, func() { + pos1 := queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { wg.Done() <-done // Block }) @@ -52,11 +54,11 @@ func TestBuildQueue_QueueWhenAtCapacity(t *testing.T) { wg.Wait() // Second build should be queued - pos2 := queue.Enqueue("build-2", CreateBuildRequest{}, func() {}) + pos2 := queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) {}) assert.Equal(t, 1, pos2, "second build should be queued at position 1") // Third build should be queued at position 2 - pos3 := queue.Enqueue("build-3", CreateBuildRequest{}, func() {}) + pos3 := queue.Enqueue(context.Background(), "build-3", CreateBuildRequest{}, func(context.Context) {}) assert.Equal(t, 2, pos3, "third build should be queued at position 2") close(done) @@ -67,7 +69,7 @@ func TestBuildQueue_DeduplicationActive(t *testing.T) { done := make(chan struct{}) // Start a build - queue.Enqueue("build-1", CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { <-done }) @@ -75,7 +77,7 @@ func TestBuildQueue_DeduplicationActive(t *testing.T) { time.Sleep(10 * time.Millisecond) // Try to enqueue the same build again - should return position 0 (active) - pos := queue.Enqueue("build-1", CreateBuildRequest{}, func() {}) + pos := queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) {}) assert.Equal(t, 0, pos, "re-enqueueing active build should return position 0") close(done) @@ -86,16 +88,16 @@ func TestBuildQueue_DeduplicationPending(t *testing.T) { done := make(chan struct{}) // Fill the queue - queue.Enqueue("build-1", CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { <-done }) // Add a second build to pending - pos1 := queue.Enqueue("build-2", CreateBuildRequest{}, func() {}) + pos1 := queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) {}) assert.Equal(t, 1, pos1) // Try to enqueue build-2 again - should return same position - pos2 := queue.Enqueue("build-2", CreateBuildRequest{}, func() {}) + pos2 := queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) {}) assert.Equal(t, 1, pos2, "re-enqueueing pending build should return same position") close(done) @@ -106,13 +108,13 @@ func TestBuildQueue_Cancel(t *testing.T) { done := make(chan struct{}) // Fill the queue - queue.Enqueue("build-1", CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { <-done }) // Add to pending - queue.Enqueue("build-2", CreateBuildRequest{}, func() {}) - queue.Enqueue("build-3", CreateBuildRequest{}, func() {}) + queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) {}) + queue.Enqueue(context.Background(), "build-3", CreateBuildRequest{}, func(context.Context) {}) // Cancel build-2 cancelled := queue.Cancel("build-2") @@ -134,11 +136,11 @@ func TestBuildQueue_GetPosition(t *testing.T) { queue := NewBuildQueue(1) done := make(chan struct{}) - queue.Enqueue("build-1", CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { <-done }) - queue.Enqueue("build-2", CreateBuildRequest{}, func() {}) - queue.Enqueue("build-3", CreateBuildRequest{}, func() {}) + queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) {}) + queue.Enqueue(context.Background(), "build-3", CreateBuildRequest{}, func(context.Context) {}) // Active build has no position (returns nil) pos1 := queue.GetPosition("build-1") @@ -168,14 +170,14 @@ func TestBuildQueue_AutoStartNextOnComplete(t *testing.T) { completionOrder := []string{} // Add builds - queue.Enqueue("build-1", CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { started <- "build-1" time.Sleep(10 * time.Millisecond) mu.Lock() completionOrder = append(completionOrder, "build-1") mu.Unlock() }) - queue.Enqueue("build-2", CreateBuildRequest{}, func() { + queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) { started <- "build-2" time.Sleep(10 * time.Millisecond) mu.Lock() @@ -208,8 +210,8 @@ func TestBuildQueue_Counts(t *testing.T) { assert.Equal(t, 0, queue.QueueLength()) done := make(chan struct{}) - queue.Enqueue("build-1", CreateBuildRequest{}, func() { <-done }) - queue.Enqueue("build-2", CreateBuildRequest{}, func() { <-done }) + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(context.Context) { <-done }) + queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(context.Context) { <-done }) // Wait for them to start time.Sleep(10 * time.Millisecond) @@ -219,7 +221,7 @@ func TestBuildQueue_Counts(t *testing.T) { assert.Equal(t, 2, queue.QueueLength()) // Add a pending one - queue.Enqueue("build-3", CreateBuildRequest{}, func() {}) + queue.Enqueue(context.Background(), "build-3", CreateBuildRequest{}, func(context.Context) {}) assert.Equal(t, 2, queue.ActiveCount()) assert.Equal(t, 1, queue.PendingCount()) @@ -233,8 +235,8 @@ func TestBuildQueue_SerialKeySerializesSameKey(t *testing.T) { release := make(chan struct{}) started := make(chan string, 3) - startFn := func(id string) func() { - return func() { + startFn := func(id string) func(context.Context) { + return func(context.Context) { started <- id <-release } @@ -242,8 +244,8 @@ func TestBuildQueue_SerialKeySerializesSameKey(t *testing.T) { // Two builds with the same serial key: only the first starts even though // a concurrency slot is free. - pos1 := queue.EnqueueSerial("build-1", CreateBuildRequest{}, "builder-a", startFn("build-1")) - pos2 := queue.EnqueueSerial("build-2", CreateBuildRequest{}, "builder-a", startFn("build-2")) + pos1 := queue.EnqueueSerial(context.Background(), "build-1", CreateBuildRequest{}, "builder-a", startFn("build-1")) + pos2 := queue.EnqueueSerial(context.Background(), "build-2", CreateBuildRequest{}, "builder-a", startFn("build-2")) assert.Equal(t, 0, pos1) assert.Equal(t, 1, pos2) @@ -258,7 +260,7 @@ func TestBuildQueue_SerialKeySerializesSameKey(t *testing.T) { assert.Equal(t, 1, queue.PendingCount()) // A build with a different key starts immediately in the free slot. - pos3 := queue.EnqueueSerial("build-3", CreateBuildRequest{}, "builder-b", startFn("build-3")) + pos3 := queue.EnqueueSerial(context.Background(), "build-3", CreateBuildRequest{}, "builder-b", startFn("build-3")) assert.Equal(t, 0, pos3) select { case id := <-started: @@ -288,17 +290,17 @@ func TestBuildQueue_SerialKeySkipsBlockedPending(t *testing.T) { release[id] = make(chan struct{}) } started := make(chan string, 4) - startFn := func(id string) func() { - return func() { + startFn := func(id string) func(context.Context) { + return func(context.Context) { started <- id <-release[id] } } - queue.EnqueueSerial("build-1", CreateBuildRequest{}, "builder-a", startFn("build-1")) - queue.EnqueueSerial("build-2", CreateBuildRequest{}, "builder-b", startFn("build-2")) - pos3 := queue.EnqueueSerial("build-3", CreateBuildRequest{}, "builder-a", startFn("build-3")) - pos4 := queue.EnqueueSerial("build-4", CreateBuildRequest{}, "builder-c", startFn("build-4")) + queue.EnqueueSerial(context.Background(), "build-1", CreateBuildRequest{}, "builder-a", startFn("build-1")) + queue.EnqueueSerial(context.Background(), "build-2", CreateBuildRequest{}, "builder-b", startFn("build-2")) + pos3 := queue.EnqueueSerial(context.Background(), "build-3", CreateBuildRequest{}, "builder-a", startFn("build-3")) + pos4 := queue.EnqueueSerial(context.Background(), "build-4", CreateBuildRequest{}, "builder-c", startFn("build-4")) assert.Equal(t, 1, pos3) assert.Equal(t, 2, pos4, "queue positions retain submission order") @@ -344,15 +346,15 @@ func TestBuildQueue_ReleaseSerialKeyStartsSameKeyPending(t *testing.T) { release := make(chan struct{}) started := make(chan string, 2) - startFn := func(id string) func() { - return func() { + startFn := func(id string) func(context.Context) { + return func(context.Context) { started <- id <-release } } - queue.EnqueueSerial("build-1", CreateBuildRequest{}, "builder-a", startFn("build-1")) - queue.EnqueueSerial("build-2", CreateBuildRequest{}, "builder-a", startFn("build-2")) + queue.EnqueueSerial(context.Background(), "build-1", CreateBuildRequest{}, "builder-a", startFn("build-1")) + queue.EnqueueSerial(context.Background(), "build-2", CreateBuildRequest{}, "builder-a", startFn("build-2")) select { case id := <-started: @@ -380,11 +382,11 @@ func TestBuildQueue_SerialKeyIntrospection(t *testing.T) { queue := NewBuildQueue(1) release := make(chan struct{}) - startFn := func() { <-release } + startFn := func(context.Context) { <-release } - queue.EnqueueSerial("build-1", CreateBuildRequest{}, "builder-a", startFn) - queue.EnqueueSerial("build-2", CreateBuildRequest{}, "builder-a", startFn) - queue.EnqueueSerial("build-3", CreateBuildRequest{}, "builder-b", startFn) + queue.EnqueueSerial(context.Background(), "build-1", CreateBuildRequest{}, "builder-a", startFn) + queue.EnqueueSerial(context.Background(), "build-2", CreateBuildRequest{}, "builder-a", startFn) + queue.EnqueueSerial(context.Background(), "build-3", CreateBuildRequest{}, "builder-b", startFn) active := queue.ActiveBuildForSerialKey("builder-a") require.NotNil(t, active) @@ -410,12 +412,12 @@ func TestBuildQueue_HasSerialKeyAcrossPendingToActiveTransition(t *testing.T) { release := make(chan struct{}) started := make(chan string, 2) - queue.EnqueueSerial("build-1", CreateBuildRequest{}, "builder-a", func() { + queue.EnqueueSerial(context.Background(), "build-1", CreateBuildRequest{}, "builder-a", func(context.Context) { started <- "build-1" <-release }) require.Equal(t, "build-1", <-started) - queue.EnqueueSerial("build-2", CreateBuildRequest{}, "builder-a", func() { + queue.EnqueueSerial(context.Background(), "build-2", CreateBuildRequest{}, "builder-a", func(context.Context) { started <- "build-2" <-release }) @@ -430,6 +432,121 @@ func TestBuildQueue_HasSerialKeyAcrossPendingToActiveTransition(t *testing.T) { } } +// TestBuildQueue_ShutdownCancelsAndAwaitsActiveBuilds verifies Shutdown +// cancels every active build's run context and blocks until their goroutines +// have returned. +func TestBuildQueue_ShutdownCancelsAndAwaitsActiveBuilds(t *testing.T) { + queue := NewBuildQueue(2) + + cancelled := make(chan string, 2) + for _, id := range []string{"build-1", "build-2"} { + queue.Enqueue(context.Background(), id, CreateBuildRequest{}, func(ctx context.Context) { + <-ctx.Done() + cancelled <- id + }) + } + require.Equal(t, 2, queue.ActiveCount()) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + require.NoError(t, queue.Shutdown(ctx)) + + // Shutdown only returns after both goroutines exited. + assert.ElementsMatch(t, []string{"build-1", "build-2"}, []string{<-cancelled, <-cancelled}) + assert.Equal(t, 0, queue.ActiveCount()) +} + +// TestBuildQueue_ShutdownLeavesPendingForRecovery verifies pending builds are +// not started during shutdown; they stay queued so startup recovery can +// re-enqueue them from disk. +func TestBuildQueue_ShutdownLeavesPendingForRecovery(t *testing.T) { + queue := NewBuildQueue(1) + + started := make(chan string, 2) + queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(ctx context.Context) { + started <- "build-1" + <-ctx.Done() + }) + require.Equal(t, "build-1", <-started) + + queue.Enqueue(context.Background(), "build-2", CreateBuildRequest{}, func(ctx context.Context) { + started <- "build-2" + }) + require.Equal(t, 1, queue.PendingCount()) + + require.NoError(t, queue.Shutdown(context.Background())) + + select { + case id := <-started: + t.Fatalf("pending build %s started during shutdown", id) + default: + } + pos := queue.GetPosition("build-2") + require.NotNil(t, pos, "pending build must stay queued for recovery") + assert.Equal(t, 1, *pos) +} + +// TestBuildQueue_EnqueueAfterShutdownStaysPending verifies work submitted +// after shutdown never starts. +func TestBuildQueue_EnqueueAfterShutdownStaysPending(t *testing.T) { + queue := NewBuildQueue(1) + require.NoError(t, queue.Shutdown(context.Background())) + + ran := make(chan struct{}, 1) + pos := queue.Enqueue(context.Background(), "build-1", CreateBuildRequest{}, func(ctx context.Context) { + ran <- struct{}{} + }) + assert.Equal(t, 1, pos, "build enqueued after shutdown is pending, not started") + select { + case <-ran: + t.Fatal("build started after shutdown") + default: + } +} + +// TestBuildQueue_ShutdownIsIdempotent verifies a second Shutdown returns +// immediately without re-cancelling or deadlocking. +func TestBuildQueue_ShutdownIsIdempotent(t *testing.T) { + queue := NewBuildQueue(1) + require.NoError(t, queue.Shutdown(context.Background())) + require.NoError(t, queue.Shutdown(context.Background())) +} + +// TestBuildQueue_ShutdownRace stresses Enqueue/Cancel/ReleaseSerialKey racing +// Shutdown; run with -race. +func TestBuildQueue_ShutdownRace(t *testing.T) { + for i := 0; i < 50; i++ { + queue := NewBuildQueue(2) + var wg sync.WaitGroup + for j := 0; j < 4; j++ { + id := fmt.Sprintf("build-%d", j) + wg.Add(1) + go func() { + defer wg.Done() + queue.EnqueueSerial(context.Background(), id, CreateBuildRequest{}, "builder-a", func(ctx context.Context) { + select { + case <-ctx.Done(): + case <-time.After(20 * time.Millisecond): + } + }) + }() + } + wg.Add(1) + go func() { + defer wg.Done() + queue.Shutdown(context.Background()) + }() + wg.Add(1) + go func() { + defer wg.Done() + queue.Cancel("build-2") + queue.ReleaseSerialKey("build-1") + }() + wg.Wait() + require.NoError(t, queue.Shutdown(context.Background())) + } +} + // TestBuildQueue_ReleaseSerialKeyStartsSuccessorAtFullCapacity releases the // serial key while the build keeps its global slot (post-build work) and // verifies the serialized successor starts anyway: at MaxConcurrentBuilds=1 @@ -440,14 +557,14 @@ func TestBuildQueue_ReleaseSerialKeyStartsSuccessorAtFullCapacity(t *testing.T) release := make(chan struct{}) started := make(chan string, 2) - queue.EnqueueSerial("build-1", CreateBuildRequest{}, "builder-a", func() { + queue.EnqueueSerial(context.Background(), "build-1", CreateBuildRequest{}, "builder-a", func(context.Context) { started <- "build-1" // VM phase done; the build keeps its global slot for post-build // work and only completes after release closes. queue.ReleaseSerialKey("build-1") <-release }) - queue.EnqueueSerial("build-2", CreateBuildRequest{}, "builder-a", func() { + queue.EnqueueSerial(context.Background(), "build-2", CreateBuildRequest{}, "builder-a", func(context.Context) { started <- "build-2" <-release }) diff --git a/lib/builds/storage.go b/lib/builds/storage.go index 2ca771b9b..12995eb0c 100644 --- a/lib/builds/storage.go +++ b/lib/builds/storage.go @@ -3,6 +3,7 @@ package builds import ( "encoding/json" "fmt" + "log/slog" "os" "sort" "time" @@ -98,8 +99,11 @@ func readMetadata(p *paths.Paths, id string) (*buildMetadata, error) { return &meta, nil } -// listAllBuilds returns all builds sorted by creation time (newest first) -func listAllBuilds(p *paths.Paths) ([]*buildMetadata, error) { +// listAllBuilds returns all builds sorted by creation time (newest first). +// Entries whose metadata cannot be read or parsed are corrupt: they are +// logged and skipped so valid builds are still returned, but the corruption +// is never silently hidden. +func listAllBuilds(p *paths.Paths, logger *slog.Logger) ([]*buildMetadata, error) { buildsDir := p.BuildsDir() entries, err := os.ReadDir(buildsDir) @@ -118,7 +122,8 @@ func listAllBuilds(p *paths.Paths) ([]*buildMetadata, error) { meta, err := readMetadata(p, entry.Name()) if err != nil { - continue // Skip invalid entries + logger.Warn("skipping corrupt build metadata", "id", entry.Name(), "error", err) + continue } metas = append(metas, meta) } @@ -133,8 +138,8 @@ func listAllBuilds(p *paths.Paths) ([]*buildMetadata, error) { // listPendingBuilds returns builds that need to be recovered on startup // Returns builds with status queued/building, sorted by created_at (oldest first for FIFO) -func listPendingBuilds(p *paths.Paths) ([]*buildMetadata, error) { - all, err := listAllBuilds(p) +func listPendingBuilds(p *paths.Paths, logger *slog.Logger) ([]*buildMetadata, error) { + all, err := listAllBuilds(p, logger) if err != nil { return nil, err } diff --git a/lib/builds/storage_test.go b/lib/builds/storage_test.go index e1e821826..2a744f7a6 100644 --- a/lib/builds/storage_test.go +++ b/lib/builds/storage_test.go @@ -1,12 +1,15 @@ package builds import ( + "bytes" + "log/slog" "os" "path/filepath" "testing" "time" "github.com/kernel/hypeman/lib/paths" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -36,3 +39,47 @@ func TestBuildMetadataReadWrite_MetadataRoundTrip(t *testing.T) { loaded.Tags["team"] = "mutated" require.Equal(t, "backend", build.Tags["team"]) } + +func TestListAllBuilds_LogsAndSkipsCorruptMetadata(t *testing.T) { + tempDir := t.TempDir() + p := paths.New(tempDir) + + // One valid build and two corrupt entries: unparsable JSON and a + // directory with no metadata file at all. + valid := &buildMetadata{ID: "valid-build", Status: StatusReady, CreatedAt: time.Now()} + require.NoError(t, writeMetadata(p, valid)) + require.NoError(t, os.MkdirAll(p.BuildDir("corrupt-json"), 0755)) + require.NoError(t, os.WriteFile(p.BuildMetadata("corrupt-json"), []byte("{not json"), 0644)) + require.NoError(t, os.MkdirAll(p.BuildDir("missing-file"), 0755)) + + var logBuf bytes.Buffer + logger := slog.New(slog.NewTextHandler(&logBuf, nil)) + + metas, err := listAllBuilds(p, logger) + require.NoError(t, err) + require.Len(t, metas, 1) + assert.Equal(t, "valid-build", metas[0].ID) + + out := logBuf.String() + assert.Contains(t, out, "corrupt-json") + assert.Contains(t, out, "missing-file") +} + +func TestListPendingBuilds_LogsAndSkipsCorruptMetadata(t *testing.T) { + tempDir := t.TempDir() + p := paths.New(tempDir) + + pending := &buildMetadata{ID: "pending-build", Status: StatusQueued, CreatedAt: time.Now()} + require.NoError(t, writeMetadata(p, pending)) + require.NoError(t, os.MkdirAll(p.BuildDir("corrupt-json"), 0755)) + require.NoError(t, os.WriteFile(p.BuildMetadata("corrupt-json"), []byte("{not json"), 0644)) + + var logBuf bytes.Buffer + logger := slog.New(slog.NewTextHandler(&logBuf, nil)) + + metas, err := listPendingBuilds(p, logger) + require.NoError(t, err) + require.Len(t, metas, 1) + assert.Equal(t, "pending-build", metas[0].ID) + assert.Contains(t, logBuf.String(), "corrupt-json") +} diff --git a/lib/builds/types.go b/lib/builds/types.go index 4ad6008e9..66df16241 100644 --- a/lib/builds/types.go +++ b/lib/builds/types.go @@ -97,9 +97,6 @@ type BuildPolicy struct { // NetworkMode controls network access during build // "isolated" = no network, "egress" = outbound allowed NetworkMode string `json:"network_mode,omitempty"` - - // AllowedDomains restricts egress to specific domains (only when NetworkMode="egress") - AllowedDomains []string `json:"allowed_domains,omitempty"` } // SecretRef references a secret to inject during build diff --git a/lib/builds/vsock_handler.go b/lib/builds/vsock_handler.go index f5aebfcd9..2b9136e68 100644 --- a/lib/builds/vsock_handler.go +++ b/lib/builds/vsock_handler.go @@ -19,6 +19,7 @@ type VsockMessage struct { Log string `json:"log,omitempty"` SecretIDs []string `json:"secret_ids,omitempty"` // For secrets request Secrets map[string]string `json:"secrets,omitempty"` // For secrets response + Error string `json:"error,omitempty"` // Set on a secrets response when the host could not provide every requested secret } // SecretsRequest is sent by the builder agent to fetch secrets diff --git a/lib/oapi/oapi.go b/lib/oapi/oapi.go index 4578c7ae7..dcfdb91fb 100644 --- a/lib/oapi/oapi.go +++ b/lib/oapi/oapi.go @@ -1762,7 +1762,8 @@ type CreateBuildMultipartBody struct { // Example: [{"id": "npm_token"}, {"id": "github_token"}] Secrets *string `json:"secrets,omitempty"` - // Source Source tarball (tar.gz) containing application code and optionally a Dockerfile + // Source Source tarball (tar.gz) containing application code and optionally a Dockerfile. + // Maximum size is 512 MiB; other form fields are limited to 1 MiB each. Source openapi_types.File `json:"source"` // Tags JSON object of tags. @@ -17851,92 +17852,93 @@ var swaggerSpec = []string{ "o/+eCswUjAkUQwnq9bS4pGsZgLHepOtVCgeYoYQnBvwZmGqGNcuV2gAAABxFSMFWct9qxR1WsmE+Fs8t", "HjeCuRn0rco2ogydPC9spsHOU/9+kiQQxBdR8p9nb14jOJX1GpjX8gghk0XAtMKAwhSuTp1Me4GDGTIX", "VVDgZ9ih4bCTXeeG6zDWVNqU+F4P7hR/0UP7xXTTpeEv/b5uylxX7qM/PptW9vVeSuKR4heEDTtfuqjw", - "YErVLB1nzz76CdqEiXVWEgRozRxz6yBJMAX4ksKJb45IzELE7SkQLRBGuQQqBq6MKcNisSx3zUN6S0E+", - "McFzBWJ8HkKw3LCzP3ThcsNOd9ghbA6/2Zi6YeeLnwL21rK5mgycZ9nlZsZEe4PB+mqwSktfz51li4uB", - "W7YBG62irBSWXsE/U5J+d/cD/9b2Z3b1g5nuPMe7MYa/c74/wguIgsZetEQ9VxAVtRuzgERO7V7t6Ln/", - "ywO9WAGJovtm0Idiz+x6zNZMfGT3YbBY+TZa6r5/YI4b3NehUnLbPwz/Pjr/ucd7bn3nZO5Cnf1Q4oBz", - "Yk1pZF5GWKIzGFPvTBvfL+DXvv2vs/0AqO884tPzfWO6o4hPUUSZDUEvBCpr9cDSEj4yUCfZdxb5xNVx", - "WTOaxL/+539hUJRN//U//6vtCvMXbPcNg9kFdR/PZwQLNSZYne+j3whJejiic+ImA5XZyJyIBdoeWJ8/", - "PCpW87ZamhyyIXtLVCpYIVTflFSRtkF7VaDnQ1lKpIWK0S/SicV7N7GNHr+N28uGlPe6o7sejD2YQWEC", - "+lR0PAAAZdQUv7SWaMfvMjVzLjlNq2GatWC91fJFkStluLdnBnhNAQMk9u07eGAnjdbOzl6s9xFYW4Yr", - "ANMfbIe8GWtG9H/IpNUyyUiUskABKhvZZPCTljv9j+w77bz+tsXvye1vy7xcw+9vnD8Am+hW4McdQIs7", - "AD/d3H2Azyl/5ADC7i5Y0HTxQLGCjvfqNDdPCiR7CGcAWnNADOBQ5QKdHh4jHIaCSLn+7+0q0DM1XJof", - "HYgzgP1/iFtrOxaoCR6TzFQrM8hjEQdv7agRdvOqVs8qnm8bpWIQjSddVhciP/Lu/vSodHqdYySv8JXz", - "2o+TZGWcHpUB198WuKUX4AQI6dSXbJ8WuWiVQ8pEAGZHzlJ1yYrn4yO3Ie/PNWW7Tln1bLgHoXhUEYgP", - "KAjLWZrFmniPiZvfZ6vowFCXeK6+LdYc3J8WdN9eLB+bPyY3Vlghm5aCBk+g8QB9RZRBEejc4ULbHjwT", - "PyPC7WpXwhRmnU3LfIoMHAJMCK7ml9u+x+aVdqavae97snyBPNfRWCzJf6goLYzdnFbLDFyzBHdp30IP", - "1zJvb+/G2zKYh8gQdjN2HmuhSIjWsFywYP3Hpfetc7QJicqNWIGywssoibCC6EgAYsnsLD22rXvQ697a", - "OCoksCI2bugxJuOdplHkrmbmRCj05vDYiIDiYbXxGSLJVhshTiwsPbfev/29R1jAIXQwC3vza3v2yS2b", - "IoazSgl298/PjzDJjLqDt0kV+4r1NxGeyASl9in/j62XER0LLBb/sfUSRwll5D+2DyKsiFTrd8Ysg/s6", - "Q+7bNHjEzKctA1omGogmNgXowBWqdPZWS23avf9dKdRm0tdSqTO6/tCq22jVRXItVaztUtypam36eKC7", - "o4zZfNSGRz9wJu7BHWk5soAzUbqfyZEmZlwqePT4kg5tpCfNOK54bLT0q+cbcunx4Vj3+KgLhITaoYBs", - "bnN67snL7sZx78qt7ff+XewH8ZhOU57KYrpQjFUwI9Km0kWkLIAfm9qdH8+Nivc3zKWD+zw67l2v/sH3", - "d6TxVxfUCG9zVbZK53dvtdX57fta5zcwgzbd0MKvd11pjvWG6EcHNNiWjUt4jPWoTN+4fLYIeq8Nldxc", - "QGBB7A/Z/9H2xx+K4PjjLy6vKR0Mtvbgd8LmH39xqU3sxLEKYVAVHFJcD14fwf3kFBAaodhSnkVZHYep", - "zgqs5+Cl/+0MpPyKtr2F5Ljwh4XUykIqkGu5hWTX4m5NpDJE/b3bSI7ffAS3QL/fp5X0DV883LsFJ9PJ", - "hAaUMAD6h2xRWYu0M5bcj5uRG2YJMnvTVwjTKWkirc3ITGqt0NDz2qL3HqJ1nBdTuW/r0ZUxfZzZDjyx", - "dQGtvZZrC80G27fGD4P7Pb3u31B7zCxmLKI66RKtdHvKdZhCNXGqILw0B/iB+F0kjFmTtdhHh1let0yT", - "hAslTbEbsBBMOcyZthB8hXHKtW58xW2goAslsjtkUO5UPzb4FBsXZGFK2VDOsqo12UxtRRhfFl25lNCD", - "bqPbV0L9dZJaKaH3vI1t5buHU0IfTHTci7p3XCooupZtDLC4xyTbyTxL06SfKJuuP6pYYiOssrkV4Mg8", - "qtYGThV3dfA3ZtwgE/nB2U4jHAA2m37NwAbZvF+DE1ZsCpJ5BY8iIgwcVJIqVzdryLLBUVaoC2yLVJzr", - "5kcpUzQ675poGsjplwizhcVEGbJSZ1gpEidasFmUHxihIIkZcaVgmB405amEt7pI8lKXCEeXeCGHTJBJ", - "RAI7NyiuKEhgkNOiqI9+5ZBIjfAUU2Zze/WbptzWT3LIzmkYkZHNgz5HVCI540IRRkIU8zmR5X4JFhEl", - "AiZxiDXlJIrxAgCJDDaboQ9PiAH9KWVbc/1vzEIKxah0z9mU94cMo63BAMUEM4koJORKPCH6K9sGgkGU", - "BvQzwmhn8Mx+VVk3AM105F/T+0UIMucBHkcLRDQXw4mo1mEBs/0FdR318k2okGa9MveirfpTWlgqXX3K", - "sItSFkD5xFTof3GBUmaPV92igBxzmKe9hCNUZGXHbEL8mARY05Pxcj8ARcaDIBW+w1EvdaEw3r+jklmY", - "3hmQyiedNB0Q7KkQ1pxxNYM9zWErrf/cwFU5U30fh4x3k3CBMCrwde5QgBrSbIrWALrrPC+5xVzVxvP1", - "n93e0dvXCgK3/Q141mM5n4CJ+GRS2oCrjyazgZclLtRZ+Hvdp4eu1mJRxIUUTxmXigZOGFarAf8wHlsb", - "j8sp6+XmCRcXRd2qzL8vubhoa32duRL3j8oIK87wG7wH0MMD8NWHvw4AZ7QxVDTT3LuBVuWvbJeC0kWV", - "dHHGHEWcTfUuyp3i9+61r1h0QZSCWu5MOecE0UbIyP5oyjXqydhieODhD2yrDy2LdO/3cBf0mitE4yQi", - "MYFyjj3DbHqxM63aVFumEs2yOnrXk5V6VxWTco0tKM31f9epQ8BXbsHWQHuvL5dXqEZ8uhqIK+vcoU55", - "kLiGzFSDJq509DnKZLBWaA3sNbqc0WAGqFxgt+r2DWgXTpLzDJB0fR+9go1cxGWFztcM2LXmNckjYsC2", - "5nF8vl8vWPjh5AQ+MoBcpjTh+T5yRQqz80Pqt4ooW3oWEZYKvbbYYWuZMQ4req6wtjez+a1b/K0cMHbI", - "fFhcjFzaBukEnRdguc4bcLmcvP2dTx9MGes2w3ybuSiOrOkIvElY2GmKsaCRH5FrczDwoc+2RAczw7hj", - "cLDaYH7n0wxivMTKOEnasq8dJnDxPI6X8DBayyUIkirkqfqbVCERAj623N3E3GgNB7a8DL7QjMqMVHIb", - "ex3YzxtJZDB/vaTSQrXT7RCWxp39P+y/5nHc6XbseApYwddQ7legrFUbrEe86JUpQKn9UMuvA5JWFvYF", - "lLTKyWHN6WaN/K154bu/WXQ+uwdkQ9APKk7cb0kFLYy37PBhHEmGEznj6nHhMllXU0Vra3bVuFn29PDC", - "1NXAaBPCcWY/PXNffgPW76rIDjdm5KZ77yEe9RE85kxYWZvNhIsqmM+q2I9vnpFub0lqU23DIT948/p+", - "vlaMmfiq8bulCU1NJJwqHmNFA6jHEcw4lwW2H5MZnlNur0rdnVXGmeDcMHamDaE/16x6bh3B51aR37dO", - "K4SLj2wfffjcBt77v3CP8i9eFuzyTOJ3nfINmNVQMFhQMkEJTiXRelUaE2RK8dsCLAQHMxTgRKWCQG0p", - "gmLKaJzGBVeDNpzEHEeISnS+GZ930ThVKMJiCnaReWjC6QUJeBwTFhLwkA3ZjOA51UadQBFWhAWLniRQ", - "k3JO8kr/2si3UTimppUgmgMpZ10UE4VDrDCoGud6x49MFs95VqbSGNaMXOXcEA6ZSNnPBmdbN3vuBnqO", - "iFR4HFE5y8qZBTgkLPCCWJ9922Ls9r3BZ0RVJ/pAcTk3kqUPGahT9Hq64XwbMTyPLBiZC7uMbcT8EqVX", - "NhuR5fQHx0b/nlvazNXN8YGueDISL9vF38bdTsZ038z9zsNf4HCBwtR0V9iVwObf661MJlCK4U6QWmmW", - "8aZXM1ndpozM15J5G5/dn8c38KZ9I5Kw22jYN1UIySf9LYhcS9UbydwHciNaX1LBK/aAItjFVD2Y+sRF", - "Qco9FnenFdhma2ZyuyidlMBgfXH2Q2xXxbYNObip2Ha+2dqlekGQU9aDKE2/BLdu3EZRbV0H/6a5IJXZ", - "FUTmg4vI/O7g3sTicSYIjWhM8CLiOPwewnSX3OAEXAiD/wCIEo8Jf7TgNSwG6INvrptJiK7LrfxwcrLe", - "JCWEWiojhHrEEqJcoz+IfUX050QIGrri4IcnRzZglkokUtZHb2IKFbsvCEnynBIA8ujr+TkkjHqZ4xLk", - "RbdDmBKLhFOmVo4if/VuBvPlRsWR71lOWqjoHxfSrS+kwbP/+MQZSBnImjATWG6ZKqwaQwFdaBxlpva5", - "1svwmKe6dS2DNJn0ek7hFJzQiMiFVCQ2cYGTNILtBmUHbFVK+51Z5S5ExeqdYxLWEiJiKiXlTA6ZzdZI", - "iNB96891+4UQJ++FgMKZfD01QvLbCJ/TgzERY1g1UQ0wi6AmfGe/s4GTZCPECjeEaNnhfcWQXkI8HJKL", - "eMwjGqCIsguJ1iJ6YcwTNJco0n+sLw2oG8F3t11z8+Y7S1P6mE24tyyZ4dmMmb+rvCor1tzF5KMTa69I", - "cbM4+QML7RdrcqVcEwRHPUVjkiHXoFTRiH4yok43QqWigUn6ySELPpzkqAVDdkKU0O9gSC6LIhIo57DZ", - "SAQPNobpYLAdJBTgz7YJDA4EXvPjGHo8PH1vEkFJzMWiO2T6H9Dwu4NTc7s7wdabUBgoI+qSiwt0vPFm", - "RYjxGZDp3zhGz0xwKXaAd8F/XAleHxGkcQ/Jhi3Kk2WmEk+++yBSq8H98Cs8Tr8CQDJls1mbChyAUixn", - "qQr5JfP7EOY8SmP9D/PH8SpgL4WD2Qd49ZvRds1wVnbjJvgoNqWdU0hM2cQHufQwBHusMauacG4KoMSU", - "ogG9p8CB+h65+/bd90U6foPXnZairiTpN7O37vvks2NwGBdFejyWbW44zc1E8eXep0tMm71PzyMeXEgL", - "hlJ0G2q7DQDG9Y85ILS9IgQ1AXIzkQURQuQqoQKQ3yoOSIO5IxFGioiYMhxtwJxNIwBt7bxYeM4ppEgH", - "EYUkNRoCalEE6HSXM8KQng04qlwDhRtdaUtLFd8pXkYqjsYk4DFxcN/rPtPt75iql1yUsbu/Fbn4rkB/", - "PR89VT3PFXDlzT1+FXz5Cb6CUOkwtRfKbkRrr3j+o3EFdRGszbCzPZDDThcNO1vxsKNX4BCDCxUrtIti", - "ylJFZB8dGf8WJMHuDZAkAWehdKjjzoO3PZBNKbGGLRvyK/fgu/tUeyxXASnf2k584kG/h/T3kLSD1oob", - "zu7JsAubLkQ8Vcbdb/eVfSskCtwj6/d+V1vYIz9s+zaS/O92+5ZkFKyyFpeFpTeSPcN+Xul1c4kaMwMO", - "Z50GAU5wQNWii3AU8SD3HqQyux3oZUMZC4IvtA3VH7K3Geq0Ta5Ah6fvu85phkIqL0wL1i/WR2/mRMh0", - "nA0OgTQwHjxYDBIOmeIowFGQRppvyWRCAsiLADBp2eBXy4Zyl4Wg8068yNeFCKP00RXc8PMErF7OFrLC", - "cRtmqTcECSJM42YoRqv6wuUvuH3HulGuj+FJZK+3AsGlRLapHonolI4je1kj++idVjlwTIYsiTBjRKBU", - "mgglPfReIoiUqUm20Q0AZJnhqC7KYVYSwZV1E0ecC2k8u5rDP5wgqUiyhM3empZPYM53VH/ANG57eiCD", - "oTKG5mPJvoL0ghhOMQTXfKSP6QcICzIDeug6BY9l478TdDolQu8KbISsuRo129qR02z6UvZIY/Gds+yt", - "dsV3slYLEeKF6OmlMBmjHHkw7FzvBtbT+QVtRFKxj66X0fGb/qhl3+XMAf8g7KOvnOX3UtP0rBCw3bZk", - "T87hj616TmHkpa1aSnpYDXHQOsvhLrMOWmMZPBiEwWNGLsClVIYmiIJvjxEG95txd9/lKR43b5WQB0oV", - "+xrSr1Zji34THHg3oKIPnHF6A1DRbyoHClAfHy4X1btRHyqnqeQHdJW/vntc0LtKZTLgoACN0ZTKZKSe", - "DSRYaih9sO+0M5Nsi9+TBm/vnq+hvzuy/7D6W5gMBWL5XXYm39phwZA4UQt3ucgnlQtAST9B2oYPTCKL", - "Ibg7DIcbXK/fHns4Pm28XP8+C3U+yP29LaRCJTo+8lTAfGR4L8U9VzpYNvSp08MimNE5aXa6l3ewJVEi", - "SC/hCVyuhIZglh7uLFNY9KefkG3e4l/Zf0ElHgAuJSEKqSCBihamKpKWCKaPnyQSXFsC8JyLhc+ZXty5", - "LwWPD+xsVpyHdk9ZZ1h+5xsveiFWuDd30maJC+0rbtrd3bYWeIgy9Oo5WiNXShi8XzTRlg+ik4ykpvSp", - "BJ5cLw54c9Dg2aSfyGg6bjPKJcjNbywyNgpSqXjs1v74CK1BJYgpYXottKo/AU02EXxOQ1PhPCfqnEeG", - "qpsNBL2u31UrFVkZD2dcmME9iA7T5kCafqJJWSyY0IXOfmdMGYbBrcRILu8pk1Cl+8OU2SJjbo3cKH4c", - "YdbyW3PGjuZEqEpkiag4N3B76z+Oucd8zBUDU92ZVjrt2pVYbher2jKE9C5AeLM45vt1W3/4dsIrqXyU", - "kZXWdT7PDNImt/m3xYKD+zsf7ttd/uERh+O/Is74LrjKoQHdoo9hfucBjlBI5iTiCVRfNu92up1URJ39", - "zkypZH9jI9LvzbhU+zvPnmx3vnz88v8HAAD//61HyZEIuwEA", + "YErVLB1nzz76CdqEiXVWEgRozRxz6yBJMAX4ksKJb45IzELE7SkQLRBGuQTqD9mJSa6x+WIS7W5uoRP6", + "/Gdk4KvhmJxQorlAC2gLZmOSe07oc8DuMVyRhcCMKcNisSwLzrOIdi34xIThFcj6eQhhd8PO/tAF3g07", + "3WGHsDn8ZqPzhp0vflra+8/mujRwMmbXpBk77g0G66thL+1KeW4/W1wx3LI12WhfZUW1NC/8mZL0u7tp", + "+Le2ZLNLJMx05zlyjnEhODf+I7zKKOj+RZvWc5lRUeAxC0jkFPjVLqP7v4bQixWQKLpvBn0o9swu2mz1", + "xUd2swaLlW+jpRcBD8xxg/s6VEoXAA/Dv4/OE+/xw1svPJm7oGk/KDkgplijHJmXEZboDMbUO9Nm/Av4", + "tW//66xIgPw7j/j0fN84AVDEpyiizAazF0KetXpgaQkfGdCU7DuLoeIqwqwZTeJf//O/MCjKpv/6n//V", + "For5C7b7hkH/ggqS5zOChRoTrM730W+EJD0c0Tlxk4Eab2ROxAJtD+ztATwq1gW3WpocsiF7S1QqWCHo", + "3xRnkbZBe+mg50NZSqQFndEv0olFjjdRkh4PkNvLhpT3uqO7HrQ+mEFhAvpUdDwAUGfUlNG0Nm3H73w1", + "cy65X6sBn7Wwv9XyRZErZbi3ZwZ4TQEDJPbtO3hgJ43Wzs5erPcR2G2GK6A6ANgOeTPWjOj/kEmrZZKR", + "KGWBAlQ2sskgMS2/Pjiy77S7P7Atfk8XCLZgzDVuEIwbCQAY3Qr8uE1ocZvgp5u7WfC5948c1NjdhR2a", + "Lh4o6tDxXp3m5kmBZA/hDEBrDtIBXLNcoNPDY4TDUBAp1/+9XQV6poZL86MDcQYFBB7i/tuOBaqLxyQz", + "1coM8ljEwVs7aoTdvKp1uIrn20aprETjSZdVmMiPvLs/PSqdXucYyWuF5bz24yRZGfFHZcD1twVu6QU4", + "AUI69SXbp0UuWuWQMrGE2ZGzVF2y4vn4yG3I+3NN2a5TVj0b7kEoHlUE4gMKwnK+Z7G63mPi5vfZKjpY", + "1SWeq2+LNQf3pwXdtxfLx+aPyY0VVsimpaBBJmg8QF8RZfAIOne40LYHz8TPiHC72hVDhVln0zKfIgOs", + "ABOCS/7ltu+xeaWd6Wva+54sXyDPdTQWS/IfKkoLYzen1TID1yzBXdq30MO1zNvbu/G2DOYhMgTwjJ3H", + "WigSojUsFyxY/3HpfescbYKrciNWoKyEM0oirCCABCBdMjtLj23rHvS6tzYiCwmsbOTKo0zrO02jyF3N", + "zIlQ6M3hsREBxcNq4zPEpK02QpxYWHpuvX/7e4+wgEMQYhZA59f27JNbNkUMZ5VS9e6fnx9huhp1B2+T", + "KvYV629iRZEJb+1T/h9bLyM6Flgs/mPrJY4Sysh/bB9EWBGp1u+MWQb3dYbct2nwiJlPWwa0TDQQTWwK", + "IIQrVOnsrZbatHv/u1KozaSvpVJndP2hVbfRqovkWqpY26W4U9Xa9PFAd0cZs/moDY9+IFbcgzvScmQB", + "saJ0P5NjVsy4VPDo8aUv2khPmnFc8dho6VfPN+TS48Ox7vFRFwgJVUgBI91mB92Tl92N496VW9vv/bvY", + "D+IxnaY8lcXEoxirYEakTcqLSFkAPza1Oz+eGxXvb5hLB/d5dNy7Xv2D7+9I468uqBHe5qpslc7v3mqr", + "89v3tc5vAAtt4qIFcu+6Ih/rDdGPDrKwLRuXkB3rUZm+cflsEfReGyq5uYDAgtgfsv+j7Y8/FMHxx19c", + "XlM6GGztwe+EzT/+4lKb2IljFcKgvjjkYh28PoL7ySlgPULZpjwfszoOU+cVWM8BVf/bGUj5FW17C8lx", + "4Q8LqZWFVCDXcgvJrsXdmkhlsPt7t5Ecv/kIbiGDv08r6Ru+eLh3C06mkwkNKGFQMgCyRWUt0s5Ycj9u", + "Rm6YJcjsTV8hTKekibQ2IzOptUJDz6uU3nuI1nFeluW+rUdXEPVxZjvwxFYYtPZari00G2zfGj8M7vf0", + "un9D7TGzmLGI6qRLtNLtKfxhSt7EqYLw0hwqCOJ3kTBmTdZiHx1med0yTRIulDRlc8BCMIU1Z9pC8JXY", + "KVfN8ZXJgdIwlMjukEHhVP3YIF1sXJCFKYpDOcvq32QztbVlfFl05aJED7qNbl8J9VdcaqWE3vM2tjX0", + "Hk4JfTDRcS/q3nGpNOlatjHA4h6TbCfzLE2TfqJsuv6oYomNsMrmVgA286haGzhV3FXU35hxg3Hkh3k7", + "jXAAKG/6NQNAZPN+DeJYsSlI5hU8iogwwFJJqlwFriHLBkdZocKwLXdxrpsfpUzR6Lxromkgp18izBYW", + "E2XISp1hpUicaMFm8YJghIIkZsSV0mN60JSnEt7qIslLXSIcXeKFHDJBJhEJ7NygTKMggcFgi6I++pVD", + "IjXCU0yZze3Vb5rCXT/JITunYURGNg/6HFGJ5IwLRRgJUcznRJb7JVhElAiYxCHWlJMoxguANjIob4Y+", + "PCEGPqiUbc31vzELKZS10j1nU94fMoy2BgMUE8wkopCQK/GE6K9sGwgGURrQzwijncEz+1Vl3QB+05F/", + "Te8XIcicB3gcLRDRXAwnolqHBcz2F1SI1Ms3oUKa9crci7Z+UGlhqXSVLsMuSlkAhRhTof/FBUqZPV51", + "iwJyzGGe9hKOUJEVMLMJ8WMSYE1Pxsv9AKgZD4JU+A5HvdSFEnv/jkpmYXpnQCqfdNJ0QLCnQlhzxtUM", + "9jSHrbT+cwNX5Uz1fRwy3k3CBcKowNe5QwGqUbMpWgMQsPO8eBdz9R/P1392e0dvXysI3PY34FmP5XwC", + "JuKTSWkDrj6azAZelrhQZ+HvdZ8euqqNRREXUjxlXCoaOGFYrSv8w3hsbTwup6yXmydcXBR1qzL/vuTi", + "oq31deaK5T8qI6w4w2/wHkAPD2BcH/46AJzRxlDRTHPvBlqVv7JdCkoXVdLFGXMUcTbVuyh3it+7175i", + "0QVRCmq5M+WcE0QbISP7oyn8qCdjy+qBhz+wrT60LNK938Nd0GuuEI2TiMQECkP2DLPpxc60alO3mUo0", + "yyryXU9W6l1VTMo1tqA01/9dpw4BX7kFWwPtvb5cXqEa8elqIK6sc4c65UHiGjJTV5q4ItTnKJPBWqE1", + "ANrockaDGaBygd2q2zegXThJzjNA0vV99Ao2chHhFTpfM7DZmtckj4gB25rH8fl+vfThh5MT+MgAcpki", + "h+f7yJU7zM4Pqd8qomzpWURYKvTaYoetZcY4rOi5wtrezOa3bvG3cujZIfNhcTFyaRukE3RegOU6b8Dl", + "cvL2dz59MGWs2wwYbuaiOLKmI/AmYWGnKcaCRn5Ers3BwIc+2xIdzAzjjsHBaoP5nU8zsPISK+Mkacu+", + "dpjAxfM4XsLDaC2XIEiqkKfqb1KFRAj42HJ3E3OjNRzYQjX4QjMqM1LJbex1YD9vJJHB/PWSSgvVTrdD", + "WBp39v+w/5rHcafbseMpYAVfQ7lfgbJWbbAe8aJXpgCl9kMtvw5IWlnYF1DSKieHNaebNfK35oXv/mbR", + "+ewekA1BP6g4cb8lFbQw3rLDh3EkGU7kjKvHhctkXU0Vra3ZVeNm2dPDC1NXTaNNCMeZ/fTMffkNWL+r", + "IjvcmJGb7r2HeNRH8JgzYWVtNhMuqmA+q2I/vnlGur0lqU21DYf84M3r+/laMWbiq+vvliY01ZVwqniM", + "FQ2gskcw41wW2H5MZnhOub0qdXdWGWeCc8PYmTaE/lyz6rl1BJ9bRX7fOq0QLj6yffThcxt47//CPcq/", + "eFmwyzOJ33XKN2BWQ+lhQckEJTiVROtVaUyQKepvS7kQHMxQgBOVCgJVqgiKKTPFTDJXgzacxBxHiEp0", + "vhmfd9E4VSjCYgp2kXlowukFCXgcExYS8JAN2YzgOdVGnUARVoQFi54kUN1yTtAlFxcRxyEY+TYKx1TH", + "EkRzIOWsi2KicIgVBlXjXO/4kcniOc8KXhrDmpGrnBvCIRMp+9ngbOtmz91AzxGRCo8jKmdZYbQAh4QF", + "XhDrs29bjN2+N/iMqOpEHygu50ay9CEDdYpeTzecbyOG55EFI3Nhl7GNmF+i9MpmI7Kc/uDY6N9zS5u5", + "ujk+0BVPRuJlu/jbuNvJmO6bud95+AscLlCYmu4KuxLY/Hu9lckESjHcCVIrzTLe9Gomq9uUkflaMm/j", + "s/vz+AbetG9EEnYbDfumCiH5pL8FkWupeiOZ+0BuROtLKnjFHlAEu5iqB1OfuChIucfi7rQC22zNTG4X", + "pZMSGKwvzn6I7arYtiEHNxXbzjdbu1QvCHLKehCl6Zfg1o3bKKqt6+DfNBekMruCyHxwEZnfHdybWDzO", + "BKERjQleRByH30OY7pIbnIALYfAfAFHiMeGPFryGxQB98M11MwnRdbmVH05O1pukhFBLZYRQj1hClKv9", + "B7GvHP+cCEFDV2b88OTIBsxSiUTK+uhNTKH29wUhSZ5TAkAefT0/h4RRL3NcgrzodghTYpFwytTKUeSv", + "3s1gvtyoOPI9y0kLFf3jQrr1hTR49h+fOAMpA1kTZgLLLVOFVWMooAuNo8zUPtd6GR7zVLeuZZAmk17P", + "KZyCExoRuZCKxCYucJJGsN2g7ICtSmm/M6vchahYvXNMwlpCREylpJzJIbPZGgkRum/9uW6/EOLkvRBQ", + "OJOvp0ZIfhvhc3owJmIMqyaqAWYR1ITv7Hc2cJJshFjhhhAtO7yvGNJLiIdDchGPeUQDFFF2IdFaRC+M", + "eYLmEkX6j/WlAXUj+O62a27efGdpSh+zCfeWJTM8mzHzd5VXZcWau5h8dGLtFSluFid/YKH9Yk2ulGuC", + "4KinaEwy5BqUKhrRT0bU6UaoVDQwST85ZMGHkxy1YMhOiBL6HQzJZVFEAuUcNhuJ4MHGMB0MtoOEAvzZ", + "NoHBgcBrfhxDj4en700iKIm5WHSHTP8DGn53cGpudyfYehMKA2VEXXJxgY433qwIMT4DMv0bx+iZCS7F", + "DvAu+I8rwesjgjTuIdmwRXmyzFTiyXcfRGo1uB9+hcfpVwBIpmw2a1OBA1CK5SxVIb9kfh/CnEdprP9h", + "/jheBeylcDD7AK9+M9quGc7KbtwEH8WmtHMKiSmb+CCXHoZgjzVmVRPOTQGUmFI0oPcUOFDfI3ffvvu+", + "SMdv8LrTUtSVJP1m9tZ9n3x2DA7jokiPx7LNDae5mSi+3Pt0iWmz9+l5xIMLacFQim5DbbcBwLj+MQeE", + "tleEoCZAbiayIEKIXCVUAPJbxQFpMHckwkgREVOGow2Ys2kEoK2dFwvPOYUU6SCikKRGQ0AtigCd7nJG", + "GNKzAUeVa6BwoyttaaniO8XLSMXRmAQ8Jg7ue91nuv0dU/WSizJ297ciF98V6K/no6eq57kCrry5x6+C", + "Lz/BVxAqHab2QtmNaO0Vz380rqAugrUZdrYHctjpomFnKx529AocYnChYoV2UUxZqojsoyPj34Ik2L0B", + "kiTgLJQOddx58LYHsikl1rBlQ37lHnx3n2qP5Sog5VvbiU886PeQ/h6SdtBaccPZPRl2YdOFiKfKuPvt", + "vrJvhUSBe2T93u9qC3vkh23fRpL/3W7fkoyCVdbisrD0RrJn2M8rvW4uUWNmwOGs0yDACQ6oWnQRjiIe", + "5N6DVGa3A71sKGNB8IW2ofpD9jZDnbbJFejw9H3XOc1QSOWFacH6xfrozZwImY6zwSGQBsaDB4tBwiFT", + "HAU4CtJI8y2ZTEgAeREAJi0b/GrZUO6yEHTeiRf5uhBhlD66ght+noDVy9lCVjhuwyz1hiBBhGncDMVo", + "VV+4/AW371g3yvUxPIns9VYguJTINtUjEZ3ScWQva2QfvdMqB47JkCURZowIlEoToaSH3ksEkTI1yTa6", + "AYAsMxzVRTnMSiK4sm7iiHMhjWdXc/iHEyQVSZaw2VvT8gnM+Y7qD5jGbU8PZDBUxtB8LNlXkF4QwymG", + "4JqP9DH9AGFBZkAPXafgsWz8d4JOp0ToXYGNkDVXo2ZbO3KaTV/KHmksvnOWvdWu+E7WaiFCvBA9vRQm", + "Y5QjD4ad693Aejq/oI1IKvbR9TI6ftMftey7nDngH4R99JWz/F5qmp4VArbbluzJOfyxVc8pjLy0VUtJ", + "D6shDlpnOdxl1kFrLIMHgzB4zMgFuJTK0ARR8O0xwuB+M+7uuzzF4+atEvJAqWJfQ/rVamzRb4ID7wZU", + "9IEzTm8AKvpN5UAB6uPD5aJ6N+pD5TSV/ICu8td3jwt6V6lMBhwUoDGaUpmM1LOBBEsNpQ/2nXZmkm3x", + "e9Lg7d3zNfR3R/YfVn8Lk6FALL/LzuRbOywYEidq4S4X+aRyASjpJ0jb8IFJZDEEd4fhcIPr9dtjD8en", + "jZfr32ehzge5v7eFVKhEx0eeCpiPDO+luOdKB8uGPnV6WAQzOifNTvfyDrYkSgTpJTyBy5XQEMzSw51l", + "Cov+9BOyzVv8K/svqMQDwKUkRCEVJFDRwlRF0hLB9PGTRIJrSwCec7HwOdOLO/el4PGBnc2K89DuKesM", + "y+9840UvxAr35k7aLHGhfcVNu7vb1gIPUYZePUdr5EoJg/eLJtryQXSSkdSUPpXAk+vFAW8OGjyb9BMZ", + "TcdtRrkEufmNRcZGQSoVj93aHx+hNagEMSVMr4VW9SegySaCz2loKpznRJ3zyFB1s4Gg1/W7aqUiK+Ph", + "jAszuAfRYdocSNNPNCmLBRO60NnvjCnDMLiVGMnlPWUSqnR/mDJbZMytkRvFjyPMWn5rztjRnAhViSwR", + "FecGbm/9xzH3mI+5YmCqO9NKp127EsvtYlVbhpDeBQhvFsd8v27rD99OeCWVjzKy0rrO55lB2uQ2/7ZY", + "cHB/58N9u8s/POJw/FfEGd8FVzk0oFv0MczvPMARCsmcRDyB6svm3U63k4qos9+ZKZXsb2xE+r0Zl2p/", + "59mT7c6Xj1/+/wAAAP//l6zTI1K7AQA=", } // GetSwagger returns the content of the embedded swagger specification file diff --git a/openapi.yaml b/openapi.yaml index a8a0855d7..b6263bc4e 100644 --- a/openapi.yaml +++ b/openapi.yaml @@ -4241,7 +4241,9 @@ paths: source: type: string format: binary - description: Source tarball (tar.gz) containing application code and optionally a Dockerfile + description: | + Source tarball (tar.gz) containing application code and optionally a Dockerfile. + Maximum size is 512 MiB; other form fields are limited to 1 MiB each. dockerfile: type: string description: Dockerfile content. Required if not included in the source tarball. From 18fafa464e6310567b018d99320eb3ed3a821482 Mon Sep 17 00:00:00 2001 From: hiroTamada <88675973+hiroTamada@users.noreply.github.com> Date: Mon, 10 Aug 2026 18:56:06 +0000 Subject: [PATCH 2/2] builds: leave shutdown-interrupted builds for startup recovery A build whose run is cancelled by service shutdown kept no useful terminal state: it was marked failed with an opaque "context canceled" error and, unlike a crash, was never recovered on the next start. Keep such builds in their non-terminal status so RecoverPendingBuilds re-runs them, matching crash semantics now that leftover source/config volume recovery makes re-runs safe. Also surface multipart part read errors in the 400 message and skip the trailing retry sleep after the final builder-agent dial attempt. --- cmd/api/api/builds.go | 2 +- cmd/api/main.go | 3 ++- lib/builds/manager.go | 22 ++++++++++++++++++++-- lib/builds/manager_test.go | 11 ++++++++--- 4 files changed, 31 insertions(+), 7 deletions(-) diff --git a/cmd/api/api/builds.go b/cmd/api/api/builds.go index 2c170f0cf..ba94292ba 100644 --- a/cmd/api/api/builds.go +++ b/cmd/api/api/builds.go @@ -32,7 +32,7 @@ var ( func readLimitedPart(part *multipart.Part, limit int64) ([]byte, error) { data, err := io.ReadAll(io.LimitReader(part, limit+1)) if err != nil { - return nil, fmt.Errorf("failed to read %s field", part.FormName()) + return nil, fmt.Errorf("failed to read %s field: %w", part.FormName(), err) } if int64(len(data)) > limit { return nil, fmt.Errorf("%s exceeds the maximum size of %d bytes", part.FormName(), limit) diff --git a/cmd/api/main.go b/cmd/api/main.go index 526d2c62d..f3c16fb4c 100644 --- a/cmd/api/main.go +++ b/cmd/api/main.go @@ -696,7 +696,8 @@ func run() error { } // Cancel in-flight builds and wait for their goroutines to return. - // Builds still pending stay queued on disk and recover on next start. + // Interrupted builds keep their status and pending builds stay + // queued on disk; both are recovered on next start. if err := app.BuildManager.Shutdown(shutdownCtx); err != nil { logger.Error("failed to shutdown build manager", "error", err) // Don't return error - continue with shutdown diff --git a/lib/builds/manager.go b/lib/builds/manager.go index 1e5229361..68c70f4c2 100644 --- a/lib/builds/manager.go +++ b/lib/builds/manager.go @@ -55,7 +55,8 @@ type Manager interface { Start(ctx context.Context) error // Shutdown cancels in-flight builds and waits for their goroutines to - // return. Pending builds are left queued for recovery on the next start. + // return. Interrupted builds keep their status and pending builds stay + // queued on disk; both are recovered on the next start. Shutdown(ctx context.Context) error // CreateBuild starts a new build job @@ -243,7 +244,8 @@ func (m *manager) ReadyForBuilds() bool { } // Shutdown cancels in-flight builds and waits for their goroutines to return. -// Pending builds stay queued on disk and are recovered on the next start. +// Interrupted builds keep their status and pending builds stay queued on +// disk; both are recovered on the next start. func (m *manager) Shutdown(ctx context.Context) error { return m.queue.Shutdown(ctx) } @@ -738,6 +740,14 @@ func (m *manager) runBuild(ctx context.Context, id string, req CreateBuildReques durationMS := duration.Milliseconds() if err != nil { + if ctx.Err() != nil { + // The queue shut down mid-build (the run context is only ever + // cancelled by Shutdown). Leave the build in its non-terminal + // status so RecoverPendingBuilds re-runs it on the next start + // instead of failing it with an opaque cancellation error. + m.logger.Info("build interrupted by shutdown, left for recovery on next start", "id", id, "duration", duration) + return + } m.logger.Error("build failed", "id", id, "error", err, "duration", duration) errMsg := err.Error() m.updateBuildComplete(id, StatusFailed, nil, &errMsg, nil, &durationMS) @@ -776,6 +786,11 @@ func (m *manager) runBuild(ctx context.Context, id string, req CreateBuildReques // Recalculate duration to include image wait time duration = time.Since(start) durationMS = duration.Milliseconds() + if ctx.Err() != nil { + // Interrupted by shutdown; leave for recovery as above. + m.logger.Info("build interrupted by shutdown during image wait, left for recovery on next start", "id", id, "duration", duration) + return + } m.logger.Error("image conversion failed after build", "id", id, "error", err, "duration", duration) errMsg := fmt.Sprintf("image conversion failed: %v", err) m.updateBuildComplete(id, StatusFailed, nil, &errMsg, &result.Provenance, &durationMS) @@ -1110,6 +1125,9 @@ func (m *manager) waitForResult(ctx context.Context, buildID string, inst *insta }, nil } + if attempt == buildAgentDialMaxAttempts-1 { + break // no retry sleep after the final attempt + } select { case <-ctx.Done(): return nil, ctx.Err() diff --git a/lib/builds/manager_test.go b/lib/builds/manager_test.go index d00e9c4bd..3578eb537 100644 --- a/lib/builds/manager_test.go +++ b/lib/builds/manager_test.go @@ -693,7 +693,9 @@ func TestCancelBuild_QueuedBuild(t *testing.T) { // TestManagerShutdown_InFlightBuild verifies the build run detaches from // request cancellation while keeping request context values, and that -// manager shutdown cancels the run and waits for its goroutine. +// manager shutdown cancels the run and waits for its goroutine. The +// interrupted build is left in a non-terminal status so startup recovery +// re-runs it rather than reporting an opaque cancellation failure. func TestManagerShutdown_InFlightBuild(t *testing.T) { mgr, instanceMgr, _, tempDir := setupTestManager(t) defer os.RemoveAll(tempDir) @@ -749,9 +751,12 @@ func TestManagerShutdown_InFlightBuild(t *testing.T) { require.NoError(t, mgr.Shutdown(shutdownCtx)) assert.Less(t, time.Since(start), 5*time.Second, "shutdown should cancel the build promptly") - finished, err := mgr.GetBuild(context.Background(), build.ID) + // The interrupted build keeps its building status (no opaque "context + // canceled" failure) so RecoverPendingBuilds picks it up on next start. + interrupted, err := mgr.GetBuild(context.Background(), build.ID) require.NoError(t, err) - assert.Equal(t, StatusFailed, finished.Status) + assert.Equal(t, StatusBuilding, interrupted.Status) + assert.Nil(t, interrupted.Error) } // TestCancelBuild_BuildingBuild verifies cancelling a running build deletes