mirror of
https://github.com/stablyai/orca.git
synced 2026-10-09 08:02:35 +00:00
feat: add canonical agent turn lifecycle reducer
This commit is contained in:
@@ -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<AgentTurnPhase, 'recovering'>
|
||||
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<AgentTurnEvidence, 'producerId' | 'eventId'>
|
||||
|
||||
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<AgentTurnOutcome, 'interrupted'>
|
||||
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[]
|
||||
}
|
||||
@@ -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<AgentTurnLifecycleEvent, { kind: 'turn-started' }>
|
||||
): 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<AgentTurnLifecycleEvent, { kind: 'current-turn-inventory' }>
|
||||
): 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<AgentTurnLifecycleEvent, { kind: 'work-started' }>
|
||||
): 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<AgentTurnLifecycleEvent, { kind: 'work-outcome-observed' }>
|
||||
): 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'
|
||||
}
|
||||
}
|
||||
@@ -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<AgentTurnWorkRecord, 'phase' | 'outcome'> {
|
||||
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
|
||||
}
|
||||
@@ -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<AgentTurnLifecycleEvent, { kind: 'turn-recovery-started' }>
|
||||
): 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<AgentTurnLifecycleEvent, { kind: 'turn-recovery-expired' }>
|
||||
): 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<AgentTurnLifecycleEvent, { kind: 'turn-recovery-abandoned' }>
|
||||
): 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
|
||||
}
|
||||
@@ -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<AgentTurnOutcome, 'interrupted'>,
|
||||
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
|
||||
}
|
||||
@@ -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<Candidate, 'owner' | 'evidence'>
|
||||
: 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<typeof state>, 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)
|
||||
})
|
||||
})
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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<string, unknown> {
|
||||
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<string, unknown>, 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<string, unknown> {
|
||||
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
|
||||
}
|
||||
}
|
||||
@@ -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'
|
||||
Reference in New Issue
Block a user