fix(ui): stream phase progress without gates
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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") {
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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(<WorkflowBar />);
|
||||
|
||||
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" } });
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -176,14 +176,26 @@ export const useSessionStore = create<SessionState>((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 };
|
||||
|
||||
@@ -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" },
|
||||
]);
|
||||
});
|
||||
@@ -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) => {
|
||||
|
||||
Reference in New Issue
Block a user