From fa2298653b69c55433716ec8ed8c87578542fa8d Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 24 Aug 2026 12:09:24 +0200 Subject: [PATCH] fix(ui): stream phase progress without gates --- PROJECT_STATE.md | 19 ++++++ backend/src/bridge/session-bridge.ts | 17 +++++- backend/test/session-bridge.test.ts | 23 ++++++++ .../contracts/workflow-observable-baseline.md | 3 + frontend/src/api/types.ts | 2 +- frontend/src/shell/WorkflowBar.test.tsx | 12 ++++ frontend/src/store/sessionStore.test.ts | 13 ++++ frontend/src/store/sessionStore.ts | 14 ++++- .../gate_phase_observability.test.js | 59 +++++++++++++++++++ harness/.pi/extensions/tht-gate.js | 29 ++++++++- 10 files changed, 186 insertions(+), 5 deletions(-) create mode 100644 harness/.pi/extensions/gate/__tests__/gate_phase_observability.test.js diff --git a/PROJECT_STATE.md b/PROJECT_STATE.md index ff04c353..226d80b6 100644 --- a/PROJECT_STATE.md +++ b/PROJECT_STATE.md @@ -1,5 +1,24 @@ # ThothII — Project State +## Modular workflow refactor candidate — live (2026-08-24) + +- **Isolation:** worktree `.worktrees/refactoring-modulare-contract-baseline`, branch + `codex/refactoring-modulare-contract-baseline`; the base stack is unchanged. +- **Candidate runtime:** Compose project `thothii-d811484ae2e4` is healthy at frontend + `127.0.0.1:18080`, core `127.0.0.1:18787`, and Qdrant `127.0.0.1:16333`. Its explicit volumes + preserve the PSD sessions, schema index, and embedding cache across image rebuilds. +- **Real data:** the VPN-backed PSD DWH is reachable; preprocessing indexed 163 tables and 2275 + columns. Manual session `85a0b758-ae30-432e-8298-837bf55abad9` completed F1-F8 on the modular + image and finalized successfully. +- **Phase progress fix:** Pi now emits a deduplicated, machine-readable phase-start notification + after persisted workflow mutations; the backend sanitizes it to `phase_started`; the frontend + advances the active workflow dot without waiting for a human widget. The real Pi RPC probe + emitted F8, and the harness/backend/frontend contract suites plus both TypeScript builds pass. +- **Testing:** the automated probe session `cccbee8c-b13d-4bcb-88a6-aa60a4cae533` is closed at F8 + after repeated cold-Resume probes preserved its 29 decisions and SQL/CTE/schema artifacts, so it + does not occupy the admin principal. The candidate stack is intentionally left running for owner + validation at `http://127.0.0.1:18080/`. + > Starting-point snapshot for new sessions. > **Requisito finale del progetto (owner, 2026-08-11):** al termine dell'ultima fase tecnica deve > essere prodotto un documento unico che guidi l'utente passo-passo su (1) come preparare il diff --git a/backend/src/bridge/session-bridge.ts b/backend/src/bridge/session-bridge.ts index f1dbf2c2..69ea9559 100644 --- a/backend/src/bridge/session-bridge.ts +++ b/backend/src/bridge/session-bridge.ts @@ -4,6 +4,15 @@ const GENERIC_MODEL_FAILURE = "Model request failed. Check provider connectivity, then Resume the session."; const SUBSCRIPTION_MODEL_FAILURE = "The selected model is unavailable for the current subscription. Choose another model and start a new session."; +const PHASE_STARTED_NOTIFICATION_PREFIX = "__tht_phase_started__:"; + +function phaseStartedNotification(message: unknown): string | null { + if (typeof message !== "string" || !message.startsWith(PHASE_STARTED_NOTIFICATION_PREFIX)) { + return null; + } + const phase = message.slice(PHASE_STARTED_NOTIFICATION_PREFIX.length); + return /^F[1-8]$/.test(phase) ? phase : ""; +} function safeModelFailure(error: unknown): string { const detail = typeof error === "string" ? error : ""; @@ -36,7 +45,7 @@ export type ClientEvent = | { type: "activity_event"; activity: ToolActivity } | { type: "usage"; usage: TokenUsage } | { type: "info"; [k: string]: any } - | { type: "system_event"; event: string }; + | { type: "system_event"; event: string; phase?: string }; export type TurnState = "idle" | "running" | "waiting" | "failed"; @@ -71,7 +80,11 @@ export class SessionBridge { }); } } else if (m.type === "extension_ui_request" && m.method === "notify") { - this.fan({ type: "info", level: m.notifyType ?? "info", text: m.message ?? "" }); + const phase = phaseStartedNotification(m.message); + if (phase) this.fan({ type: "system_event", event: "phase_started", phase }); + else if (phase === null) { + this.fan({ type: "info", level: m.notifyType ?? "info", text: m.message ?? "" }); + } } else if (m.type === "message_update" && m.assistantMessageEvent?.type === "text_delta") { this.fan({ type: "text_delta", text: m.assistantMessageEvent.delta ?? "" }); } else if (m.type === "message_update" && m.assistantMessageEvent?.type === "thinking_delta") { diff --git a/backend/test/session-bridge.test.ts b/backend/test/session-bridge.test.ts index f0a41341..87637866 100644 --- a/backend/test/session-bridge.test.ts +++ b/backend/test/session-bridge.test.ts @@ -42,6 +42,29 @@ test("real Pi thinking_delta becomes a dedicated activity_delta to the FE", () = expect(seen).toEqual([{ type: "activity_delta", text: "Valuto le ambiguità" }]); }); +test("a reserved phase notification becomes a structured phase_started event", () => { + const { rpc, fire } = fakeRpc(); + const bridge = new SessionBridge(rpc); + const seen: any[] = []; + bridge.onClientEvent((event) => seen.push(event)); + + fire({ + type: "extension_ui_request", + method: "notify", + notifyType: "info", + message: "__tht_phase_started__:F2", + }); + fire({ + type: "extension_ui_request", + method: "notify", + notifyType: "info", + message: "__tht_phase_started__:F9_DO_NOT_FORWARD", + }); + + expect(seen).toEqual([{ type: "system_event", event: "phase_started", phase: "F2" }]); + expect(JSON.stringify(seen)).not.toContain("DO_NOT_FORWARD"); +}); + test("assistant message_end exposes sanitized token usage with the configured context window", () => { const { rpc, fire } = fakeRpc(); const bridge = new SessionBridge(rpc); diff --git a/docs/contracts/workflow-observable-baseline.md b/docs/contracts/workflow-observable-baseline.md index 72be860d..c201e64a 100644 --- a/docs/contracts/workflow-observable-baseline.md +++ b/docs/contracts/workflow-observable-baseline.md @@ -24,6 +24,7 @@ The gate suite fixes: - F3 rewritten question and assumptions, including mutation failure ordering; - F4 Evidence used, accepted, rejected, and legacy-without-corpus projections; - F8 Memory promotion accepted, declined, and absent, including mutation failure ordering; +- a newly folded phase is announced once in RPC mode even when it requires no human gate; - resume reconstruction for the touched F1, F2, F3, F4, and F8 states; - artifact payload compatibility, anti-bypass behavior, and final phase closing. @@ -73,6 +74,7 @@ These suites fix: - new-session versus resume Pi prompts; - refusal to resume finalized, archived, foreign, unavailable, or read-only sessions; - Pi RPC to client event mapping, SSE replay/reset behavior, and runtime replacement ordering; +- reserved Pi phase notifications mapped to sanitized `phase_started` client events; - failure persistence and sanitization before a client-visible response. ### Frontend client @@ -92,6 +94,7 @@ These suites fix: - widget registry and gate response payloads; - `ui_request`, `text_delta`, activity, usage, and lifecycle event reduction; +- phase progress and the active workflow dot advancing on `phase_started` without a `ui_request`; - stream replacement, cursor reset, reconnection, and pending-text flush behavior; - session document projections shown to the reviewer. diff --git a/frontend/src/api/types.ts b/frontend/src/api/types.ts index 61d1b4f1..15f1c53d 100644 --- a/frontend/src/api/types.ts +++ b/frontend/src/api/types.ts @@ -98,7 +98,7 @@ export type StreamEvent = } | { type: "usage"; usage: TokenUsage } | { type: "info"; level?: "info" | "warning" | "error"; text: string } - | { type: "system_event"; event: string }; + | { type: "system_event"; event: string; phase?: string }; export interface SessionSummary { id: string; diff --git a/frontend/src/shell/WorkflowBar.test.tsx b/frontend/src/shell/WorkflowBar.test.tsx index 59090edb..86b6b5b4 100644 --- a/frontend/src/shell/WorkflowBar.test.tsx +++ b/frontend/src/shell/WorkflowBar.test.tsx @@ -29,6 +29,18 @@ test("dots carry done / running / pending states around the active phase", () => expect(screen.getByTestId("phase-F5")).toHaveAttribute("data-state", "pending"); }); +test("phase_started colors a phase even when no human gate was rendered", () => { + const store = useSessionStore.getState(); + store.setPhase("F1"); + store.applyEvent({ type: "system_event", event: "phase_started", phase: "F2" }); + + render(); + + expect(screen.getByTestId("phase-F1")).toHaveAttribute("data-state", "done"); + expect(screen.getByTestId("phase-F2")).toHaveAttribute("data-state", "running"); + expect(screen.getByTestId("phase-F3")).toHaveAttribute("data-state", "pending"); +}); + test("the active dot turns to the error state when phaseError flags it", () => { const st = useSessionStore.getState(); st.applyEvent({ type: "ui_request", ui_request: { id: "u", widget: "select", phase: "F3_x" } }); diff --git a/frontend/src/store/sessionStore.test.ts b/frontend/src/store/sessionStore.test.ts index 4cf65eae..e073ea4b 100644 --- a/frontend/src/store/sessionStore.test.ts +++ b/frontend/src/store/sessionStore.test.ts @@ -184,6 +184,19 @@ test("ui_request without a phase keeps the existing currentPhase", () => { expect(useSessionStore.getState().currentPhase).toBe("F2"); }); +test("phase_started advances progress even when the phase has no human gate", () => { + const store = useSessionStore.getState(); + store.setPhase("F1"); + + store.applyEvent({ + type: "system_event", + event: "phase_started", + phase: "F2_memory", + }); + + expect(useSessionStore.getState().currentPhase).toBe("F2"); +}); + test("agentActive lifecycle: off by default, on with activity, off on agent_end", () => { const st = useSessionStore.getState(); expect(useSessionStore.getState().agentActive).toBe(false); diff --git a/frontend/src/store/sessionStore.ts b/frontend/src/store/sessionStore.ts index 926bff99..8222dbcf 100644 --- a/frontend/src/store/sessionStore.ts +++ b/frontend/src/store/sessionStore.ts @@ -176,14 +176,26 @@ export const useSessionStore = create((set) => ({ : { stepMessages, activityLog }; } if (e.type === "system_event") { + const currentPhase = e.event === "phase_started" + ? phaseOf(e.phase, st.currentPhase) + : st.currentPhase; const activityLog = [ ...st.activityLog, { kind: "lifecycle" as const, - phase: st.currentPhase, + phase: currentPhase, text: lifecycleText(e.event), }, ]; + if (e.event === "phase_started") { + return { + lastSystemEvent: e, + activityLog, + currentPhase, + phaseError: null, + agentActive: true, + }; + } return e.event === "agent_end" ? { lastSystemEvent: e, activityLog, agentActive: false } : { lastSystemEvent: e, activityLog }; diff --git a/harness/.pi/extensions/gate/__tests__/gate_phase_observability.test.js b/harness/.pi/extensions/gate/__tests__/gate_phase_observability.test.js new file mode 100644 index 00000000..82bd7afd --- /dev/null +++ b/harness/.pi/extensions/gate/__tests__/gate_phase_observability.test.js @@ -0,0 +1,59 @@ +const test = require("node:test"); +const assert = require("node:assert"); +const cp = require("node:child_process"); +const path = require("node:path"); +const { createRequire } = require("node:module"); +const { createFakePi } = require("./fake_pi_runtime.js"); + +const GATE = path.join(__dirname, "..", "..", "tht-gate.js"); +globalThis.require = createRequire(GATE); + +let phase = 2; +cp.execFileSync = (_file, args) => { + if (args.join(" ") === "phase meta --json") { + return JSON.stringify({ + max_phase: 8, + phases: Array.from({ length: 8 }, (_, index) => ({ + num: index + 1, + id: `F${index + 1}`, + name: `phase-${index + 1}`, + emits: [], + })), + }); + } + if (args.slice(0, 4).join(" ") === "phase show --session session-1") { + return `Fase corrente: ${phase}\n`; + } + return ""; +}; + +const installGate = require(GATE).default ?? require(GATE); + +test("tool completion announces each newly active phase even without a reviewer gate", async () => { + const { pi, ctx } = createFakePi(); + ctx.mode = "rpc"; + installGate(pi); + await pi.emit("session_start", {}); + + await pi.emit("tool_result", { + toolName: "reviewer_select", + input: { session: "session-1" }, + isError: false, + }); + phase = 3; + await pi.emit("tool_result", { + toolName: "reviewer_decide", + input: { session: "session-1" }, + isError: false, + }); + await pi.emit("tool_result", { + toolName: "reviewer_decide", + input: { session: "session-1" }, + isError: false, + }); + + assert.deepEqual(ctx.notifications, [ + { message: "__tht_phase_started__:F2", level: "info" }, + { message: "__tht_phase_started__:F3", level: "info" }, + ]); +}); diff --git a/harness/.pi/extensions/tht-gate.js b/harness/.pi/extensions/tht-gate.js index 6cf8179d..32a71b39 100644 --- a/harness/.pi/extensions/tht-gate.js +++ b/harness/.pi/extensions/tht-gate.js @@ -52,6 +52,11 @@ import { } from "./gate/disambiguation/index.js"; import { isReserved } from "./gate/core/reserved-labels.mjs"; +// Machine-readable RPC notification consumed by SessionBridge. Pi's extension API +// exposes notify/input but no custom client-event emitter, so this reserved prefix is +// the narrow transport contract for phase progress that does not require a human gate. +export const PHASE_STARTED_NOTIFICATION_PREFIX = "__tht_phase_started__:"; + // Load the workflow contract once when Pi loads the extension. Asking the model to // discover/read the skill as its first action proved unreliable with remote models: // they can spend the whole RPC turn exploring the repository before opening the exact @@ -596,6 +601,18 @@ export default function (pi) { let lastSteered = false; let pendingKickoff = null; let activeSessionId = null; + const announcedPhases = new Map(); + const announceCurrentPhase = async (ctx, session) => { + if (ctx.mode !== "rpc" || typeof session !== "string" || !session.trim()) return; + try { + const phase = phaseId(ctx, currentPhase(ctx, session)); + if (!/^F[1-8]$/.test(phase) || announcedPhases.get(session) === phase) return; + await ctx.ui.notify(`${PHASE_STARTED_NOTIFICATION_PREFIX}${phase}`, "info"); + announcedPhases.set(session, phase); + } catch { + // Progress reporting is observational: it must never block the workflow. + } + }; const memoryGate = createMemoryGate({ workflow: { activate: () => { @@ -748,7 +765,8 @@ export default function (pi) { // 3) BEFORE_AGENT_START: one-shot kickoff injection (appends the operational // instructions to the system prompt on the turn that starts/resumes a session). - pi.on("before_agent_start", (event, ctx) => { + pi.on("before_agent_start", async (event, ctx) => { + await announceCurrentPhase(ctx, process.env.THT_SESSION); if (!pendingKickoff) return undefined; const sessionId = process.env.THT_SESSION; const isProvidedNewSession = @@ -789,9 +807,18 @@ export default function (pi) { lastSteered = false; } activeSessionId = null; + announcedPhases.clear(); _phaseMetaCache = null; }); + // A reviewer tool is the ownership boundary that persists a decision and may + // advance the folded workflow phase. Observe the persisted phase after the tool + // completes, including auto-approved phases that never open a human widget. + pi.on("tool_result", async (event, ctx) => { + const session = typeof event.input?.session === "string" ? event.input.session : null; + if (session) await announceCurrentPhase(ctx, session); + }); + // 5) AGENT_END prose safety net: if the model emits prose instead of a reviewer_* // tool call (and the lock is active), nudge it back to the gate tools. pi.on("agent_end", async (event) => {