{"version":3,"file":"reconcile.d.ts","sourceRoot":"","sources":["../../../../src/harness/runtime/drive/reconcile.ts"],"names":[],"mappings":"AASA,OAAO,KAAK,EAAE,IAAI,EAAE,MAAM,YAAY,CAAC;AACvC,OAAO,KAAK,EAAE,KAAK,EAAE,eAAe,EAAE,MAAM,aAAa,CAAC;AAwH1D,6EAA6E;AAC7E,wBAAsB,kBAAkB,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EAC3E,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,GACV,OAAO,CAAC,eAAe,CAAC,CAsC1B","sourcesContent":["import type { DeferredHandle } from \"@earendil-works/pi-ai\";\nimport type { HarnessEvent } from \"../../agent-harness.ts\";\nimport { getTelemetryContext } from \"../../context.ts\";\nimport { SessionInvariantError } from \"../../session/session.ts\";\nimport type {\n\tDeferredEffectPendingOperation,\n\tDeferredSuspendedOperation,\n\tOperationState,\n} from \"../../session/types.ts\";\nimport type { Lane } from \"../lane.ts\";\nimport type { Drive, ProcedureResult } from \"../types.ts\";\nimport { readDeferredSourceHandle } from \"./deferred.ts\";\nimport { recoverCancelledAssistantEffect } from \"./recovery.ts\";\nimport { operationCleanupWrites, operationResultRecord } from \"./terminal.ts\";\nimport { runTools } from \"./tools.ts\";\n\nasync function cancelDeferredBestEffort<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tdeferred: DeferredSuspendedOperation | DeferredEffectPendingOperation,\n\thandle: DeferredHandle,\n): Promise<void> {\n\tconst identity = deferred.configuration.model;\n\tconst model = lane.models.getModel(identity.provider, identity.modelId);\n\tif (model === undefined) return;\n\ttry {\n\t\tawait lane.models.cancelDeferred(model, handle, {\n\t\t\tsignal: drive.closeSignal,\n\t\t\ttelemetryContext: getTelemetryContext(drive.context),\n\t\t\ttimeoutMs: deferred.streamOptions.timeoutMs,\n\t\t\tmaxRetries: deferred.streamOptions.maxRetries,\n\t\t\tmaxRetryDelayMs: deferred.streamOptions.maxRetryDelayMs,\n\t\t\theaders: deferred.streamOptions.headers,\n\t\t});\n\t} catch {\n\t\t// Remote cancellation is best-effort; durable local reconciliation must continue.\n\t}\n}\n\nasync function readDeferredHandle<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tdeferred: DeferredSuspendedOperation | DeferredEffectPendingOperation,\n): Promise<DeferredHandle> {\n\treturn lane.settleOperation(\n\t\tdeferred,\n\t\tasync (_state, _current, _meta, reader) => ({\n\t\t\tkind: \"return\",\n\t\t\tresult: await readDeferredSourceHandle(reader, deferred, drive.context),\n\t\t}),\n\t\tdrive.context,\n\t);\n}\n\nasync function publishAbortedTerminal<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tcapability: OperationState,\n): Promise<ProcedureResult> {\n\treturn lane.settleOperation(\n\t\tcapability,\n\t\tasync (state, current, meta, reader) => {\n\t\t\tif (current.control.status !== \"cancel_requested\") {\n\t\t\t\tthrow new SessionInvariantError(\"Cancellation reconciliation requires cancelled durable control\");\n\t\t\t}\n\t\t\tconst record = operationResultRecord(meta, \"aborted\", state.tipId);\n\t\t\tconst cleanup = await operationCleanupWrites(reader, drive.operationId, current, drive.context);\n\t\t\tconst events: HarnessEvent[] = [];\n\t\t\tif (meta.intent.kind === \"run\") {\n\t\t\t\tswitch (current.at) {\n\t\t\t\t\tcase \"summary.deciding\":\n\t\t\t\t\tcase \"summary.ready\":\n\t\t\t\t\tcase \"summary.effect_pending\":\n\t\t\t\t\tcase \"summary.retry_wait\":\n\t\t\t\t\t\tif (current.task.boundary.kind !== \"resume_checkpoint\" || current.task.reason === undefined) {\n\t\t\t\t\t\t\tthrow new SessionInvariantError(\"Cancelled run summary has an invalid result boundary\");\n\t\t\t\t\t\t}\n\t\t\t\t\t\tevents.push({\n\t\t\t\t\t\t\ttype: \"compaction_end\",\n\t\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\t\treason: current.task.reason,\n\t\t\t\t\t\t\tstatus: \"aborted\",\n\t\t\t\t\t\t\tendedAt: record.endedAt,\n\t\t\t\t\t\t});\n\t\t\t\t\t\tbreak;\n\t\t\t\t\tdefault:\n\t\t\t\t\t\tbreak;\n\t\t\t\t}\n\t\t\t\tevents.push({\n\t\t\t\t\ttype: \"run_end\",\n\t\t\t\t\tlane: lane.name,\n\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\tstatus: \"aborted\",\n\t\t\t\t\tfromTipId: meta.sourceTipId,\n\t\t\t\t\ttipId: state.tipId,\n\t\t\t\t\tendedAt: record.endedAt,\n\t\t\t\t});\n\t\t\t} else if (meta.intent.kind === \"compaction\") {\n\t\t\t\tevents.push({\n\t\t\t\t\ttype: \"compaction_end\",\n\t\t\t\t\tlane: lane.name,\n\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\treason: \"manual\",\n\t\t\t\t\tstatus: \"aborted\",\n\t\t\t\t\tendedAt: record.endedAt,\n\t\t\t\t});\n\t\t\t} else {\n\t\t\t\tevents.push({\n\t\t\t\t\ttype: \"navigation_end\",\n\t\t\t\t\tlane: lane.name,\n\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\tstatus: \"aborted\",\n\t\t\t\t\tfromTipId: meta.sourceTipId,\n\t\t\t\t\ttipId: state.tipId,\n\t\t\t\t\tendedAt: record.endedAt,\n\t\t\t\t});\n\t\t\t}\n\t\t\treturn {\n\t\t\t\tkind: \"finish\",\n\t\t\t\twrites: cleanup,\n\t\t\t\trecord,\n\t\t\t\tmaterialize: () => ({ kind: \"settled\", outcome: record }) as const,\n\t\t\t\tevents: () => events,\n\t\t\t};\n\t\t},\n\t\tdrive.context,\n\t);\n}\n\n/** Advance one cancelled durable leaf without starting new ordinary work. */\nexport async function reconcileOperation<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n): Promise<ProcedureResult> {\n\tconst operation = lane.state.operation;\n\tif (operation === null || operation.meta.operationId !== drive.operationId) {\n\t\tthrow new SessionInvariantError(`Drive ${drive.operationId} has no matching operation to reconcile`);\n\t}\n\tif (operation.state.control.status !== \"cancel_requested\") {\n\t\tthrow new SessionInvariantError(`Operation ${drive.operationId} is not cancelled`);\n\t}\n\tdrive.beginAbort(Promise.resolve());\n\tdrive.signalAbort();\n\n\tconst state = operation.state;\n\tswitch (state.at) {\n\t\tcase \"assistant.effect_pending\":\n\t\t\treturn recoverCancelledAssistantEffect(lane, drive, state);\n\t\tcase \"tools\":\n\t\t\treturn runTools(lane, drive, state);\n\t\tcase \"deferred.suspended\": {\n\t\t\tconst handle = await readDeferredHandle(lane, drive, state);\n\t\t\tawait cancelDeferredBestEffort(lane, drive, state, handle);\n\t\t\treturn publishAbortedTerminal(lane, drive, state);\n\t\t}\n\t\tcase \"deferred.effect_pending\": {\n\t\t\tconst handle = await readDeferredHandle(lane, drive, state);\n\t\t\tawait cancelDeferredBestEffort(lane, drive, state, handle);\n\t\t\treturn recoverCancelledAssistantEffect(lane, drive, state);\n\t\t}\n\t\tcase \"starting\":\n\t\tcase \"checkpoint\":\n\t\tcase \"assistant.ready\":\n\t\tcase \"assistant.retry_wait\":\n\t\tcase \"summary.deciding\":\n\t\tcase \"summary.ready\":\n\t\tcase \"summary.effect_pending\":\n\t\tcase \"summary.retry_wait\":\n\t\tcase \"navigation.ready_to_commit\":\n\t\t\treturn publishAbortedTerminal(lane, drive, state);\n\t}\n}\n"]}