refactor(native-chat): structured chat failures always reach the diagnostics log (#24312)

* refactor(native-chat): give the structured chat host one required logger

The structured chat runtime took an optional onError callback that the
desktop never passed, so a late dispatch settlement, an unanswered-dispatch
release, a journal event-sink write and a provider lifecycle delivery that
failed were dropped with no trace. Other host failures went to scattered
console.warn calls, which reach nothing in a packaged desktop build.

The runtime and host now take one required logger (warn/error with a scope
and fields). The production logger writes each entry as a failed span to
<userData>/logs/main.trace.ndjson, which the diagnostic bundle collects, and
to the console (stderr under a supervised headless host). The runtime and the
host wrap it so a logger that throws never fails what it reports, and the
install refuses without one. Sites that deliberately kept a recovery-capsule
error out of the log still log no error object.

* refactor(native-chat): hand the chat host's collaborators the logger, and give orcad its trace file

The delivery loop, idle sweep, queued-message drain, lease renewer, event
sink, conversation map and provider start/exit settlement each took an
internal error callback that the host mapped onto the logger. They now take
the logger itself and log under their own scope. The event sink keeps one
onFailed hook, which decides whether to stop the provider, not whether to
report. The dead-generation settlement returns its failure so each caller
logs it under its own scope.

orcad now installs the desktop's local trace sink under its own data root, so
a headless host's chat failures reach <data-root>/logs/main.trace.ndjson as
well as stderr.

Also passes the logger in the test fixtures the first commit missed, which
tc:node caught.

* fix(native-chat): keep repeated chat failures from flooding the trace file, and record their causes

- The production structured-chat logger writes a repeated failure (same level, scope, session,
  message and error text) once per 5 minutes, carrying how many repeats it swallowed; the
  tracked set is capped at 256.
- Trace entries now carry the error's code (and SQLite errcode) and up to three causes by name and
  message.
- A chat read whose conversation will not open is logged through the host's logger
  (open-for-read), and so are the runtime's chat-tab bookkeeping failures that already hold the
  host.
- orcad writes its own orcad.trace.ndjson, closes it after every quit handler, and flushes it on
  process exit; a trace file that cannot be opened leaves tracing off instead of stopping the app
  or orcad.
- Tests: the desktop wiring test proves the logger reaches the trace sink, and the privacy tests
  read every level the logger received.

* fix(native-chat): log a created chat's tab-publication and launch-prompt failures through the host's logger

* fix(native-chat): key a repeated chat failure on everything its entry writes

The repeat suppression keyed on the message and the error's text, so two refusals with the same
code but different causes, a plain error and a refusal of one code, or two object-valued errors
shared a key and the second was swallowed for five minutes. The key is now the entry's whole
written content (fields, code, errcode, refusal reason, cause chain, a stable rendering of a
non-error value) plus the error's name and message; a refusal's reason is also written.

* test(native-chat): pin that an error's name keeps two repeated failures apart

* test(native-chat): build the refusal in the repeat-key test as the wire does

* fix(native-chat): read an error's code and a refusal's reason by narrowing, not Reflect.get
This commit is contained in:
Brennan Benson
2026-10-01 11:21:50 -07:00
committed by GitHub
parent 3fbdaba262
commit c6cfcc034e
221 changed files with 2435 additions and 546 deletions
+6 -1
View File
@@ -165,7 +165,12 @@ An external supervisor (systemd, launchd, a process manager). orcad conforms to
to **stdout**; the supervisor owns capture and rotation. The daemon, being detached, writes
its own NDJSON lifecycle log to `<data-root>/logs/daemon.log` (suppressed by
`ORCA_DIAGNOSTICS_DISABLED=1`). Rotation of that file is not implemented — see
[What is not covered](#what-is-not-covered).
[What is not covered](#what-is-not-covered). orcad records every trace span it emits
(git commands, worktree paths, terminal spawns, structured-chat failures and the rest) to
`<data-root>/logs/orcad.trace.ndjson`, rotated at 10 MB × 10 files, private to its user and
redacted for secret-shaped strings. It stays on the host: a desktop's diagnostics bundle does
not collect it. `ORCA_DIAGNOSTICS_DISABLED=1` turns it off, and a logs folder orcad cannot
open leaves it off with one stderr warning rather than stopping orcad.
### orcad supervising the daemon
@@ -63,6 +63,7 @@ import {
noUpstreamError,
workingEvent
} from './first-work-branch-rename-test-harness'
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
function makeDeps(overrides: Partial<FirstWorkBranchRenameDeps> = {}) {
return makeBranchRenameDeps(vi.fn, overrides)
@@ -117,6 +118,7 @@ describe('maybeAutoRenameBranchOnFirstWork', () => {
}
})
const feed = new StructuredAgentSessionStatusFeed({
logger: createStructuredAgentSessionLogger(),
sessions: new Map([
[
'session',
@@ -215,6 +217,7 @@ describe('maybeAutoRenameBranchOnFirstWork', () => {
}
const pending: Promise<void>[] = []
const feed = new StructuredAgentSessionStatusFeed({
logger: createStructuredAgentSessionLogger(),
sessions: new Map([['session', { journal, params: { location, provider: 'codex' } }]]),
getRecord: () => null,
now: () => 1,
@@ -18,6 +18,7 @@ import {
} from '../native-chat/agent-session-wire/structured-agent-session-host-test-data'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
const SWEEP_MS = 5
const RETRY_GAP_MS = 10 * 60_000
@@ -77,6 +78,7 @@ beforeEach(async () => {
setOption: async () => undefined
}
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter,
journalDatabase: openTestJournalHostDatabase(root),
+2 -1
View File
@@ -12,6 +12,7 @@ import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/ag
import { unhandledProviderFrameJournalItem } from '../native-chat/agent-session-wire/unhandled-provider-frame'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const IDENTITY: AgentSessionJournalIdentity = {
sessionId: 'session-1',
@@ -54,7 +55,7 @@ async function statusRowsFor(frames: Record<string, unknown>[]) {
now: () => 1_700_000_000_000,
mintEpoch: () => 'epoch-1'
})
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind({ journal, fence: 1, publish: vi.fn() })
const translator = createClaudeJournalTranslator({ sink: deferred.sink, fallbackIdPrefix: '1' })
for (const frame of frames) {
@@ -15,6 +15,7 @@ import { backgroundWakeCapture } from './claude-captured-fold-steer-frames.test-
import { MOVED_TO_BACKGROUND } from './claude-captured-task-frames.test-fixture'
import { hostWithParent, parent, producer } from './claude-child-work-producer-harness.test-fixture'
import { PROVIDER_SESSION_ID } from './claude-structured-session-test-support'
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
vi.mock('../telemetry/client', () => ({ track: vi.fn() }))
vi.mock('../telemetry/cohort-classifier', () => ({ getCohortAtEmit: vi.fn(() => ({})) }))
@@ -40,6 +41,7 @@ async function wiredSession() {
})
}
const feed = new StructuredAgentSessionStatusFeed({
logger: createStructuredAgentSessionLogger(),
sessions: new Map([
[
parent.sessionId,
@@ -10,6 +10,7 @@ import { settleStaleStructuredAgentSessionState } from '../native-chat/agent-ses
import { readAgentJournalTurn } from '../../shared/agent-session-turn-record'
import { bindClaudeContextUsageCapture } from './claude-context-usage'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const journals = createTrackedJournalOpener()
let root: string
@@ -112,7 +113,7 @@ async function openJournal(): Promise<AgentSessionJournal> {
/** One acquisition: a fresh sink and translator over the session's journal. */
function acquire(journal: AgentSessionJournal) {
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const translator = createClaudeJournalTranslator({ sink: deferred.sink, coalesceMs: 0 })
deferred.bind({ journal, fence: 1, publish: () => {} })
const answers: ((value: unknown) => void)[] = []
@@ -27,6 +27,7 @@ import {
userFrame
} from './claude-context-usage-test-support'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const SESSION = 'orca-session'
const journals = createTrackedJournalOpener()
@@ -56,7 +57,7 @@ async function openJournal(): Promise<AgentSessionJournal> {
}
function translate(journal: AgentSessionJournal) {
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const translator = createClaudeJournalTranslator({ sink: deferred.sink, coalesceMs: 0 })
deferred.bind({ journal, fence: 1, publish: () => {} })
const settle = async (): Promise<void> => {
@@ -21,6 +21,7 @@ import { readClaudeStructuredSessionOptions } from './claude-structured-session-
import type { ClaudeSession } from './claude-structured-session-state'
import { CLAUDE_STRUCTURED_BASE_OPTIONS } from './claude-structured-launch-resolution'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
// These drive the real SDK against the scripted fake CLI, so every assertion is
// about the environment, argv and frames a real child actually saw.
@@ -385,7 +386,7 @@ describe('Claude stream-json connection', () => {
now: () => 1_700_000_000_000,
mintEpoch: () => 'epoch-1'
})
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind({ journal, fence: 1, publish: vi.fn() })
const translator = createClaudeJournalTranslator({ sink: deferred.sink })
let settled = false
@@ -15,6 +15,7 @@ import {
import type { AgentSessionJournal } from '../native-chat/agent-session-journal/journal-store'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import { blockOf } from './claude-background-task-row-test-support'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
// The frames below are the ones the reported session actually carried: two real
// failures that printed THREE red rows whose visible text was the wire opcode,
@@ -179,6 +180,7 @@ describe('claude journal translation — background task rows', () => {
return appendItem(...args)
})
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedOperations: 1,
maxQueuedOperations: 4,
@@ -252,7 +254,7 @@ describe('claude journal translation — background task rows', () => {
it('coalesces an unbound overflow patch and aliased final notification', async () => {
const persisted = new Map<string, AgentJournalItemBody>()
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const translator = createClaudeJournalTranslator({
sink: deferred.sink,
fallbackIdPrefix: 'test'
@@ -292,7 +294,7 @@ describe('claude journal translation — background task rows', () => {
it('reconciles an aliased notification after a parentless overflow patch was persisted', async () => {
const persisted = new Map<string, AgentJournalItemBody>()
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(persistedTarget(persisted))
const first = createClaudeJournalTranslator({ sink: deferred.sink, fallbackIdPrefix: 'first' })
fillLiveTaskRows(first)
@@ -353,7 +355,7 @@ describe('claude journal translation — background task rows', () => {
it('resolves a queued restart identity after the sink rebinds', async () => {
const persisted = new Map<string, AgentJournalItemBody>()
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(persistedTarget(persisted))
const first = createClaudeJournalTranslator({ sink: deferred.sink, fallbackIdPrefix: 'first' })
@@ -371,7 +373,7 @@ describe('claude journal translation — background task rows', () => {
first.dispose()
await deferred.drained()
const restarted = createDeferredStructuredAgentSessionEventSink()
const restarted = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const second = createClaudeJournalTranslator({
sink: restarted.sink,
fallbackIdPrefix: 'second'
@@ -397,7 +399,7 @@ describe('claude journal translation — background task rows', () => {
it('keeps pending writes from distinct runs when a translator is recreated', async () => {
const persisted = new Map<string, AgentJournalItemBody>()
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const first = createClaudeJournalTranslator({ sink: deferred.sink, fallbackIdPrefix: 'first' })
spawnToolCall(first, 'toolu-first')
@@ -23,6 +23,7 @@ import type { ClaudePendingPrompt } from './claude-structured-prompt-replies'
import { readAgentJournalTurn } from '../../shared/agent-session-turn-record'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
function sinkState() {
const items: { identity: AgentJournalItemIdentity; body: AgentJournalItemBody }[] = []
@@ -258,7 +259,7 @@ describe('Claude structured journal translation', () => {
now: () => 1_700_000_000_000,
mintEpoch: () => 'epoch-1'
})
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind({ journal, fence: 1, publish: vi.fn() })
let scheduled: (() => void) | null = null
const translator = createClaudeJournalTranslator({
@@ -317,7 +318,7 @@ describe('Claude structured journal translation', () => {
now: () => 1_700_000_000_000,
mintEpoch: () => 'epoch-1'
})
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind({ journal, fence: 1, publish: vi.fn() })
const translator = createClaudeJournalTranslator({ sink: deferred.sink })
const approval = prompt({
@@ -11,6 +11,7 @@ import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/ag
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const IDENTITY: AgentSessionJournalIdentity = {
sessionId: 'session-1',
@@ -80,7 +81,7 @@ describe('Claude provider fallback', () => {
now: () => 1_700_000_000_000,
mintEpoch: () => 'epoch-1'
})
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind({
journal,
fence: 1,
@@ -20,6 +20,7 @@ import {
identityFor,
PROVIDER_SESSION_ID
} from './claude-structured-session-test-support'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
function controlledSink(): {
sink: StructuredAgentSessionEventSink
@@ -212,6 +213,7 @@ describe('Claude structured reading control', () => {
const persisted = new Map<string, AgentJournalItemBody>()
const target = persistedTarget(persisted)
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedOperations: 1,
maxQueuedOperations: 4,
@@ -16,6 +16,7 @@ import { createTrackedJournalOpener } from '../native-chat/agent-session-journal
import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { createClaudeJournalTranslator } from './claude-structured-journal-translation'
import { claudeSubagentGroupBody, claudeSubagentGroupIdentity } from './claude-subagent-group-row'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
// Frame orders are real sessions', scrubbed. A resumed agent's frames still name its ORIGINAL
// spawn call while the announcement names the message call that resumed it, and a new provider
@@ -137,7 +138,7 @@ async function openJournal(): Promise<AgentSessionJournal> {
/** One provider process: a fresh sink and translator over the session's journal. */
function acquire(journal: AgentSessionJournal) {
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const translator = createClaudeJournalTranslator({ sink: deferred.sink, coalesceMs: 0 })
deferred.bind({ journal, fence: 1, publish: () => {} })
const settle = async (): Promise<void> => {
@@ -190,7 +191,7 @@ const rosterIds = (journal: AgentSessionJournal): Set<string> =>
/** An older build re-rostered a resumed child in the later turn's row, so two rows list it. */
async function journalAnOlderBuildListedTwice(): Promise<AgentSessionJournal> {
const journal = await openJournal()
const older = createDeferredStructuredAgentSessionEventSink()
const older = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
older.bind({ journal, fence: 1, publish: () => {} })
const listed = (id: string, state: NativeChatSubagentEntry['state']) => ({
id,
@@ -14,6 +14,7 @@ import {
type StructuredAgentSessionEventSink
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { ClaudeSubagentRoster } from './claude-subagent-roster'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const TURN_1 = 'claude-session:turn-1'
@@ -462,7 +463,7 @@ describe('ClaudeSubagentRoster — through the real sink queue', () => {
},
appendTombstone: async () => ({ epoch: 'e', sequence: 0 })
} as unknown as AgentSessionJournal
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind({
journal,
fence: 1,
@@ -18,6 +18,7 @@ import {
} from '../native-chat/agent-session-wire/structured-agent-session-host-test-data'
import { createCodexJournalTranslator } from './codex-structured-journal-translation'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
const SWEEP_MS = 5
const RETRY_GAP_MS = 10 * 60_000
@@ -75,6 +76,7 @@ beforeEach(async () => {
setOption: async () => undefined
}
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter,
journalDatabase: openTestJournalHostDatabase(root),
@@ -12,6 +12,7 @@ import {
} from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { CodexJournalGoals } from './codex-structured-journal-goals'
import { MAX_CODEX_GOAL_THREADS } from './codex-structured-journal-limits'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const THREAD = '01a08cc2-f96e-76d0-bb74-88b9bc0b03fc'
@@ -34,7 +35,7 @@ function goalFrame(goal: Record<string, unknown> = {}): Record<string, unknown>
}
function goalJournal(
options: Parameters<typeof createDeferredStructuredAgentSessionEventSink>[0] = {}
options: Partial<Parameters<typeof createDeferredStructuredAgentSessionEventSink>[0]> = {}
) {
let rowSequence = 0
let publishes = 0
@@ -43,7 +44,10 @@ function goalJournal(
let visitedItems = 0
const rows = new Map<string, AgentJournalRenderItem>()
const writes: string[] = []
const deferred = createDeferredStructuredAgentSessionEventSink(options)
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
...options
})
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: a fake journal exposing only the members the goal translator and deferred sink call.
const journal = {
get epoch() {
@@ -14,6 +14,7 @@ import { createTrackedJournalOpener } from '../native-chat/agent-session-journal
import { readAgentSessionHistory } from '../native-chat/agent-session-wire/agent-session-history-page'
import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { CodexJournalGoals } from './codex-structured-journal-goals'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const THREAD = '01a08cc2-f96e-76d0-bb74-88b9bc0b03fc'
const IDENTITY: AgentSessionJournalIdentity = {
@@ -56,7 +57,7 @@ function goalFrame(goal: Record<string, unknown> = {}) {
/** The host's own sink bound to a real journal. Its publish is what a subscriber
* receives: the page after the cursor it had caught up to. */
function journalSink(journal: AgentSessionJournal) {
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const published: AgentSessionHistoryPage[] = []
let subscriberCursor: AgentJournalCursor | null = null
deferred.bind({
@@ -21,6 +21,7 @@ import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/ag
import { createCodexJournalTranslator } from './codex-structured-journal-translation'
import type { CodexThreadItem } from './codex-thread-item-identity'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const SESSION = 'session-codex-children'
const PARENT = 'thread-parent'
@@ -52,7 +53,7 @@ async function openJournal(root: string): Promise<AgentSessionJournal> {
async function session() {
const root = await mkdtemp(join(tmpdir(), 'orca-codex-children-'))
let journal = await openJournal(root)
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const activities: (AgentSessionTurnActivity | null)[] = []
deferred.bind({
journal,
@@ -25,6 +25,7 @@ import {
CODEX_USER_INPUT_METHOD
} from './codex-structured-prompt-replies'
import type { CodexStructuredSessionEvent } from './codex-structured-session-adapter'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const SESSION_ID = 'session-1'
const THREAD_ID = 'thread-abc'
@@ -129,6 +130,7 @@ function deferredTarget(
function hardWatermarkDeferred() {
return createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedBytes: 1,
maxQueuedBytes: 1,
@@ -26,6 +26,7 @@ import {
import { createCodexStructuredNotificationRetry } from './codex-structured-notification-retry'
import type { CodexStructuredSessionEvent } from './codex-structured-session-adapter'
import type { CodexSession } from './codex-structured-session-state'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const SESSION_ID = 'session-1'
const THREAD_ID = 'thread-abc'
@@ -223,7 +224,7 @@ describe('codex turn lifecycle rows', () => {
now: () => 9_000,
stateDirectory: join(root, SESSION_ID)
})
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const translator = createCodexJournalTranslator({
sink: deferred.sink,
sessionId: SESSION_ID,
@@ -18,6 +18,7 @@ import {
CODEX_USER_INPUT_METHOD
} from './codex-structured-prompt-replies'
import type { CodexStructuredSessionEvent } from './codex-structured-session-adapter'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const SESSION_ID = 'session-1'
const THREAD_ID = 'thread-abc'
@@ -127,6 +128,7 @@ function deferredTarget(
function hardWatermarkDeferred() {
return createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedBytes: 1,
maxQueuedBytes: 1,
@@ -713,6 +715,7 @@ describe('codex journal translation', () => {
const publishes: string[] = []
const readingControl = { pauseReading: vi.fn(), resumeReading: vi.fn() }
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedBytes: 1,
maxQueuedBytes: 1,
@@ -13,6 +13,7 @@ import {
import { createTrackedJournalOpener } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { createDeferredStructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink'
import { createCodexJournalTranslator } from './codex-structured-journal-translation'
import { testEventSinkLogging } from '../native-chat/agent-session-wire/structured-agent-session-logger-test-support'
const SESSION = 'session-codex-failed-turn'
const THREAD = 'thread-abc'
@@ -39,7 +40,7 @@ async function session() {
stateDirectory: root,
now: () => 1_000
})
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind({ journal, fence: 1, publish: () => {} })
cleanups.push(async () => {
deferred.close()
@@ -36,6 +36,7 @@ import {
structuredQuestionTranscript
} from '../../renderer/src/components/native-chat/structured-agent-question-projection'
import { openTestJournalHostDatabase } from '../native-chat/agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from '../native-chat/agent-session-wire/structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -93,6 +94,7 @@ beforeEach(async () => {
}
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter,
journalDatabase: openTestJournalHostDatabase(root),
@@ -1,7 +1,7 @@
import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
import { agentSessionRefusalError } from '../../../shared/agent-session-wire-refusals'
import { NO_LEGACY_JOURNAL_RECORDS, openJournalDatabase } from './journal-database'
import {
@@ -14,6 +14,7 @@ import {
import { AgentSessionJournalError } from './journal-write-guards'
import { journalDatabasePath } from './journal-host-database'
import { replayJournal } from './journal-open'
import { recordingStructuredAgentSessionLogger } from '../agent-session-wire/structured-agent-session-logger-test-support'
let root: string
@@ -111,30 +112,39 @@ describe('classifyJournalOpenFailure', () => {
describe('journalOpenReadRefusal', () => {
it('names the reason, keeps the message the code and the storage text only as the cause', () => {
vi.spyOn(console, 'warn').mockImplementation(() => undefined)
const log = recordingStructuredAgentSessionLogger()
const storage = nodeSqliteError(26)
const refusal = journalOpenReadRefusal(storage)
const refusal = journalOpenReadRefusal(storage, log.logger, 'session-1')
expect(refusal.message).toBe('agent_session_journal_unreadable')
expect(refusal.refusal).toMatchObject({
code: 'agent_session_journal_unreadable',
details: { reason: 'journalCorrupt' }
})
expect(refusal.cause).toBe(storage)
vi.restoreAllMocks()
expect(log.entries).toEqual([
{
level: 'warn',
message: 'opening the conversation for a read failed',
fields: { scope: 'open-for-read', sessionId: 'session-1', error: storage }
}
])
})
it('passes a refusal the open already raised through unchanged', () => {
const raised = agentSessionRefusalError('agent_session_identity_required', {
reason: 'recordMissing'
})
expect(journalOpenReadRefusal(raised)).toBe(raised)
expect(
journalOpenReadRefusal(raised, recordingStructuredAgentSessionLogger().logger, 's')
).toBe(raised)
})
})
describe('createJournalOpenReadRefusals', () => {
it('logs a session once per failure until it opens, and each session on its own', () => {
const warn = vi.spyOn(console, 'warn').mockImplementation(() => undefined)
const refusals = createJournalOpenReadRefusals()
const log = recordingStructuredAgentSessionLogger()
const logged = () => log.entries.map((entry) => entry.fields.sessionId)
const refusals = createJournalOpenReadRefusals(log.logger)
const denied = systemError('EACCES', -13)
const corrupt = nodeSqliteError(26)
@@ -143,17 +153,21 @@ describe('createJournalOpenReadRefusals', () => {
expect(refusal.refusal).toMatchObject({ details: { reason: 'journalUnavailable' } })
expect(refusal.cause).toBe(denied)
}
expect(warn).toHaveBeenCalledTimes(1)
expect(logged()).toEqual(['session-1'])
refusals.refusal('session-2', denied)
expect(warn).toHaveBeenCalledTimes(2)
expect(logged()).toEqual(['session-1', 'session-2'])
expect(refusals.refusal('session-1', corrupt).refusal).toMatchObject({
details: { reason: 'journalCorrupt' }
})
expect(warn).toHaveBeenCalledTimes(3)
expect(logged()).toEqual(['session-1', 'session-2', 'session-1'])
refusals.forget('session-1')
refusals.refusal('session-1', corrupt)
expect(warn).toHaveBeenCalledTimes(4)
vi.restoreAllMocks()
expect(log.scopes()).toEqual([
'open-for-read',
'open-for-read',
'open-for-read',
'open-for-read'
])
})
})
@@ -176,13 +190,12 @@ describe('a journal a newer Orca wrote', () => {
})
it('refuses a read with the same reason', () => {
vi.spyOn(console, 'warn').mockImplementation(() => undefined)
expect(journalOpenReadRefusal(readOnly()).refusal).toMatchObject({
const { logger } = recordingStructuredAgentSessionLogger()
expect(journalOpenReadRefusal(readOnly(), logger, 'session-1').refusal).toMatchObject({
details: { reason: 'journalWrittenByNewerOrca' }
})
expect(createJournalOpenReadRefusals().refusal('session-1', readOnly()).refusal).toMatchObject({
details: { reason: 'journalWrittenByNewerOrca' }
})
vi.restoreAllMocks()
expect(
createJournalOpenReadRefusals(logger).refusal('session-1', readOnly()).refusal
).toMatchObject({ details: { reason: 'journalWrittenByNewerOrca' } })
})
})
@@ -13,6 +13,7 @@ import {
} from '../../../shared/agent-session-wire-refusals'
import { isSqliteCorruption } from '../../sqlite/sqlite-read-failure'
import { AgentSessionJournalError } from './journal-write-guards'
import type { StructuredAgentSessionLogger } from '../agent-session-wire/structured-agent-session-logger'
type JournalRefusalReason = AgentSessionRefusalReason<'agent_session_journal_unreadable'>
@@ -107,11 +108,15 @@ export function journalOpenRefusalError(error: unknown): AgentSessionRefusalErro
* (a path, "file is not a database") goes to the log only; the reader gets the classified refusal,
* whose message stays the bare code.
*/
export function journalOpenReadRefusal(error: unknown): AgentSessionRefusalError {
export function journalOpenReadRefusal(
error: unknown,
logger: StructuredAgentSessionLogger,
sessionId: string
): AgentSessionRefusalError {
if (isAgentSessionRefusalError(error)) {
return error
}
return unreadableRefusal(error, journalRefusalReason(error), true)
return unreadableRefusal(error, journalRefusalReason(error), { logger, sessionId })
}
const MAX_LOGGED_SESSIONS = 256
@@ -120,7 +125,7 @@ const MAX_LOGGED_SESSIONS = 256
* The read door's refusals for one host. A reader reconnects on a timer while an open can clear,
* so a session's failure is logged once until that session opens or the failure changes.
*/
export function createJournalOpenReadRefusals() {
export function createJournalOpenReadRefusals(logger: StructuredAgentSessionLogger) {
const logged = new Map<string, string>()
return {
refusal: (sessionId: string, error: unknown): AgentSessionRefusalError => {
@@ -134,7 +139,7 @@ export function createJournalOpenReadRefusals() {
if (!repeat && (logged.has(sessionId) || logged.size < MAX_LOGGED_SESSIONS)) {
logged.set(sessionId, failure)
}
return unreadableRefusal(error, reason, !repeat)
return unreadableRefusal(error, reason, repeat ? null : { logger, sessionId })
},
/** The session opened or closed: its next failure is news. */
forget: (sessionId: string): void => {
@@ -146,11 +151,13 @@ export function createJournalOpenReadRefusals() {
function unreadableRefusal(
error: unknown,
reason: JournalRefusalReason,
log: boolean
log: { logger: StructuredAgentSessionLogger; sessionId: string } | null
): AgentSessionRefusalError {
if (log) {
console.warn('[agent-session] opening the conversation for a read failed:', error)
}
log?.logger.warn('opening the conversation for a read failed', {
scope: 'open-for-read',
sessionId: log.sessionId,
error
})
const code = 'agent_session_journal_unreadable'
return new AgentSessionRefusalError(refuse(code, { reason }, code), { cause: error })
}
@@ -28,6 +28,7 @@ import {
import { journalDirectoryFor, legacyJournalDatabaseFile } from './journal-paths'
import { importPerSessionJournal } from './journal-per-session-import'
import { readJournalSessionEpoch, type JournalStoredRow } from './journal-row-table'
import { createStructuredAgentSessionLogger } from '../agent-session-wire/structured-agent-session-logger'
vi.mock('node:fs', async (importOriginal) => {
const actual = await importOriginal<typeof NodeFs>()
@@ -221,6 +222,7 @@ describe('importing a per-chat journal', () => {
const journal = await openChat()
const withdrawal = createStructuredAgentSessionRestartOfferWithdrawal({
logger: createStructuredAgentSessionLogger(),
sessions: new Map([[IDENTITY.sessionId, { journal, child: null }]]),
now: () => clock,
enqueue: (operation) => operation()
@@ -45,6 +45,7 @@ import {
import { agentSessionFailureFact } from '../../../shared/agent-session-failure'
import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -76,6 +77,7 @@ const spawnChild: StructuredAgentSessionAdapter['acquire'] = async ({ fence, spa
async function startHost(): Promise<void> {
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire,
@@ -20,6 +20,7 @@ import { openTestAttachConversation } from './structured-agent-session-attach-te
import { performAttach } from './structured-agent-session-attach-flow'
import type { AgentSessionCreatePhaseRecorder } from '../../observability/agent-session-instrumentation'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const NOW = 1_800_000_000_000
const SESSION = 'legacy-session'
@@ -151,6 +152,7 @@ describe('structured session acquisition options', () => {
let firstJournal: AgentSessionJournal | undefined
const first = await performAttach({
logger: createStructuredAgentSessionLogger(),
store: initialStore,
adapter: withHistory('created'),
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -188,6 +190,7 @@ describe('structured session acquisition options', () => {
const releasedFence = store.getRecord(SESSION)?.lease.runtimeFence ?? 0
const second = await performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: withHistory('resumed'),
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -219,6 +222,7 @@ describe('structured session acquisition options', () => {
const recordPhase = vi.fn<AgentSessionCreatePhaseRecorder>()
const created = await performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: sessionAdapter,
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -248,6 +252,7 @@ describe('structured session acquisition options', () => {
const sessionAdapter = adapter({ origin: 'created' })
const attempt = async (options: Readonly<Record<string, string>>, spawnToken: string) =>
performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: sessionAdapter,
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -278,6 +283,7 @@ describe('structured session acquisition options', () => {
const store = await openTestAgentSessionRecordStore(root)
const created = await performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: adapter({ origin: 'created' }),
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -308,6 +314,7 @@ describe('structured session acquisition options', () => {
})
const releasedFence = resumedStore.getRecord(SESSION)?.lease.runtimeFence ?? 0
const resumed = await performAttach({
logger: createStructuredAgentSessionLogger(),
store: resumedStore,
adapter: adapter({
origin: 'resumed',
@@ -350,6 +357,7 @@ describe('structured session acquisition options', () => {
})
const created = await performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: sessionAdapter,
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -386,6 +394,7 @@ describe('structured session acquisition options', () => {
await expect(
performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: failingAdapter,
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -479,6 +488,7 @@ describe('structured session acquisition options', () => {
fence: number | null
) =>
performAttach({
logger: createStructuredAgentSessionLogger(),
store: target,
adapter: failingAdapter,
openConversation: openTestAttachConversation(
@@ -580,6 +590,7 @@ describe('the tab a create reserves', () => {
function attachWith(store: AgentSessionRecordStore, surfaceTabId?: string) {
return performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: adapter({ origin: 'created' }),
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -19,6 +19,7 @@ import { AgentSessionJournal } from '../agent-session-journal/journal-store'
import * as legacyImport from '../agent-session-journal/journal-legacy-import'
import { StructuredAgentSessionHost } from './structured-agent-session-host'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const NOW = 1_800_000_000_000
const SESSION = 'codex_adopting_session'
@@ -122,6 +123,7 @@ async function attach(
) {
store ??= await openTestAgentSessionRecordStore(root!)
return performAttach({
logger: createStructuredAgentSessionLogger(),
store,
adapter: sessionAdapter,
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root!)),
@@ -233,6 +235,7 @@ describe('adopting a provider conversation on create', () => {
await writeCodexRollout(transcriptPath, 'valid source')
store = await openTestAgentSessionRecordStore(root)
const host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: adapter(),
journalDatabase: openTestJournalHostDatabase(root),
@@ -79,7 +79,11 @@ export function ensureStructuredAgentSessionAgentForOperation(
): Promise<StructuredAgentSessionResumeOutcome> {
return ensureStructuredAgentSessionAgent(context, sessionId).catch((error: unknown) => {
// The error is Orca's own and goes to the log; the refusal says only that the start failed.
console.warn('[agent-session] starting the agent for an operation failed:', error)
context.deps.logger.warn('starting the agent for an operation failed', {
scope: 'operation-agent-start',
sessionId,
error
})
return {
ok: false,
refusal: refuseUnclassified(
@@ -26,6 +26,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
const EXIT_REASON = 'Claude Code is not signed in. Sign in with the Claude CLI'
@@ -128,6 +129,7 @@ beforeEach(async () => {
}))
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire,
@@ -18,6 +18,7 @@ import type { AgentSessionRecordStore } from '../../runtime/agent-session-record
import { openTestAgentSessionRecordStore } from '../../runtime/agent-session-record-store-test-harness'
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
import { StructuredAgentSessionHost } from './structured-agent-session-host'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
import {
HOST_TEST_LOCATION,
HOST_TEST_NOW,
@@ -87,6 +88,7 @@ async function openHost(catalog = catalogFor(workspace)): Promise<void> {
adapter: adapter(catalog),
journalDatabase: openTestJournalHostDatabase(directory),
claimKeyId: 'key',
logger: recordingStructuredAgentSessionLogger().logger,
now: () => clock,
probeOwner: async () => ({ outcome: 'pid-absent' })
})
@@ -44,10 +44,12 @@ import {
} from '../../observability/agent-session-instrumentation'
import type { ProviderHistoryWindow } from '../agent-session-journal/journal-submission-reconciler'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
export type AttachFlowInput = {
store: AgentSessionRecordStore
adapter: StructuredAgentSessionAdapter
logger: StructuredAgentSessionLogger
authority: AgentSessionAttachAuthority
callerKey: string
params: AgentSessionAttachParams
@@ -206,7 +208,11 @@ export async function performAttach(
const thrown = failed ? error : preSpawnFailureInWords(error, wording)
if (failed || thrown !== error) {
// The answer carries only its sentence, so what failed is kept here.
console.warn('[agent-session] provider start failed:', error)
input.logger.warn('starting the provider failed', {
scope: 'provider-start',
sessionId,
error
})
}
return (
failed ?? {
@@ -135,6 +135,7 @@ async function runAttach(
const attached = await performAttach({
store: context.deps.store,
adapter: context.deps.adapter,
logger: context.deps.logger,
eventSink: attemptSink.sink,
// The superseded child's writes settle into its own journal before a new child starts.
onAcquiring: async () => {
@@ -201,7 +202,7 @@ async function runAttach(
}
}
await recoverStructuredRewind(
context.deps.store,
context.deps,
sessionId,
attached.journal,
fence,
@@ -3,6 +3,7 @@ import type { JournalHostDatabase } from '../agent-session-journal/journal-host-
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type { StructuredAgentSessionAdapter } from './structured-agent-session-adapter'
import { openStructuredAgentSessionConversationJournal } from './structured-agent-session-conversation-open'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
/** For a test that attaches without a host: the conversation's journal, through the one open a
* host would take, so the attach under test adopts it the way it adopts the host's. */
@@ -12,8 +13,12 @@ export function openTestAttachConversation(
): (record: AgentSessionRecord) => Promise<AgentSessionJournal> {
return async (record) =>
(
await openStructuredAgentSessionConversationJournal({ journalDatabase, adapter }, record, {
acquisition: true
})
await openStructuredAgentSessionConversationJournal(
{ journalDatabase, adapter, logger: recordingStructuredAgentSessionLogger().logger },
record,
{
acquisition: true
}
)
).session.journal
}
@@ -10,6 +10,7 @@ import type {
StructuredAgentSessionHostSession
} from './structured-agent-session-host-types'
import type { AgentSessionSubscribers } from './structured-agent-session-subscribers'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
const child: AgentChildWorkView = {
id: 'child-1',
@@ -31,7 +32,7 @@ function channelOver(
) {
const sessions = new StructuredAgentSessionConversations({
deliver: () => {},
onDeliveryError: () => {},
logger: recordingStructuredAgentSessionLogger().logger,
now: () => 1
})
const sent = vi.fn()
@@ -108,7 +108,12 @@ export function mutateWithChatStop<TValue>(
{ sessionId, adapter: context.deps.adapter },
windDown,
() => context.stopAgent(sessionId),
(error) => context.deps.onEventSinkError?.({ sessionId, error })
(error) =>
context.deps.logger.warn("ending a stopped chat's provider session failed", {
scope: 'chat-stop',
sessionId,
error
})
)
}
})
@@ -28,6 +28,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -70,6 +71,7 @@ beforeEach(async () => {
})
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
// The production router is what declares create support; the bare adapter only knows locations.
adapter: Object.assign(adapter, { supportsCreate: () => true }),
@@ -27,6 +27,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -66,6 +67,7 @@ beforeEach(async () => {
})
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: Object.assign(adapter, { supportsCreate: () => true }),
journalDatabase: openTestJournalHostDatabase(root),
@@ -16,6 +16,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-claude' }
const CLAUDE_SESSION = '019fd532-7c11-7a90-b6de-4e1a2c3d5f61'
@@ -83,6 +84,7 @@ beforeEach(async () => {
optionWritable = Promise.resolve()
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: adapter(),
journalDatabase: openTestJournalHostDatabase(root),
@@ -29,6 +29,7 @@ import {
hostTestOperationId,
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
// As Claude Code 2.1.280 advertises them on a turn's system/init frame.
@@ -82,6 +83,7 @@ beforeEach(async () => {
})
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: Object.assign(adapter, { supportsCreate: () => true }),
journalDatabase: openTestJournalHostDatabase(root),
@@ -17,6 +17,7 @@ import type { AgentSessionAttachParams } from './structured-agent-session-attach
import { stopStructuredAgentSessionAgentUnderSerialize } from './structured-agent-session-host-lifetime'
import { StructuredAgentSessionHostRuntimeState } from './structured-agent-session-host-runtime-state'
import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const NOW = 1_788_727_031_330
const roots: string[] = []
@@ -121,7 +122,8 @@ describe('Claude root-exit stop', () => {
adapter,
// The one database the store and the journal share, as the runtime installs them.
journalDatabase: openTestJournalHostDatabase(stateDirectory),
claimKeyId: 'key-1'
claimKeyId: 'key-1',
logger: createStructuredAgentSessionLogger()
}
const runtimeState = new StructuredAgentSessionHostRuntimeState(deps)
@@ -32,6 +32,7 @@ import { openTestAgentSessionRecordStore } from '../../runtime/agent-session-rec
import { structuredClaudeLifecycleEvent } from '../../runtime/structured-claude-runtime-adapter'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { StructuredAgentSessionHost } from './structured-agent-session-host'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
import {
HOST_TEST_NOW as NOW,
HOST_TEST_SESSION as SESSION,
@@ -52,7 +53,7 @@ let store: AgentSessionRecordStore
let queued: string[]
let claude: ReturnType<typeof fakeClaude>
let events: ClaudeStructuredSessionEvent[]
let sinkErrors: unknown[]
let logs: ReturnType<typeof recordingStructuredAgentSessionLogger>
// The host's status row and child records, as the app's hook server holds them.
let server: AgentHookServer
let statusSubject: AgentStatusStructuredSessionSubject | undefined
@@ -62,7 +63,7 @@ beforeEach(async () => {
resetHostTestOperationIds()
queued = []
events = []
sinkErrors = []
logs = recordingStructuredAgentSessionLogger()
server = new AgentHookServer()
statusSubject = undefined
claude = fakeClaude({
@@ -107,7 +108,7 @@ beforeEach(async () => {
journalDatabase: openTestJournalHostDatabase(root),
claimKeyId: 'key-1',
mintSpawnToken: () => 'spawn-a',
onEventSinkError: ({ error }) => sinkErrors.push(error),
logger: logs.logger,
statusSink: {
publish: (summary, subject) => {
statusSubject = subject
@@ -601,10 +602,15 @@ it('still answers the Stop, with its row, when the child cannot be proven gone;
await laneDrained()
expect(await statusTexts()).toEqual(['Cancellation requested.'])
expect(sinkErrors).toEqual([
expect(logs.entries).toEqual([
expect.objectContaining({
name: 'StructuredAgentSessionEvictionError',
step: 'stop-provider-child'
fields: expect.objectContaining({
scope: 'chat-stop',
error: expect.objectContaining({
name: 'StructuredAgentSessionEvictionError',
step: 'stop-provider-child'
})
})
})
])
connection.close = close
@@ -660,7 +666,7 @@ it('keeps a second Stop pressed while the first ends the child quiet', async ()
expect(connection.closed).toBe(true)
expect(connection.calls.filter((call) => call.subtype === 'interrupt')).toHaveLength(1)
expect(await statusTexts()).toEqual(['Cancellation requested.'])
expect(sinkErrors).toEqual([])
expect(logs.entries).toEqual([])
})
const BRANCH_QUESTION = {
@@ -31,6 +31,7 @@ import {
childEndCauseOfEndedEvent,
turnVerdictForChildEnd
} from './structured-agent-session-stale-turn-verdict'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CUT_TURN = { provider: 'codex' as const, threadId: THREAD, turnId: 'cut-turn', ordinal: 1 }
@@ -51,6 +52,7 @@ beforeEach(() => {
exitObservedFirst = false
closeCalls = 0
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store: state.store,
adapter: {
...adapter(),
@@ -24,6 +24,7 @@ import {
hostTestOperationId,
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -51,6 +52,7 @@ beforeEach(async () => {
}
const store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: Object.assign(adapterFor(codex), { supportsCreate: () => true }),
journalDatabase: openTestJournalHostDatabase(root),
@@ -317,8 +317,14 @@ describe('the wind-down retry with a message queued (P2-31)', () => {
stopAgent,
stopStartingAgent: stopAgent,
closeConversation: vi.fn(async () => false),
onError: (_id, error) => {
throw error
// A failed step fails the test.
logger: {
warn: (_message, fields) => {
throw fields.error
},
error: (_message, fields) => {
throw fields.error
}
}
})
await sweep.tick()
@@ -24,6 +24,7 @@ import type { StructuredAgentSessionHostSession } from './structured-agent-sessi
import { StructuredAgentSessionIdleSweep } from './structured-agent-session-idle-sweep'
import { AGENT_SESSION_NOT_ATTACHED } from './structured-agent-session-mutation-admission'
import { adapterSupportsRecord } from './structured-agent-session-provider-support'
import { deferredStructuredAgentSessionLogger } from './structured-agent-session-logger'
export type StructuredAgentSessionConversationLifetime = ReturnType<
typeof createStructuredAgentSessionConversationLifetime
@@ -45,7 +46,9 @@ export function createStructuredAgentSessionConversationLifetime(host: {
let disposed = false
const { sessions, serialize } = host
const deps = () => host.context().deps
const readRefusals = createJournalOpenReadRefusals()
const readRefusals = createJournalOpenReadRefusals(
deferredStructuredAgentSessionLogger(() => deps().logger)
)
// The owed copy fails as an open does: the reader gets the classified refusal, never its text.
const whenImported = (sessionId: string, session: StructuredAgentSessionHostSession) =>
session.journal.whenImported().catch((error: unknown) => {
@@ -88,7 +91,7 @@ export function createStructuredAgentSessionConversationLifetime(host: {
cause: 'host-stop'
}),
closeConversation,
onError: (sessionId, error) => deps().onEventSinkError?.({ sessionId, error }),
logger: deps().logger,
...deps().idleSweep
})
@@ -38,7 +38,7 @@ export type StructuredAgentSessionConversationOpenDeps = {
store: Pick<AgentSessionRecordStore, 'getRecord'>
adapter: Pick<StructuredAgentSessionAdapter, 'historyFilePath'>
journalDatabase: JournalHostDatabase
onEventSinkError?: StructuredAgentSessionHostDeps['onEventSinkError']
logger: StructuredAgentSessionHostDeps['logger']
}
/** An acquisition's own open: its reserve cleared the record's death evidence, so it settles
@@ -104,7 +104,11 @@ export async function openStructuredAgentSessionConversationJournal(
// one is only doubt, which provider history decides under a won lease.
await opened.journal.markPendingSubmissionsUnknown(fence)
} catch (error) {
deps.onEventSinkError?.({ sessionId, error })
deps.logger.warn('marking pending sends unknown on open failed', {
scope: 'open-pending-unknown',
sessionId,
error
})
}
// No child in this process writes to a journal nobody had open, so whatever it shows running
// belongs to a generation that is gone, whatever the lease still claims. Settled before any
@@ -135,7 +139,7 @@ export async function resettleOpenStructuredAgentSessionConversation(
}
async function settleGoneGeneration(
deps: Pick<StructuredAgentSessionConversationOpenDeps, 'onEventSinkError'>,
deps: Pick<StructuredAgentSessionConversationOpenDeps, 'logger'>,
record: AgentSessionRecord,
journal: AgentSessionJournal
): Promise<void> {
@@ -150,7 +154,11 @@ async function settleGoneGeneration(
})
} catch (error) {
// Best effort: the next open or acquire re-derives it.
deps.onEventSinkError?.({ sessionId: record.sessionId, error })
deps.logger.warn("settling a gone agent's work on open failed", {
scope: 'open-dead-generation',
sessionId: record.sessionId,
error
})
}
}
@@ -26,6 +26,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -55,6 +56,7 @@ beforeEach(async () => {
stopEndsSession = false
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire: async ({ fence, spawnToken, events: sink }) => {
@@ -9,6 +9,7 @@ import { createTrackedJournalOpener } from '../agent-session-journal/journal-hos
import { StructuredAgentSessionConversations } from './structured-agent-session-conversations'
import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types'
import { hostTestAttachParams } from './structured-agent-session-host-test-data'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
const IDENTITY: AgentSessionJournalIdentity = {
sessionId: 'session-1',
@@ -55,7 +56,7 @@ describe('a conversation delivers what its journal commits', () => {
const deliver = vi.fn()
const conversations = new StructuredAgentSessionConversations({
deliver,
onDeliveryError: vi.fn(),
logger: recordingStructuredAgentSessionLogger().logger,
now: () => 0
})
const journal = await openJournal('a')
@@ -74,7 +75,7 @@ describe('a conversation delivers what its journal commits', () => {
const sessions: Map<string, StructuredAgentSessionHostSession> =
new StructuredAgentSessionConversations({
deliver,
onDeliveryError: vi.fn(),
logger: recordingStructuredAgentSessionLogger().logger,
now: () => 0
})
const journal = await openJournal('a')
@@ -89,7 +90,7 @@ describe('a conversation delivers what its journal commits', () => {
const deliver = vi.fn()
const conversations = new StructuredAgentSessionConversations({
deliver,
onDeliveryError: vi.fn(),
logger: recordingStructuredAgentSessionLogger().logger,
now: () => 0
})
const journal = await openJournal('a')
@@ -104,7 +105,7 @@ describe('a conversation delivers what its journal commits', () => {
const deliver = vi.fn()
const conversations = new StructuredAgentSessionConversations({
deliver,
onDeliveryError: vi.fn(),
logger: recordingStructuredAgentSessionLogger().logger,
now: () => 0
})
const replaced = await openJournal('a')
@@ -123,7 +124,7 @@ describe('a conversation delivers what its journal commits', () => {
const deliver = vi.fn()
const conversations = new StructuredAgentSessionConversations({
deliver,
onDeliveryError: vi.fn(),
logger: recordingStructuredAgentSessionLogger().logger,
now: () => 0
})
const journal = await openJournal('a')
@@ -137,12 +138,12 @@ describe('a conversation delivers what its journal commits', () => {
it('reports a reader failure without failing the durable write', async () => {
const failure = new Error('reader failed')
const onDeliveryError = vi.fn()
const log = recordingStructuredAgentSessionLogger()
const conversations = new StructuredAgentSessionConversations({
deliver: () => {
throw failure
},
onDeliveryError,
logger: log.logger,
now: () => 0
})
const journal = await openJournal('a')
@@ -152,7 +153,9 @@ describe('a conversation delivers what its journal commits', () => {
itemId: expect.any(String)
})
expect(onDeliveryError).toHaveBeenCalledExactlyOnceWith('session-1', failure)
expect(log.entries.map((entry) => entry.fields)).toEqual([
{ scope: 'journal-delivery', sessionId: 'session-1', error: failure }
])
expect(journal.snapshot().items.map((item) => item.body)).toContainEqual({
kind: 'status',
text: 'durable'
@@ -1,5 +1,6 @@
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
/**
* The host's open conversations. A journal handle becomes a conversation's when it is set here,
@@ -22,7 +23,7 @@ export class StructuredAgentSessionConversations extends Map<
constructor(
private readonly delivery: {
deliver: (sessionId: string, journal: AgentSessionJournal) => void
onDeliveryError: (sessionId: string, error: unknown) => void
logger: StructuredAgentSessionLogger
/** A conversation became held: state that waited on it (queued drafts) re-derives. */
onOpened?: (sessionId: string) => void
now: () => number
@@ -47,7 +48,11 @@ export class StructuredAgentSessionConversations extends Map<
try {
this.delivery.deliver(sessionId, journal)
} catch (error) {
this.delivery.onDeliveryError(sessionId, error)
this.delivery.logger.warn('delivering a journal commit failed', {
scope: 'journal-delivery',
sessionId,
error
})
}
})
})
@@ -25,6 +25,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
const CHILD_PID = 4321
@@ -92,6 +93,7 @@ function host(
overrides: Partial<StructuredAgentSessionHostDeps> = {}
): StructuredAgentSessionHost {
return new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter,
journalDatabase: openTestJournalHostDatabase(generationRoot(generation)),
@@ -45,6 +45,8 @@ import {
HOST_TEST_LOCATION as LOCATION,
HOST_TEST_SESSION as SESSION
} from './structured-agent-session-host-test-data'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
const PROVIDER_SESSION = 'provider-session-alpha-1'
/** The tool call's row: the last thing the provider wrote before the crash. */
@@ -155,6 +157,7 @@ async function seedClaudeToolTurn(): Promise<void> {
function openHost(overrides: Partial<StructuredAgentSessionHostDeps>): void {
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire: vi.fn(),
@@ -440,11 +443,11 @@ describe('a turn a read reached before the reconcile proved its owner dead', ()
it('stays unverifiable when the revision cannot be written, and a later open revises it', async () => {
let now = RELAUNCHED_AT
const onEventSinkError = vi.fn()
const log = recordingStructuredAgentSessionLogger()
openHost({
probeOwner: async () => ({ outcome: 'pid-absent' }),
now: () => now,
onEventSinkError
logger: log.logger
})
await host.history({ sessionId: SESSION, direction: 'tail' })
const { journal } = host.collaboratorsForTests().sessions.get(SESSION)!
@@ -453,7 +456,7 @@ describe('a turn a read reached before the reconcile proved its owner dead', ()
await host.reconcileRestartLeases()
await drainSession()
expect(onEventSinkError).toHaveBeenCalledOnce()
expect(log.scopes()).toEqual(['open-dead-generation'])
expect(await settledTurn()).toEqual(UNVERIFIABLE_TURN)
// The proof is durable on the record, so the next open converges.
now += STRUCTURED_AGENT_SESSION_IDLE_MS + 1
@@ -100,7 +100,7 @@ describe('dead structured-session generation settlement', () => {
verdict: { state: 'unverifiable' },
showUnexpectedExitOutcome: false
})
).resolves.toBe(true)
).resolves.toEqual({ ok: true })
const snapshot = journal.snapshot()
expect(snapshot.submissions).toEqual([
@@ -135,9 +135,9 @@ describe('dead structured-session generation settlement', () => {
showUnexpectedExitOutcome: true
}
await expect(settleStructuredAgentSessionDeadGeneration(input)).resolves.toBe(true)
await expect(settleStructuredAgentSessionDeadGeneration(input)).resolves.toEqual({ ok: true })
const settledCursor = journal.cursor()
await expect(settleStructuredAgentSessionDeadGeneration(input)).resolves.toBe(true)
await expect(settleStructuredAgentSessionDeadGeneration(input)).resolves.toEqual({ ok: true })
expect(journal.cursor()).toEqual(settledCursor)
expect(
@@ -166,7 +166,7 @@ describe('dead structured-session generation settlement', () => {
detail: providerDiagnostic('stack frame '.repeat(4_000), 'log')
})
})
).resolves.toBe(true)
).resolves.toEqual({ ok: true })
const statuses = journal
.snapshot()
@@ -258,7 +258,7 @@ describe('dead structured-session generation settlement', () => {
verdict: { state: 'interrupted', completedAt: 1_000 },
showUnexpectedExitOutcome: false
})
).resolves.toBe(true)
).resolves.toEqual({ ok: true })
})
it('settles a live unknown submission even when no unfinished item remains', async () => {
@@ -285,7 +285,7 @@ describe('dead structured-session generation settlement', () => {
verdict: { state: 'interrupted', completedAt: 1_000 },
showUnexpectedExitOutcome: false
})
).resolves.toBe(true)
).resolves.toEqual({ ok: true })
expect(journal.submissions()).toEqual([
expect.objectContaining({
@@ -463,7 +463,7 @@ describe('whether a dead generation interrupted anything', () => {
1_000
)
})
).resolves.toBe(true)
).resolves.toEqual({ ok: true })
const snapshot = journal.snapshot()
expect(snapshot.items.some((item) => item.body.kind === 'status')).toBe(false)
@@ -102,6 +102,11 @@ export function unfinishedStructuredAgentSessionWorkWasInterrupted(
return outcomeItems.some((item) => !isCleanlySettled(currentItems.get(item.itemId)))
}
/** Whether the settlement was written, and what stopped it when it was not. */
export type StructuredAgentSessionDeadGenerationSettlement =
| { ok: true }
| { ok: false; error: unknown }
export async function settleStructuredAgentSessionDeadGeneration(input: {
journal: DeadGenerationJournal
sessionId: string
@@ -117,13 +122,12 @@ export async function settleStructuredAgentSessionDeadGeneration(input: {
/** The provider never finished starting: the start that failed, keyed by the child's
* generation. Its row is the one the delivery loop writes for the same start. */
exitedDuringStartup?: { generation: string | null }
onError?: (sessionId: string, error: unknown) => void
}): Promise<boolean> {
}): Promise<StructuredAgentSessionDeadGenerationSettlement> {
try {
const hasUnfinishedWork = hasUnfinishedStructuredAgentSessionWork(input.journal)
const showUnexpectedExitOutcome = input.showUnexpectedExitOutcome ?? hasUnfinishedWork
if (!showUnexpectedExitOutcome && !hasUnfinishedWork) {
return true
return { ok: true }
}
// A queued message is the delivery loop's to settle: it was never handed to this child. A
// child that never proved its start accepted nothing either — input is written only after it
@@ -189,10 +193,10 @@ export async function settleStructuredAgentSessionDeadGeneration(input: {
mutations: chunk.mutations
})
}
return true
return { ok: true }
} catch (error) {
input.onError?.(input.sessionId, error)
return false
// Returned rather than logged: each caller logs it under its own scope.
return { ok: false, error }
}
}
@@ -40,6 +40,7 @@ import {
import { failedProviderChildStart } from './structured-agent-session-provider-child'
import { handOverSubmission } from './structured-agent-session-turns'
import { structuredAgentSessionCommandRunning } from './structured-agent-session-command-turn'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
export type StructuredAgentSessionDeliveryLoopDeps = {
sessions: ReadonlyMap<string, StructuredAgentSessionHostSession>
@@ -62,7 +63,7 @@ export type StructuredAgentSessionDeliveryLoopDeps = {
) => Promise<boolean>
/** Who the chat's failure sentences name. */
failureTextContext: (sessionId: string) => AgentSessionFailureWordsContext
onError: (sessionId: string, error: unknown) => void
logger: StructuredAgentSessionLogger
record: (sessionId: string) => AgentSessionRecord | null
readChildWork: (sessionId: string) => readonly AgentChildWorkView[] | undefined
flushStreamedEvents: (sessionId: string) => Promise<void>
@@ -139,14 +140,22 @@ export class StructuredAgentSessionDeliveryLoop {
}
} catch (error) {
// The error is Orca's own and goes to the log; the chat says only that Orca failed.
this.deps.onError(sessionId, error)
this.deps.logger.warn('delivering a queued message failed', {
scope: 'delivery-loop',
sessionId,
error
})
const cause = { hostFault: true } as const
await this.deps
.serialize(sessionId, () => this.fail(sessionId, { startKey: null, cause }))
.catch((failure: unknown) => {
// Rows left queued are rejected by the next open, or by the next loop an accept wakes.
this.running.delete(sessionId)
this.deps.onError(sessionId, failure)
this.deps.logger.warn('recording a failed delivery failed', {
scope: 'delivery-loop-fail',
sessionId,
error: failure
})
})
}
}
@@ -22,10 +22,13 @@ export class StructuredAgentSessionEventRecovery {
publishStatus?: (sessionId: string) => void
serialize: <T>(sessionId: string, task: () => Promise<T>) => Promise<T>
now: () => number
onBarrierError: (sessionId: string, error: unknown) => void
}
) {}
private get exitContext() {
return { ...this.context, logger: this.context.deps.logger }
}
recoverAfterSinkFailure(sessionId: string, error: unknown): void {
if (this.sinkFailures.has(sessionId)) {
return
@@ -56,7 +59,16 @@ export class StructuredAgentSessionEventRecovery {
} as const
})
.then((event) => (event ? this.handle(event) : undefined))
.catch((recoveryError) => this.context.onBarrierError(sessionId, recoveryError))
.catch((error: unknown) =>
this.context.deps.logger.warn(
'stopping a provider after its journal failed did not finish',
{
scope: 'sink-failure-recovery',
sessionId,
error
}
)
)
.finally(() => this.sinkFailures.delete(sessionId))
}
@@ -66,6 +78,6 @@ export class StructuredAgentSessionEventRecovery {
if (event.type === 'started') {
return settleStructuredAgentSessionProviderStarted(this.context, event)
}
await settleUnexpectedStructuredAgentSessionExit(this.context, event)
await settleUnexpectedStructuredAgentSessionExit(this.exitContext, event)
}
}
@@ -45,7 +45,8 @@ export class StructuredAgentSessionSinkQueue {
constructor(
private readonly deps: {
watermarks: StructuredAgentSessionSinkWatermarks
onError?: (error: unknown) => void
/** The queue just failed for good; it accepts and runs nothing more. */
onFailed?: (error: unknown) => void
readingControl?: StructuredAgentSessionReadingControl
onBackpressureChange?: (
backpressured: boolean,
@@ -212,7 +213,7 @@ export class StructuredAgentSessionSinkQueue {
private fail = (error: unknown): void => {
if (this.failure === null) {
this.failure = { error }
this.deps.onError?.(error)
this.deps.onFailed?.(error)
}
this.queue.length = 0
this.queuedBytes = 0
@@ -15,6 +15,8 @@ import {
type StructuredAgentSessionEventTarget
} from './structured-agent-session-event-sink'
import { StructuredAgentSessionHostRuntimeState } from './structured-agent-session-host-runtime-state'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
import { testEventSinkLogging } from './structured-agent-session-logger-test-support'
const BODY: AgentJournalItemBody = {
kind: 'message',
@@ -91,7 +93,7 @@ function target(
describe('deferred structured agent-session event sink', () => {
it('buffers writes made before the journal exists and drains them in arrival order', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.sink.appendItem(identity(0), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
deferred.sink.appendItem(identity(1), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
@@ -110,7 +112,7 @@ describe('deferred structured agent-session event sink', () => {
it('writes at the fence bound at submission time, so a rebind cannot backdate a write', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(target(1, log))
deferred.sink.appendItem(identity(0), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
@@ -127,7 +129,7 @@ describe('deferred structured agent-session event sink', () => {
it('buffers replacement-acquisition events while unbound', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(target(1, log))
deferred.unbind()
@@ -141,7 +143,7 @@ describe('deferred structured agent-session event sink', () => {
it('resolves a lifecycle transition after journal bind and skips an existing state', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
expect(
deferred.sink.tryAppendLifecycleTransition?.(identity(0), BODY, () => identity(1), {
@@ -166,7 +168,7 @@ describe('deferred structured agent-session event sink', () => {
it('drops buffered and later writes once closed, and refuses to rebind', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.sink.appendItem(identity(0), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
deferred.close()
@@ -182,7 +184,8 @@ describe('deferred structured agent-session event sink', () => {
const errors: unknown[] = []
const readingControl = { pauseReading: vi.fn(), resumeReading: vi.fn() }
const deferred = createDeferredStructuredAgentSessionEventSink({
onError: (error) => errors.push(error),
...testEventSinkLogging(),
onFailed: (error) => errors.push(error),
readingControl
})
deferred.bind(target(4, log, 0))
@@ -206,7 +209,11 @@ describe('deferred structured agent-session event sink', () => {
})
it('replaces a failed cached sink before recovery drain', async () => {
const runtime = new StructuredAgentSessionHostRuntimeState({ store: {} } as never)
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: the cached-sink path reads only the logger, on the failed drain.
const runtime = new StructuredAgentSessionHostRuntimeState({
store: {},
logger: createStructuredAgentSessionLogger()
} as never)
const failed = runtime.eventSinkFor('session-1')
failed.bind(target(1, [], 0))
failed.sink.appendItem(identity(0), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
@@ -226,6 +233,7 @@ describe('deferred structured agent-session event sink', () => {
const changes: boolean[] = []
const readingControl = { pauseReading: vi.fn(), resumeReading: vi.fn() }
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
maxQueuedBytes: 1_000_000,
lowQueuedBytes: 0,
@@ -262,6 +270,7 @@ describe('deferred structured agent-session event sink', () => {
it('admits a resolved append and publication as one bounded operation', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedOperations: 1,
maxQueuedOperations: 2,
@@ -299,6 +308,7 @@ describe('deferred structured agent-session event sink', () => {
const changes: boolean[] = []
const readingControl = { pauseReading: vi.fn(), resumeReading: vi.fn() }
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedBytes: 1,
maxQueuedBytes: 1_000_000,
@@ -329,7 +339,8 @@ describe('deferred structured agent-session event sink', () => {
const log: Recorded[] = []
const errors: unknown[] = []
const deferred = createDeferredStructuredAgentSessionEventSink({
onError: (error) => errors.push(error),
...testEventSinkLogging(),
onFailed: (error) => errors.push(error),
watermarks: {
pauseQueuedBytes: 1,
maxQueuedBytes: 1,
@@ -364,6 +375,7 @@ describe('deferred structured agent-session event sink', () => {
const firstControl = { pauseReading: vi.fn(), resumeReading: vi.fn() }
const secondControl = { pauseReading: vi.fn(), resumeReading: vi.fn() }
const deferred = createDeferredStructuredAgentSessionEventSink({
...testEventSinkLogging(),
watermarks: {
pauseQueuedBytes: 1,
maxQueuedBytes: 1_000_000,
@@ -395,7 +407,7 @@ describe('deferred structured agent-session event sink', () => {
it('replaces a queued same-item checkpoint before it runs', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const options = { coalescingKey: 'checkpoint:item-1', turnScope: AGENT_JOURNAL_THREAD_SCOPE }
deferred.sink.appendItem(identity(0), BODY, options)
@@ -411,7 +423,7 @@ describe('deferred structured agent-session event sink', () => {
// The journal applies a settlement id once and skips any later batch with it,
// so the queue must not let a later batch replace one it has not run yet.
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const batch = (ordinal: number) => [
{
kind: 'item' as const,
@@ -440,7 +452,7 @@ describe('deferred structured agent-session event sink', () => {
it('keeps a replacement checkpoint after distinct intervening operations', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
const options = { coalescingKey: 'checkpoint:item-1', turnScope: AGENT_JOURNAL_THREAD_SCOPE }
deferred.sink.appendItem(identity(0), BODY, options)
@@ -457,7 +469,7 @@ describe('deferred structured agent-session event sink', () => {
it('coalesces provider activity as a publication without a journal write', async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.sink.setActivity?.({ turnId: 'turn-1', text: 'Thinking' })
deferred.sink.setActivity?.({ turnId: 'turn-1', text: 'Checking the result' })
@@ -486,7 +498,7 @@ describe('producer linkage reaches the journal through every append path', () =>
it('forwards the whole bundle on the plain and try append paths', async () => {
for (const append of ['appendItem', 'tryAppendItem'] as const) {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(target(5, log))
deferred.sink[append]?.(identity(1), BODY, {
...LINKAGE,
@@ -510,7 +522,7 @@ describe('producer linkage reaches the journal through every append path', () =>
'tryAppendLifecycleTransition'
] as const) {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(target(5, log))
deferred.sink[append]?.(identity(1), BODY, () => identity(1), {
...LINKAGE,
@@ -530,7 +542,7 @@ describe('producer linkage reaches the journal through every append path', () =>
// would stamp whoever opened the batch onto all of them. Both callers are
// single-producer today; a mixed batch would have to stamp per mutation.
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(target(5, log))
deferred.sink.appendLifecycleBatch?.(
'settle-1',
@@ -547,7 +559,7 @@ describe('producer linkage reaches the journal through every append path', () =>
it("writes no linkage keys at all for the session's own agent", async () => {
const log: Recorded[] = []
const deferred = createDeferredStructuredAgentSessionEventSink()
const deferred = createDeferredStructuredAgentSessionEventSink(testEventSinkLogging())
deferred.bind(target(5, log))
deferred.sink.appendItem(identity(1), BODY, { turnScope: AGENT_JOURNAL_THREAD_SCOPE })
await deferred.drained()
@@ -12,6 +12,7 @@ import { estimateStructuredAgentSessionItemBytes } from './structured-agent-sess
import { StructuredAgentSessionSinkQueue } from './structured-agent-session-event-sink-queue'
import { structuredAgentSessionJournalAppendOptions } from './structured-agent-session-journal-append-options'
import { createStructuredAgentSessionResolvedAppend } from './structured-agent-session-resolved-append'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
export type StructuredAgentSessionSinkAdmission =
| { accepted: true }
@@ -193,18 +194,28 @@ const DEFAULT_WATERMARKS: StructuredAgentSessionSinkWatermarks = {
maxLifecycleQueuedOperations: 1_024
}
export function createDeferredStructuredAgentSessionEventSink(
deps: {
onError?: (error: unknown) => void
watermarks?: Partial<StructuredAgentSessionSinkWatermarks>
readingControl?: StructuredAgentSessionReadingControl
onBackpressureChange?: (backpressured: boolean, state: StructuredAgentSessionSinkState) => void
} = {}
): DeferredStructuredAgentSessionEventSink {
export function createDeferredStructuredAgentSessionEventSink(deps: {
/** The session this sink writes for, named in every failure it logs. */
sessionId: string
logger: StructuredAgentSessionLogger
/** The sink failed for good; the owner decides what that costs the provider. */
onFailed?: (error: unknown) => void
watermarks?: Partial<StructuredAgentSessionSinkWatermarks>
readingControl?: StructuredAgentSessionReadingControl
onBackpressureChange?: (backpressured: boolean, state: StructuredAgentSessionSinkState) => void
}): DeferredStructuredAgentSessionEventSink {
const watermarks = { ...DEFAULT_WATERMARKS, ...deps.watermarks }
const failed = (error: unknown): void => {
deps.logger.error('writing provider events to the chat journal failed', {
scope: 'journal-event-sink',
sessionId: deps.sessionId,
error
})
deps.onFailed?.(error)
}
const queue = new StructuredAgentSessionSinkQueue({
watermarks,
...(deps.onError ? { onError: deps.onError } : {}),
onFailed: failed,
...(deps.readingControl ? { readingControl: deps.readingControl } : {}),
...(deps.onBackpressureChange ? { onBackpressureChange: deps.onBackpressureChange } : {})
})
@@ -280,7 +291,7 @@ export function createDeferredStructuredAgentSessionEventSink(
appendLifecycleBatch: (settlementId, mutations, options = {}) => {
const admission = appendLifecycleBatch(settlementId, mutations, options)
if (!admission.accepted) {
deps.onError?.(
failed(
new Error(
`lifecycle journal batch ${settlementId} rejected by sink ${admission.reason}`
)
@@ -10,12 +10,15 @@ import {
AgentSessionPreSpawnError
} from './structured-agent-session-adapter'
import { StructuredAgentSessionHostRuntimeState } from './structured-agent-session-host-runtime-state'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
function context(): StructuredAgentSessionEvictionContext & { order: string[] } {
const order: string[] = []
return {
order,
sessionId: 'session-1',
logger: createStructuredAgentSessionLogger(),
eventSink: {
unbind: vi.fn(() => order.push('unbind')),
drained: vi.fn(async () => {
@@ -44,9 +47,11 @@ function context(): StructuredAgentSessionEvictionContext & { order: string[] }
}
function runtimeState(): StructuredAgentSessionHostRuntimeState {
// oxlint-disable-next-line typescript/consistent-type-assertions -- SAFETY: eviction against the sink cache reads only the sinks and the logger; store and adapter are never reached.
return new StructuredAgentSessionHostRuntimeState({
store: {} as never,
adapter: {} as never
store: {},
adapter: {},
logger: recordingStructuredAgentSessionLogger().logger
} as never)
}
@@ -146,6 +151,7 @@ describe('rows the provider emits while closing', () => {
await evictStructuredAgentSession({
sessionId,
logger: recordingStructuredAgentSessionLogger().logger,
eventSink: sink,
adapter: {
closeSession: async () => {
@@ -227,6 +233,7 @@ describe('eviction against the real sink cache', () => {
const sessionId = 'session-reattach'
await evictStructuredAgentSession({
sessionId,
logger: recordingStructuredAgentSessionLogger().logger,
eventSink: state.eventSinkFor(sessionId),
adapter: { closeSession: async () => true } as never,
discardSink: () => state.discardEventSink(sessionId),
@@ -25,6 +25,7 @@ import { stopAgentSessionProviderRoot } from './structured-agent-session-provide
import type { DeferredStructuredAgentSessionEventSink } from './structured-agent-session-event-sink'
import type { StructuredAgentSessionStopVerdict } from './structured-agent-session-host-types'
import { withTimeout } from '../../../shared/promise-timeout-fallback'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
export type StructuredAgentSessionEvictionContext = {
sessionId: string
@@ -33,6 +34,7 @@ export type StructuredAgentSessionEvictionContext = {
hasProviderChild?: boolean
eventSink: DeferredStructuredAgentSessionEventSink
adapter: StructuredAgentSessionAdapter
logger: StructuredAgentSessionLogger
/** Tells the adapter the released lease is done with, so it drops this child's route and index.
* The conversation stays: stopping the agent never closes its journal. */
acknowledgeRelease: () => Promise<void> | void
@@ -76,7 +78,10 @@ export const STRUCTURED_AGENT_SESSION_EVICTION_STEPS: readonly StructuredAgentSe
try {
context.beforeProviderChildStop()
} catch {
console.warn('[structured-agent-session] capturing recovery witness failed')
context.logger.warn('capturing a recovery witness before a stop failed', {
scope: 'recovery-witness',
sessionId: context.sessionId
})
}
}
},
@@ -23,6 +23,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
const EXIT_REASON = 'claude stream-json exited (code 1): stderr tail'
@@ -51,6 +52,7 @@ beforeEach(async () => {
}))
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire,
@@ -21,6 +21,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
const EXIT_REASON = 'claude stream-json exited (code 1): claude: not signed in'
@@ -47,6 +48,7 @@ beforeEach(async () => {
}))
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire,
@@ -24,6 +24,7 @@ import { attachStructuredAgentSession } from './structured-agent-session-attach-
import type { StructuredAgentSessionAttachContext } from './structured-agent-session-attach-context'
import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types'
import { StructuredAgentSessionStatusFeed } from './structured-agent-session-status-feed'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
// Everything before the journal is out of scope here; what matters is what the orchestration does
// when the attach throws after acquisition.
@@ -146,6 +147,7 @@ async function workingSession(): Promise<{
const server = new AgentHookServer()
const records = new Map([[SESSION, ownerRecord()]])
const feed = new StructuredAgentSessionStatusFeed({
logger: createStructuredAgentSessionLogger(),
sessions,
getRecord: (sessionId) => records.get(sessionId) ?? null,
now: () => 1,
@@ -24,6 +24,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -133,6 +134,7 @@ beforeEach(async () => {
answerPrompt = vi.fn(async ({ commit }) => commit())
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: adapter(),
journalDatabase: openTestJournalHostDatabase(root),
@@ -78,7 +78,7 @@ export function createStructuredAgentSessionConversationDelivery(input: {
deps.store.getRecord(sessionId),
sessions.get(sessionId)?.journal
),
onError: (sessionId, error) => deps.onEventSinkError?.({ sessionId, error }),
logger: deps.logger,
record: (sessionId) => deps.store.getRecord(sessionId),
readChildWork: input.clientDelivery.readChildWork,
flushStreamedEvents: input.flushStreamedEvents,
@@ -116,7 +116,11 @@ export function createStructuredAgentSessionConversationDelivery(input: {
})
.catch((error: unknown) => {
wakesQueued.delete(sessionId)
deps.onEventSinkError?.({ sessionId, error })
deps.logger.warn('waking the delivery loop after a commit failed', {
scope: 'delivery-wake',
sessionId,
error
})
})
}
// A chat open before its owner's death was proven revises what its open settled. Queued, never
@@ -129,7 +133,13 @@ export function createStructuredAgentSessionConversationDelivery(input: {
resettleOpenStructuredAgentSessionConversation(deps, sessionId, sessions.get(sessionId))
)
)
.catch((error: unknown) => deps.onEventSinkError?.({ sessionId, error }))
.catch((error: unknown) =>
deps.logger.warn('resettling an open chat after its owner died failed', {
scope: 'death-evidence-resettle',
sessionId,
error
})
)
}
})
return {
@@ -158,8 +168,12 @@ async function settleInterruptedCommands(
): Promise<void> {
const fence = structuredAgentSessionConversationFence(deps.store, sessionId)
try {
await recoverStructuredRewind(deps.store, sessionId, session.journal, fence)
await recoverStructuredRewind(deps, sessionId, session.journal, fence)
} catch (error) {
deps.onEventSinkError?.({ sessionId, error })
deps.logger.warn('settling an interrupted rewind on open failed', {
scope: 'rewind-recovery',
sessionId,
error
})
}
}
@@ -47,7 +47,7 @@ export type StructuredAgentSessionLifetimeContext = {
}
}
type ConversationCloseDeps = Pick<StructuredAgentSessionHostDeps, 'onEventSinkError'> & {
type ConversationCloseDeps = Pick<StructuredAgentSessionHostDeps, 'logger'> & {
store: Pick<StructuredAgentSessionHostDeps['store'], 'getRecord'>
}
@@ -70,7 +70,11 @@ export async function abandonQueuedStructuredAgentSessionMessages(
.then(
() => true,
(error: unknown) => {
deps.onEventSinkError?.({ sessionId, error })
deps.logger.warn('rejecting queued messages of a closed chat failed', {
scope: 'queued-abandon',
sessionId,
error
})
return false
}
)
@@ -113,7 +117,6 @@ export async function stopStructuredAgentSessionAgentUnderSerialize(
const owed = owedProviderChildWindDown(session)
session.owesProviderChildWindDown = owed ? { ...owed, cause } : undefined
const stopping = session.child
let settlementError: unknown
const eviction: StructuredAgentSessionEvictionContext = {
sessionId,
// The retry must not re-stop a child the adapter already proved gone, so this stays honest.
@@ -121,6 +124,7 @@ export async function stopStructuredAgentSessionAgentUnderSerialize(
owesProviderChildWindDown: owed !== undefined,
eventSink: context.runtimeState.eventSinkFor(sessionId),
adapter: context.deps.adapter,
logger: context.deps.logger,
// The adapter settles its own open turn with this, so who asked travels with the stop.
stopCause: cause,
...(context.restartWitness
@@ -153,15 +157,16 @@ export async function stopStructuredAgentSessionAgentUnderSerialize(
pendingSubmissionReason: 'provider_closed_before_acknowledgement',
// Only a turn no adapter settled: one with no close, or whose settle threw.
verdict: turnVerdictForChildEnd(cause, context.now()),
showUnexpectedExitOutcome: false,
onError: (id, error) => {
settlementError = error
context.deps.onEventSinkError?.({ sessionId: id, error })
}
showUnexpectedExitOutcome: false
})
if (!settled) {
if (!settled.ok) {
context.deps.logger.warn("settling a closed agent's work failed", {
scope: 'close-settlement',
sessionId,
error: settled.error
})
// Without the cause the log names the step and nothing else.
throw new Error('dead generation work settlement failed', { cause: settlementError })
throw new Error('dead generation work settlement failed', { cause: settled.error })
}
},
releaseLease: async () => {
@@ -13,6 +13,8 @@ import {
} from '../agent-session-journal/journal-host-database-test-support'
import { StructuredAgentSessionHostRuntimeState } from './structured-agent-session-host-runtime-state'
import type { StructuredAgentSessionHostDeps } from './structured-agent-session-host'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const NOW = 1_800_000_000_000
@@ -68,7 +70,8 @@ function runtimeState(
adapter: {},
journalDatabase: openTestJournalHostDatabase(stateDirectory),
claimKeyId: 'key-1',
probeOwner
probeOwner,
logger: createStructuredAgentSessionLogger()
} as StructuredAgentSessionHostDeps
return new StructuredAgentSessionHostRuntimeState(deps)
}
@@ -129,7 +132,7 @@ describe('host runtime-state owner probe', () => {
it('does not force-close a provider for transient lease probe errors', async () => {
const onEventSinkFailure = vi.fn()
const onEventSinkError = vi.fn()
const log = recordingStructuredAgentSessionLogger()
const probeOwner = vi.fn(async () => {
throw new Error('lease probe unavailable')
})
@@ -144,7 +147,7 @@ describe('host runtime-state owner probe', () => {
journalDatabase: openTestJournalHostDatabase(stateDirectory),
claimKeyId: 'key-1',
probeOwner,
onEventSinkError
logger: log.logger
} as unknown as StructuredAgentSessionHostDeps
const state = new StructuredAgentSessionHostRuntimeState(deps, onEventSinkFailure)
@@ -152,10 +155,15 @@ describe('host runtime-state owner probe', () => {
state as unknown as { leaseRenewer: { renewNow: () => Promise<void> } }
).leaseRenewer.renewNow()
expect(onEventSinkError).toHaveBeenCalledWith({
sessionId: record.sessionId,
error: expect.any(Error)
})
expect(log.entries).toContainEqual(
expect.objectContaining({
fields: expect.objectContaining({
scope: 'lease-renewal',
sessionId: record.sessionId,
error: expect.any(Error)
})
})
)
expect(onEventSinkFailure).not.toHaveBeenCalled()
})
})
@@ -26,7 +26,7 @@ export class StructuredAgentSessionHostRuntimeState {
now: () => deps.now?.() ?? Date.now(),
// Lease/ownership failures are transient and stay on the visible lease-error path.
// Only deferred sink I/O failures are terminal and may force-close a provider.
onError: ({ sessionId, error }) => deps.onEventSinkError?.({ sessionId, error })
logger: deps.logger
})
}
@@ -67,8 +67,9 @@ export class StructuredAgentSessionHostRuntimeState {
mintEventSink(sessionId: string): DeferredStructuredAgentSessionEventSink {
const minted: DeferredStructuredAgentSessionEventSink =
createDeferredStructuredAgentSessionEventSink({
onError: (error) => {
this.deps.onEventSinkError?.({ sessionId, error })
sessionId,
logger: this.deps.logger,
onFailed: (error) => {
// Only the session's own sink may force its provider down; an attempt's never is.
if (this.eventSinks.get(sessionId) === minted) {
this.onEventSinkFailure?.(sessionId, error)
@@ -1,4 +1,5 @@
import type { AgentSessionRecord } from '../../../shared/agent-session-record'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
/**
* Chat-tab visibility is the deletion funnel: every path that removes a chat as a user-facing
@@ -9,6 +10,7 @@ import type { AgentSessionRecord } from '../../../shared/agent-session-record'
export function setStructuredAgentSessionTabVisibility(
host: {
deps: {
logger: StructuredAgentSessionLogger
store: {
setSessionTabVisibility: (
sessionId: string,
@@ -25,7 +27,10 @@ export function setStructuredAgentSessionTabVisibility(
): Promise<void> {
if (!visible) {
void host.restartResume.dismiss([sessionId]).catch(() => {
console.warn('[structured-agent-session] forgetting recovery records on chat close failed')
host.deps.logger.warn('forgetting recovery records on chat close failed', {
scope: 'tab-close-recovery-dismiss',
sessionId
})
})
}
return host.deps.store.setSessionTabVisibility(sessionId, visible, tabId)
@@ -13,12 +13,15 @@ import {
RESUME_MARKER_RECORD_TIMEOUT_MS,
structuredAgentSessionHostTeardownPhases
} from './structured-agent-session-host-teardown'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
const noop = async (): Promise<void> => undefined
describe('structured agent-session host teardown', () => {
it('names every phase, so the quit-path order is pinned rather than incidental', () => {
const phases = structuredAgentSessionHostTeardownPhases({
logger: createStructuredAgentSessionLogger(),
idleSweep: { dispose: noop },
runtimeState: { stopLeaseRenewal: () => undefined, flushAllEventSinks: noop },
tasks: { drainAttaches: noop },
@@ -62,6 +65,7 @@ describe('structured agent-session host teardown', () => {
it('ends child eviction as soon as every chat has closed, not at its bound', async () => {
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
const evict = structuredAgentSessionHostTeardownPhases({
logger: createStructuredAgentSessionLogger(),
idleSweep: { dispose: noop },
runtimeState: { stopLeaseRenewal: () => undefined, flushAllEventSinks: noop },
tasks: { drainAttaches: noop },
@@ -86,11 +90,12 @@ describe('structured agent-session host teardown', () => {
it('bounds stalled recovery publication without preventing later cleanup', async () => {
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
const warning = vi.spyOn(console, 'warn').mockImplementation(() => {})
const log = recordingStructuredAgentSessionLogger()
const pending = Promise.withResolvers<void>()
const cleaned = vi.fn(async () => {})
const flush = vi.fn(async () => cleaned())
const phases = structuredAgentSessionHostTeardownPhases({
logger: log.logger,
idleSweep: { dispose: cleaned },
runtimeState: { stopLeaseRenewal: () => {}, flushAllEventSinks: flush },
tasks: { drainAttaches: cleaned },
@@ -107,13 +112,10 @@ describe('structured agent-session host teardown', () => {
await vi.advanceTimersByTimeAsync(2000)
await teardown
expect(cleaned).toHaveBeenCalledTimes(4)
expect(warning).toHaveBeenCalledWith(
'[structured-agent-session] recording recovery capsule failed'
)
expect(log.scopes()).toEqual(['teardown-recovery-capsule'])
expect(vi.getTimerCount()).toBe(0)
} finally {
pending.resolve()
warning.mockRestore()
vi.useRealTimers()
}
})
@@ -17,6 +17,7 @@ import {
} from './structured-agent-session-host-lifetime'
import { withTimeout } from '../../../shared/promise-timeout-fallback'
import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
export type StructuredAgentSessionTeardownPhase = {
name: string
@@ -66,6 +67,7 @@ export function structuredAgentSessionHostTeardownPhases(collaborators: {
/** Opens this teardown's witnesses; each session's own is taken as eviction stops its child. */
beginResumeMarkers: () => void
recordResumeMarkers: () => Promise<void>
logger: StructuredAgentSessionLogger
}): StructuredAgentSessionTeardownPhase[] {
return [
{
@@ -74,7 +76,9 @@ export function structuredAgentSessionHostTeardownPhases(collaborators: {
try {
collaborators.beginResumeMarkers()
} catch {
console.warn('[structured-agent-session] capturing recovery witnesses failed')
collaborators.logger.warn('capturing recovery witnesses for teardown failed', {
scope: 'teardown-recovery-witnesses'
})
}
}
},
@@ -90,7 +94,9 @@ export function structuredAgentSessionHostTeardownPhases(collaborators: {
run: () =>
withPhaseTimeout(collaborators.recordResumeMarkers, RESUME_MARKER_RECORD_TIMEOUT_MS).catch(
() => {
console.warn('[structured-agent-session] recording recovery capsule failed')
collaborators.logger.warn('recording the recovery capsule at teardown failed', {
scope: 'teardown-recovery-capsule'
})
}
)
},
@@ -168,7 +174,8 @@ export async function flushStructuredAgentSessionHost(
retainSessionIds
),
beginResumeMarkers: () => context.restartResume.beginTeardown(context.trigger),
recordResumeMarkers: context.restartResume.recordMarkers
recordResumeMarkers: context.restartResume.recordMarkers,
logger: context.deps.logger
}),
sessions: context.sessions,
retainSessionIds,
@@ -19,6 +19,7 @@ import {
hostTestMessage
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -66,6 +67,7 @@ describe('abandoning a structured agent-session host', () => {
setOption: async () => undefined
}
const host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter,
probeOwner: async () => ({
@@ -29,6 +29,7 @@ import {
hostTestOperationId,
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { recordingProductionStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
const journals = createTrackedJournalOpener()
@@ -62,6 +63,7 @@ let root: string
let recoveryCapsule: TrackedTestRecoveryCapsule
let store: AgentSessionRecordStore
let host: StructuredAgentSessionHost
let log: ReturnType<typeof recordingProductionStructuredAgentSessionLogger>
let acquire: Mock<StructuredAgentSessionAdapter['acquire']>
let releaseAcquisition: Mock<NonNullable<StructuredAgentSessionAdapter['releaseAcquisition']>>
let dispatch: Mock<StructuredAgentSessionAdapter['dispatch']>
@@ -152,7 +154,9 @@ beforeEach(async () => {
setOption = vi.fn(async () => undefined)
store = await openTestAgentSessionRecordStore(root)
recoveryCapsule = new TrackedTestRecoveryCapsule(root)
log = recordingProductionStructuredAgentSessionLogger()
host = new StructuredAgentSessionHost({
logger: log.logger,
store,
adapter: adapter(),
journalDatabase: openTestJournalHostDatabase(root),
@@ -200,6 +204,8 @@ export function hostTestState() {
root,
store,
host,
/** Every entry the beforeEach host logged; a host a test builds itself logs elsewhere. */
log,
acquire,
releaseAcquisition,
dispatch,
@@ -17,6 +17,7 @@ import type {
import type { AgentSessionAttachParams } from './structured-agent-session-attach'
import type { StructuredAgentSessionStatusSink } from './structured-agent-session-status-feed'
import type { AgentModelCatalogService } from '../agent-model-catalog/agent-model-catalog-service'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
export type StructuredAgentSessionCaller = { callerKey: string }
@@ -121,11 +122,9 @@ export type StructuredAgentSessionHostDeps = {
idleSweep?: { intervalMs?: number; idleMs?: number }
/** Whether an orchestration dispatch still owns this session's worker; absent answers no. */
hasOpenDispatch?: (record: AgentSessionRecord) => boolean
onEventSinkError?: (input: { sessionId: string; error: unknown }) => void
/** Lease bookkeeping run for startup or a read (the reconcile, or resolving a chat's recovery)
* that refused or threw, once per distinct failure. Startup and the read carry on: the next
* attach or send reconciles and resolves recovery again before it acts. */
onLeaseReconcileFailure?: (failure: unknown) => void
/** Where every failure the host carries on past is reported. Required: a host without one would
* drop exactly the failures nobody sees in the UI. */
logger: StructuredAgentSessionLogger
/** Every status projection this host publishes. `replay` marks a re-projection of state the host
* already knew (restore, an arriving subscriber) rather than a fresh journal edge. */
onSessionStatusChanged?: (
@@ -28,6 +28,7 @@ import {
hostTestMessage
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
let root: string
let store: AgentSessionRecordStore
@@ -132,6 +133,7 @@ describe('attach', () => {
}
}))
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: { ...adapter(), acquire },
journalDatabase: openTestJournalHostDatabase(root),
@@ -562,6 +564,7 @@ describe('restart', () => {
) {
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: { ...adapter(), ...adapterOverrides },
journalDatabase: openTestJournalHostDatabase(root),
@@ -56,6 +56,7 @@ import { structuredAgentSessionRestartResumeSurfaces } from './structured-agent-
import { createStructuredAgentSessionConversationDelivery } from './structured-agent-session-host-delivery'
import { structuredAgentSessionConversationFence } from './structured-agent-session-provider-child'
import { wireStructuredAgentSessionQueuedMessages } from './structured-agent-session-queued-wiring'
import * as sessionLogger from './structured-agent-session-logger'
export type { StructuredAgentSessionHostDeps } from './structured-agent-session-host-types'
export class StructuredAgentSessionHost {
@@ -68,7 +69,7 @@ export class StructuredAgentSessionHost {
this.subscribers.publish(sessionId, journal)
this.conversationDelivery.afterCommit(sessionId, journal)
},
onDeliveryError: (sessionId, error) => this.deps.onEventSinkError?.({ sessionId, error }),
logger: sessionLogger.deferredStructuredAgentSessionLogger(() => this.deps.logger),
onOpened: (sessionId) => this.queued.drain.schedule(sessionId),
now: () => this.now()
})
@@ -101,6 +102,8 @@ export class StructuredAgentSessionHost {
readonly restartResume: StructuredAgentSessionRestartResume
constructor(readonly deps: StructuredAgentSessionHostDeps) {
// Every collaborator reads this copy, so a logger that throws cannot fail what it reports.
this.deps = deps = sessionLogger.withNeverThrowingLogger(deps)
this.clientDelivery.watchAtRestCommands(deps.adapter)
this.backgroundTasks = new StructuredAgentSessionBackgroundTaskChannel(
deps,
@@ -158,8 +161,7 @@ export class StructuredAgentSessionHost {
),
publishStatus: this.clientDelivery.publishStatusAndSettlement,
serialize: (sessionId, task) => this.tasks.trackAttach(this.serialize(sessionId, task)),
now: () => this.now(),
onBarrierError: (sessionId, error) => deps.onEventSinkError?.({ sessionId, error })
now: () => this.now()
})
this.restartResume = createStructuredAgentSessionRestartResume(deps, this.sessions, {
...structuredAgentSessionRestartResumeSurfaces(this, this.now),
@@ -320,8 +320,14 @@ describe('the idle sweep with no child running (P2-22 ii)', () => {
stopAgent,
stopStartingAgent: stopAgent,
closeConversation,
onError: (_id, error) => {
throw error
// A failed step fails the test.
logger: {
warn: (_message, fields) => {
throw fields.error
},
error: (_message, fields) => {
throw fields.error
}
}
})
await sweep.tick()
@@ -14,6 +14,7 @@ import { isQueuedAgentJournalSubmission } from '../../../shared/agent-session-qu
import type { AgentJournalRenderItem } from '../../../shared/agent-session-journal-types'
import type { AgentChildWorkView } from '../../../shared/agent-status-child-work-view'
import type { StructuredAgentSessionHostSession } from './structured-agent-session-host-types'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
export const STRUCTURED_AGENT_SESSION_IDLE_SWEEP_INTERVAL_MS = 5 * 60_000
export const STRUCTURED_AGENT_SESSION_IDLE_MS = 30 * 60_000
@@ -37,7 +38,7 @@ export type StructuredAgentSessionIdleSweepDeps = {
stopAgent: (sessionId: string) => Promise<void>
stopStartingAgent: (sessionId: string) => Promise<void>
closeConversation: (sessionId: string) => Promise<boolean>
onError: (sessionId: string, error: unknown) => void
logger: StructuredAgentSessionLogger
intervalMs?: number
idleMs?: number
}
@@ -87,7 +88,13 @@ export class StructuredAgentSessionIdleSweep {
[...this.deps.sessions.keys()].map((sessionId) =>
this.deps
.serialize(sessionId, () => this.tickUnderSerialize(sessionId))
.catch((error: unknown) => this.deps.onError(sessionId, error))
.catch((error: unknown) =>
this.deps.logger.warn('an idle sweep step failed', {
scope: 'idle-sweep',
sessionId,
error
})
)
)
)
} finally {
@@ -25,6 +25,7 @@ import {
HOST_TEST_SESSION as SESSION,
hostTestMessage
} from './structured-agent-session-host-test-data'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const { copy } = vi.hoisted(() => ({ copy: { mismatched: false } }))
@@ -93,6 +94,7 @@ async function relaunchWithMismatchedCopy(): Promise<StructuredAgentSessionHost>
})
const store = await openTestAgentSessionRecordStore(relaunched)
const host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: adapter(),
journalDatabase: openTestJournalHostDatabase(relaunched),
@@ -16,6 +16,9 @@ import {
} from './structured-agent-session-host-test-harness'
import { hostTestMessage } from './structured-agent-session-host-test-data'
const STOP_LEDGER_ROW_FAILED =
"[agent-session] stop-ledger-row: writing Stop's ledger row failed; Stop runs without it"
let root: string
let store: AgentSessionRecordStore
let host: StructuredAgentSessionHost
@@ -70,9 +73,10 @@ it('refuses a send as corrupt when SQLite reports damage, and still stops the ag
// the caller sees, after the fact.
await expect(stop()).rejects.toBe(damaged)
expect(cancelTurn).toHaveBeenCalledTimes(1)
expect(warn).toHaveBeenCalledWith("[agent-session] Stop's ledger row skipped:", {
expect(warn).toHaveBeenCalledWith(STOP_LEDGER_ROW_FAILED, {
scope: 'stop-ledger-row',
sessionId: expect.any(String),
error: 'database disk image is malformed'
error: expect.objectContaining({ message: 'database disk image is malformed' })
})
// The restart-offer withdrawal the attach started holds its lock until it ends.
await hostTestRecoveryCapsuleSettled()
@@ -93,9 +97,10 @@ it.each([
value: { turnId: 'turn-1', cancelled: true }
})
expect(cancelTurn).toHaveBeenCalledTimes(1)
expect(warn).toHaveBeenCalledWith("[agent-session] Stop's ledger row skipped:", {
expect(warn).toHaveBeenCalledWith(STOP_LEDGER_ROW_FAILED, {
scope: 'stop-ledger-row',
sessionId: expect.any(String),
error: error.message
error
})
})
@@ -29,6 +29,7 @@ import {
import { agentSessionFailureFact } from '../../../shared/agent-session-failure'
import { agentSessionFailureWords } from '../../../shared/agent-session-failure-words'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -95,6 +96,7 @@ beforeEach(async () => {
closeSession = vi.fn(async () => true)
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire: vi.fn(async ({ fence }) => ({
@@ -9,6 +9,7 @@ import {
import type { AgentSessionRecordStore } from '../../runtime/agent-session-record-store'
import { openTestAgentSessionRecordStore } from '../../runtime/agent-session-record-store-test-harness'
import { StructuredAgentSessionLeaseRenewer } from './structured-agent-session-lease-renewer'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
const NOW = 1_800_000_000_000
const roots: string[] = []
@@ -103,6 +104,7 @@ describe('structured agent-session lease renewal', () => {
)
)
const renewer = new StructuredAgentSessionLeaseRenewer({
logger: recordingStructuredAgentSessionLogger().logger,
store: { listRecords: () => records, renewLeases } as unknown as AgentSessionRecordStore,
probe: vi.fn(),
probeMany,
@@ -144,7 +146,7 @@ describe('structured agent-session lease renewal', () => {
}
return records[0]!
})
const onError = vi.fn()
const log = recordingStructuredAgentSessionLogger()
const renewer = new StructuredAgentSessionLeaseRenewer({
store: {
listRecords: () => records,
@@ -156,14 +158,15 @@ describe('structured agent-session lease renewal', () => {
matchedOn: ['spawn-token' as const]
}),
now: () => NOW + 10_000,
onError
logger: log.logger
})
await renewer.renewNow()
expect(renewLeases).toHaveBeenCalledOnce()
expect(renewLease).toHaveBeenCalledTimes(2)
expect(onError).toHaveBeenCalledWith({
expect(log.entries.map((entry) => entry.fields)).toContainEqual({
scope: 'lease-renewal',
sessionId: 'session-b',
error: expect.objectContaining({ message: 'agent_session_checkpoint_stale' })
})
@@ -174,6 +177,7 @@ describe('structured agent-session lease renewal', () => {
const store = await liveStore()
let now = NOW
const renewer = new StructuredAgentSessionLeaseRenewer({
logger: recordingStructuredAgentSessionLogger().logger,
store,
probe: async () => ({
outcome: 'identity-matched',
@@ -203,6 +207,7 @@ describe('structured agent-session lease renewal', () => {
releaseProbe = resolve
})
const renewer = new StructuredAgentSessionLeaseRenewer({
logger: recordingStructuredAgentSessionLogger().logger,
store,
probe: async () => {
await probing
@@ -228,6 +233,7 @@ describe('structured agent-session lease renewal', () => {
matchedOn: ['process-start-time' as const]
}))
const renewer = new StructuredAgentSessionLeaseRenewer({
logger: recordingStructuredAgentSessionLogger().logger,
store,
probe,
now: () => NOW + 10_000
@@ -241,18 +247,19 @@ describe('structured agent-session lease renewal', () => {
it('stops extending the lease when child proof is no longer sufficient', async () => {
const store = await liveStore()
const onError = vi.fn()
const log = recordingStructuredAgentSessionLogger()
const renewer = new StructuredAgentSessionLeaseRenewer({
store,
probe: async () => ({ outcome: 'indeterminate', reason: 'probe unavailable' }),
now: () => NOW + 10_000,
onError
logger: log.logger
})
await renewer.renewNow()
expect(store.getRecord('session-renewal')?.lease.lastRenewedAt).toBe(NOW)
expect(onError).toHaveBeenCalledWith({
expect(log.entries.map((entry) => entry.fields)).toContainEqual({
scope: 'lease-renewal',
sessionId: 'session-renewal',
error: expect.any(Error)
})
@@ -271,6 +278,7 @@ describe('structured agent-session lease renewal', () => {
matchedOn: ['process-start-time' as const]
}))
const renewer = new StructuredAgentSessionLeaseRenewer({
logger: recordingStructuredAgentSessionLogger().logger,
store,
probe,
now: () => NOW + 10_000
@@ -4,13 +4,14 @@ import {
AGENT_SESSION_LEASE_TTL_MS,
type AgentSessionRecordStore
} from '../../runtime/agent-session-record-store'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
const RENEW_INTERVAL_MS = Math.floor(AGENT_SESSION_LEASE_TTL_MS / 3)
export class StructuredAgentSessionLeaseRenewer {
private timer: ReturnType<typeof setInterval> | null = null
private running = false
/** The tick in flight. Never rejects: the timer path reports renewal failures through `onError`,
/** The tick in flight. Never rejects: the timer path logs renewal failures,
* and stopping must not turn one into a teardown failure as well. */
private inFlight: Promise<void> = Promise.resolve()
@@ -22,7 +23,7 @@ export class StructuredAgentSessionLeaseRenewer {
records: readonly AgentSessionRecord[]
) => Promise<Map<string, AgentSessionOwnerProbe>>
now: () => number
onError?: (input: { sessionId: string; error: unknown }) => void
logger: StructuredAgentSessionLogger
intervalMs?: number
}
) {}
@@ -106,7 +107,7 @@ export class StructuredAgentSessionLeaseRenewer {
if (result.status === 'rejected') {
const renewal = renewals[index]
if (renewal) {
this.input.onError?.({ sessionId: renewal.sessionId, error: result.reason })
this.reportFailure(renewal.sessionId, result.reason)
}
}
})
@@ -126,15 +127,24 @@ export class StructuredAgentSessionLeaseRenewer {
if (result.status === 'fulfilled') {
probes.set(record.sessionId, result.value)
} else {
this.input.onError?.({ sessionId: record.sessionId, error: result.reason })
this.reportFailure(record.sessionId, result.reason)
}
}
return probes
} catch (error) {
for (const record of records) {
this.input.onError?.({ sessionId: record.sessionId, error })
this.reportFailure(record.sessionId, error)
}
return new Map()
}
}
/** Lease and ownership failures are transient: the next tick, attach or send retries them. */
private reportFailure(sessionId: string, error: unknown): void {
this.input.logger.warn('renewing a chat lease failed', {
scope: 'lease-renewal',
sessionId,
error
})
}
}
@@ -24,6 +24,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
@@ -42,6 +43,7 @@ let dispatch: Mock<StructuredAgentSessionAdapter['dispatch']>
function openHost(): void {
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire,
@@ -0,0 +1,54 @@
// Keeps one failure that recurs on a timer (lease renewal, idle sweep) or on every publication
// (status feed) from flooding the shared, rotated trace file and erasing the history around it.
/** How long a repeated entry stays quiet after it is written. */
export const STRUCTURED_AGENT_SESSION_LOG_REPEAT_WINDOW_MS = 5 * 60_000
const MAX_TRACKED_REPEATS = 256
type TrackedRepeat = { writtenAt: number; suppressed: number }
export type StructuredAgentSessionLogRepeats = {
/** `null` swallows this entry; a number writes it, carrying the repeats swallowed since the last. */
admit: (key: string) => number | null
}
export function createStructuredAgentSessionLogRepeats(options?: {
now?: () => number
windowMs?: number
maxTracked?: number
}): StructuredAgentSessionLogRepeats {
const now = options?.now ?? Date.now
const windowMs = options?.windowMs ?? STRUCTURED_AGENT_SESSION_LOG_REPEAT_WINDOW_MS
const maxTracked = options?.maxTracked ?? MAX_TRACKED_REPEATS
// Insertion order is write order, so the first key is the one written longest ago.
const tracked = new Map<string, TrackedRepeat>()
const makeRoom = (at: number): void => {
for (const [key, repeat] of tracked) {
if (at - repeat.writtenAt >= windowMs) {
tracked.delete(key)
}
}
const oldest = tracked.keys().next()
if (tracked.size >= maxTracked && !oldest.done) {
tracked.delete(oldest.value)
}
}
return {
admit: (key) => {
const at = now()
const repeat = tracked.get(key)
if (repeat && at - repeat.writtenAt < windowMs) {
repeat.suppressed += 1
return null
}
tracked.delete(key)
if (tracked.size >= maxTracked) {
makeRoom(at)
}
tracked.set(key, { writtenAt: at, suppressed: 0 })
return repeat?.suppressed ?? 0
}
}
}
@@ -0,0 +1,58 @@
import {
createStructuredAgentSessionLogger,
type StructuredAgentSessionLogFields,
type StructuredAgentSessionLogger
} from './structured-agent-session-logger'
export type RecordedStructuredAgentSessionLog = {
level: keyof StructuredAgentSessionLogger
message: string
fields: StructuredAgentSessionLogFields
}
/** A logger that keeps every entry, for tests that assert a failure was reported. */
export function recordingStructuredAgentSessionLogger(): {
logger: StructuredAgentSessionLogger
entries: RecordedStructuredAgentSessionLog[]
scopes: () => string[]
} {
const entries: RecordedStructuredAgentSessionLog[] = []
return {
logger: {
warn: (message, fields) => entries.push({ level: 'warn', message, fields }),
error: (message, fields) => entries.push({ level: 'error', message, fields })
},
entries,
scopes: () => entries.map((entry) => entry.fields.scope)
}
}
/** A recording logger that also prints as production does, for a harness whose tests read either:
* the entries hold every level and every field, whatever the console mirror keeps. */
export function recordingProductionStructuredAgentSessionLogger(): ReturnType<
typeof recordingStructuredAgentSessionLogger
> {
const recording = recordingStructuredAgentSessionLogger()
const production = createStructuredAgentSessionLogger()
return {
...recording,
logger: {
warn: (message, fields) => {
recording.logger.warn(message, fields)
production.warn(message, fields)
},
error: (message, fields) => {
recording.logger.error(message, fields)
production.error(message, fields)
}
}
}
}
/** What an event sink built outside a host needs: the session it writes for, and a logger. */
export function testEventSinkLogging(sessionId = 'session-1'): {
sessionId: string
logger: StructuredAgentSessionLogger
} {
return { sessionId, logger: recordingStructuredAgentSessionLogger().logger }
}
@@ -0,0 +1,212 @@
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'
import { _resetTracerForTests, setActiveSink } from '../../observability/tracer'
import {
AgentSessionRefusalError,
agentSessionRefusalError,
refuse
} from '../../../shared/agent-session-wire-refusals'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
import {
STRUCTURED_AGENT_SESSION_LOG_REPEAT_WINDOW_MS as WINDOW_MS,
createStructuredAgentSessionLogRepeats
} from './structured-agent-session-log-repeats'
let push: Mock<(record: unknown) => void>
let clock: number
function written(): { name: string; attributes: Record<string, unknown>; exit: unknown }[] {
// As the sink writes them: one JSON line per span.
return JSON.parse(JSON.stringify(push.mock.calls.map(([span]) => span)))
}
beforeEach(() => {
push = vi.fn<(record: unknown) => void>()
clock = 1_000_000
setActiveSink({ push, flush: () => {}, close: () => {} })
vi.spyOn(console, 'warn').mockImplementation(() => undefined)
vi.spyOn(console, 'error').mockImplementation(() => undefined)
})
afterEach(() => {
_resetTracerForTests()
vi.restoreAllMocks()
})
describe('a failure that repeats', () => {
it('is written once per window, then once more carrying how many repeats it swallowed', () => {
const logger = createStructuredAgentSessionLogger({ now: () => clock })
const renew = (): void =>
logger.warn('renewing a chat lease failed', {
scope: 'lease-renewal',
sessionId: 'session-1',
error: new Error('database is locked')
})
for (let tick = 0; tick < 30; tick += 1) {
renew()
clock += 10_000
}
expect(written()).toHaveLength(1)
expect(written()[0]?.attributes).not.toHaveProperty('suppressed')
expect(console.warn).toHaveBeenCalledOnce()
clock += WINDOW_MS
renew()
expect(written()).toHaveLength(2)
expect(written()[1]?.attributes).toMatchObject({ suppressed: 29, sessionId: 'session-1' })
expect(console.warn).toHaveBeenLastCalledWith(
'[agent-session] lease-renewal: renewing a chat lease failed',
expect.objectContaining({ suppressed: 29 })
)
renew()
expect(written()).toHaveLength(2)
})
it('never merges different sessions, scopes, levels, messages or errors', () => {
const logger = createStructuredAgentSessionLogger({ now: () => clock })
for (let round = 0; round < 3; round += 1) {
logger.warn('renewing a chat lease failed', { scope: 'lease-renewal', sessionId: 'a' })
logger.warn('renewing a chat lease failed', { scope: 'lease-renewal', sessionId: 'b' })
logger.warn('renewing a chat lease failed', { scope: 'idle-sweep', sessionId: 'a' })
logger.error('renewing a chat lease failed', { scope: 'lease-renewal', sessionId: 'a' })
logger.warn('a different step failed', { scope: 'lease-renewal', sessionId: 'a' })
logger.warn('a different step failed', {
scope: 'lease-renewal',
sessionId: 'a',
error: new Error('disk full')
})
}
expect(written().map((span) => [span.name, span.attributes['sessionId']])).toEqual([
['agentSession.lease-renewal', 'a'],
['agentSession.lease-renewal', 'b'],
['agentSession.idle-sweep', 'a'],
['agentSession.lease-renewal', 'a'],
['agentSession.lease-renewal', 'a'],
['agentSession.lease-renewal', 'a']
])
})
it('writes a failure whose cause, refusal reason or value differs, and swallows a true repeat', () => {
const logger = createStructuredAgentSessionLogger({ now: () => clock })
// A refusal's message is its bare code: only the cause tells these two apart.
const unreadable = (cause: Error): Error =>
new AgentSessionRefusalError(
refuse(
'agent_session_journal_unreadable',
{ reason: 'journalUnavailable' },
'agent_session_journal_unreadable'
),
{ cause }
)
const report = (error: unknown): void =>
logger.warn('recording an opened chat tab failed', {
scope: 'tab-visibility-open',
sessionId: 'session-1',
error
})
const denied = Object.assign(new Error('EACCES: permission denied'), { code: 'EACCES' })
const failing = Object.assign(new Error('EIO: i/o error'), { code: 'EIO' })
report(unreadable(denied))
report(unreadable(denied))
report(unreadable(failing))
report(
agentSessionRefusalError('agent_session_journal_unreadable', { reason: 'journalCorrupt' })
)
report({ code: 'busy', attempt: 1 })
report({ attempt: 1, code: 'busy' })
report({ code: 'locked', attempt: 1 })
expect(
written().map((span) => span.attributes['errorCause'] ?? span.attributes['error'])
).toEqual([
['Error: EACCES: permission denied'],
['Error: EIO: i/o error'],
undefined,
{ code: 'busy', attempt: 1 },
{ code: 'locked', attempt: 1 }
])
expect(written()[2]?.attributes).toMatchObject({ refusalReason: 'journalCorrupt' })
expect(written()[0]?.attributes).toMatchObject({ refusalReason: 'journalUnavailable' })
})
it('writes one message under another error name, or as a refusal, as separate failures', () => {
const logger = createStructuredAgentSessionLogger({ now: () => clock })
const renew = (error: Error): void =>
logger.warn('renewing a chat lease failed', { scope: 'lease-renewal', sessionId: 'a', error })
const unreadableRecord = (): Error => new Error('execution_owner_reconciling')
const reconciling = (): Error =>
agentSessionRefusalError('execution_owner_reconciling', { reason: 'hostReconciling' })
for (let tick = 0; tick < 3; tick += 1) {
renew(unreadableRecord())
renew(reconciling())
renew(new TypeError('execution_owner_reconciling'))
}
expect(written().map((span) => span.exit)).toEqual([
expect.objectContaining({ cause: expect.stringMatching(/^Error: /) }),
expect.objectContaining({ cause: expect.stringMatching(/^AgentSessionRefusalError: /) }),
expect.objectContaining({ cause: expect.stringMatching(/^TypeError: /) })
])
})
it('tracks a bounded number of repeating entries', () => {
const repeats = createStructuredAgentSessionLogRepeats({ now: () => clock, maxTracked: 2 })
expect(repeats.admit('a')).toBe(0)
expect(repeats.admit('b')).toBe(0)
expect(repeats.admit('a')).toBeNull()
// A third key evicts the one written longest ago, so `a` is written again.
expect(repeats.admit('c')).toBe(0)
expect(repeats.admit('a')).toBe(0)
expect(repeats.admit('c')).toBeNull()
})
})
describe('what an Error leaves in the trace file', () => {
it('carries its code and its causes, depth-capped, by name and message only', () => {
const denied = Object.assign(new Error('EACCES: permission denied, open'), {
code: 'EACCES',
path: '/private/payload'
})
const locked = Object.assign(new Error('database is locked'), {
code: 'ERR_SQLITE_ERROR',
errcode: 5
})
const chain = new Error('agent_session_store_corrupt', {
cause: new Error('read failed', {
cause: new Error('third', { cause: new Error('fourth', { cause: new Error('fifth') }) })
})
})
const logger = createStructuredAgentSessionLogger({ now: () => clock })
logger.warn('importing legacy records failed', { scope: 'legacy-record-import', error: chain })
logger.warn('opening the journal failed', {
scope: 'journal-database-open',
error: new Error('opening failed', { cause: denied })
})
logger.error('writing events failed', { scope: 'journal-event-sink', error: locked })
const [corrupt, open, sqlite] = written()
expect(corrupt?.attributes['errorCause']).toEqual([
'Error: read failed',
'Error: third',
'Error: fourth'
])
expect(open?.attributes['errorCause']).toEqual(['Error: EACCES: permission denied, open'])
expect(JSON.stringify(written())).not.toContain('/private/payload')
expect(sqlite?.attributes).toMatchObject({ errorCode: 'ERR_SQLITE_ERROR', errorErrcode: 5 })
expect(sqlite?.exit).toMatchObject({ cause: expect.stringContaining('database is locked') })
})
it('names the code of a cause whose message does not', () => {
const logger = createStructuredAgentSessionLogger({ now: () => clock })
logger.warn('x failed', {
scope: 'x',
error: new Error('x', { cause: Object.assign(new Error('no space'), { code: 'ENOSPC' }) })
})
expect(written()[0]?.attributes['errorCause']).toEqual(['Error: no space [ENOSPC]'])
})
})
@@ -0,0 +1,173 @@
// Where the structured chat host and its runtime report a failure they carry on past.
//
// One required dependency rather than a callback per failure: a host built without it does not
// compile, so no failure path can quietly drop what went wrong. The default writes each entry to
// the host's local trace file (the desktop's `<userData>/logs/main.trace.ndjson`, collected by the
// diagnostic bundle; orcad's `<data-root>/logs/orcad.trace.ndjson`) and to the console.
import { isAgentSessionRefusalError } from '../../../shared/agent-session-wire-refusals'
import { startSpan } from '../../observability/tracer'
import { createStructuredAgentSessionLogRepeats } from './structured-agent-session-log-repeats'
export type StructuredAgentSessionLogFields = {
/** The step that failed; a stable name a log search can find. */
readonly scope: string
readonly sessionId?: string
readonly error?: unknown
readonly [field: string]: unknown
}
export type StructuredAgentSessionLogger = {
warn: (message: string, fields: StructuredAgentSessionLogFields) => void
error: (message: string, fields: StructuredAgentSessionLogFields) => void
}
type LogLevel = keyof StructuredAgentSessionLogger
const guarded = new WeakSet<StructuredAgentSessionLogger>()
/** Reporting is bookkeeping: a logger that throws must never fail the operation it reports. */
export function neverThrowingStructuredAgentSessionLogger(
logger: StructuredAgentSessionLogger
): StructuredAgentSessionLogger {
if (guarded.has(logger)) {
return logger
}
const call =
(level: LogLevel) =>
(message: string, fields: StructuredAgentSessionLogFields): void => {
try {
logger[level](message, fields)
} catch (loggerError) {
try {
console.warn(`[agent-session] ${message}`, { ...fields, loggerError })
} catch {
// Nothing is left to report to.
}
}
}
const safe: StructuredAgentSessionLogger = { warn: call('warn'), error: call('error') }
guarded.add(safe)
return safe
}
/** A copy of `deps` whose logger cannot throw; every collaborator built from it inherits that. */
export function withNeverThrowingLogger<T extends { logger: StructuredAgentSessionLogger }>(
deps: T
): T {
return { ...deps, logger: neverThrowingStructuredAgentSessionLogger(deps.logger) }
}
/** The logger `resolve` returns at each call, for a collaborator built before the host's deps. */
export function deferredStructuredAgentSessionLogger(
resolve: () => StructuredAgentSessionLogger
): StructuredAgentSessionLogger {
return {
warn: (message, fields) => resolve().warn(message, fields),
error: (message, fields) => resolve().error(message, fields)
}
}
/** The production logger: a failed span in the local trace file, plus the console. A repeat of
* the same entry is written at most once per window, carrying how many were swallowed. */
export function createStructuredAgentSessionLogger(options?: {
now?: () => number
}): StructuredAgentSessionLogger {
const repeats = createStructuredAgentSessionLogRepeats({ now: options?.now })
const write =
(level: LogLevel) =>
(message: string, fields: StructuredAgentSessionLogFields): void => {
const { scope, error, ...rest } = fields
const attributes = { level, message, ...rest, ...errorAttributes(error) }
// Keyed on all the entry writes but its stack: only a failure writing the same entry repeats.
const suppressed = repeats.admit(
stableJson([
scope,
attributes,
error instanceof Error ? `${error.name}: ${error.message}` : null
])
)
if (suppressed === null) {
return
}
const span = startSpan(`agentSession.${scope}`, {
attributes: suppressed > 0 ? { ...attributes, suppressed } : attributes
})
span.fail(error instanceof Error ? error : message)
const print = level === 'error' ? console.error : console.warn
print(
`[agent-session] ${scope}: ${message}`,
suppressed > 0 ? { ...fields, suppressed } : fields
)
}
return neverThrowingStructuredAgentSessionLogger({ warn: write('warn'), error: write('error') })
}
const MAX_LOGGED_CAUSES = 3
/** The span's failure carries an Error's name, message and stack; its code and causes ride here.
* A value that is not an Error is written whole, as it would serialize. */
function errorAttributes(error: unknown): Record<string, unknown> {
if (error === undefined) {
return {}
}
if (!(error instanceof Error)) {
return { error }
}
const causes: string[] = []
let cause: unknown = error.cause
// Name and message only: a cause's own fields can carry what the failing step was handling.
while (cause !== undefined && causes.length < MAX_LOGGED_CAUSES) {
if (cause instanceof Error) {
causes.push(`${cause.name}: ${cause.message}${codeSuffix(cause)}`)
cause = cause.cause
} else {
causes.push(typeof cause === 'object' && cause !== null ? '[non-error cause]' : String(cause))
cause = undefined
}
}
const code = errorCode(error)
// node:sqlite's `code` is one generic value; `errcode` tells SQLITE_BUSY from SQLITE_FULL.
const sqliteCode = codeValue('errcode' in error ? error.errcode : undefined)
// A refusal's message is its bare code; its reason is wire-safe and tells refusals apart.
const details = isAgentSessionRefusalError(error) ? error.refusal.details : undefined
const reason = details && 'reason' in details ? details.reason : undefined
return {
...(code !== undefined ? { errorCode: code } : {}),
...(sqliteCode !== undefined ? { errorErrcode: sqliteCode } : {}),
...(typeof reason === 'string' ? { refusalReason: reason } : {}),
...(causes.length > 0 ? { errorCause: causes } : {})
}
}
/** The same value always renders the same: object keys sorted, and no throw on a cycle. */
function stableJson(value: unknown): string {
try {
return JSON.stringify(value, (_key, inner: unknown) => {
if (typeof inner === 'bigint') {
return String(inner)
}
if (typeof inner !== 'object' || inner === null || Array.isArray(inner)) {
return inner
}
return Object.fromEntries(
Object.entries(inner).sort(([left], [right]) => (left < right ? -1 : 1))
)
})
} catch {
return String(value)
}
}
function errorCode(error: Error): string | number | undefined {
return codeValue('code' in error ? error.code : undefined)
}
function codeValue(value: unknown): string | number | undefined {
return typeof value === 'string' || typeof value === 'number' ? value : undefined
}
function codeSuffix(error: Error): string {
const code = errorCode(error)
return code !== undefined && !error.message.includes(String(code)) ? ` [${code}]` : ''
}
@@ -37,6 +37,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
const PROVIDER_ROW = { provider: 'codex' as const, threadId: THREAD, turnId: 'turn-1' }
@@ -62,6 +63,7 @@ beforeEach(async () => {
dispatch = vi.fn(async () => ({ state: 'admitted' as const }))
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: {
acquire: async (input) => {
@@ -38,6 +38,7 @@ import type { MutationPlan } from './structured-agent-session-mutation-plans'
import { runSettledAgentSessionMutation } from './structured-agent-session-operation-settlement'
import { resolveAgentSessionReplayOutcome } from './structured-agent-session-replay-outcome'
import type { AgentSessionTurnContext } from './structured-agent-session-turns'
import type { StructuredAgentSessionLogger } from './structured-agent-session-logger'
// The code is shared with the client so a read that refuses this way can be told apart from a
// transcript that failed to load; the two must never drift apart.
@@ -61,6 +62,7 @@ export type AgentSessionMutationSessionPreparation =
export type AgentSessionMutationRequest<TValue> = {
store: AgentSessionRecordStore
adapter: StructuredAgentSessionAdapter
logger: StructuredAgentSessionLogger
callerKey: string
envelope: AgentSessionMutationEnvelope
plan: MutationPlan<TValue>
@@ -127,7 +129,7 @@ export async function admitAndRunAgentSessionMutation<TValue>(
admitted = await request.store.admitMutationOperation(operation)
} catch (error) {
if (plan.runsWithoutLedgerRow) {
admitted = admitWithoutLedgerRow(request.store, operation, error)
admitted = admitWithoutLedgerRow(request, operation, error)
ledgerRowWritten = false
} else if (
isAgentSessionRefusalError(error) ||
@@ -198,13 +200,14 @@ export async function admitAndRunAgentSessionMutation<TValue>(
/** The committed ledger's admission, placing nothing: a failed commit left memory as it was. */
function admitWithoutLedgerRow(
store: AgentSessionRecordStore,
{ store, logger }: Pick<AgentSessionMutationRequest<unknown>, 'store' | 'logger'>,
operation: AgentSessionMutationOperationAdmission,
error: unknown
): AgentSessionMutationOperationDecision {
console.warn("[agent-session] Stop's ledger row skipped:", {
logger.warn("writing Stop's ledger row failed; Stop runs without it", {
scope: 'stop-ledger-row',
sessionId: operation.envelope.sessionId,
error: error instanceof Error ? error.message : String(error)
error
})
const evaluated = store.evaluateMutationOperation(operation)
if (!evaluated) {
@@ -231,6 +234,7 @@ function turnContext<TValue>(
journal,
fence,
adapter: request.adapter,
logger: request.logger,
...(persistedOptions ? { persistedOptions } : {}),
persistOptions: (options) =>
request.store
@@ -55,6 +55,7 @@ export function mutateStructuredAgentSession<TValue>(
admitAndRunAgentSessionMutation({
store: context.deps.store,
adapter: context.deps.adapter,
logger: context.deps.logger,
callerKey: caller.callerKey,
envelope,
plan,
@@ -17,9 +17,11 @@ import {
} from './structured-agent-session-host-test-data'
import type { AgentSessionTurnContext } from './structured-agent-session-turns'
import { sendPlan } from './structured-agent-session-mutation-plans'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
async function context(): Promise<AgentSessionTurnContext> {
return {
logger: createStructuredAgentSessionLogger(),
sessionId: SESSION,
journal: await journals.open({
identity: {
@@ -52,14 +52,23 @@ export async function runSettledAgentSessionMutation<TValue>(input: {
if (input.plan.markUnknownBeforeRun && error instanceof AgentSessionPreDispatchError) {
throw error
}
const { logger, sessionId } = input.context
try {
await settle({ status: 'unknown' })
} catch {
// Bookkeeping must not replace the operation's proof of whether dispatch began.
console.warn('[structured-agent-session] operation uncertainty persistence failed')
logger.warn('recording an operation as unknown failed', {
scope: 'operation-unknown-settlement',
sessionId,
operationId: input.envelope.clientOperationId
})
}
if (outcome && !outcome.ok) {
console.warn('[structured-agent-session] refused operation settlement failed')
logger.warn('recording a refused operation failed', {
scope: 'operation-refused-settlement',
sessionId,
operationId: input.envelope.clientOperationId
})
return outcome
}
throw error
@@ -6,14 +6,16 @@ import type { AgentSessionRecord } from '../../../shared/agent-session-record'
import type { StructuredAgentSessionAttachContext } from './structured-agent-session-attach-context'
import { ensureStructuredAgentSessionAgentForOperation } from './structured-agent-session-agent-start'
import { recordStructuredAgentSessionOptionIntent } from './structured-agent-session-options-read'
import { recordingStructuredAgentSessionLogger } from './structured-agent-session-logger-test-support'
const SESSION = 'session-1'
describe('an operation whose agent start throws', () => {
it('is refused as a failed restart, with the error only in the log', async () => {
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {})
const log = recordingStructuredAgentSessionLogger()
const cause = new Error('EACCES: permission denied, open /Users/me/.orca/leases.json')
const context = {
deps: { logger: log.logger },
sessions: new Map(),
reconcileLeases: () => Promise.reject(cause)
}
@@ -31,8 +33,9 @@ describe('an operation whose agent start throws', () => {
message: "The agent couldn't restart."
}
})
expect(warn).toHaveBeenCalledWith(expect.stringContaining('starting the agent'), cause)
warn.mockRestore()
expect(log.entries.map((entry) => entry.fields)).toEqual([
{ scope: 'operation-agent-start', sessionId: SESSION, error: cause }
])
})
})
@@ -21,6 +21,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
const DEFAULT_MODEL = 'gpt-default'
@@ -124,6 +125,7 @@ beforeEach(async () => {
store = await openTestAgentSessionRecordStore(root)
router = adapter()
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: router,
journalDatabase: openTestJournalHostDatabase(root),
@@ -117,7 +117,7 @@ export async function readStructuredAgentSessionOptions(
const { adapter, store } = context.deps
const live = await context.serialize(sessionId, async () => {
const session = await context.openConversation(sessionId).catch((error: unknown) => {
throw journalOpenReadRefusal(error)
throw journalOpenReadRefusal(error, context.deps.logger, sessionId)
})
const child = session?.child
if (!child) {
@@ -32,6 +32,7 @@ import {
resetHostTestOperationIds
} from './structured-agent-session-host-test-data'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const CALLER = { callerKey: 'client-1' }
const SWEEP_MS = 5
@@ -93,6 +94,7 @@ beforeEach(async () => {
})
store = await openTestAgentSessionRecordStore(root)
host = new StructuredAgentSessionHost({
logger: createStructuredAgentSessionLogger(),
store,
adapter: Object.assign(adapter, { supportsCreate: () => true }),
journalDatabase: openTestJournalHostDatabase(root),
@@ -17,6 +17,7 @@ import {
import { openTestAttachConversation } from './structured-agent-session-attach-test-conversation'
import { performAttach } from './structured-agent-session-attach-flow'
import { openTestJournalHostDatabase } from '../agent-session-journal/journal-host-database-test-support'
import { createStructuredAgentSessionLogger } from './structured-agent-session-logger'
const NOW = 1_800_000_000_000
const SESSION = 'session-alpha'
@@ -82,6 +83,7 @@ async function firstAnswerAndReplay(thrown: AgentSessionPreSpawnError) {
const input = {
store,
adapter,
logger: createStructuredAgentSessionLogger(),
journalDatabase: openTestJournalHostDatabase(root),
openConversation: openTestAttachConversation(openTestJournalHostDatabase(root)),
authority: {
@@ -157,8 +159,11 @@ describe('a create that fails before any process spawns', () => {
}
// What failed is kept for the log.
expect(warn).toHaveBeenCalledWith(
'[agent-session] provider start failed:',
expect.objectContaining({ message: raw })
'[agent-session] provider-start: starting the provider failed',
expect.objectContaining({
scope: 'provider-start',
error: expect.objectContaining({ message: raw })
})
)
})

Some files were not shown because too many files have changed in this diff Show More