From aa077e534d9c414c3ddb9e80c2a99fe86dc0bed4 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Mon, 14 Sep 2026 17:57:49 -0700 Subject: [PATCH] feat: add canonical agent turn lifecycle reducer --- src/shared/agent-turn-lifecycle-contract.ts | 222 ++++++++++++ .../agent-turn-lifecycle-reducer-events.ts | 304 ++++++++++++++++ ...agent-turn-lifecycle-reducer-operations.ts | 297 ++++++++++++++++ .../agent-turn-lifecycle-reducer-recovery.ts | 98 +++++ ...gent-turn-lifecycle-reducer-transitions.ts | 271 ++++++++++++++ .../agent-turn-lifecycle-reducer.test.ts | 336 ++++++++++++++++++ src/shared/agent-turn-lifecycle-reducer.ts | 56 +++ src/shared/agent-turn-lifecycle-snapshot.ts | 43 +++ src/shared/agent-turn-lifecycle-state.ts | 261 ++++++++++++++ src/shared/agent-turn-lifecycle.ts | 48 +++ 10 files changed, 1936 insertions(+) create mode 100644 src/shared/agent-turn-lifecycle-contract.ts create mode 100644 src/shared/agent-turn-lifecycle-reducer-events.ts create mode 100644 src/shared/agent-turn-lifecycle-reducer-operations.ts create mode 100644 src/shared/agent-turn-lifecycle-reducer-recovery.ts create mode 100644 src/shared/agent-turn-lifecycle-reducer-transitions.ts create mode 100644 src/shared/agent-turn-lifecycle-reducer.test.ts create mode 100644 src/shared/agent-turn-lifecycle-reducer.ts create mode 100644 src/shared/agent-turn-lifecycle-snapshot.ts create mode 100644 src/shared/agent-turn-lifecycle-state.ts create mode 100644 src/shared/agent-turn-lifecycle.ts diff --git a/src/shared/agent-turn-lifecycle-contract.ts b/src/shared/agent-turn-lifecycle-contract.ts new file mode 100644 index 00000000000..386d3141eef --- /dev/null +++ b/src/shared/agent-turn-lifecycle-contract.ts @@ -0,0 +1,222 @@ +import type { AgentHookSource } from './agent-hook-relay' +import type { + AgentStatusExecutionAttachment, + AgentStatusRunId, + AgentStatusRunVerdict +} from './agent-status-run' +import type { AgentStatusExecutionScope } from './agent-status-subject' + +export type AgentTurnId = string +export type AgentTurnWorkId = string +export type AgentTurnDispatchId = string + +export type AgentTurnOwner = AgentStatusExecutionScope & { + runId: AgentStatusRunId + attachment: AgentStatusExecutionAttachment + provider: AgentHookSource +} + +export type AgentTurnEvidence = { + /** Stable within one producer so replay is idempotent. */ + eventId: string + /** Names the concrete hook, journal, inventory, or host observer. */ + producerId: string + observedAt: number + /** Opaque provider cursor retained for diagnostics and recovery handoff. */ + cursor?: string +} + +export type AgentTurnOutcome = 'completed' | 'failed' | 'interrupted' +export type AgentTurnPhase = 'active' | 'recovering' | 'settled' | 'unresolved' | 'abandoned' +export type AgentTurnInterruptState = 'none' | 'requested' | 'acknowledged' + +export type AgentTurnRecord = { + turnId: AgentTurnId + phase: AgentTurnPhase + outcome: AgentTurnOutcome | null + interrupt: AgentTurnInterruptState + /** Writing Escape/Ctrl+C is delivery evidence, not an interrupt acknowledgement. */ + interruptInputWrittenAt: number | null + startedAt: number | null + settledAt: number | null + lastEvidence: AgentTurnEvidence +} + +export type AgentTurnWorkKind = 'joined-child' | 'resident-background' + +export type AgentTurnWorkRecord = { + turnId: AgentTurnId + workId: AgentTurnWorkId + kind: AgentTurnWorkKind + phase: Exclude + outcome: AgentTurnOutcome | null + startedAt: number | null + settledAt: number | null + lastEvidence: AgentTurnEvidence +} + +export type AgentTurnDispatchReceipt = 'unobserved' | 'received' | 'rejected' +export type AgentTurnDispatchOutcome = AgentTurnOutcome | 'unresolved' | 'abandoned' + +export type AgentTurnDispatchRecord = { + dispatchId: AgentTurnDispatchId + turnId: AgentTurnId + receipt: AgentTurnDispatchReceipt + outcome: AgentTurnDispatchOutcome | null + settledAt: number | null + lastEvidence: AgentTurnEvidence +} + +export type AgentTurnRecoveryCustody = { + custodyId: string + turnId: AgentTurnId + deadlineAt: number + lastEvidence: AgentTurnEvidence +} + +export type AgentTurnIntegrityIssue = { + kind: 'capacity-overflow' | 'conflicting-outcome' | 'conflicting-dispatch-turn' + turnId?: AgentTurnId + workId?: AgentTurnWorkId + dispatchId?: AgentTurnDispatchId + observedAt: number +} + +export type AgentTurnAppliedEvent = Pick + +export type AgentTurnLifecycleState = { + /** Host-owned projection; structured session journals remain their durable source. */ + version: 1 + owner: AgentTurnOwner + currentTurnId: AgentTurnId | null + turns: AgentTurnRecord[] + work: AgentTurnWorkRecord[] + dispatches: AgentTurnDispatchRecord[] + recoveries: AgentTurnRecoveryCustody[] + executionVerdict: AgentStatusRunVerdict + integrityIssues: AgentTurnIntegrityIssue[] + appliedEvents: AgentTurnAppliedEvent[] +} + +export type AgentTurnInventoryWork = { + workId: AgentTurnWorkId + phase: 'active' | 'settled' | 'unresolved' + outcome?: AgentTurnOutcome + startedAt?: number + settledAt?: number +} + +export type AgentCurrentTurnInventory = { + /** A complete host inventory can re-derive active work and resolve omissions as unknown. */ + turnId: AgentTurnId + startedAt?: number + joinedChildren: AgentTurnInventoryWork[] + residentBackground: AgentTurnInventoryWork[] +} + +type AgentTurnEventBase = { + owner: AgentTurnOwner + evidence: AgentTurnEvidence +} + +export type AgentTurnLifecycleEvent = + | (AgentTurnEventBase & { kind: 'turn-started'; turnId: AgentTurnId; startedAt?: number }) + | (AgentTurnEventBase & { + kind: 'turn-outcome-observed' + turnId: AgentTurnId + outcome: Exclude + settledAt?: number + recordKind: 'event' | 'terminal-record' + }) + | (AgentTurnEventBase & { kind: 'turn-interrupt-requested'; turnId: AgentTurnId }) + | (AgentTurnEventBase & { + kind: 'turn-interrupt-input-written' + turnId: AgentTurnId + writtenAt?: number + }) + | (AgentTurnEventBase & { + kind: 'turn-interrupt-acknowledged' + turnId: AgentTurnId + settledAt?: number + }) + | (AgentTurnEventBase & { + kind: 'work-started' + turnId: AgentTurnId + workId: AgentTurnWorkId + workKind: AgentTurnWorkKind + startedAt?: number + }) + | (AgentTurnEventBase & { + kind: 'work-outcome-observed' + turnId: AgentTurnId + workId: AgentTurnWorkId + workKind: AgentTurnWorkKind + outcome: AgentTurnOutcome + settledAt?: number + }) + | (AgentTurnEventBase & { + kind: 'current-turn-inventory' + complete: true + currentTurn: AgentCurrentTurnInventory | null + }) + | (AgentTurnEventBase & { + kind: 'turn-recovery-started' + turnId: AgentTurnId + custodyId: string + deadlineAt: number + }) + | (AgentTurnEventBase & { + kind: 'turn-recovery-expired' + turnId: AgentTurnId + custodyId: string + }) + | (AgentTurnEventBase & { + kind: 'turn-recovery-abandoned' + turnId: AgentTurnId + custodyId: string + }) + | (AgentTurnEventBase & { + kind: 'execution-verdict-observed' + verdict: AgentStatusRunVerdict + }) + | (AgentTurnEventBase & { + kind: 'dispatch-associated' + dispatchId: AgentTurnDispatchId + turnId: AgentTurnId + }) + | (AgentTurnEventBase & { + kind: 'dispatch-received' + dispatchId: AgentTurnDispatchId + turnId: AgentTurnId + }) + | (AgentTurnEventBase & { + kind: 'dispatch-rejected' + dispatchId: AgentTurnDispatchId + turnId: AgentTurnId + }) + | (AgentTurnEventBase & { + kind: 'dispatch-abandoned' + dispatchId: AgentTurnDispatchId + turnId: AgentTurnId + }) + +export type AgentTurnCommittedOutcome = { + completionId: string + turnId: AgentTurnId + outcome: AgentTurnOutcome +} + +export type AgentTurnCommittedDispatch = { + settlementId: string + dispatchId: AgentTurnDispatchId + turnId: AgentTurnId + outcome: AgentTurnDispatchOutcome +} + +export type AgentTurnLifecycleReduction = { + state: AgentTurnLifecycleState + disposition: 'accepted' | 'duplicate' | 'ignored' + reason?: 'invalid-event' | 'owner-mismatch' | 'capacity' | 'conflict' | 'stale' + committedOutcomes: AgentTurnCommittedOutcome[] + committedDispatches: AgentTurnCommittedDispatch[] +} diff --git a/src/shared/agent-turn-lifecycle-reducer-events.ts b/src/shared/agent-turn-lifecycle-reducer-events.ts new file mode 100644 index 00000000000..2509cacca1c --- /dev/null +++ b/src/shared/agent-turn-lifecycle-reducer-events.ts @@ -0,0 +1,304 @@ +import type { + AgentTurnLifecycleEvent, + AgentTurnLifecycleReduction, + AgentTurnLifecycleState +} from './agent-turn-lifecycle-contract' +import { + addTurn, + addWork, + applyInventoryWork, + eventEvidence, + findTurn, + issue, + turnFromInventory, + unresolvedTurn +} from './agent-turn-lifecycle-reducer-operations' +import { + applyDispatch, + applyInterruptAcknowledgement, + applyOutcome +} from './agent-turn-lifecycle-reducer-transitions' +import { + recoveryAbandoned, + recoveryExpired, + recoveryStarted +} from './agent-turn-lifecycle-reducer-recovery' + +type Reason = AgentTurnLifecycleReduction['reason'] | undefined + +function turnStarted( + state: AgentTurnLifecycleState, + event: Extract +): Reason { + const existing = findTurn(state, event.turnId) + if (existing) { + if (existing.phase === 'active' || existing.phase === 'recovering') { + existing.startedAt = existing.startedAt ?? event.startedAt ?? event.evidence.observedAt + existing.lastEvidence = eventEvidence(event) + state.currentTurnId = event.turnId + state.recoveries = state.recoveries.filter((entry) => entry.turnId !== event.turnId) + } + return undefined + } + if ( + !addTurn(state, { + turnId: event.turnId, + phase: 'active', + outcome: null, + interrupt: 'none', + interruptInputWrittenAt: null, + startedAt: event.startedAt ?? event.evidence.observedAt, + settledAt: null, + lastEvidence: eventEvidence(event) + }) + ) { + issue(state, { + kind: 'capacity-overflow', + turnId: event.turnId, + observedAt: event.evidence.observedAt + }) + return 'capacity' + } + if (state.currentTurnId && state.currentTurnId !== event.turnId) { + unresolvedTurn(state, state.currentTurnId, event.evidence.observedAt) + } + state.currentTurnId = event.turnId + return undefined +} + +function inventory( + state: AgentTurnLifecycleState, + event: Extract +): Reason { + const priorCurrent = state.currentTurnId + if (priorCurrent && (!event.currentTurn || priorCurrent !== event.currentTurn.turnId)) { + unresolvedTurn(state, priorCurrent, event.evidence.observedAt) + } + if (event.currentTurn && !turnFromInventory(state, event.currentTurn, event)) { + issue(state, { + kind: 'capacity-overflow', + turnId: event.currentTurn.turnId, + observedAt: event.evidence.observedAt + }) + return 'capacity' + } + if (!event.currentTurn) { + state.currentTurnId = null + return undefined + } + const turn = findTurn(state, event.currentTurn.turnId) + if (!turn || turn.phase !== 'active') { + state.currentTurnId = null + return undefined + } + if ( + !applyInventoryWork( + state, + event.currentTurn.turnId, + event.currentTurn.joinedChildren, + 'joined-child', + event + ) || + !applyInventoryWork( + state, + event.currentTurn.turnId, + event.currentTurn.residentBackground, + 'resident-background', + event + ) + ) { + return 'capacity' + } + return undefined +} + +function workStarted( + state: AgentTurnLifecycleState, + event: Extract +): Reason { + const existing = state.work.find( + (item) => item.turnId === event.turnId && item.workId === event.workId + ) + if (existing) { + if (existing.kind !== event.workKind) { + issue(state, { + kind: 'conflicting-outcome', + turnId: event.turnId, + workId: event.workId, + observedAt: event.evidence.observedAt + }) + return 'conflict' + } + if (existing.phase === 'active') { + existing.lastEvidence = eventEvidence(event) + } + return undefined + } + if ( + !addWork(state, { + turnId: event.turnId, + workId: event.workId, + kind: event.workKind, + phase: 'active', + outcome: null, + startedAt: event.startedAt ?? event.evidence.observedAt, + settledAt: null, + lastEvidence: eventEvidence(event) + }) + ) { + issue(state, { + kind: 'capacity-overflow', + turnId: event.turnId, + workId: event.workId, + observedAt: event.evidence.observedAt + }) + return 'capacity' + } + return undefined +} + +function workOutcome( + state: AgentTurnLifecycleState, + event: Extract +): Reason { + const existing = state.work.find( + (item) => item.turnId === event.turnId && item.workId === event.workId + ) + if (existing?.phase === 'abandoned') { + return undefined + } + if (existing && existing.kind !== event.workKind) { + issue(state, { + kind: 'conflicting-outcome', + turnId: event.turnId, + workId: event.workId, + observedAt: event.evidence.observedAt + }) + return 'conflict' + } + if (existing && existing.outcome !== null && existing.outcome !== event.outcome) { + issue(state, { + kind: 'conflicting-outcome', + turnId: event.turnId, + workId: event.workId, + observedAt: event.evidence.observedAt + }) + return 'conflict' + } + if (existing) { + existing.phase = 'settled' + existing.outcome = event.outcome + existing.settledAt = event.settledAt ?? event.evidence.observedAt + existing.lastEvidence = eventEvidence(event) + return undefined + } + if ( + !addWork(state, { + turnId: event.turnId, + workId: event.workId, + kind: event.workKind, + phase: 'settled', + outcome: event.outcome, + startedAt: null, + settledAt: event.settledAt ?? event.evidence.observedAt, + lastEvidence: eventEvidence(event) + }) + ) { + issue(state, { + kind: 'capacity-overflow', + turnId: event.turnId, + workId: event.workId, + observedAt: event.evidence.observedAt + }) + return 'capacity' + } + return undefined +} + +export function reduceAgentTurnEvent( + state: AgentTurnLifecycleState, + event: AgentTurnLifecycleEvent +): Reason { + switch (event.kind) { + case 'turn-started': + return turnStarted(state, event) + case 'turn-outcome-observed': + return applyOutcome( + state, + event.turnId, + event.outcome, + event.settledAt ?? event.evidence.observedAt, + event, + event.recordKind === 'terminal-record' + ) + ? undefined + : event.recordKind === 'event' && !findTurn(state, event.turnId) + ? 'stale' + : 'conflict' + case 'turn-interrupt-requested': { + const turn = findTurn(state, event.turnId) + if (turn && turn.phase === 'active') { + turn.interrupt = 'requested' + turn.lastEvidence = eventEvidence(event) + } + return undefined + } + case 'turn-interrupt-input-written': { + const turn = findTurn(state, event.turnId) + if (turn && (turn.phase === 'active' || turn.phase === 'recovering')) { + turn.interruptInputWrittenAt = event.writtenAt ?? event.evidence.observedAt + turn.lastEvidence = eventEvidence(event) + } + return undefined + } + case 'turn-interrupt-acknowledged': + return applyInterruptAcknowledgement( + state, + event.turnId, + event.settledAt ?? event.evidence.observedAt, + event + ) + ? undefined + : 'stale' + case 'work-started': + return workStarted(state, event) + case 'work-outcome-observed': + return workOutcome(state, event) + case 'current-turn-inventory': + return inventory(state, event) + case 'turn-recovery-started': + return recoveryStarted(state, event) + case 'turn-recovery-expired': + return recoveryExpired(state, event) + case 'turn-recovery-abandoned': + return recoveryAbandoned(state, event) + case 'execution-verdict-observed': + if (state.executionVerdict !== 'exited') { + state.executionVerdict = event.verdict + } + if (event.verdict === 'exited') { + for (const turn of state.turns) { + if (turn.phase === 'active' || turn.phase === 'recovering') { + unresolvedTurn(state, turn.turnId, event.evidence.observedAt) + } + } + } + return undefined + case 'dispatch-associated': + return applyDispatch(state, event.dispatchId, event.turnId, 'unobserved', null, event) + ? undefined + : 'capacity' + case 'dispatch-received': + return applyDispatch(state, event.dispatchId, event.turnId, 'received', null, event) + ? undefined + : 'conflict' + case 'dispatch-rejected': + return applyDispatch(state, event.dispatchId, event.turnId, 'rejected', 'failed', event) + ? undefined + : 'conflict' + case 'dispatch-abandoned': + return applyDispatch(state, event.dispatchId, event.turnId, 'rejected', 'abandoned', event) + ? undefined + : 'conflict' + } +} diff --git a/src/shared/agent-turn-lifecycle-reducer-operations.ts b/src/shared/agent-turn-lifecycle-reducer-operations.ts new file mode 100644 index 00000000000..135d94dd41d --- /dev/null +++ b/src/shared/agent-turn-lifecycle-reducer-operations.ts @@ -0,0 +1,297 @@ +import type { + AgentCurrentTurnInventory, + AgentTurnEvidence, + AgentTurnIntegrityIssue, + AgentTurnLifecycleEvent, + AgentTurnDispatchRecord, + AgentTurnLifecycleState, + AgentTurnRecord, + AgentTurnWorkKind, + AgentTurnWorkRecord +} from './agent-turn-lifecycle-contract' +import { + AGENT_TURN_MAX_APPLIED_EVENTS, + AGENT_TURN_MAX_DISPATCHES, + AGENT_TURN_MAX_INTEGRITY_ISSUES, + AGENT_TURN_MAX_TURNS, + AGENT_TURN_MAX_WORK_ITEMS +} from './agent-turn-lifecycle-state' + +export const eventEvidence = (event: AgentTurnLifecycleEvent): AgentTurnEvidence => ({ + ...event.evidence +}) + +export function copyLifecycleState(state: AgentTurnLifecycleState): AgentTurnLifecycleState { + return { + ...state, + turns: state.turns.map((turn) => ({ ...turn, lastEvidence: { ...turn.lastEvidence } })), + work: state.work.map((item) => ({ ...item, lastEvidence: { ...item.lastEvidence } })), + dispatches: state.dispatches.map((dispatch) => ({ + ...dispatch, + lastEvidence: { ...dispatch.lastEvidence } + })), + recoveries: state.recoveries.map((recovery) => ({ + ...recovery, + lastEvidence: { ...recovery.lastEvidence } + })), + integrityIssues: state.integrityIssues.map((entry) => ({ ...entry })), + appliedEvents: state.appliedEvents.map((entry) => ({ ...entry })) + } +} + +export function findTurn( + state: AgentTurnLifecycleState, + turnId: string +): AgentTurnRecord | undefined { + return state.turns.find((turn) => turn.turnId === turnId) +} + +export function findWork( + state: AgentTurnLifecycleState, + turnId: string, + workId: string +): AgentTurnWorkRecord | undefined { + return state.work.find((item) => item.turnId === turnId && item.workId === workId) +} + +export function findDispatch( + state: AgentTurnLifecycleState, + dispatchId: string +): AgentTurnDispatchRecord | undefined { + return state.dispatches.find((dispatch) => dispatch.dispatchId === dispatchId) +} + +export function issue(state: AgentTurnLifecycleState, entry: AgentTurnIntegrityIssue): void { + if ( + state.integrityIssues.some( + (existing) => + existing.kind === entry.kind && + existing.turnId === entry.turnId && + existing.workId === entry.workId && + existing.dispatchId === entry.dispatchId + ) + ) { + return + } + state.integrityIssues.push(entry) + if (state.integrityIssues.length > AGENT_TURN_MAX_INTEGRITY_ISSUES) { + state.integrityIssues.splice(0, state.integrityIssues.length - AGENT_TURN_MAX_INTEGRITY_ISSUES) + } +} + +export function rememberEvent( + state: AgentTurnLifecycleState, + event: AgentTurnLifecycleEvent +): void { + state.appliedEvents.push({ + producerId: event.evidence.producerId, + eventId: event.evidence.eventId + }) + if (state.appliedEvents.length > AGENT_TURN_MAX_APPLIED_EVENTS) { + state.appliedEvents.splice(0, state.appliedEvents.length - AGENT_TURN_MAX_APPLIED_EVENTS) + } +} + +export function wasApplied( + state: AgentTurnLifecycleState, + event: AgentTurnLifecycleEvent +): boolean { + return state.appliedEvents.some( + (applied) => + applied.producerId === event.evidence.producerId && applied.eventId === event.evidence.eventId + ) +} + +export function addTurn(state: AgentTurnLifecycleState, turn: AgentTurnRecord): boolean { + if (state.turns.some((existing) => existing.turnId === turn.turnId)) { + return true + } + const removable = state.turns.findIndex( + (existing) => existing.phase === 'settled' || existing.phase === 'abandoned' + ) + if (state.turns.length >= AGENT_TURN_MAX_TURNS && removable === -1) { + return false + } + if (removable !== -1 && state.turns.length >= AGENT_TURN_MAX_TURNS) { + state.turns.splice(removable, 1) + } + state.turns.push(turn) + return true +} + +export function addWork(state: AgentTurnLifecycleState, work: AgentTurnWorkRecord): boolean { + if ( + state.work.some( + (existing) => existing.turnId === work.turnId && existing.workId === work.workId + ) + ) { + return true + } + const removable = state.work.findIndex( + (existing) => existing.phase === 'settled' || existing.phase === 'abandoned' + ) + if (state.work.length >= AGENT_TURN_MAX_WORK_ITEMS && removable === -1) { + return false + } + if (removable !== -1 && state.work.length >= AGENT_TURN_MAX_WORK_ITEMS) { + state.work.splice(removable, 1) + } + state.work.push(work) + return true +} + +export function addDispatch( + state: AgentTurnLifecycleState, + dispatch: AgentTurnDispatchRecord +): boolean { + if (state.dispatches.some((existing) => existing.dispatchId === dispatch.dispatchId)) { + return true + } + const removable = state.dispatches.findIndex((existing) => existing.outcome !== null) + if (state.dispatches.length >= AGENT_TURN_MAX_DISPATCHES && removable === -1) { + return false + } + if (removable !== -1 && state.dispatches.length >= AGENT_TURN_MAX_DISPATCHES) { + state.dispatches.splice(removable, 1) + } + state.dispatches.push(dispatch) + return true +} + +export function unresolvedTurn(state: AgentTurnLifecycleState, turnId: string, at: number): void { + const turn = findTurn(state, turnId) + if (turn && (turn.phase === 'active' || turn.phase === 'recovering')) { + turn.phase = 'unresolved' + turn.settledAt = at + turn.lastEvidence = { ...turn.lastEvidence, observedAt: at } + } + for (const item of state.work) { + if (item.turnId === turnId && item.phase === 'active') { + item.phase = 'unresolved' + item.settledAt = at + item.lastEvidence = { ...item.lastEvidence, observedAt: at } + } + } + if (state.currentTurnId === turnId) { + state.currentTurnId = null + } +} + +export function turnFromInventory( + state: AgentTurnLifecycleState, + inventory: AgentCurrentTurnInventory, + event: AgentTurnLifecycleEvent +): boolean { + const existing = findTurn(state, inventory.turnId) + if (existing?.phase === 'settled' || existing?.phase === 'abandoned') { + return true + } + if (existing) { + existing.phase = 'active' + existing.outcome = null + existing.interrupt = 'none' + existing.interruptInputWrittenAt = null + existing.settledAt = null + existing.startedAt = existing.startedAt ?? inventory.startedAt ?? event.evidence.observedAt + existing.lastEvidence = eventEvidence(event) + } else if ( + !addTurn(state, { + turnId: inventory.turnId, + phase: 'active', + outcome: null, + interrupt: 'none', + interruptInputWrittenAt: null, + startedAt: inventory.startedAt ?? event.evidence.observedAt, + settledAt: null, + lastEvidence: eventEvidence(event) + }) + ) { + return false + } + state.currentTurnId = inventory.turnId + state.recoveries = state.recoveries.filter((entry) => entry.turnId !== inventory.turnId) + return true +} + +function normaliseInventoryWork( + item: AgentCurrentTurnInventory['joinedChildren'][number] +): Pick { + if (item.phase === 'settled' && item.outcome === undefined) { + // A complete inventory without an outcome proves settlement, not success. + return { phase: 'unresolved', outcome: null } + } + return { + phase: item.phase, + outcome: item.phase === 'settled' ? (item.outcome ?? null) : null + } +} + +export function applyInventoryWork( + state: AgentTurnLifecycleState, + turnId: string, + items: AgentCurrentTurnInventory['joinedChildren'], + kind: AgentTurnWorkKind, + event: AgentTurnLifecycleEvent +): boolean { + const ids: string[] = [] + for (const item of items) { + ids.push(item.workId) + const existing = findWork(state, turnId, item.workId) + const normalised = normaliseInventoryWork(item) + if (existing) { + if (existing.kind !== kind) { + issue(state, { + kind: 'conflicting-outcome', + turnId, + workId: item.workId, + observedAt: event.evidence.observedAt + }) + continue + } + if (existing.phase !== 'settled' && existing.phase !== 'abandoned') { + existing.phase = normalised.phase + existing.outcome = normalised.outcome + existing.startedAt = existing.startedAt ?? item.startedAt ?? null + existing.settledAt = + normalised.phase === 'active' ? null : (item.settledAt ?? event.evidence.observedAt) + existing.lastEvidence = eventEvidence(event) + } + continue + } + if ( + !addWork(state, { + turnId, + workId: item.workId, + kind, + phase: normalised.phase, + outcome: normalised.outcome, + startedAt: item.startedAt ?? null, + settledAt: + normalised.phase === 'active' ? null : (item.settledAt ?? event.evidence.observedAt), + lastEvidence: eventEvidence(event) + }) + ) { + issue(state, { + kind: 'capacity-overflow', + turnId, + workId: item.workId, + observedAt: event.evidence.observedAt + }) + return false + } + } + for (const existing of state.work) { + if ( + existing.turnId === turnId && + existing.kind === kind && + existing.phase === 'active' && + !ids.includes(existing.workId) + ) { + existing.phase = 'unresolved' + existing.outcome = null + existing.settledAt = event.evidence.observedAt + existing.lastEvidence = eventEvidence(event) + } + } + return true +} diff --git a/src/shared/agent-turn-lifecycle-reducer-recovery.ts b/src/shared/agent-turn-lifecycle-reducer-recovery.ts new file mode 100644 index 00000000000..894fb296b60 --- /dev/null +++ b/src/shared/agent-turn-lifecycle-reducer-recovery.ts @@ -0,0 +1,98 @@ +import type { + AgentTurnLifecycleEvent, + AgentTurnLifecycleReduction, + AgentTurnLifecycleState +} from './agent-turn-lifecycle-contract' +import { + addTurn, + eventEvidence, + findTurn, + issue, + unresolvedTurn +} from './agent-turn-lifecycle-reducer-operations' +import { markRecovery } from './agent-turn-lifecycle-reducer-transitions' + +type Reason = AgentTurnLifecycleReduction['reason'] | undefined + +export function recoveryStarted( + state: AgentTurnLifecycleState, + event: Extract +): Reason { + const current = findTurn(state, event.turnId) + if (current?.phase === 'settled' || current?.phase === 'abandoned') { + return undefined + } + if (!markRecovery(state, event.turnId, event.custodyId, event.deadlineAt, event)) { + issue(state, { + kind: 'capacity-overflow', + turnId: event.turnId, + observedAt: event.evidence.observedAt + }) + return 'capacity' + } + const turn = findTurn(state, event.turnId) + if (turn) { + turn.phase = 'recovering' + turn.lastEvidence = eventEvidence(event) + return undefined + } + if ( + !addTurn(state, { + turnId: event.turnId, + phase: 'recovering', + outcome: null, + interrupt: 'none', + interruptInputWrittenAt: null, + startedAt: null, + settledAt: null, + lastEvidence: eventEvidence(event) + }) + ) { + issue(state, { + kind: 'capacity-overflow', + turnId: event.turnId, + observedAt: event.evidence.observedAt + }) + return 'capacity' + } + return undefined +} + +export function recoveryExpired( + state: AgentTurnLifecycleState, + event: Extract +): Reason { + const recovery = state.recoveries.find( + (entry) => entry.custodyId === event.custodyId && entry.turnId === event.turnId + ) + if (!recovery || event.evidence.observedAt < recovery.deadlineAt) { + return 'stale' + } + state.recoveries = state.recoveries.filter((entry) => entry.custodyId !== event.custodyId) + unresolvedTurn(state, event.turnId, event.evidence.observedAt) + return undefined +} + +export function recoveryAbandoned( + state: AgentTurnLifecycleState, + event: Extract +): Reason { + const recovery = state.recoveries.find( + (entry) => entry.custodyId === event.custodyId && entry.turnId === event.turnId + ) + if (!recovery) { + return 'stale' + } + state.recoveries = state.recoveries.filter((entry) => entry.custodyId !== event.custodyId) + const turn = findTurn(state, event.turnId) + if (turn && turn.phase !== 'settled') { + turn.phase = 'abandoned' + turn.outcome = null + turn.settledAt = event.evidence.observedAt + turn.lastEvidence = eventEvidence(event) + } + if (state.currentTurnId === event.turnId) { + state.currentTurnId = null + } + return undefined +} diff --git a/src/shared/agent-turn-lifecycle-reducer-transitions.ts b/src/shared/agent-turn-lifecycle-reducer-transitions.ts new file mode 100644 index 00000000000..ad056667bd5 --- /dev/null +++ b/src/shared/agent-turn-lifecycle-reducer-transitions.ts @@ -0,0 +1,271 @@ +import type { + AgentTurnCommittedDispatch, + AgentTurnCommittedOutcome, + AgentTurnDispatchRecord, + AgentTurnLifecycleEvent, + AgentTurnLifecycleState, + AgentTurnOutcome +} from './agent-turn-lifecycle-contract' +import { + agentTurnCompletionId, + agentTurnDispatchSettlementId, + AGENT_TURN_MAX_RECOVERIES +} from './agent-turn-lifecycle-state' +import { + addDispatch, + addTurn, + eventEvidence, + findDispatch, + findTurn, + issue +} from './agent-turn-lifecycle-reducer-operations' + +export function applyOutcome( + state: AgentTurnLifecycleState, + turnId: string, + outcome: Exclude, + settledAt: number, + event: AgentTurnLifecycleEvent, + allowMissedStart: boolean +): boolean { + const existing = findTurn(state, turnId) + if (existing?.outcome !== null && existing?.outcome !== undefined) { + if (existing.outcome !== outcome) { + issue(state, { kind: 'conflicting-outcome', turnId, observedAt: event.evidence.observedAt }) + return false + } + return true + } + if (existing?.phase === 'abandoned') { + return false + } + if (!existing && !allowMissedStart) { + return false + } + if (existing) { + existing.phase = 'settled' + existing.outcome = outcome + existing.settledAt = settledAt + existing.lastEvidence = eventEvidence(event) + } else if ( + !addTurn(state, { + turnId, + phase: 'settled', + outcome, + interrupt: 'none', + interruptInputWrittenAt: null, + startedAt: null, + settledAt, + lastEvidence: eventEvidence(event) + }) + ) { + return false + } + state.recoveries = state.recoveries.filter((entry) => entry.turnId !== turnId) + if (state.currentTurnId === turnId) { + state.currentTurnId = null + } + return true +} + +export function applyInterruptAcknowledgement( + state: AgentTurnLifecycleState, + turnId: string, + settledAt: number, + event: AgentTurnLifecycleEvent +): boolean { + const existing = findTurn(state, turnId) + if (existing?.outcome !== null && existing?.outcome !== undefined) { + return existing.outcome === 'interrupted' + } + if (existing?.phase === 'abandoned') { + return false + } + if (existing) { + existing.phase = 'settled' + existing.outcome = 'interrupted' + existing.interrupt = 'acknowledged' + existing.settledAt = settledAt + existing.lastEvidence = eventEvidence(event) + } else { + return false + } + state.recoveries = state.recoveries.filter((entry) => entry.turnId !== turnId) + if (state.currentTurnId === turnId) { + state.currentTurnId = null + } + return true +} + +export function applyDispatch( + state: AgentTurnLifecycleState, + dispatchId: string, + turnId: string, + receipt: AgentTurnDispatchRecord['receipt'], + outcome: AgentTurnDispatchRecord['outcome'], + event: AgentTurnLifecycleEvent +): boolean { + const existing = findDispatch(state, dispatchId) + if (existing && existing.turnId !== turnId) { + issue(state, { + kind: 'conflicting-dispatch-turn', + dispatchId, + observedAt: event.evidence.observedAt + }) + return false + } + if (existing) { + if (existing.outcome !== null) { + return existing.outcome === outcome || outcome === null + } + if ( + (existing.receipt === 'rejected' && receipt === 'received') || + (existing.receipt === 'received' && receipt === 'rejected') + ) { + return false + } + if (existing.receipt === 'unobserved') { + existing.receipt = receipt + } + if (outcome !== null) { + existing.outcome = outcome + existing.settledAt = event.evidence.observedAt + } + existing.lastEvidence = eventEvidence(event) + return true + } + return addDispatch(state, { + dispatchId, + turnId, + receipt, + outcome, + settledAt: outcome === null ? null : event.evidence.observedAt, + lastEvidence: eventEvidence(event) + }) +} + +function deriveDispatchOutcome( + state: AgentTurnLifecycleState, + dispatch: AgentTurnDispatchRecord +): AgentTurnDispatchRecord['outcome'] { + if (dispatch.receipt !== 'received' || dispatch.outcome !== null) { + return dispatch.outcome + } + const turn = findTurn(state, dispatch.turnId) + if (!turn) { + return null + } + if (turn.phase === 'abandoned') { + return 'abandoned' + } + if (turn.phase === 'unresolved' || turn.phase === 'recovering') { + return 'unresolved' + } + if (turn.phase !== 'settled' || turn.outcome === null) { + return null + } + if (turn.outcome !== 'completed') { + return turn.outcome + } + if ( + state.integrityIssues.some( + (entry) => entry.kind === 'capacity-overflow' && entry.turnId === dispatch.turnId + ) + ) { + return 'unresolved' + } + const joined = state.work.filter( + (item) => item.turnId === dispatch.turnId && item.kind === 'joined-child' + ) + if (joined.some((item) => item.phase === 'active')) { + return null + } + if (joined.some((item) => item.phase === 'unresolved' || item.phase === 'abandoned')) { + return 'unresolved' + } + const childFailure = joined.find( + (item) => item.outcome === 'failed' || item.outcome === 'interrupted' + ) + return childFailure?.outcome ?? 'completed' +} + +export function reconcileDispatches(state: AgentTurnLifecycleState, observedAt: number): void { + for (const dispatch of state.dispatches) { + const outcome = deriveDispatchOutcome(state, dispatch) + if (outcome !== null && dispatch.outcome === null) { + dispatch.outcome = outcome + dispatch.settledAt = observedAt + } + } +} + +export function committedOutcomes( + before: AgentTurnLifecycleState, + after: AgentTurnLifecycleState +): AgentTurnCommittedOutcome[] { + const committed: AgentTurnCommittedOutcome[] = [] + for (const turn of after.turns) { + if (turn.phase !== 'settled' || turn.outcome === null) { + continue + } + const prior = findTurn(before, turn.turnId) + if (prior?.phase === 'settled' && prior.outcome === turn.outcome) { + continue + } + committed.push({ + completionId: agentTurnCompletionId(after.owner, turn.turnId, turn.outcome), + turnId: turn.turnId, + outcome: turn.outcome + }) + } + return committed +} + +export function committedDispatches( + before: AgentTurnLifecycleState, + after: AgentTurnLifecycleState +): AgentTurnCommittedDispatch[] { + const committed: AgentTurnCommittedDispatch[] = [] + for (const dispatch of after.dispatches) { + if (dispatch.outcome === null) { + continue + } + const prior = findDispatch(before, dispatch.dispatchId) + if (prior?.outcome === dispatch.outcome) { + continue + } + committed.push({ + settlementId: agentTurnDispatchSettlementId( + after.owner, + dispatch.dispatchId, + dispatch.turnId + ), + dispatchId: dispatch.dispatchId, + turnId: dispatch.turnId, + outcome: dispatch.outcome + }) + } + return committed +} + +export function markRecovery( + state: AgentTurnLifecycleState, + turnId: string, + custodyId: string, + deadlineAt: number, + event: AgentTurnLifecycleEvent +): boolean { + if (state.recoveries.some((entry) => entry.custodyId === custodyId)) { + return true + } + if (state.recoveries.length >= AGENT_TURN_MAX_RECOVERIES) { + return false + } + state.recoveries.push({ + custodyId, + turnId, + deadlineAt, + lastEvidence: eventEvidence(event) + }) + return true +} diff --git a/src/shared/agent-turn-lifecycle-reducer.test.ts b/src/shared/agent-turn-lifecycle-reducer.test.ts new file mode 100644 index 00000000000..f87b801171b --- /dev/null +++ b/src/shared/agent-turn-lifecycle-reducer.test.ts @@ -0,0 +1,336 @@ +import { describe, expect, it } from 'vitest' +import type { AgentTurnLifecycleEvent, AgentTurnOwner } from './agent-turn-lifecycle-contract' +import { + createAgentTurnLifecycleState, + isAgentTurnLifecycleEvent, + reduceAgentTurnLifecycle, + readAgentTurnLifecycleSnapshot +} from './agent-turn-lifecycle' + +const owner: AgentTurnOwner = { + executionHostId: 'local', + wslDistro: null, + workspaceId: 'workspace-c1', + workspaceKind: 'folder', + runId: 'run-c1', + attachment: { executionId: 'execution-c1' }, + provider: 'claude' +} + +let eventNumber = 0 +type EventInput = AgentTurnLifecycleEvent extends infer Candidate + ? Candidate extends AgentTurnLifecycleEvent + ? Omit + : never + : never + +function event(input: EventInput): AgentTurnLifecycleEvent { + eventNumber += 1 + // oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: EventInput is a discriminated union with owner/evidence intentionally omitted; this constructor restores both fields. + return { + ...input, + owner, + evidence: { + eventId: `event-${eventNumber}`, + producerId: 'c1-test', + observedAt: eventNumber + } + } as AgentTurnLifecycleEvent +} + +function state() { + eventNumber = 0 + return createAgentTurnLifecycleState(owner) +} + +function apply(current: ReturnType, nextEvent: AgentTurnLifecycleEvent) { + return reduceAgentTurnLifecycle(current, nextEvent) +} + +describe('canonical agent turn lifecycle reducer', () => { + it('folds child-before-parent completion without letting a child settle the root', () => { + let current = state() + current = apply( + current, + event({ + kind: 'work-outcome-observed', + turnId: 'turn-1', + workId: 'child-1', + workKind: 'joined-child', + outcome: 'completed' + }) + ).state + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-1' })).state + expect(current.work[0]).toMatchObject({ turnId: 'turn-1', workId: 'child-1', phase: 'settled' }) + expect(current.turns[0]).toMatchObject({ turnId: 'turn-1', phase: 'active', outcome: null }) + }) + + it('keeps a newer current turn when a late completion arrives for the prior turn', () => { + let current = state() + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-1' })).state + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-2' })).state + const late = apply( + current, + event({ + kind: 'turn-outcome-observed', + turnId: 'turn-1', + outcome: 'completed', + recordKind: 'event' + }) + ) + expect(late.state.currentTurnId).toBe('turn-2') + expect(late.state.turns).toEqual( + expect.arrayContaining([ + expect.objectContaining({ turnId: 'turn-1', phase: 'settled', outcome: 'completed' }), + expect.objectContaining({ turnId: 'turn-2', phase: 'active' }) + ]) + ) + }) + + it('does not gate root dispatch settlement on resident background work', () => { + let current = state() + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-1' })).state + current = apply( + current, + event({ + kind: 'work-started', + turnId: 'turn-1', + workId: 'monitor-1', + workKind: 'resident-background' + }) + ).state + current = apply( + current, + event({ kind: 'dispatch-associated', dispatchId: 'dispatch-1', turnId: 'turn-1' }) + ).state + current = apply( + current, + event({ kind: 'dispatch-received', dispatchId: 'dispatch-1', turnId: 'turn-1' }) + ).state + const settled = apply( + current, + event({ + kind: 'turn-outcome-observed', + turnId: 'turn-1', + outcome: 'completed', + recordKind: 'event' + }) + ) + expect(settled.committedDispatches).toEqual([ + expect.objectContaining({ dispatchId: 'dispatch-1', outcome: 'completed' }) + ]) + expect(settled.state.work).toEqual([ + expect.objectContaining({ kind: 'resident-background', phase: 'active' }) + ]) + }) + + it('allows an attributable terminal record to recover a missed start, but not an ordinary end', () => { + const ordinary = apply( + state(), + event({ + kind: 'turn-outcome-observed', + turnId: 'turn-missed', + outcome: 'completed', + recordKind: 'event' + }) + ) + expect(ordinary.disposition).toBe('ignored') + expect(ordinary.reason).toBe('stale') + const recovered = apply( + state(), + event({ + kind: 'turn-outcome-observed', + turnId: 'turn-missed', + outcome: 'completed', + recordKind: 'terminal-record' + }) + ) + expect(recovered.committedOutcomes).toEqual([ + expect.objectContaining({ turnId: 'turn-missed', outcome: 'completed' }) + ]) + expect(recovered.state.turns[0]).toMatchObject({ startedAt: null, phase: 'settled' }) + }) + + it('requires dispatch receipt before settling a dispatch', () => { + let current = state() + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-1' })).state + current = apply( + current, + event({ kind: 'dispatch-associated', dispatchId: 'dispatch-1', turnId: 'turn-1' }) + ).state + current = apply( + current, + event({ + kind: 'turn-outcome-observed', + turnId: 'turn-1', + outcome: 'completed', + recordKind: 'event' + }) + ).state + expect(current.dispatches[0].outcome).toBeNull() + const received = apply( + current, + event({ kind: 'dispatch-received', dispatchId: 'dispatch-1', turnId: 'turn-1' }) + ) + expect(received.committedDispatches).toEqual([ + expect.objectContaining({ dispatchId: 'dispatch-1', outcome: 'completed' }) + ]) + }) + + it('keeps interrupt request separate from acknowledgement', () => { + let current = state() + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-1' })).state + current = apply(current, event({ kind: 'turn-interrupt-requested', turnId: 'turn-1' })).state + expect(current.turns[0]).toMatchObject({ + phase: 'active', + outcome: null, + interrupt: 'requested' + }) + current = apply( + current, + event({ kind: 'turn-interrupt-input-written', turnId: 'turn-1' }) + ).state + expect(current.turns[0]).toMatchObject({ + phase: 'active', + outcome: null, + interruptInputWrittenAt: 3 + }) + const acknowledged = apply( + current, + event({ kind: 'turn-interrupt-acknowledged', turnId: 'turn-1' }) + ) + expect(acknowledged.committedOutcomes).toEqual([ + expect.objectContaining({ turnId: 'turn-1', outcome: 'interrupted' }) + ]) + expect(acknowledged.state.turns[0]).toMatchObject({ + phase: 'settled', + interrupt: 'acknowledged' + }) + }) + + it('marks active turns unresolved when execution exits without declaring success', () => { + let current = state() + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-1' })).state + const exited = apply(current, event({ kind: 'execution-verdict-observed', verdict: 'exited' })) + expect(exited.committedOutcomes).toEqual([]) + expect(exited.state.executionVerdict).toBe('exited') + expect(exited.state.turns[0]).toMatchObject({ phase: 'unresolved', outcome: null }) + }) + + it('reconciles a complete current-turn inventory and preserves uncertainty', () => { + let current = state() + current = apply(current, event({ kind: 'turn-started', turnId: 'turn-1' })).state + current = apply( + current, + event({ + kind: 'work-started', + turnId: 'turn-1', + workId: 'child-1', + workKind: 'joined-child' + }) + ).state + const inventory = apply( + current, + event({ + kind: 'current-turn-inventory', + complete: true, + currentTurn: { + turnId: 'turn-1', + joinedChildren: [], + residentBackground: [{ workId: 'monitor-1', phase: 'settled' }] + } + }) + ) + expect(inventory.state.work).toEqual( + expect.arrayContaining([ + expect.objectContaining({ workId: 'child-1', phase: 'unresolved', outcome: null }), + expect.objectContaining({ workId: 'monitor-1', phase: 'unresolved', outcome: null }) + ]) + ) + expect(readAgentTurnLifecycleSnapshot(inventory.state).currentTurn?.turnId).toBe('turn-1') + }) + + it('bounds recovery custody and supports expiry versus explicit abandon', () => { + let current = state() + current = apply( + current, + event({ + kind: 'turn-recovery-started', + turnId: 'turn-1', + custodyId: 'custody-1', + deadlineAt: 20 + }) + ).state + const early = apply( + current, + event({ kind: 'turn-recovery-expired', turnId: 'turn-1', custodyId: 'custody-1' }) + ) + expect(early.disposition).toBe('ignored') + expect(early.reason).toBe('stale') + const expiryEvent = event({ + kind: 'turn-recovery-expired', + turnId: 'turn-1', + custodyId: 'custody-1' + }) + const expired = apply(early.state, { + ...expiryEvent, + evidence: { ...expiryEvent.evidence, observedAt: 20 } + }) + expect(expired.state.recoveries).toHaveLength(0) + expect(expired.state.turns[0]).toMatchObject({ phase: 'unresolved', outcome: null }) + const abandoned = apply( + expired.state, + event({ + kind: 'turn-recovery-started', + turnId: 'turn-2', + custodyId: 'custody-2', + deadlineAt: 30 + }) + ) + const released = apply( + abandoned.state, + event({ kind: 'turn-recovery-abandoned', turnId: 'turn-2', custodyId: 'custody-2' }) + ) + expect(released.state.turns).toEqual( + expect.arrayContaining([ + expect.objectContaining({ turnId: 'turn-2', phase: 'abandoned', outcome: null }) + ]) + ) + }) + + it('deduplicates replay and keeps a stable completion identity', () => { + const first = apply( + state(), + event({ + kind: 'turn-outcome-observed', + turnId: 'turn-1', + outcome: 'failed', + recordKind: 'terminal-record' + }) + ) + const replay = apply(first.state, { + ...event({ + kind: 'turn-outcome-observed', + turnId: 'turn-1', + outcome: 'failed', + recordKind: 'terminal-record' + }), + evidence: { ...first.state.appliedEvents[0], observedAt: first.state.turns[0].settledAt ?? 1 } + }) + expect(first.committedOutcomes[0].completionId).toBeTruthy() + expect(replay.disposition).toBe('duplicate') + expect(replay.committedOutcomes).toEqual([]) + }) + + it('rejects malformed anonymous lifecycle identity before reducing', () => { + const malformed = { + kind: 'turn-outcome-observed', + turnId: '', + outcome: 'completed', + recordKind: 'event', + owner, + evidence: { eventId: 'event-malformed', producerId: 'c1-test', observedAt: 1 } + } + expect(isAgentTurnLifecycleEvent(malformed)).toBe(false) + }) +}) diff --git a/src/shared/agent-turn-lifecycle-reducer.ts b/src/shared/agent-turn-lifecycle-reducer.ts new file mode 100644 index 00000000000..a39e17f57af --- /dev/null +++ b/src/shared/agent-turn-lifecycle-reducer.ts @@ -0,0 +1,56 @@ +import type { + AgentTurnLifecycleEvent, + AgentTurnLifecycleReduction, + AgentTurnLifecycleState +} from './agent-turn-lifecycle-contract' +import { + agentTurnOwnersEqual, + isAgentTurnEvidence, + isAgentTurnLifecycleEvent, + isAgentTurnOwner +} from './agent-turn-lifecycle-state' +import { + copyLifecycleState, + rememberEvent, + wasApplied +} from './agent-turn-lifecycle-reducer-operations' +import { + committedDispatches, + committedOutcomes, + reconcileDispatches +} from './agent-turn-lifecycle-reducer-transitions' +import { reduceAgentTurnEvent } from './agent-turn-lifecycle-reducer-events' + +function ignored( + state: AgentTurnLifecycleState, + reason: AgentTurnLifecycleReduction['reason'] +): AgentTurnLifecycleReduction { + return { state, disposition: 'ignored', reason, committedOutcomes: [], committedDispatches: [] } +} + +/** Applies one host-attributed fact; provider adapters must not reduce status themselves. */ +export function reduceAgentTurnLifecycle( + state: AgentTurnLifecycleState, + event: AgentTurnLifecycleEvent +): AgentTurnLifecycleReduction { + if (!isAgentTurnLifecycleEvent(event) || !isAgentTurnEvidence(event.evidence)) { + return ignored(state, 'invalid-event') + } + if (!isAgentTurnOwner(state.owner) || !agentTurnOwnersEqual(state.owner, event.owner)) { + return ignored(state, 'owner-mismatch') + } + if (wasApplied(state, event)) { + return { state, disposition: 'duplicate', committedOutcomes: [], committedDispatches: [] } + } + const next = copyLifecycleState(state) + rememberEvent(next, event) + const reason = reduceAgentTurnEvent(next, event) + reconcileDispatches(next, event.evidence.observedAt) + return { + state: next, + disposition: reason ? 'ignored' : 'accepted', + ...(reason ? { reason } : {}), + committedOutcomes: committedOutcomes(state, next), + committedDispatches: committedDispatches(state, next) + } +} diff --git a/src/shared/agent-turn-lifecycle-snapshot.ts b/src/shared/agent-turn-lifecycle-snapshot.ts new file mode 100644 index 00000000000..87e32f6edf4 --- /dev/null +++ b/src/shared/agent-turn-lifecycle-snapshot.ts @@ -0,0 +1,43 @@ +import type { + AgentTurnDispatchRecord, + AgentTurnLifecycleState, + AgentTurnRecord, + AgentTurnWorkRecord +} from './agent-turn-lifecycle-contract' + +export type AgentTurnLifecycleSnapshot = Readonly<{ + owner: AgentTurnLifecycleState['owner'] + currentTurnId: AgentTurnLifecycleState['currentTurnId'] + currentTurn: AgentTurnRecord | null + turns: readonly AgentTurnRecord[] + joinedChildren: readonly AgentTurnWorkRecord[] + residentBackground: readonly AgentTurnWorkRecord[] + dispatches: readonly AgentTurnDispatchRecord[] + recoveries: AgentTurnLifecycleState['recoveries'] + executionVerdict: AgentTurnLifecycleState['executionVerdict'] + integrityIssues: AgentTurnLifecycleState['integrityIssues'] +}> + +export function readAgentTurnLifecycleSnapshot( + state: AgentTurnLifecycleState +): AgentTurnLifecycleSnapshot { + const currentTurn = state.currentTurnId + ? (state.turns.find((turn) => turn.turnId === state.currentTurnId) ?? null) + : null + return { + owner: state.owner, + currentTurnId: state.currentTurnId, + currentTurn, + turns: state.turns, + joinedChildren: state.work.filter((item) => item.kind === 'joined-child'), + residentBackground: state.work.filter((item) => item.kind === 'resident-background'), + dispatches: state.dispatches, + recoveries: state.recoveries, + executionVerdict: state.executionVerdict, + integrityIssues: state.integrityIssues + } +} + +export function isAgentTurnDispatchSettled(dispatch: AgentTurnDispatchRecord): boolean { + return dispatch.outcome !== null +} diff --git a/src/shared/agent-turn-lifecycle-state.ts b/src/shared/agent-turn-lifecycle-state.ts new file mode 100644 index 00000000000..906cf7c0f81 --- /dev/null +++ b/src/shared/agent-turn-lifecycle-state.ts @@ -0,0 +1,261 @@ +import { isAgentHookSource } from './agent-hook-relay' +import { isAgentStatusExecutionId, isAgentStatusRunId } from './agent-status-run' +import { parseAgentStatusExecutionScope } from './agent-status-subject' +import type { + AgentCurrentTurnInventory, + AgentTurnEvidence, + AgentTurnLifecycleEvent, + AgentTurnLifecycleState, + AgentTurnOwner, + AgentTurnOutcome +} from './agent-turn-lifecycle-contract' + +export const AGENT_TURN_MAX_TURNS = 128 +export const AGENT_TURN_MAX_WORK_ITEMS = 512 +export const AGENT_TURN_MAX_DISPATCHES = 256 +export const AGENT_TURN_MAX_RECOVERIES = 128 +export const AGENT_TURN_MAX_APPLIED_EVENTS = 1024 +export const AGENT_TURN_MAX_INTEGRITY_ISSUES = 128 + +const MAX_LIFECYCLE_ID_LENGTH = 512 + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +function isFiniteTimestamp(value: unknown): value is number { + return typeof value === 'number' && Number.isFinite(value) && value >= 0 +} + +export function isAgentTurnLifecycleId(value: unknown): value is string { + if ( + typeof value !== 'string' || + value.length === 0 || + value.length > MAX_LIFECYCLE_ID_LENGTH || + value !== value.trim() + ) { + return false + } + for (let index = 0; index < value.length; index += 1) { + const code = value.charCodeAt(index) + if (code <= 0x1f || code === 0x7f) { + return false + } + } + return true +} + +export function isAgentTurnEvidence(value: unknown): value is AgentTurnEvidence { + if (!isRecord(value)) { + return false + } + return ( + isAgentTurnLifecycleId(value.eventId) && + isAgentTurnLifecycleId(value.producerId) && + isFiniteTimestamp(value.observedAt) && + (value.cursor === undefined || isAgentTurnLifecycleId(value.cursor)) + ) +} + +export function isAgentTurnOwner(value: unknown): value is AgentTurnOwner { + if (!isRecord(value) || !isRecord(value.attachment)) { + return false + } + const scope = { + executionHostId: value.executionHostId, + wslDistro: value.wslDistro, + workspaceId: value.workspaceId, + workspaceKind: value.workspaceKind + } + return ( + parseAgentStatusExecutionScope(scope) !== null && + isAgentStatusRunId(value.runId) && + isAgentStatusExecutionId(value.attachment.executionId) && + isAgentHookSource(value.provider) + ) +} + +export function agentTurnOwnersEqual(left: AgentTurnOwner, right: AgentTurnOwner): boolean { + return ( + left.executionHostId === right.executionHostId && + left.wslDistro === right.wslDistro && + left.workspaceId === right.workspaceId && + left.workspaceKind === right.workspaceKind && + left.runId === right.runId && + left.attachment.executionId === right.attachment.executionId && + left.provider === right.provider + ) +} + +export function createAgentTurnLifecycleState(owner: AgentTurnOwner): AgentTurnLifecycleState { + if (!isAgentTurnOwner(owner)) { + throw new Error('Invalid agent turn owner') + } + return { + version: 1, + owner, + currentTurnId: null, + turns: [], + work: [], + dispatches: [], + recoveries: [], + executionVerdict: 'unverifiable', + integrityIssues: [], + appliedEvents: [] + } +} + +function completionOwnerTuple(owner: AgentTurnOwner): readonly (string | null)[] { + return [ + owner.executionHostId, + owner.wslDistro, + owner.workspaceId, + owner.workspaceKind, + owner.runId, + owner.attachment.executionId, + owner.provider + ] +} + +export function agentTurnCompletionId( + owner: AgentTurnOwner, + turnId: string, + outcome: AgentTurnOutcome +): string { + return `agent-turn-completion-v1:${JSON.stringify([...completionOwnerTuple(owner), turnId, outcome])}` +} + +export function agentTurnDispatchSettlementId( + owner: AgentTurnOwner, + dispatchId: string, + turnId: string +): string { + return `agent-turn-dispatch-v1:${JSON.stringify([ + ...completionOwnerTuple(owner), + dispatchId, + turnId + ])}` +} + +function isTurnOutcome(value: unknown): value is AgentTurnOutcome { + return value === 'completed' || value === 'failed' || value === 'interrupted' +} + +function isOptionalTimestamp(record: Record, key: string): boolean { + return record[key] === undefined || isFiniteTimestamp(record[key]) +} + +function isInventoryWork( + value: unknown +): value is AgentCurrentTurnInventory['joinedChildren'][number] { + if (!isRecord(value) || !isAgentTurnLifecycleId(value.workId)) { + return false + } + if (value.phase !== 'active' && value.phase !== 'settled' && value.phase !== 'unresolved') { + return false + } + if (!isOptionalTimestamp(value, 'startedAt') || !isOptionalTimestamp(value, 'settledAt')) { + return false + } + if (value.outcome !== undefined && !isTurnOutcome(value.outcome)) { + return false + } + if ((value.phase === 'active' || value.phase === 'unresolved') && value.outcome !== undefined) { + return false + } + return true +} + +function isInventory(value: unknown): value is AgentCurrentTurnInventory { + if (!isRecord(value) || !isAgentTurnLifecycleId(value.turnId)) { + return false + } + if (!isOptionalTimestamp(value, 'startedAt')) { + return false + } + if ( + !Array.isArray(value.joinedChildren) || + !Array.isArray(value.residentBackground) || + value.joinedChildren.length > AGENT_TURN_MAX_WORK_ITEMS || + value.residentBackground.length > AGENT_TURN_MAX_WORK_ITEMS + ) { + return false + } + return ( + value.joinedChildren.every(isInventoryWork) && value.residentBackground.every(isInventoryWork) + ) +} + +function hasEventBase(value: unknown): value is Record { + return isRecord(value) && isAgentTurnOwner(value.owner) && isAgentTurnEvidence(value.evidence) +} + +export function isAgentTurnLifecycleEvent(value: unknown): value is AgentTurnLifecycleEvent { + if (!hasEventBase(value)) { + return false + } + const kind = value.kind + if ( + !isAgentTurnLifecycleId(value.turnId) && + kind !== 'current-turn-inventory' && + kind !== 'execution-verdict-observed' + ) { + return false + } + switch (kind) { + case 'turn-started': + return isOptionalTimestamp(value, 'startedAt') + case 'turn-outcome-observed': + return ( + (value.outcome === 'completed' || value.outcome === 'failed') && + (value.recordKind === 'event' || value.recordKind === 'terminal-record') && + isOptionalTimestamp(value, 'settledAt') + ) + case 'turn-interrupt-requested': + return true + case 'turn-interrupt-input-written': + return isOptionalTimestamp(value, 'writtenAt') + case 'turn-interrupt-acknowledged': + return isOptionalTimestamp(value, 'settledAt') + case 'work-started': + return ( + isAgentTurnLifecycleId(value.workId) && + (value.workKind === 'joined-child' || value.workKind === 'resident-background') && + isOptionalTimestamp(value, 'startedAt') + ) + case 'work-outcome-observed': + return ( + isAgentTurnLifecycleId(value.workId) && + (value.workKind === 'joined-child' || value.workKind === 'resident-background') && + isTurnOutcome(value.outcome) && + isOptionalTimestamp(value, 'settledAt') + ) + case 'current-turn-inventory': + return ( + value.complete === true && (value.currentTurn === null || isInventory(value.currentTurn)) + ) + case 'turn-recovery-started': + if (!isAgentTurnEvidence(value.evidence)) { + return false + } + return ( + isAgentTurnLifecycleId(value.custodyId) && + isFiniteTimestamp(value.deadlineAt) && + value.deadlineAt > value.evidence.observedAt + ) + case 'turn-recovery-expired': + case 'turn-recovery-abandoned': + return isAgentTurnLifecycleId(value.custodyId) + case 'execution-verdict-observed': + return ( + value.verdict === 'live' || value.verdict === 'unverifiable' || value.verdict === 'exited' + ) + case 'dispatch-associated': + case 'dispatch-received': + case 'dispatch-rejected': + case 'dispatch-abandoned': + return isAgentTurnLifecycleId(value.dispatchId) + default: + return false + } +} diff --git a/src/shared/agent-turn-lifecycle.ts b/src/shared/agent-turn-lifecycle.ts new file mode 100644 index 00000000000..56f9c1b5296 --- /dev/null +++ b/src/shared/agent-turn-lifecycle.ts @@ -0,0 +1,48 @@ +export type { + AgentCurrentTurnInventory, + AgentTurnAppliedEvent, + AgentTurnCommittedDispatch, + AgentTurnCommittedOutcome, + AgentTurnDispatchOutcome, + AgentTurnDispatchReceipt, + AgentTurnDispatchRecord, + AgentTurnEvidence, + AgentTurnId, + AgentTurnInterruptState, + AgentTurnIntegrityIssue, + AgentTurnLifecycleEvent, + AgentTurnLifecycleReduction, + AgentTurnLifecycleState, + AgentTurnOwner, + AgentTurnOutcome, + AgentTurnPhase, + AgentTurnRecord, + AgentTurnRecoveryCustody, + AgentTurnDispatchId, + AgentTurnWorkId, + AgentTurnWorkKind, + AgentTurnWorkRecord +} from './agent-turn-lifecycle-contract' +export { reduceAgentTurnLifecycle } from './agent-turn-lifecycle-reducer' +export { + agentTurnCompletionId, + agentTurnDispatchSettlementId, + createAgentTurnLifecycleState, + isAgentTurnEvidence, + isAgentTurnLifecycleEvent, + isAgentTurnLifecycleId, + isAgentTurnOwner +} from './agent-turn-lifecycle-state' +export { + AGENT_TURN_MAX_APPLIED_EVENTS, + AGENT_TURN_MAX_DISPATCHES, + AGENT_TURN_MAX_INTEGRITY_ISSUES, + AGENT_TURN_MAX_RECOVERIES, + AGENT_TURN_MAX_TURNS, + AGENT_TURN_MAX_WORK_ITEMS +} from './agent-turn-lifecycle-state' +export { + isAgentTurnDispatchSettled, + readAgentTurnLifecycleSnapshot, + type AgentTurnLifecycleSnapshot +} from './agent-turn-lifecycle-snapshot'