feat: harden runtime readiness and session workflow

This commit is contained in:
User
2026-07-14 10:27:25 +02:00
parent 6d7738d538
commit 6dbf93fff9
52 changed files with 1464 additions and 169 deletions
+7 -2
View File
@@ -11,6 +11,7 @@ import { metaRoutes, type ListModelsFn } from "./routes/meta.js";
import { settingsRoutes, effectiveSettings } from "./routes/settings.js";
import { createPiModelLister } from "./pi/list-models.js";
import { loadSettings, type Settings } from "./settings/settings-store.js";
import { ReadinessManager } from "./runtime/readiness-manager.js";
export interface BuildAppDeps {
thtRunner?: ThtRunner;
@@ -18,6 +19,7 @@ export interface BuildAppDeps {
spawnFn?: () => any;
listModels?: ListModelsFn;
getSettings?: () => Settings;
readiness?: ReadinessManager;
}
export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstance {
@@ -38,6 +40,10 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
});
const mgr = deps?.mgr ?? new PiProcessManager(config, deps?.spawnFn ? { spawnFn: deps.spawnFn } : undefined);
const hub = new SseHub();
const readiness = deps?.readiness ?? new ReadinessManager(
tht as ThtRunner,
Math.round(config.ollamaEnsureTimeoutMs / 1000),
);
const listModels = deps?.listModels ?? createPiModelLister(config);
const getSettings = deps?.getSettings ?? (() => effectiveSettings(config, loadSettings(config)));
@@ -50,8 +56,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
});
app.get("/health", async () => ({ status: "ok" }));
sessionRoutes(app, {
mgr, tht: tht as ThtRunner, hub, getSettings,
ollamaEnsureTimeoutSec: Math.round(config.ollamaEnsureTimeoutMs / 1000),
mgr, tht: tht as ThtRunner, hub, getSettings, readiness,
});
sqlRoutes(app, { tht: tht as ThtRunner });
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
+34 -12
View File
@@ -12,6 +12,15 @@ export interface SessionRuntime {
child: ChildProcessWithoutNullStreams;
}
export interface RuntimeOptions {
provider?: string;
model?: string;
thinking?: string;
author?: string;
question?: string;
mode?: "new" | "resume";
}
/** Injectable child-process boundary; callbacks may ignore arguments in simpler tests. */
type SpawnFn = (
command: string,
@@ -68,17 +77,13 @@ export class PiProcessManager {
get(id: string): SessionRuntime | undefined { return this.runtimes.get(id); }
async spawnFor(
sessionId: string,
o: { provider?: string; model?: string; thinking?: string; author?: string; question?: string; mode?: "new" | "resume" },
): Promise<SessionRuntime> {
// Idempotent per session id: tear down any existing runtime for this id
// first (before the cap check) so a resume/respawn neither leaks the old
// child nor falsely hits the process cap.
/** Spawn and register a runtime synchronously, without starting a model turn. */
createFor(sessionId: string, o: RuntimeOptions = {}): SessionRuntime {
// A duplicate start must never tear down a live session: that used to send
// SIGTERM to the in-flight Pi process and lose its pending gate.
const existing = this.runtimes.get(sessionId);
if (existing) {
existing.child.kill();
this.runtimes.delete(sessionId);
throw new Error(`session runtime already active: ${sessionId}`);
}
if (this.runtimes.size >= this.cfg.maxPiProcesses) {
throw new Error("max Pi processes reached");
@@ -103,24 +108,41 @@ export class PiProcessManager {
level: "error",
text: `Pi process exited unexpectedly (code ${code ?? "?"})`,
});
rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" });
}
});
return rt;
}
/** Configure model and thinking. Safe to run alongside deterministic retrieval. */
async configure(rt: SessionRuntime, o: RuntimeOptions = {}): Promise<void> {
const provider = canonicalPiProvider(o.provider ?? this.cfg.defaults.provider);
const model = o.model ?? this.cfg.defaults.model;
const thinking = o.thinking ?? this.cfg.defaults.thinking;
if (provider && model) {
await rpc.request({ type: "set_model", provider, modelId: model } as object & { type: string });
await rt.rpc.request({ type: "set_model", provider, modelId: model } as object & { type: string });
}
if (thinking) {
await rpc.request({ type: "set_thinking_level", level: thinking } as object & { type: string });
await rt.rpc.request({ type: "set_thinking_level", level: thinking } as object & { type: string });
}
}
/** Start the first turn only after callers have attached the runtime bridge. */
start(sessionId: string, rt: SessionRuntime, o: RuntimeOptions = {}): void {
if (this.runtimes.get(sessionId) !== rt) throw new Error("session runtime is no longer active");
const message = o.mode === "resume"
? `/riprendi-sessione ${sessionId}`
: `/nuova-domanda ${JSON.stringify(o.question ?? "")}`;
rpc.send({ type: "prompt", message });
rt.rpc.send({ type: "prompt", message });
}
async spawnFor(sessionId: string, o: RuntimeOptions = {}): Promise<SessionRuntime> {
const rt = this.createFor(sessionId, o);
await this.configure(rt, o);
this.start(sessionId, rt, o);
return rt;
}
+85 -11
View File
@@ -4,15 +4,57 @@ import type { ThtRunner } from "../tht/tht-runner.js";
import type { SseHub } from "../sse/sse-hub.js";
import type { Settings } from "../settings/settings-store.js";
import { getUser } from "../auth/auth.js";
import type { ReadinessManager } from "../runtime/readiness-manager.js";
export function sessionRoutes(
app: FastifyInstance,
d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; ollamaEnsureTimeoutSec: number },
d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; readiness: ReadinessManager },
) {
const info = (id: string, text: string, level = "info") =>
d.hub.publish(id, "info", { type: "info", level, text });
const bindRuntime = (id: string, rt: ReturnType<PiProcessManager["createFor"]>) =>
rt.bridge.onClientEvent((e) => {
if (e.type === "system_event" && e.event === "session_failed") {
void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
}
d.hub.publish(id, e.type, e);
});
const bootstrap = (
id: string,
rt: ReturnType<PiProcessManager["createFor"]>,
configure: Promise<void>,
retrieval: Promise<void> | null,
start: () => void,
) => {
void (async () => {
try {
if (retrieval) info(id, "Preparing retrieval context");
await Promise.all([configure, retrieval]);
info(id, "Starting model");
start();
} catch (error) {
d.mgr.teardown(id);
const text = error instanceof Error ? error.message : String(error);
void d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
rt.bridge.emitClientEvent({ type: "info", level: "error", text: `Session bootstrap failed: ${text}` });
rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
rt.bridge.emitClientEvent({ type: "system_event", event: "agent_end" });
}
})();
};
app.post("/runtime/prewarm", async (_req, reply) => {
const workspace = d.getSettings().workspace ?? "";
void d.readiness.ensure(workspace).catch(() => undefined);
return reply.code(202).send({ status: "warming" });
});
app.post("/sessions", async (req, reply) => {
const b = req.body as { question: string; name?: string };
const s = d.getSettings();
const ensure = await d.tht.ollamaEnsure(s.workspace ?? "", d.ollamaEnsureTimeoutSec);
const ensure = await d.readiness.ensure(s.workspace ?? "");
if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" });
// Settings (global) supply workspace/provider/model/thinking. The new-question
// form sends only the question text. `workspace` selects the tht `-c <config>`.
@@ -24,14 +66,22 @@ export function sessionRoutes(
model: s.model,
thinking: s.thinking,
});
const rt = await d.mgr.spawnFor(id, {
const options = {
provider: s.provider,
model: s.model,
thinking: s.thinking,
// Keep the saved preference in the manifest; F1 starts tool-first.
thinking: "off",
author: getUser(req).id,
question: b.question,
});
rt.bridge.onClientEvent((e) => d.hub.publish(id, e.type, e));
};
const rt = d.mgr.createFor(id, options);
bindRuntime(id, rt);
info(id, "Session created");
bootstrap(
id, rt, d.mgr.configure(rt, options),
d.tht.searchPack(b.question, id, s.workspace),
() => d.mgr.start(id, rt, options),
);
return { id };
});
app.get("/sessions", async () => d.tht.sessionList(d.getSettings().workspace));
@@ -51,20 +101,43 @@ export function sessionRoutes(
});
app.post("/sessions/:id/resume", async (req, reply) => {
const id = (req.params as any).id;
if (d.mgr.get(id)) {
// Idempotent resume: reattach the browser to the existing Pi runtime.
// Do not respawn it (which would discard a pending reviewer widget).
return reply.code(200).send({ id, alreadyActive: true });
}
const manifest = (await d.tht.sessionShow(id, d.getSettings().workspace)) as { status?: string; archived?: boolean } | null;
if (manifest?.status === "finalized" || manifest?.archived) {
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
}
const ensure = await d.tht.ollamaEnsure(d.getSettings().workspace ?? "", d.ollamaEnsureTimeoutSec);
const settings = d.getSettings();
const ensure = await d.readiness.ensure(settings.workspace ?? "");
if (!ensure.ok) return reply.code(503).send({ error: ensure.error ?? "Ollama/embeddings non disponibili" });
const rt = await d.mgr.resume(id, d.tht);
rt.bridge.onClientEvent((e) => d.hub.publish(id, e.type, e));
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
const options = {
provider: saved?.provider,
model: saved?.model,
// Phase 1 must reach a widget instead of exposing a long reasoning trace.
// The session keeps its saved preference for later turns.
thinking: "off",
author: getUser(req).id,
mode: "resume" as const,
};
await d.tht.reopenSession(id, settings.workspace);
const rt = d.mgr.createFor(id, options);
bindRuntime(id, rt);
info(id, "Resuming session");
bootstrap(id, rt, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
return reply.code(200).send({ id });
});
app.post("/sessions/:id/close", async (req) => {
const id = (req.params as { id: string }).id;
d.mgr.teardown(id);
d.hub.clear(id);
try {
await d.tht.closeSession(id, d.getSettings().workspace);
} finally {
d.mgr.teardown(id);
d.hub.clear(id);
}
return { closed: true };
});
app.get("/sessions/:id/events", (req, reply) => {
@@ -76,6 +149,7 @@ export function sessionRoutes(
reply.raw.writeHead(200, {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
Connection: "keep-alive",
"Access-Control-Allow-Origin": origin,
"Access-Control-Allow-Credentials": "true",
+44
View File
@@ -0,0 +1,44 @@
import type { OllamaEnsureResult, ThtRunner } from "../tht/tht-runner.js";
interface ReadyEntry {
expiresAt: number;
result: OllamaEnsureResult;
}
/**
* Deduplicates embedding readiness checks and keeps only short-lived successes.
* Failures are deliberately not cached so a submit can retry after a transient outage.
*/
export class ReadinessManager {
private inFlight = new Map<string, Promise<OllamaEnsureResult>>();
private ready = new Map<string, ReadyEntry>();
constructor(
private tht: ThtRunner,
private timeoutSec: number,
private ttlMs = 60_000,
private now: () => number = Date.now,
) {}
ensure(workspace = ""): Promise<OllamaEnsureResult> {
const cached = this.ready.get(workspace);
if (cached && cached.expiresAt > this.now()) return Promise.resolve(cached.result);
if (cached) this.ready.delete(workspace);
const current = this.inFlight.get(workspace);
if (current) return current;
const pending = this.tht.ollamaEnsure(workspace, this.timeoutSec)
.then((result) => {
if (result.ok) {
this.ready.set(workspace, { result, expiresAt: this.now() + this.ttlMs });
}
return result;
})
.finally(() => {
if (this.inFlight.get(workspace) === pending) this.inFlight.delete(workspace);
});
this.inFlight.set(workspace, pending);
return pending;
}
}
+12 -2
View File
@@ -83,8 +83,8 @@ export class ThtRunner {
return JSON.parse(stdout) as T;
}
private async ok(args: string[]): Promise<void> {
const { code, stderr } = await this.run(args);
private async ok(args: string[], workspace?: string): Promise<void> {
const { code, stderr } = await this.run(args, workspace);
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
}
@@ -109,6 +109,13 @@ export class ThtRunner {
return this.json<{ id: string }>(a, o.workspace);
}
/** Build and persist the deterministic F1 retrieval pack for a new session. */
async searchPack(question: string, sessionId: string, workspace?: string): Promise<void> {
const args = ["search", "pack", question, "--session", sessionId];
const { code, stderr } = await this.run(args, workspace);
if (code !== 0) throw new Error(`tht ${args.join(" ")} exit ${code}: ${stderr.trim()}`);
}
sessionList(workspace?: string) {
return this.json<SessionRow[]>(["session", "list", "--json"], workspace);
}
@@ -137,6 +144,9 @@ export class ThtRunner {
return { path: stdout.trim() };
}
closeSession(id: string, workspace?: string) { return this.ok(["session", "close", id], workspace); }
failSession(id: string, workspace?: string) { return this.ok(["session", "fail", id], workspace); }
reopenSession(id: string, workspace?: string) { return this.ok(["session", "reopen", id], workspace); }
setName(id: string, name: string) { return this.ok(["session", "set-name", id, "--name", name]); }
setGroup(id: string, group: string) { return this.ok(["session", "set-group", id, "--group", group]); }
archive(id: string) { return this.ok(["session", "archive", id]); }