212 lines
7.3 KiB
TypeScript
212 lines
7.3 KiB
TypeScript
import { create } from "zustand";
|
|
import type { ActivityEntry, StreamEvent, TokenUsage, WidgetDescriptor } from "../api/types";
|
|
|
|
interface Entry {
|
|
role: "assistant";
|
|
text: string;
|
|
}
|
|
|
|
interface SessionState {
|
|
pendingWidget: WidgetDescriptor | null;
|
|
transcript: Entry[];
|
|
activityLog: ActivityEntry[];
|
|
toasts: { level: string; text: string }[];
|
|
stepMessages: { level: string; text: string }[];
|
|
lastUserEntry: { kind: "input" | "choice"; text: string } | null;
|
|
lastSystemEvent: StreamEvent | null;
|
|
currentPhase: string | null;
|
|
phaseError: string | null;
|
|
seenGateIds: Set<string>;
|
|
// True while a Pi turn is in flight. Turned off by the backend-forwarded `agent_end`
|
|
// 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;
|
|
setPhase: (phase: string | null) => void;
|
|
pushToast: (toast: { level: string; text: string }) => void;
|
|
setLastUserEntry: (e: { kind: "input" | "choice"; text: string }) => void;
|
|
recordLifecycle: (text: string) => void;
|
|
setAgentActive: (v: boolean) => void;
|
|
}
|
|
|
|
function phaseOf(value: unknown, fallback: string | null): string | null {
|
|
return typeof value === "string" && value ? value.split("_")[0] : fallback;
|
|
}
|
|
|
|
function appendStream(
|
|
log: ActivityEntry[],
|
|
entry: ActivityEntry & { kind: "thinking" | "assistant" },
|
|
): ActivityEntry[] {
|
|
const next = [...log];
|
|
const last = next.at(-1);
|
|
if (last?.kind === entry.kind) next[next.length - 1] = { ...last, text: last.text + entry.text };
|
|
else next.push(entry);
|
|
return next;
|
|
}
|
|
|
|
function lifecycleText(event: string): string {
|
|
if (event === "agent_start") return "Agent started";
|
|
if (event === "agent_end") return "Agent finished";
|
|
const words = event.replaceAll("_", " ").trim();
|
|
return words ? words[0].toUpperCase() + words.slice(1) : "System event";
|
|
}
|
|
|
|
const empty = {
|
|
pendingWidget: null,
|
|
transcript: [] as Entry[],
|
|
activityLog: [] as ActivityEntry[],
|
|
toasts: [] as { level: string; text: string }[],
|
|
stepMessages: [] as { level: string; text: string }[],
|
|
lastUserEntry: null as { kind: "input" | "choice"; text: string } | null,
|
|
lastSystemEvent: null,
|
|
currentPhase: null as string | null,
|
|
phaseError: null as string | null,
|
|
seenGateIds: new Set<string>(),
|
|
agentActive: false,
|
|
tokenUsage: null as TokenUsage | null,
|
|
};
|
|
|
|
export const useSessionStore = create<SessionState>((set) => ({
|
|
...empty,
|
|
applyEvent: (e) =>
|
|
set((st) => {
|
|
if (e.type === "ui_request") {
|
|
if (st.seenGateIds.has(e.ui_request.id)) return {};
|
|
const currentPhase = phaseOf(e.ui_request.phase, st.currentPhase);
|
|
return {
|
|
pendingWidget: e.ui_request,
|
|
seenGateIds: new Set([...st.seenGateIds, e.ui_request.id]),
|
|
// The turn is blocked on the gate, not ended: keep the agent marked active.
|
|
agentActive: true,
|
|
currentPhase,
|
|
// A new gate means the phase moved on (or re-presented): clear any error flag.
|
|
phaseError: null,
|
|
activityLog: [
|
|
...st.activityLog,
|
|
{
|
|
kind: "gate" as const,
|
|
phase: currentPhase,
|
|
text: e.ui_request.title ?? "Review requested",
|
|
},
|
|
],
|
|
};
|
|
}
|
|
if (e.type === "text_delta") {
|
|
const t = [...st.transcript];
|
|
const last = t.at(-1);
|
|
if (last && !st.pendingWidget) t[t.length - 1] = { role: "assistant", text: last.text + e.text };
|
|
else t.push({ role: "assistant", text: e.text });
|
|
// Streamed text means the turn is alive (also covers reattaching mid-turn).
|
|
return {
|
|
transcript: t,
|
|
activityLog: appendStream(st.activityLog, {
|
|
kind: "assistant",
|
|
phase: st.currentPhase,
|
|
text: e.text,
|
|
}),
|
|
agentActive: true,
|
|
};
|
|
}
|
|
if (e.type === "activity_delta") {
|
|
return {
|
|
activityLog: appendStream(st.activityLog, {
|
|
kind: "thinking",
|
|
phase: st.currentPhase,
|
|
text: e.text,
|
|
}),
|
|
agentActive: true,
|
|
};
|
|
}
|
|
if (e.type === "activity_event") {
|
|
const entry: ActivityEntry = {
|
|
kind: "tool",
|
|
phase: st.currentPhase,
|
|
text: e.activity.toolName,
|
|
toolCallId: e.activity.toolCallId,
|
|
status: e.activity.status,
|
|
};
|
|
if (e.activity.status === "running") {
|
|
return { activityLog: [...st.activityLog, entry], agentActive: true };
|
|
}
|
|
|
|
const activityLog = [...st.activityLog];
|
|
let matchingIndex = -1;
|
|
for (let index = activityLog.length - 1; index >= 0; index -= 1) {
|
|
if (activityLog[index].kind === "tool" && activityLog[index].toolCallId === e.activity.toolCallId) {
|
|
matchingIndex = index;
|
|
break;
|
|
}
|
|
}
|
|
if (matchingIndex >= 0) {
|
|
activityLog[matchingIndex] = { ...activityLog[matchingIndex], status: e.activity.status };
|
|
} else {
|
|
activityLog.push(entry);
|
|
}
|
|
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 }];
|
|
const activityLog = [
|
|
...st.activityLog,
|
|
{ kind: "status" as const, phase: st.currentPhase, text: e.text, level },
|
|
];
|
|
// An error during the active phase marks that phase red until the next gate.
|
|
return e.level === "error"
|
|
? { stepMessages, activityLog, phaseError: st.currentPhase }
|
|
: { stepMessages, activityLog };
|
|
}
|
|
if (e.type === "system_event") {
|
|
const activityLog = [
|
|
...st.activityLog,
|
|
{
|
|
kind: "lifecycle" as const,
|
|
phase: st.currentPhase,
|
|
text: lifecycleText(e.event),
|
|
},
|
|
];
|
|
return e.event === "agent_end"
|
|
? { lastSystemEvent: e, activityLog, agentActive: false }
|
|
: { lastSystemEvent: e, activityLog };
|
|
}
|
|
return {};
|
|
}),
|
|
clearPending: () => set({ pendingWidget: null }),
|
|
resetSession: () => set({ ...empty }),
|
|
setPhase: (phase) => set({ currentPhase: phase }),
|
|
pushToast: (toast) => set((st) => ({ toasts: [...st.toasts, toast] })),
|
|
// A user entry (question, gate choice, steer) hands the ball back to the harness.
|
|
setLastUserEntry: (e) =>
|
|
set((st) => ({
|
|
lastUserEntry: e,
|
|
stepMessages: [],
|
|
activityLog: [
|
|
...st.activityLog,
|
|
{ kind: "prompt", phase: st.currentPhase, text: e.text },
|
|
],
|
|
agentActive: true,
|
|
})),
|
|
recordLifecycle: (text) =>
|
|
set((st) => ({
|
|
activityLog: [
|
|
...st.activityLog,
|
|
{ kind: "lifecycle", phase: st.currentPhase, text },
|
|
],
|
|
})),
|
|
setAgentActive: (v) => set({ agentActive: v }),
|
|
}));
|