From 2ff63d371f272759732c5210bd0de1c7811380b4 Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 21 Jul 2026 12:14:26 +0200 Subject: [PATCH] feat: harden workflow gates and expose token usage --- backend/src/bridge/session-bridge.ts | 57 +++++++++--- backend/src/pi/pi-process-manager.ts | 5 +- backend/test/pi-process-manager.test.ts | 37 ++++++++ backend/test/session-bridge.test.ts | 37 ++++++++ frontend/src/api/types.ts | 9 ++ frontend/src/shell/SteerInput.test.tsx | 50 ++++++++++- frontend/src/shell/SteerInput.tsx | 71 +++++++++++---- frontend/src/store/sessionStore.test.ts | 20 +++++ frontend/src/store/sessionStore.ts | 15 +++- frontend/src/stream/useSessionStream.test.tsx | 65 +++++++++++++- frontend/src/stream/useSessionStream.ts | 67 +++++++++++++- frontend/src/viewers/ArtifactView.test.tsx | 5 ++ frontend/src/viewers/ArtifactView.tsx | 3 + .../__tests__/gate_agent_end_safety.test.js | 33 +++++++ .../gate/__tests__/gate_antibypass.test.js | 26 ++++++ .../__tests__/gate_confirm_autoclose.test.js | 39 +++++++- .../__tests__/gate_resume_kickoff.test.js | 2 + .../__tests__/gate_select_decision.test.js | 18 +++- harness/.pi/extensions/tht-gate.js | 90 ++++++++++++++++--- harness/.pi/skills/tht-sessione/SKILL.md | 8 +- 20 files changed, 604 insertions(+), 53 deletions(-) create mode 100644 harness/.pi/extensions/gate/__tests__/gate_agent_end_safety.test.js diff --git a/backend/src/bridge/session-bridge.ts b/backend/src/bridge/session-bridge.ts index 4fc2493c..1e65926a 100644 --- a/backend/src/bridge/session-bridge.ts +++ b/backend/src/bridge/session-bridge.ts @@ -7,11 +7,20 @@ export type ToolActivity = { status: "running" | "completed" | "failed"; }; +export type TokenUsage = { + input: number; + cacheRead: number; + output: number; + totalTokens: number; + contextWindow: number; +}; + export type ClientEvent = | { type: "ui_request"; ui_request: any } | { type: "text_delta"; text: string } | { type: "activity_delta"; text: string } | { type: "activity_event"; activity: ToolActivity } + | { type: "usage"; usage: TokenUsage } | { type: "info"; [k: string]: any } | { type: "system_event"; event: string }; @@ -25,6 +34,7 @@ export class SessionBridge { // del descriptor (che viaggia opaco nel `title`). Va memorizzato e rimandato indietro, // altrimenti Pi scarta la risposta e ctx.ui.input non si risolve mai (stuck senza output). private pendingPiId: string | null = null; + private contextWindow = 0; private cbs = new Set<(e: ClientEvent) => void>(); constructor(private rpc: RpcClient) { @@ -36,17 +46,16 @@ export class SessionBridge { this.pendingPiId = m.id as string | null; this.state = "waiting"; this.fan({ type: "ui_request", ui_request: descriptor }); - } else if ( - m.type === "message_end" && - m.message?.role === "assistant" && - m.message?.stopReason === "error" - ) { - this.markFailed(); - this.fan({ - type: "info", - level: "error", - text: "Model request failed. Check provider connectivity, then Resume the session.", - }); + } else if (m.type === "message_end" && m.message?.role === "assistant") { + this.emitUsage(m.message.usage); + if (m.message.stopReason === "error") { + this.markFailed(); + this.fan({ + type: "info", + level: "error", + text: "Model request failed. Check provider connectivity, then Resume the session.", + }); + } } else if (m.type === "extension_ui_request" && m.method === "notify") { this.fan({ type: "info", level: m.notifyType ?? "info", text: m.message ?? "" }); } else if (m.type === "message_update" && m.assistantMessageEvent?.type === "text_delta") { @@ -93,12 +102,38 @@ export class SessionBridge { }); } + private emitUsage(usage: any): void { + if (!usage || this.contextWindow <= 0) return; + const number = (value: unknown): number => + typeof value === "number" && Number.isFinite(value) && value >= 0 ? value : 0; + const input = number(usage.input); + const cacheRead = number(usage.cacheRead); + const output = number(usage.output); + const reportedTotal = number(usage.totalTokens); + this.fan({ + type: "usage", + usage: { + input, + cacheRead, + output, + totalTokens: reportedTotal || input + cacheRead + output, + contextWindow: this.contextWindow, + }, + }); + } + private fan(e: ClientEvent) { for (const cb of this.cbs) cb(e); } turnState(): TurnState { return this.state; } beginTurn(): void { this.state = "running"; } + setContextWindow(value: unknown): void { + if (typeof value === "number" && Number.isFinite(value) && value > 0) { + this.contextWindow = value; + } + } + /** Record a backend-detected failure; callers own any sanitized client message. */ markFailed(): void { this.state = "failed"; } diff --git a/backend/src/pi/pi-process-manager.ts b/backend/src/pi/pi-process-manager.ts index aa3e49a2..9b626c30 100644 --- a/backend/src/pi/pi-process-manager.ts +++ b/backend/src/pi/pi-process-manager.ts @@ -166,7 +166,10 @@ export class PiProcessManager { const thinking = o.thinking ?? this.cfg.defaults.thinking; if (provider && model) { - await rt.rpc.request({ type: "set_model", provider, modelId: model } as object & { type: string }); + const response = await rt.rpc.request( + { type: "set_model", provider, modelId: model } as object & { type: string }, + ); + rt.bridge.setContextWindow(response?.data?.contextWindow); } if (thinking) { await rt.rpc.request({ type: "set_thinking_level", level: thinking } as object & { type: string }); diff --git a/backend/test/pi-process-manager.test.ts b/backend/test/pi-process-manager.test.ts index 7d898370..9a552292 100644 --- a/backend/test/pi-process-manager.test.ts +++ b/backend/test/pi-process-manager.test.ts @@ -207,6 +207,43 @@ test("createFor does not prompt until start is called", async () => { mgr.teardown("sid-deferred"); }); +test("configure gives the bridge the context window returned by set_model", async () => { + const child = recordingChild(); + child.stdin.write = (data: unknown) => { + const request = JSON.parse(String(data)); + child._writes.push(String(data)); + if (request.type === "set_model") { + queueMicrotask(() => child.stdout.emit("data", `${JSON.stringify({ + type: "response", + id: request.id, + success: true, + data: { contextWindow: 200_000 }, + })}\n`)); + } + return true; + }; + const mgr = new PiProcessManager(loadConfig({}), { spawnFn: () => child as any }); + const rt = mgr.createFor("context-window", {}); + const seen: any[] = []; + rt.bridge.onClientEvent((event) => seen.push(event)); + + await mgr.configure(rt, { provider: "zai", model: "glm-5.2" }); + child.stdout.emit("data", `${JSON.stringify({ + type: "message_end", + message: { + role: "assistant", + stopReason: "stop", + usage: { input: 10, cacheRead: 20, output: 5, totalTokens: 35 }, + }, + })}\n`); + + expect(seen).toContainEqual({ + type: "usage", + usage: { input: 10, cacheRead: 20, output: 5, totalTokens: 35, contextWindow: 200_000 }, + }); + mgr.teardown("context-window"); +}); + test("a created runtime is active during configure/bootstrap", () => { const child = recordingChild(); const mgr = new PiProcessManager(loadConfig({}), { spawnFn: () => child as any }); diff --git a/backend/test/session-bridge.test.ts b/backend/test/session-bridge.test.ts index e555ef91..fbfe1d50 100644 --- a/backend/test/session-bridge.test.ts +++ b/backend/test/session-bridge.test.ts @@ -42,6 +42,43 @@ test("real Pi thinking_delta becomes a dedicated activity_delta to the FE", () = expect(seen).toEqual([{ type: "activity_delta", text: "Valuto le ambiguità" }]); }); +test("assistant message_end exposes sanitized token usage with the configured context window", () => { + const { rpc, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + const seen: any[] = []; + bridge.onClientEvent((event) => seen.push(event)); + bridge.setContextWindow(200_000); + + fire({ + type: "message_end", + message: { + role: "assistant", + stopReason: "stop", + usage: { + input: 110_348, + cacheRead: 633_344, + cacheWrite: 12, + output: 33_652, + totalTokens: 777_356, + cost: { total: 99 }, + }, + content: "DO_NOT_FORWARD", + }, + }); + + expect(seen).toEqual([{ + type: "usage", + usage: { + input: 110_348, + cacheRead: 633_344, + output: 33_652, + totalTokens: 777_356, + contextWindow: 200_000, + }, + }]); + expect(JSON.stringify(seen)).not.toMatch(/DO_NOT_FORWARD|cost|cacheWrite/); +}); + test("assistant provider errors are sanitized and leave the turn failed", () => { const { rpc, fire } = fakeRpc(); const bridge = new SessionBridge(rpc); diff --git a/frontend/src/api/types.ts b/frontend/src/api/types.ts index 9901e7d8..4172cc9e 100644 --- a/frontend/src/api/types.ts +++ b/frontend/src/api/types.ts @@ -75,6 +75,14 @@ export interface ActivityEntry { status?: "running" | "completed" | "failed"; } +export interface TokenUsage { + input: number; + cacheRead: number; + output: number; + totalTokens: number; + contextWindow: number; +} + export type StreamEvent = | { type: "ui_request"; ui_request: WidgetDescriptor } | { type: "text_delta"; text: string } @@ -88,6 +96,7 @@ export type StreamEvent = status: "running" | "completed" | "failed"; }; } + | { type: "usage"; usage: TokenUsage } | { type: "info"; level?: "info" | "warning" | "error"; text: string } | { type: "system_event"; event: string }; diff --git a/frontend/src/shell/SteerInput.test.tsx b/frontend/src/shell/SteerInput.test.tsx index 4f8033be..efa2d656 100644 --- a/frontend/src/shell/SteerInput.test.tsx +++ b/frontend/src/shell/SteerInput.test.tsx @@ -2,8 +2,10 @@ import { render, screen, waitFor } from "@testing-library/react"; import userEvent from "@testing-library/user-event"; import { http, HttpResponse } from "msw"; +import { QueryClient, QueryClientProvider } from "@tanstack/react-query"; import { server } from "../test/msw"; -import { SteerInput } from "./SteerInput"; +import { useSessionStore } from "../store/sessionStore"; +import { ComposerFooter, ContextGauge, SteerInput } from "./SteerInput"; beforeEach(() => { server.use( @@ -83,6 +85,52 @@ test("marks the composer as awaiting input only when requested", () => { expect(input).toHaveClass("thot-awaiting-input"); }); +test("the context gauge uses green, yellow, and red at the requested thresholds", () => { + const { rerender } = render(); + expect(screen.getByLabelText("Context usage 70%")).toHaveAttribute("data-level", "green"); + + rerender(); + expect(screen.getByLabelText("Context usage 70%")).toHaveAttribute("data-level", "yellow"); + + rerender(); + expect(screen.getByLabelText("Context usage 85%")).toHaveAttribute("data-level", "yellow"); + + rerender(); + expect(screen.getByLabelText("Context usage 85%")).toHaveAttribute("data-level", "red"); +}); + +test("footer shows cumulative k-token counters after workspace and context gauge after thinking", async () => { + server.use( + http.get("http://localhost:8787/settings", () => HttpResponse.json({ + workspace: "psd", provider: "zai", model: "glm-5.2", thinking: "medium", + })), + http.get("http://localhost:8787/workspaces", () => HttpResponse.json([{ name: "psd" }])), + http.get("http://localhost:8787/models", () => HttpResponse.json({ + models: [{ provider: "zai", id: "glm-5.2", name: "GLM-5.2", reasoning: true }], + })), + ); + useSessionStore.getState().resetSession(); + useSessionStore.getState().applyEvent({ + type: "usage", + usage: { + input: 110_348, + cacheRead: 633_344, + output: 33_652, + totalTokens: 55_935, + contextWindow: 200_000, + }, + }); + const client = new QueryClient({ defaultOptions: { queries: { retry: false } } }); + render(); + + const counters = screen.getByText("110k/633k/34k"); + const workspace = await screen.findByRole("combobox", { name: "Workspace" }); + const gauge = screen.getByLabelText("Context usage 28%"); + const thinking = await screen.findByRole("combobox", { name: "Thinking level" }); + expect(workspace.compareDocumentPosition(counters) & Node.DOCUMENT_POSITION_FOLLOWING).toBeTruthy(); + expect(thinking.compareDocumentPosition(gauge) & Node.DOCUMENT_POSITION_FOLLOWING).toBeTruthy(); +}); + test("pulses the stop dot only while the harness is working", () => { const { rerender } = render(); const dot = () => diff --git a/frontend/src/shell/SteerInput.tsx b/frontend/src/shell/SteerInput.tsx index 7909c971..f07d59d0 100644 --- a/frontend/src/shell/SteerInput.tsx +++ b/frontend/src/shell/SteerInput.tsx @@ -151,9 +151,8 @@ export function SteerInput({ /** * Status strip beneath the composer. Left: the active workspace selector. - * Right: live model + thinking-level selectors (these replace the old Settings - * dialog and persist through PUT /settings) plus a context-usage gauge - * (placeholder until the backend reports token usage). + * Right: live context usage plus model + thinking-level selectors (these replace + * the old Settings dialog and persist through PUT /settings). */ export function ComposerFooter() { const qc = useQueryClient(); @@ -161,6 +160,7 @@ export function ComposerFooter() { const { data: workspaces = [] } = useQuery({ queryKey: ["workspaces"], queryFn: listWorkspaces }); const { data: modelsData } = useQuery({ queryKey: ["models"], queryFn: listModels }); const models = modelsData?.models ?? []; + const tokenUsage = useSessionStore((state) => state.tokenUsage); const workspace = settings?.workspace ?? ""; const model = settings?.model ?? ""; @@ -178,20 +178,31 @@ export function ComposerFooter() { } const knownModel = models.some((m) => m.id === model); + const contextPct = tokenUsage && tokenUsage.contextWindow > 0 + ? tokenUsage.totalTokens / tokenUsage.contextWindow + : 0; return (
- update({ workspace: v })}> - {workspaces.length === 0 ? ( - - ) : ( - workspaces.map((w) => ( - - )) - )} - +
+ update({ workspace: v })}> + {workspaces.length === 0 ? ( + + ) : ( + workspaces.map((w) => ( + + )) + )} + + + {formatTokenUsage(tokenUsage)} + +
@@ -215,7 +226,7 @@ export function ComposerFooter() { ))} - +
); @@ -249,14 +260,36 @@ function FooterSelect({ ); } +function formatThousands(tokens: number): string { + return `${Math.round(Math.max(tokens, 0) / 1_000)}k`; +} + +export function formatTokenUsage(usage: { + input: number; + cacheRead: number; + output: number; +} | null): string { + return [usage?.input ?? 0, usage?.cacheRead ?? 0, usage?.output ?? 0] + .map(formatThousands) + .join("/"); +} + /** Circular context-usage gauge. `pct` in [0,1]; 0 renders an empty track. */ -function ContextGauge({ pct = 0 }: { pct?: number }) { +export function ContextGauge({ pct = 0 }: { pct?: number }) { const r = 6; const circ = 2 * Math.PI * r; + const clamped = Math.min(Math.max(pct, 0), 1); + const level = clamped < 0.7 ? "green" : clamped <= 0.85 ? "yellow" : "red"; + const progressClass = level === "green" + ? "stroke-emerald-500" + : level === "yellow" + ? "stroke-amber-400" + : "stroke-red-500"; return ( @@ -266,11 +299,11 @@ function ContextGauge({ pct = 0 }: { pct?: number }) { cy="8" r={r} fill="none" - stroke="oklch(var(--primary))" + className={progressClass} strokeWidth="2.25" strokeLinecap="round" strokeDasharray={circ} - strokeDashoffset={circ * (1 - Math.min(Math.max(pct, 0), 1))} + strokeDashoffset={circ * (1 - clamped)} transform="rotate(-90 8 8)" /> diff --git a/frontend/src/store/sessionStore.test.ts b/frontend/src/store/sessionStore.test.ts index 091c0b79..b8dcb5d7 100644 --- a/frontend/src/store/sessionStore.test.ts +++ b/frontend/src/store/sessionStore.test.ts @@ -38,6 +38,26 @@ test("text_delta accumulates into transcript", () => { expect(useSessionStore.getState().transcript.at(-1)?.text).toBe("Analisi"); }); +test("usage events accumulate billed token categories but keep the latest context occupancy", () => { + const store = useSessionStore.getState(); + store.applyEvent({ + type: "usage", + usage: { input: 60_000, cacheRead: 400_000, output: 20_000, totalTokens: 80_000, contextWindow: 200_000 }, + }); + store.applyEvent({ + type: "usage", + usage: { input: 50_348, cacheRead: 233_344, output: 13_652, totalTokens: 55_935, contextWindow: 200_000 }, + }); + + expect(useSessionStore.getState().tokenUsage).toEqual({ + input: 110_348, + cacheRead: 633_344, + output: 33_652, + totalTokens: 55_935, + contextWindow: 200_000, + }); +}); + test("folds prompt, lifecycle, streams, tool updates, and gates in chronological order", () => { const store = useSessionStore.getState(); store.setPhase("F1"); diff --git a/frontend/src/store/sessionStore.ts b/frontend/src/store/sessionStore.ts index 3db2317d..d2fcbe1c 100644 --- a/frontend/src/store/sessionStore.ts +++ b/frontend/src/store/sessionStore.ts @@ -1,5 +1,5 @@ import { create } from "zustand"; -import type { ActivityEntry, StreamEvent, WidgetDescriptor } from "../api/types"; +import type { ActivityEntry, StreamEvent, TokenUsage, WidgetDescriptor } from "../api/types"; interface Entry { role: "assistant"; @@ -21,6 +21,7 @@ interface SessionState { // system event — the only end-of-turn signal on the final workflow step, where no // follow-up gate arrives to release the spinner. agentActive: boolean; + tokenUsage: TokenUsage | null; applyEvent: (e: StreamEvent) => void; clearPending: () => void; resetSession: () => void; @@ -65,6 +66,7 @@ const empty = { phaseError: null as string | null, seenGateIds: new Set(), agentActive: false, + tokenUsage: null as TokenUsage | null, }; export const useSessionStore = create((set) => ({ @@ -145,6 +147,17 @@ export const useSessionStore = create((set) => ({ } return { activityLog, agentActive: true }; } + if (e.type === "usage") { + return { + tokenUsage: { + input: (st.tokenUsage?.input ?? 0) + e.usage.input, + cacheRead: (st.tokenUsage?.cacheRead ?? 0) + e.usage.cacheRead, + output: (st.tokenUsage?.output ?? 0) + e.usage.output, + totalTokens: e.usage.totalTokens, + contextWindow: e.usage.contextWindow, + }, + }; + } if (e.type === "info") { const level = e.level ?? "info"; const stepMessages = [...st.stepMessages, { level, text: e.text }]; diff --git a/frontend/src/stream/useSessionStream.test.tsx b/frontend/src/stream/useSessionStream.test.tsx index 58363e27..4e669f36 100644 --- a/frontend/src/stream/useSessionStream.test.tsx +++ b/frontend/src/stream/useSessionStream.test.tsx @@ -1,7 +1,7 @@ import { renderHook } from "@testing-library/react"; import { act, useLayoutEffect } from "react"; import { FakeEventSource } from "../test/fakeEventSource"; -import { useSessionStream } from "./useSessionStream"; +import { createStreamEventCoalescer, useSessionStream } from "./useSessionStream"; import { useSessionStore } from "../store/sessionStore"; beforeEach(() => { @@ -10,6 +10,49 @@ beforeEach(() => { useSessionStore.getState().resetSession(); }); +test("coalesces a burst of text deltas into bounded store updates", () => { + vi.useFakeTimers(); + try { + const applied: Array<{ type: string; text?: string }> = []; + const coalescer = createStreamEventCoalescer((event) => applied.push(event), 100); + + for (let i = 0; i < 100; i += 1) { + coalescer.push({ type: "text_delta", text: "x" }); + } + + expect(applied).toEqual([{ type: "text_delta", text: "x" }]); + vi.advanceTimersByTime(100); + expect(applied).toEqual([ + { type: "text_delta", text: "x" }, + { type: "text_delta", text: "x".repeat(99) }, + ]); + coalescer.dispose(); + } finally { + vi.useRealTimers(); + } +}); + +test("flushes pending stream text before a structural event", () => { + vi.useFakeTimers(); + try { + const applied: Array<{ type: string; text?: string }> = []; + const coalescer = createStreamEventCoalescer((event) => applied.push(event), 100); + + coalescer.push({ type: "text_delta", text: "a" }); + coalescer.push({ type: "text_delta", text: "b" }); + coalescer.push({ type: "system_event", event: "turn_end" }); + + expect(applied).toEqual([ + { type: "text_delta", text: "a" }, + { type: "text_delta", text: "b" }, + { type: "system_event", event: "turn_end" }, + ]); + coalescer.dispose(); + } finally { + vi.useRealTimers(); + } +}); + test("opens an EventSource and feeds NAMED events to the store", () => { renderHook(() => useSessionStream("s1")); const es = FakeEventSource.instances[0]; @@ -56,6 +99,26 @@ test("feeds named activity_event events to the tool activity timeline", () => { }); }); +test("feeds named usage events to the session token counters", () => { + renderHook(() => useSessionStream("s1")); + const es = FakeEventSource.instances[0]; + + act(() => + es.emitNamed("usage", { + type: "usage", + usage: { input: 1_000, cacheRead: 2_000, output: 300, totalTokens: 3_300, contextWindow: 200_000 }, + }) + ); + + expect(useSessionStore.getState().tokenUsage).toEqual({ + input: 1_000, + cacheRead: 2_000, + output: 300, + totalTokens: 3_300, + contextWindow: 200_000, + }); +}); + test("also feeds UNNAMED (default message) events to the store", () => { renderHook(() => useSessionStream("s1")); const es = FakeEventSource.instances[0]; diff --git a/frontend/src/stream/useSessionStream.ts b/frontend/src/stream/useSessionStream.ts index ca6f7d63..62f28ccd 100644 --- a/frontend/src/stream/useSessionStream.ts +++ b/frontend/src/stream/useSessionStream.ts @@ -4,6 +4,68 @@ import { joinBackendPath } from "../api/runtime-config"; import { useSessionStore } from "../store/sessionStore"; import type { StreamEvent } from "../api/types"; +const STREAM_UPDATE_INTERVAL_MS = 100; + +function isStreamingDelta( + event: StreamEvent, +): event is Extract { + return event.type === "text_delta" || event.type === "activity_delta"; +} + +export function createStreamEventCoalescer( + applyEvent: (event: StreamEvent) => void, + intervalMs = STREAM_UPDATE_INTERVAL_MS, +) { + let pending: Extract | null = null; + let timer: ReturnType | null = null; + let disposed = false; + + const flushPending = () => { + if (!pending) return; + const event = pending; + pending = null; + applyEvent(event); + }; + + const schedule = () => { + timer = setTimeout(() => { + timer = null; + if (!pending || disposed) return; + flushPending(); + schedule(); + }, intervalMs); + }; + + return { + push(event: StreamEvent) { + if (disposed) return; + if (!isStreamingDelta(event)) { + flushPending(); + if (timer) clearTimeout(timer); + timer = null; + applyEvent(event); + return; + } + if (!timer) { + applyEvent(event); + schedule(); + return; + } + if (pending && pending.type !== event.type) flushPending(); + pending = pending + ? { ...event, text: pending.text + event.text } + : event; + }, + dispose() { + if (disposed) return; + if (timer) clearTimeout(timer); + timer = null; + flushPending(); + disposed = true; + }, + }; +} + export function useSessionStream( sessionId: string | null, generation = 0, @@ -36,6 +98,7 @@ export function useSessionStream( : ""; const es = new EventSource(joinBackendPath(BASE, `/sessions/${sessionId}/events${query}`)); const identity = { source: es, sessionId, cursorResetEpoch }; + const coalescer = createStreamEventCoalescer(applyEvent); activeSource.current = identity; es.onopen = () => { if (activeSource.current === identity) setConnected(true); @@ -48,7 +111,7 @@ export function useSessionStream( if (activeSource.current !== identity) return; if (ev.lastEventId) cursor.current.lastEventId = ev.lastEventId; try { - applyEvent(JSON.parse(ev.data) as StreamEvent); + coalescer.push(JSON.parse(ev.data) as StreamEvent); } catch { /* ignore malformed */ } @@ -62,6 +125,7 @@ export function useSessionStream( "text_delta", "activity_delta", "activity_event", + "usage", "info", "system_event", ] as const; @@ -71,6 +135,7 @@ export function useSessionStream( for (const name of namedEvents) es.removeEventListener(name, handle); es.onmessage = null; es.close(); + coalescer.dispose(); if (activeSource.current === identity) { activeSource.current = null; setConnected(false); diff --git a/frontend/src/viewers/ArtifactView.test.tsx b/frontend/src/viewers/ArtifactView.test.tsx index 0f34be7a..405d4c36 100644 --- a/frontend/src/viewers/ArtifactView.test.tsx +++ b/frontend/src/viewers/ArtifactView.test.tsx @@ -21,6 +21,11 @@ test("sql artifact accepts a { content } wrapper", async () => { await screen.findByTestId("hl"); }); +test("an empty artifact shows an explicit fallback instead of a blank panel", () => { + render(); + expect(screen.getByText("No artifact content available.")).toBeInTheDocument(); +}); + test("cte_plan renders an ordered list of names", () => { render(); expect(screen.getByText("a_cte")).toBeInTheDocument(); diff --git a/frontend/src/viewers/ArtifactView.tsx b/frontend/src/viewers/ArtifactView.tsx index 5f6efa41..f38cedb7 100644 --- a/frontend/src/viewers/ArtifactView.tsx +++ b/frontend/src/viewers/ArtifactView.tsx @@ -135,6 +135,9 @@ function StructuredValue({ value, depth = 0 }: { value: unknown; depth?: number ); } const entries = Object.entries(value as Record); + if (entries.length === 0) { + return

No artifact content available.

; + } return (
{entries.map(([k, v], i) => ( diff --git a/harness/.pi/extensions/gate/__tests__/gate_agent_end_safety.test.js b/harness/.pi/extensions/gate/__tests__/gate_agent_end_safety.test.js new file mode 100644 index 00000000..6d5d49d6 --- /dev/null +++ b/harness/.pi/extensions/gate/__tests__/gate_agent_end_safety.test.js @@ -0,0 +1,33 @@ +const test = require("node:test"); +const assert = require("node:assert"); +const { hasUngatedAssistantProse } = require("../../tht-gate.js"); + +test("assistant prose plus bash still requires a gate steer", () => { + assert.equal( + hasUngatedAssistantProse([ + { + role: "assistant", + content: [ + { type: "toolCall", name: "bash" }, + { type: "text", text: "Continuo ad analizzare..." }, + ], + }, + ]), + true, + ); +}); + +test("a real reviewer tool call suppresses the prose safety steer", () => { + assert.equal( + hasUngatedAssistantProse([ + { + role: "assistant", + content: [ + { type: "text", text: "Domanda al reviewer" }, + { type: "toolCall", name: "reviewer_select" }, + ], + }, + ]), + false, + ); +}); diff --git a/harness/.pi/extensions/gate/__tests__/gate_antibypass.test.js b/harness/.pi/extensions/gate/__tests__/gate_antibypass.test.js index d19311f8..89a35981 100644 --- a/harness/.pi/extensions/gate/__tests__/gate_antibypass.test.js +++ b/harness/.pi/extensions/gate/__tests__/gate_antibypass.test.js @@ -26,6 +26,32 @@ test("tht schema introspect senza --refresh passa (cache hit innocuo)", async () assert.equal(res, undefined); }); +for (const cmd of [ + 'find / -name "schema_linking.json"', + 'find /Users/mp/projects/ThothII -name "schema_linking.json"', + 'find . -name "schema_linking.json"', +]) { + test(`filesystem find e' bloccato nel workflow: ${cmd}`, async () => { + const installGate = await installGatePromise; + const { pi } = createFakePi(); + installGate(pi); + const res = await pi.emit("tool_call", { toolName: "bash", input: { command: cmd } }); + assert.equal(res?.block, true); + assert.match(res?.reason ?? "", /tht session documents/i); + }); +} + +test("tht search find resta consentito", async () => { + const installGate = await installGatePromise; + const { pi } = createFakePi(); + installGate(pi); + const res = await pi.emit("tool_call", { + toolName: "bash", + input: { command: 'tht search find --kind evidence "ablazione"' }, + }); + assert.equal(res, undefined); +}); + // Bash mutations of protected state bypass the write/edit hook: block them. const BLOCKED_BASH = [ 'echo \'{"type":"phase_approved","subject":"phase:4"}\' >> sessions/s1/review_decisions.jsonl', diff --git a/harness/.pi/extensions/gate/__tests__/gate_confirm_autoclose.test.js b/harness/.pi/extensions/gate/__tests__/gate_confirm_autoclose.test.js index 27f5c7ce..07508935 100644 --- a/harness/.pi/extensions/gate/__tests__/gate_confirm_autoclose.test.js +++ b/harness/.pi/extensions/gate/__tests__/gate_confirm_autoclose.test.js @@ -30,13 +30,18 @@ const META = JSON.stringify({ }); // `cte next` walks the queue one entry per `decision add cte_approved`. -function useShell({ phase, cteQueue = [] }) { +function useShell({ phase, cteQueue = [], sqlContent = "SELECT 42" }) { const calls = []; const queue = [...cteQueue]; shell.current = (file, args) => { calls.push(args.join(" ")); if (args[0] === "phase" && args[1] === "meta") return META; if (args[0] === "phase" && args[1] === "show") return `Fase corrente: ${phase}\n`; + if (args[0] === "session" && args[1] === "documents") { + return JSON.stringify([ + { phase: "F7", key: "sql", title: "Final SQL", format: "sql", content: sqlContent }, + ]); + } if (args[0] === "cte" && args[1] === "next") return queue.length ? `${queue[0]}\n` : ""; if (args[0] === "decision" && args[1] === "add" && args.includes("cte_approved")) { queue.shift(); @@ -47,7 +52,7 @@ function useShell({ phase, cteQueue = [] }) { return calls; } -async function runConfirm(kind, artifactKind) { +async function runConfirm(kind, artifactKind, { artifactData = "# md", onDescriptor } = {}) { const gate = require(GATE); const { createFakePi } = require("./fake_pi_runtime.js"); const { pi, ctx, tools } = createFakePi(); @@ -57,11 +62,12 @@ async function runConfirm(kind, artifactKind) { await pi.emit("session_start", {}); ctx.ui.input = async (title) => { const d = JSON.parse(title); + onDescriptor?.(d); return JSON.stringify({ id: d.id, kind: "artifact-gate", choices: ["approve"] }); }; return tools.get("reviewer_confirm").def.execute( "call-1", - { session: "s1", kind, title: "t", artifact: { kind: artifactKind, data: "# md" } }, + { session: "s1", kind, title: "t", artifact: { kind: artifactKind, data: artifactData } }, null, null, ctx, @@ -104,3 +110,30 @@ test("approving the SQL closes F7 automatically (no separate phase gate)", async ); assert.ok(!calls.some((c) => c.startsWith("session finalize")), "F7 must not finalize"); }); + +test("SQL approval hydrates the gate from persisted sql_final instead of trusting empty model data", async () => { + useShell({ phase: 7, sqlContent: "SELECT persisted_sql" }); + let descriptor; + + await runConfirm("sql", "sql", { + artifactData: {}, + onDescriptor: (value) => { descriptor = value; }, + }); + + assert.equal(descriptor.artifact.kind, "sql"); + assert.equal(descriptor.artifact.data, "SELECT persisted_sql"); +}); + +test("SQL approval refuses to show a gate when persisted sql_final is empty", async () => { + const calls = useShell({ phase: 7, sqlContent: "" }); + let shown = false; + + const result = await runConfirm("sql", "sql", { + artifactData: {}, + onDescriptor: () => { shown = true; }, + }); + + assert.equal(shown, false); + assert.match(result.content[0].text, /sql_final.*vuoto/i); + assert.ok(!calls.some((c) => c.includes("decision add"))); +}); diff --git a/harness/.pi/extensions/gate/__tests__/gate_resume_kickoff.test.js b/harness/.pi/extensions/gate/__tests__/gate_resume_kickoff.test.js index 78ee091f..cac6f3ff 100644 --- a/harness/.pi/extensions/gate/__tests__/gate_resume_kickoff.test.js +++ b/harness/.pi/extensions/gate/__tests__/gate_resume_kickoff.test.js @@ -21,6 +21,8 @@ test("the resume kickoff injects the bootstrap steps and forces in-turn action", assert.match(text, //); assert.match(text, /# Thoth session workflow \(phases 1-8\)/); assert.match(text, /non esplorare il repository/i); + assert.match(text, /tht session documents --json/); + assert.match(text, /non cercare i file fisici con `find`/i); assert.doesNotMatch(text, /Carica la skill leggendo/); }); diff --git a/harness/.pi/extensions/gate/__tests__/gate_select_decision.test.js b/harness/.pi/extensions/gate/__tests__/gate_select_decision.test.js index 140e2bdf..a378faa4 100644 --- a/harness/.pi/extensions/gate/__tests__/gate_select_decision.test.js +++ b/harness/.pi/extensions/gate/__tests__/gate_select_decision.test.js @@ -1,6 +1,10 @@ const test = require("node:test"); const assert = require("node:assert"); -const { resolveSelectOutcome, decisionAddArgs } = require("../../tht-gate.js"); +const { + resolveSelectOutcome, + decisionAddArgs, + decisionRecordedResultText, +} = require("../../tht-gate.js"); // Workstream F — single-select answers auto-confirm. A concrete reviewer_select choice // that carries a `decision` payload IS the confirmation: the gate persists it directly @@ -46,3 +50,15 @@ test("decisionAddArgs omits optional detail/rationale when absent", () => { ["decision", "add", "--session", "s1", "--type", "t", "--subject", "sub"], ); }); + +test("a persisted select decision immediately directs the model to the next gate tool", () => { + const text = decisionRecordedResultText( + { type: "concept_clarified" }, + { label: "Procedure invasive" }, + ); + + assert.match(text, /Decisione registrata \(concept_clarified\): Procedure invasive\./); + assert.match(text, /tht session show/); + assert.match(text, /prossimo tool reviewer_/); + assert.match(text, /Non scrivere analisi o spiegazioni visibili/); +}); diff --git a/harness/.pi/extensions/tht-gate.js b/harness/.pi/extensions/tht-gate.js index 05817c9f..9a5dcb1b 100644 --- a/harness/.pi/extensions/tht-gate.js +++ b/harness/.pi/extensions/tht-gate.js @@ -165,6 +165,14 @@ export function isProtectedBashMutation(cmd) { return PROTECTED_PATH_TOKEN.test(cmd) && BASH_MUTATION.test(cmd); } +// Session artifacts are exposed by deterministic `tht session documents`; shell +// discovery is both unnecessary and dangerous (a model once escalated to `find /` +// and wedged the whole RPC turn). Match shell command boundaries so the supported +// `tht search find` subcommand remains available. +export function isFilesystemFind(cmd) { + return /(?:^|[;&|\n])\s*(?:(?:sudo|command)\s+)?find(?:\s|$)/.test(cmd); +} + // --- kickoff payloads (verbatim from source L184-212, load-bearing model prose) - const NUOVA_DOMANDA_KICKOFF = "AZIONE IMMEDIATA OBBLIGATORIA: la skill canonica è già inclusa integralmente nel " + @@ -480,6 +488,35 @@ export function decisionAddArgs(session, d) { return args; } +// Keep the model inside the gate workflow immediately after a reviewer selection. +// The agent_end prose safety net is too late for providers that keep streaming a +// single, very long assistant turn: put the continuation contract in the tool result +// that unlocks the model after the reviewer response. +export function decisionRecordedResultText(decision, option) { + return ( + `Decisione registrata (${decision.type}): ${option.label}. ` + + "Ora rileggi lo stato persistito con `tht session show` e invoca immediatamente " + + "il prossimo tool reviewer_ richiesto dal workflow. " + + "Non scrivere analisi o spiegazioni visibili." + ); +} + +export function hasUngatedAssistantProse(messages) { + const last = messages?.[messages.length - 1]; + if (!last || last.role !== "assistant" || !Array.isArray(last.content)) + return false; + const hasText = last.content.some( + (block) => block.type === "text" && (block.text ?? "").trim().length > 0, + ); + const hasReviewerTool = last.content.some( + (block) => + block.type === "toolCall" && + typeof block.name === "string" && + block.name.startsWith("reviewer_"), + ); + return hasText && !hasReviewerTool; +} + // Workstream F: classifies a reviewer_select response into the action the gate takes. // A concrete choice that carries a `decision` payload auto-confirms (persist directly, // no second gate); a bare choice stays ask-only; back/exit/Other never persist. @@ -618,6 +655,15 @@ export default function (pi) { lastSteered = false; return; } + if (isFilesystemFind(cmd)) { + return { + block: true, + reason: + "Non cercare file con `find`: gli artefatti persistiti della sessione " + + "si leggono con `tht session documents --json`; per catalogo ed " + + "evidence usa `tht schema render` / `tht search find`.", + }; + } if (FORBIDDEN.some((re) => re.test(cmd))) { return { block: true, @@ -739,17 +785,7 @@ export default function (pi) { // tool call (and the lock is active), nudge it back to the gate tools. pi.on("agent_end", async (event) => { if (!lockActive) return; - const msgs = event.messages ?? []; - const last = msgs[msgs.length - 1]; - const isProse = - last && - last.role === "assistant" && - Array.isArray(last.content) && - last.content.some( - (b) => b.type === "text" && (b.text ?? "").trim().length > 0, - ) && - !last.content.some((b) => b.type === "toolCall"); - if (isProse && !lastSteered) { + if (hasUngatedAssistantProse(event.messages ?? []) && !lastSteered) { lastSteered = true; await pi.sendUserMessage( "Le risposte del reviewer arrivano solo dai widget del gate. Riproponi la " + @@ -832,7 +868,7 @@ export default function (pi) { if (err) return err; if (advance) advanceIfReady(ctx, session); return textResult( - `Decisione registrata (${outcome.decision.type}): ${outcome.option.label}.`, + decisionRecordedResultText(outcome.decision, outcome.option), ); } return textResult( @@ -1159,6 +1195,36 @@ export default function (pi) { const curNum = currentPhase(ctx, session); const phase = phaseId(ctx, curNum); + // F7's review target is the persisted artifact, never the model-authored + // descriptor payload. A model can (and did) call reviewer_confirm with + // `artifact.data: {}` even though write_final_sql had already persisted a + // non-empty sql_final.sql; trusting that payload rendered a blank approval + // dialog. Hydrate from the repository boundary and fail closed before any + // widget is shown if the artifact is unavailable or empty. + if (kind === "sql") { + let docs; + try { + docs = JSON.parse(tht(ctx, ["session", "documents", session, "--json"])); + } catch (e) { + const msg = (e.stderr || e.message || String(e)).toString().trim(); + return textResult( + `Impossibile leggere sql_final.sql dalla sessione ${session}: ${msg}. ` + + "Verifica l'artefatto persistito e riprova.", + ); + } + const sqlDoc = Array.isArray(docs) + ? docs.find((doc) => doc && doc.key === "sql" && doc.format === "sql") + : null; + const sql = typeof sqlDoc?.content === "string" ? sqlDoc.content : ""; + if (!sql.trim()) { + return textResult( + `sql_final.sql assente o vuoto per la sessione ${session}: ` + + "scrivi un SQL finale valido prima di ripresentare il gate F7.", + ); + } + artifact = { ...artifact, kind: "sql", data: sql }; + } + // Sessione gia' finalizzata: non c'e' piu' nulla da approvare — mai // ripresentare il phase-gate (la form "approve") a workflow chiuso. if (kind === "phase" && curNum >= phaseMeta(ctx).max_phase) { diff --git a/harness/.pi/skills/tht-sessione/SKILL.md b/harness/.pi/skills/tht-sessione/SKILL.md index 38463e48..167e9f05 100644 --- a/harness/.pi/skills/tht-sessione/SKILL.md +++ b/harness/.pi/skills/tht-sessione/SKILL.md @@ -151,9 +151,13 @@ persisted state is your only context. Bootstrap before doing anything else: 1. `tht session show --json` → read `phase` (the current phase N), `status`, and the manifest (`question`, `database`, `schema`). -2. Load the artifacts produced so far, as needed for phase N: `question.md` (revised - question), `schema_linking.json` (F4 output), `ctes/*.sql` + `cte_tests.json` (F6), +2. Load the artifacts produced so far with `tht session documents --json`, which + returns their keys and contents directly: `question.md` (revised question), + `schema_linking.json` (F4 output), `ctes/*.sql` + `cte_tests.json` (F6), `sql_final.sql` (F7). The decision ledger is summarized by `tht session show`. + **Non cercare i file fisici con `find`, `ls`, `cat` o il generic read tool**: the + session repository may live outside the checkout and the documents command is the + canonical read boundary. 3. **Resume at phase N reviewing the existing artifacts** (same discipline as rollback, §Disciplines 11). Do NOT restart from Phase 1, do NOT re-run `tht` commands for artifacts that already exist and are valid, and do NOT treat this as a new question.