From 818563c4085c9ebce69dca2cbc15ed71fe13f381 Mon Sep 17 00:00:00 2001 From: Codex Date: Tue, 8 Sep 2026 13:37:10 +0200 Subject: [PATCH] fix(core): keep workspace runtime available --- backend/src/operator-command.ts | 23 +++++--- backend/src/routes/sessions.ts | 36 ++++++++--- backend/test/operator-command.test.ts | 18 ++++++ backend/test/routes-sessions.test.ts | 59 +++++++++++++++++++ .../internal/modelprojection/projection.go | 36 ++++++++--- .../modelprojection/projection_test.go | 44 ++++++++++++++ 6 files changed, 192 insertions(+), 24 deletions(-) diff --git a/backend/src/operator-command.ts b/backend/src/operator-command.ts index d1408311..bfaf3945 100644 --- a/backend/src/operator-command.ts +++ b/backend/src/operator-command.ts @@ -88,17 +88,22 @@ async function workflowDiagnostics(config: AppConfig): Promise<{ ready: true; wo if (revisions.length === 0) throw new Error("workflow diagnostics unavailable"); return await withOperatorRunner(config, async (runner) => { for (const revision of revisions) { - const result = await runner.run(["doctor", "--json"], revision.snapshotPath); - let payload: unknown; + const runtime = await runner.acquireWorkspaceRuntime(revision.snapshotPath); try { - payload = JSON.parse(result.stdout); - } catch { - throw new Error("workflow diagnostics failed"); + const result = await runner.run(["doctor", "--json"], runtime.path); + let payload: unknown; + try { + payload = JSON.parse(result.stdout); + } catch { + throw new Error("workflow diagnostics failed"); + } + if ( + result.code !== 0 || !payload || typeof payload !== "object" + || (payload as { ok?: unknown }).ok !== true + ) throw new Error("workflow diagnostics failed"); + } finally { + runtime.release(); } - if ( - result.code !== 0 || !payload || typeof payload !== "object" - || (payload as { ok?: unknown }).ok !== true - ) throw new Error("workflow diagnostics failed"); } return { ready: true, workspaces: revisions.length }; }); diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 7923fdc2..4eeca147 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -150,13 +150,34 @@ export function sessionRoutes( }); app.addHook("onResponse", async (req) => { admissionLeases.get(req)?.(); }); + type SessionRevisionScan = { + revisions: Awaited>; + retainedComplete: boolean; + }; + let retainedSnapshotWarning: string | null = null; + /** Include retained historical descriptors so removed workspaces remain resumable. */ - const sessionRevisions = async () => { + const sessionRevisions = async (): Promise => { const registry = d.workspaceRegistry as Partial; if (typeof registry.listRetainedSnapshots === "function") { - return await registry.listRetainedSnapshots(); + try { + const revisions = await registry.listRetainedSnapshots(); + retainedSnapshotWarning = null; + return { revisions, retainedComplete: true }; + } catch (error) { + // A legacy/corrupt historical snapshot must not make active sessions (and their SSE + // reviewer gates) unreachable. The registry still rejects that snapshot; this fallback + // exposes only descriptors from the verified active state and deliberately disables + // retention reconciliation because the resulting session view is incomplete. + const detail = error instanceof Error ? error.message : "unknown error"; + if (retainedSnapshotWarning !== detail) { + console.warn("[sessions] retained snapshot discovery failed; using active snapshots:", detail); + retainedSnapshotWarning = detail; + } + return { revisions: await d.workspaceRegistry.list(), retainedComplete: false }; + } } - return await d.workspaceRegistry.list(); + return { revisions: await d.workspaceRegistry.list(), retainedComplete: true }; }; const isNotFound = (error: unknown) => @@ -200,7 +221,7 @@ export function sessionRoutes( }; let revisions: Awaited>; try { - revisions = await sessionRevisions(); + ({ revisions } = await sessionRevisions()); } catch (registryError) { // Sessions created before revision pinning still live under the installation's legacy // default config. Keep that compatibility path available when a fresh installation has @@ -559,8 +580,8 @@ export function sessionRoutes( ? { ...principal, isAdmin: false } : ownershipPrincipal(principal, "session.read_all"); const runner = runnerFor(scopedPrincipal); - const revisions = await sessionRevisions(); - const lists = await Promise.all(revisions + const revisionScan = await sessionRevisions(); + const lists = await Promise.all(revisionScan.revisions .map((revision) => runner.sessionList(revision.snapshotPath) as Promise)); const sessions = new Map(); for (const row of lists.flat()) { @@ -570,7 +591,8 @@ export function sessionRoutes( // Only an administrator-visible complete list (or the single local principal) is safe // input for retention. A remote per-user view can never discard another principal's pin. const reconcileSnapshotRetention = (d.workspaceRegistry as Partial).reconcileSnapshotRetention; - const hasCompleteRetentionView = (scope === "all" || principal.issuer === "local") + const hasCompleteRetentionView = revisionScan.retainedComplete + && (scope === "all" || principal.issuer === "local") && hasPermission(principal, "session.read_all"); if (hasCompleteRetentionView && typeof reconcileSnapshotRetention === "function") { const retained = [...new Set(list diff --git a/backend/test/operator-command.test.ts b/backend/test/operator-command.test.ts index 5661d222..1a7cfa72 100644 --- a/backend/test/operator-command.test.ts +++ b/backend/test/operator-command.test.ts @@ -5,6 +5,8 @@ const fakes = vi.hoisted(() => ({ catalogRepository: { close: vi.fn(async () => {}) }, createCatalogRepository: vi.fn(), runnerConfig: undefined as Record | undefined, + acquireWorkspaceRuntime: vi.fn(), + releaseWorkspaceRuntime: vi.fn(), run: vi.fn(async () => ({ code: 0, stdout: JSON.stringify({ ok: true }), @@ -24,6 +26,8 @@ vi.mock("../src/tht/tht-runner.js", () => ({ run = fakes.run; + acquireWorkspaceRuntime = fakes.acquireWorkspaceRuntime; + withPrincipal() { return this; } @@ -67,6 +71,12 @@ beforeEach(() => { fakes.createCatalogRepository.mockReset(); fakes.createCatalogRepository.mockReturnValue(fakes.catalogRepository); fakes.run.mockClear(); + fakes.acquireWorkspaceRuntime.mockReset(); + fakes.acquireWorkspaceRuntime.mockResolvedValue({ + path: "/data/workspace-registry/snapshots/runtime/workspace.yaml", + release: fakes.releaseWorkspaceRuntime, + }); + fakes.releaseWorkspaceRuntime.mockClear(); fakes.runnerConfig = undefined; }); @@ -78,5 +88,13 @@ test("workflow doctor gives schema-v4 runtime rendering a live Catalog repositor expect(fakes.createCatalogRepository).toHaveBeenCalledWith(config.catalogDatabase); expect(fakes.runnerConfig?.catalogRepository).toBe(fakes.catalogRepository); + expect(fakes.acquireWorkspaceRuntime).toHaveBeenCalledWith( + "/data/workspace-registry/snapshots/revision/workspace.yaml", + ); + expect(fakes.run).toHaveBeenCalledWith( + ["doctor", "--json"], + "/data/workspace-registry/snapshots/runtime/workspace.yaml", + ); + expect(fakes.releaseWorkspaceRuntime).toHaveBeenCalledOnce(); expect(fakes.catalogRepository.close).toHaveBeenCalledOnce(); }); diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 93320887..29a86711 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -364,6 +364,65 @@ test("retention scans a removed workspace's retained snapshot", async () => { expect(response.json()).toEqual([expect.objectContaining({ id: "resumable" })]); }); +test("active sessions remain available when retained snapshot discovery is unreadable", async () => { + const activeSnapshot = "/registry/snapshots/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa/active.yaml"; + const retainedFailure = "invalid retained snapshot /registry/snapshots/secret/legacy.yaml"; + const reconcileSnapshotRetention = vi.fn(async () => {}); + const subscribed = vi.fn(); + const warning = vi.spyOn(console, "warn").mockImplementation(() => {}); + const manifest = { + id: "live-session", status: "open", archived: false, + workspace_id: "active", workspace_revision: "a".repeat(40), + }; + const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), { + thtRunner: { + withPrincipal: () => ({ + sessionList: async (snapshotPath: string) => { + expect(snapshotPath).toBe(activeSnapshot); + return [manifest]; + }, + sessionShow: async (id: string, snapshotPath?: string) => { + if (id === manifest.id && snapshotPath === activeSnapshot) return manifest; + throw new Error("session not found"); + }, + }), + } as any, + mgr: { get: () => undefined } as any, + hub: { + subscribe: (_id: string, _send: unknown, options: { close?: () => void }) => { + subscribed(); + options.close?.(); + return () => {}; + }, + } as any, + workspaceRegistry: { + listRetainedSnapshots: async () => { throw new Error(retainedFailure); }, + list: async () => [{ + id: "active", commit: "a".repeat(40), blob: "b".repeat(40), snapshotPath: activeSnapshot, + }], + reconcileSnapshotRetention, + } as any, + }); + + try { + const list = await app.inject({ method: "GET", url: "/sessions", headers: aliceHeaders }); + const detail = await app.inject({ method: "GET", url: `/sessions/${manifest.id}`, headers: aliceHeaders }); + const events = await app.inject({ method: "GET", url: `/sessions/${manifest.id}/events`, headers: aliceHeaders }); + + expect(list.statusCode).toBe(200); + expect(list.json()).toEqual([expect.objectContaining({ id: manifest.id, active: false })]); + expect(detail.statusCode).toBe(200); + expect(detail.json()).toMatchObject(manifest); + expect(events.statusCode).toBe(200); + expect(subscribed).toHaveBeenCalledOnce(); + expect(reconcileSnapshotRetention).not.toHaveBeenCalled(); + expect(list.body + detail.body + events.body).not.toContain(retainedFailure); + } finally { + warning.mockRestore(); + await app.close(); + } +}); + test("the single local installation listing reconciles its resumable workspace pins", async () => { const retained = vi.fn(async () => {}); const retainedRevision = "d".repeat(40); diff --git a/tools/tht/internal/modelprojection/projection.go b/tools/tht/internal/modelprojection/projection.go index 12225f72..85eb8b62 100644 --- a/tools/tht/internal/modelprojection/projection.go +++ b/tools/tht/internal/modelprojection/projection.go @@ -140,7 +140,8 @@ func Render(installation config.Installation) (map[string][]byte, error) { }, nil } -// Generate publishes all adapters as one directory generation. A failed replacement restores the +// Generate publishes all adapters as one directory generation. An unchanged generation stays in +// place so active file bind mounts keep their source identity. A failed replacement restores the // previous directory, so callers never observe a successfully returned mixed generation. func Generate(installation config.Installation) error { artifacts, err := Render(installation) @@ -152,6 +153,17 @@ func Generate(installation config.Installation) error { if err := os.MkdirAll(parent, 0o755); err != nil { return fmt.Errorf("create model projection parent: %w", err) } + info, statErr := os.Lstat(target) + if statErr == nil { + if !info.IsDir() || info.Mode()&os.ModeSymlink != 0 { + return fmt.Errorf("current model projection path is not a regular directory") + } + if projectionMatches(target, artifacts) { + return nil + } + } else if !os.IsNotExist(statErr) { + return fmt.Errorf("inspect current model projection generation: %w", statErr) + } candidate, err := os.MkdirTemp(parent, ".model-projections-candidate-*") if err != nil { return fmt.Errorf("create model projection candidate: %w", err) @@ -161,19 +173,12 @@ func Generate(installation config.Installation) error { return err } - info, statErr := os.Lstat(target) if os.IsNotExist(statErr) { if err := renameProjectionDirectory(candidate, target); err != nil { return fmt.Errorf("publish model projection generation: %w", err) } return nil } - if statErr != nil { - return fmt.Errorf("inspect current model projection generation: %w", statErr) - } - if !info.IsDir() || info.Mode()&os.ModeSymlink != 0 { - return fmt.Errorf("current model projection path is not a regular directory") - } previous, err := absentTemporaryPath(parent) if err != nil { return fmt.Errorf("reserve previous model projection generation: %w", err) @@ -191,6 +196,21 @@ func Generate(installation config.Installation) error { return nil } +func projectionMatches(directory string, artifacts map[string][]byte) bool { + for _, relative := range sortedArtifactPaths(artifacts) { + path := filepath.Join(directory, filepath.FromSlash(relative)) + info, err := os.Lstat(path) + if err != nil || !info.Mode().IsRegular() { + return false + } + actual, err := os.ReadFile(path) + if err != nil || !bytes.Equal(actual, artifacts[relative]) { + return false + } + } + return true +} + func writeProjectionCandidate(directory string, artifacts map[string][]byte) error { if err := os.Chmod(directory, 0o755); err != nil { return fmt.Errorf("protect model projection candidate: %w", err) diff --git a/tools/tht/internal/modelprojection/projection_test.go b/tools/tht/internal/modelprojection/projection_test.go index 1253d9ad..de6bb3bd 100644 --- a/tools/tht/internal/modelprojection/projection_test.go +++ b/tools/tht/internal/modelprojection/projection_test.go @@ -56,6 +56,50 @@ func TestRenderProducesDeterministicCatalogPiAndComposeProjections(t *testing.T) } } +func TestGenerateDoesNotReplaceUnchangedProjection(t *testing.T) { + installation := projectionFixture(t) + if err := Generate(installation); err != nil { + t.Fatalf("Generate() initial error = %v", err) + } + paths := []string{ + installation.GeneratedModelCatalogPath(), + installation.GeneratedPiModelsPath(), + installation.GeneratedPiSettingsPath(), + installation.ModelProjectionComposePath(), + } + before := make(map[string]os.FileInfo, len(paths)) + for _, path := range paths { + info, err := os.Stat(path) + if err != nil { + t.Fatal(err) + } + before[path] = info + } + + originalRename := renameProjectionDirectory + t.Cleanup(func() { renameProjectionDirectory = originalRename }) + renames := 0 + renameProjectionDirectory = func(oldPath, newPath string) error { + renames++ + return os.Rename(oldPath, newPath) + } + if err := Generate(installation); err != nil { + t.Fatalf("Generate() repeated error = %v", err) + } + if renames != 0 { + t.Fatalf("Generate() replaced an unchanged generation with %d renames", renames) + } + for _, path := range paths { + after, err := os.Stat(path) + if err != nil { + t.Fatal(err) + } + if !os.SameFile(before[path], after) { + t.Fatalf("Generate() replaced unchanged artifact %q", path) + } + } +} + func TestGenerateRestoresWholePreviousGenerationWhenPublishFails(t *testing.T) { installation := projectionFixture(t) if err := Generate(installation); err != nil {