Something went wrong. Try again.
watches turns and automatically creates and resolves tasks that get pinned into main agent context
Something went wrong. Try again.
prime-commitment-observer index.ts
8.2 kB · 240 lines
TypeScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241import { randomUUID } from "node:crypto";import type { AgentEndEvent, BeforeAgentStartEvent, CustomEntry, ExtensionAPI, ExtensionContext, InputSource, SessionEntry,} from "@earendil-works/pi-coding-agent";import { CLEARED_CONTEXT, nextContextRevision, REFRESHING_CONTEXT } from "./src/inject.ts";import { onLiveBranch } from "./src/liveness.ts";import { configFromEnv, runObserver } from "./src/observer.ts";import { isCommitmentEvent, liveThreads, observationChangesState, replay } from "./src/state.ts";import type { CommitmentEvent, Registry } from "./src/types.ts";
const COMMITMENT_EVENT_TYPE = "commitment-event-v2";const CONTEXT_CUSTOM_TYPE = "commitment-threads";const LIST_CUSTOM_TYPE = "commitment-threads-list";type BeforeAgentStartWithSource = BeforeAgentStartEvent & { source: InputSource | "internal" };
function readEvents(ctx: ExtensionContext): CommitmentEvent[] { return ctx.sessionManager .getBranch() .filter( (entry): entry is CustomEntry => entry.type === "custom" && entry.customType === COMMITMENT_EVENT_TYPE, ) .map((entry) => entry.data) .filter(isCommitmentEvent);}
const registryFrom = (ctx: ExtensionContext): Registry => replay(readEvents(ctx));
function appendEvent(pi: ExtensionAPI, event: CommitmentEvent): void { pi.appendEntry(COMMITMENT_EVENT_TYPE, event);}
function latestContext(entries: readonly SessionEntry[]): string | undefined { for (let index = entries.length - 1; index >= 0; index -= 1) { const entry = entries[index]; if (entry.type !== "custom_message" || entry.customType !== CONTEXT_CUSTOM_TYPE) continue; return typeof entry.content === "string" ? entry.content : undefined; } return undefined;}
function formatThread(thread: Registry["threads"][number]): string { return `[${thread.id}] ${thread.status} ${thread.canonicalRequest}`;}
function formatRegistry(registry: Registry): string { const open = liveThreads(registry); const openIds = new Set(open.map((thread) => thread.id)); const closed = registry.threads.filter((thread) => !openIds.has(thread.id)); const sections: string[] = []; if (open.length > 0) sections.push(`Open requests:\n${open.map(formatThread).join("\n")}`); if (closed.length > 0) sections.push(`Closed requests:\n${closed.map(formatThread).join("\n")}`); if (registry.unreconciled.length > 0) { sections.push(`Unreconciled user messages: ${registry.unreconciled.length}`); } return sections.join("\n\n") || "No tracked requests.";}
export default function commitmentObserver(pi: ExtensionAPI): void { const cfg = configFromEnv(); let observerQueue = Promise.resolve(); let pendingReconciliations = 0; let consecutiveFailures = 0;
// The failure path has no owner to catch it: a replaced session has nowhere to record the // failure, and the queue that holds this work discards the rejection. const recordFailure = (ctx: ExtensionContext, messageIds: string[], error: unknown): void => { const text = (error instanceof Error ? error.message : String(error)).slice(0, 2_000); try { appendEvent(pi, { kind: "observer_failure", messageIds, error: text, at: Date.now() }); consecutiveFailures += 1; if (consecutiveFailures === 1 && ctx.hasUI) { ctx.ui.notify(`Commitment observer failed: ${text}`, "warning"); } } catch { // The session was replaced under this reconciliation. } };
const liveRevision = (registry: Registry): string => JSON.stringify( liveThreads(registry).map(({ id, canonicalRequest, status, updatedAt }) => ({ id, canonicalRequest, status, updatedAt, })), );
const reconcile = async ( ctx: ExtensionContext, messages: AgentEndEvent["messages"], targetMessageIds: ReadonlySet<string>, boundaryEntryId: string | null, ): Promise<void> => { if (!onLiveBranch(ctx, boundaryEntryId)) return; const registry = registryFrom(ctx); const currentTargets = new Set( registry.unreconciled .map((message) => message.messageId) .filter((messageId) => targetMessageIds.has(messageId)), ); if (currentTargets.size === 0 && liveThreads(registry).length === 0) return; const revision = liveRevision(registry); const result = await runObserver(registry, messages, currentTargets, ctx.modelRegistry, cfg); if (!onLiveBranch(ctx, boundaryEntryId) || liveRevision(registryFrom(ctx)) !== revision) return; if (!result.ok) { recordFailure(ctx, [...currentTargets], result.error); return; } consecutiveFailures = 0; if (currentTargets.size === 0 && !observationChangesState(registry, result.output)) return; appendEvent(pi, { kind: "observation", reconciled: [...currentTargets], output: result.output, at: Date.now(), }); };
const scheduleReconcile = ( ctx: ExtensionContext, messages: AgentEndEvent["messages"] = [], ): Promise<void> => { const targetMessageIds = new Set(registryFrom(ctx).unreconciled.map((message) => message.messageId)); const boundaryEntryId = ctx.sessionManager.getLeafId(); pendingReconciliations += 1; observerQueue = observerQueue .then(() => reconcile(ctx, messages, targetMessageIds, boundaryEntryId)) .catch((error) => { if (onLiveBranch(ctx, boundaryEntryId)) recordFailure(ctx, [...targetMessageIds], error); }) .finally(() => { pendingReconciliations -= 1; }); return observerQueue; };
pi.on("input", async (event) => { if (event.text.trim().length === 0 || event.text.trimStart().startsWith("/")) return { action: "continue" }; appendEvent(pi, { kind: "user_message", messageId: randomUUID(), text: event.text, at: Date.now(), }); return { action: "continue" }; });
pi.on("agent_end", async (event, ctx) => { void scheduleReconcile(ctx, event.messages); });
pi.on("session_start", async (event, ctx) => { if (event.reason !== "new") void scheduleReconcile(ctx); });
pi.on("context", async (event) => ({ messages: event.messages.filter( (message) => message.role !== "custom" || message.customType !== LIST_CUSTOM_TYPE, ), }));
pi.on("before_agent_start", async (rawEvent, ctx) => { const event = rawEvent as BeforeAgentStartWithSource; if (event.source === "internal") return undefined; const content = nextContextRevision( event.source, registryFrom(ctx), latestContext(ctx.sessionManager.getBranch()), pendingReconciliations > 0, ); if (!content) return undefined; return { message: { customType: CONTEXT_CUSTOM_TYPE, content, display: false, details: { state: content === CLEARED_CONTEXT ? "clear" : content === REFRESHING_CONTEXT ? "refreshing" : "open", }, }, }; });
pi.registerCommand("threads", { description: "Inspect and edit tracked user requests", getArgumentCompletions: (prefix: string) => ["done", "dismiss", "reopen", "drop", "edit", "reconcile"] .filter((value) => value.startsWith(prefix)) .map((value) => ({ value, label: value })), handler: async (args, ctx) => { const show = (content: string): void => { pi.sendMessage({ customType: LIST_CUSTOM_TYPE, content, display: true }); }; const [command = "", id = "", ...rest] = args.trim().split(/\s+/); const registry = registryFrom(ctx); if (!command) { show(formatRegistry(registry)); return; } if (command === "reconcile") { await scheduleReconcile(ctx); return; } const thread = registry.threads.find((candidate) => candidate.id === id); if (!thread) { show(`Unknown request ${id || "id"}.`); return; } const at = Date.now(); if (command === "done" || command === "dismiss" || command === "reopen" || command === "drop") { appendEvent(pi, { kind: "user_action", threadId: id, action: command === "done" ? "accept" : command, reason: rest.join(" ") || undefined, at, }); return; } if (command === "edit") { const canonicalRequest = rest.join(" ").trim(); if (!canonicalRequest) { show("Usage: /threads edit <id> <request>"); return; } appendEvent(pi, { kind: "correction", threadId: id, canonicalRequest, at }); return; } show("Usage: /threads [done|dismiss|reopen|drop|edit|reconcile]"); }, });}