diff --git a/docs/deploy.md b/docs/deploy.md index 61b24e8..95217d0 100644 --- a/docs/deploy.md +++ b/docs/deploy.md @@ -82,8 +82,8 @@ sandboxd reads one JSON file (`-config`, default | `preview_listen` | (off) | address for a preview HTTP server that serves guest ports under signed URLs; needs `preview_secret` | | `preview_secret` | — | cluster-shared HMAC secret signing preview tokens (all nodes share one) | | `preview_advertise` | = `preview_listen` | the base URL a browser/proxy reaches this node's preview server at | -| `checkpoint_dir` | `/checkpoints` | where checkpoints and promoted templates live. Point it at a shared FUSE mount (JuiceFS over object storage, NFS) and every node sharing the mount can branch every checkpoint — the path's filesystem is the operator's choice. One contract on any shared root (mount or bucket): a template key has a single writer — promotes go to the sandbox's owner node, and operators must not race promotes of one name from different nodes (checkpoint ids are node-generated and never collide). A checkpoint deleted on one node while another is mid-branch from it fails that branch visibly | -| `checkpoint_store` | dir | checkpoint AND promoted-template backend (both live in one store root, id-namespaced ck_/tp_): `{"kind": "s3", "s3": {"bucket": "…", "prefix": "ck/", "endpoint": "…", "region": "…", "force_path_style": true}}` stores checkpoints in object storage (any node claims any checkpoint, no shared mount needed). Credentials come from the standard AWS chain (env/IAM role), never this file. A crash between upload and the meta.json commit marker leaves orphan objects invisible to listings — add an S3 lifecycle rule to reclaim them. Absent = the dir backend at `checkpoint_dir` | +| `checkpoint_dir` | `/checkpoints` | where checkpoints and promoted templates live. Point it at a shared FUSE mount (JuiceFS over object storage, NFS) and every node sharing the mount can branch every checkpoint — the filesystem must provide working cross-node POSIX advisory `flock`; local-only or ignored locks are unsupported. A fixed, bounded hash-striped lock set keeps a fetched directory generation read-locked through clone, so replace/delete waits for that reader. One contract on any shared root (mount or bucket): a template key has a single writer — promotes go to the sandbox's owner node, and operators must not race promotes of one name from different nodes (checkpoint ids are node-generated and never collide) | +| `checkpoint_store` | dir | checkpoint AND promoted-template backend (both live in one store root, id-namespaced ck_/tp_): `{"kind": "s3", "s3": {"bucket": "…", "prefix": "ck/", "endpoint": "…", "region": "…", "force_path_style": true}}` stores checkpoints in object storage (any node claims any checkpoint, no shared mount needed). Credentials come from the standard AWS chain (env/IAM role), never this file. Re-publish retains prior export generations until Delete so an in-flight fetch that selected old metadata can finish; budget storage for those generations. An explicit S3 Delete can still make a concurrent fetch that has not finished materializing fail visibly. A crash between upload and the meta.json commit marker leaves orphan objects invisible to listings — add an S3 lifecycle rule to reclaim them. Absent = the dir backend at `checkpoint_dir` | | `checkpoint_ttl_hours` | 0 (keep forever) | ages out checkpoints older than this; the sweep runs hourly and at startup. Explicit deletes never wait for it. Must be nonzero and match fleet-wide when `checkpoint_peer_heal` is on — it is the expiry eligibility point for a healed replica a delete broadcast missed, after which its next successful hourly sweep removes it; persistent sweep failure extends retention until one succeeds, so it is not a hard ceiling | | `checkpoint_peer_heal` | false | on a cluster, lets a node pull a checkpoint it lacks from a peer — found via a live probe, not gossip — rather than failing the branch; see [placement lifecycle](cluster.md#checkpoints-on-a-cluster). Three requirements, all enforced at config load: a nonempty `api_token` (the blob transfer between peers authenticates with it; without one the raw record stream would be open), `mesh.cluster_key` set (the pull presents the fleet `api_token` to an address learned from the peer probe, so the gossip layer carrying that address must itself be authenticated), and `checkpoint_ttl_hours` nonzero (a replica a delete broadcast missed becomes eligible for expiry after it, and its next successful hourly sweep removes it — so it is the finite eligibility point, not an exact ceiling). A shared checkpoint store (`checkpoint_store` kind `s3`) ignores this setting — every node already resolves every checkpoint directly, so there is nothing to heal | | `warm_max` (pool entry) | 0 (static) | turns on the demand-adaptive watermark for that pool: the warm target rises from `warm` toward `warm_max` while claims arrive faster than the measured provision lead covers, and decays back over ~a minute of silence | diff --git a/docs/sandboxd-api.md b/docs/sandboxd-api.md index 7e178ce..598586b 100644 --- a/docs/sandboxd-api.md +++ b/docs/sandboxd-api.md @@ -39,9 +39,14 @@ Success: ```json {"id": "sb_…", "token": "…", "deadline": "2026-07-06T00:05:00Z", - "owner_addr": "10.0.0.5:7777"} + "owner_addr": "10.0.0.5:7777", "template_digest": "sha256:…"} ``` +A claim cloned from a promoted template carries `template_digest`, the exact +export generation fetched for that clone. It is absent for configured pools, +cold image boots, forks, checkpoints, and templates published by an older +sandboxd until they are re-promoted. + A claim branched from a checkpoint (fork children included) additionally carries `"from_checkpoint": "ck_…"` — the lineage edge for reconstructing the checkpoint tree. @@ -136,15 +141,24 @@ Publishes the sandbox's current state as a node-local template under from it, provision-on-demand — no warm pool unless the node config adds one. Re-promoting to the same name replaces the template. A hibernated sandbox is promoted from its memory image without waking. 200 returns the template's -full key. On the default local-disk backend a template is node-local, so a -cluster client claims from and deletes on this node (name-based calls route -via gossip); a shared checkpoint store makes every node resolve it. Under -exactly this key: +full key and immutable content identity. On the default local-disk backend a +template is node-local, so a cluster client claims from and deletes on this +node (name-based calls route via gossip); a shared checkpoint store makes +every node resolve it. Under exactly this key: ```json -{"key": {"template": "myproj:v1", "net": "none", "size": "small"}} +{"key": {"template": "myproj:v1", "net": "none", "size": "small"}, + "content_digest": "sha256:…"} ``` +`content_digest` is SHA-256 over a versioned canonical stream of the published +export's regular files: slash-relative path, byte length, and bytes, ordered +lexically. Directory entries, modes, mtimes, and the template's ownership/ +creation metadata do not affect it. The digest is computed once while +promoting, stored in `meta.json`, and therefore has identical semantics on the +directory and S3 backends. Re-promoting unchanged export bytes keeps the +digest; changing any exported path or bytes changes it. + 400 invalid name, 401 bad api token, 409 when the name collides with a configured pool, the template is owned by another tenant, or the sandbox is on the egress lane (see [egress](egress.md)), 404 unknown id or wrong diff --git a/docs/sdk-python.md b/docs/sdk-python.md index 2afa87c..8d9e60b 100644 --- a/docs/sdk-python.md +++ b/docs/sdk-python.md @@ -98,9 +98,10 @@ sb = client.new("ghcr.io/cocoonstack/sandbox/rt:24.04", `new` returns when the sandbox's silkd answers: a warm hit is milliseconds, a cold key can take the full boot. The handle exposes `sb.id`, `sb.token`, `sb.owner`, `sb.deadline`, and `sb.from_checkpoint` (the lineage edge when -branched). `Sandbox` is a context manager; `sb.close()` releases it -(releasing one already gone is not an error — double-release and reap races -stay silent). +branched). `sb.template_digest` is the exact content identity when the claim +cloned a promoted template; it is empty for other sources. `Sandbox` is a +context manager; `sb.close()` releases it (releasing one already gone is not +an error — double-release and reap races stay silent). ## Hibernating @@ -138,6 +139,7 @@ All-or-nothing: on error no child survived. Count is capped at the node's ```python tpl = sb.promote("myproj:v1") # publish current state child = tpl.new() # clones the promoted state +assert tpl.content_digest and tpl.content_digest == child.template_digest tpl.delete() # caller owns the lifecycle ``` @@ -149,6 +151,11 @@ bound there, so its `new`/`delete` always reach it. The name-based calls cluster-wide via template gossip and lag a promote/delete by about a gossip tick — prefer the handle right after promoting (see [Templates on a cluster](cluster.md#templates-on-a-cluster)). +`tpl.content_digest` identifies the published export bytes. A caller pinning +the mutable name can compare a claim's `template_digest` with its expected +value and close/refuse a mismatch. +Templates published by an older sandboxd have empty digests until they are +re-promoted after the node is upgraded. ## Checkpoints — branching and time travel diff --git a/docs/sdk.md b/docs/sdk.md index b14b62f..0795c04 100644 --- a/docs/sdk.md +++ b/docs/sdk.md @@ -143,11 +143,12 @@ defer sb.Close() `New` returns when the sandbox's silkd answers: a warm hit is milliseconds, a cold key can take the full boot. `Sandbox.ID`, `Sandbox.Deadline`, and -`Sandbox.FromCheckpoint` (the lineage edge when branched) are exported, -`Owner()` names the owning node, and `Token()` returns the per-sandbox -bearer to persist with `ID` for a later `Lookup`; `Close()` releases the -sandbox (releasing one already gone is not an error, and `Close` is bounded -internally so it stays defer-friendly). +`Sandbox.FromCheckpoint` (the lineage edge when branched) are exported. +`Sandbox.TemplateDigest` is the exact content identity when the claim cloned a +promoted template; it is empty for other sources. `Owner()` names the owning +node, and `Token()` returns the per-sandbox bearer to persist with `ID` for a +later `Lookup`; `Close()` releases the sandbox (releasing one already gone is +not an error, and `Close` is bounded internally so it stays defer-friendly). ## Hibernating @@ -190,6 +191,7 @@ client needs `WithAPIToken` — a sandbox handle alone cannot amplify. ```go tpl, err := sb.Promote(ctx, "myproj:v1") // publish current state child, err := tpl.New(ctx) // clones the promoted state +fmt.Println(tpl.ContentDigest != "" && tpl.ContentDigest == child.TemplateDigest) // true err = tpl.Delete(ctx) // caller owns the lifecycle ``` @@ -197,7 +199,12 @@ err = tpl.Delete(ctx) // caller owns the lifecycle keyed by (name, the sandbox's network lane, its size). Claims clone on demand (~a golden-clone's latency); there is no warm pool for promoted templates unless the node's config adds one. Re-promoting to the same name -replaces the template. +replaces the template. `Template.ContentDigest` identifies the published +export bytes; a claim from that exact generation carries the same value in +`Sandbox.TemplateDigest`. A caller pinning a mutable template name can compare +the claim's value with its expected digest and close/refuse a mismatch. +Templates published by an older sandboxd have empty digests until they are +re-promoted after the node is upgraded. **On the default local-disk backend templates live on one node**, and on a cluster the parent claim may have been redirected — the returned `Template` diff --git a/e2e/e2e_test.go b/e2e/e2e_test.go index d2af823..c0de464 100644 --- a/e2e/e2e_test.go +++ b/e2e/e2e_test.go @@ -121,12 +121,18 @@ func TestPromoteEndToEnd(t *testing.T) { if err != nil { t.Fatalf("Promote: %v", err) } + if tpl.ContentDigest == "" { + t.Fatal("Promote returned an empty content digest") + } // Both claim surfaces must work: the owner-bound handle and the // name-based Client call on the (single) node that holds the template. child, err := tpl.New(t.Context()) if err != nil { t.Fatalf("claim via template handle: %v", err) } + if child.TemplateDigest != tpl.ContentDigest { + t.Errorf("claim template digest %q, want %q", child.TemplateDigest, tpl.ContentDigest) + } if out, err := child.Exec(t.Context(), "echo", "tpl"); err != nil || out != "tpl\n" { t.Errorf("exec on promoted claim: %q, %v", out, err) } @@ -135,6 +141,9 @@ func TestPromoteEndToEnd(t *testing.T) { if nameErr != nil { t.Fatalf("claim promoted template by name: %v", nameErr) } + if byName.TemplateDigest != tpl.ContentDigest { + t.Errorf("name claim template digest %q, want %q", byName.TemplateDigest, tpl.ContentDigest) + } _ = byName.Close() if err := tpl.Delete(t.Context()); err != nil { diff --git a/sandboxd/pool/archive.go b/sandboxd/pool/archive.go index 090cc55..29e5bc4 100644 --- a/sandboxd/pool/archive.go +++ b/sandboxd/pool/archive.go @@ -4,6 +4,9 @@ import ( "context" "errors" "fmt" + "io/fs" + "os" + "path/filepath" "time" "github.com/projecteru2/core/log" @@ -12,6 +15,8 @@ import ( "github.com/cocoonstack/sandbox/sandboxd/types" ) +const archiveDeleteDir = "archive-deletes" + // archive*For resolve a key's thresholds (pool's, else node default); callers hold m.mu. func (m *Manager) archiveAfterFor(key types.PoolKey) time.Duration { if p, ok := m.activePool(key); ok { @@ -92,6 +97,9 @@ func (m *Manager) archive(ctx context.Context, sb *types.Sandbox) error { m.pendingCks[ckID] = struct{}{} m.mu.Unlock() defer m.untrack(m.pendingCks, ckID) + if err := m.markArchiveCk(ckID); err != nil { + return fmt.Errorf("track archive checkpoint: %w", err) + } // A hibernated source is copied from its wake image, so no VM starts. ck, srcSnap, err := m.publishCheckpoint(ctx, sb, ckID, "archive", sb.Tenant, true) if err != nil { @@ -174,8 +182,8 @@ func (m *Manager) wakeArchived(ctx context.Context, sb *types.Sandbox) (string, if err != nil { return "", fmt.Errorf("wake %s: fetch archive: %w", sb.ID, err) } - defer release() built, err := m.provision(ctx, sb.Key, dir) + release() if err != nil { return "", fmt.Errorf("wake %s: %w", sb.ID, err) } @@ -192,6 +200,9 @@ func (m *Manager) wakeArchived(ctx context.Context, sb *types.Sandbox) (string, log.WithFunc("pool.wakeArchived").Warnf(ctx, "delete consumed archive ck %s: %v", ck, delErr) } else { consumed = true + if clearErr := m.clearArchiveCk(ck); clearErr != nil { + log.WithFunc("pool.wakeArchived").Warnf(ctx, "clear archive ck %s: %v", ck, clearErr) + } } // Only the none lane reaches here (the egress guard above fails closed), so // this rebinds the none-lane proxy; there is no NIC to re-lock. @@ -236,7 +247,105 @@ func (m *Manager) commitWake(ctx context.Context, sb *types.Sandbox, vmName, soc // deleteOrphanArchiveCk drops the published ck when archive() aborts pre-commit. func (m *Manager) deleteOrphanArchiveCk(ctx context.Context, ckID string) { - if err := m.ckpts.Delete(ctx, ckID); err != nil { + if err := m.deleteArchiveCk(ctx, ckID); err != nil { log.WithFunc("pool.deleteOrphanArchiveCk").Warnf(ctx, "delete orphaned archive ck %s: %v", ckID, err) } } + +func (m *Manager) markArchiveCk(ckID string) error { + if !store.CheckpointIDRe.MatchString(ckID) { + return fmt.Errorf("invalid checkpoint id %q", ckID) + } + dir := filepath.Join(m.dataDir, archiveDeleteDir) + if err := os.MkdirAll(dir, 0o750); err != nil { + return err + } + return os.WriteFile(filepath.Join(dir, ckID), nil, 0o600) +} + +func (m *Manager) clearArchiveCk(ckID string) error { + err := os.Remove(filepath.Join(m.dataDir, archiveDeleteDir, ckID)) + if errors.Is(err, fs.ErrNotExist) { + return nil + } + return err +} + +func (m *Manager) deleteArchiveCk(ctx context.Context, ckID string) error { + if err := m.markArchiveCk(ckID); err != nil { + return fmt.Errorf("track: %w", err) + } + if err := m.deleteCkLocked(ctx, ckID); err != nil { + return err + } + if err := m.clearArchiveCk(ckID); err != nil { + return fmt.Errorf("clear: %w", err) + } + return nil +} + +func (m *Manager) retryArchiveDeletes(ctx context.Context) { + if !m.archiveDeleteSweep.CompareAndSwap(false, true) { + return + } + defer m.archiveDeleteSweep.Store(false) + entries, err := os.ReadDir(filepath.Join(m.dataDir, archiveDeleteDir)) + if errors.Is(err, fs.ErrNotExist) { + return + } + if err != nil { + log.WithFunc("pool.retryArchiveDeletes").Warnf(ctx, "list: %v", err) + return + } + if len(entries) == 0 { + return + } + pinned := m.pinnedArchiveCks() + ids := make([]string, 0, len(entries)) + for _, entry := range entries { + if !entry.Type().IsRegular() || !store.CheckpointIDRe.MatchString(entry.Name()) { + continue + } + if _, ok := pinned[entry.Name()]; !ok { + ids = append(ids, entry.Name()) + } + } + m.runBounded(ctx, len(ids), func(ctx context.Context, i int) { + m.retryArchiveDelete(ctx, ids[i]) + }).Wait() +} + +func (m *Manager) retryArchiveDelete(ctx context.Context, ckID string) { + l := m.recLock(ckID) + l.Lock() + if m.archiveCkPinned(ckID) { + l.Unlock() + m.recDone(ckID) + return + } + err := m.ckpts.Delete(ctx, ckID) + l.Unlock() + if err != nil { + m.recDone(ckID) + log.WithFunc("pool.retryArchiveDeletes").Warnf(ctx, "delete %s: %v", ckID, err) + return + } + m.recDoneEvict(ckID) + if err := m.clearArchiveCk(ckID); err != nil { + log.WithFunc("pool.retryArchiveDeletes").Warnf(ctx, "clear %s: %v", ckID, err) + } +} + +func (m *Manager) archiveCkPinned(ckID string) bool { + m.mu.Lock() + defer m.mu.Unlock() + if _, ok := m.pendingCks[ckID]; ok { + return true + } + for _, sb := range m.claimed { + if sb.ArchiveCk == ckID { + return true + } + } + return false +} diff --git a/sandboxd/pool/archive_test.go b/sandboxd/pool/archive_test.go index 6744a8f..f98823c 100644 --- a/sandboxd/pool/archive_test.go +++ b/sandboxd/pool/archive_test.go @@ -14,6 +14,11 @@ import ( "github.com/cocoonstack/sandbox/sandboxd/types" ) +const ( + testArchiveRelease = "release" + testArchiveReap = "reap" +) + // TestArchivePublishWindowPinsCheckpoint guards the pre-pin: between an // archive's checkpoint publish and the ArchiveCk commit, the checkpoint must // be invisible to listings and refuse deletion — a delete landing in that @@ -31,6 +36,10 @@ func TestArchivePublishWindowPinsCheckpoint(t *testing.T) { archived := make(chan error, 1) go func() { archived <- m.archive(t.Context(), sb) }() ck := <-stall.published + m.retryArchiveDeletes(t.Context()) + if !ckExists(t, m, ck) { + t.Fatal("archive delete retry removed a checkpoint still being published") + } if ckpts, err := m.Checkpoints(t.Context(), ""); err != nil || len(ckpts) != 0 { t.Errorf("mid-publish checkpoint visible: %v, %v", ckpts, err) @@ -356,6 +365,204 @@ func TestReleaseArchivedDeletesCheckpoint(t *testing.T) { } } +func TestArchiveDeleteRetryAfterRestart(t *testing.T) { + for _, action := range []string{testArchiveRelease, testArchiveReap} { + t.Run(action, func(t *testing.T) { + eng := newFakeEngine() + dataDir := t.TempDir() + m := newTestManagerAt(t, eng, dataDir, archivePool(3600)) + sb := mustClaim(t, m, testKey) + id, token := sb.ID, sb.Token + mustArchive(t, m, sb) + ck := sb.ArchiveCk + failed := &failingDeleteStore{Store: m.ckpts, attempts: make(chan string, 1)} + m.ckpts = failed + + switch action { + case testArchiveRelease: + if err := m.Release(t.Context(), id, Cred{Token: token}); err != nil { + t.Fatalf("release: %v", err) + } + case testArchiveReap: + m.mu.Lock() + sb.Deadline = time.Now().Add(-time.Second) + m.mu.Unlock() + m.reapOnce(t.Context()) + select { + case <-failed.attempts: + case <-time.After(3 * time.Second): + t.Fatal("reap did not attempt archive deletion") + } + waitFor(t, func() bool { + m.mu.Lock() + defer m.mu.Unlock() + _, pending := m.pendingCks[ck] + return !pending + }) + } + if _, ok := m.claim(id, token); ok { + t.Fatal("removed archive claim survived") + } + if !ckExists(t, m, ck) || !archiveCkMarked(m, ck) { + t.Fatal("failed deletion did not retain the checkpoint and cleanup marker") + } + + m2 := newTestManagerAt(t, eng, dataDir, archivePool(3600)) + if err := m2.Reconcile(t.Context()); err != nil { + t.Fatalf("Reconcile: %v", err) + } + if ckExists(t, m2, ck) || archiveCkMarked(m2, ck) { + t.Fatal("restart did not finish archive deletion") + } + if _, ok := m2.claim(id, token); ok { + t.Fatal("restart revived the removed claim") + } + }) + } +} + +func TestReconcileAdoptsLegacyArchiveMarker(t *testing.T) { + eng := newFakeEngine() + dataDir := t.TempDir() + m := newTestManagerAt(t, eng, dataDir, archivePool(3600)) + sb := mustClaim(t, m, testKey) + mustArchive(t, m, sb) + id, token, ck := sb.ID, sb.Token, sb.ArchiveCk + if err := m.clearArchiveCk(ck); err != nil { + t.Fatalf("clear marker: %v", err) + } + + m2 := newTestManagerAt(t, eng, dataDir, archivePool(3600)) + if err := m2.Reconcile(t.Context()); err != nil { + t.Fatalf("Reconcile: %v", err) + } + if !ckExists(t, m2, ck) || !archiveCkMarked(m2, ck) { + t.Fatal("reconcile did not retain and mark the live archive") + } + if _, err := m2.WakeAgentSocket(t.Context(), id, token); err != nil { + t.Fatalf("wake: %v", err) + } + if archiveCkMarked(m2, ck) { + t.Fatal("wake left the archive cleanup marker") + } +} + +func TestArchiveRemovalCommitPinsCheckpoint(t *testing.T) { + for _, action := range []string{testArchiveRelease, testArchiveReap} { + t.Run(action, func(t *testing.T) { + eng := newFakeEngine() + m := newTestManager(t, eng, archivePool(3600)) + sb := mustClaim(t, m, testKey) + mustArchive(t, m, sb) + id, token, ck := sb.ID, sb.Token, sb.ArchiveCk + if action == testArchiveReap { + m.mu.Lock() + sb.Deadline = time.Now().Add(-time.Second) + m.mu.Unlock() + } + + m.store.mu.Lock() + locked := true + t.Cleanup(func() { + if locked { + m.store.mu.Unlock() + } + }) + done := make(chan error, 1) + go func() { + if action == testArchiveRelease { + done <- m.Release(t.Context(), id, Cred{Token: token}) + return + } + m.reapOnce(t.Context()) + done <- nil + }() + waitFor(t, func() bool { + m.mu.Lock() + defer m.mu.Unlock() + _, claimed := m.claimed[id] + _, pending := m.pendingCks[ck] + return !claimed && pending + }) + m.retryArchiveDeletes(t.Context()) + if !ckExists(t, m, ck) { + t.Error("retry deleted the checkpoint before the claims commit") + } + m.store.path = filepath.Join(t.TempDir(), "gone", "claims.json") + m.store.mu.Unlock() + locked = false + if err := <-done; action == testArchiveRelease && err == nil { + t.Fatal("release succeeded despite a persist failure") + } + waitFor(t, func() bool { + m.mu.Lock() + defer m.mu.Unlock() + _, claimed := m.claimed[id] + _, pending := m.pendingCks[ck] + return claimed && !pending + }) + m.retryArchiveDeletes(t.Context()) + if !ckExists(t, m, ck) { + t.Error("retry deleted the checkpoint restored by rollback") + } + healStore(t, m) + if _, err := m.WakeAgentSocket(t.Context(), id, token); err != nil { + t.Fatalf("wake after rollback: %v", err) + } + }) + } +} + +func TestArchiveDeleteRetryRechecksWakeRollback(t *testing.T) { + eng := newFakeEngine() + m := newTestManager(t, eng, archivePool(3600)) + sb := mustClaim(t, m, testKey) + mustArchive(t, m, sb) + id, token, ck := sb.ID, sb.Token, sb.ArchiveCk + + m.store.mu.Lock() + locked := true + t.Cleanup(func() { + if locked { + m.store.mu.Unlock() + } + }) + woke := make(chan error, 1) + go func() { + _, err := m.WakeAgentSocket(t.Context(), id, token) + woke <- err + }() + waitFor(t, func() bool { + m.mu.Lock() + defer m.mu.Unlock() + return sb.ArchiveCk == "" + }) + retried := make(chan struct{}) + go func() { + m.retryArchiveDeletes(t.Context()) + close(retried) + }() + waitFor(t, func() bool { + m.recLocksMu.Lock() + defer m.recLocksMu.Unlock() + return m.recRefs[ck] >= 2 + }) + m.store.path = filepath.Join(t.TempDir(), "gone", "claims.json") + m.store.mu.Unlock() + locked = false + if err := <-woke; err == nil { + t.Fatal("wake succeeded despite a persist failure") + } + <-retried + m.mu.Lock() + archived := sb.ArchiveCk == ck && sb.VMName == "" + m.mu.Unlock() + if !archived || !ckExists(t, m, ck) { + t.Fatal("retry deleted the archive restored by wake rollback") + } + healStore(t, m) +} + // TestIdleOnceSkipsArchived guards the fix where the idle sweep tried to // hibernate an archived (VM-less) claim, corrupting its record. func TestIdleOnceSkipsArchived(t *testing.T) { @@ -583,6 +790,19 @@ func (s *stallingStore) Publish(ctx context.Context, staging, id string) error { return nil } +type failingDeleteStore struct { + store.Store + attempts chan string +} + +func (s *failingDeleteStore) Delete(_ context.Context, id string) error { + select { + case s.attempts <- id: + default: + } + return errors.New("delete failed") +} + // breakStore points the claim store at a path whose parent is missing, so the // next save fails — the fault used to exercise the persist-failure paths. func breakStore(t *testing.T, m *Manager) { @@ -640,6 +860,11 @@ func ckExists(t *testing.T, m *Manager, ck string) bool { return true } +func archiveCkMarked(m *Manager, ck string) bool { + _, err := os.Stat(filepath.Join(m.dataDir, archiveDeleteDir, ck)) + return err == nil +} + // pinnedHidden reports whether the store holds exactly one checkpoint and it is // hidden from listings — i.e. a live claim's ArchiveCk still pins it. func pinnedHidden(t *testing.T, m *Manager, ck string) bool { diff --git a/sandboxd/pool/checkpoint.go b/sandboxd/pool/checkpoint.go index 7cb1908..af14978 100644 --- a/sandboxd/pool/checkpoint.go +++ b/sandboxd/pool/checkpoint.go @@ -364,7 +364,7 @@ func (m *Manager) vetoIfHealPending(ckptID string) { } // pinnedArchiveCks is the set of checkpoint ids backing a live archived claim -// or an archive publish in flight (pendingCks): wake images the listing hides +// or an archive transition in flight (pendingCks): wake images the listing hides // and delete/TTL must spare (deleting one would strand its sandbox). func (m *Manager) pinnedArchiveCks() map[string]struct{} { m.mu.Lock() @@ -392,7 +392,7 @@ func (m *Manager) sweepExpiredCheckpoints(ctx context.Context) { logger := log.WithFunc("pool.sweepExpiredCheckpoints") // Checkpoints already hides archive images backing a live claim, so the // TTL never reaches one; their retention is the claim's own Deadline - // (reapPurge). An orphaned archive ck carries no live reference and ages out. + // (reapPurge). ckpts, err := m.Checkpoints(ctx, "") if err != nil { logger.Error(ctx, err, "list for retention") diff --git a/sandboxd/pool/claim.go b/sandboxd/pool/claim.go index a83c497..b640f4d 100644 --- a/sandboxd/pool/claim.go +++ b/sandboxd/pool/claim.go @@ -67,15 +67,16 @@ func (m *Manager) ClaimProvision(ctx context.Context, key types.PoolKey, ttl tim if err := m.overQuota(1, tenant); err != nil { return nil, err } - golden, release, err := m.resolveGolden(ctx, key) + golden, templateDigest, release, err := m.resolveGolden(ctx, key) if err != nil { return nil, fmt.Errorf("resolve template: %w", err) } - defer release() sb, err := m.provision(ctx, key, golden) + release() if err != nil { return nil, err } + sb.TemplateDigest = templateDigest sb.Tenant = tenant sb.ClaimRef = claimRef out, err := m.finalize(ctx, sb, ttl) @@ -161,12 +162,16 @@ func (m *Manager) releaseResolved(ctx context.Context, id string, sb *types.Sand delete(m.claimed, id) m.tenantDelta(sb.Tenant, -1) snap, ck, vmName := sb.HibernateSnap, sb.ArchiveCk, sb.VMName + if ck != "" { + m.pendingCks[ck] = struct{}{} + } js := m.store.snapshot(m.claimed) m.mu.Unlock() if saveErr := m.store.commit(js); saveErr != nil { m.mu.Lock() - m.claimed[id] = sb // roll back so memory matches the still-durable claim; ck stays pinned + m.claimed[id] = sb // roll back so memory matches the still-durable claim; the restored claim re-pins ck m.tenantDelta(sb.Tenant, 1) + delete(m.pendingCks, ck) rb := m.store.snapshot(m.claimed) m.mu.Unlock() m.recommit(ctx, rb) @@ -177,6 +182,7 @@ func (m *Manager) releaseResolved(ctx context.Context, id string, sb *types.Sand ctx = context.WithoutCancel(ctx) if ck != "" { m.purgeArchiveCk(ctx, id, ck, sb.Tenant) // archived: no local VM + m.untrack(m.pendingCks, ck) } var err error if vmName != "" && !m.removeOrRetry(ctx, vmName, id, "") { @@ -312,6 +318,7 @@ func (m *Manager) reapOnce(ctx context.Context) { switch { case sb.ArchiveCk != "": expired = append(expired, victim{action: reapPurge, id: id, ck: sb.ArchiveCk, tenant: sb.Tenant, sb: sb}) + m.pendingCks[sb.ArchiveCk] = struct{}{} delete(m.claimed, id) m.tenantDelta(sb.Tenant, -1) case sb.HibernateSnap != "" && m.archiveEnabledFor(sb.Key): @@ -344,6 +351,7 @@ func (m *Manager) reapOnce(ctx context.Context) { } else { m.claimed[v.id] = v.sb // roll back so memory matches the still-durable claim m.tenantDelta(v.tenant, 1) + delete(m.pendingCks, v.ck) } } rb := m.store.snapshot(m.claimed) @@ -361,6 +369,7 @@ func (m *Manager) reapOnce(ctx context.Context) { case reapPurge: m.disarmEgress(v.id, true) m.purgeArchiveCk(ctx, v.id, v.ck, v.tenant) + m.untrack(m.pendingCks, v.ck) logger.Infof(ctx, "purged archived sandbox %s", v.id) case reapArchive: logSweepResult(ctx, logger, m.archive(ctx, v.sb), "archived expired sandbox "+v.id, "archive expired sandbox "+v.id) @@ -377,7 +386,7 @@ func (m *Manager) reapOnce(ctx context.Context) { // purgeArchiveCk deletes an archived claim's store checkpoint and records the // retention billing event; shared by Release and reapOnce's purge. func (m *Manager) purgeArchiveCk(ctx context.Context, id, ck, tenant string) { - if err := m.deleteCkLocked(ctx, ck); err != nil { + if err := m.deleteArchiveCk(ctx, ck); err != nil { log.WithFunc("pool.purgeArchiveCk").Warnf(ctx, "delete archive ck %s: %v", ck, err) } m.counters.archiveDeletes.Add(1) diff --git a/sandboxd/pool/egress_test.go b/sandboxd/pool/egress_test.go index 5af85ca..cd1506e 100644 --- a/sandboxd/pool/egress_test.go +++ b/sandboxd/pool/egress_test.go @@ -450,7 +450,7 @@ func TestEgressLaneCannotForkOrCheckpoint(t *testing.T) { if _, err := m.Checkpoint(t.Context(), sb.ID, Cred{Token: "tok"}, "", ""); !errors.Is(err, ErrNoEgressFork) { t.Errorf("Checkpoint on egress lane: got %v, want ErrNoEgressFork", err) } - if _, err := m.Promote(t.Context(), sb.ID, Cred{Token: "tok"}, "tpl", ""); !errors.Is(err, ErrNoEgressFork) { + if _, _, err := m.Promote(t.Context(), sb.ID, Cred{Token: "tok"}, "tpl", ""); !errors.Is(err, ErrNoEgressFork) { t.Errorf("Promote on egress lane: got %v, want ErrNoEgressFork", err) } } diff --git a/sandboxd/pool/intercept_test.go b/sandboxd/pool/intercept_test.go index f9899ac..5bd01e8 100644 --- a/sandboxd/pool/intercept_test.go +++ b/sandboxd/pool/intercept_test.go @@ -124,7 +124,7 @@ func TestInterceptPoolAllowsPromote(t *testing.T) { m.mu.Unlock() // The cluster root is shared, so an interception sandbox's disk carries no // node-private material: promote/checkpoint are unrestricted. - if _, err := m.Promote(t.Context(), sb.ID, Cred{Token: "tok"}, "tpl:x", ""); err != nil { + if _, _, err := m.Promote(t.Context(), sb.ID, Cred{Token: "tok"}, "tpl:x", ""); err != nil { t.Errorf("Promote of an interception-pool sandbox: %v, want success", err) } } diff --git a/sandboxd/pool/pool.go b/sandboxd/pool/pool.go index 00981e3..c9781c6 100644 --- a/sandboxd/pool/pool.go +++ b/sandboxd/pool/pool.go @@ -255,10 +255,11 @@ type Manager struct { archiveDeleteDefault time.Duration archiveEnabled bool archiveSweep atomic.Bool + archiveDeleteSweep atomic.Bool // archiving holds ids with an archive() export in flight, so the reap tick // and archive sweep don't both re-export the same sandbox; pendingCks pins - // checkpoint ids an archive is publishing before ArchiveCk can (a delete - // in that window would strand the claim). Both guarded by m.mu. + // checkpoint ids during archive publish and removal commits. Both guarded + // by m.mu. archiving map[string]struct{} pendingCks map[string]struct{} // egressListeners holds the per-sandbox egress proxy accept point, keyed by @@ -473,9 +474,6 @@ func NewManager(ctx context.Context, cfg *config.Config, eng Engine, secrets *eg if err := m.adoptPersistedPools(ctx); err != nil { return nil, err } - if m.archiveEnabled && m.ckptTTL == 0 { - log.WithFunc("pool.NewManager").Warn(ctx, "archive enabled with checkpoint_ttl_hours=0: a checkpoint whose delete fails is not reclaimed") - } return m, nil } @@ -517,6 +515,7 @@ func (m *Manager) Run(ctx context.Context) { m.reapOnce(ctx) m.idleOnce(ctx) m.archiveOnce(ctx) + go m.retryArchiveDeletes(ctx) case <-ckptSweep: m.sweepExpiredCheckpoints(ctx) } diff --git a/sandboxd/pool/pool_test.go b/sandboxd/pool/pool_test.go index fd2d855..a71d36f 100644 --- a/sandboxd/pool/pool_test.go +++ b/sandboxd/pool/pool_test.go @@ -765,6 +765,7 @@ type fakeEngine struct { hibernates, restores, snapRemoves []string snapSaves, exports, snapshots []string + exportContent []byte caInstalls []string // vsock sockets InstallCACert was called on staleReconciles []string // VM names ReconcileStaleCreate was called on installCAErr error @@ -863,8 +864,15 @@ func (f *fakeEngine) SnapshotSave(_ context.Context, _, snapName string) error { func (f *fakeEngine) SnapshotExport(_ context.Context, snapName, toDir string) error { f.mu.Lock() f.exports = append(f.exports, snapName) + content := slices.Clone(f.exportContent) f.mu.Unlock() - return os.MkdirAll(toDir, 0o750) + if err := os.MkdirAll(toDir, 0o750); err != nil { + return err + } + if content == nil { + return nil + } + return os.WriteFile(filepath.Join(toDir, "state.bin"), content, 0o600) } func (f *fakeEngine) SnapshotList(_ context.Context) ([]string, error) { diff --git a/sandboxd/pool/promote_test.go b/sandboxd/pool/promote_test.go index 5bffcfd..548290d 100644 --- a/sandboxd/pool/promote_test.go +++ b/sandboxd/pool/promote_test.go @@ -1,6 +1,7 @@ package pool import ( + "encoding/json" "errors" "os" "path/filepath" @@ -18,7 +19,7 @@ func TestPromoteThenClaimClonesFromTemplate(t *testing.T) { m := newTestManager(t, eng) parent := mustClaim(t, m, testKey) - gotKey, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:x", "") + gotKey, gotDigest, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:x", "") if err != nil { t.Fatalf("Promote: %v", err) } @@ -29,6 +30,9 @@ func TestPromoteThenClaimClonesFromTemplate(t *testing.T) { if gotKey != key { t.Errorf("returned key %+v, want %+v (the parent's axes)", gotKey, key) } + if gotDigest == "" { + t.Fatal("Promote returned an empty content digest") + } golden := filepath.Join(m.dataDir, "checkpoints", "tp_"+key.Hash(), "export") if fi, statErr := os.Stat(golden); statErr != nil || !fi.IsDir() { t.Fatalf("template export %s missing: %v", golden, statErr) @@ -44,6 +48,118 @@ func TestPromoteThenClaimClonesFromTemplate(t *testing.T) { if child.Key != key { t.Errorf("child key %+v, want %+v", child.Key, key) } + if child.TemplateDigest != gotDigest { + t.Errorf("claim template digest %q, want promoted content digest %q", child.TemplateDigest, gotDigest) + } + raw, err := m.tpls.ReadMeta(t.Context(), store.TemplateID(key.Hash())) + if err != nil { + t.Fatalf("read template metadata: %v", err) + } + var rec templateRecord + if err := json.Unmarshal(raw, &rec); err != nil { + t.Fatalf("decode template metadata: %v", err) + } + if rec.ContentDigest != gotDigest { + t.Errorf("persisted content digest %q, want %q", rec.ContentDigest, gotDigest) + } +} + +func TestRepromoteContentDigestTracksExportBytes(t *testing.T) { + eng := newFakeEngine() + eng.exportContent = []byte("generation one") + m := newTestManager(t, eng) + parent := mustClaim(t, m, testKey) + + _, firstDigest, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:digest", "") + if err != nil { + t.Fatalf("first Promote: %v", err) + } + _, secondDigest, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:digest", "") + if err != nil { + t.Fatalf("same-byte Promote: %v", err) + } + if secondDigest != firstDigest { + t.Errorf("same export bytes changed digest: %q then %q", firstDigest, secondDigest) + } + + eng.exportContent = []byte("generation two") + thirdKey, thirdDigest, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:digest", "") + if err != nil { + t.Fatalf("changed-byte Promote: %v", err) + } + if thirdDigest == firstDigest { + t.Errorf("changed export bytes kept digest %q", thirdDigest) + } + child, err := claimAny(t.Context(), m, thirdKey, 0) + if err != nil { + t.Fatalf("claim latest generation: %v", err) + } + if child.TemplateDigest != thirdDigest { + t.Errorf("claim digest %q, want latest %q", child.TemplateDigest, thirdDigest) + } +} + +func TestExportContentDigestCanonicalTree(t *testing.T) { + a, b := t.TempDir(), t.TempDir() + for _, root := range []string{a, b} { + if err := os.Mkdir(filepath.Join(root, "nested"), 0o750); err != nil { + t.Fatalf("mkdir tree: %v", err) + } + } + if err := os.Mkdir(filepath.Join(b, "ignored-empty-dir"), 0o750); err != nil { + t.Fatalf("mkdir empty dir: %v", err) + } + if err := os.WriteFile(filepath.Join(a, "z.bin"), []byte("z"), 0o600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(a, "nested", "a.bin"), []byte("a"), 0o600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(b, "nested", "a.bin"), []byte("a"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(b, "z.bin"), []byte("z"), 0o644); err != nil { + t.Fatal(err) + } + if err := os.Chtimes(filepath.Join(b, "z.bin"), time.Unix(1, 0), time.Unix(2, 0)); err != nil { + t.Fatal(err) + } + + da, err := exportContentDigest(a) + if err != nil { + t.Fatalf("digest tree a: %v", err) + } + db, err := exportContentDigest(b) + if err != nil { + t.Fatalf("digest tree b: %v", err) + } + if da != db { + t.Errorf("creation order/mode/mtime changed digest: %q != %q", da, db) + } + const want = "sha256:dfcaba733ae86b82c0760b4d2bb7828a595475e098e99a9769090aeda1390c0a" + if da != want { + t.Errorf("canonical digest %q, want %q", da, want) + } + if writeErr := os.WriteFile(filepath.Join(b, "z.bin"), []byte("changed"), 0o644); writeErr != nil { + t.Fatal(writeErr) + } + changed, err := exportContentDigest(b) + if err != nil { + t.Fatalf("digest changed tree: %v", err) + } + if changed == da { + t.Errorf("content change kept digest %q", changed) + } +} + +func TestExportContentDigestRejectsNonRegularEntry(t *testing.T) { + root := t.TempDir() + if err := os.Symlink("target", filepath.Join(root, "link")); err != nil { + t.Fatal(err) + } + if _, err := exportContentDigest(root); err == nil { + t.Fatal("digest accepted a symbolic link") + } } func TestPromoteHibernatedUsesWakeImage(t *testing.T) { @@ -54,7 +170,7 @@ func TestPromoteHibernatedUsesWakeImage(t *testing.T) { t.Fatalf("Hibernate: %v", err) } - if _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:hib", ""); err != nil { + if _, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:hib", ""); err != nil { t.Fatalf("Promote: %v", err) } if len(eng.snapSaves) != 0 { @@ -73,15 +189,15 @@ func TestPromoteValidations(t *testing.T) { m := newTestManager(t, eng, config.PoolSpec{PoolKey: testKey, Warm: 0}) parent := mustClaim(t, m, testKey) - if _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "_bad", ""); !errors.Is(err, ErrBadKey) { + if _, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "_bad", ""); !errors.Is(err, ErrBadKey) { t.Errorf("bad name: %v, want ErrBadKey", err) } - if _, err := m.Promote(t.Context(), parent.ID, Cred{Token: "wrong"}, "tpl:x", ""); !errors.Is(err, ErrUnknownSandbox) { + if _, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: "wrong"}, "tpl:x", ""); !errors.Is(err, ErrUnknownSandbox) { t.Errorf("bad token: %v, want ErrUnknownSandbox", err) } // Same template/net/size as the configured pool: the golden path would // collide with the pool's own. - if _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, testKey.Template, ""); !errors.Is(err, ErrPooledTemplate) { + if _, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, testKey.Template, ""); !errors.Is(err, ErrPooledTemplate) { t.Errorf("pooled key: %v, want ErrPooledTemplate", err) } if len(eng.snapSaves) != 0 { @@ -93,7 +209,7 @@ func TestDeleteTemplate(t *testing.T) { eng := newFakeEngine() m := newTestManager(t, eng, config.PoolSpec{PoolKey: testKey, Warm: 0}) parent := mustClaim(t, m, testKey) - if _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:del", ""); err != nil { + if _, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:del", ""); err != nil { t.Fatalf("Promote: %v", err) } key := types.PoolKey{Template: "tpl:del", Net: testKey.Net, Size: testKey.Size, Engine: testKey.Engine} @@ -147,10 +263,10 @@ func TestResolveGoldenSkipsPromotedEgressTemplate(t *testing.T) { if err = os.MkdirAll(filepath.Join(staging, store.ExportDir), 0o750); err != nil { t.Fatalf("mkdir export: %v", err) } - if err = m.commitTemplate(t.Context(), staging, id, ""); err != nil { + if _, err = m.commitTemplate(t.Context(), staging, id, ""); err != nil { t.Fatalf("seed template: %v", err) } - dir, release, err := m.resolveGolden(t.Context(), egKey) + dir, _, release, err := m.resolveGolden(t.Context(), egKey) if err != nil { t.Fatalf("resolveGolden: %v", err) } @@ -167,7 +283,7 @@ func TestTemplateTenantScopedDelete(t *testing.T) { if err != nil { t.Fatalf("claim: %v", err) } - key, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:tenant", "acme") + key, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:tenant", "acme") if err != nil { t.Fatalf("Promote: %v", err) } @@ -178,7 +294,7 @@ func TestTemplateTenantScopedDelete(t *testing.T) { if err := m.DeleteTemplate(t.Context(), key, "acme"); err != nil { t.Errorf("own delete: %v", err) } - if _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:tenant", "acme"); err != nil { + if _, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, "tpl:tenant", "acme"); err != nil { t.Fatalf("re-promote: %v", err) } if err := m.DeleteTemplate(t.Context(), key, ""); err != nil { @@ -203,10 +319,10 @@ func TestCommitTemplateRechecksOwnerUnderLock(t *testing.T) { } return staging } - if err := m.commitTemplate(t.Context(), stage(), id, "acme"); err != nil { + if _, err := m.commitTemplate(t.Context(), stage(), id, "acme"); err != nil { t.Fatalf("acme publish: %v", err) } - if err := m.commitTemplate(t.Context(), stage(), id, "beta"); !errors.Is(err, ErrTemplateOwned) { + if _, err := m.commitTemplate(t.Context(), stage(), id, "beta"); !errors.Is(err, ErrTemplateOwned) { t.Errorf("beta publish over acme: %v, want ErrTemplateOwned", err) } } @@ -223,7 +339,7 @@ func TestPromoteFailsClosedOnMetaError(t *testing.T) { if err != nil { t.Fatalf("claim: %v", err) } - if _, err := m.Promote(t.Context(), a.ID, Cred{Token: a.Token}, "shared:v1", "acme"); err != nil { + if _, _, err := m.Promote(t.Context(), a.ID, Cred{Token: a.Token}, "shared:v1", "acme"); err != nil { t.Fatalf("promote: %v", err) } key := types.PoolKey{Template: "shared:v1", Net: testKey.Net, Size: testKey.Size, Engine: testKey.Engine} @@ -232,7 +348,7 @@ func TestPromoteFailsClosedOnMetaError(t *testing.T) { t.Fatalf("chmod: %v", err) } t.Cleanup(func() { _ = os.Chmod(meta, 0o600) }) - if _, err := m.Promote(t.Context(), a.ID, Cred{Token: a.Token}, "shared:v1", "beta"); err == nil { + if _, _, err := m.Promote(t.Context(), a.ID, Cred{Token: a.Token}, "shared:v1", "beta"); err == nil { t.Fatal("promote succeeded despite an unreadable owner record") } } @@ -249,15 +365,15 @@ func TestPromoteRefusesCrossTenantOverwrite(t *testing.T) { return sb } a := claim("acme") - if _, err := m.Promote(t.Context(), a.ID, Cred{Token: a.Token}, "shared:v1", "acme"); err != nil { + if _, _, err := m.Promote(t.Context(), a.ID, Cred{Token: a.Token}, "shared:v1", "acme"); err != nil { t.Fatalf("acme promote: %v", err) } b := claim("beta") - if _, err := m.Promote(t.Context(), b.ID, Cred{Token: b.Token}, "shared:v1", "beta"); !errors.Is(err, ErrTemplateOwned) { + if _, _, err := m.Promote(t.Context(), b.ID, Cred{Token: b.Token}, "shared:v1", "beta"); !errors.Is(err, ErrTemplateOwned) { t.Errorf("beta overwrite: %v, want ErrTemplateOwned", err) } r := claim("") // root may replace anything - if _, err := m.Promote(t.Context(), r.ID, Cred{Token: r.Token}, "shared:v1", ""); err != nil { + if _, _, err := m.Promote(t.Context(), r.ID, Cred{Token: r.Token}, "shared:v1", ""); err != nil { t.Errorf("root replace: %v, want ok", err) } } @@ -267,7 +383,7 @@ func TestTemplateHashesSortedForMeshCompare(t *testing.T) { m := newTestManager(t, eng) parent := mustClaim(t, m, testKey) for _, name := range []string{"tpl:a", "tpl:b", "tpl:c", "tpl:d"} { - if _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, name, ""); err != nil { + if _, _, err := m.Promote(t.Context(), parent.ID, Cred{Token: parent.Token}, name, ""); err != nil { t.Fatalf("Promote %s: %v", name, err) } } diff --git a/sandboxd/pool/reconcile.go b/sandboxd/pool/reconcile.go index d069e67..61a03ff 100644 --- a/sandboxd/pool/reconcile.go +++ b/sandboxd/pool/reconcile.go @@ -27,6 +27,13 @@ func (m *Manager) Reconcile(ctx context.Context) error { if err != nil { return err } + for _, sb := range claims { + if sb.ArchiveCk != "" { + if markErr := m.markArchiveCk(sb.ArchiveCk); markErr != nil { + return fmt.Errorf("track archive checkpoint %s: %w", sb.ArchiveCk, markErr) + } + } + } vms, err := m.eng.List(ctx) if err != nil { return fmt.Errorf("list vms: %w", err) @@ -95,6 +102,7 @@ func (m *Manager) Reconcile(ctx context.Context) error { logger := log.WithFunc("pool.Reconcile") m.reclaimOrphanArchiveCks(ctx, claims) + m.retryArchiveDeletes(ctx) removed := m.sweepStaleVMs(ctx, live, owned) // A hibernate snapshot no adopted claim references is an orphan; @@ -268,7 +276,7 @@ func (m *Manager) reclaimOrphanArchiveCks(ctx context.Context, claims map[string if !mine || orig.ArchiveCk == ckpt.ID { continue } - if err := m.deleteCkLocked(ctx, ckpt.ID); err != nil { + if err := m.deleteArchiveCk(ctx, ckpt.ID); err != nil { logger.Warnf(ctx, "reclaim %s: %v", ckpt.ID, err) } else { logger.Infof(ctx, "reclaimed orphan archive ck %s", ckpt.ID) diff --git a/sandboxd/pool/template.go b/sandboxd/pool/template.go index b8cdf77..4f26200 100644 --- a/sandboxd/pool/template.go +++ b/sandboxd/pool/template.go @@ -2,12 +2,18 @@ package pool import ( "context" + "crypto/sha256" + "encoding/binary" + "encoding/hex" "encoding/json" "errors" "fmt" + "io" + "io/fs" "os" "path/filepath" "slices" + "strings" "sync" "time" @@ -20,36 +26,37 @@ import ( // from the requested key, never the reverse. Tenant attributes the promote // and scopes deletion; empty means the operator (root). type templateRecord struct { - ID string `json:"id"` - Tenant string `json:"tenant,omitempty"` - CreatedAt time.Time `json:"created_at"` + ID string `json:"id"` + Tenant string `json:"tenant,omitempty"` + ContentDigest string `json:"content_digest,omitempty"` + CreatedAt time.Time `json:"created_at"` } // Promote publishes a claimed sandbox as a template under (template, parent // net, parent size); later claims for that key clone from it. Re-promoting the // same name replaces it, and the caller owns its lifecycle (DeleteTemplate). // tenant attributes the record; empty means the operator (root). -func (m *Manager) Promote(ctx context.Context, id string, cred Cred, template, tenant string) (types.PoolKey, error) { +func (m *Manager) Promote(ctx context.Context, id string, cred Cred, template, tenant string) (types.PoolKey, string, error) { sb, ok := m.resolve(id, cred) if !ok { - return types.PoolKey{}, ErrUnknownSandbox + return types.PoolKey{}, "", ErrUnknownSandbox } if !types.NameRe.MatchString(template) { - return types.PoolKey{}, fmt.Errorf("%w: template %q must match %s", ErrBadKey, template, types.NameRe) + return types.PoolKey{}, "", fmt.Errorf("%w: template %q must match %s", ErrBadKey, template, types.NameRe) } if !sb.Key.Capturable() { - return types.PoolKey{}, ErrNoEgressFork + return types.PoolKey{}, "", ErrNoEgressFork } key := types.PoolKey{Template: template, Net: sb.Key.Net, Size: sb.Key.Size, Engine: sb.Key.Engine} if m.pooledHash(key.Hash()) { // A configured pool owns this key — promoting over it would // silently change what refills produce. - return types.PoolKey{}, ErrPooledTemplate + return types.PoolKey{}, "", ErrPooledTemplate } // Fast-fail before the export; commitTemplate re-checks under the // template lock, closing the check-then-publish race. if err := m.checkTemplateOwner(ctx, store.TemplateID(key.Hash()), tenant); err != nil { - return types.PoolKey{}, err + return types.PoolKey{}, "", err } // See Fork: the transition lock pins the source snapshot, and a started // promote must finish even if the caller hangs up. @@ -59,18 +66,19 @@ func (m *Manager) Promote(ctx context.Context, id string, cred Cred, template, t snap, cleanup, err := m.sourceSnap(ctx, sb) if err != nil { - return types.PoolKey{}, fmt.Errorf("promote %s: %w", sb.ID, err) + return types.PoolKey{}, "", fmt.Errorf("promote %s: %w", sb.ID, err) } defer cleanup() - if err := m.publishTemplate(ctx, snap, key, tenant); err != nil { - return types.PoolKey{}, fmt.Errorf("promote %s: %w", sb.ID, err) + digest, err := m.publishTemplate(ctx, snap, key, tenant) + if err != nil { + return types.PoolKey{}, "", fmt.Errorf("promote %s: %w", sb.ID, err) } if m.notifyTemplates != nil { m.notifyTemplates() } m.counters.promotes.Add(1) m.recordUsage(ctx, usageEvent{Event: "promote", ID: sb.ID, VMName: sb.VMName, Reference: key.Template}) - return key, nil + return key, digest, nil } // DeleteTemplate removes a promoted template. Configured pools are refused: @@ -254,7 +262,7 @@ func (m *Manager) checkTemplateOwner(ctx context.Context, id, tenant string) err // golden (no release), else a promoted template fetched from the store; // empty dir cold-boots. Only a true absence cold-boots — a backend failure // propagates rather than silently booting a template name as an image ref. -func (m *Manager) resolveGolden(ctx context.Context, key types.PoolKey) (string, func(), error) { +func (m *Manager) resolveGolden(ctx context.Context, key types.PoolKey) (string, string, func(), error) { m.mu.Lock() var dir string if p := m.pools[key]; p != nil { @@ -262,62 +270,141 @@ func (m *Manager) resolveGolden(ctx context.Context, key types.PoolKey) (string, } m.mu.Unlock() if dir != "" { - return dir, func() {}, nil + return dir, "", func() {}, nil } if key.Net == types.NetEgress { - return "", func() {}, nil // never resume a live-captured template on the egress lane; cold-boot instead + return "", "", func() {}, nil // never resume a live-captured template on the egress lane; cold-boot instead } id := store.TemplateID(key.Hash()) l := m.recLock(id) l.RLock() - dir, _, release, err := m.tpls.Fetch(ctx, id) + dir, meta, release, err := m.tpls.Fetch(ctx, id) if err != nil { l.RUnlock() m.recDone(id) if errors.Is(err, store.ErrNotFound) { - return "", func() {}, nil + return "", "", func() {}, nil } - return "", func() {}, err + return "", "", func() {}, err + } + var rec templateRecord + if err := json.Unmarshal(meta, &rec); err != nil { + release() + l.RUnlock() + m.recDone(id) + return "", "", func() {}, fmt.Errorf("decode template metadata: %w", err) } - return dir, func() { release(); l.RUnlock(); m.recDone(id) }, nil + return dir, rec.ContentDigest, func() { release(); l.RUnlock(); m.recDone(id) }, nil } // publishTemplate exports snap into the store under the key's template id. -func (m *Manager) publishTemplate(ctx context.Context, snap string, key types.PoolKey, tenant string) error { +func (m *Manager) publishTemplate(ctx context.Context, snap string, key types.PoolKey, tenant string) (string, error) { id := store.TemplateID(key.Hash()) staging, err := m.tpls.Stage(id) if err != nil { - return fmt.Errorf("stage template: %w", err) + return "", fmt.Errorf("stage template: %w", err) } defer func() { _ = os.RemoveAll(staging) }() if err = m.eng.SnapshotExport(ctx, snap, filepath.Join(staging, store.ExportDir)); err != nil { - return fmt.Errorf("export template: %w", err) + return "", fmt.Errorf("export template: %w", err) } return m.commitTemplate(ctx, staging, id, tenant) } -// commitTemplate writes the meta record, publishes the staged template, and -// registers it in the gossip set — owner re-checked and swap serialized -// under the template lock. -func (m *Manager) commitTemplate(ctx context.Context, staging, id, tenant string) error { - meta, err := json.Marshal(templateRecord{ID: id, Tenant: tenant, CreatedAt: time.Now()}) +// commitTemplate hashes the staged export once, then re-checks the owner and +// swaps meta+generation in, all serialized under the template lock. +func (m *Manager) commitTemplate(ctx context.Context, staging, id, tenant string) (string, error) { + digest, err := exportContentDigest(filepath.Join(staging, store.ExportDir)) if err != nil { - return err - } - if err := os.WriteFile(filepath.Join(staging, store.MetaFile), meta, 0o600); err != nil { - return err + return "", fmt.Errorf("digest template export: %w", err) } l := m.recLock(id) l.Lock() defer func() { l.Unlock(); m.recDone(id) }() - if err := m.checkTemplateOwner(ctx, id, tenant); err != nil { - return err + if ownerErr := m.checkTemplateOwner(ctx, id, tenant); ownerErr != nil { + return "", ownerErr + } + meta, err := json.Marshal(templateRecord{ID: id, Tenant: tenant, ContentDigest: digest, CreatedAt: time.Now()}) + if err != nil { + return "", err + } + if err := os.WriteFile(filepath.Join(staging, store.MetaFile), meta, 0o600); err != nil { + return "", err } if err := m.tpls.Publish(ctx, staging, id); err != nil { - return fmt.Errorf("publish template: %w", err) + return "", fmt.Errorf("publish template: %w", err) } m.tplMu.Lock() m.tplSet[id] = struct{}{} m.tplMu.Unlock() - return nil + return digest, nil +} + +// exportContentDigest identifies a snapshot export by a versioned canonical +// regular-file stream: slash-relative path, type, length, and content, globally +// sorted by path. Directories, ownership metadata, modes, and mtimes are excluded. +func exportContentDigest(root string) (string, error) { + type digestEntry struct { + path string + rel string + size int64 + } + var entries []digestEntry + err := filepath.WalkDir(root, func(path string, entry fs.DirEntry, walkErr error) error { + if walkErr != nil { + return walkErr + } + if path == root { + return nil + } + rel, err := filepath.Rel(root, path) + if err != nil { + return err + } + rel = filepath.ToSlash(rel) + info, err := entry.Info() + if err != nil { + return err + } + if entry.IsDir() { + return nil + } + if !info.Mode().IsRegular() { + return fmt.Errorf("unsupported export entry %s (%s)", rel, info.Mode().Type()) + } + entries = append(entries, digestEntry{path: path, rel: rel, size: info.Size()}) + return nil + }) + if err != nil { + return "", err + } + slices.SortFunc(entries, func(a, b digestEntry) int { return strings.Compare(a.rel, b.rel) }) + + h := sha256.New() + _, _ = io.WriteString(h, "sandbox-template-export-v1\x00") + for _, entry := range entries { + var frame [9]byte + frame[0] = 'f' + binary.BigEndian.PutUint64(frame[1:], uint64(len(entry.rel))) + _, _ = h.Write(frame[:]) + _, _ = io.WriteString(h, entry.rel) + binary.BigEndian.PutUint64(frame[:8], uint64(entry.size)) //nolint:gosec // regular file sizes are non-negative + _, _ = h.Write(frame[:8]) + f, err := os.Open(entry.path) //nolint:gosec // path is walked from private export staging + if err != nil { + return "", err + } + n, copyErr := io.Copy(h, f) + closeErr := f.Close() + if copyErr != nil { + return "", copyErr + } + if closeErr != nil { + return "", closeErr + } + if n != entry.size { + return "", fmt.Errorf("export entry %s changed size while hashing", entry.rel) + } + } + return "sha256:" + hex.EncodeToString(h.Sum(nil)), nil } diff --git a/sandboxd/server/server.go b/sandboxd/server/server.go index bd81e45..da6bd58 100644 --- a/sandboxd/server/server.go +++ b/sandboxd/server/server.go @@ -64,7 +64,7 @@ type Manager interface { Hibernate(ctx context.Context, id string, cred pool.Cred) error Wake(ctx context.Context, id string, cred pool.Cred) error Fork(ctx context.Context, id string, cred pool.Cred, count int, ttl time.Duration) ([]*types.Sandbox, error) - Promote(ctx context.Context, id string, cred pool.Cred, template, tenant string) (types.PoolKey, error) + Promote(ctx context.Context, id string, cred pool.Cred, template, tenant string) (types.PoolKey, string, error) DeleteTemplate(ctx context.Context, key types.PoolKey, tenant string) error Checkpoint(ctx context.Context, id string, cred pool.Cred, name, tenant string) (types.Checkpoint, error) Counters() pool.Counters @@ -351,9 +351,9 @@ func (s *Server) handlePromote(w http.ResponseWriter, r *http.Request) { return } id := r.PathValue("id") - key, err := s.mgr.Promote(r.Context(), id, s.bodyCred(r, req.Token), req.Template, tenantFrom(r.Context())) + key, digest, err := s.mgr.Promote(r.Context(), id, s.bodyCred(r, req.Token), req.Template, tenantFrom(r.Context())) writeResult(w, r, "promote", id, "promote failed", err, func() { - writeJSON(w, http.StatusOK, types.PromoteResponse{Key: key}) + writeJSON(w, http.StatusOK, types.PromoteResponse{Key: key, ContentDigest: digest}) }) } @@ -640,6 +640,6 @@ func (s *Server) resolveScope(r *http.Request) (string, bool) { func (s *Server) claimResponse(sb *types.Sandbox) types.ClaimResponse { return types.ClaimResponse{ ID: sb.ID, Token: sb.Token, Deadline: sb.Deadline, - OwnerAddr: s.advertise, FromCheckpoint: sb.FromCheckpoint, + OwnerAddr: s.advertise, FromCheckpoint: sb.FromCheckpoint, TemplateDigest: sb.TemplateDigest, } } diff --git a/sandboxd/server/server_test.go b/sandboxd/server/server_test.go index f8d813e..17066eb 100644 --- a/sandboxd/server/server_test.go +++ b/sandboxd/server/server_test.go @@ -27,7 +27,10 @@ func TestClaimHappyPath(t *testing.T) { mgr := &fakeManager{ claim: func(_ context.Context, key types.PoolKey, ttl time.Duration) (*types.Sandbox, error) { gotKey, gotTTL = key, ttl - return &types.Sandbox{ID: "sb_1", Token: "tok", Deadline: time.Unix(42, 0).UTC()}, nil + return &types.Sandbox{ + ID: "sb_1", Token: "tok", Deadline: time.Unix(42, 0).UTC(), + TemplateDigest: "sha256:claim-digest", + }, nil }, } ts := newTestServer(t, "", mgr, nil) @@ -48,6 +51,9 @@ func TestClaimHappyPath(t *testing.T) { if cr.ID != "sb_1" || cr.Token != "tok" || !cr.Deadline.Equal(time.Unix(42, 0)) { t.Errorf("got %+v", cr) } + if cr.TemplateDigest != "sha256:claim-digest" { + t.Errorf("template digest %q, want sha256:claim-digest", cr.TemplateDigest) + } want := types.PoolKey{Template: "rt:24.04", Net: types.NetNone, Size: types.SizeSmall, Engine: types.EngineCH} if gotKey != want { t.Errorf("key %+v, want defaults %+v", gotKey, want) @@ -562,6 +568,7 @@ func TestForkFlow(t *testing.T) { func TestPromoteAndDeleteTemplateFlow(t *testing.T) { var gotKey types.PoolKey mgr := &fakeManager{ + promoteContentDigest: "sha256:promoted-digest", promote: func(id, token, template string) error { switch { case token != "tok": @@ -583,7 +590,7 @@ func TestPromoteAndDeleteTemplateFlow(t *testing.T) { } ts := newTestServer(t, "sekret", mgr, nil) - promote := func(auth, body string) int { + promote := func(auth, body string) (int, types.PromoteResponse) { req, _ := http.NewRequestWithContext(t.Context(), http.MethodPost, ts.URL+"/v1/sandboxes/sb_1/promote", strings.NewReader(body)) if auth != "" { req.Header.Set("Authorization", auth) @@ -593,7 +600,13 @@ func TestPromoteAndDeleteTemplateFlow(t *testing.T) { t.Fatalf("do: %v", err) } defer resp.Body.Close() - return resp.StatusCode + var result types.PromoteResponse + if resp.StatusCode == http.StatusOK { + if err := json.NewDecoder(resp.Body).Decode(&result); err != nil { + t.Fatalf("decode promote response: %v", err) + } + } + return resp.StatusCode, result } for _, tt := range []struct { name, auth, body string @@ -606,9 +619,13 @@ func TestPromoteAndDeleteTemplateFlow(t *testing.T) { {"missing api token", "", `{"token":"tok","template":"tpl:x"}`, http.StatusUnauthorized}, } { t.Run("promote/"+tt.name, func(t *testing.T) { - if got := promote(tt.auth, tt.body); got != tt.want { + got, result := promote(tt.auth, tt.body) + if got != tt.want { t.Errorf("status %d, want %d", got, tt.want) } + if got == http.StatusOK && result.ContentDigest != "sha256:promoted-digest" { + t.Errorf("content digest %q, want sha256:promoted-digest", result.ContentDigest) + } }) } @@ -1391,16 +1408,17 @@ func credToken(cred pool.Cred) string { // the claim hook stands in for the provision result. Tenant-scoped methods // record the tenant they were handed in gotTenant. type fakeManager struct { - ckptDir string - hasCheckpoint map[string]bool - claim func(ctx context.Context, key types.PoolKey, ttl time.Duration) (*types.Sandbox, error) - release func(id, token string) error - releaseOp func(id string) error - socket func(id, token string) (string, error) - hibernate func(id, token string) error - wake func(id, token string) error - fork func(id, token string, count int, ttl time.Duration) ([]*types.Sandbox, error) - promote func(id, token, template string) error + ckptDir string + hasCheckpoint map[string]bool + claim func(ctx context.Context, key types.PoolKey, ttl time.Duration) (*types.Sandbox, error) + release func(id, token string) error + releaseOp func(id string) error + socket func(id, token string) (string, error) + hibernate func(id, token string) error + wake func(id, token string) error + fork func(id, token string, count int, ttl time.Duration) ([]*types.Sandbox, error) + promote func(id, token, template string) error + promoteContentDigest string deleteGolden func(key types.PoolKey) error hasGolden bool @@ -1478,15 +1496,15 @@ func (f *fakeManager) Fork(_ context.Context, id string, cred pool.Cred, count i return f.fork(id, credToken(cred), count, ttl) } -func (f *fakeManager) Promote(_ context.Context, id string, cred pool.Cred, template, tenant string) (types.PoolKey, error) { +func (f *fakeManager) Promote(_ context.Context, id string, cred pool.Cred, template, tenant string) (types.PoolKey, string, error) { f.gotTenant = tenant if f.promote == nil { - return types.PoolKey{}, pool.ErrUnknownSandbox + return types.PoolKey{}, "", pool.ErrUnknownSandbox } if err := f.promote(id, credToken(cred), template); err != nil { - return types.PoolKey{}, err + return types.PoolKey{}, "", err } - return types.PoolKey{Template: template, Net: types.NetNone, Size: types.SizeSmall}, nil + return types.PoolKey{Template: template, Net: types.NetNone, Size: types.SizeSmall}, f.promoteContentDigest, nil } func (f *fakeManager) DeleteTemplate(_ context.Context, key types.PoolKey, tenant string) error { diff --git a/sandboxd/store/dir/dir.go b/sandboxd/store/dir/dir.go index edb7a33..375afa1 100644 --- a/sandboxd/store/dir/dir.go +++ b/sandboxd/store/dir/dir.go @@ -5,6 +5,7 @@ package dir import ( "context" + "crypto/sha256" "errors" "fmt" "io/fs" @@ -13,17 +14,24 @@ import ( "regexp" "strings" + "golang.org/x/sys/unix" + "github.com/cocoonstack/sandbox/sandboxd/store" ) -const oldSuffix = ".old" +const ( + oldSuffix = ".old" + lockDir = ".locks" + recordLockStripes = 256 +) var _ store.Store = (*Store)(nil) // Store keeps records as //{export,meta.json}; staging dirs are // /-*.tmp siblings so publish is one rename. idRe names the // instance's id namespace — two instances (checkpoints, templates) share a -// root without seeing each other's records. +// root without seeing each other's records. A fixed SHA-256 lock stripe set +// coordinates readers and writers across processes without growing per id. type Store struct { root string idRe *regexp.Regexp @@ -33,6 +41,9 @@ func New(root string, idRe *regexp.Regexp) (*Store, error) { if err := os.MkdirAll(root, 0o750); err != nil { return nil, fmt.Errorf("create store dir: %w", err) } + if err := os.MkdirAll(filepath.Join(root, lockDir), 0o750); err != nil { + return nil, fmt.Errorf("create store lock dir: %w", err) + } return &Store{root: root, idRe: idRe}, nil } @@ -41,6 +52,12 @@ func (d *Store) Stage(id string) (string, error) { } func (d *Store) Publish(_ context.Context, staging, id string) error { + lock, err := d.lockRecord(id, unix.LOCK_EX) + if err != nil { + return fmt.Errorf("lock record: %w", err) + } + defer unlockRecord(lock) + final := filepath.Join(d.root, id) old := final + oldSuffix // Re-publish (re-promote) replaces by swap, not delete-then-rename: the @@ -73,18 +90,27 @@ func (d *Store) Publish(_ context.Context, staging, id string) error { } func (d *Store) Fetch(ctx context.Context, id string) (string, []byte, func(), error) { + final := filepath.Join(d.root, id) + lock, err := d.lockRecord(id, unix.LOCK_SH) + if err != nil { + return "", nil, nil, fmt.Errorf("lock record: %w", err) + } + release := func() { unlockRecord(lock) } // Meta is the commit marker: a half-published record stays invisible. meta, err := d.ReadMeta(ctx, id) if err != nil { + release() return "", nil, nil, err } - dir := filepath.Join(d.root, id, store.ExportDir) + dir := filepath.Join(final, store.ExportDir) if _, err := os.Stat(dir); os.IsNotExist(err) { + release() return "", nil, nil, store.ErrNotFound } else if err != nil { + release() return "", nil, nil, err } - return dir, meta, func() {}, nil + return dir, meta, release, nil } func (d *Store) ReadMeta(_ context.Context, id string) ([]byte, error) { @@ -121,6 +147,12 @@ func (d *Store) Metas(ctx context.Context) ([][]byte, error) { // Delete also clears an unswept .old so a sweep cannot resurrect the record. func (d *Store) Delete(_ context.Context, id string) error { + lock, err := d.lockRecord(id, unix.LOCK_EX) + if err != nil { + return fmt.Errorf("lock record: %w", err) + } + defer unlockRecord(lock) + final := filepath.Join(d.root, id) return errors.Join(os.RemoveAll(final), os.RemoveAll(final+oldSuffix)) } @@ -155,3 +187,30 @@ func (d *Store) SweepStaging() error { } return nil } + +func (d *Store) lockRecord(id string, op int) (*os.File, error) { + // Stripe files stay for the store lifetime so waiters never split across + // unlinked inodes. SHA-256 keeps one id on the same stripe across processes. + lock, err := os.OpenFile( //nolint:gosec // fixed path under our root + d.recordLockPath(id), os.O_CREATE|os.O_RDWR, 0o600, + ) + if err != nil { + return nil, err + } + if err := unix.Flock(int(lock.Fd()), op); err != nil { //nolint:gosec // unix fds fit int + _ = lock.Close() + return nil, err + } + return lock, nil +} + +func (d *Store) recordLockPath(id string) string { + sum := sha256.Sum256([]byte(id)) + stripe := int(sum[0]) % recordLockStripes + return filepath.Join(d.root, lockDir, fmt.Sprintf("%02x", stripe)) +} + +func unlockRecord(lock *os.File) { + _ = unix.Flock(int(lock.Fd()), unix.LOCK_UN) //nolint:gosec // unix fds fit int + _ = lock.Close() +} diff --git a/sandboxd/store/dir/dir_test.go b/sandboxd/store/dir/dir_test.go index 7d9b5cd..5e75b5e 100644 --- a/sandboxd/store/dir/dir_test.go +++ b/sandboxd/store/dir/dir_test.go @@ -1,10 +1,14 @@ package dir import ( + "errors" + "fmt" "os" "path/filepath" "testing" + "golang.org/x/sys/unix" + "github.com/cocoonstack/sandbox/sandboxd/store" "github.com/cocoonstack/sandbox/sandboxd/store/storetest" ) @@ -17,6 +21,88 @@ func TestDirBackendContract(t *testing.T) { storetest.RunContract(t, st) } +func TestFetchPinsGenerationAcrossStoreInstances(t *testing.T) { + const id = "ck_00000000000000aa" + root := t.TempDir() + reader, err := New(root, store.CheckpointIDRe) + if err != nil { + t.Fatalf("new reader: %v", err) + } + writer, err := New(root, store.CheckpointIDRe) + if err != nil { + t.Fatalf("new writer: %v", err) + } + staging, err := writer.Stage(id) + if err != nil { + t.Fatalf("stage first: %v", err) + } + seedRecord(t, staging, "first") + if err = writer.Publish(t.Context(), staging, id); err != nil { + t.Fatalf("publish first: %v", err) + } + + dir, meta, release, err := reader.Fetch(t.Context(), id) + if err != nil { + t.Fatalf("fetch first: %v", err) + } + if string(meta) != `{"id":"first"}` { + t.Fatalf("fetch meta %q, want first generation", meta) + } + if lock, lockErr := writer.lockRecord(id, unix.LOCK_EX|unix.LOCK_NB); !errors.Is(lockErr, unix.EWOULDBLOCK) { + if lockErr == nil { + unlockRecord(lock) + } + t.Fatalf("writer lock while fetch is live: %v, want EWOULDBLOCK", lockErr) + } + if _, statErr := os.Stat(dir); statErr != nil { + t.Fatalf("pinned export: %v", statErr) + } + release() + + staging, err = writer.Stage(id) + if err != nil { + t.Fatalf("stage second: %v", err) + } + seedRecord(t, staging, "second") + if err = writer.Publish(t.Context(), staging, id); err != nil { + t.Fatalf("publish second: %v", err) + } + _, meta, release, err = reader.Fetch(t.Context(), id) + if err != nil { + t.Fatalf("fetch second: %v", err) + } + defer release() + if string(meta) != `{"id":"second"}` { + t.Fatalf("fetch meta %q, want second generation", meta) + } +} + +func TestRecordLockFilesStayBoundedAcrossIDChurn(t *testing.T) { + root := t.TempDir() + st, err := New(root, store.CheckpointIDRe) + if err != nil { + t.Fatalf("New: %v", err) + } + want := map[string]struct{}{} + for i := range recordLockStripes * 16 { + id := fmt.Sprintf("ck_%016x", i) + want[filepath.Base(st.recordLockPath(id))] = struct{}{} + if _, _, _, err = st.Fetch(t.Context(), id); !errors.Is(err, store.ErrNotFound) { + t.Fatalf("Fetch missing %s: %v, want store.ErrNotFound", id, err) + } + } + locks, err := os.ReadDir(filepath.Join(root, lockDir)) + if err != nil { + t.Fatalf("read lock stripes: %v", err) + } + if len(locks) != len(want) { + t.Fatalf("lock files after churn = %d, want %d used stripes", len(locks), len(want)) + } + if len(locks) > recordLockStripes { + t.Fatalf("lock files after churn = %d, want at most %d", len(locks), recordLockStripes) + } +} + // TestSweepRecoversInterruptedPublish guards the re-publish swap: a crash // between moving the old generation aside and renaming the new one in leaves // only .old, and the startup sweep must restore it — the delete-then- diff --git a/sandboxd/store/s3/s3.go b/sandboxd/store/s3/s3.go index 1971395..f3bbf5c 100644 --- a/sandboxd/store/s3/s3.go +++ b/sandboxd/store/s3/s3.go @@ -102,7 +102,6 @@ func (s *Store) Publish(ctx context.Context, staging, id string) error { return fmt.Errorf("staging has no %s: %w", store.MetaFile, err) } gen := exportGen(metaRaw) - fresh := map[string]struct{}{s.key(id, store.MetaFile): {}} g, gctx := errgroup.WithContext(ctx) g.SetLimit(4) // files in parallel; each already multiparts internally err = filepath.WalkDir(staging, func(path string, d os.DirEntry, err error) error { @@ -117,7 +116,6 @@ func (s *Store) Publish(ctx context.Context, staging, id string) error { return nil } key := s.key(id, gen+strings.TrimPrefix(rel, store.ExportDir)) - fresh[key] = struct{}{} g.Go(func() error { return s.upload(gctx, key, path) }) return nil }) @@ -130,22 +128,7 @@ func (s *Store) Publish(ctx context.Context, staging, id string) error { if err = s.uploadReader(ctx, s.key(id, store.MetaFile), bytes.NewReader(metaRaw)); err != nil { return err } - // A re-publish (re-promote) may ship a different export file set: - // after the new meta commits, sweep keys the new generation did not - // write, or Fetch would download the union of generations. - keys, err := s.list(ctx, s.key(id, "")+"/") - if err != nil { - return err - } - var stale []string - for _, key := range keys { - if _, ok := fresh[key]; !ok { - stale = append(stale, key) - } - } - if err := s.deleteKeys(ctx, stale); err != nil { - return err - } + // Keep old generations because another node may have selected the previous meta; Delete reclaims them. return os.RemoveAll(staging) } @@ -229,16 +212,17 @@ func (s *Store) Metas(ctx context.Context) ([][]byte, error) { func (s *Store) Delete(ctx context.Context, id string) error { _ = os.RemoveAll(filepath.Join(s.staging, "cache", id)) //nolint:gosec // id pinned by idRe - // Uncommit first: dropping meta.json makes the record invisible to - // loads before any export object disappears under a concurrent fetch. - if err := s.deleteKeys(ctx, []string{s.key(id, store.MetaFile)}); err != nil { - return err - } + metaKey := s.key(id, store.MetaFile) keys, err := s.list(ctx, s.key(id, "")+"/") if err != nil { return err } - return s.deleteKeys(ctx, keys) + // Keep the commit marker so a failed export cleanup remains discoverable on retry. + keys = slices.DeleteFunc(keys, func(key string) bool { return key == metaKey }) + if err := s.deleteKeys(ctx, keys); err != nil { + return err + } + return s.deleteKeys(ctx, []string{metaKey}) } // SweepStaging clears local staging residue AND stale cache generations — diff --git a/sandboxd/store/s3/s3_test.go b/sandboxd/store/s3/s3_test.go index 4bf7dba..1f435a5 100644 --- a/sandboxd/store/s3/s3_test.go +++ b/sandboxd/store/s3/s3_test.go @@ -82,6 +82,90 @@ func TestDeleteRetryConverges(t *testing.T) { } } +func TestDeleteRetainsMetaUntilExportsAreGone(t *testing.T) { + const id = "ck_00000000000000bb" + metaKey := "ck/" + id + "/" + store.MetaFile + exportKeys := []string{ + "ck/" + id + "/export-first/disk.img", + "ck/" + id + "/export-second/disk.img", + } + t.Setenv("AWS_ACCESS_KEY_ID", "test") + t.Setenv("AWS_SECRET_ACCESS_KEY", "test") + t.Setenv("AWS_REQUEST_CHECKSUM_CALCULATION", "when_required") + t.Setenv("AWS_RESPONSE_CHECKSUM_VALIDATION", "when_required") + + for _, tc := range []struct { + name string + failList bool + failDelete bool + }{ + {name: "list failure", failList: true}, + {name: "delete failure", failDelete: true}, + } { + t.Run(tc.name, func(t *testing.T) { + fake := &fakeS3{ + objects: map[string][]byte{ + metaKey: []byte(`{"id":"` + id + `"}`), + exportKeys[0]: []byte("first"), + exportKeys[1]: []byte("second"), + }, + failList: tc.failList, + failDelete: tc.failDelete, + } + ts := httptest.NewServer(fake) + t.Cleanup(ts.Close) + st, err := New(t.Context(), Config{ + Bucket: "testbucket", Prefix: "ck/", Endpoint: ts.URL, + Region: "us-east-1", ForcePathStyle: true, + }, t.TempDir(), store.CheckpointIDRe) + if err != nil { + t.Fatalf("New: %v", err) + } + + if err = st.Delete(t.Context(), id); err == nil { + t.Fatal("first Delete succeeded with an injected failure") + } + fake.mu.Lock() + if _, ok := fake.objects[metaKey]; !ok { + fake.mu.Unlock() + t.Fatal("first Delete removed meta before export cleanup succeeded") + } + for _, batch := range fake.deleteBatches { + if slices.Contains(batch, metaKey) { + fake.mu.Unlock() + t.Fatalf("first Delete included meta in export batch %v", batch) + } + } + fake.failList = false + fake.failDelete = false + fake.deleteBatches = nil + fake.mu.Unlock() + + if err = st.Delete(t.Context(), id); err != nil { + t.Fatalf("second Delete: %v", err) + } + fake.mu.Lock() + defer fake.mu.Unlock() + if len(fake.objects) != 0 { + t.Errorf("objects left after retry: %v", fake.objects) + } + if len(fake.deleteBatches) < 2 { + t.Fatalf("delete batches = %v, want exports then meta", fake.deleteBatches) + } + last := len(fake.deleteBatches) - 1 + for i, batch := range fake.deleteBatches { + hasMeta := slices.Contains(batch, metaKey) + if (i == last) != hasMeta { + t.Errorf("delete batch %d = %v, meta must appear only in final batch", i, batch) + } + } + if batch := fake.deleteBatches[last]; len(batch) != 1 || batch[0] != metaKey { + t.Errorf("final delete batch = %v, want only %s", batch, metaKey) + } + }) + } +} + // TestFetchLegacyExportLayout: records published before per-generation // export prefixes keep flat export/ keys; Fetch must fall back to them. func TestFetchLegacyExportLayout(t *testing.T) { @@ -119,6 +203,53 @@ func TestFetchLegacyExportLayout(t *testing.T) { } } +func TestRepublishRetainsGenerationSelectedByAnotherStore(t *testing.T) { + const id = "ck_00000000000000aa" + fake := &fakeS3{objects: map[string][]byte{}} + ts := httptest.NewServer(fake) + t.Cleanup(ts.Close) + + t.Setenv("AWS_ACCESS_KEY_ID", "test") + t.Setenv("AWS_SECRET_ACCESS_KEY", "test") + t.Setenv("AWS_REQUEST_CHECKSUM_CALCULATION", "when_required") + t.Setenv("AWS_RESPONSE_CHECKSUM_VALIDATION", "when_required") + cfg := Config{ + Bucket: "testbucket", Prefix: "ck/", Endpoint: ts.URL, + Region: "us-east-1", ForcePathStyle: true, + } + reader, err := New(t.Context(), cfg, t.TempDir(), store.CheckpointIDRe) + if err != nil { + t.Fatalf("new reader: %v", err) + } + writer, err := New(t.Context(), cfg, t.TempDir(), store.CheckpointIDRe) + if err != nil { + t.Fatalf("new writer: %v", err) + } + publishRecord(t, writer, id, []byte(`{"id":"`+id+`","gen":1}`), "first") + selected, err := reader.ReadMeta(t.Context(), id) + if err != nil { + t.Fatalf("select first generation: %v", err) + } + publishRecord(t, writer, id, []byte(`{"id":"`+id+`","gen":2}`), "second") + + gen := filepath.Join(reader.staging, "selected-first") + if err = reader.populate(t.Context(), id, selected, gen); err != nil { + t.Fatalf("fetch selected first generation after re-publish: %v", err) + } + got, err := os.ReadFile(filepath.Join(gen, store.ExportDir, "disk.img")) //nolint:gosec // test path + if err != nil || string(got) != "first" { + t.Fatalf("selected generation bytes: %q, %v, want first", got, err) + } + if err = writer.Delete(t.Context(), id); err != nil { + t.Fatalf("delete all generations: %v", err) + } + fake.mu.Lock() + defer fake.mu.Unlock() + if len(fake.objects) != 0 { + t.Errorf("objects left after Delete: %v", fake.objects) + } +} + // TestS3BackendContractRealEndpoint runs the same contract against a real // S3 implementation (MinIO on a testbed) when SANDBOX_S3_E2E names its // endpoint — real list pagination, checksums, and path-style behavior the @@ -141,11 +272,35 @@ func TestS3BackendContractRealEndpoint(t *testing.T) { storetest.RunContract(t, st) } +func publishRecord(t *testing.T, st *Store, id string, meta []byte, content string) { + t.Helper() + staging, err := st.Stage(id) + if err != nil { + t.Fatalf("stage record: %v", err) + } + export := filepath.Join(staging, store.ExportDir) + if err := os.MkdirAll(export, 0o750); err != nil { + t.Fatalf("mkdir export: %v", err) + } + if err := os.WriteFile(filepath.Join(export, "disk.img"), []byte(content), 0o600); err != nil { + t.Fatalf("write export: %v", err) + } + if err := os.WriteFile(filepath.Join(staging, store.MetaFile), meta, 0o600); err != nil { + t.Fatalf("write meta: %v", err) + } + if err := st.Publish(t.Context(), staging, id); err != nil { + t.Fatalf("publish record: %v", err) + } +} + // fakeS3 implements just enough of the S3 REST surface (path-style) for // the backend: PutObject, GetObject, DeleteObjects, ListObjectsV2. type fakeS3 struct { - mu sync.Mutex - objects map[string][]byte // key -> body + mu sync.Mutex + objects map[string][]byte // key -> body + failList bool + failDelete bool + deleteBatches [][]string } func (f *fakeS3) ServeHTTP(w http.ResponseWriter, r *http.Request) { @@ -157,6 +312,10 @@ func (f *fakeS3) ServeHTTP(w http.ResponseWriter, r *http.Request) { body, _ := io.ReadAll(r.Body) f.objects[key] = body case r.Method == http.MethodGet && r.URL.Query().Get("list-type") == "2": + if f.failList { + http.Error(w, "injected list failure", http.StatusBadRequest) + return + } prefix := r.URL.Query().Get("prefix") delim := r.URL.Query().Get("delimiter") type object struct { @@ -233,16 +392,30 @@ func (f *fakeS3) ServeHTTP(w http.ResponseWriter, r *http.Request) { http.Error(w, "bad delete xml", http.StatusBadRequest) return } + batch := make([]string, len(req.Objects)) + for i := range req.Objects { + batch[i] = req.Objects[i].Key + } + f.deleteBatches = append(f.deleteBatches, batch) // Strict-backend emulation: absent keys answer a per-entry NoSuchKey // (AWS succeeds silently) so the client's tolerance stays exercised. type delErr struct { - Key string `xml:"Key"` - Code string `xml:"Code"` + Key string `xml:"Key"` + Code string `xml:"Code"` + Message string `xml:"Message,omitempty"` } var result struct { XMLName xml.Name `xml:"DeleteResult"` Errors []delErr `xml:"Error"` } + if f.failDelete { + result.Errors = append(result.Errors, delErr{ + Key: req.Objects[0].Key, Code: "AccessDenied", Message: "injected delete failure", + }) + w.Header().Set("Content-Type", "application/xml") + _ = xml.NewEncoder(w).Encode(result) + return + } for _, o := range req.Objects { if _, ok := f.objects[o.Key]; !ok { result.Errors = append(result.Errors, delErr{Key: o.Key, Code: "NoSuchKey"}) diff --git a/sandboxd/store/store.go b/sandboxd/store/store.go index 8d485b5..dba75af 100644 --- a/sandboxd/store/store.go +++ b/sandboxd/store/store.go @@ -40,17 +40,14 @@ type Store interface { // Stage returns a writable staging directory whose Publish is atomic. Stage(id string) (string, error) // Publish turns a staged directory into the record, replacing any - // previous generation atomically for listers (the dir backend renames; - // the s3 backend commits meta.json last). Same-id Publish/Fetch/Delete - // are serialized by the caller (the pool's per-template lock); two - // processes publishing one id remain unserialized. Request-path callers - // pass an uncancelable ctx so a started publish finishes. + // previous generation atomically for listers and never disturbing one a + // concurrent Fetch pinned (locks or retention, per backend). The caller + // serializes same-id operations in-process; request-path callers pass an + // uncancelable ctx so a started publish finishes. Publish(ctx context.Context, staging, id string) error // Fetch materializes a record's snapshot export as a local directory // cocoon can clone from, plus the meta it resolved on the way, and a - // release to call when the clone is done. The dir backend returns its - // path with a no-op release; the s3 backend serves a local cache - // generation. + // release pinning that exact generation — hold it until the clone is done. Fetch(ctx context.Context, id string) (dir string, meta []byte, release func(), err error) // ReadMeta returns a record's metadata, or an error when the record // does not exist. diff --git a/sandboxd/types/api.go b/sandboxd/types/api.go index de4ad7c..b56a4b4 100644 --- a/sandboxd/types/api.go +++ b/sandboxd/types/api.go @@ -44,6 +44,9 @@ type ClaimResponse struct { Token string `json:"token,omitempty"` Deadline time.Time `json:"deadline,omitzero"` OwnerAddr string `json:"owner_addr,omitempty"` + // TemplateDigest is the content identity of the promoted-template export + // this claim was cloned from; empty for any other source. + TemplateDigest string `json:"template_digest,omitempty"` // FromCheckpoint names the checkpoint a branched claim was born from, // so clients can reconstruct the checkpoint tree. @@ -102,11 +105,12 @@ type PromoteRequest struct { Template string `json:"template"` } -// PromoteResponse returns the template's full key: templates are node-local, -// so a cluster client needs the exact key (and this node's address) to claim -// from or delete the template later. +// PromoteResponse returns the template's full key and stable identity: +// templates are node-local, so a cluster client needs the exact key (and this +// node's address) to claim from or delete the template later. type PromoteResponse struct { - Key PoolKey `json:"key"` + Key PoolKey `json:"key"` + ContentDigest string `json:"content_digest"` } // PreviewRequest is the wire body of POST /v1/sandboxes/{id}/preview; auth diff --git a/sandboxd/types/types.go b/sandboxd/types/types.go index 4595af2..83dc5a4 100644 --- a/sandboxd/types/types.go +++ b/sandboxd/types/types.go @@ -184,6 +184,9 @@ type Sandbox struct { // FromCheckpoint names the checkpoint this sandbox branched from, for // lineage; empty for pool and template claims. FromCheckpoint string `json:"from_checkpoint,omitempty"` + // TemplateDigest identifies the promoted-template export this sandbox was + // initially cloned from; empty for every other provisioning source. + TemplateDigest string `json:"template_digest,omitempty"` // StaleSnap names a consumed wake snapshot a lagging journal still // references; dropped once a later write lands. Guarded by Transition. diff --git a/sdk/go/client.go b/sdk/go/client.go index d6c339e..637531f 100644 --- a/sdk/go/client.go +++ b/sdk/go/client.go @@ -18,6 +18,13 @@ import ( "time" ) +const ( + templateQueryParam = "template" + netQueryParam = "net" + sizeQueryParam = "size" + noRedirectQueryParam = "no_redirect" +) + // ClientOption configures Connect. type ClientOption func(*Client) @@ -95,7 +102,7 @@ func (c *Client) DeleteTemplate(ctx context.Context, template string, opts ...Op for _, opt := range opts { opt(&claim) } - u := url.Values{"template": {claim.Template}, "net": {claim.Net}, "size": {claim.Size}} + u := url.Values{templateQueryParam: {claim.Template}, netQueryParam: {claim.Net}, sizeQueryParam: {claim.Size}} redirect, err := c.deleteTemplates(ctx, c.addr, u) if err != nil || len(redirect) == 0 { return err @@ -103,7 +110,7 @@ func (c *Client) DeleteTemplate(ctx context.Context, template string, opts ...Op // The entry node doesn't hold the template but gossip named its owners. // The retry carries no_redirect, mirroring the claim protocol: the owner // answers for itself, never a second hop. - u.Set("no_redirect", "1") + u.Set(noRedirectQueryParam, "1") if tryErr := tryEach(redirect, func(addr string) error { _, retryErr := c.deleteTemplates(ctx, addr, u) return retryErr @@ -128,7 +135,10 @@ func (c *Client) ownerAt(ctx context.Context, addr, id, token string) (string, e // handleFrom builds a sandbox handle, defaulting the data-plane owner to the // node that answered when a single-node deployment omits owner_addr. func (c *Client) handleFrom(dialed string, cr claimResponse) *Sandbox { - return &Sandbox{ID: cr.ID, Deadline: cr.Deadline, FromCheckpoint: cr.FromCheckpoint, c: c, token: cr.Token, owner: cmp.Or(cr.OwnerAddr, dialed)} + return &Sandbox{ + ID: cr.ID, Deadline: cr.Deadline, FromCheckpoint: cr.FromCheckpoint, TemplateDigest: cr.TemplateDigest, + c: c, token: cr.Token, owner: cmp.Or(cr.OwnerAddr, dialed), + } } func (c *Client) claimAt(ctx context.Context, addr string, body []byte) (claimResponse, error) { @@ -430,6 +440,7 @@ type claimResponse struct { Deadline time.Time `json:"deadline"` OwnerAddr string `json:"owner_addr,omitempty"` FromCheckpoint string `json:"from_checkpoint,omitempty"` + TemplateDigest string `json:"template_digest,omitempty"` Redirect []string `json:"redirect,omitempty"` } @@ -454,6 +465,7 @@ type promoteResponse struct { Net string `json:"net"` Size string `json:"size"` } `json:"key"` + ContentDigest string `json:"content_digest"` } type errorResponse struct { diff --git a/sdk/go/client_test.go b/sdk/go/client_test.go index 7268903..5f2a648 100644 --- a/sdk/go/client_test.go +++ b/sdk/go/client_test.go @@ -45,7 +45,9 @@ func TestNewSendsClaim(t *testing.T) { if err := json.NewDecoder(r.Body).Decode(&gotBody); err != nil { t.Errorf("decode body: %v", err) } - _ = json.NewEncoder(w).Encode(claimResponse{ID: "sb_1", Token: "tok", Deadline: time.Unix(42, 0)}) + _ = json.NewEncoder(w).Encode(claimResponse{ + ID: "sb_1", Token: "tok", Deadline: time.Unix(42, 0), TemplateDigest: "sha256:template", + }) })) t.Cleanup(ts.Close) @@ -65,6 +67,9 @@ func TestNewSendsClaim(t *testing.T) { if sb.ID != "sb_1" || sb.token != "tok" { t.Errorf("handle %+v", sb) } + if sb.TemplateDigest != "sha256:template" { + t.Errorf("template digest %q, want sha256:template", sb.TemplateDigest) + } } func TestNewSurfacesServerError(t *testing.T) { diff --git a/sdk/go/sandbox.go b/sdk/go/sandbox.go index ed858fc..c683a1c 100644 --- a/sdk/go/sandbox.go +++ b/sdk/go/sandbox.go @@ -55,6 +55,9 @@ func (e *ExitError) Error() string { type Sandbox struct { ID string Deadline time.Time + // TemplateDigest is the content identity of the promoted-template export + // this sandbox was cloned from; empty for any other source. + TemplateDigest string // FromCheckpoint names the checkpoint this sandbox branched from; empty // for pool and template claims. @@ -176,7 +179,10 @@ func (s *Sandbox) Promote(ctx context.Context, template string) (*Template, erro if err != nil { return nil, err } - return &Template{Name: pr.Key.Template, c: s.c, addr: s.owner, net: pr.Key.Net, size: pr.Key.Size}, nil + return &Template{ + Name: pr.Key.Template, ContentDigest: pr.ContentDigest, + c: s.c, addr: s.owner, net: pr.Key.Net, size: pr.Key.Size, + }, nil } // Hibernate atomically snapshots the sandbox and stops its VM, freeing its diff --git a/sdk/go/template.go b/sdk/go/template.go index b0cad7e..f475f19 100644 --- a/sdk/go/template.go +++ b/sdk/go/template.go @@ -9,7 +9,8 @@ import ( // and Delete dial the owner directly — no gossip involved — so unlike the // name-based Client calls they are usable the instant Promote returns. type Template struct { - Name string + Name string + ContentDigest string c *Client addr string @@ -47,7 +48,12 @@ func (t *Template) New(ctx context.Context, opts ...Option) (*Sandbox, error) { // template is gone there) and gossip about same-name templates elsewhere is // never consulted. func (t *Template) Delete(ctx context.Context) error { - u := url.Values{"template": {t.Name}, "net": {t.net}, "size": {t.size}, "no_redirect": {"1"}} + u := url.Values{ + templateQueryParam: {t.Name}, + netQueryParam: {t.net}, + sizeQueryParam: {t.size}, + noRedirectQueryParam: {"1"}, + } _, err := t.c.deleteTemplates(ctx, t.addr, u) return err } diff --git a/sdk/go/template_test.go b/sdk/go/template_test.go new file mode 100644 index 0000000..09d10fb --- /dev/null +++ b/sdk/go/template_test.go @@ -0,0 +1,30 @@ +package sandbox + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "testing" +) + +func TestPromoteReturnsContentDigest(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost || r.URL.Path != "/v1/sandboxes/sb_1/promote" { + t.Errorf("got %s %s", r.Method, r.URL.Path) + } + _ = json.NewEncoder(w).Encode(map[string]any{ + "key": map[string]string{"template": "task:v1", "net": "none", "size": "small"}, + "content_digest": "sha256:promoted", + }) + })) + t.Cleanup(ts.Close) + + c := testClient(t, ts) + tpl, err := (&Sandbox{ID: "sb_1", token: "tok", owner: c.addr, c: c}).Promote(t.Context(), "task:v1") + if err != nil { + t.Fatalf("Promote: %v", err) + } + if tpl.Name != "task:v1" || tpl.ContentDigest != "sha256:promoted" { + t.Errorf("template %+v, want name and content digest", tpl) + } +} diff --git a/sdk/python/cocoonsandbox/client.py b/sdk/python/cocoonsandbox/client.py index cb94243..2c3701f 100644 --- a/sdk/python/cocoonsandbox/client.py +++ b/sdk/python/cocoonsandbox/client.py @@ -109,6 +109,7 @@ def _handle_from(self, dialed: str, reply: dict) -> Sandbox: owner=reply.get("owner_addr") or dialed, deadline=reply.get("deadline", ""), from_checkpoint=reply.get("from_checkpoint", ""), + template_digest=reply.get("template_digest", ""), ) def _post_json(self, addr: str, path: str, body: dict, verb: str) -> dict: diff --git a/sdk/python/cocoonsandbox/sandbox.py b/sdk/python/cocoonsandbox/sandbox.py index 5b9ccaa..a76f580 100644 --- a/sdk/python/cocoonsandbox/sandbox.py +++ b/sdk/python/cocoonsandbox/sandbox.py @@ -25,13 +25,14 @@ class Sandbox: """One claimed microVM.""" def __init__(self, client: Client, id: str, token: str, owner: str, - deadline: str = "", from_checkpoint: str = ""): + deadline: str = "", from_checkpoint: str = "", template_digest: str = ""): self._client = client self.id = id self.token = token self.owner = owner self.deadline = deadline self.from_checkpoint = from_checkpoint + self.template_digest = template_digest def __enter__(self) -> Sandbox: return self @@ -241,7 +242,10 @@ def promote(self, template: str) -> Template: reply = self._client._post_json(self.owner, f"/v1/sandboxes/{self.id}/promote", {"token": self.token, "template": template}, "promote") key = reply["key"] - return Template(self._client, self.owner, key["template"], key.get("net", ""), key.get("size", "")) + return Template( + self._client, self.owner, key["template"], key.get("net", ""), key.get("size", ""), + reply.get("content_digest", ""), + ) def start_lsp(self, language: str, root: str = "") -> Lsp: """Spawns the language server the flavor image provides for language; diff --git a/sdk/python/cocoonsandbox/template.py b/sdk/python/cocoonsandbox/template.py index 2fa3f2a..5aff4c6 100644 --- a/sdk/python/cocoonsandbox/template.py +++ b/sdk/python/cocoonsandbox/template.py @@ -15,12 +15,13 @@ class Template: """A promoted template on its owner node.""" - def __init__(self, client: Client, addr: str, name: str, net: str, size: str): + def __init__(self, client: Client, addr: str, name: str, net: str, size: str, content_digest: str = ""): self._client = client self._addr = addr self.name = name self.net = net self.size = size + self.content_digest = content_digest def new(self, ttl_seconds: int = 0) -> Sandbox: """Claims a sandbox cloned from the template, on the template's diff --git a/sdk/python/tests/test_client.py b/sdk/python/tests/test_client.py index 8e171c6..0bf34fe 100644 --- a/sdk/python/tests/test_client.py +++ b/sdk/python/tests/test_client.py @@ -56,9 +56,24 @@ def node(): def test_claim_happy_path(node): FakeNode.routes[("POST", "/v1/claim")] = lambda body, path: ( - 200, {"id": "sb_1", "token": "tok", "owner_addr": node}) + 200, {"id": "sb_1", "token": "tok", "owner_addr": node, "template_digest": "sha256:task"}) sb = Client(node).new("rt:24.04") assert sb.id == "sb_1" and sb.owner == node + assert sb.template_digest == "sha256:task" + + +def test_promote_returns_content_digest(node): + FakeNode.routes[("POST", "/v1/claim")] = lambda body, path: ( + 200, {"id": "sb_1", "token": "tok", "owner_addr": node}) + FakeNode.routes[("POST", "/v1/sandboxes/sb_1/promote")] = lambda body, path: ( + 200, { + "key": {"template": "task:v1", "net": "none", "size": "small"}, + "content_digest": "sha256:promoted", + }) + + tpl = Client(node).new("rt:24.04").promote("task:v1") + assert tpl.name == "task:v1" + assert tpl.content_digest == "sha256:promoted" def test_claim_follows_redirect_with_no_redirect(node):