diff --git a/api/openapi.yaml b/api/openapi.yaml index 32be8a0..594134e 100644 --- a/api/openapi.yaml +++ b/api/openapi.yaml @@ -209,10 +209,16 @@ components: UpdateScrollRequest: type: object + required: [artifact] properties: artifact: type: string - description: Optional target artifact. If omitted, the daemon refreshes the runtime's current artifact. + pattern: '^[^@\s]+@sha256:[a-f0-9]{64}$' + description: Explicitly accepted immutable target. Missing or mutable tag references are rejected before stopping the runtime. + expected_installed_digest: + type: string + pattern: '^sha256:[a-f0-9]{64}$' + description: Optional precondition checked against the actual installed descriptor under the maintenance lock. registry_credentials: type: array items: @@ -482,6 +488,33 @@ paths: '404': description: Runtime scroll not found + /api/v1/scrolls/{id}/release: + get: + operationId: getInstalledRelease + tags: [scrolls] + summary: Read the installed release descriptor from the runtime volume + parameters: + - name: id + in: path + required: true + schema: + type: string + responses: + '200': + description: Installed canonical repository and last successful release digest + content: + application/json: + schema: + type: object + required: [repository, digest] + properties: + repository: + type: string + digest: + type: string + '404': + description: Runtime scroll not found + /api/v1/scrolls/{id}/start: post: operationId: startScroll @@ -532,7 +565,7 @@ paths: schema: type: string requestBody: - required: false + required: true content: application/json: schema: @@ -544,6 +577,10 @@ paths: application/json: schema: $ref: '#/components/schemas/RuntimeScroll' + '400': + description: An explicitly accepted SHA256 artifact reference is required + '409': + description: Installed release changed since the update was checked '404': description: Runtime scroll not found diff --git a/apps/druid/adapters/cli/client/update.go b/apps/druid/adapters/cli/client/update.go index 45e8f96..3fb36cf 100644 --- a/apps/druid/adapters/cli/client/update.go +++ b/apps/druid/adapters/cli/client/update.go @@ -3,14 +3,11 @@ package client import "github.com/spf13/cobra" var UpdateCommand = &cobra.Command{ - Use: "update [artifact]", - Short: "Update a daemon-managed scroll runtime", - Args: cobra.RangeArgs(1, 2), + Use: "update ", + Short: "Apply an explicitly accepted immutable Scroll revision", + Args: cobra.ExactArgs(2), RunE: func(cmd *cobra.Command, args []string) error { - artifact := "" - if len(args) == 2 { - artifact = args[1] - } + artifact := args[1] daemon, err := runtimeDaemonClient() if err != nil { return err diff --git a/apps/druid/adapters/cli/worker_pull.go b/apps/druid/adapters/cli/worker_pull.go index 4040629..20b5695 100644 --- a/apps/druid/adapters/cli/worker_pull.go +++ b/apps/druid/adapters/cli/worker_pull.go @@ -16,6 +16,7 @@ import ( "github.com/highcard-dev/daemon/internal/core/ports" coreservices "github.com/highcard-dev/daemon/internal/core/services" "github.com/highcard-dev/daemon/internal/core/services/registry" + v1 "github.com/opencontainers/image-spec/specs-go/v1" "github.com/spf13/cobra" "github.com/spf13/viper" ) @@ -52,7 +53,7 @@ func init() { WorkerPullCommand.Flags().StringVar(&workerPullAction.CallbackURL, "callback-url", "", "Daemon worker callback URL") WorkerPullCommand.Flags().StringVar(&workerPullAction.TokenFile, "callback-token-file", "", "Projected ServiceAccount token file for callbacks") WorkerPullCommand.Flags().BoolVar(&workerPullAction.PreserveReleaseManifest, "preserve-release-manifest", false, "Restore the release manifest.json carried by a backup") - WorkerPullCommand.Flags().StringVar(&workerPullMode, "mode", string(ports.RuntimeWorkerModeCreate), "Pull mode: create, update, or restore") + WorkerPullCommand.Flags().StringVar(&workerPullMode, "mode", string(ports.RuntimeWorkerModeCreate), "Worker mode: create, update, restore, or inspect") WorkerPullCommand.MarkFlagRequired("artifact") WorkerPullCommand.MarkFlagRequired("runtime-id") } @@ -67,6 +68,9 @@ func runWorkerPull(action ports.RuntimeWorkerAction) ports.RuntimeWorkerResult { if root == "" { root = "/scroll" } + if action.Mode == ports.RuntimeWorkerModeInspect { + return inspectInstalledRelease(root) + } oci := registry.NewOciClient(loadWorkerRegistryStore()) digest, err := oci.ResolveDigest(action.Artifact) if err == nil { @@ -97,6 +101,39 @@ func runWorkerPull(action ports.RuntimeWorkerAction) ports.RuntimeWorkerResult { return result } +// Inspection reads the installed descriptor from the volume, never from a +// mutable registry tag or a separately persisted release baseline. +func inspectInstalledRelease(root string) ports.RuntimeWorkerResult { + result := ports.RuntimeWorkerResult{} + file, err := os.Open(filepath.Join(root, "manifest.json")) + if err != nil { + result.Error = err.Error() + return result + } + defer file.Close() + var descriptor v1.Descriptor + if err := json.NewDecoder(io.LimitReader(file, 4*1024*1024)).Decode(&descriptor); err != nil { + result.Error = fmt.Sprintf("invalid installed release descriptor: %v", err) + return result + } + if descriptor.Digest.Validate() != nil || descriptor.Digest.Algorithm() != "sha256" { + result.Error = "installed release descriptor requires a valid sha256 digest" + return result + } + yaml, err := os.ReadFile(filepath.Join(root, "scroll.yaml")) + if err != nil { + result.Error = err.Error() + return result + } + if _, err := domain.NewScrollFromBytes(root, yaml); err != nil { + result.Error = err.Error() + return result + } + result.ArtifactDigest = descriptor.Digest.String() + result.ScrollYAML = string(yaml) + return result +} + func loadWorkerRegistryStore() *registry.CredentialStore { var config struct { Registries []domain.RegistryCredential `json:"registries"` @@ -161,7 +198,86 @@ func pullWorkerUpdate(root string, artifact string, oci ports.OciRegistryInterfa } skipData := map[string]bool{} collectSkipUpdatePaths(skipData, "", scroll.Chunks) - return mergePulledRoot(tmp, root, skipData) + installedYAML, err := os.ReadFile(filepath.Join(root, "scroll.yaml")) + if err != nil { + return fmt.Errorf("read installed update protection: %w", err) + } + installed, err := domain.NewScrollFromBytes(root, installedYAML) + if err != nil { + return fmt.Errorf("parse installed update protection: %w", err) + } + // Honor both sides of this transition, including protected chunks removed + // by the candidate. Do not rewrite the candidate's immutable Scroll metadata. + collectSkipUpdatePaths(skipData, "", installed.Chunks) + if err := os.MkdirAll(root, 0755); err != nil { + return err + } + stage, err := os.MkdirTemp(root, ".druid-worker-update-stage-*") + if err != nil { + return err + } + defer os.RemoveAll(stage) + if err := copyPath(tmp, stage); err != nil { + return err + } + if err := preserveSkippedUpdateData(root, stage, skipData); err != nil { + return err + } + return replaceRestoredRoot(root, stage) +} + +// preserveSkippedUpdateData copies only the paths a Scroll explicitly marks as +// skip_update from the installed root into the staged candidate. Candidate +// content is otherwise complete, so unprotected files absent from the +// candidate are removed when the staged root replaces the installed root. +func preserveSkippedUpdateData(root string, stage string, skipData map[string]bool) error { + for skip := range skipData { + skip = filepath.ToSlash(filepath.Clean(skip)) + if !filepath.IsLocal(skip) { + return fmt.Errorf("skip_update path must stay inside runtime data: %q", skip) + } + if skip == "." { + skip = "" + } + source := filepath.Join(root, domain.RuntimeDataDir, filepath.FromSlash(skip)) + target := filepath.Join(stage, domain.RuntimeDataDir, filepath.FromSlash(skip)) + for _, base := range []string{root, stage} { + if err := rejectSymlinkParents(base, filepath.Join(domain.RuntimeDataDir, filepath.FromSlash(skip))); err != nil { + return err + } + } + if err := os.RemoveAll(target); err != nil { + return err + } + if _, err := os.Lstat(source); os.IsNotExist(err) { + continue + } else if err != nil { + return err + } + if err := copyPath(source, target); err != nil { + return err + } + } + return nil +} + +func rejectSymlinkParents(root, relative string) error { + parts := strings.Split(filepath.Clean(relative), string(filepath.Separator)) + current := root + for _, part := range parts[:len(parts)-1] { + current = filepath.Join(current, part) + info, err := os.Lstat(current) + if os.IsNotExist(err) { + return nil + } + if err != nil { + return err + } + if info.Mode()&os.ModeSymlink != 0 { + return fmt.Errorf("protected update path traverses symlink: %s", relative) + } + } + return nil } func pullWorkerRestore(root string, artifact string, oci ports.OciRegistryInterface) error { @@ -200,7 +316,12 @@ func replaceRestoredRootWithRename(root string, stage string, rename func(string if err != nil { return err } - defer os.RemoveAll(rollback) + keepRecovery := false + defer func() { + if !keepRecovery { + _ = os.RemoveAll(rollback) + } + }() stageName := filepath.Base(stage) rollbackName := filepath.Base(rollback) @@ -242,7 +363,8 @@ func replaceRestoredRootWithRename(root string, stage string, rename func(string for _, entry := range original { if err := rename(filepath.Join(root, entry), filepath.Join(rollback, entry)); err != nil { if rollbackErr := restoreOriginal(nil); rollbackErr != nil { - return fmt.Errorf("%s failed to restore original runtime after moving %s: %w", restoreRootUnsafePrefix, entry, rollbackErr) + keepRecovery = true + return fmt.Errorf("%s recovery files retained at %s; failed to restore original runtime after moving %s: %w", restoreRootUnsafePrefix, rollback, entry, rollbackErr) } return fmt.Errorf("restore transaction rolled back while moving original entry %s: %w", entry, err) } @@ -253,7 +375,8 @@ func replaceRestoredRootWithRename(root string, stage string, rename func(string name := entry.Name() if err := rename(filepath.Join(stage, name), filepath.Join(root, name)); err != nil { if rollbackErr := restoreOriginal(installed); rollbackErr != nil { - return fmt.Errorf("%s failed to restore original runtime after staging %s: %w", restoreRootUnsafePrefix, name, rollbackErr) + keepRecovery = true + return fmt.Errorf("%s recovery files retained at %s; failed to restore original runtime after staging %s: %w", restoreRootUnsafePrefix, rollback, name, rollbackErr) } return fmt.Errorf("restore transaction rolled back while staging %s: %w", name, err) } @@ -278,79 +401,20 @@ func collectSkipUpdatePaths(out map[string]bool, parent string, chunks []*domain } } -func mergePulledRoot(src string, dst string, skipData map[string]bool) error { - if err := os.MkdirAll(dst, 0755); err != nil { - return err - } - entries, err := os.ReadDir(src) +func copyPath(src string, dst string) error { + info, err := os.Lstat(src) if err != nil { return err } - for _, entry := range entries { - name := entry.Name() - srcPath := filepath.Join(src, name) - dstPath := filepath.Join(dst, name) - if name == domain.RuntimeDataDir { - if err := copyDataUpdate(srcPath, dstPath, skipData); err != nil { - return err - } - continue - } - if err := os.RemoveAll(dstPath); err != nil { - return err - } - if err := copyPath(srcPath, dstPath); err != nil { - return err - } - } - return nil -} - -func copyDataUpdate(srcData string, dstData string, skipData map[string]bool) error { - return filepath.WalkDir(srcData, func(srcPath string, entry os.DirEntry, err error) error { + if info.Mode()&os.ModeSymlink != 0 { + target, err := os.Readlink(src) if err != nil { return err } - rel, err := filepath.Rel(srcData, srcPath) - if err != nil { + if err := os.MkdirAll(filepath.Dir(dst), 0755); err != nil { return err } - if rel == "." { - return os.MkdirAll(dstData, 0755) - } - rel = filepath.ToSlash(rel) - if shouldSkipWorkerUpdate(rel, skipData) { - if entry.IsDir() { - return filepath.SkipDir - } - return nil - } - target := filepath.Join(dstData, filepath.FromSlash(rel)) - if entry.IsDir() { - info, err := entry.Info() - if err != nil { - return err - } - return os.MkdirAll(target, info.Mode().Perm()) - } - return copyPath(srcPath, target) - }) -} - -func shouldSkipWorkerUpdate(rel string, skipData map[string]bool) bool { - rel = filepath.ToSlash(filepath.Clean(rel)) - for skip := range skipData { - if skip == "" || rel == skip || strings.HasPrefix(rel, skip+"/") { - return true - } - } - return false -} - -func copyPath(src string, dst string) error { - info, err := os.Stat(src) - if err != nil { - return err + return os.Symlink(target, dst) } if info.IsDir() { return filepath.WalkDir(src, func(path string, entry os.DirEntry, walkErr error) error { diff --git a/apps/druid/adapters/cli/worker_release_test.go b/apps/druid/adapters/cli/worker_release_test.go new file mode 100644 index 0000000..c3d291e --- /dev/null +++ b/apps/druid/adapters/cli/worker_release_test.go @@ -0,0 +1,69 @@ +package cli + +import ( + "os" + "path/filepath" + "strings" + "testing" + + "github.com/highcard-dev/daemon/internal/core/ports" +) + +func TestInspectWorkerReadsInstalledReleaseNotRegistryReference(t *testing.T) { + root := t.TempDir() + want := "sha256:" + strings.Repeat("a", 64) + raw := "{\n \"digest\": \"" + want + "\"\n}" + mustWrite(t, filepath.Join(root, "manifest.json"), raw) + mustWrite(t, filepath.Join(root, "scroll.yaml"), "name: registry.local/owned/example\napp_version: '1'\ncommands: {}\n") + result := runWorkerPull(ports.RuntimeWorkerAction{Mode: ports.RuntimeWorkerModeInspect, MountPath: root, Artifact: "unreachable.invalid/backup:latest"}) + if result.Error != "" || result.ArtifactDigest != want { + t.Fatalf("inspection = %#v", result) + } + assertFile(t, filepath.Join(root, "manifest.json"), raw) +} + +func TestInspectWorkerRejectsMissingOrInvalidInstalledDescriptor(t *testing.T) { + for _, raw := range []string{"", "{}", `{"digest":"sha256:bad"}`} { + t.Run(raw, func(t *testing.T) { + root := t.TempDir() + if raw != "" { + mustWrite(t, filepath.Join(root, "manifest.json"), raw) + } + if result := inspectInstalledRelease(root); result.Error == "" { + t.Fatal("invalid installed release was accepted") + } + }) + } +} + +func TestPreserveSkippedUpdateDataRejectsTraversal(t *testing.T) { + root := t.TempDir() + stage, err := os.MkdirTemp(root, ".stage-") + if err != nil { + t.Fatal(err) + } + if err := preserveSkippedUpdateData(root, stage, map[string]bool{"../../outside": true}); err == nil { + t.Fatal("unsafe skip_update path accepted") + } +} + +func TestPreserveSkippedUpdateDataDoesNotFollowSymlinks(t *testing.T) { + root, stage, outside := t.TempDir(), t.TempDir(), t.TempDir() + mustWrite(t, filepath.Join(outside, "save"), "outside runtime") + if err := os.MkdirAll(filepath.Join(root, "data"), 0755); err != nil { + t.Fatal(err) + } + if err := os.Symlink(outside, filepath.Join(root, "data", "link")); err != nil { + t.Fatal(err) + } + if err := preserveSkippedUpdateData(root, stage, map[string]bool{"link/save": true}); err == nil { + t.Fatal("followed a protected path through an external symlink") + } + if err := preserveSkippedUpdateData(root, stage, map[string]bool{"link": true}); err != nil { + t.Fatal(err) + } + if target, err := os.Readlink(filepath.Join(stage, "data", "link")); err != nil || target != outside { + t.Fatalf("protected symlink was dereferenced: %q (%v)", target, err) + } + assertFile(t, filepath.Join(outside, "save"), "outside runtime") +} diff --git a/apps/druid/adapters/cli/worker_test.go b/apps/druid/adapters/cli/worker_test.go index ac18897..cec8ff3 100644 --- a/apps/druid/adapters/cli/worker_test.go +++ b/apps/druid/adapters/cli/worker_test.go @@ -85,24 +85,73 @@ func TestReportWorkerResultUsesTokenOnlyWhenProvided(t *testing.T) { } } -func TestWorkerUpdateMergePreservesSkipUpdateAndExtraFiles(t *testing.T) { - src := t.TempDir() - dst := t.TempDir() - mustWrite(t, filepath.Join(src, "scroll.yaml"), "name: next\n") - mustWrite(t, filepath.Join(src, "data", "keep", "state.txt"), "new") - mustWrite(t, filepath.Join(src, "data", "overwrite.txt"), "new") - mustWrite(t, filepath.Join(dst, "scroll.yaml"), "name: old\n") - mustWrite(t, filepath.Join(dst, "data", "keep", "state.txt"), "old") - mustWrite(t, filepath.Join(dst, "data", "overwrite.txt"), "old") - mustWrite(t, filepath.Join(dst, "data", "extra.txt"), "extra") - - if err := mergePulledRoot(src, dst, map[string]bool{"keep": true}); err != nil { +func TestStagedUpdateReplacesUnprotectedAndPreservesSkipUpdate(t *testing.T) { + root := t.TempDir() + mustWrite(t, filepath.Join(root, "scroll.yaml"), "name: old\n") + mustWrite(t, filepath.Join(root, "data", "keep", "state.txt"), "old") + mustWrite(t, filepath.Join(root, "data", "extra.txt"), "old") + stage, err := os.MkdirTemp(root, ".druid-worker-update-stage-*") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(stage) + mustWrite(t, filepath.Join(stage, "scroll.yaml"), "name: new\n") + mustWrite(t, filepath.Join(stage, "data", "keep", "state.txt"), "new") + mustWrite(t, filepath.Join(stage, "data", "replace.txt"), "new") + + if err := preserveSkippedUpdateData(root, stage, map[string]bool{"keep": true}); err != nil { + t.Fatal(err) + } + if err := replaceRestoredRoot(root, stage); err != nil { + t.Fatal(err) + } + + assertFile(t, filepath.Join(root, "scroll.yaml"), "name: new\n") + assertFile(t, filepath.Join(root, "data", "keep", "state.txt"), "old") + assertFile(t, filepath.Join(root, "data", "replace.txt"), "new") + if _, err := os.Stat(filepath.Join(root, "data", "extra.txt")); !os.IsNotExist(err) { + t.Fatalf("unprotected destination-only file should be removed, stat err = %v", err) + } +} + +func TestPreserveSkippedUpdateDataKeepsMissingPathsMissing(t *testing.T) { + root := t.TempDir() + stage, err := os.MkdirTemp(root, ".druid-worker-update-stage-*") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(stage) + mustWrite(t, filepath.Join(stage, "data", "keep", "state.txt"), "candidate") + + if err := preserveSkippedUpdateData(root, stage, map[string]bool{"keep": true}); err != nil { t.Fatal(err) } - assertFile(t, filepath.Join(dst, "scroll.yaml"), "name: next\n") - assertFile(t, filepath.Join(dst, "data", "keep", "state.txt"), "old") - assertFile(t, filepath.Join(dst, "data", "overwrite.txt"), "new") - assertFile(t, filepath.Join(dst, "data", "extra.txt"), "extra") + if _, err := os.Stat(filepath.Join(stage, "data", "keep")); !os.IsNotExist(err) { + t.Fatalf("missing installed skip_update path should remain missing, stat err = %v", err) + } +} + +func TestWorkerUpdatePreservesRemovedProtectedChunk(t *testing.T) { + for _, chunks := range []string{ + " - name: world\n path: world\n skip_update: true\n", + " - name: server\n path: .\n chunks:\n - name: world\n path: world\n skip_update: true\n", + } { + t.Run(chunks, func(t *testing.T) { + root := t.TempDir() + candidate := t.TempDir() + mustWrite(t, filepath.Join(root, "scroll.yaml"), "name: example\nchunks:\n"+chunks) + mustWrite(t, filepath.Join(root, "data", "world", "save.dat"), "user world") + mustWrite(t, filepath.Join(root, "data", "obsolete.txt"), "old release") + mustWrite(t, filepath.Join(candidate, "scroll.yaml"), "name: example\n") + if err := pullWorkerUpdate(root, candidate, nil); err != nil { + t.Fatal(err) + } + assertFile(t, filepath.Join(root, "data", "world", "save.dat"), "user world") + if _, err := os.Stat(filepath.Join(root, "data", "obsolete.txt")); !os.IsNotExist(err) { + t.Fatalf("unprotected obsolete file retained: %v", err) + } + }) + } } func TestWorkerRestoreStagesBeforeReplacingRoot(t *testing.T) { @@ -182,6 +231,11 @@ func TestReplaceRestoredRootMarksFailedRollbackUnsafe(t *testing.T) { if err == nil || !strings.HasPrefix(err.Error(), restoreRootUnsafePrefix) { t.Fatalf("restore error = %v, want unsafe rollback marker", err) } + recovery, err := filepath.Glob(filepath.Join(root, ".druid-worker-restore-rollback-*", "original.txt")) + if err != nil || len(recovery) != 1 { + t.Fatalf("original data must survive failed rollback: files=%v err=%v", recovery, err) + } + assertFile(t, recovery[0], "original") } func TestWorkerCollectSkipUpdatePaths(t *testing.T) { diff --git a/apps/druid/adapters/daemonclient/openapi_client.go b/apps/druid/adapters/daemonclient/openapi_client.go index 0f4179c..88d1de2 100644 --- a/apps/druid/adapters/daemonclient/openapi_client.go +++ b/apps/druid/adapters/daemonclient/openapi_client.go @@ -80,10 +80,7 @@ func (c *OpenAPIClient) CreateScroll(ctx context.Context, name string, artifact } func (c *OpenAPIClient) UpdateScroll(ctx context.Context, id string, artifact string, registryCredentials []api.RegistryCredential) (*api.RuntimeScroll, error) { - request := api.UpdateScrollJSONRequestBody{} - if artifact != "" { - request.Artifact = &artifact - } + request := api.UpdateScrollJSONRequestBody{Artifact: artifact} if len(registryCredentials) > 0 { request.RegistryCredentials = ®istryCredentials } diff --git a/apps/druid/adapters/http/handlers/routes.go b/apps/druid/adapters/http/handlers/routes.go index 4e4a072..adb9fc9 100644 --- a/apps/druid/adapters/http/handlers/routes.go +++ b/apps/druid/adapters/http/handlers/routes.go @@ -63,6 +63,8 @@ func RegisterPublicRoutes(app *fiber.App, handlers RouteHandlers) { app.Get("/:id/api/v1/health", handlers.Server.GetHealthAuth) app.Get("/:id/api/v1/token", handlers.Server.CreateDaemonToken) app.Get("/:id/api/v1/scroll", handlers.Server.GetDaemonScroll) + app.Get("/:id/api/v1/release", func(c *fiber.Ctx) error { return handlers.Server.GetInstalledRelease(c, c.Params("id")) }) + app.Post("/:id/api/v1/update", func(c *fiber.Ctx) error { return handlers.Server.UpdateScroll(c, c.Params("id")) }) app.Put("/:id/api/v1/scroll/commands/:command", handlers.Server.AddDaemonCommand) app.Delete("/:id/api/v1/scroll/commands/:command", handlers.Server.RemoveDaemonCommand) app.Post("/:id/api/v1/command", handlers.Server.RunDaemonCommand) diff --git a/apps/druid/adapters/http/handlers/routes_test.go b/apps/druid/adapters/http/handlers/routes_test.go index a66a273..11fca1d 100644 --- a/apps/druid/adapters/http/handlers/routes_test.go +++ b/apps/druid/adapters/http/handlers/routes_test.go @@ -3,9 +3,13 @@ package handlers import ( "net/http" "net/http/httptest" + "strings" "testing" "github.com/gofiber/fiber/v2" + appservices "github.com/highcard-dev/daemon/apps/druid/core/services" + "github.com/highcard-dev/daemon/internal/core/domain" + "github.com/highcard-dev/daemon/internal/core/ports" ) func TestRouteSplitKeepsManagementAndPublicSurfacesSeparate(t *testing.T) { @@ -68,3 +72,63 @@ func requestStatus(t *testing.T, app *fiber.App, path string) int { defer resp.Body.Close() return resp.StatusCode } + +type ownerRouteStore struct{ ports.RuntimeScrollStore } + +func (ownerRouteStore) GetScroll(id string) (*domain.RuntimeScroll, error) { + return &domain.RuntimeScroll{ID: id, OwnerID: "owner"}, nil +} + +type ownerRouteAuth struct { + ports.AuthorizerServiceInterface +} + +func (ownerRouteAuth) CheckHeader(c *fiber.Ctx) (*ports.AuthContext, error) { + if c.Get("Authorization") == "" { + return nil, nil + } + return &ports.AuthContext{Subject: strings.TrimPrefix(c.Get("Authorization"), "Bearer ")}, nil +} + +type ownerRouteBackend struct{ ports.RuntimeBackendInterface } + +func TestPublicReleaseRoutesEnforceOwnerBeforeReadingOrUpdating(t *testing.T) { + supervisor, err := appservices.NewRuntimeSupervisor(ownerRouteStore{}, nil, ports.RuntimeBackendFactoryFunc(func(ports.ProcedureStatusObserver) (ports.RuntimeBackendInterface, error) { + return ownerRouteBackend{}, nil + })) + if err != nil { + t.Fatal(err) + } + app := fiber.New() + RegisterPublicRoutes(app, RouteHandlers{Server: NewRuntimeServer(NewHealthHandler(), NewScrollHandler(supervisor, ownerRouteAuth{})), Websocket: &WebsocketHandler{}}) + for _, route := range []struct{ method, path string }{{http.MethodGet, "/owned/api/v1/release"}, {http.MethodPost, "/owned/api/v1/update"}} { + for _, auth := range []struct { + header string + status int + }{{"", http.StatusUnauthorized}, {"Bearer other-owner", http.StatusForbidden}} { + req := httptest.NewRequest(route.method, route.path, strings.NewReader(`{"artifact":"registry/scroll:mutable"}`)) + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", auth.header) + res, err := app.Test(req) + if err != nil { + t.Fatal(err) + } + res.Body.Close() + if res.StatusCode != auth.status { + t.Fatalf("%s %s: got %d, want %d", route.method, route.path, res.StatusCode, auth.status) + } + } + } + // An owner reaches the explicit-update validator, not a 404 or another API. + req := httptest.NewRequest(http.MethodPost, "/owned/api/v1/update", strings.NewReader(`{"artifact":"registry/scroll:mutable"}`)) + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer owner") + res, err := app.Test(req) + if err != nil { + t.Fatal(err) + } + defer res.Body.Close() + if res.StatusCode != http.StatusBadRequest { + t.Fatalf("owner mutable update status = %d, want 400", res.StatusCode) + } +} diff --git a/apps/druid/adapters/http/handlers/scroll_handler.go b/apps/druid/adapters/http/handlers/scroll_handler.go index 68fc4cc..3b5e906 100644 --- a/apps/druid/adapters/http/handlers/scroll_handler.go +++ b/apps/druid/adapters/http/handlers/scroll_handler.go @@ -124,6 +124,17 @@ func (h *ScrollHandler) GetScroll(c *fiber.Ctx, id string) error { return c.JSON(runtimeScroll) } +func (h *ScrollHandler) GetInstalledRelease(c *fiber.Ctx, id string) error { + if _, err := h.getScroll(id); err != nil { + return err + } + release, err := h.supervisor.InstalledRelease(c.UserContext(), id) + if err != nil { + return err + } + return c.JSON(release) +} + func (h *ScrollHandler) DeleteScroll(c *fiber.Ctx, id string) error { runtimeScroll, err := h.getScroll(id) if err != nil { @@ -170,12 +181,14 @@ func (h *ScrollHandler) UpdateScroll(c *fiber.Ctx, id string) error { return fiber.NewError(fiber.StatusBadRequest, err.Error()) } } - artifact := "" - if request.Artifact != nil { - artifact = *request.Artifact - } - runtimeScroll, err := h.supervisor.Update(id, artifact, registryCredentials(request.RegistryCredentials)) + runtimeScroll, err := h.supervisor.Update(id, request.Artifact, registryCredentials(request.RegistryCredentials), request.ExpectedInstalledDigest) if err != nil { + if errors.Is(err, appservices.ErrInstalledReleaseChanged) { + return fiber.NewError(fiber.StatusConflict, err.Error()) + } + if errors.Is(err, appservices.ErrUnacceptedUpdate) { + return fiber.NewError(fiber.StatusBadRequest, err.Error()) + } return err } return c.JSON(runtimeScroll) diff --git a/apps/druid/core/services/runtime_access.go b/apps/druid/core/services/runtime_access.go index f3c6786..5503735 100644 --- a/apps/druid/core/services/runtime_access.go +++ b/apps/druid/core/services/runtime_access.go @@ -7,6 +7,7 @@ import ( "github.com/highcard-dev/daemon/internal/core/domain" "github.com/highcard-dev/daemon/internal/core/ports" + coreservices "github.com/highcard-dev/daemon/internal/core/services" ) func (s *RuntimeSupervisor) Run(id string, command string) (*domain.RuntimeScroll, error) { @@ -14,20 +15,36 @@ func (s *RuntimeSupervisor) Run(id string, command string) (*domain.RuntimeScrol } func (s *RuntimeSupervisor) RunWithContext(ctx context.Context, id string, command string) (*domain.RuntimeScroll, error) { + unlock := s.lockRuntimeOperation(id) session, err := s.sessionFor(id) + if err != nil { + unlock() + return nil, err + } + session.Start() + longRunning, err := session.beginRun(command) + unlock() if err != nil { return nil, err } - return session.RunWithContext(ctx, command) + return session.finishRun(ctx, longRunning) } // RunAndWait waits for the requested command rather than every command in the runtime queue. func (s *RuntimeSupervisor) RunAndWait(id string, command string) (*domain.RuntimeScroll, error) { + unlock := s.lockRuntimeOperation(id) session, err := s.sessionFor(id) + if err != nil { + unlock() + return nil, err + } + session.Start() + wait, err := session.enqueueQueueItem(command, coreservices.AddItemOptions{Wait: true}, nil) + unlock() if err != nil { return nil, err } - if err := session.AddTempItemWithWait(command); err != nil { + if err := wait(); err != nil { return nil, err } return s.store.GetScroll(id) diff --git a/apps/druid/core/services/runtime_command_maintenance_test.go b/apps/druid/core/services/runtime_command_maintenance_test.go new file mode 100644 index 0000000..de72aa7 --- /dev/null +++ b/apps/druid/core/services/runtime_command_maintenance_test.go @@ -0,0 +1,94 @@ +package services + +import ( + "testing" + "time" + + "github.com/highcard-dev/daemon/internal/core/domain" + "github.com/highcard-dev/daemon/internal/core/ports" + coreservices "github.com/highcard-dev/daemon/internal/core/services" +) + +func TestCommandAdmissionWaitsForMaintenance(t *testing.T) { + for _, sync := range []bool{false, true} { + t.Run(map[bool]string{false: "run", true: "run-and-wait"}[sync], func(t *testing.T) { + store := newTestStateStore(t) + fixture := &domain.RuntimeScroll{ID: "command-maintenance", Root: "runtime://fixture", ScrollYAML: cachedScrollYAML("start"), Status: domain.RuntimeScrollStatusStopped} + if err := store.CreateScroll(fixture); err != nil { + t.Fatal(err) + } + supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), &fakeWorkerBackend{}) + unlock := supervisor.lockRuntimeOperation(fixture.ID) + defer func() { + if unlock != nil { + unlock() + } + }() + done := make(chan error, 1) + go func() { + var err error + if sync { + _, err = supervisor.RunAndWait(fixture.ID, "missing") + } else { + _, err = supervisor.Run(fixture.ID, "missing") + } + done <- err + }() + select { + case err := <-done: + t.Fatalf("command reached session during maintenance: %v", err) + case <-time.After(50 * time.Millisecond): + } + unlock() + unlock = nil + select { + case err := <-done: + if err == nil { + t.Fatal("missing command accepted") + } + case <-time.After(time.Second): + t.Fatal("command admission did not resume") + } + }) + } +} + +func TestCommandWaitDoesNotHoldMaintenanceLock(t *testing.T) { + store := newTestStateStore(t) + fixture := &domain.RuntimeScroll{ID: "command-wait", Root: "runtime://fixture", ScrollYAML: cachedScrollYAML("start"), Status: domain.RuntimeScrollStatusStopped} + if err := store.CreateScroll(fixture); err != nil { + t.Fatal(err) + } + started, release := make(chan struct{}), make(chan struct{}) + backend := &fakeWorkerBackend{runCommand: func(ports.RuntimeCommand) (*int, error) { + close(started) + <-release + code := 0 + return &code, nil + }} + supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), backend) + done := make(chan error, 1) + go func() { _, err := supervisor.RunAndWait(fixture.ID, "start"); done <- err }() + select { + case <-started: + case <-time.After(time.Second): + t.Fatal("command did not start") + } + acquired := make(chan struct{}) + go func() { unlock := supervisor.lockRuntimeOperation(fixture.ID); close(acquired); unlock() }() + select { + case <-acquired: + case <-time.After(time.Second): + close(release) + t.Fatal("command wait blocks maintenance admission") + } + close(release) + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(time.Second): + t.Fatal("command waiter did not complete") + } +} diff --git a/apps/druid/core/services/runtime_release.go b/apps/druid/core/services/runtime_release.go new file mode 100644 index 0000000..e6f63b4 --- /dev/null +++ b/apps/druid/core/services/runtime_release.go @@ -0,0 +1,36 @@ +package services + +import ( + "context" + "fmt" + + "github.com/highcard-dev/daemon/internal/core/domain" + "github.com/highcard-dev/daemon/internal/core/ports" +) + +// InstalledRelease reads the actual runtime volume under the same operation +// lock as update/restore, so it cannot observe a partially replaced root. +func (s *RuntimeSupervisor) InstalledRelease(ctx context.Context, id string) (map[string]string, error) { + unlock := s.lockRuntimeOperation(id) + defer unlock() + runtime, err := s.store.GetScroll(id) + if err != nil { + return nil, err + } + return s.installedRelease(ctx, runtime) +} + +func (s *RuntimeSupervisor) installedRelease(ctx context.Context, runtime *domain.RuntimeScroll) (map[string]string, error) { + installed, err := s.runPullWorker(ctx, s.runtimeBackend, ports.RuntimeWorkerModeInspect, runtime.ID, runtime.Artifact, runtime.Root, nil, "") + if err != nil { + return nil, err + } + scroll, err := domain.NewScrollFromBytes(runtime.Root, installed.ScrollYAML) + if err != nil { + return nil, err + } + if !acceptedUpdateReference.MatchString(scroll.Name + "@" + installed.ArtifactDigest) { + return nil, fmt.Errorf("invalid installed release identity") + } + return map[string]string{"repository": scroll.Name, "digest": installed.ArtifactDigest}, nil +} diff --git a/apps/druid/core/services/runtime_session_commands.go b/apps/druid/core/services/runtime_session_commands.go index 949d993..0b2e8ff 100644 --- a/apps/druid/core/services/runtime_session_commands.go +++ b/apps/druid/core/services/runtime_session_commands.go @@ -121,34 +121,42 @@ func (s *RuntimeSession) Run(command string) (*domain.RuntimeScroll, error) { } func (s *RuntimeSession) RunWithContext(ctx context.Context, command string) (*domain.RuntimeScroll, error) { + longRunning, err := s.beginRun(command) + if err != nil { + return nil, err + } + return s.finishRun(ctx, longRunning) +} + +func (s *RuntimeSession) beginRun(command string) (bool, error) { s.refreshCommandState() targetCommand, err := s.scrollService.GetCommand(command) if err != nil { s.markError(err) - return nil, err + return false, err } longRunning := targetCommand.Run == domain.RunModeRestart || targetCommand.Run == domain.RunModePersistent s.rememberDoneDependencies(targetCommand, map[string]bool{}) if err := s.AddTempItem(command); err != nil { s.markError(err) - return nil, err + return false, err } + return longRunning, nil +} + +func (s *RuntimeSession) finishRun(ctx context.Context, longRunning bool) (*domain.RuntimeScroll, error) { if !longRunning { if err := s.WaitUntilEmptyContext(ctx); err != nil { - s.markError(err) return nil, err } } s.mu.Lock() - s.runtimeScroll.Status = deriveRuntimeScrollStatus(s.runtimeScroll.Procedures, s.scrollService.GetFile().Commands) - err = s.store.UpdateScroll(s.runtimeScroll) id := s.runtimeScroll.ID s.mu.Unlock() - if err != nil { - return nil, err - } + // Queue callbacks own status updates; a waiter must not overwrite state + // after concurrent maintenance has installed a different session/release. return s.store.GetScroll(id) } diff --git a/apps/druid/core/services/runtime_session_queue.go b/apps/druid/core/services/runtime_session_queue.go index 874002b..d110419 100644 --- a/apps/druid/core/services/runtime_session_queue.go +++ b/apps/druid/core/services/runtime_session_queue.go @@ -57,11 +57,21 @@ func (s *RuntimeSession) addQueueItem(cmd string, options coreservices.AddItemOp } func (s *RuntimeSession) addQueueItemWithEnv(cmd string, options coreservices.AddItemOptions, procedureEnv map[string]map[string]string) error { + wait, err := s.enqueueQueueItem(cmd, options, procedureEnv) + if err != nil { + return err + } + return wait() +} + +// enqueueQueueItem separates admission from waiting. Supervisors serialize +// admission with maintenance without holding a lock while a command runs. +func (s *RuntimeSession) enqueueQueueItem(cmd string, options coreservices.AddItemOptions, procedureEnv map[string]map[string]string) (func() error, error) { logger.Log().Debug("Running command", zap.String("cmd", cmd)) command, err := s.scrollService.GetCommand(cmd) if err != nil { - return err + return nil, err } s.queueMu.Lock() @@ -71,12 +81,12 @@ func (s *RuntimeSession) addQueueItemWithEnv(cmd string, options coreservices.Ad if item != nil { if currentStatus != domain.ScrollLockStatusDone && currentStatus != domain.ScrollLockStatusError { s.queueMu.Unlock() - return coreservices.ErrAlreadyInQueue + return nil, coreservices.ErrAlreadyInQueue } } if hasCurrentStatus && currentStatus == domain.ScrollLockStatusDone && command.Run == domain.RunModeOnce && !options.Force { s.queueMu.Unlock() - return coreservices.ErrCommandDoneOnce + return nil, coreservices.ErrCommandDoneOnce } var doneChan chan struct{} @@ -90,21 +100,19 @@ func (s *RuntimeSession) addQueueItemWithEnv(cmd string, options coreservices.Ad s.triggerRunQueue() - if options.Wait { - <-doneChan - s.queueMu.Lock() - item := s.queue[cmd] - var itemErr error - if item != nil { - itemErr = item.err - } - s.queueMu.Unlock() - if itemErr != nil { - return itemErr + return func() error { + if options.Wait { + <-doneChan + s.queueMu.Lock() + itemErr := item.err + s.queueMu.Unlock() + if itemErr != nil { + return itemErr + } } - } - return nil + return nil + }, nil } func (s *RuntimeSession) HydrateFromState(statuses domain.ProcedureStatusMap) error { diff --git a/apps/druid/core/services/runtime_supervisor.go b/apps/druid/core/services/runtime_supervisor.go index 6bc64ec..064e7d5 100644 --- a/apps/druid/core/services/runtime_supervisor.go +++ b/apps/druid/core/services/runtime_supervisor.go @@ -213,15 +213,9 @@ func (s *RuntimeSupervisor) Ensure(options EnsureOptions) (*domain.RuntimeScroll if runtimeScroll.Status == domain.RuntimeScrollStatusError && (options.Artifact == "" || options.Artifact == runtimeScroll.Artifact) { return s.persistEnsureOptions(runtimeScroll, options) } - if options.Artifact != "" { - nextDigest := resolveArtifactDigest(options.Artifact, options.RegistryCredentials) - artifactChanged := options.Artifact != runtimeScroll.Artifact - digestChanged := nextDigest != "" && nextDigest != runtimeScroll.ArtifactDigest - if artifactChanged || digestChanged { - applyEnsureOptions(runtimeScroll, options) - return s.updateExistingScroll(runtimeScroll, options.Artifact, nextDigest, options.RegistryCredentials, false) - } - } + // Reconciliation does not accept releases on the user's behalf. In + // particular it must neither follow a moved tag nor undo a restore. + // Existing workloads change only through the explicit Update operation. return s.persistEnsureOptions(runtimeScroll, options) } if !errors.Is(err, domain.ErrRuntimeScrollNotFound) { diff --git a/apps/druid/core/services/runtime_supervisor_test.go b/apps/druid/core/services/runtime_supervisor_test.go index 17fc378..16b7b06 100644 --- a/apps/druid/core/services/runtime_supervisor_test.go +++ b/apps/druid/core/services/runtime_supervisor_test.go @@ -960,7 +960,7 @@ func TestRuntimeSupervisorEnsureDoesNotRetryExistingError(t *testing.T) { } } -func TestRuntimeSupervisorEnsureUpdatesChangedArtifact(t *testing.T) { +func TestRuntimeSupervisorEnsureDoesNotApplyUnacceptedArtifact(t *testing.T) { store := newTestStateStore(t) root := "k8s://druid/druid-update-scroll-data" existing := &domain.RuntimeScroll{ @@ -999,26 +999,11 @@ func TestRuntimeSupervisorEnsureUpdatesChangedArtifact(t *testing.T) { t.Fatal(err) } - if backend.stopRoot != root { - t.Fatalf("stop root = %s, want %s", backend.stopRoot, root) + if backend.spawnCount != 0 || backend.stopRoot != "" { + t.Fatalf("reconciliation mutated installed runtime: worker=%#v stop=%s", backend.action, backend.stopRoot) } - if backend.action.Mode != ports.RuntimeWorkerModeUpdate || backend.action.Artifact != "registry.local/lab:2.0" || backend.action.RootRef != root { - t.Fatalf("worker action = %#v", backend.action) - } - if updated.Artifact != "registry.local/lab:2.0" || updated.ScrollName != "updated-scroll" { - t.Fatalf("updated scroll = %#v", updated) - } - if updated.Status != domain.RuntimeScrollStatusStopped { - t.Fatalf("status = %s, want stopped", updated.Status) - } - if len(updated.Procedures) != 0 { - t.Fatalf("procedures = %#v, want cleared", updated.Procedures) - } - if len(updated.Routing) != 1 || updated.Routing[0].PortName != "main" { - t.Fatalf("routing = %#v, want matching route preserved", updated.Routing) - } - if !strings.Contains(updated.ScrollYAML, "updated-scroll") { - t.Fatalf("scroll yaml = %q", updated.ScrollYAML) + if updated.Artifact != existing.Artifact || updated.ScrollYAML != existing.ScrollYAML || updated.Status != domain.RuntimeScrollStatusRunning { + t.Fatalf("reconciliation changed installed release: %#v", updated) } } @@ -1046,19 +1031,20 @@ func TestRuntimeSupervisorUpdateUsesPullWorkerWhenAvailable(t *testing.T) { ) supervisor.SetWorkerCallbacks(callbacks, "http://druid-cli:8083") - updated, err := supervisor.Ensure(EnsureOptions{Artifact: "registry.local/lab:2.0", Name: "update-worker"}) + accepted := "registry.local/lab@sha256:" + strings.Repeat("a", 64) + updated, err := supervisor.Update("update-worker", accepted, nil, nil) if err != nil { t.Fatal(err) } - if backend.action.Mode != ports.RuntimeWorkerModeUpdate || backend.action.RootRef != root { + if backend.action.Mode != ports.RuntimeWorkerModeUpdate || backend.action.RootRef != root || backend.action.Artifact != accepted { t.Fatalf("worker action = %#v", backend.action) } - if updated.Artifact != "registry.local/lab:2.0" || updated.ArtifactDigest != "sha256:updated" || updated.ScrollName != "updated-worker" { + if updated.Artifact != accepted || updated.ArtifactDigest != "sha256:updated" || updated.ScrollName != "updated-worker" { t.Fatalf("updated scroll = %#v", updated) } } -func TestRuntimeSupervisorUpdateRefreshesCurrentArtifactAndRestartsRunningScroll(t *testing.T) { +func TestRuntimeSupervisorUpdateAppliesAcceptedDigestAndRestartsRunningScroll(t *testing.T) { store := newTestStateStore(t) root := "runtime://refresh-worker" existing := &domain.RuntimeScroll{ @@ -1078,7 +1064,8 @@ func TestRuntimeSupervisorUpdateRefreshesCurrentArtifactAndRestartsRunningScroll supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), backend) supervisor.SetWorkerCallbacks(callbacks, "http://druid-cli:8083") - updated, err := supervisor.Update("refresh-worker", "", nil) + accepted := "registry.local/lab@sha256:" + strings.Repeat("b", 64) + updated, err := supervisor.Update("refresh-worker", accepted, nil, nil) if err != nil { t.Fatal(err) } @@ -1089,7 +1076,7 @@ func TestRuntimeSupervisorUpdateRefreshesCurrentArtifactAndRestartsRunningScroll if backend.stopRoot != root { t.Fatalf("stop root = %s, want %s", backend.stopRoot, root) } - if backend.action.Mode != ports.RuntimeWorkerModeUpdate || backend.action.Artifact != "registry.local/lab:1.0" { + if backend.action.Mode != ports.RuntimeWorkerModeUpdate || backend.action.Artifact != accepted { t.Fatalf("worker action = %#v", backend.action) } if updated.Status != domain.RuntimeScrollStatusRunning { @@ -1191,14 +1178,14 @@ func TestRuntimeSupervisorBackupFailureRestartsPriorRunningScroll(t *testing.T) Artifact: "registry.local/lab:1.0", Root: "runtime://backup-recovery", ScrollName: "backup-recovery", - ScrollYAML: cachedScrollYAML("start"), + ScrollYAML: strings.Replace(cachedScrollYAML("start"), "run: once", "run: persistent", 1), Status: domain.RuntimeScrollStatusRunning, Procedures: domain.ProcedureStatusMap{}, } if err := store.CreateScroll(runtimeScroll); err != nil { t.Fatal(err) } - backend := &fakeWorkerBackend{backupErr: errors.New("registry unavailable")} + backend := &fakeWorkerBackend{backupErr: errors.New("registry unavailable"), procedureStatusUpdates: []ports.ProcedureStatusUpdate{{Procedure: "start.0", Status: domain.ScrollLockStatusRunning}}} supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), backend) if _, err := supervisor.Backup("backup-recovery", "registry.local/backups:1", nil); err == nil { @@ -1220,7 +1207,7 @@ func TestRuntimeSupervisorRestoreFailureRestartsPriorRunningScroll(t *testing.T) Artifact: "registry.local/lab:1.0", Root: "runtime://restore-recovery", ScrollName: "restore-recovery", - ScrollYAML: cachedScrollYAML("start"), + ScrollYAML: strings.Replace(cachedScrollYAML("start"), "run: once", "run: persistent", 1), Status: domain.RuntimeScrollStatusRunning, Procedures: domain.ProcedureStatusMap{}, } @@ -1228,7 +1215,7 @@ func TestRuntimeSupervisorRestoreFailureRestartsPriorRunningScroll(t *testing.T) t.Fatal(err) } callbacks := NewWorkerCallbackManager() - backend := &fakeWorkerBackend{callbacks: callbacks, workerErr: errors.New("backup pull failed")} + backend := &fakeWorkerBackend{callbacks: callbacks, workerErr: errors.New("backup pull failed"), procedureStatusUpdates: []ports.ProcedureStatusUpdate{{Procedure: "start.0", Status: domain.ScrollLockStatusRunning}}} supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), backend) supervisor.SetWorkerCallbacks(callbacks, "http://druid-cli:8083") diff --git a/apps/druid/core/services/runtime_update.go b/apps/druid/core/services/runtime_update.go index ff2823b..6c237e3 100644 --- a/apps/druid/core/services/runtime_update.go +++ b/apps/druid/core/services/runtime_update.go @@ -3,6 +3,8 @@ package services import ( "context" "errors" + "regexp" + "strings" "github.com/highcard-dev/daemon/internal/core/domain" "github.com/highcard-dev/daemon/internal/core/ports" @@ -11,17 +13,33 @@ import ( "go.uber.org/zap" ) -func (s *RuntimeSupervisor) Update(id string, artifact string, registryCredentials []domain.RegistryCredential) (*domain.RuntimeScroll, error) { +var ErrUnacceptedUpdate = errors.New("update requires an explicitly accepted sha256 artifact reference") +var ErrInstalledReleaseChanged = errors.New("installed release changed; check updates again") +var acceptedUpdateReference = regexp.MustCompile(`^[^@\s]+@sha256:[a-f0-9]{64}$`) + +func (s *RuntimeSupervisor) Update(id string, artifact string, registryCredentials []domain.RegistryCredential, expectedDigest *string) (*domain.RuntimeScroll, error) { + if !acceptedUpdateReference.MatchString(artifact) { + return nil, ErrUnacceptedUpdate + } unlock := s.lockRuntimeOperation(id) defer unlock() runtimeScroll, err := s.store.GetScroll(id) if err != nil { return nil, err } - if artifact == "" { - artifact = runtimeScroll.Artifact + if expectedDigest != nil { + if !acceptedUpdateReference.MatchString("expected@" + *expectedDigest) { + return nil, ErrUnacceptedUpdate + } + installed, err := s.installedRelease(context.Background(), runtimeScroll) + if err != nil { + return nil, err + } + if installed["digest"] != *expectedDigest { + return nil, ErrInstalledReleaseChanged + } } - knownDigest := resolveArtifactDigest(artifact, registryCredentials) + knownDigest := artifact[strings.LastIndex(artifact, "@")+1:] return s.updateExistingScroll(runtimeScroll, artifact, knownDigest, registryCredentials, true) } diff --git a/apps/druid/core/services/runtime_update_test.go b/apps/druid/core/services/runtime_update_test.go new file mode 100644 index 0000000..90cbc87 --- /dev/null +++ b/apps/druid/core/services/runtime_update_test.go @@ -0,0 +1,52 @@ +package services + +import ( + "errors" + "strings" + "testing" + + "github.com/highcard-dev/daemon/internal/core/domain" + "github.com/highcard-dev/daemon/internal/core/ports" + coreservices "github.com/highcard-dev/daemon/internal/core/services" +) + +func TestRuntimeUpdateRejectsUnacceptedReferenceBeforeStopping(t *testing.T) { + for _, artifact := range []string{"", "registry.local/example:latest", "registry.local/example@sha256:bad"} { + t.Run(artifact, func(t *testing.T) { + store := newTestStateStore(t) + original := &domain.RuntimeScroll{ID: "accepted-update", Root: "runtime://accepted-update", Artifact: "registry.local/example:v1", ScrollYAML: cachedScrollYAML("start"), Status: domain.RuntimeScrollStatusRunning} + if err := store.CreateScroll(original); err != nil { + t.Fatal(err) + } + backend := &fakeWorkerBackend{} + supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), backend) + _, err := supervisor.Update(original.ID, artifact, nil, nil) + if err == nil || !strings.Contains(err.Error(), "accepted sha256") { + t.Fatalf("error = %v, want accepted digest validation", err) + } + if backend.stopRoot != "" || backend.spawnCount != 0 { + t.Fatal("invalid selection stopped or modified the workload") + } + }) + } +} + +func TestUpdateRejectsChangedInstalledDescriptorBeforeStopping(t *testing.T) { + store := newTestStateStore(t) + fixture := &domain.RuntimeScroll{ID: "stale-update", Root: "runtime://fixture", Artifact: "registry.local/deployment:v1", ScrollYAML: cachedScrollYAML("start"), Status: domain.RuntimeScrollStatusRunning} + if err := store.CreateScroll(fixture); err != nil { + t.Fatal(err) + } + callbacks := NewWorkerCallbackManager() + backend := &fakeWorkerBackend{callbacks: callbacks, scrollYAML: "name: registry.local/reusable\n", digest: "sha256:" + strings.Repeat("b", 64)} + supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), backend) + supervisor.SetWorkerCallbacks(callbacks, "http://worker-callback") + expected := "sha256:" + strings.Repeat("a", 64) + _, err := supervisor.Update(fixture.ID, "registry.local/deployment@sha256:"+strings.Repeat("c", 64), nil, &expected) + if !errors.Is(err, ErrInstalledReleaseChanged) { + t.Fatalf("stale update error: %v", err) + } + if backend.stopRoot != "" || backend.spawnCount != 1 || backend.action.Mode != ports.RuntimeWorkerModeInspect { + t.Fatal("stale acceptance mutated runtime") + } +} diff --git a/docs/runtime-lifecycle-verification.md b/docs/runtime-lifecycle-verification.md new file mode 100644 index 0000000..cbfdcfb --- /dev/null +++ b/docs/runtime-lifecycle-verification.md @@ -0,0 +1,35 @@ +# Runtime lifecycle acceptance + +Verified locally on 2026-09-28. This is runtime acceptance, not complete product acceptance. + +## Contract implemented + +- Updates require an explicitly accepted `repository@sha256:<64 hex>` reference. Empty references and tags fail before stopping the workload. +- Reconciliation no longer changes an installed release. Explicit updates and restores own those transitions. +- The installed release endpoint reads the actual volume's `manifest.json` and canonical Scroll name through a read-only worker. It does not resolve the deployment tag or add another persisted baseline. +- Callers can send `expected_installed_digest` with an accepted update. The runtime checks the actual installed descriptor under its maintenance lock and rejects a stale approval with HTTP 409 before stopping or changing the workload. +- The owner-authenticated public listener exposes installed-release reads and explicit updates, using the same runtime handlers as management. Tests reject missing and cross-owner authorization before either operation. +- Updates stage the whole candidate, preserve protected paths, and remove obsolete unprotected files. Both installed and candidate `skip_update` declarations protect the current transition, including nested declarations removed by the candidate. +- Protection longevity is unresolved: if a release removes a protected declaration, this transition retains its data, but indefinite retention across later releases is not guaranteed. No hidden persisted protection list or finalized-metadata rewrite was introduced. +- Failed rollback retains recovery files and keeps the workload stopped. Recovery locations are included in the error. +- Backup mode snapshots all runtime data, including files outside explicit release chunk selections, and preserves the installed descriptor bytes. +- Command admission shares the maintenance lock. Waiting for a command does not hold that lock, and waiters do not overwrite state after a release transition. +- Protected-path copying preserves symlinks without dereferencing them and rejects paths traversing symlink parents. + +## Evidence + +- Full `go test ./... -timeout=180s`: passed. +- Regressions were observed failing before fixes for removed protected chunks, rollback-file retention, omitted runtime-created backup files, command admission during maintenance, and symlink traversal. +- `TestKubernetesBackendCLIComplexLifecycle`, with rebuilt daemon/worker binaries and image, passed three times (100.42s, 103.78s, 104.95s). The final run includes command-admission, complete-backup, and symlink-hardening changes. +- Focused command-admission/maintenance tests passed with Go's race detector enabled. +- The installed-digest precondition passed focused runtime/HTTP/client tests and a further rebuilt-image Kubernetes lifecycle run (102.55s), including successful v2 acceptance, rejection of the stale original baseline with HTTP 409, unchanged v2 contents, and backup restoration. +- The Kubernetes test uses a unique namespace, PVC, registry, management ports, and daemon socket. It does not replace the shared daemon/operator or consume canonical ReservedPorts. +- Verified real workload start and command execution; stopped backup; installed-descriptor read; acceptance of v2 followed by moving the tag to v3; installation of v2; preservation of protected runtime data; restoration of v1, runtime-created data outside declared chunks, and byte-identical installed descriptor. +- Test harness failures were diagnosed separately: registry:2 needs an OCI Accept header, and explicit fixture chunks must include the version marker. +- Recovery-restart tests use persistent workload fixtures. The previous finite `true` command legitimately stopped immediately and was not evidence of a failed restart. + +## Remaining product work + +- Server/UI check-and-apply wiring is implemented in the separate monorepo worktree but shared-stack integration, Team publication, and browser acceptance remain outside this runtime milestone. +- The shared browser fixture is blocked by all canonical local game ports being allocated. No existing deployment was retired. +- This changes the update API/CLI contract; callers must send the accepted immutable reference. Deploy the matching callers with this runtime revision. diff --git a/internal/api/generated.go b/internal/api/generated.go index 914574b..7d08a31 100644 --- a/internal/api/generated.go +++ b/internal/api/generated.go @@ -229,9 +229,12 @@ type RuntimeUIPackages map[string]RuntimeUIPackage // UpdateScrollRequest defines model for UpdateScrollRequest. type UpdateScrollRequest struct { - // Artifact Optional target artifact. If omitted, the daemon refreshes the runtime's current artifact. - Artifact *string `json:"artifact,omitempty"` - RegistryCredentials *[]RegistryCredential `json:"registry_credentials,omitempty"` + // Artifact Explicitly accepted immutable target. Missing or mutable tag references are rejected before stopping the runtime. + Artifact string `json:"artifact"` + + // ExpectedInstalledDigest Optional precondition checked against the actual installed descriptor under the maintenance lock. + ExpectedInstalledDigest *string `json:"expected_installed_digest,omitempty"` + RegistryCredentials *[]RegistryCredential `json:"registry_credentials,omitempty"` } // RunScrollCommandParams defines parameters for RunScrollCommand. @@ -379,6 +382,9 @@ type ClientInterface interface { // GetScrollQueue request GetScrollQueue(ctx context.Context, id string, reqEditors ...RequestEditorFn) (*http.Response, error) + // GetInstalledRelease request + GetInstalledRelease(ctx context.Context, id string, reqEditors ...RequestEditorFn) (*http.Response, error) + // RestoreScrollWithBody request with any body RestoreScrollWithBody(ctx context.Context, id string, contentType string, body io.Reader, reqEditors ...RequestEditorFn) (*http.Response, error) @@ -592,6 +598,18 @@ func (c *Client) GetScrollQueue(ctx context.Context, id string, reqEditors ...Re return c.Client.Do(req) } +func (c *Client) GetInstalledRelease(ctx context.Context, id string, reqEditors ...RequestEditorFn) (*http.Response, error) { + req, err := NewGetInstalledReleaseRequest(c.Server, id) + if err != nil { + return nil, err + } + req = req.WithContext(ctx) + if err := c.applyEditors(ctx, req, reqEditors); err != nil { + return nil, err + } + return c.Client.Do(req) +} + func (c *Client) RestoreScrollWithBody(ctx context.Context, id string, contentType string, body io.Reader, reqEditors ...RequestEditorFn) (*http.Response, error) { req, err := NewRestoreScrollRequestWithBody(c.Server, id, contentType, body) if err != nil { @@ -1184,6 +1202,40 @@ func NewGetScrollQueueRequest(server string, id string) (*http.Request, error) { return req, nil } +// NewGetInstalledReleaseRequest generates requests for GetInstalledRelease +func NewGetInstalledReleaseRequest(server string, id string) (*http.Request, error) { + var err error + + var pathParam0 string + + pathParam0, err = runtime.StyleParamWithLocation("simple", false, "id", runtime.ParamLocationPath, id) + if err != nil { + return nil, err + } + + serverURL, err := url.Parse(server) + if err != nil { + return nil, err + } + + operationPath := fmt.Sprintf("/api/v1/scrolls/%s/release", pathParam0) + if operationPath[0] == '/' { + operationPath = "." + operationPath + } + + queryURL, err := serverURL.Parse(operationPath) + if err != nil { + return nil, err + } + + req, err := http.NewRequest("GET", queryURL.String(), nil) + if err != nil { + return nil, err + } + + return req, nil +} + // NewRestoreScrollRequest calls the generic RestoreScroll builder with application/json body func NewRestoreScrollRequest(server string, id string, body RestoreScrollJSONRequestBody) (*http.Request, error) { var bodyReader io.Reader @@ -1600,6 +1652,9 @@ type ClientWithResponsesInterface interface { // GetScrollQueueWithResponse request GetScrollQueueWithResponse(ctx context.Context, id string, reqEditors ...RequestEditorFn) (*GetScrollQueueResponse, error) + // GetInstalledReleaseWithResponse request + GetInstalledReleaseWithResponse(ctx context.Context, id string, reqEditors ...RequestEditorFn) (*GetInstalledReleaseResponse, error) + // RestoreScrollWithBodyWithResponse request with any body RestoreScrollWithBodyWithResponse(ctx context.Context, id string, contentType string, body io.Reader, reqEditors ...RequestEditorFn) (*RestoreScrollResponse, error) @@ -1898,6 +1953,31 @@ func (r GetScrollQueueResponse) StatusCode() int { return 0 } +type GetInstalledReleaseResponse struct { + Body []byte + HTTPResponse *http.Response + JSON200 *struct { + Digest string `json:"digest"` + Repository string `json:"repository"` + } +} + +// Status returns HTTPResponse.Status +func (r GetInstalledReleaseResponse) Status() string { + if r.HTTPResponse != nil { + return r.HTTPResponse.Status + } + return http.StatusText(0) +} + +// StatusCode returns HTTPResponse.StatusCode +func (r GetInstalledReleaseResponse) StatusCode() int { + if r.HTTPResponse != nil { + return r.HTTPResponse.StatusCode + } + return 0 +} + type RestoreScrollResponse struct { Body []byte HTTPResponse *http.Response @@ -2206,6 +2286,15 @@ func (c *ClientWithResponses) GetScrollQueueWithResponse(ctx context.Context, id return ParseGetScrollQueueResponse(rsp) } +// GetInstalledReleaseWithResponse request returning *GetInstalledReleaseResponse +func (c *ClientWithResponses) GetInstalledReleaseWithResponse(ctx context.Context, id string, reqEditors ...RequestEditorFn) (*GetInstalledReleaseResponse, error) { + rsp, err := c.GetInstalledRelease(ctx, id, reqEditors...) + if err != nil { + return nil, err + } + return ParseGetInstalledReleaseResponse(rsp) +} + // RestoreScrollWithBodyWithResponse request with arbitrary body returning *RestoreScrollResponse func (c *ClientWithResponses) RestoreScrollWithBodyWithResponse(ctx context.Context, id string, contentType string, body io.Reader, reqEditors ...RequestEditorFn) (*RestoreScrollResponse, error) { rsp, err := c.RestoreScrollWithBody(ctx, id, contentType, body, reqEditors...) @@ -2629,6 +2718,35 @@ func ParseGetScrollQueueResponse(rsp *http.Response) (*GetScrollQueueResponse, e return response, nil } +// ParseGetInstalledReleaseResponse parses an HTTP response from a GetInstalledReleaseWithResponse call +func ParseGetInstalledReleaseResponse(rsp *http.Response) (*GetInstalledReleaseResponse, error) { + bodyBytes, err := io.ReadAll(rsp.Body) + defer func() { _ = rsp.Body.Close() }() + if err != nil { + return nil, err + } + + response := &GetInstalledReleaseResponse{ + Body: bodyBytes, + HTTPResponse: rsp, + } + + switch { + case strings.Contains(rsp.Header.Get("Content-Type"), "json") && rsp.StatusCode == 200: + var dest struct { + Digest string `json:"digest"` + Repository string `json:"repository"` + } + if err := json.Unmarshal(bodyBytes, &dest); err != nil { + return nil, err + } + response.JSON200 = &dest + + } + + return response, nil +} + // ParseRestoreScrollResponse parses an HTTP response from a RestoreScrollWithResponse call func ParseRestoreScrollResponse(rsp *http.Response) (*RestoreScrollResponse, error) { bodyBytes, err := io.ReadAll(rsp.Body) @@ -2875,6 +2993,9 @@ type ServerInterface interface { // Get runtime queue state // (GET /api/v1/scrolls/{id}/queue) GetScrollQueue(c *fiber.Ctx, id string) error + // Read the installed release descriptor from the runtime volume + // (GET /api/v1/scrolls/{id}/release) + GetInstalledRelease(c *fiber.Ctx, id string) error // Execute runtime restore // (POST /api/v1/scrolls/{id}/restore) RestoreScroll(c *fiber.Ctx, id string) error @@ -3084,6 +3205,22 @@ func (siw *ServerInterfaceWrapper) GetScrollQueue(c *fiber.Ctx) error { return siw.Handler.GetScrollQueue(c, id) } +// GetInstalledRelease operation middleware +func (siw *ServerInterfaceWrapper) GetInstalledRelease(c *fiber.Ctx) error { + + var err error + + // ------------- Path parameter "id" ------------- + var id string + + err = runtime.BindStyledParameterWithOptions("simple", "id", c.Params("id"), &id, runtime.BindStyledParameterOptions{Explode: false, Required: true}) + if err != nil { + return fiber.NewError(fiber.StatusBadRequest, fmt.Errorf("Invalid format for parameter id: %w", err).Error()) + } + + return siw.Handler.GetInstalledRelease(c, id) +} + // RestoreScroll operation middleware func (siw *ServerInterfaceWrapper) RestoreScroll(c *fiber.Ctx) error { @@ -3265,6 +3402,8 @@ func RegisterHandlersWithOptions(router fiber.Router, si ServerInterface, option router.Get(options.BaseURL+"/api/v1/scrolls/:id/queue", wrapper.GetScrollQueue) + router.Get(options.BaseURL+"/api/v1/scrolls/:id/release", wrapper.GetInstalledRelease) + router.Post(options.BaseURL+"/api/v1/scrolls/:id/restore", wrapper.RestoreScroll) router.Post(options.BaseURL+"/api/v1/scrolls/:id/routing", wrapper.ApplyScrollRouting) @@ -3286,54 +3425,59 @@ func RegisterHandlersWithOptions(router fiber.Router, si ServerInterface, option // Base64 encoded, gzipped, json marshaled Swagger object var swaggerSpec = []string{ - "H4sIAAAAAAAC/+xbX3MbNw7/KhzezdzLWnKuTR98T65z7blNJz47mTy0GQ1FQhKrXZIhuZZ1Hn33G/7Z", - "1f7hSpZiN3anL4klkiDwAwiAIHSPqSyUFCCswWf32NAFFMT/ea5Uvr6WpeVifg2fSzDWfa20VKAtBz+J", - "GMPnoqiWcwuF/+PvGmb4DP9tvCU/jrTH16WwvABHGs7r9XiTYbtWgM8w0Zqs8WaTYQ2fS66B4bNfW1t9", - "qufK6e9A/eILDcTCDdUyz4f51ZbPCPUjDAzVXFkuBT7D7y4uUTWKNMxAg6CApEa5pCRHxhNGitgFzjDc", - "kULlgdmwxoyYLjkbzedjC8b6f87cP7jm1VjNxdzxylmfgTegNFBigSGSc2LQTGokSAEj9M7PcUxYMs0B", - "6YAg4mwcJlzOkCy4tcAyZBeAGIFCCjQHAZpYMIgIxNmoxfjvcmpSvDmKCXgeh4V/hTFuVE7WXjpkLM9z", - "RGUBBs20LCLSozUp8odzbBShCbZ/LqegBbj961ke2Ip/DUaWmoIZocu5kBoYmq6RkOKksXRK6BIEM6PU", - "7nIlQE9SGo2GjvwMxBkqDTC/Oy2NlQXokxmhXMyRdmcBkdIupOb/I259ci8Nc26sXk+oBgbCcpIfcO7i", - "4ot67f4zVx2X1IF7AzlYYOHE9Y9aQKQngrHEln7CVrEsUOpL3GGHuymRQIqjfwtT6kNcQI+7anDC+Dyu", - "Hji8g+fmL/N8Hub5HyC5XVyDUVIY6NtBIVlCIxel1iAsWvjVKBgb8nObrkguU/IrLecajOmTvYojSIGm", - "ICyZBz3nkjCHsGPM4+oc3Ezqglh8hme5JC5+FOSOF2WBz16dnma44CJ8Oq1ZEGUxBR2Pl7YTRmxCto8L", - "EE3f7Of6Y1fv6BaeOKvAGRZlnjtfj8+sLmHf2fQQpfTwVtLlTX3o2zqAO24nNCpiYD8uLMyDcDkxdhJU", - "MqELIuZ+Xc08F/a7b3FqYcPpCIfcr1iXQjgxMsyk8LrVWmqc4RXhLuNpiDIgcKSZ5CqFw5XUCW/U0tAh", - "XkVFcrVtfPf69TevG9bxKgWE0tJKKvMmFAtrFc78fz68UvepZGo/BJ65yEqDdlJ6LSkw5509UL8Q5X0x", - "YzzkFVdtHz3w/S7/0bCzTYKBPkflNOdm8eHyitAlmcNgwPAp346EyIebEy2lPdGQE8tvAY1WxBQ+WRyh", - "NzAjZW4NshIpzW+JhTHjxo6JUmGe1Eg5bmj7+1EyIPYESTjOngwLORDMFDFmJXU6pJUG9IABdizB028s", - "aBBOWUMMPefRf7+rvN9xQXso7LAAPD6bkdxA9nhhyG3pnWeDnamUORBxWIyKODjXMOQip7IULLVP5kGf", - "cJXExI9VPqLvB5YA6jznt/Bek9mM0yQN79gItfyW2/WE2JazbUaKw71W0jMFB5Fe1/Bbff3fTaZrG+B6", - "SDDwGVWSku2h0YA7Dh60V7VGLnfTXHHB5CrN0yHSDTjoGtu+s86ihW2FrxHaYbHdy3sislvnCvJd9nl4", - "wJsMj+4ykOBcdxyHUudpH7dLfi7m74meQ0L6h90FDjkd+4Q/9uwYyIFaqXdF3b5JdlExoG85hcnDgsWA", - "VU4G0okO+R1WOXQV3Rk9qK8bsYP828AV0DvMkEimhptXsWEd7o1QiVQqRCTQt8C8lT/81uWzUr+csHci", - "X3eS723Ak3Ig+IaT8OjVvwyHxGrY6vtJfVQlzhrpvbFSKf9dleFX1YZPCcWWfKJCOvhQQer80aedpWIH", - "GlOqxFHba8S9jUW2vXo0bLe1944zUvM7nOj2UTlYqh0utSmtm5ThWFM9kP+jLwo9IFIe7YNn5vhqcnU7", - "sD5C1JXlwTqphpkGswDjv4zln38YRGM9oibwteouHYS806el5nZ94wjFZBWIBn1eBjsKn36ozOWnj+/9", - "6Wvi9NPH98jKJYhQ+uWeA7tGSstbzkB703fkXdrkyW3l9/dWx5lbX+3ZJn+zkNqeuDyXoc8l6HW1mdTo", - "I0xvJF2CRVQKAbSqvnC30E/GVTYSttjuTBT/GRwsLhSImXQbUylsMIVNV8g3uuQMXby9RDkpBV34YjhD", - "BRHOjJFfyQXoE1/IY9VTA1Eq5zRUhTKU8yX8Jua+Yu4cvTYZYsSSKTFgMk9wBdNqbPSbZ5dbX62qGcAZ", - "dqOBrdPRq9Gpj0sKBFEcn+Fv/FfhRHqFjoni49tX41AOc9/EfKct4Y9gK0NuFc6wJx7udpcsTAx1Oa8v", - "H7V8ec5v9s/T0wrJmFM2IBj/bkKJJNjtPqvuVP+8qto8hxk+0rw+/eYP3PgmZDOoFOSW8FDy8uepLAqi", - "1xHOLo6WzI2/aQdNZDjgjT+5pZWaguWYhp7a8L/lxt7EOV8I/iHBPuZlfbfSw6aqSVeCtHFx7NelcVPL", - "UUETR5rYuGzSJIBoPhbiEJPA2O8lWz+aIaTeIzftAOhSrU1PD68ejYUO/PvgRlX+1EY9CNLBfTfsfZMc", - "g3+b8TE0qZHm280TaST1PPQgjZx+NY0E1LoaCYJ0NILgjhsbQouM6Ue+DkV+c7C67jnbBD/vkuW+usLj", - "X60uRTQpwIJ2W9yHGBrTuhhCfWLbBjprgNbNEz89oRLaD5f7lVBdGDYZ/vb02+GHtDhdSItmvqbS1lrY", - "9qBzlKXd+I9gXybyB5r/lyLu4ugXui13DsYuLyvVsO/63o8/sUoe3x/uq8I/N98YYEZuaTyQba94B7Rs", - "HLCotaM0TmVREMHM+D7+tRnW/nUpAs8XYepTWECWJELrDQ+i1HkVJtz6C5G/eAbVA0ORNrKyBhxNYSa1", - "b0pQUjAu5qOB+5JZC4qbXHSfYnqPJl/V64TLPuv6ii/0Ptel6IborcKOskkx4/PB3L4OChdh3jMMDd0a", - "Qk8RV0Sb7QU4Ctz36So17VhMjczBPAjVMPMZ4pqugrUKucOYV4KhJaxDe9G2MN+HPr52GyqVdxI1KEeA", - "X1eqdyN/JUMC++xgP+QK3HjiPfgajBxQVSXg0VOiFvWj9Pi5hBL26/G/ftoLy1hTjy19fXnRPIawA+/P", - "jVlboD9HWCrAh3HWYKzcdYe+DhP+SkSf+pYScH5wJlop7qjT1XhgS2vdN+rHqkac+3JUn/qVwXNT91B6", - "2NL5FWjDjY2dnFKfhN8rAIutXUjXuukbgX/y3msC4/CY9ICQ2WpPePGxs91s8ZDwGRagCq9EGhN+xRA7", - "pyvd1AuO0FHdD5Y+pDdu+E9atLkJfcS7z4ef9CjVGGOl2gW0VH9anH0rwz6cpepmeCupl7kkzKDVgueA", - "VGgWcRbPiCXHqaHk42arxG6H1Hi1f5laabZ5JO6toZcYGPpwiSpU3FXK35NSN9jEAvTh+q35Yl2M7/2e", - "m3HcYvikRKY7CvrjqlcBm110qr6e2DWNq06+VGf6E+UnQ13im5ikPJN3ohW3CxQbaJomVYAl/oh3kpUg", - "FSICkVwDYeuTaclzi0KnwIdL9PH85peKypE2qapfoaTNr9lh84IS1lRj0Nc2hqepXwaq3Vgy43nsX+le", - "ZFO2EXtDK60O9cWcX13i2DKGx9ipLhLtNfQEJkLrTOEbo8S2Vg3+3uVmNpxMhKL3Q6+4JlzKtwS3S8PF", - "PFEwTzUOoVArX4LvIYoUVjA1fmaCypXUFnERmum4FDWkZYOACt2Z98nWFUQXQJcmuTB2ifSX/lLmlp9E", - "XVaqTUlfabNP4k1o9Mn5DOia5unl0QT6q39wCciKWLpw6YfjncEt5FJ5bcYf2lX4uWkJGudCSBtQc+aI", - "CKVgGtKTetzgzafN/wMAAP//k6kvPuU+AAA=", + "H4sIAAAAAAAC/+xbW3PbtrP/Khicvh1ZctokM/V5qev04jaZ+NjJ5CFxNRCwklCBAAKAtvX36Lv/BxdS", + "FAnKluK0dqcvbSwAi8VvF3vD8hZTVWglQTqLj26xpXMoSPjnsdZiea5Kx+XsHD6XYJ3/WRulwTgOYRKx", + "ls9kUS3nDorwj28MTPER/p/Rmvwo0R6dl9LxAjxpOK7X49UAu6UGfISJMWSJV6sBNvC55AYYPvq4sdVl", + "PVdN/gQaFp8YIA4uqFFC9PNrHJ8SGkYYWGq4dlxJfITfnpyiahQZmIIBSQEpg4SiRCAbCCNN3BwPMNyQ", + "QovIbFxjh8yUnA1ns5ED68J/jvx/cM2rdYbLmeeVsy4Dr0AboMQBQ0RwYtFUGSRJAUP0NszxTDgyEYBM", + "RBBxNooTTqdIFdw5YAPk5oAYgUJJNAMJhjiwiEjE2XCD8T/VxOZ48xQz8DwMC/8Xx7jVgizD6ZB1XAhE", + "VQEWTY0qEtLDJSnE/Tm2mtAM27+XEzAS/P71rABsxb8Bq0pDwQ7R6UwqAwxNlkgqedBYOiF0AZLZYW53", + "dS3BjHMSTYqOwgzEGSotsLA7La1TBZiDKaFczpDxdwGR0s2V4f8hfn12LwMzbp1ZjqkBBtJxIna4d2nx", + "Sb327jtXXZfchXsFAhyweOO6Vy0i0jmCdcSVYcJasCxS6p64xQ73UxKBHEc/SVuaXUxAh7tqcMz4LK3u", + "uby99+Zf9Xwc6vkrEOHm52C1kha6elAolpHISWkMSIfmYTWKyobC3KYpUovc+bVRMwPWdsmepRGkwVCQ", + "jsyinIUizCPsGQu4egM3VaYgDh/hqVDE+4+C3PCiLPDRs8PDAS64jH8d1izIspiASdfLuDEjLnO2D3OQ", + "Tdsc5oZrV+/oFx54rcADLEshvK3HR86UcNfdDBDl5PBa0cVFfek3ZQA33I1pEkTPflw6mMXDCWLdOIpk", + "TOdEzsK6mnku3cvnOLewYXSkR+4jNqWU/hgDzJQMsjVGGTzA14T7iKdxlJ4DJ5pZrnI4nCmTsUYbEtrF", + "quhErtaNly9efPeioR3PckBoo5yiSjShmDun8SD8L7hX6v8qmb4bgsBcYqVBO3t6oygwb50DUG+IDraY", + "MR7jirNNG93z+zb70dCzVYaBLkflRHA7f396RuiCzKDXYYSQb0tAFNzNgVHKHRgQxPErQMNrYosQLA7R", + "K5iSUjiLnELa8CviYMS4dSOidZynDNKeG7r5+zDrEDsHyRjOzhnmqseZaWLttTJ5l1ZaMD0K2NKEQL+x", + "oEE4pw3J9Rwn+/22sn77Oe0+t8Mi8PhoSoSFwcO5Ib9lMJ4NdiZKCSByNx+VcPCmoc9ETlQpWW6fQQB9", + "zHUWkzBW2YiuHVgA6GPBr+CdIdMpp1kawbAR6vgVd8sxcRvGtukpdrdaWcsUDUR+XcNudeV/M54sXYTr", + "Ps4gRFRZSq6DRgPuNLjTXtUatdhO85pLpq7zPO1yuh4DXWPbNdaDpGHrw9cIbdHYdvKe8ezOmwKxTT93", + "d3jj/tFtChKN65brUBqRt3Hbzs/l7B0xM8ic/n65wC63467D73t3LAigTpltXrerkm1ULJgrTmF8P2fR", + "o5XjnnCiRX6LVvalolu9Bw11I7aTfetJAYPBjIFkbriZivXL8E4PlQmlokcCcwUsaPn9s64QlYblhL2V", + "YtkKvtcOT6ke5xtvwoNX/wY4Blb9Wt8N6pMo8aAR3luntA6/VRF+VW24zAi25GMdw8H7HqSOH0PYWWq2", + "ozLlShy1vibcN7EYrFOPhu5u7L3ljtT89ge6XVR2PtUWk9o8rZ80wKmmuiP/eycKHSByFu19YGbvavJP", + "N1pwyp1YIkIpaAcM8aIoY/HUBbcxRG+4tSH7N2g9NFtXny0iBpABzxQwNIGpMoCCRvtlPpVPhaJhhNH7", + "XHyE//j4xw+fPtnL//3Bzsm3L14efSQH08OD7y9vXz5ffZOTFtzosMeYS+uIEMAaFbCexEcboEpG/BGd", + "A10AQ2RGPInAG6GuJALVJFFFSBlUSgYmzCqI93OSSApIKLpoHeW+J3hEtafgEGlpuFte+A1SIA/EgDku", + "4x2Lf/1cXaXfPrwLlqkJ9G8f3iGnFiBjWZwHztwSaaOuOAMTzIIn70PKQG6NS8jp/RH8+mrPTfIXc2Xc", + "gc8BGPpcgllWmymDPsDkQtEFOESVlECryhT3C8NkXEVqcYv1zkTz38HD5d2knCq/MVXSxWuyah/ylSk5", + "QyevT5EgpaTz8FDAUEGkv+IorOQSzEEocrLqGYZof79ixWyABF/AJzkLrwneCRo7QIw4MiEW7CAQvIZJ", + "NTb8FNjlLlTyagbwAPvRyNbh8NnwMPhsDZJojo/wd+GnaK2CQEdE89HVs1EsFfpfUiy4ecJfwFUFt42i", + "Ig7EY957yuLEWLMM8goePZQuw2bfHh5WSKZ4uwHB6E8by0dRn+/S9lZlNIhqk+c4I3jhF4ff/YUbX8RI", + "D5WSXBEey4HhPpVFQcwywdnG0ZGZDVWIKIkBjnjjS7+0ElPUHNuQ0yb8r7l1F2nOF4K/SyCUYtauuelg", + "U9Xrq4Ns4uLZr58NbH2OCpo00sTGR9o2A0TzIRVHowfW/ajY8sEUIfdWu9q0sD4MXXXk8OzBWGjBfxfc", + "qIotN1GPB2nhvh32rkqOILxbhfgiK5Hmu9ZXkkju6exeEjn82yQSUWtLJB6kJREEN9y66FpUil/EMj6A", + "2J3FdcvZKtp5n0h0xRUfRmtxaWJIAQ6M3+I2+tAU8iYXGoL+TaAHDdDaMfTlVxTC5qPu3UKokqnVAD8/", + "fN7/yJimS+XQNNSbNqUWt93pHg3yZvwXcE8T+R3V/0sR9370C82WvwcjH5eVut92/RjGv7JIHt4e3vVC", + "8dhsY4QZ+aXpQm5axRugZeOCJantJXGqioJIZke36V+rfumflzLyfBKnfg0NGGSJ0HrDnSi1XswJdyEh", + "Cql2FD0wlGgjp2rAq+Q8agDjcjbsyZfsUlLc5KL9TNV5UPpbrU4shLC2rfhC63NeyraLXgtsL52UUz7r", + "je1rp3AS5z1C19CuIXQEcUaMXSfA6cBdm65z0/bF1CoB9l6oxpmPENd8hXCjyN2PeXUwtIBlbL1aP1p0", + "oU+dAJYqHYxEDcoe4NdV/O3In6kYwD462HdJgRvP3zunwcgDVVUCHjwk2qC+lxw/l1DC3XL8/zDtiUWs", + "uYeorrzC0QKGsAXvz41Za6A/J1gqwPtxNiCA2K1In1Z16PM09xHC3WoO629CNaCV5U6Z5d2vLI25g4rm", + "5T2sXw0YokQqySkRaE0rJNOCWIdsSSlYOy39cIAWpV2+NEoAwkLotX5CqDdYPyWE6njjLQRdKVEWTTWq", + "ymHb1Mc6ta0Ecx4n/JvHfO0kN+J870SmEtxexrnxdp2XevgGJhXF0tynI/rcBzyPTdx92cWGzM/AWG5d", + "apJW5iB+CgQsdU0iU8umqwShm+ROFRjFJ9l7RFwbnT9PPvTa7GO6T/QVF6AKr0wUHB+y00cJlWzqBXvI", + "qG61zF/SCz/8D635XcQW/e33I0x6kGKedUpvA1rpfyzOoUvoLpyVbicI18oshCLMous5F4B07MPyGs+I", + "I/uJoeSjZhfSdoPUaIh5mlJpdlBlyh6xTR8Yen+KKlR8Jh7S7FwBJLMAvT9/bb9YFqPbsOdqlLbovymJ", + "6ZaA/rriZ8RmG52qZS59kICrJtncRx9fKT7p+wBjlYKUR/LMeM3dHKXetKZKFeBIuOKtYCWeChGJiDBA", + "2PJgUnLhUGw0eX+KPhxfvKmo7KmTuvrAK69+zea1JxSw5nrunkrA6hPdw26ieywRZFoBL349/vbFy9z3", + "5tyi+rS7Z89+wffdBaed9Dl+p8aQ5X5TnzxHnULXxFa9fC3Fjidv+78pF6llq127yelzahWvNLGvFez4", + "7BSnDlI8wl7dEtFOD1tkInaLFSBd4KR6noGQK/qZDcOYxNX57jOtiXWoNcH10liLyrwR5XrlUHweWkBo", + "m0sUrmFiw8wMlTNlHOIy9tZyJWtIywYBHZu1b7PdWlFsNrswNUZ1l74pheMHSZaVaHOnr6TZJfEq9rYJ", + "PgW6pCK/PKlAd/XPPmi6Jo7OfcjkeWdwBULpIM303W2Fn5+WoXEspXIRNa+O4ZrZxulJPW7x6nL13wAA", + "AP//pU/ObfRCAAA=", } // GetSwagger returns the content of the embedded swagger specification file diff --git a/internal/core/ports/services_ports.go b/internal/core/ports/services_ports.go index 6b16a12..a8efa86 100644 --- a/internal/core/ports/services_ports.go +++ b/internal/core/ports/services_ports.go @@ -159,6 +159,7 @@ const ( RuntimeWorkerModeCreate RuntimeWorkerMode = "create" RuntimeWorkerModeUpdate RuntimeWorkerMode = "update" RuntimeWorkerModeRestore RuntimeWorkerMode = "restore" + RuntimeWorkerModeInspect RuntimeWorkerMode = "inspect" ) type RuntimeWorkerAction struct { diff --git a/internal/core/services/registry/oci.go b/internal/core/services/registry/oci.go index 965dac6..b90d726 100644 --- a/internal/core/services/registry/oci.go +++ b/internal/core/services/registry/oci.go @@ -837,7 +837,9 @@ func (c *OciClient) PushWithOptions(folder string, repo string, tag string, over dataExists, _ := utils.FileExists(dataDir) if dataExists { var explicitChunks []*domain.Chunks - if scrollFile != nil { + // A backup is a complete runtime snapshot, not the release's curated + // chunk selection. Include files created since installation as well. + if scrollFile != nil && !options.PreserveReleaseManifest { explicitChunks = scrollFile.Chunks } chunks, err := utils.AutoChunkDataDir(dataDir, explicitChunks) diff --git a/internal/core/services/registry/oci_test.go b/internal/core/services/registry/oci_test.go index 8efc65c..6aa7168 100644 --- a/internal/core/services/registry/oci_test.go +++ b/internal/core/services/registry/oci_test.go @@ -308,7 +308,16 @@ func TestPreserveReleaseManifestRoundTrip(t *testing.T) { } client := &OciClient{credentialStore: NewCredentialStore(nil), plainHTTP: true} repo := registryHost + "/test/preserve-release" - if _, err := client.PushWithOptions(folder, repo, "backup", nil, false, nil, TransferOptions{PreserveReleaseManifest: true}); err != nil { + if err := os.MkdirAll(filepath.Join(folder, "data"), 0755); err != nil { + t.Fatal(err) + } + for _, name := range []string{"declared.txt", "runtime-created.txt"} { + if err := os.WriteFile(filepath.Join(folder, "data", name), []byte(name), 0644); err != nil { + t.Fatal(err) + } + } + scrollFile := &domain.File{Chunks: []*domain.Chunks{{Name: "declared", Path: "declared.txt"}}} + if _, err := client.PushWithOptions(folder, repo, "backup", nil, false, scrollFile, TransferOptions{PreserveReleaseManifest: true}); err != nil { t.Fatalf("push backup: %v", err) } pullDir := t.TempDir() @@ -322,6 +331,12 @@ func TestPreserveReleaseManifestRoundTrip(t *testing.T) { if string(got) != string(want) { t.Fatalf("restored release manifest = %s, want %s", got, want) } + for _, name := range []string{"declared.txt", "runtime-created.txt"} { + got, err := os.ReadFile(filepath.Join(pullDir, "data", name)) + if err != nil || string(got) != name { + t.Fatalf("backup omitted runtime data %s: %s (%v)", name, got, err) + } + } } func TestPreserveReleaseManifestRequiresSnapshotDescriptor(t *testing.T) { diff --git a/internal/runtime/docker/workers.go b/internal/runtime/docker/workers.go index 7b635ed..4b818b5 100644 --- a/internal/runtime/docker/workers.go +++ b/internal/runtime/docker/workers.go @@ -113,8 +113,10 @@ func (b *Backend) SpawnPullWorker(ctx context.Context, action ports.RuntimeWorke if err := b.pullImage(ctx, b.config.WorkerImage); err != nil { return nil, err } - if err := b.prepareWritableRoot(ctx, root); err != nil { - return nil, err + if action.Mode != ports.RuntimeWorkerModeInspect { + if err := b.prepareWritableRoot(ctx, root); err != nil { + return nil, err + } } registryConfig, err := json.Marshal(struct { Registries []domain.RegistryCredential `json:"registries"` @@ -122,7 +124,7 @@ func (b *Backend) SpawnPullWorker(ctx context.Context, action ports.RuntimeWorke if err != nil { return nil, err } - rootMount, err := DockerMount(root, action.MountPath, false, "") + rootMount, err := DockerMount(root, action.MountPath, action.Mode == ports.RuntimeWorkerModeInspect, "") if err != nil { return nil, err } diff --git a/internal/runtime/kubernetes/resources.go b/internal/runtime/kubernetes/resources.go index 5d3b40a..2fcf36d 100644 --- a/internal/runtime/kubernetes/resources.go +++ b/internal/runtime/kubernetes/resources.go @@ -88,11 +88,20 @@ func workerPullJobSpec(namespace string, jobName string, pvc string, image strin if action.PreserveReleaseManifest { command = append(command, "--preserve-release-manifest") } + if action.Mode == ports.RuntimeWorkerModeInspect { + // No credential shell or recursive ownership changes for read-only inspection. + command = append([]string{"druid"}, command[4:]...) + } job := helperJobSpec(namespace, jobName, pvc, image, command, imagePullSecret, map[string]string{ labelComponent: "worker-pull", labelRuntimeID: runtimeLabel(action.RuntimeID), }) container := &job.Spec.Template.Spec.Containers[0] + if action.Mode == ports.RuntimeWorkerModeInspect { + for i := range container.VolumeMounts { + container.VolumeMounts[i].ReadOnly = true + } + } runAsRoot := int64(0) runAsNonRoot := false container.SecurityContext = &corev1.SecurityContext{ diff --git a/internal/runtime/kubernetes/resources_test.go b/internal/runtime/kubernetes/resources_test.go index 766468a..18e43c4 100644 --- a/internal/runtime/kubernetes/resources_test.go +++ b/internal/runtime/kubernetes/resources_test.go @@ -487,6 +487,21 @@ func TestWorkerPullJobSpecRunsDruidWorkerPull(t *testing.T) { } } +func TestInstalledReleaseInspectionMountsReadOnlyWithoutChown(t *testing.T) { + action := ports.RuntimeWorkerAction{Mode: ports.RuntimeWorkerModeInspect, RuntimeID: "inspect", Artifact: "registry.local/source:tag", MountPath: "/scroll"} + job := workerPullJobSpec("fixture", "inspect", "runtime-pvc", "druid:test", action, "", "", true, "druid-cli") + container := job.Spec.Template.Spec.Containers[0] + command := strings.Join(container.Command, " ") + if !strings.HasPrefix(command, "druid worker pull ") || strings.Contains(command, "chown") || !strings.Contains(command, "--mode inspect") { + t.Fatalf("unsafe inspection command: %s", command) + } + for _, mount := range container.VolumeMounts { + if !mount.ReadOnly { + t.Fatalf("inspection mount is writable: %#v", mount) + } + } +} + func TestSpawnPullWorkerCreateUsesFinalPVCAndWorkerJob(t *testing.T) { client := fake.NewSimpleClientset() backend := NewWithClient(Config{Namespace: "druid", PullImage: "druid-cli:test"}, client) diff --git a/test/integration/internal/e2e/harness.go b/test/integration/internal/e2e/harness.go index 7c4f673..4e15b91 100644 --- a/test/integration/internal/e2e/harness.go +++ b/test/integration/internal/e2e/harness.go @@ -462,7 +462,7 @@ func DockerHostAddress(t *testing.T) string { return gateway } -func UnixJSONRequest(t *testing.T, socket string, method string, path string, body string) string { +func UnixJSONRequest(t *testing.T, socket string, method string, path string, body string, expectedStatus ...int) string { t.Helper() transport := &http.Transport{ DialContext: func(ctx context.Context, network string, addr string) (net.Conn, error) { @@ -490,7 +490,7 @@ func UnixJSONRequest(t *testing.T, socket string, method string, path string, bo if err != nil { t.Fatal(err) } - if resp.StatusCode >= 400 { + if (len(expectedStatus) == 0 && resp.StatusCode >= 400) || (len(expectedStatus) > 0 && resp.StatusCode != expectedStatus[0]) { t.Fatalf("%s %s failed with %d: %s", method, path, resp.StatusCode, data) } return string(data) diff --git a/test/integration/kubernetes/kubernetes_cli_test.go b/test/integration/kubernetes/kubernetes_cli_test.go index 182b20e..31fbaba 100644 --- a/test/integration/kubernetes/kubernetes_cli_test.go +++ b/test/integration/kubernetes/kubernetes_cli_test.go @@ -4,6 +4,7 @@ package kubernetes_test import ( "context" + "encoding/json" "fmt" "net/http" "os" @@ -29,6 +30,25 @@ func TestKubernetesBackendCLIComplexLifecycle(t *testing.T) { namespace := "druid-cli-e2e-" + suffix name := "k8s-cli-" + suffix fixture := e2e.WriteFixture(t, filepath.Join(t.TempDir(), "scroll"), name, port, routePort) + fixtureYAML, err := os.ReadFile(filepath.Join(fixture.Dir, "scroll.yaml")) + if err != nil { + t.Fatal(err) + } + fixtureYAML = append(fixtureYAML, []byte("\nchunks:\n - name: finite\n path: finite.txt\n skip_update: true\n - name: version\n path: version.txt\n")...) + if err := os.WriteFile(filepath.Join(fixture.Dir, "scroll.yaml"), fixtureYAML, 0644); err != nil { + t.Fatal(err) + } + if err := os.MkdirAll(filepath.Join(fixture.Dir, "data"), 0755); err != nil { + t.Fatal(err) + } + writeFixtureData := func(file, content string) { + t.Helper() + if err := os.WriteFile(filepath.Join(fixture.Dir, "data", file), []byte(content), 0644); err != nil { + t.Fatal(err) + } + } + writeFixtureData("finite.txt", "initial") + writeFixtureData("version.txt", "v1") workerImage := e2e.BuildDockerImage(t, "druid-cli-e2e:"+name) importImageIntoCurrentCluster(t, workerImage) pushArtifact := fmt.Sprintf("127.0.0.1:%d/druid-e2e/%s:v1", registryPort, name) @@ -120,6 +140,58 @@ func TestKubernetesBackendCLIComplexLifecycle(t *testing.T) { t.Fatalf("stopped status = %s, want stopped", stopped.Status) } waitKubernetesResourcesGone(t, namespace, pvc, "statefulset,job,pod") + + // Exercise actual changed workers against an isolated PVC/registry, without + // consuming any canonical local ReservedPorts or touching user deployments. + releasePath := "/api/v1/scrolls/" + created.ID + "/release" + readRelease := func() map[string]string { + t.Helper() + var release map[string]string + if err := json.Unmarshal([]byte(e2e.UnixJSONRequest(t, socket, http.MethodGet, releasePath, "")), &release); err != nil { + t.Fatal(err) + } + return release + } + originalRelease := readRelease() + originalManifest := readPVCFile(t, namespace, pvc, "manifest.json") + backupArtifact := strings.TrimSuffix(runtimeArtifact, ":v1") + "-backup:v1" + e2e.UnixJSONRequest(t, socket, http.MethodPost, "/api/v1/scrolls/"+created.ID+"/backup", fmt.Sprintf(`{"artifact":%q}`, backupArtifact)) + writeFixtureData("version.txt", "v2") + writeFixtureData("finite.txt", "upstream replacement must not win") + e2e.RunEnv(t, []string{"DRUID_REGISTRY_PLAIN_HTTP=true"}, bins.Druid, "push", pushArtifact, fixture.Dir) + acceptedDigest := resolveRegistryCandidate(t, fmt.Sprintf("http://127.0.0.1:%d/v2/druid-e2e/%s/manifests/v1", registryPort, name)) + // Move the tag again after acceptance: update must still install v2. + writeFixtureData("version.txt", "v3-unaccepted") + e2e.RunEnv(t, []string{"DRUID_REGISTRY_PLAIN_HTTP=true"}, bins.Druid, "push", pushArtifact, fixture.Dir) + acceptedArtifact := strings.TrimSuffix(runtimeArtifact, ":v1") + "@" + acceptedDigest + updatePath := "/api/v1/scrolls/" + created.ID + "/update" + acceptedBody := fmt.Sprintf(`{"artifact":%q,"expected_installed_digest":%q}`, acceptedArtifact, originalRelease["digest"]) + e2e.UnixJSONRequest(t, socket, http.MethodPost, updatePath, acceptedBody) + // A concurrent restore/update invalidates an earlier acceptance. The + // precondition is checked against the actual PVC before any mutation. + e2e.UnixJSONRequest(t, socket, http.MethodPost, updatePath, acceptedBody, http.StatusConflict) + if got := readPVCFile(t, namespace, pvc, "data/version.txt"); strings.TrimSpace(got) != "v2" { + t.Fatalf("update followed moved tag: %q", got) + } + if got := readPVCFile(t, namespace, pvc, "data/finite.txt"); !strings.Contains(got, "finite-ok") { + t.Fatalf("update overwrote protected state: %q", got) + } + if got := readRelease(); got["digest"] != acceptedDigest { + t.Fatalf("installed release=%v, want %s", got, acceptedDigest) + } + e2e.UnixJSONRequest(t, socket, http.MethodPost, "/api/v1/scrolls/"+created.ID+"/restore", fmt.Sprintf(`{"artifact":%q,"restart":false}`, backupArtifact)) + if got := readPVCFile(t, namespace, pvc, "data/version.txt"); strings.TrimSpace(got) != "v1" { + t.Fatalf("restore contents=%q", got) + } + if got := readPVCFile(t, namespace, pvc, "data/record-env.txt"); !strings.Contains(got, "USER_ENV=finite") { + t.Fatalf("backup omitted runtime-created data outside release chunks: %q", got) + } + if got := readPVCFile(t, namespace, pvc, "manifest.json"); got != originalManifest { + t.Fatal("restore changed original installed descriptor bytes") + } + if got := readRelease(); got["digest"] != originalRelease["digest"] { + t.Fatalf("restored release=%v, want %v", got, originalRelease) + } deleted := e2e.RunClient(t, bins, socket, "delete", created.ID) if !strings.Contains(deleted, `"status": "deleted"`) { t.Fatalf("delete response = %s, want deleted status", deleted) diff --git a/test/integration/kubernetes/registry_resolution_test.go b/test/integration/kubernetes/registry_resolution_test.go new file mode 100644 index 0000000..1a2f051 --- /dev/null +++ b/test/integration/kubernetes/registry_resolution_test.go @@ -0,0 +1,44 @@ +//go:build integration && kubernetes + +package kubernetes_test + +import ( + "fmt" + "io" + "net/http" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/highcard-dev/daemon/test/integration/internal/e2e" +) + +func TestRegistryCandidateResolution(t *testing.T) { + bins := e2e.BuildBinaries(t) + port := e2e.StartRegistry(t) + fixture := e2e.WriteFixture(t, filepath.Join(t.TempDir(), "scroll"), "registry-resolution", 8080, 8081) + artifact := fmt.Sprintf("127.0.0.1:%d/druid-e2e/registry-resolution:v1", port) + e2e.RunEnv(t, []string{"DRUID_REGISTRY_PLAIN_HTTP=true"}, bins.Druid, "push", artifact, fixture.Dir) + resolveRegistryCandidate(t, fmt.Sprintf("http://127.0.0.1:%d/v2/druid-e2e/registry-resolution/manifests/v1", port)) +} + +func resolveRegistryCandidate(t *testing.T, url string) string { + t.Helper() + request, err := http.NewRequest(http.MethodGet, url, nil) + if err != nil { + t.Fatal(err) + } + request.Header.Set("Accept", "application/vnd.oci.image.manifest.v1+json") + response, err := (&http.Client{Timeout: 10 * time.Second}).Do(request) + if err != nil { + t.Fatal(err) + } + defer response.Body.Close() + body, _ := io.ReadAll(io.LimitReader(response.Body, 1024)) + digest := response.Header.Get("Docker-Content-Digest") + if response.StatusCode != http.StatusOK || !strings.HasPrefix(digest, "sha256:") { + t.Fatalf("candidate lookup: status=%d body=%s", response.StatusCode, body) + } + return digest +}