import type { AssistantMessageFrame } from "@earendil-works/pi-ai"; import type { AgentToolResult } from "../../types.ts"; import type { Context } from "../context.ts"; import type { SessionReader, Write } from "../session/types.ts"; import { appendList, pendingAssistantFrames, pendingToolOutput, setValue } from "../session/values.ts"; import type { Lane } from "./lane.ts"; import type { Drive, LaneState } from "./types.ts"; export interface ProgressChannel { write(item: T): void; seal(): void; drain(): Promise; } export async function readAssistantFrames( reader: SessionReader, operationId: string, responseEntryId: string, context: Context, ): Promise { const frames: AssistantMessageFrame[] = []; let cursor: { seq: number } | undefined; for (;;) { const page = await reader.readList( pendingAssistantFrames(operationId, responseEntryId), { order: "asc", limit: 1_000, ...(cursor === undefined ? {} : { cursor }) }, context, ); frames.push(...page.map(({ value }) => value)); if (page.length < 1_000) return frames; cursor = { seq: page[page.length - 1]!.seq }; } } function openProgress( lane: Lane, drive: Drive, commitWrite: (item: T) => Write, stillOwns: (state: LaneState) => boolean, ): ProgressChannel { let sealed = false; let latest: Promise = Promise.resolve(); return { write(item) { if (sealed) return; const write = lane .command((projection) => { if (!stillOwns(projection)) return { kind: "return", result: undefined }; return { kind: "commit", writes: [commitWrite(item)], next: projection, materialize: () => undefined, }; }, drive.context) .then(() => undefined); latest = write; void write.catch(() => {}); }, seal() { sealed = true; }, async drain() { await latest; }, }; } export function openFrameProgress( lane: Lane, drive: Drive, responseEntryId: string, ): ProgressChannel { const address = pendingAssistantFrames(drive.operationId, responseEntryId); return openProgress( lane, drive, (frame) => appendList(address, frame), (state) => { const run = state.operation?.state; if (run === undefined) return false; return ( (run.at === "assistant.effect_pending" || run.at === "deferred.effect_pending") && run.responseEntryId === responseEntryId ); }, ); } export function openToolProgress( lane: Lane, drive: Drive, turnId: string, sourceIndex: number, invocationId: string, ): ProgressChannel> { const address = pendingToolOutput(drive.operationId, invocationId); return openProgress( lane, drive, (snapshot) => setValue(address, snapshot), (state) => { const operation = state.operation; if (operation?.state.at !== "tools") return false; const batch = operation.state.batch; return ( batch.turnId === turnId && batch.calls.some( (call) => call.sourceIndex === sourceIndex && call.resultEntryId === invocationId && call.status === "effect_pending", ) ); }, ); }