fix(core): keep workspace runtime available

This commit is contained in:
Codex
2026-09-08 13:37:10 +02:00
parent 50c546e42d
commit 67b02ce10e
6 changed files with 192 additions and 24 deletions
+6 -1
View File
@@ -88,7 +88,9 @@ 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);
const runtime = await runner.acquireWorkspaceRuntime(revision.snapshotPath);
try {
const result = await runner.run(["doctor", "--json"], runtime.path);
let payload: unknown;
try {
payload = JSON.parse(result.stdout);
@@ -99,6 +101,9 @@ async function workflowDiagnostics(config: AppConfig): Promise<{ ready: true; wo
result.code !== 0 || !payload || typeof payload !== "object"
|| (payload as { ok?: unknown }).ok !== true
) throw new Error("workflow diagnostics failed");
} finally {
runtime.release();
}
}
return { ready: true, workspaces: revisions.length };
});
+29 -7
View File
@@ -150,13 +150,34 @@ export function sessionRoutes(
});
app.addHook("onResponse", async (req) => { admissionLeases.get(req)?.(); });
type SessionRevisionScan = {
revisions: Awaited<ReturnType<typeof d.workspaceRegistry.list>>;
retainedComplete: boolean;
};
let retainedSnapshotWarning: string | null = null;
/** Include retained historical descriptors so removed workspaces remain resumable. */
const sessionRevisions = async () => {
const sessionRevisions = async (): Promise<SessionRevisionScan> => {
const registry = d.workspaceRegistry as Partial<WorkspaceRegistry>;
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 await d.workspaceRegistry.list();
return { revisions: await d.workspaceRegistry.list(), retainedComplete: false };
}
}
return { revisions: await d.workspaceRegistry.list(), retainedComplete: true };
};
const isNotFound = (error: unknown) =>
@@ -200,7 +221,7 @@ export function sessionRoutes(
};
let revisions: Awaited<ReturnType<typeof d.workspaceRegistry.list>>;
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<SessionRow[]>));
const sessions = new Map<string, SessionRow>();
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<WorkspaceRegistry>).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
+18
View File
@@ -5,6 +5,8 @@ const fakes = vi.hoisted(() => ({
catalogRepository: { close: vi.fn(async () => {}) },
createCatalogRepository: vi.fn(),
runnerConfig: undefined as Record<string, unknown> | 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();
});
+59
View File
@@ -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);
@@ -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)
@@ -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 {