fix(agent-status): close canonical store race windows

This commit is contained in:
Brennan Benson
2026-09-15 13:34:47 -07:00
parent 7178f20565
commit a90aacd83e
9 changed files with 264 additions and 60 deletions
@@ -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
}
@@ -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', () => {
@@ -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'
@@ -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
}
+97 -18
View File
@@ -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<typeof makeStructuredAgentStatusSubject>,
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<typeof makeStructuredAgentStatusSubject>, 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)
})
})
+3 -3
View File
@@ -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 ||
+27 -10
View File
@@ -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<string>
): void {
state.children.delete(childWorkId)
addTombstone(state, 'child', childWorkId, revision)
removedChildWorkIds.add(childWorkId)
}
function removeAliasesForChildren(
state: AgentStatusStoreState,
removedChildWorkIds: ReadonlySet<string>,
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<string>
): 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<string>
): 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<string>()
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
}
@@ -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<AgentStatusSnapshotRecord, number>()
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<AgentStatusSnapshotRecord>): 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<AgentStatusSnapshotRecord>[] = [
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
}
+2 -4
View File
@@ -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(