fix(backend): robustness pass — spawn leak, timeouts, workspace fail-loud, 409 order, respond guard
Audit findings 4.1-4.6. - spawnFor: a rejected configure/start no longer leaks a registered runtime with a live Pi child (identity-checked teardown + rethrow); every later start used to hit "session runtime already active". - ThtRunner.run: default 60s timeout on every tht child (SIGKILL backstop), 120s for DWH-touching calls (sql preview/export, search pack); a dropped VPN mid-call no longer wedges the HTTP request forever. - configArg: a NAMED workspace whose yaml is missing now throws instead of silently falling back to the default config (operations were silently targeting the wrong workspace). - resume: the finalized/archived 409 is evaluated BEFORE the alreadyActive fast-path — the manifest is the truth even with a lingering runtime. - ollamaEnsure: exit-0 with non-JSON stdout is a failed check, not ok:true. - SessionBridge.respond: only the response matching the pending descriptor is forwarded to Pi; stale/duplicate submissions return 409 instead of being sent with the current gate's RPC id. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -107,13 +107,22 @@ export class SessionBridge {
|
|||||||
/** Eventi generati dal backend stesso (es. exit inatteso del child Pi), non da Pi. */
|
/** Eventi generati dal backend stesso (es. exit inatteso del child Pi), non da Pi. */
|
||||||
emitClientEvent(e: ClientEvent): void { this.fan(e); }
|
emitClientEvent(e: ClientEvent): void { this.fan(e); }
|
||||||
|
|
||||||
respond(uiResponse: object & { id: string }): void {
|
/**
|
||||||
|
* Deliver a reviewer response to Pi. Accepts ONLY a response matching the pending
|
||||||
|
* descriptor: a stale/duplicate response (old gate id, double submit) would otherwise
|
||||||
|
* be sent with the CURRENT gate's RPC id, flip the state to running, and leave the
|
||||||
|
* real gate waiting. Returns false when rejected so the route can 409.
|
||||||
|
*/
|
||||||
|
respond(uiResponse: object & { id: string }): boolean {
|
||||||
|
if (!this.pending || uiResponse.id !== (this.pending as { id?: unknown }).id) return false;
|
||||||
// Correla sull'id RPC di Pi; `value` porta l'uiResponse (con l'id del descriptor) cosi'
|
// Correla sull'id RPC di Pi; `value` porta l'uiResponse (con l'id del descriptor) cosi'
|
||||||
// il check interno del gate (resp.id === descriptor.id) regge.
|
// il check interno del gate (resp.id === descriptor.id) regge.
|
||||||
const piId = this.pendingPiId ?? uiResponse.id;
|
const piId = this.pendingPiId ?? uiResponse.id;
|
||||||
this.state = "running";
|
this.state = "running";
|
||||||
this.rpc.send({ type: "extension_ui_response", id: piId, value: JSON.stringify(uiResponse) });
|
this.rpc.send({ type: "extension_ui_response", id: piId, value: JSON.stringify(uiResponse) });
|
||||||
if (this.pending && uiResponse.id === this.pending.id) { this.pending = null; this.pendingPiId = null; }
|
this.pending = null;
|
||||||
|
this.pendingPiId = null;
|
||||||
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
steer(text: string): void {
|
steer(text: string): void {
|
||||||
|
|||||||
@@ -176,8 +176,15 @@ export class PiProcessManager {
|
|||||||
|
|
||||||
async spawnFor(sessionId: string, o: RuntimeOptions = {}): Promise<SessionRuntime> {
|
async spawnFor(sessionId: string, o: RuntimeOptions = {}): Promise<SessionRuntime> {
|
||||||
const rt = this.createFor(sessionId, o);
|
const rt = this.createFor(sessionId, o);
|
||||||
|
try {
|
||||||
await this.configure(rt, o);
|
await this.configure(rt, o);
|
||||||
this.start(sessionId, rt, o);
|
this.start(sessionId, rt, o);
|
||||||
|
} catch (error) {
|
||||||
|
// A rejected configure/start must not leak a registered runtime with a live
|
||||||
|
// child: every later start would see "session runtime already active".
|
||||||
|
this.teardownIfCurrent(sessionId, rt);
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
return rt;
|
return rt;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -246,7 +246,9 @@ export function sessionRoutes(
|
|||||||
} catch { return storageFailure(reply); }
|
} catch { return storageFailure(reply); }
|
||||||
const rt = d.mgr.get(id);
|
const rt = d.mgr.get(id);
|
||||||
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
|
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
|
||||||
rt.bridge.respond((req.body as any).ui_response);
|
if (!rt.bridge.respond((req.body as any).ui_response)) {
|
||||||
|
return reply.code(409).send({ error: "risposta non corrispondente al gate in attesa" });
|
||||||
|
}
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/steer", async (req, reply) => {
|
app.post("/sessions/:id/steer", async (req, reply) => {
|
||||||
@@ -273,6 +275,11 @@ export function sessionRoutes(
|
|||||||
} catch { return storageFailure(reply); }
|
} catch { return storageFailure(reply); }
|
||||||
if (!manifest) return reply.code(404).send({ error: "session not found" });
|
if (!manifest) return reply.code(404).send({ error: "session not found" });
|
||||||
const runner = runnerFor(principal);
|
const runner = runnerFor(principal);
|
||||||
|
// Read-only contract FIRST: a finalized/archived session must refuse resume even
|
||||||
|
// when a lingering runtime still looks active — the manifest is the truth.
|
||||||
|
if (manifest?.status === "finalized" || manifest?.archived) {
|
||||||
|
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
||||||
|
}
|
||||||
// This check belongs inside the per-session lock: a preceding cold Resume may have
|
// This check belongs inside the per-session lock: a preceding cold Resume may have
|
||||||
// installed a running runtime while this request was waiting.
|
// installed a running runtime while this request was waiting.
|
||||||
const existing = d.mgr.get(id);
|
const existing = d.mgr.get(id);
|
||||||
@@ -282,9 +289,6 @@ export function sessionRoutes(
|
|||||||
return reply.code(200).send({ id, alreadyActive: true });
|
return reply.code(200).send({ id, alreadyActive: true });
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (manifest?.status === "finalized" || manifest?.archived) {
|
|
||||||
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
|
||||||
}
|
|
||||||
const ensure = await d.readiness.ensure(settings.workspace ?? "", principal);
|
const ensure = await d.readiness.ensure(settings.workspace ?? "", principal);
|
||||||
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
||||||
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
|
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
|
||||||
|
|||||||
@@ -44,12 +44,15 @@ export class ThtRunner {
|
|||||||
withPrincipal(principal: PrincipalContext): ThtRunner { return new ThtRunner(this.cfg, principal); }
|
withPrincipal(principal: PrincipalContext): ThtRunner { return new ThtRunner(this.cfg, principal); }
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Resolve the `-c <config>` args. If `workspace` is given AND a matching
|
* Resolve the `-c <config>` args. A named workspace MUST exist: silently falling
|
||||||
* `workspaces/<workspace>.yaml` exists under harnessDir, select it; otherwise
|
* back to the default config would point every operation at the wrong workspace
|
||||||
* fall back to the default configPath.
|
* (wrong DB, wrong sessions dir) — fail loud instead.
|
||||||
*/
|
*/
|
||||||
private configArg(workspace?: string): string[] {
|
private configArg(workspace?: string): string[] {
|
||||||
if (workspace && existsSync(join(this.cfg.harnessDir, "workspaces", `${workspace}.yaml`))) {
|
if (workspace) {
|
||||||
|
if (!existsSync(join(this.cfg.harnessDir, "workspaces", `${workspace}.yaml`))) {
|
||||||
|
throw new Error(`workspace non trovato: workspaces/${workspace}.yaml (harness: ${this.cfg.harnessDir})`);
|
||||||
|
}
|
||||||
return ["-c", `workspaces/${workspace}.yaml`];
|
return ["-c", `workspaces/${workspace}.yaml`];
|
||||||
}
|
}
|
||||||
return ["-c", this.cfg.configPath];
|
return ["-c", this.cfg.configPath];
|
||||||
@@ -64,8 +67,14 @@ export class ThtRunner {
|
|||||||
return [...args, ...this.configArg(workspace)];
|
return [...args, ...this.configArg(workspace)];
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Every route awaits these children; without a ceiling, one hung DWH/vector call
|
||||||
|
// (dropped VPN mid-connect) wedges its HTTP request forever. Session/file commands
|
||||||
|
// get the default; DWH-touching commands pass a wider explicit budget.
|
||||||
|
static readonly DEFAULT_TIMEOUT_MS = 60_000;
|
||||||
|
static readonly DWH_TIMEOUT_MS = 120_000;
|
||||||
|
|
||||||
run(
|
run(
|
||||||
args: string[], workspace?: string, timeoutMs?: number,
|
args: string[], workspace?: string, timeoutMs: number = ThtRunner.DEFAULT_TIMEOUT_MS,
|
||||||
): Promise<{ code: number; stdout: string; stderr: string }> {
|
): Promise<{ code: number; stdout: string; stderr: string }> {
|
||||||
return new Promise((resolve) => {
|
return new Promise((resolve) => {
|
||||||
const env: NodeJS.ProcessEnv = { ...process.env };
|
const env: NodeJS.ProcessEnv = { ...process.env };
|
||||||
@@ -100,8 +109,8 @@ export class ThtRunner {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
private async json<T>(args: string[], workspace?: string): Promise<T> {
|
private async json<T>(args: string[], workspace?: string, timeoutMs?: number): Promise<T> {
|
||||||
const { code, stdout, stderr } = await this.run(args, workspace);
|
const { code, stdout, stderr } = await this.run(args, workspace, timeoutMs);
|
||||||
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
|
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
|
||||||
return JSON.parse(stdout) as T;
|
return JSON.parse(stdout) as T;
|
||||||
}
|
}
|
||||||
@@ -135,7 +144,7 @@ export class ThtRunner {
|
|||||||
/** Build and persist the deterministic F1 retrieval pack for a new session. */
|
/** Build and persist the deterministic F1 retrieval pack for a new session. */
|
||||||
async searchPack(question: string, sessionId: string, workspace?: string): Promise<void> {
|
async searchPack(question: string, sessionId: string, workspace?: string): Promise<void> {
|
||||||
const args = ["search", "pack", question, "--session", sessionId];
|
const args = ["search", "pack", question, "--session", sessionId];
|
||||||
const { code, stderr } = await this.run(args, workspace);
|
const { code, stderr } = await this.run(args, workspace, ThtRunner.DWH_TIMEOUT_MS);
|
||||||
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
|
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -169,11 +178,13 @@ export class ThtRunner {
|
|||||||
rows: unknown[][];
|
rows: unknown[][];
|
||||||
execution_ms: number;
|
execution_ms: number;
|
||||||
truncated: boolean;
|
truncated: boolean;
|
||||||
}>(a, workspace);
|
}>(a, workspace, ThtRunner.DWH_TIMEOUT_MS);
|
||||||
}
|
}
|
||||||
|
|
||||||
async sqlExport(id: string, workspace?: string) {
|
async sqlExport(id: string, workspace?: string) {
|
||||||
const { code, stdout, stderr } = await this.run(["sql", "export", "--session", id], workspace);
|
const { code, stdout, stderr } = await this.run(
|
||||||
|
["sql", "export", "--session", id], workspace, ThtRunner.DWH_TIMEOUT_MS,
|
||||||
|
);
|
||||||
if (code !== 0) throw new Error(`tht sql export exit ${code}: ${stderr.trim()}`);
|
if (code !== 0) throw new Error(`tht sql export exit ${code}: ${stderr.trim()}`);
|
||||||
return { path: stdout.trim() };
|
return { path: stdout.trim() };
|
||||||
}
|
}
|
||||||
@@ -197,17 +208,25 @@ export class ThtRunner {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async ollamaEnsure(workspace: string, timeoutSec: number): Promise<OllamaEnsureResult> {
|
async ollamaEnsure(workspace: string, timeoutSec: number): Promise<OllamaEnsureResult> {
|
||||||
|
// Process budget wider than the CLI's own --timeout so the CLI reports its
|
||||||
|
// failure itself; SIGKILL is only the backstop for a wedged child.
|
||||||
const { code, stdout, stderr } = await this.run(
|
const { code, stdout, stderr } = await this.run(
|
||||||
["ollama", "ensure", "--json", "--timeout", String(timeoutSec)],
|
["ollama", "ensure", "--json", "--timeout", String(timeoutSec)],
|
||||||
workspace,
|
workspace,
|
||||||
|
timeoutSec * 1000 + 30_000,
|
||||||
);
|
);
|
||||||
let parsed: Partial<OllamaEnsureResult> = {};
|
let parsed: Partial<OllamaEnsureResult> | null = null;
|
||||||
try { parsed = JSON.parse(stdout.trim() || "{}"); } catch { /* leave {} */ }
|
try { parsed = JSON.parse(stdout.trim()); } catch { /* not JSON */ }
|
||||||
if (code === 0) return { ok: true, ...parsed };
|
// Exit 0 with unparseable output is NOT a verified readiness: --json promises
|
||||||
|
// pristine JSON, so treat the violation as a failed check, never as ok.
|
||||||
|
if (code === 0 && parsed !== null) return { ok: true, ...parsed };
|
||||||
|
if (code === 0) {
|
||||||
|
return { ok: false, error: `tht ollama ensure: output non-JSON: ${stdout.trim().slice(0, 200)}` };
|
||||||
|
}
|
||||||
return {
|
return {
|
||||||
ok: false,
|
ok: false,
|
||||||
stage: parsed.stage,
|
stage: parsed?.stage,
|
||||||
error: parsed.error ?? (stderr.trim() || `tht ollama ensure exit ${code}`),
|
error: parsed?.error ?? (stderr.trim() || `tht ollama ensure exit ${code}`),
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1349,7 +1349,7 @@ test("POST resume tears down a created runtime when bridge binding fails", async
|
|||||||
expect(delivered).toEqual(["before", "post-failure probe"]);
|
expect(delivered).toEqual(["before", "post-failure probe"]);
|
||||||
});
|
});
|
||||||
|
|
||||||
test("POST /sessions/:id/response inoltra al bridge (no error)", async () => {
|
test("POST /sessions/:id/response senza gate pendente risponde 409 (risposta stantia)", async () => {
|
||||||
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), {
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), {
|
||||||
thtRunner: {
|
thtRunner: {
|
||||||
ollamaEnsure: async () => ({ ok: true }),
|
ollamaEnsure: async () => ({ ok: true }),
|
||||||
@@ -1361,9 +1361,11 @@ test("POST /sessions/:id/response inoltra al bridge (no error)", async () => {
|
|||||||
spawnFn: () => nodeSpawn("node", [FAKE, SCRIPT]) as any,
|
spawnFn: () => nodeSpawn("node", [FAKE, SCRIPT]) as any,
|
||||||
});
|
});
|
||||||
await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
||||||
|
// The fake Pi never emitted a ui_request: the bridge has no pending descriptor, so a
|
||||||
|
// response (stale UI, double submit) must be rejected instead of forwarded to Pi.
|
||||||
const res = await app.inject({ method: "POST", url: "/sessions/s1/response",
|
const res = await app.inject({ method: "POST", url: "/sessions/s1/response",
|
||||||
payload: { ui_response: { id: "u1", choices: ["a"] } } });
|
payload: { ui_response: { id: "u1", choices: ["a"] } } });
|
||||||
expect(res.statusCode).toBe(204);
|
expect(res.statusCode).toBe(409);
|
||||||
});
|
});
|
||||||
|
|
||||||
test("POST /sessions/:id/rename calls setName", async () => {
|
test("POST /sessions/:id/rename calls setName", async () => {
|
||||||
|
|||||||
@@ -146,13 +146,32 @@ test("respond correla sull'id RPC di Pi (non sull'id del descriptor) e azzera il
|
|||||||
// Pi emette la richiesta con il SUO id RPC ("pi-req-1"); il descriptor nel title ha id "u1".
|
// Pi emette la richiesta con il SUO id RPC ("pi-req-1"); il descriptor nel title ha id "u1".
|
||||||
fire({ type: "extension_ui_request", id: "pi-req-1", method: "input", title: JSON.stringify({ id: "u1", widget: "select" }) });
|
fire({ type: "extension_ui_request", id: "pi-req-1", method: "input", title: JSON.stringify({ id: "u1", widget: "select" }) });
|
||||||
// Il frontend rimanda l'id del descriptor ("u1").
|
// Il frontend rimanda l'id del descriptor ("u1").
|
||||||
b.respond({ id: "u1", choices: ["a"] });
|
expect(b.respond({ id: "u1", choices: ["a"] })).toBe(true);
|
||||||
// Pi correla la risposta sul SUO id ("pi-req-1") per risolvere ctx.ui.input; il value
|
// Pi correla la risposta sul SUO id ("pi-req-1") per risolvere ctx.ui.input; il value
|
||||||
// continua a portare l'id del descriptor, cosi' il check interno del gate regge.
|
// continua a portare l'id del descriptor, cosi' il check interno del gate regge.
|
||||||
expect(sent.at(-1)).toEqual({ type: "extension_ui_response", id: "pi-req-1", value: JSON.stringify({ id: "u1", choices: ["a"] }) });
|
expect(sent.at(-1)).toEqual({ type: "extension_ui_response", id: "pi-req-1", value: JSON.stringify({ id: "u1", choices: ["a"] }) });
|
||||||
expect(b.pendingWidget()).toBeNull();
|
expect(b.pendingWidget()).toBeNull();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("respond rifiuta risposte senza gate pendente o con id non corrispondente", () => {
|
||||||
|
const { rpc, sent, fire } = fakeRpc();
|
||||||
|
const b = new SessionBridge(rpc);
|
||||||
|
// Nessun gate pendente: la risposta non parte e lo stato non cambia.
|
||||||
|
expect(b.respond({ id: "u0", choices: ["a"] })).toBe(false);
|
||||||
|
expect(sent).toEqual([]);
|
||||||
|
|
||||||
|
fire({ type: "extension_ui_request", id: "pi-req-1", method: "input", title: JSON.stringify({ id: "u1", widget: "select" }) });
|
||||||
|
// Risposta stantia per un ALTRO gate: rifiutata, il gate vero resta pendente in waiting.
|
||||||
|
expect(b.respond({ id: "u0", choices: ["a"] })).toBe(false);
|
||||||
|
expect(sent).toEqual([]);
|
||||||
|
expect(b.turnState()).toBe("waiting");
|
||||||
|
expect(b.pendingWidget()).toEqual({ id: "u1", widget: "select" });
|
||||||
|
// Doppio submit: il primo passa, il secondo (pendente ormai nullo) viene rifiutato.
|
||||||
|
expect(b.respond({ id: "u1", choices: ["a"] })).toBe(true);
|
||||||
|
expect(b.respond({ id: "u1", choices: ["a"] })).toBe(false);
|
||||||
|
expect(sent).toHaveLength(1);
|
||||||
|
});
|
||||||
|
|
||||||
test("agent_end di Pi diventa un system_event agent_end per il FE", () => {
|
test("agent_end di Pi diventa un system_event agent_end per il FE", () => {
|
||||||
const { rpc, fire } = fakeRpc();
|
const { rpc, fire } = fakeRpc();
|
||||||
const b = new SessionBridge(rpc);
|
const b = new SessionBridge(rpc);
|
||||||
|
|||||||
@@ -114,16 +114,14 @@ test("run with exit != 0 propagates error with stderr", async () => {
|
|||||||
await expect(r.sessionList()).rejects.toThrow(/boom/);
|
await expect(r.sessionList()).rejects.toThrow(/boom/);
|
||||||
});
|
});
|
||||||
|
|
||||||
test("sessionNew with missing workspace file falls back to default configPath argv", async () => {
|
test("sessionNew with a missing workspace file fails loud (no silent default fallback)", async () => {
|
||||||
// harnessDir "/nope" has no workspaces/foo.yaml -> configArg falls back to default.
|
// harnessDir "/nope" has no workspaces/foo.yaml. Silently falling back to the default
|
||||||
|
// config would target the WRONG workspace (wrong DB, wrong sessions dir): must throw.
|
||||||
(spawn as any).mockClear();
|
(spawn as any).mockClear();
|
||||||
const r = new ThtRunner({ thtBin: "tht", harnessDir: "/nope", configPath: "config/tht.yaml" });
|
const r = new ThtRunner({ thtBin: "tht", harnessDir: "/nope", configPath: "config/tht.yaml" });
|
||||||
await r.sessionNew({ question: "q", workspace: "foo" });
|
await expect(r.sessionNew({ question: "q", workspace: "foo" }))
|
||||||
const [bin, argv] = (spawn as any).mock.calls[0];
|
.rejects.toThrow(/workspace non trovato: workspaces\/foo\.yaml/);
|
||||||
expect(bin).toBe("tht");
|
expect((spawn as any).mock.calls).toHaveLength(0);
|
||||||
// `--config`/`-c` is a PER-COMMAND option in tht (no global -c): it MUST follow
|
|
||||||
// the subcommand, never precede it. (Prepending it caused a live 500 "No such option: -c".)
|
|
||||||
expect(argv).toEqual(["session", "new", "q", "--json", "-c", "config/tht.yaml"]);
|
|
||||||
});
|
});
|
||||||
|
|
||||||
test("buildArgv appends -c AFTER the subcommand (never a global -c)", () => {
|
test("buildArgv appends -c AFTER the subcommand (never a global -c)", () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user