{"version":3,"file":"progress.d.ts","sourceRoot":"","sources":["../../../src/harness/runtime/progress.ts"],"names":[],"mappings":"AAAA,OAAO,KAAK,EAAE,qBAAqB,EAAE,MAAM,uBAAuB,CAAC;AACnE,OAAO,KAAK,EAAE,eAAe,EAAE,MAAM,gBAAgB,CAAC;AACtD,OAAO,KAAK,EAAE,OAAO,EAAE,MAAM,eAAe,CAAC;AAC7C,OAAO,KAAK,EAAE,aAAa,EAAS,MAAM,qBAAqB,CAAC;AAEhE,OAAO,KAAK,EAAE,IAAI,EAAE,MAAM,WAAW,CAAC;AACtC,OAAO,KAAK,EAAE,KAAK,EAAa,MAAM,YAAY,CAAC;AAEnD,MAAM,WAAW,eAAe,CAAC,CAAC;IACjC,KAAK,CAAC,IAAI,EAAE,CAAC,GAAG,IAAI,CAAC;IACrB,IAAI,IAAI,IAAI,CAAC;IACb,KAAK,IAAI,OAAO,CAAC,IAAI,CAAC,CAAC;CACvB;AAED,wBAAsB,mBAAmB,CACxC,MAAM,EAAE,aAAa,EACrB,WAAW,EAAE,MAAM,EACnB,eAAe,EAAE,MAAM,EACvB,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,qBAAqB,EAAE,CAAC,CAalC;AAoCD,wBAAgB,iBAAiB,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACpE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,eAAe,EAAE,MAAM,GACrB,eAAe,CAAC,qBAAqB,CAAC,CAexC;AAED,wBAAgB,gBAAgB,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACnE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,MAAM,EAAE,MAAM,EACd,WAAW,EAAE,MAAM,EACnB,YAAY,EAAE,MAAM,GAClB,eAAe,CAAC,eAAe,CAAC,OAAO,CAAC,CAAC,CAqB3C","sourcesContent":["import type { AssistantMessageFrame } from \"@earendil-works/pi-ai\";\nimport type { AgentToolResult } from \"../../types.ts\";\nimport type { Context } from \"../context.ts\";\nimport type { SessionReader, Write } from \"../session/types.ts\";\nimport { appendList, pendingAssistantFrames, pendingToolOutput, setValue } from \"../session/values.ts\";\nimport type { Lane } from \"./lane.ts\";\nimport type { Drive, LaneState } from \"./types.ts\";\n\nexport interface ProgressChannel<T> {\n\twrite(item: T): void;\n\tseal(): void;\n\tdrain(): Promise<void>;\n}\n\nexport async function readAssistantFrames(\n\treader: SessionReader,\n\toperationId: string,\n\tresponseEntryId: string,\n\tcontext: Context,\n): Promise<AssistantMessageFrame[]> {\n\tconst frames: AssistantMessageFrame[] = [];\n\tlet cursor: { seq: number } | undefined;\n\tfor (;;) {\n\t\tconst page = await reader.readList(\n\t\t\tpendingAssistantFrames(operationId, responseEntryId),\n\t\t\t{ order: \"asc\", limit: 1_000, ...(cursor === undefined ? {} : { cursor }) },\n\t\t\tcontext,\n\t\t);\n\t\tframes.push(...page.map(({ value }) => value));\n\t\tif (page.length < 1_000) return frames;\n\t\tcursor = { seq: page[page.length - 1]!.seq };\n\t}\n}\n\nfunction openProgress<TContext extends object | undefined, T>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tcommitWrite: (item: T) => Write,\n\tstillOwns: (state: LaneState) => boolean,\n): ProgressChannel<T> {\n\tlet sealed = false;\n\tlet latest: Promise<void> = Promise.resolve();\n\treturn {\n\t\twrite(item) {\n\t\t\tif (sealed) return;\n\t\t\tconst write = lane\n\t\t\t\t.command((projection) => {\n\t\t\t\t\tif (!stillOwns(projection)) return { kind: \"return\", result: undefined };\n\t\t\t\t\treturn {\n\t\t\t\t\t\tkind: \"commit\",\n\t\t\t\t\t\twrites: [commitWrite(item)],\n\t\t\t\t\t\tnext: projection,\n\t\t\t\t\t\tmaterialize: () => undefined,\n\t\t\t\t\t};\n\t\t\t\t}, drive.context)\n\t\t\t\t.then(() => undefined);\n\t\t\tlatest = write;\n\t\t\tvoid write.catch(() => {});\n\t\t},\n\t\tseal() {\n\t\t\tsealed = true;\n\t\t},\n\t\tasync drain() {\n\t\t\tawait latest;\n\t\t},\n\t};\n}\n\nexport function openFrameProgress<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tresponseEntryId: string,\n): ProgressChannel<AssistantMessageFrame> {\n\tconst address = pendingAssistantFrames(drive.operationId, responseEntryId);\n\treturn openProgress(\n\t\tlane,\n\t\tdrive,\n\t\t(frame) => appendList(address, frame),\n\t\t(state) => {\n\t\t\tconst run = state.operation?.state;\n\t\t\tif (run === undefined) return false;\n\t\t\treturn (\n\t\t\t\t(run.at === \"assistant.effect_pending\" || run.at === \"deferred.effect_pending\") &&\n\t\t\t\trun.responseEntryId === responseEntryId\n\t\t\t);\n\t\t},\n\t);\n}\n\nexport function openToolProgress<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tturnId: string,\n\tsourceIndex: number,\n\tinvocationId: string,\n): ProgressChannel<AgentToolResult<unknown>> {\n\tconst address = pendingToolOutput(drive.operationId, invocationId);\n\treturn openProgress(\n\t\tlane,\n\t\tdrive,\n\t\t(snapshot) => setValue(address, snapshot),\n\t\t(state) => {\n\t\t\tconst operation = state.operation;\n\t\t\tif (operation?.state.at !== \"tools\") return false;\n\t\t\tconst batch = operation.state.batch;\n\t\t\treturn (\n\t\t\t\tbatch.turnId === turnId &&\n\t\t\t\tbatch.calls.some(\n\t\t\t\t\t(call) =>\n\t\t\t\t\t\tcall.sourceIndex === sourceIndex &&\n\t\t\t\t\t\tcall.resultEntryId === invocationId &&\n\t\t\t\t\t\tcall.status === \"effect_pending\",\n\t\t\t\t)\n\t\t\t);\n\t\t},\n\t);\n}\n"]}