From a90aacd83eaa132bed336d4e924d3355b3c3d117 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Tue, 15 Sep 2026 13:29:31 -0700 Subject: [PATCH] fix(agent-status): close canonical store race windows --- ...-status-child-work-admission-operations.ts | 28 ++--- .../agent-status-child-work-admission.test.ts | 40 +++++- .../agent-status-child-work-admission.ts | 2 +- src/shared/agent-status-child-work-stop.ts | 22 ++++ src/shared/agent-status-store-bounds.test.ts | 115 +++++++++++++++--- src/shared/agent-status-store-codec.ts | 6 +- src/shared/agent-status-store-mutation.ts | 37 ++++-- .../agent-status-store-snapshot-budget.ts | 68 +++++++++++ src/shared/agent-status-store-state.ts | 6 +- 9 files changed, 264 insertions(+), 60 deletions(-) create mode 100644 src/shared/agent-status-child-work-stop.ts create mode 100644 src/shared/agent-status-store-snapshot-budget.ts diff --git a/src/shared/agent-status-child-work-admission-operations.ts b/src/shared/agent-status-child-work-admission-operations.ts index beee1f717ba..57b35736df4 100644 --- a/src/shared/agent-status-child-work-admission-operations.ts +++ b/src/shared/agent-status-child-work-admission-operations.ts @@ -2,8 +2,7 @@ import { serializeAgentChildWorkAliasKey } from './agent-status-child-work-alias import { AGENT_CHILD_WORK_INVOCATION_HISTORY_MAX, agentChildWorkFencesEqual, - type AgentChildWorkId, - type AgentChildWorkRecord + type AgentChildWorkId } from './agent-status-child-work' import { agentChildWorkAliasesForChild, @@ -21,8 +20,7 @@ import type { AgentChildWorkAdoptRequest, AgentChildWorkAnnounceRequest, AgentChildWorkReparentRequest, - AgentChildWorkResumeRequest, - AgentChildWorkStopRequest + AgentChildWorkResumeRequest } from './agent-status-child-work-admission' import { parseAgentChildWorkInput, @@ -85,6 +83,9 @@ export function announceAgentChildWork( if (!child) { return rejectAgentChildWorkAdmission('ambiguous') } + if (!agentChildWorkFencesEqual(child.invocation, fence)) { + return rejectAgentChildWorkAdmission('stale-invocation') + } const aliases = buildAgentChildWorkAliases( parent, request.provider, @@ -219,11 +220,13 @@ export function resumeAgentChildWork( } ].slice(-AGENT_CHILD_WORK_INVOCATION_HISTORY_MAX) const retainedFences = [nextFence, ...previousInvocations.map((entry) => entry.fence)] + const nextAliasKeys = new Set(aliases.map(serializeAgentChildWorkAliasKey)) const removeAliases = agentChildWorkAliasesForChild(store, child.childWorkId) .filter( (alias) => !retainedFences.some((fence) => agentChildWorkFencesEqual(alias.fence, fence)) ) .map(serializeAgentChildWorkAliasKey) + .filter((key) => !nextAliasKeys.has(key)) const resumed = buildAgentChildWork( request, child.childWorkId, @@ -283,20 +286,3 @@ export function reparentAgentChildWork( ) : rejectAgentChildWorkAdmission('invalid') } - -export function authorizeAgentChildWorkStop( - store: AgentStatusStore, - request: AgentChildWorkStopRequest -): AgentChildWorkRecord | null { - const child = findAgentChildWork(store, request.childWorkId) - if ( - !child || - !agentStatusSubjectsEqual(child.parent, request.parent) || - !agentChildWorkFencesEqual(child.invocation, request.expectedFence) || - child.membership !== 'live' || - child.stoppable !== true - ) { - return null - } - return child -} diff --git a/src/shared/agent-status-child-work-admission.test.ts b/src/shared/agent-status-child-work-admission.test.ts index 46f4006b5b4..091405d9aec 100644 --- a/src/shared/agent-status-child-work-admission.test.ts +++ b/src/shared/agent-status-child-work-admission.test.ts @@ -189,14 +189,15 @@ describe('agent child-work admission', () => { generation <= AGENT_CHILD_WORK_INVOCATION_HISTORY_MAX + 2; generation += 1 ) { + const reusesOldestAlias = generation === AGENT_CHILD_WORK_INVOCATION_HISTORY_MAX + 2 expect( admission.resume({ ...announce(parent, { aliases: [ { - segmentId: `segment-${generation}`, + segmentId: reusesOldestAlias ? 'segment-1' : `segment-${generation}`, aliasKind: 'task_id', - alias: `task-${generation}` + alias: reusesOldestAlias ? 'task-1' : `task-${generation}` } ], observedAt: 10 + generation @@ -213,7 +214,40 @@ describe('agent child-work admission', () => { const aliases = store.getAliasesForChild('child-1') expect(aliases).toHaveLength(AGENT_CHILD_WORK_INVOCATION_HISTORY_MAX + 1) - expect(aliases.some((entry) => entry.alias === 'task-1')).toBe(false) + expect(aliases.some((entry) => entry.fence.generation === 1)).toBe(false) + expect(aliases.find((entry) => entry.alias === 'task-1')?.fence.generation).toBe( + AGENT_CHILD_WORK_INVOCATION_HISTORY_MAX + 2 + ) + }) + + it('rejects delayed observations resolved through a previous invocation alias', () => { + const parent = subject() + const { admission, store } = setup([parent]) + expect(admission.announce(announce(parent)).accepted).toBe(true) + expect( + admission.resume({ + ...announce(parent, { + aliases: [{ segmentId: 'segment-2', aliasKind: 'task_id', alias: 'task-2' }], + observedAt: 20 + }), + childWorkId: 'child-1', + expectedFence: { invocationId: 'invocation-1', generation: 1 }, + nextFence: { invocationId: 'invocation-2', generation: 2 } + }) + ).toMatchObject({ accepted: true, childWorkId: 'child-1' }) + const before = store.getSnapshot() + + expect( + admission.announce( + announce(parent, { + state: 'done', + membership: 'settled', + outcome: 'succeeded', + observedAt: 30 + }) + ) + ).toEqual({ accepted: false, reason: 'stale-invocation' }) + expect(store.getSnapshot()).toEqual(before) }) it('mints a distinct id for proven reuse and rejects delayed predecessor updates and stops', () => { diff --git a/src/shared/agent-status-child-work-admission.ts b/src/shared/agent-status-child-work-admission.ts index 2cf29699580..554191db895 100644 --- a/src/shared/agent-status-child-work-admission.ts +++ b/src/shared/agent-status-child-work-admission.ts @@ -13,10 +13,10 @@ import type { import { adoptAgentChildWork, announceAgentChildWork, - authorizeAgentChildWorkStop, reparentAgentChildWork, resumeAgentChildWork } from './agent-status-child-work-admission-operations' +import { authorizeAgentChildWorkStop } from './agent-status-child-work-stop' import type { AgentStatusStore } from './agent-status-store' import type { AgentStatusSubject } from './agent-status-subject' diff --git a/src/shared/agent-status-child-work-stop.ts b/src/shared/agent-status-child-work-stop.ts new file mode 100644 index 00000000000..6a5ed05be9e --- /dev/null +++ b/src/shared/agent-status-child-work-stop.ts @@ -0,0 +1,22 @@ +import { agentChildWorkFencesEqual, type AgentChildWorkRecord } from './agent-status-child-work' +import type { AgentChildWorkStopRequest } from './agent-status-child-work-admission' +import { findAgentChildWork } from './agent-status-child-work-admission-core' +import type { AgentStatusStore } from './agent-status-store' +import { agentStatusSubjectsEqual } from './agent-status-subject' + +export function authorizeAgentChildWorkStop( + store: AgentStatusStore, + request: AgentChildWorkStopRequest +): AgentChildWorkRecord | null { + const child = findAgentChildWork(store, request.childWorkId) + if ( + !child || + !agentStatusSubjectsEqual(child.parent, request.parent) || + !agentChildWorkFencesEqual(child.invocation, request.expectedFence) || + child.membership !== 'live' || + child.stoppable !== true + ) { + return null + } + return child +} diff --git a/src/shared/agent-status-store-bounds.test.ts b/src/shared/agent-status-store-bounds.test.ts index 4fd89d31010..45ba1de57d7 100644 --- a/src/shared/agent-status-store-bounds.test.ts +++ b/src/shared/agent-status-store-bounds.test.ts @@ -1,8 +1,11 @@ -import { describe, expect, it } from 'vitest' +import { describe, expect, it, vi } from 'vitest' +import type { AgentChildWorkAliasInput } from './agent-status-child-work-alias' import type { AgentChildWorkInput } from './agent-status-child-work' import { serializeAgentStatusRunAliasIndex } from './agent-status-run-alias-index' import { createAgentStatusStore } from './agent-status-store' import { AGENT_STATUS_STORE_LIMITS } from './agent-status-store-contract' +import { applyAgentStatusStoreMutation } from './agent-status-store-mutation' +import { agentStatusStoreStateFromSnapshot } from './agent-status-store-state' import { makePtyRunAgentStatusSubject, makeStructuredAgentStatusSubject, @@ -16,24 +19,31 @@ const scope: AgentStatusExecutionScope = { workspaceKind: 'git-worktree' } +function child( + parent: ReturnType, + childIndex: number, + description = 'x'.repeat(8_000) +) { + return { + childWorkId: `child-${childIndex}`, + parent, + provider: 'claude', + kind: 'agent', + state: 'working', + membership: 'live', + firstObservedAt: 1, + observedAt: 1, + stoppable: true, + invocation: { invocationId: `invocation-${childIndex}`, generation: 1 }, + provenance: { source: 'structured-session', producerId: 'journal-1' }, + description + } satisfies AgentChildWorkInput +} + function children(parent: ReturnType, offset: number) { - return Array.from({ length: AGENT_STATUS_STORE_LIMITS.mutationEntries / 2 }, (_, index) => { - const childIndex = offset + index - return { - childWorkId: `child-${childIndex}`, - parent, - provider: 'claude', - kind: 'agent', - state: 'working', - membership: 'live', - firstObservedAt: 1, - observedAt: 1, - stoppable: true, - invocation: { invocationId: `invocation-${childIndex}`, generation: 1 }, - provenance: { source: 'structured-session', producerId: 'journal-1' }, - description: 'x'.repeat(8_000) - } satisfies AgentChildWorkInput - }) + return Array.from({ length: AGENT_STATUS_STORE_LIMITS.mutationEntries / 2 }, (_, index) => + child(parent, offset + index) + ) } describe('AgentStatusStore bounds', () => { @@ -77,4 +87,73 @@ describe('AgentStatusStore bounds', () => { expect(() => serializeAgentStatusRunAliasIndex(store.getRunAliasIndex())).not.toThrow() }) + + it('reuses immutable record byte measurements when one child changes', () => { + const parent = makeStructuredAgentStatusSubject(scope, 'session-1') + const store = createAgentStatusStore({ epoch: 'epoch-a', mode: 'authority' }) + const first = child(parent, 0, 'first-description') + const unchanged = child(parent, 1, 'unchanged-description') + expect( + store.applyMutation({ parent: { subject: parent }, children: [first, unchanged] }) + ).not.toBeNull() + const stringify = vi.spyOn(JSON, 'stringify') + + try { + expect( + store.applyMutation({ children: [{ ...first, observedAt: first.observedAt + 1 }] }) + ).not.toBeNull() + const serializedValues = stringify.mock.results + .map((result) => result.value) + .filter((value): value is string => typeof value === 'string') + expect(serializedValues.some((value) => value.includes('unchanged-description'))).toBe(false) + } finally { + stringify.mockRestore() + } + }) + + it('removes aliases for a child batch with one alias-map pass', () => { + const parent = makeStructuredAgentStatusSubject(scope, 'session-1') + const store = createAgentStatusStore({ epoch: 'epoch-a', mode: 'authority' }) + const childRecords = Array.from({ length: 4 }, (_, index) => child(parent, index, 'brief')) + const aliases: AgentChildWorkAliasInput[] = childRecords.map((record, index) => ({ + parent, + provider: record.provider, + segmentId: `segment-${index}`, + kind: record.kind, + aliasKind: 'task_id', + alias: `task-${index}`, + childWorkId: record.childWorkId, + fence: record.invocation + })) + expect( + store.applyMutation({ parent: { subject: parent }, children: childRecords, aliases }) + ).not.toBeNull() + const state = agentStatusStoreStateFromSnapshot(store.getSnapshot(), 'epoch-a') + expect(state).not.toBeNull() + if (!state) { + throw new Error('Expected a valid store state') + } + let childWorkIdReads = 0 + for (const [key, record] of state.aliases) { + const measured = { ...record } + Object.defineProperty(measured, 'childWorkId', { + enumerable: true, + get: () => { + childWorkIdReads += 1 + return record.childWorkId + } + }) + state.aliases.set(key, measured) + } + childWorkIdReads = 0 + + const next = applyAgentStatusStoreMutation( + state, + { removeChildren: childRecords.map((record) => record.childWorkId) }, + state.revision + 1 + ) + + expect(next).not.toBeNull() + expect(childWorkIdReads).toBe(aliases.length) + }) }) diff --git a/src/shared/agent-status-store-codec.ts b/src/shared/agent-status-store-codec.ts index fb385242151..5a0d6fc2878 100644 --- a/src/shared/agent-status-store-codec.ts +++ b/src/shared/agent-status-store-codec.ts @@ -50,7 +50,7 @@ function isRevision(value: unknown): value is number { return typeof value === 'number' && Number.isSafeInteger(value) && value >= 0 } -export function isAgentStatusStoreValueWithinSerializedLimit(value: unknown): boolean { +function isWithinSerializedLimit(value: unknown): boolean { try { return !measureUtf8ByteLength(JSON.stringify(value), { stopAfterBytes: AGENT_STATUS_STORE_LIMITS.serializedBytes @@ -153,7 +153,7 @@ export function parseAgentStatusStoreMutation(value: unknown): AgentStatusStoreM 'tombstones' ] if ( - !isAgentStatusStoreValueWithinSerializedLimit(value) || + !isWithinSerializedLimit(value) || !isRecord(value) || !hasOnlyKeys(value, [], optionalKeys) || Object.keys(value).length === 0 @@ -239,7 +239,7 @@ export function parseAgentStatusStoreSnapshot(value: unknown): AgentStatusStoreS 'tombstones' ] if ( - !isAgentStatusStoreValueWithinSerializedLimit(value) || + !isWithinSerializedLimit(value) || !isRecord(value) || !hasOnlyKeys(value, keys) || value.version !== AGENT_STATUS_STORE_SNAPSHOT_VERSION || diff --git a/src/shared/agent-status-store-mutation.ts b/src/shared/agent-status-store-mutation.ts index 7302a1e5dac..b8338523e64 100644 --- a/src/shared/agent-status-store-mutation.ts +++ b/src/shared/agent-status-store-mutation.ts @@ -63,11 +63,24 @@ function removeFact(state: AgentStatusStoreState, key: string, revision: number) addTombstone(state, 'fact', key, revision) } -function removeChild(state: AgentStatusStoreState, childWorkId: string, revision: number): void { +function removeChild( + state: AgentStatusStoreState, + childWorkId: string, + revision: number, + removedChildWorkIds: Set +): void { state.children.delete(childWorkId) addTombstone(state, 'child', childWorkId, revision) + removedChildWorkIds.add(childWorkId) +} + +function removeAliasesForChildren( + state: AgentStatusStoreState, + removedChildWorkIds: ReadonlySet, + revision: number +): void { for (const [key, alias] of state.aliases) { - if (alias.childWorkId === childWorkId) { + if (removedChildWorkIds.has(alias.childWorkId)) { removeAlias(state, key, revision) } } @@ -76,14 +89,15 @@ function removeChild(state: AgentStatusStoreState, childWorkId: string, revision function removeParent( state: AgentStatusStoreState, subject: AgentStatusSubject, - revision: number + revision: number, + removedChildWorkIds: Set ): void { const key = serializeAgentStatusSubject(subject) state.parents.delete(key) addTombstone(state, 'parent', key, revision) for (const child of state.children.values()) { if (agentStatusSubjectsEqual(child.parent, subject)) { - removeChild(state, child.childWorkId, revision) + removeChild(state, child.childWorkId, revision, removedChildWorkIds) } } for (const [factMapKey, fact] of state.facts) { @@ -96,18 +110,19 @@ function removeParent( function applyExplicitTombstone( state: AgentStatusStoreState, tombstone: { entity: AgentStatusTombstoneEntity; key: string }, - revision: number + revision: number, + removedChildWorkIds: Set ): boolean { if (tombstone.entity === 'parent') { const subject = deserializeAgentStatusSubject(tombstone.key) if (!subject) { return false } - removeParent(state, subject, revision) + removeParent(state, subject, revision, removedChildWorkIds) return true } if (tombstone.entity === 'child') { - removeChild(state, tombstone.key, revision) + removeChild(state, tombstone.key, revision, removedChildWorkIds) } else if (tombstone.entity === 'alias') { if (!deserializeAgentChildWorkAliasKey(tombstone.key)) { return false @@ -209,11 +224,12 @@ export function applyAgentStatusStoreMutation( ): AgentStatusStoreState | null { const next = cloneAgentStatusStoreState(current) next.revision = revision + const removedChildWorkIds = new Set() if (mutation.removeParent) { - removeParent(next, mutation.removeParent, revision) + removeParent(next, mutation.removeParent, revision, removedChildWorkIds) } for (const childWorkId of mutation.removeChildren ?? []) { - removeChild(next, childWorkId, revision) + removeChild(next, childWorkId, revision, removedChildWorkIds) } for (const key of mutation.removeAliases ?? []) { if (!deserializeAgentChildWorkAliasKey(key)) { @@ -225,10 +241,11 @@ export function applyAgentStatusStoreMutation( removeFact(next, agentStatusFactMapKey(identity), revision) } for (const tombstone of mutation.tombstones ?? []) { - if (!applyExplicitTombstone(next, tombstone, revision)) { + if (!applyExplicitTombstone(next, tombstone, revision, removedChildWorkIds)) { return null } } + removeAliasesForChildren(next, removedChildWorkIds, revision) if (mutation.parent && !upsertParent(next, mutation.parent, revision)) { return null } diff --git a/src/shared/agent-status-store-snapshot-budget.ts b/src/shared/agent-status-store-snapshot-budget.ts new file mode 100644 index 00000000000..6ae4dc8f040 --- /dev/null +++ b/src/shared/agent-status-store-snapshot-budget.ts @@ -0,0 +1,68 @@ +import { + AGENT_STATUS_STORE_LIMITS, + AGENT_STATUS_STORE_SNAPSHOT_VERSION +} from './agent-status-store-contract' +import type { AgentChildWorkAliasRecord } from './agent-status-child-work-alias' +import type { AgentChildWorkRecord } from './agent-status-child-work' +import type { AgentStatusFactRecord } from './agent-status-store-contract' +import type { AgentStatusParentRecord } from './agent-status-store-parent' +import type { AgentStatusStoreState } from './agent-status-store-state' +import type { AgentStatusTombstoneRecord } from './agent-status-store-contract' +import { getUtf8ByteLength } from './utf8-byte-limits' + +type AgentStatusSnapshotRecord = + | AgentStatusParentRecord + | AgentChildWorkRecord + | AgentChildWorkAliasRecord + | AgentStatusFactRecord + | AgentStatusTombstoneRecord + +const snapshotRecordByteLengths = new WeakMap() + +function getSnapshotRecordByteLength(record: AgentStatusSnapshotRecord): number { + const cached = snapshotRecordByteLengths.get(record) + if (cached !== undefined) { + return cached + } + const byteLength = getUtf8ByteLength(JSON.stringify(record)) + snapshotRecordByteLengths.set(record, byteLength) + return byteLength +} + +function getSnapshotCollectionEntryBytes(records: Iterable): number { + let byteLength = 0 + let count = 0 + for (const record of records) { + byteLength += getSnapshotRecordByteLength(record) + count += 1 + } + return byteLength + Math.max(0, count - 1) +} + +export function isAgentStatusStoreSnapshotWithinByteLimit(state: AgentStatusStoreState): boolean { + const emptySnapshot = { + version: AGENT_STATUS_STORE_SNAPSHOT_VERSION, + epoch: state.epoch, + revision: state.revision, + parents: [], + children: [], + aliases: [], + facts: [], + tombstones: [] + } + let byteLength = getUtf8ByteLength(JSON.stringify(emptySnapshot)) + const collections: Iterable[] = [ + state.parents.values(), + state.children.values(), + state.aliases.values(), + state.facts.values(), + state.tombstones.values() + ] + for (const records of collections) { + byteLength += getSnapshotCollectionEntryBytes(records) + if (byteLength > AGENT_STATUS_STORE_LIMITS.serializedBytes) { + return false + } + } + return true +} diff --git a/src/shared/agent-status-store-state.ts b/src/shared/agent-status-store-state.ts index 4ce51b5ae89..42749f270ca 100644 --- a/src/shared/agent-status-store-state.ts +++ b/src/shared/agent-status-store-state.ts @@ -21,7 +21,6 @@ import { type AgentStatusTombstoneRecord } from './agent-status-store-contract' import { - isAgentStatusStoreValueWithinSerializedLimit, parseAgentStatusStoreSnapshot, parseAgentStatusTombstoneRecord } from './agent-status-store-codec' @@ -34,6 +33,7 @@ import { parseAgentStatusParentRecord, type AgentStatusParentRecord } from './agent-status-store-parent' +import { isAgentStatusStoreSnapshotWithinByteLimit } from './agent-status-store-snapshot-budget' import { deserializeAgentStatusSubject, serializeAgentStatusSubject } from './agent-status-subject' export type AgentStatusStoreState = { @@ -182,9 +182,7 @@ export function validateAgentStatusStoreState(state: AgentStatusStoreState): boo return false } } - return isAgentStatusStoreValueWithinSerializedLimit( - snapshotCandidateFromAgentStatusStoreState(state) - ) + return isAgentStatusStoreSnapshotWithinByteLimit(state) } export function snapshotFromAgentStatusStoreState(