{"version":3,"file":"response.d.ts","sourceRoot":"","sources":["../../../../src/harness/runtime/drive/response.ts"],"names":[],"mappings":"AAQA,OAAO,KAAK,EAAE,OAAO,EAAE,MAAM,kBAAkB,CAAC;AAChD,OAAO,KAAK,EAAE,yBAAyB,EAAE,uBAAuB,EAAE,MAAM,8BAA8B,CAAC;AAGvG,OAAO,EACN,KAAK,+BAA+B,EAEpC,KAAK,8BAA8B,EAGnC,KAAK,cAAc,EACnB,KAAK,cAAc,EAEnB,KAAK,uBAAuB,EAI5B,MAAM,wBAAwB,CAAC;AAEhC,OAAO,KAAK,EAAE,IAAI,EAAE,MAAM,YAAY,CAAC;AAEvC,OAAO,KAAK,EAAE,KAAK,EAAE,eAAe,EAAE,MAAM,aAAa,CAAC;AAK1D,MAAM,MAAM,0BAA0B,GAAG;IACxC,QAAQ,EAAE,uBAAuB,CAAC;IAClC,aAAa,CACZ,OAAO,EAAE,uBAAuB,EAChC,QAAQ,EAAE,yBAAyB,EACnC,OAAO,EAAE,OAAO,GACd,OAAO,CAAC,uBAAuB,CAAC,CAAC;IACpC,KAAK,IAAI,OAAO,CAAC,IAAI,CAAC,CAAC;CACvB,CAAC;AAEF,wBAAgB,qBAAqB,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACxE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,eAAe,EAAE,MAAM,EACvB,QAAQ,UAAQ,GACd,0BAA0B,CA8C5B;AAED,KAAK,cAAc,GAAG,+BAA+B,GAAG,8BAA8B,CAAC;AACvF,KAAK,yBAAyB,GAAG,OAAO,CACvC,cAAc,EACd;IAAE,EAAE,EAAE,iBAAiB,GAAG,sBAAsB,GAAG,oBAAoB,GAAG,yBAAyB,CAAA;CAAE,CACrG,CAAC;AAEF,2FAA2F;AAC3F,wBAAsB,2BAA2B,CAChD,QAAQ,SAAS,MAAM,GAAG,SAAS,EACnC,MAAM,SAAS,yBAAyB,EACvC,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EAAE,KAAK,EAAE,KAAK,EAAE,UAAU,EAAE,MAAM,EAAE,KAAK,EAAE,cAAc,GAAG,OAAO,CAAC,eAAe,CAAC,CA6BzG;AA2CD,yFAAyF;AACzF,wBAAsB,eAAe,CAAC,QAAQ,SAAS,MAAM,GAAG,SAAS,EACxE,IAAI,EAAE,IAAI,CAAC,QAAQ,CAAC,EACpB,KAAK,EAAE,KAAK,EACZ,MAAM,EAAE,cAAc,EACtB,QAAQ,EAAE,uBAAuB,EACjC,OAAO,GAAE;IAAE,QAAQ,CAAC,EAAE,IAAI,CAAA;CAAO,GAC/B,OAAO,CAAC,eAAe,CAAC,CAyS1B","sourcesContent":["import {\n\ttype AssistantMessageEvent,\n\tAssistantMessageFrameEncoder,\n\tisContextOverflow,\n\tisRecoverableLength,\n\tisRetryableAssistantError,\n} from \"@earendil-works/pi-ai\";\nimport type { HarnessEvent } from \"../../agent-harness.ts\";\nimport type { Context } from \"../../context.ts\";\nimport type { AssistantResponseMetadata, AssistantStreamObserver } from \"../../execution/assistant.ts\";\nimport { insertEntry, insertUsage } from \"../../session/commit.ts\";\nimport { SessionInvariantError } from \"../../session/session.ts\";\nimport {\n\ttype AssistantEffectPendingOperation,\n\ttype CommitResult,\n\ttype DeferredEffectPendingOperation,\n\ttype MessageEntry,\n\ttype NewEntry,\n\ttype OperationError,\n\ttype OperationState,\n\toperationScopeOf,\n\ttype SettledAssistantMessage,\n\ttype SummaryDecidingOperation,\n\ttype ToolCall,\n\ttype UsageRow,\n} from \"../../session/types.ts\";\nimport { branchTip, deleteList, operationPreparation, pendingAssistantFrames, setValue } from \"../../session/values.ts\";\nimport type { Lane } from \"../lane.ts\";\nimport { openFrameProgress } from \"../progress.ts\";\nimport type { Drive, ProcedureResult } from \"../types.ts\";\nimport { retryDelay, retryNotBefore } from \"./retry.ts\";\nimport { prepareOverflowCompaction } from \"./structural.ts\";\nimport { operationCleanupWrites, operationResultRecord } from \"./terminal.ts\";\n\nexport type AssistantResponseLifecycle = {\n\tobserver: AssistantStreamObserver;\n\tafterResponse(\n\t\tmessage: SettledAssistantMessage,\n\t\tmetadata: AssistantResponseMetadata,\n\t\tcontext: Context,\n\t): Promise<SettledAssistantMessage>;\n\tclose(): Promise<void>;\n};\n\nexport function openAssistantResponse<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tresponseEntryId: string,\n\trecovery = false,\n): AssistantResponseLifecycle {\n\tconst progress = openFrameProgress(lane, drive, responseEntryId);\n\tconst frameEncoder = new AssistantMessageFrameEncoder();\n\tconst eventContext = {\n\t\tlane: lane.name,\n\t\trunId: drive.operationId,\n\t\t...(recovery ? { recovery: true as const } : {}),\n\t};\n\tconst close = async () => {\n\t\tprogress.seal();\n\t\tawait progress.drain();\n\t};\n\treturn {\n\t\tobserver: {\n\t\t\tstart(message, event, context) {\n\t\t\t\tconst frame = frameEncoder.encode(event);\n\t\t\t\tif (frame !== undefined) progress.write(frame);\n\t\t\t\treturn lane.emitBatch([{ type: \"message_start\", ...eventContext, message }], context);\n\t\t\t},\n\t\t\tupdate(message, event: AssistantMessageEvent, context) {\n\t\t\t\tconst frame = frameEncoder.encode(event);\n\t\t\t\tif (frame !== undefined) progress.write(frame);\n\t\t\t\treturn lane.emitBatch(\n\t\t\t\t\t[{ type: \"message_update\", ...eventContext, message, event, ...(frame === undefined ? {} : { frame }) }],\n\t\t\t\t\tcontext,\n\t\t\t\t);\n\t\t\t},\n\t\t\tend(message, context) {\n\t\t\t\treturn lane.emitBatch(\n\t\t\t\t\t[{ type: \"message_end\", ...eventContext, message, entryId: responseEntryId }],\n\t\t\t\t\tcontext,\n\t\t\t\t);\n\t\t\t},\n\t\t},\n\t\tasync afterResponse(message, metadata, context) {\n\t\t\tawait close();\n\t\t\tconst result = await lane.hooks.runWithGate(\n\t\t\t\t\"after_response\",\n\t\t\t\t{ lane: lane.name, runId: drive.operationId, ...metadata, message },\n\t\t\t\tdrive.gate,\n\t\t\t\tcontext,\n\t\t\t);\n\t\t\treturn result?.message ?? message;\n\t\t},\n\t\tclose,\n\t};\n}\n\ntype ResponseIntent = AssistantEffectPendingOperation | DeferredEffectPendingOperation;\ntype ConfigurationFailureState = Extract<\n\tOperationState,\n\t{ at: \"assistant.ready\" | \"assistant.retry_wait\" | \"deferred.suspended\" | \"deferred.effect_pending\" }\n>;\n\n/** Publish a non-retryable request-configuration failure before reserving response ids. */\nexport async function publishConfigurationFailure<\n\tTContext extends object | undefined,\n\tTState extends ConfigurationFailureState,\n>(lane: Lane<TContext>, drive: Drive, capability: TState, error: OperationError): Promise<ProcedureResult> {\n\tconst result = await lane.continueOperation(\n\t\tcapability,\n\t\tasync (state, current, meta, reader) => {\n\t\t\tif (state.tipId === null) throw new SessionInvariantError(\"Failed run has no Branch tip\");\n\t\t\tconst record = operationResultRecord(meta, \"failed\", state.tipId, error);\n\t\t\tconst cleanup = await operationCleanupWrites(reader, drive.operationId, current, drive.context);\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: () => [\n\t\t\t\t\t{\n\t\t\t\t\t\ttype: \"run_end\",\n\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\tstatus: \"failed\",\n\t\t\t\t\t\terror,\n\t\t\t\t\t\tfromTipId: meta.sourceTipId,\n\t\t\t\t\t\ttipId: state.tipId,\n\t\t\t\t\t\tendedAt: record.endedAt,\n\t\t\t\t\t},\n\t\t\t\t],\n\t\t\t};\n\t\t},\n\t\tdrive.context,\n\t);\n\treturn result.kind === \"cancel_requested\" ? { kind: \"continue\" } : result.value;\n}\n\nfunction uuidV7Timestamp(id: string): number {\n\tconst timestamp = Number.parseInt(id.slice(0, 8) + id.slice(9, 13), 16);\n\tif (!Number.isSafeInteger(timestamp)) throw new SessionInvariantError(`Invalid reserved UUIDv7 ${id}`);\n\treturn timestamp;\n}\n\nfunction providerError(source: \"assistant\" | \"deferred\", message: SettledAssistantMessage): OperationError {\n\treturn {\n\t\tcode: \"assistant_error\",\n\t\tmessage:\n\t\t\tmessage.errorMessage ??\n\t\t\t`${source === \"assistant\" ? \"Assistant\" : \"Deferred\"} request ended with ${message.stopReason}`,\n\t};\n}\n\nfunction normalizeError(message: SettledAssistantMessage, errorMessage: string): SettledAssistantMessage {\n\treturn { ...message, stopReason: \"error\", errorMessage };\n}\n\nfunction normalizeAborted(source: \"assistant\" | \"deferred\", message: SettledAssistantMessage): SettledAssistantMessage {\n\treturn {\n\t\t...message,\n\t\tstopReason: \"aborted\",\n\t\terrorMessage:\n\t\t\tmessage.errorMessage ?? `${source === \"assistant\" ? \"Assistant\" : \"Deferred\"} request was cancelled`,\n\t};\n}\n\nfunction deferredHandleIsValid(message: SettledAssistantMessage, generation: AssistantEffectPendingOperation): boolean {\n\tconst handle = message.deferred;\n\tconst identity = generation.generationContext.configuration.model;\n\treturn (\n\t\tmessage.stopReason === \"deferred\" &&\n\t\thandle !== undefined &&\n\t\thandle.id.length !== 0 &&\n\t\thandle.provider === identity.provider &&\n\t\thandle.modelId === identity.modelId &&\n\t\thandle.api === message.api\n\t);\n}\n\n/** Classify and atomically settle one assistant-generation or deferred-poll response. */\nexport async function publishResponse<TContext extends object | undefined>(\n\tlane: Lane<TContext>,\n\tdrive: Drive,\n\tintent: ResponseIntent,\n\tresponse: SettledAssistantMessage,\n\toptions: { recovery?: true } = {},\n): Promise<ProcedureResult> {\n\tconst overflow =\n\t\tintent.at === \"assistant.effect_pending\" &&\n\t\t(isContextOverflow(response, intent.contextWindow) || isRecoverableLength(response, intent.intendedOutputLimit));\n\tconst overflowPreparation =\n\t\toverflow && !intent.generationContext.overflowRecoveryUsed\n\t\t\t? await prepareOverflowCompaction(lane, drive, intent)\n\t\t\t: undefined;\n\treturn lane.settleOperation<ResponseIntent, ProcedureResult>(\n\t\tintent,\n\t\tasync (state, current, meta, reader) => {\n\t\t\tconst source = current.at === \"assistant.effect_pending\" ? \"assistant\" : \"deferred\";\n\t\t\tconst responseEntryId = intent.responseEntryId;\n\t\t\tconst configuration =\n\t\t\t\tcurrent.at === \"assistant.effect_pending\" ? current.generationContext.configuration : current.configuration;\n\t\t\tconst turnId =\n\t\t\t\tcurrent.at === \"assistant.effect_pending\"\n\t\t\t\t\t? current.generationContext.stepId\n\t\t\t\t\t: `${current.stepId}:poll:${current.poll}`;\n\t\t\tconst scope = { ...operationScopeOf(current), latestAssistantEntryId: responseEntryId };\n\t\t\tlet committed = response;\n\t\t\tlet settled: OperationState | undefined;\n\t\t\tlet failure: OperationError | undefined;\n\n\t\t\tif (current.control.status === \"cancel_requested\") {\n\t\t\t\tcommitted = normalizeAborted(source, response);\n\t\t\t\tsettled = {\n\t\t\t\t\t...scope,\n\t\t\t\t\tat: \"checkpoint\",\n\t\t\t\t\tcontinuation: { kind: \"may_finish\", includeFinalAssistant: true },\n\t\t\t\t\ttriggerEntryId: responseEntryId,\n\t\t\t\t};\n\t\t\t} else if (response.stopReason === \"aborted\") {\n\t\t\t\tthrow new SessionInvariantError(\n\t\t\t\t\t`${source === \"assistant\" ? \"Assistant\" : \"Deferred\"} response is aborted while durable control is running`,\n\t\t\t\t);\n\t\t\t} else if (current.at === \"assistant.effect_pending\" && overflow) {\n\t\t\t\tcommitted = normalizeError(\n\t\t\t\t\tresponse,\n\t\t\t\t\tresponse.errorMessage ?? \"Assistant request exceeded the context window\",\n\t\t\t\t);\n\t\t\t\tif (current.generationContext.overflowRecoveryUsed || overflowPreparation === undefined) {\n\t\t\t\t\tfailure = providerError(source, committed);\n\t\t\t\t} else {\n\t\t\t\t\tconst structural: SummaryDecidingOperation = {\n\t\t\t\t\t\t...scope,\n\t\t\t\t\t\tat: \"summary.deciding\",\n\t\t\t\t\t\ttask: {\n\t\t\t\t\t\t\ttaskId: overflowPreparation.taskId,\n\t\t\t\t\t\t\treason: \"overflow\",\n\t\t\t\t\t\t\tboundary: {\n\t\t\t\t\t\t\t\tkind: \"resume_checkpoint\",\n\t\t\t\t\t\t\t\tresumeAfter: {\n\t\t\t\t\t\t\t\t\tcontinuation: { kind: \"need_assistant\", overflowRecoveryUsed: true },\n\t\t\t\t\t\t\t\t\ttriggerEntryId: current.generationContext.triggerEntryId,\n\t\t\t\t\t\t\t\t},\n\t\t\t\t\t\t\t},\n\t\t\t\t\t\t},\n\t\t\t\t\t};\n\t\t\t\t\tsettled = structural;\n\t\t\t\t}\n\t\t\t} else if (response.stopReason === \"deferred\") {\n\t\t\t\tif (current.at === \"assistant.effect_pending\") {\n\t\t\t\t\tif (deferredHandleIsValid(response, current)) {\n\t\t\t\t\t\tsettled = {\n\t\t\t\t\t\t\t...scope,\n\t\t\t\t\t\t\tat: \"deferred.suspended\",\n\t\t\t\t\t\t\tstepId: current.generationContext.stepId,\n\t\t\t\t\t\t\tsourceEntryId: responseEntryId,\n\t\t\t\t\t\t\tpoll: 0,\n\t\t\t\t\t\t\tconfiguration,\n\t\t\t\t\t\t\tstreamOptions: current.generationContext.streamOptions,\n\t\t\t\t\t\t};\n\t\t\t\t\t} else {\n\t\t\t\t\t\tcommitted = normalizeError(response, \"Provider returned an invalid deferred handle\");\n\t\t\t\t\t\tfailure = providerError(source, committed);\n\t\t\t\t\t}\n\t\t\t\t} else {\n\t\t\t\t\tsettled = {\n\t\t\t\t\t\t...scope,\n\t\t\t\t\t\tat: \"deferred.suspended\",\n\t\t\t\t\t\tstepId: current.stepId,\n\t\t\t\t\t\tsourceEntryId: responseEntryId,\n\t\t\t\t\t\tpoll: current.poll,\n\t\t\t\t\t\tconfiguration,\n\t\t\t\t\t\tstreamOptions: current.streamOptions,\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t} else if (response.stopReason === \"error\") {\n\t\t\t\tif (\n\t\t\t\t\tcurrent.at === \"assistant.effect_pending\" &&\n\t\t\t\t\t(options.recovery === true || isRetryableAssistantError(response)) &&\n\t\t\t\t\tcurrent.attempt < current.generationContext.retryPolicy.maxAttempts\n\t\t\t\t) {\n\t\t\t\t\tsettled = {\n\t\t\t\t\t\t...scope,\n\t\t\t\t\t\tat: \"assistant.retry_wait\",\n\t\t\t\t\t\tgenerationContext: current.generationContext,\n\t\t\t\t\t\tnextAttempt: current.attempt + 1,\n\t\t\t\t\t\tnotBefore: retryNotBefore(current.generationContext.retryPolicy.baseDelayMs, current.attempt),\n\t\t\t\t\t\terrorMessage: response.errorMessage ?? \"Assistant request failed\",\n\t\t\t\t\t};\n\t\t\t\t} else {\n\t\t\t\t\tfailure = providerError(source, response);\n\t\t\t\t}\n\t\t\t} else {\n\t\t\t\tconst calls = response.content.flatMap((content, sourceIndex) =>\n\t\t\t\t\tcontent.type === \"toolCall\" ? [{ sourceIndex }] : [],\n\t\t\t\t);\n\t\t\t\tif (calls.length !== 0) {\n\t\t\t\t\tconst timestamp = uuidV7Timestamp(responseEntryId);\n\t\t\t\t\tconst planned: ToolCall[] = calls.map(({ sourceIndex }) => ({\n\t\t\t\t\t\tstatus: \"planned\",\n\t\t\t\t\t\tsourceIndex,\n\t\t\t\t\t\tresultEntryId: lane.session.idGenerator.next(timestamp),\n\t\t\t\t\t}));\n\t\t\t\t\tsettled = {\n\t\t\t\t\t\t...scope,\n\t\t\t\t\t\tat: \"tools\",\n\t\t\t\t\t\tbatch: { assistantEntryId: responseEntryId, configuration, turnId, calls: planned },\n\t\t\t\t\t};\n\t\t\t\t} else if (response.stopReason === \"toolUse\") {\n\t\t\t\t\tcommitted = normalizeError(response, \"Provider reported tool use without any tool calls\");\n\t\t\t\t\tfailure = providerError(source, committed);\n\t\t\t\t} else {\n\t\t\t\t\tsettled = {\n\t\t\t\t\t\t...scope,\n\t\t\t\t\t\tat: \"checkpoint\",\n\t\t\t\t\t\tcontinuation: { kind: \"may_finish\", includeFinalAssistant: true },\n\t\t\t\t\t\ttriggerEntryId: responseEntryId,\n\t\t\t\t\t};\n\t\t\t\t}\n\t\t\t}\n\n\t\t\tconst responseEntry: NewEntry<MessageEntry> = {\n\t\t\t\tid: responseEntryId,\n\t\t\t\tparentId: state.tipId,\n\t\t\t\ttype: \"message\",\n\t\t\t\tmessage: committed,\n\t\t\t};\n\t\t\tconst usageRow: Omit<UsageRow, \"seq\"> = {\n\t\t\t\tid: intent.usageId,\n\t\t\t\tusage: committed.usage,\n\t\t\t\tentryId: responseEntryId,\n\t\t\t\tadjustment: false,\n\t\t\t};\n\t\t\tif (settled === undefined && failure === undefined) {\n\t\t\t\tthrow new SessionInvariantError(\"Response settlement has no durable disposition\");\n\t\t\t}\n\t\t\tconst record =\n\t\t\t\tfailure === undefined ? undefined : operationResultRecord(meta, \"failed\", responseEntryId, failure);\n\t\t\tconst cleanup =\n\t\t\t\trecord === undefined ? [] : await operationCleanupWrites(reader, drive.operationId, current, drive.context);\n\t\t\tconst writes = [\n\t\t\t\tinsertEntry(responseEntry),\n\t\t\t\tinsertUsage(usageRow),\n\t\t\t\tsetValue(branchTip(lane.name), responseEntryId),\n\t\t\t\t...(record === undefined\n\t\t\t\t\t? [deleteList(pendingAssistantFrames(drive.operationId, responseEntryId))]\n\t\t\t\t\t: cleanup),\n\t\t\t\t...(settled?.at === \"summary.deciding\" && overflowPreparation !== undefined\n\t\t\t\t\t? [\n\t\t\t\t\t\t\tsetValue(\n\t\t\t\t\t\t\t\toperationPreparation(drive.operationId, overflowPreparation.taskId),\n\t\t\t\t\t\t\t\toverflowPreparation.preparation,\n\t\t\t\t\t\t\t),\n\t\t\t\t\t\t]\n\t\t\t\t\t: []),\n\t\t\t];\n\t\t\tconst materializeEntry = (commit: CommitResult): MessageEntry => ({\n\t\t\t\t...responseEntry,\n\t\t\t\tseq: commit.seqs[0]!,\n\t\t\t\ttimestamp: commit.timestamp,\n\t\t\t});\n\t\t\tconst events = (commit: CommitResult): HarnessEvent[] => {\n\t\t\t\tconst entry = materializeEntry(commit);\n\t\t\t\tconst batch: HarnessEvent[] = [\n\t\t\t\t\t{ type: \"entry_added\", lane: lane.name, entry, ...options },\n\t\t\t\t\t{\n\t\t\t\t\t\ttype: \"usage\",\n\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\trow: { ...usageRow, seq: commit.seqs[1]! },\n\t\t\t\t\t\ttotals: commit.stats.usage,\n\t\t\t\t\t},\n\t\t\t\t];\n\t\t\t\tif (current.at === \"assistant.effect_pending\") {\n\t\t\t\t\tif (options.recovery !== true && current.attempt > 1 && settled?.at !== \"assistant.retry_wait\") {\n\t\t\t\t\t\tconst success = committed.stopReason !== \"error\" && committed.stopReason !== \"aborted\";\n\t\t\t\t\t\tbatch.push({\n\t\t\t\t\t\t\ttype: \"retry_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\tstep: turnId,\n\t\t\t\t\t\t\tattempt: current.attempt,\n\t\t\t\t\t\t\tsuccess,\n\t\t\t\t\t\t\t...(success\n\t\t\t\t\t\t\t\t? {}\n\t\t\t\t\t\t\t\t: {\n\t\t\t\t\t\t\t\t\t\tfinalError:\n\t\t\t\t\t\t\t\t\t\t\tcommitted.errorMessage ?? `Assistant request ended with ${committed.stopReason}`,\n\t\t\t\t\t\t\t\t\t}),\n\t\t\t\t\t\t});\n\t\t\t\t\t}\n\t\t\t\t\tif (options.recovery !== true && settled?.at === \"assistant.retry_wait\") {\n\t\t\t\t\t\tbatch.push({\n\t\t\t\t\t\t\ttype: \"retry_scheduled\",\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\tstep: turnId,\n\t\t\t\t\t\t\tattempt: settled.nextAttempt,\n\t\t\t\t\t\t\tmaxAttempts: settled.generationContext.retryPolicy.maxAttempts,\n\t\t\t\t\t\t\tdelayMs: retryDelay(current.generationContext.retryPolicy.baseDelayMs, current.attempt),\n\t\t\t\t\t\t\tnotBefore: settled.notBefore,\n\t\t\t\t\t\t\terrorMessage: settled.errorMessage,\n\t\t\t\t\t\t});\n\t\t\t\t\t}\n\t\t\t\t\tif (options.recovery !== true && settled?.at !== \"tools\" && settled?.at !== \"assistant.retry_wait\") {\n\t\t\t\t\t\tbatch.push({\n\t\t\t\t\t\t\ttype: \"turn_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\tturnId,\n\t\t\t\t\t\t\tmessage: committed,\n\t\t\t\t\t\t\ttoolResults: [],\n\t\t\t\t\t\t});\n\t\t\t\t\t}\n\t\t\t\t\tif (settled?.at === \"summary.deciding\") {\n\t\t\t\t\t\tbatch.push({\n\t\t\t\t\t\t\ttype: \"compaction_start\",\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: \"overflow\",\n\t\t\t\t\t\t\tstartedAt: commit.timestamp,\n\t\t\t\t\t\t});\n\t\t\t\t\t}\n\t\t\t\t} else if (settled?.at !== \"tools\") {\n\t\t\t\t\tbatch.push({\n\t\t\t\t\t\ttype: \"turn_end\",\n\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\tturnId,\n\t\t\t\t\t\tmessage: committed,\n\t\t\t\t\t\ttoolResults: [],\n\t\t\t\t\t\t...options,\n\t\t\t\t\t});\n\t\t\t\t}\n\t\t\t\tif (\n\t\t\t\t\t(current.at === \"deferred.effect_pending\" || options.recovery !== true) &&\n\t\t\t\t\tsettled?.at === \"deferred.suspended\" &&\n\t\t\t\t\tcommitted.deferred !== undefined\n\t\t\t\t) {\n\t\t\t\t\tbatch.push({\n\t\t\t\t\t\ttype: \"run_suspend\",\n\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\treason: \"deferred\",\n\t\t\t\t\t\tdeferred: committed.deferred,\n\t\t\t\t\t\tpoll: settled.poll,\n\t\t\t\t\t\t...(current.at === \"deferred.effect_pending\" ? options : {}),\n\t\t\t\t\t});\n\t\t\t\t}\n\t\t\t\tif (record !== undefined && failure !== undefined) {\n\t\t\t\t\tbatch.push({\n\t\t\t\t\t\ttype: \"run_end\",\n\t\t\t\t\t\tlane: lane.name,\n\t\t\t\t\t\trunId: drive.operationId,\n\t\t\t\t\t\tstatus: \"failed\",\n\t\t\t\t\t\terror: failure,\n\t\t\t\t\t\tfromTipId: meta.sourceTipId,\n\t\t\t\t\t\ttipId: responseEntryId,\n\t\t\t\t\t\tendedAt: record.endedAt,\n\t\t\t\t\t});\n\t\t\t\t}\n\t\t\t\treturn batch;\n\t\t\t};\n\t\t\tif (record !== undefined) {\n\t\t\t\treturn {\n\t\t\t\t\tkind: \"finish\",\n\t\t\t\t\twrites,\n\t\t\t\t\trecord,\n\t\t\t\t\tlane: { tipId: responseEntryId },\n\t\t\t\t\tmaterialize: () => ({ kind: \"settled\", outcome: record }) as const,\n\t\t\t\t\tevents,\n\t\t\t\t};\n\t\t\t}\n\t\t\tif (settled === undefined) throw new SessionInvariantError(\"Response settlement is missing its next state\");\n\t\t\treturn {\n\t\t\t\tkind: \"commit\",\n\t\t\t\twrites,\n\t\t\t\toperationState: settled,\n\t\t\t\tlane: { tipId: responseEntryId },\n\t\t\t\tmaterialize: () => ({ kind: \"continue\" }) as const,\n\t\t\t\tevents,\n\t\t\t};\n\t\t},\n\t\tdrive.context,\n\t);\n}\n"]}