From 2a99dcd825b266829b296cb2b856c9932856c233 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marc=20Schottst=C3=A4dt?= Date: Sun, 20 Sep 2026 02:23:00 +0200 Subject: [PATCH 1/4] fix(runtime): stage scroll updates --- apps/druid/adapters/cli/worker_pull.go | 112 ++++++++++--------------- apps/druid/adapters/cli/worker_test.go | 58 +++++++++---- 2 files changed, 84 insertions(+), 86 deletions(-) diff --git a/apps/druid/adapters/cli/worker_pull.go b/apps/druid/adapters/cli/worker_pull.go index 4040629..4f90622 100644 --- a/apps/druid/adapters/cli/worker_pull.go +++ b/apps/druid/adapters/cli/worker_pull.go @@ -161,7 +161,48 @@ func pullWorkerUpdate(root string, artifact string, oci ports.OciRegistryInterfa } skipData := map[string]bool{} collectSkipUpdatePaths(skipData, "", scroll.Chunks) - return mergePulledRoot(tmp, root, skipData) + 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 skip == "." { + skip = "" + } + source := filepath.Join(root, domain.RuntimeDataDir, filepath.FromSlash(skip)) + target := filepath.Join(stage, domain.RuntimeDataDir, filepath.FromSlash(skip)) + 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 pullWorkerRestore(root string, artifact string, oci ports.OciRegistryInterface) error { @@ -278,75 +319,6 @@ 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) - 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 err != nil { - return err - } - rel, err := filepath.Rel(srcData, srcPath) - if 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 { diff --git a/apps/druid/adapters/cli/worker_test.go b/apps/druid/adapters/cli/worker_test.go index ac18897..a7d1f78 100644 --- a/apps/druid/adapters/cli/worker_test.go +++ b/apps/druid/adapters/cli/worker_test.go @@ -85,24 +85,50 @@ 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) } - 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") + 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) + } + 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 TestWorkerRestoreStagesBeforeReplacingRoot(t *testing.T) { From f53ef133e8876ea3e6925f742893e096557ce7ec Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marc=20Schottst=C3=A4dt?= Date: Mon, 28 Sep 2026 01:59:34 +0200 Subject: [PATCH 2/4] feat(runtime)!: apply accepted release digests Require explicit immutable update references; reconciliation must not undo restores or follow moved tags. Read the installed descriptor from the runtime volume. Preserve rollback recovery files, all backup data, and old/new protected paths during updates. Serialize command admission with maintenance without holding locks during command waits. Preserve symlinks without following protected parent paths. BREAKING CHANGE: runtime update callers must supply repository@sha256. Old tag/empty update requests fail before stopping the workload. Verified full Go suite, focused race tests, and rebuilt-image isolated Kubernetes lifecycle. See docs/runtime-lifecycle-verification.md for scope and remaining product work. --- api/openapi.yaml | 35 ++- apps/druid/adapters/cli/client/update.go | 11 +- apps/druid/adapters/cli/worker_pull.go | 102 +++++++- .../druid/adapters/cli/worker_release_test.go | 69 +++++ apps/druid/adapters/cli/worker_test.go | 28 ++ .../adapters/daemonclient/openapi_client.go | 5 +- apps/druid/adapters/http/handlers/routes.go | 1 + .../adapters/http/handlers/scroll_handler.go | 20 +- apps/druid/core/services/runtime_access.go | 21 +- .../runtime_command_maintenance_test.go | 94 +++++++ apps/druid/core/services/runtime_release.go | 32 +++ .../core/services/runtime_session_commands.go | 24 +- .../core/services/runtime_session_queue.go | 40 +-- .../druid/core/services/runtime_supervisor.go | 12 +- .../core/services/runtime_supervisor_test.go | 47 ++-- apps/druid/core/services/runtime_update.go | 13 +- .../core/services/runtime_update_test.go | 30 +++ docs/runtime-lifecycle-verification.md | 32 +++ internal/api/generated.go | 239 ++++++++++++++---- internal/core/ports/services_ports.go | 1 + internal/core/services/registry/oci.go | 4 +- internal/core/services/registry/oci_test.go | 17 +- internal/runtime/docker/workers.go | 8 +- internal/runtime/kubernetes/resources.go | 9 + internal/runtime/kubernetes/resources_test.go | 15 ++ .../kubernetes/kubernetes_cli_test.go | 67 +++++ .../kubernetes/registry_resolution_test.go | 44 ++++ 27 files changed, 873 insertions(+), 147 deletions(-) create mode 100644 apps/druid/adapters/cli/worker_release_test.go create mode 100644 apps/druid/core/services/runtime_command_maintenance_test.go create mode 100644 apps/druid/core/services/runtime_release.go create mode 100644 apps/druid/core/services/runtime_update_test.go create mode 100644 docs/runtime-lifecycle-verification.md create mode 100644 test/integration/kubernetes/registry_resolution_test.go diff --git a/api/openapi.yaml b/api/openapi.yaml index 32be8a0..4c3edb3 100644 --- a/api/openapi.yaml +++ b/api/openapi.yaml @@ -209,10 +209,12 @@ 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. registry_credentials: type: array items: @@ -482,6 +484,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 +561,7 @@ paths: schema: type: string requestBody: - required: false + required: true content: application/json: schema: @@ -544,6 +573,8 @@ paths: application/json: schema: $ref: '#/components/schemas/RuntimeScroll' + '400': + description: An explicitly accepted SHA256 artifact reference is required '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 4f90622..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,6 +198,17 @@ func pullWorkerUpdate(root string, artifact string, oci ports.OciRegistryInterfa } skipData := map[string]bool{} collectSkipUpdatePaths(skipData, "", scroll.Chunks) + 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 } @@ -185,11 +233,19 @@ func pullWorkerUpdate(root string, artifact string, oci ports.OciRegistryInterfa 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 } @@ -205,6 +261,25 @@ func preserveSkippedUpdateData(root string, stage string, skipData map[string]bo 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 { if err := os.MkdirAll(root, 0755); err != nil { return err @@ -241,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) @@ -283,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) } @@ -294,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) } @@ -320,10 +402,20 @@ func collectSkipUpdatePaths(out map[string]bool, parent string, chunks []*domain } func copyPath(src string, dst string) error { - info, err := os.Stat(src) + info, err := os.Lstat(src) if err != nil { return err } + if info.Mode()&os.ModeSymlink != 0 { + target, err := os.Readlink(src) + if err != nil { + return err + } + if err := os.MkdirAll(filepath.Dir(dst), 0755); 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 { if walkErr != nil { 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 a7d1f78..cec8ff3 100644 --- a/apps/druid/adapters/cli/worker_test.go +++ b/apps/druid/adapters/cli/worker_test.go @@ -131,6 +131,29 @@ func TestPreserveSkippedUpdateDataKeepsMissingPathsMissing(t *testing.T) { } } +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) { root := t.TempDir() mustWrite(t, filepath.Join(root, "scroll.yaml"), "name: old\n") @@ -208,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..e715a6b 100644 --- a/apps/druid/adapters/http/handlers/routes.go +++ b/apps/druid/adapters/http/handlers/routes.go @@ -63,6 +63,7 @@ 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.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/scroll_handler.go b/apps/druid/adapters/http/handlers/scroll_handler.go index 68fc4cc..6d10c05 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,11 @@ 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)) if err != nil { + 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..e7d90f9 --- /dev/null +++ b/apps/druid/core/services/runtime_release.go @@ -0,0 +1,32 @@ +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 + } + installed, err := s.runPullWorker(ctx, s.runtimeBackend, ports.RuntimeWorkerModeInspect, 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..5672fa2 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) 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) 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..0261325 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,20 @@ import ( "go.uber.org/zap" ) +var ErrUnacceptedUpdate = errors.New("update requires an explicitly accepted sha256 artifact reference") +var acceptedUpdateReference = regexp.MustCompile(`^[^@\s]+@sha256:[a-f0-9]{64}$`) + func (s *RuntimeSupervisor) Update(id string, artifact string, registryCredentials []domain.RegistryCredential) (*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 - } - 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..470f908 --- /dev/null +++ b/apps/druid/core/services/runtime_update_test.go @@ -0,0 +1,30 @@ +package services + +import ( + "strings" + "testing" + + "github.com/highcard-dev/daemon/internal/core/domain" + 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) + 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") + } + }) + } +} diff --git a/docs/runtime-lifecycle-verification.md b/docs/runtime-lifecycle-verification.md new file mode 100644 index 0000000..5849006 --- /dev/null +++ b/docs/runtime-lifecycle-verification.md @@ -0,0 +1,32 @@ +# 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. +- 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 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, shared Team publication, and shared-stack 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..858d758 100644 --- a/internal/api/generated.go +++ b/internal/api/generated.go @@ -229,8 +229,8 @@ 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"` + // Artifact Explicitly accepted immutable target. Missing or mutable tag references are rejected before stopping the runtime. + Artifact string `json:"artifact"` RegistryCredentials *[]RegistryCredential `json:"registry_credentials,omitempty"` } @@ -379,6 +379,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 +595,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 +1199,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 +1649,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 +1950,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 +2283,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 +2715,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 +2990,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 +3202,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 +3399,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 +3422,57 @@ 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/+xbXXPbttL+Kxi8vXtpyWmTzBydm7pJP9wmEx87mVwkrgYCVhIqEGAA0LaOR//9DD5I", + "kSIoWbLT2p3etLEALHafXewuFstbTFVeKAnSGjy6xYbOISf+nydFIZbnqrRczs7hSwnGup8LrQrQloOf", + "RIzhM5lXy7mF3P/jGw1TPML/N1yTH0baw/NSWp6DIw0n9Xq8yrBdFoBHmGhNlni1yrCGLyXXwPDoU2ur", + "y3qumvwB1C9+pYFYuKBaCdHPr7Z8SqgfYWCo5oXlSuIRfvfqFFWjSMMUNEgKSGkkFCUCGU8YFcTOcYbh", + "huSFCMyGNWbAdMnZYDYbWjDW/2fk/oNrXo3VXM4cr5x1GXgNhQZKLDBEBCcGTZVGkuQwQO/8HMeEJRMB", + "SAcEEWfDMOF0ilTOrQWWITsHxAjkSqIZSNDEgkFEIs4GLcb/UBOT4s1RTMDzMCz8O4xxUwiy9NIhY7kQ", + "iKocDJpqlUekB0uSi7tzbApCE2z/Vk5AS3D717M8sBX/GowqNQUzQKczqTQwNFkiqeRRY+mE0AVIZgap", + "3dW1BD1OaTQaOvIzEGeoNMD87rQ0VuWgj6aEcjlD2p0FREo7V5r/l7j1yb00zLixejmmGhhIy4nY49zF", + "xa/qtbvPXHVcUgfuNQiwwMKJ6x61gEhHBGOJLf2EtWJZoNSVeIMd7qZEAimOfpSm1Pu4gA531eCY8Vlc", + "3XN4e8/NP+b5OMzzFyDCzs/BFEoa6NpBrlhCI69KrUFaNPerUTA25Oc2XZFapOQvtJppMKZL9iyOoAI0", + "BWnJLOhZKMIcwo4xj6tzcFOlc2LxCE+FIi5+5OSG52WOR8+OjzOccxn+Oq5ZkGU+AR2Pl7ZjRmxCto9z", + "kE3f7Of6Y1fv6BYeOavAGZalEM7X45HVJew6mx6ilB7eKLq4qA99Wwdww+2YRkX07MelhVkQThBjx0El", + "YzoncubX1cxzaV8+x6mFDacjHXKfsC6ldGJkmCnpdau10jjD14S7jKchSo/AkWaSqxQOZ0onvFFLQ/t4", + "lSKSq23j5YsX371oWMezFBCFVlZRJZpQzK0tcOb/58MrdX+VrNgNgWcustKgnZReKwrMeWcP1FtSeF/M", + "GA95xVnbR/f8vs1/NOxslWCgy1E5EdzMP5yeEbogM+gNGD7l25IQ+XBzpJWyRxoEsfwK0OCamNwniwP0", + "GqakFNYgq1Ch+RWxMGTc2CEpijBPaVQ4bmj790EyIHYESTjOjgxz1RPMCmLMtdLpkFYa0D0GuGEJnn5j", + "QYNwyhpi6DmJ/vtd5f0OC9p9YYcF4PFoSoSB7OHCkNvSO88GOxOlBBC5X4yKODjX0OciJ6qULLVP5kEf", + "8yKJiR+rfETXDywAihPBr+C9JtMpp0ka3rERavkVt8sxsS1n24wU+3utpGcKDiK9ruG3uvq/GU+WNsB1", + "l2DgM6okJdtBowF3HNxrr2qNWmynec0lU9dpnvaRrsdB19h2nXUWLWwtfI3QFovdvLwnIrt1rkBss8/9", + "A964f3SbgQTnuuU4lFqkfdw2+bmcvSd6Bgnp73YX2Od07BL+0LNjQAC1Sm+Lul2T3ETFgL7iFMZ3CxY9", + "VjnuSSc2yG+xyr6r6NboQX3diO3l33qugN5hhkQyNdy8ivXrcGeESqRSISKBvgLmrfzuty6flfrlhL2T", + "YrmRfK8DnlI9wTechAev/mU4JFb9Vt9N6qMqcdZI741VReF/qzL8qtpwmVBsycdFSAfvKkidP/q0syzY", + "nsaUKnHU9hpxb2ORra8eDdtt7b3ljNT89ie6XVT2lmqLS21K6yZlONZU9+T/4ItCB4iUR/vgmTm4mvzj", + "TSE45VYsEaEUCgsM8TwvQ/HU+rAxQG+5Mf72r9F6aLauPhtENCANjilgaAJTpQF5i3bL3FU+FooGAUYX", + "c/EI//7p9+8/fzaX//+9mZNvX7wcfSJH0+Ojf13evny++uaRV258OKGl5nZ54TaIaTAQDfqkDBYa/vqp", + "MsRfP77357qpgV8/vkdWLUCGojL3nNklKrS64gy0P1SOvEvIPLk1Lv5G7ERw66s92+Qv5krbI5dBM/Sl", + "BL2sNlMafYTJhaILsIgqKYFWdR3uFvrJuMpzwhbrnUnBfwMHlwsycqrcxlRJG4xstSnka11yhl69OUWC", + "lJLOfZmdoZxId0CQX8kl6CNfImTVIwYpnHWGelOGBF/AZznztXgXQrTJECOWTIgBk3mC1zCpxgafPbvc", + "+jpYzQDOsBsNbB0Png2OfcQrQJKC4xH+zv8UzrpX6JAUfHj1bBgKbe6XmEm1JfwZbFWuapXksCcebo2n", + "LEwMFT+vLx8PfeHPb/bt8XGFZMxWGxAM/zCh+BLseZe1b9QVvaraPIcZPoa9OP7uT9z4IuRJqJTkivBQ", + "TPPnqcxzopcRzk0cLZkZf4cPmshwwBtfuqWVmoLlmIae2vC/4cZexDn3BH+fNCJmfF1308GmqnZXgrRx", + "cezXRXdTy1FBE0ea2Lg81SSAaD5D4uD0wNgfFFs+mCGkXjpXbQ/rkrhVRw/PHoyFDfh3wY2qzKyNehBk", + "A/ftsHdNcgj+1cdH56RGmq9CX0kjqYenO2nk+C/TSEBtUyNBkA2NILjhxobQomLZUyzD84HZW123nK2C", + "n3dpeFdd4VmxVldBNMnBgnZb3IYYGhPGGEJ9ytwGOmuAtpmBXn5FJbSfRHcrobqKrDL8/Ph5/xNdnC6V", + "RVNfrWlrLWy71znK0m78Z7BPE/k9zf++iLs4ek+35c7B0OVlZdHvu37w419ZJQ/vD3fV9x+bbwwwI7c0", + "Hsi2V7wBWjYOWNTaQRqnKs+JZGZ4G/+16tf+eSkDz6/C1K9hAVmSCK033IvSxnsz4dZfiPxFNageGIq0", + "kVU14NXVNlgA43I26LkvmaWkuMnF5iNP5znmL/U6oYzANn3FPb3PeSk3Q/RaYQfZpJzyWW9uXweFV2He", + "IwwNmzWEjiLOiDbrC3AUuOvTi9S0QzE1SoC5E6ph5iPENV1fa5WI+zGvBEMLWIbGpXXJvwt9fEc3VBXe", + "SdSgHAB+XQPfjvyZCgnso4N9nytw4/F472swckBVlYAHT4la1A/S45cSStitx//4aU8sY00943T15UXz", + "GMIWvL80Zq2B/hJhqQDvx1mDAGK2In0qjSVCADuPcx8h3ButVf0tnBoKZbhVern7jaIxN6toXt7B+9WA", + "IUqkkpwSgda0/GVaEGORKSkFY6alG/bQorjLfbMEIMynXrzmpN4g0lM6VMcbLwnoSokyb5pRVQ7bZj7G", + "qm0lmPMw4Z97zNe+5Aac73yRqRR3kHNuvPymte6/IIlFsTj36ag+9fnLY1N33+2ipfMz0IYbG1uMlT4K", + "H9IAiz2HSNe66RqB78XYaQLD8KB5h4yr1Tfz5FOvdhfQXbKvsABVeCWy4PAMHFv6K93UCw7QUd2omD6k", + "F274b1rzuwgN7tvPh5/0IMU8Y1WxDWhV/G1x9j02u3BWxeYF4VrphVCEGXQ95wJQEbqYnMUzYslhaij5", + "sNnDs90hNdpJnqZWmv1HibJHaHIHhj6cogoVdxP31+xUASSxAH04f2PurYvhrd9zNYxb9J+UyPSGgv68", + "4mfAZhudquEstvPjqsU09cnEV8pP+j5fWMUk5ZE8M15zO0exs6tpUjlY4o/4RrISpEJEIiI0ELY8mpRc", + "WBQaTT6coo8nF28rKgfaZFF9HpU2v2br1xNKWFMda08lYXUX3ePuRfdEIkg00l38cvLti5epr7W5QbW0", + "97w9B1Y3A9aUi9hjtVlsSRlg7IyuTKevd+vk7BTHhkk8xM4+ItFO01lgIrR35SCt56R6TwF/uXMzG54s", + "4tv5zDGuCYWjNcH10lA8SjzqpJrbUHjPWYDvc4sUrmFi/MwElTOlLeIytJJyJWtIywaBIvQm3ybbqxCd", + "A12Y5MLYydRd+rYUlh9FXVaqTUlfabNL4nVoRhN8CnRJRXp5NIHu6p9clnNNLJ27HMfxzuAKhCq8NuNn", + "phV+blqCxomUygbUnDn6c2Ea0pN63ODV5ep/AQAA//8lqbo540EAAA==", } // 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/kubernetes/kubernetes_cli_test.go b/test/integration/kubernetes/kubernetes_cli_test.go index 182b20e..a8d1d0e 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,53 @@ 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 + e2e.RunClient(t, bins, socket, "update", created.ID, acceptedArtifact) + 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 +} From 417ab01dd30276515e5e62cb299a1879ad961587 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marc=20Schottst=C3=A4dt?= Date: Mon, 28 Sep 2026 02:21:17 +0200 Subject: [PATCH 3/4] fix(runtime): reject stale update approvals Check the installed descriptor under the maintenance lock before stopping. Keep the accepted baseline unchanged on conflict and return HTTP 409. Verified with the full Go suite and rebuilt-image Kubernetes lifecycle. --- api/openapi.yaml | 6 + .../adapters/http/handlers/scroll_handler.go | 5 +- apps/druid/core/services/runtime_release.go | 6 +- .../core/services/runtime_supervisor_test.go | 4 +- apps/druid/core/services/runtime_update.go | 15 ++- .../core/services/runtime_update_test.go | 24 +++- docs/runtime-lifecycle-verification.md | 4 +- internal/api/generated.go | 111 +++++++++--------- test/integration/internal/e2e/harness.go | 4 +- .../kubernetes/kubernetes_cli_test.go | 7 +- 10 files changed, 123 insertions(+), 63 deletions(-) diff --git a/api/openapi.yaml b/api/openapi.yaml index 4c3edb3..594134e 100644 --- a/api/openapi.yaml +++ b/api/openapi.yaml @@ -215,6 +215,10 @@ components: type: string 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: @@ -575,6 +579,8 @@ paths: $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/http/handlers/scroll_handler.go b/apps/druid/adapters/http/handlers/scroll_handler.go index 6d10c05..3b5e906 100644 --- a/apps/druid/adapters/http/handlers/scroll_handler.go +++ b/apps/druid/adapters/http/handlers/scroll_handler.go @@ -181,8 +181,11 @@ func (h *ScrollHandler) UpdateScroll(c *fiber.Ctx, id string) error { return fiber.NewError(fiber.StatusBadRequest, err.Error()) } } - runtimeScroll, err := h.supervisor.Update(id, request.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()) } diff --git a/apps/druid/core/services/runtime_release.go b/apps/druid/core/services/runtime_release.go index e7d90f9..e6f63b4 100644 --- a/apps/druid/core/services/runtime_release.go +++ b/apps/druid/core/services/runtime_release.go @@ -17,7 +17,11 @@ func (s *RuntimeSupervisor) InstalledRelease(ctx context.Context, id string) (ma if err != nil { return nil, err } - installed, err := s.runPullWorker(ctx, s.runtimeBackend, ports.RuntimeWorkerModeInspect, id, runtime.Artifact, runtime.Root, nil, "") + 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 } diff --git a/apps/druid/core/services/runtime_supervisor_test.go b/apps/druid/core/services/runtime_supervisor_test.go index 5672fa2..16b7b06 100644 --- a/apps/druid/core/services/runtime_supervisor_test.go +++ b/apps/druid/core/services/runtime_supervisor_test.go @@ -1032,7 +1032,7 @@ func TestRuntimeSupervisorUpdateUsesPullWorkerWhenAvailable(t *testing.T) { supervisor.SetWorkerCallbacks(callbacks, "http://druid-cli:8083") accepted := "registry.local/lab@sha256:" + strings.Repeat("a", 64) - updated, err := supervisor.Update("update-worker", accepted, nil) + updated, err := supervisor.Update("update-worker", accepted, nil, nil) if err != nil { t.Fatal(err) } @@ -1065,7 +1065,7 @@ func TestRuntimeSupervisorUpdateAppliesAcceptedDigestAndRestartsRunningScroll(t supervisor.SetWorkerCallbacks(callbacks, "http://druid-cli:8083") accepted := "registry.local/lab@sha256:" + strings.Repeat("b", 64) - updated, err := supervisor.Update("refresh-worker", accepted, nil) + updated, err := supervisor.Update("refresh-worker", accepted, nil, nil) if err != nil { t.Fatal(err) } diff --git a/apps/druid/core/services/runtime_update.go b/apps/druid/core/services/runtime_update.go index 0261325..6c237e3 100644 --- a/apps/druid/core/services/runtime_update.go +++ b/apps/druid/core/services/runtime_update.go @@ -14,9 +14,10 @@ import ( ) 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) (*domain.RuntimeScroll, error) { +func (s *RuntimeSupervisor) Update(id string, artifact string, registryCredentials []domain.RegistryCredential, expectedDigest *string) (*domain.RuntimeScroll, error) { if !acceptedUpdateReference.MatchString(artifact) { return nil, ErrUnacceptedUpdate } @@ -26,6 +27,18 @@ func (s *RuntimeSupervisor) Update(id string, artifact string, registryCredentia if err != nil { return nil, err } + 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 := 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 index 470f908..90cbc87 100644 --- a/apps/druid/core/services/runtime_update_test.go +++ b/apps/druid/core/services/runtime_update_test.go @@ -1,10 +1,12 @@ 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" ) @@ -18,7 +20,7 @@ func TestRuntimeUpdateRejectsUnacceptedReferenceBeforeStopping(t *testing.T) { } backend := &fakeWorkerBackend{} supervisor := newRuntimeSupervisorForTest(t, store, coreservices.NewRuntimeScrollManager(store), backend) - _, err := supervisor.Update(original.ID, artifact, nil) + _, 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) } @@ -28,3 +30,23 @@ func TestRuntimeUpdateRejectsUnacceptedReferenceBeforeStopping(t *testing.T) { }) } } + +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 index 5849006..9352a69 100644 --- a/docs/runtime-lifecycle-verification.md +++ b/docs/runtime-lifecycle-verification.md @@ -7,6 +7,7 @@ Verified locally on 2026-09-28. This is runtime acceptance, not complete product - 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. - 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. @@ -20,6 +21,7 @@ Verified locally on 2026-09-28. This is runtime acceptance, not complete product - 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. @@ -27,6 +29,6 @@ Verified locally on 2026-09-28. This is runtime acceptance, not complete product ## Remaining product work -- Server/UI check-and-apply wiring, shared Team publication, and shared-stack browser acceptance remain outside this runtime milestone. +- 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 858d758..7d08a31 100644 --- a/internal/api/generated.go +++ b/internal/api/generated.go @@ -230,8 +230,11 @@ type RuntimeUIPackages map[string]RuntimeUIPackage // UpdateScrollRequest defines model for UpdateScrollRequest. type UpdateScrollRequest struct { // Artifact Explicitly accepted immutable target. Missing or mutable tag references are rejected before stopping the runtime. - Artifact string `json:"artifact"` - RegistryCredentials *[]RegistryCredential `json:"registry_credentials,omitempty"` + 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. @@ -3422,57 +3425,59 @@ func RegisterHandlersWithOptions(router fiber.Router, si ServerInterface, option // Base64 encoded, gzipped, json marshaled Swagger object var swaggerSpec = []string{ - "H4sIAAAAAAAC/+xbXXPbttL+Kxi8vXtpyWmTzBydm7pJP9wmEx87mVwkrgYCVhIqEGAA0LaOR//9DD5I", - "kSIoWbLT2p3etLEALHafXewuFstbTFVeKAnSGjy6xYbOISf+nydFIZbnqrRczs7hSwnGup8LrQrQloOf", - "RIzhM5lXy7mF3P/jGw1TPML/N1yTH0baw/NSWp6DIw0n9Xq8yrBdFoBHmGhNlni1yrCGLyXXwPDoU2ur", - "y3qumvwB1C9+pYFYuKBaCdHPr7Z8SqgfYWCo5oXlSuIRfvfqFFWjSMMUNEgKSGkkFCUCGU8YFcTOcYbh", - "huSFCMyGNWbAdMnZYDYbWjDW/2fk/oNrXo3VXM4cr5x1GXgNhQZKLDBEBCcGTZVGkuQwQO/8HMeEJRMB", - "SAcEEWfDMOF0ilTOrQWWITsHxAjkSqIZSNDEgkFEIs4GLcb/UBOT4s1RTMDzMCz8O4xxUwiy9NIhY7kQ", - "iKocDJpqlUekB0uSi7tzbApCE2z/Vk5AS3D717M8sBX/GowqNQUzQKczqTQwNFkiqeRRY+mE0AVIZgap", - "3dW1BD1OaTQaOvIzEGeoNMD87rQ0VuWgj6aEcjlD2p0FREo7V5r/l7j1yb00zLixejmmGhhIy4nY49zF", - "xa/qtbvPXHVcUgfuNQiwwMKJ6x61gEhHBGOJLf2EtWJZoNSVeIMd7qZEAimOfpSm1Pu4gA531eCY8Vlc", - "3XN4e8/NP+b5OMzzFyDCzs/BFEoa6NpBrlhCI69KrUFaNPerUTA25Oc2XZFapOQvtJppMKZL9iyOoAI0", - "BWnJLOhZKMIcwo4xj6tzcFOlc2LxCE+FIi5+5OSG52WOR8+OjzOccxn+Oq5ZkGU+AR2Pl7ZjRmxCto9z", - "kE3f7Of6Y1fv6BYeOavAGZalEM7X45HVJew6mx6ilB7eKLq4qA99Wwdww+2YRkX07MelhVkQThBjx0El", - "YzoncubX1cxzaV8+x6mFDacjHXKfsC6ldGJkmCnpdau10jjD14S7jKchSo/AkWaSqxQOZ0onvFFLQ/t4", - "lSKSq23j5YsX371oWMezFBCFVlZRJZpQzK0tcOb/58MrdX+VrNgNgWcustKgnZReKwrMeWcP1FtSeF/M", - "GA95xVnbR/f8vs1/NOxslWCgy1E5EdzMP5yeEbogM+gNGD7l25IQ+XBzpJWyRxoEsfwK0OCamNwniwP0", - "GqakFNYgq1Ch+RWxMGTc2CEpijBPaVQ4bmj790EyIHYESTjOjgxz1RPMCmLMtdLpkFYa0D0GuGEJnn5j", - "QYNwyhpi6DmJ/vtd5f0OC9p9YYcF4PFoSoSB7OHCkNvSO88GOxOlBBC5X4yKODjX0OciJ6qULLVP5kEf", - "8yKJiR+rfETXDywAihPBr+C9JtMpp0ka3rERavkVt8sxsS1n24wU+3utpGcKDiK9ruG3uvq/GU+WNsB1", - "l2DgM6okJdtBowF3HNxrr2qNWmynec0lU9dpnvaRrsdB19h2nXUWLWwtfI3QFovdvLwnIrt1rkBss8/9", - "A964f3SbgQTnuuU4lFqkfdw2+bmcvSd6Bgnp73YX2Od07BL+0LNjQAC1Sm+Lul2T3ETFgL7iFMZ3CxY9", - "VjnuSSc2yG+xyr6r6NboQX3diO3l33qugN5hhkQyNdy8ivXrcGeESqRSISKBvgLmrfzuty6flfrlhL2T", - "YrmRfK8DnlI9wTechAev/mU4JFb9Vt9N6qMqcdZI741VReF/qzL8qtpwmVBsycdFSAfvKkidP/q0syzY", - "nsaUKnHU9hpxb2ORra8eDdtt7b3ljNT89ie6XVT2lmqLS21K6yZlONZU9+T/4ItCB4iUR/vgmTm4mvzj", - "TSE45VYsEaEUCgsM8TwvQ/HU+rAxQG+5Mf72r9F6aLauPhtENCANjilgaAJTpQF5i3bL3FU+FooGAUYX", - "c/EI//7p9+8/fzaX//+9mZNvX7wcfSJH0+Ojf13evny++uaRV258OKGl5nZ54TaIaTAQDfqkDBYa/vqp", - "MsRfP77357qpgV8/vkdWLUCGojL3nNklKrS64gy0P1SOvEvIPLk1Lv5G7ERw66s92+Qv5krbI5dBM/Sl", - "BL2sNlMafYTJhaILsIgqKYFWdR3uFvrJuMpzwhbrnUnBfwMHlwsycqrcxlRJG4xstSnka11yhl69OUWC", - "lJLOfZmdoZxId0CQX8kl6CNfImTVIwYpnHWGelOGBF/AZznztXgXQrTJECOWTIgBk3mC1zCpxgafPbvc", - "+jpYzQDOsBsNbB0Png2OfcQrQJKC4xH+zv8UzrpX6JAUfHj1bBgKbe6XmEm1JfwZbFWuapXksCcebo2n", - "LEwMFT+vLx8PfeHPb/bt8XGFZMxWGxAM/zCh+BLseZe1b9QVvaraPIcZPoa9OP7uT9z4IuRJqJTkivBQ", - "TPPnqcxzopcRzk0cLZkZf4cPmshwwBtfuqWVmoLlmIae2vC/4cZexDn3BH+fNCJmfF1308GmqnZXgrRx", - "cezXRXdTy1FBE0ea2Lg81SSAaD5D4uD0wNgfFFs+mCGkXjpXbQ/rkrhVRw/PHoyFDfh3wY2qzKyNehBk", - "A/ftsHdNcgj+1cdH56RGmq9CX0kjqYenO2nk+C/TSEBtUyNBkA2NILjhxobQomLZUyzD84HZW123nK2C", - "n3dpeFdd4VmxVldBNMnBgnZb3IYYGhPGGEJ9ytwGOmuAtpmBXn5FJbSfRHcrobqKrDL8/Ph5/xNdnC6V", - "RVNfrWlrLWy71znK0m78Z7BPE/k9zf++iLs4ek+35c7B0OVlZdHvu37w419ZJQ/vD3fV9x+bbwwwI7c0", - "Hsi2V7wBWjYOWNTaQRqnKs+JZGZ4G/+16tf+eSkDz6/C1K9hAVmSCK033IvSxnsz4dZfiPxFNageGIq0", - "kVU14NXVNlgA43I26LkvmaWkuMnF5iNP5znmL/U6oYzANn3FPb3PeSk3Q/RaYQfZpJzyWW9uXweFV2He", - "IwwNmzWEjiLOiDbrC3AUuOvTi9S0QzE1SoC5E6ph5iPENV1fa5WI+zGvBEMLWIbGpXXJvwt9fEc3VBXe", - "SdSgHAB+XQPfjvyZCgnso4N9nytw4/F472swckBVlYAHT4la1A/S45cSStitx//4aU8sY00943T15UXz", - "GMIWvL80Zq2B/hJhqQDvx1mDAGK2In0qjSVCADuPcx8h3ButVf0tnBoKZbhVern7jaIxN6toXt7B+9WA", - "IUqkkpwSgda0/GVaEGORKSkFY6alG/bQorjLfbMEIMynXrzmpN4g0lM6VMcbLwnoSokyb5pRVQ7bZj7G", - "qm0lmPMw4Z97zNe+5Aac73yRqRR3kHNuvPymte6/IIlFsTj36ag+9fnLY1N33+2ipfMz0IYbG1uMlT4K", - "H9IAiz2HSNe66RqB78XYaQLD8KB5h4yr1Tfz5FOvdhfQXbKvsABVeCWy4PAMHFv6K93UCw7QUd2omD6k", - "F274b1rzuwgN7tvPh5/0IMU8Y1WxDWhV/G1x9j02u3BWxeYF4VrphVCEGXQ95wJQEbqYnMUzYslhaij5", - "sNnDs90hNdpJnqZWmv1HibJHaHIHhj6cogoVdxP31+xUASSxAH04f2PurYvhrd9zNYxb9J+UyPSGgv68", - "4mfAZhudquEstvPjqsU09cnEV8pP+j5fWMUk5ZE8M15zO0exs6tpUjlY4o/4RrISpEJEIiI0ELY8mpRc", - "WBQaTT6coo8nF28rKgfaZFF9HpU2v2br1xNKWFMda08lYXUX3ePuRfdEIkg00l38cvLti5epr7W5QbW0", - "97w9B1Y3A9aUi9hjtVlsSRlg7IyuTKevd+vk7BTHhkk8xM4+ItFO01lgIrR35SCt56R6TwF/uXMzG54s", - "4tv5zDGuCYWjNcH10lA8SjzqpJrbUHjPWYDvc4sUrmFi/MwElTOlLeIytJJyJWtIywaBIvQm3ybbqxCd", - "A12Y5MLYydRd+rYUlh9FXVaqTUlfabNL4nVoRhN8CnRJRXp5NIHu6p9clnNNLJ27HMfxzuAKhCq8NuNn", - "phV+blqCxomUygbUnDn6c2Ea0pN63ODV5ep/AQAA//8lqbo540EAAA==", + "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/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 a8d1d0e..31fbaba 100644 --- a/test/integration/kubernetes/kubernetes_cli_test.go +++ b/test/integration/kubernetes/kubernetes_cli_test.go @@ -164,7 +164,12 @@ func TestKubernetesBackendCLIComplexLifecycle(t *testing.T) { 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 - e2e.RunClient(t, bins, socket, "update", created.ID, acceptedArtifact) + 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) } From 1e6f057f677adb348b8af8f6d098f41ea4d878c6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Marc=20Schottst=C3=A4dt?= Date: Mon, 28 Sep 2026 02:29:45 +0200 Subject: [PATCH 4/4] feat(runtime): expose owner-checked updates Use the authenticated public listener for cross-cluster accepted updates. Reject missing and cross-owner tokens before reading or changing runtime state; keep operator-only management authorization unchanged. --- apps/druid/adapters/http/handlers/routes.go | 1 + .../adapters/http/handlers/routes_test.go | 64 +++++++++++++++++++ docs/runtime-lifecycle-verification.md | 1 + 3 files changed, 66 insertions(+) diff --git a/apps/druid/adapters/http/handlers/routes.go b/apps/druid/adapters/http/handlers/routes.go index e715a6b..adb9fc9 100644 --- a/apps/druid/adapters/http/handlers/routes.go +++ b/apps/druid/adapters/http/handlers/routes.go @@ -64,6 +64,7 @@ func RegisterPublicRoutes(app *fiber.App, handlers RouteHandlers) { 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/docs/runtime-lifecycle-verification.md b/docs/runtime-lifecycle-verification.md index 9352a69..cbfdcfb 100644 --- a/docs/runtime-lifecycle-verification.md +++ b/docs/runtime-lifecycle-verification.md @@ -8,6 +8,7 @@ Verified locally on 2026-09-28. This is runtime acceptance, not complete product - 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.