-
Notifications
You must be signed in to change notification settings - Fork 5.5k
feat(agent-status): host-owned child producer for structured sessions #21279
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Closed
Closed
Changes from all commits
Commits
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
73 changes: 73 additions & 0 deletions
73
src/main/agent-hooks/server/server-ingest-structured-children.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,73 @@ | ||
| import { randomUUID } from 'node:crypto' | ||
|
|
||
| import { createAgentChildWorkAdmission } from '../../../shared/agent-status-child-work-admission' | ||
| import { | ||
| reconcileStructuredChildWork, | ||
| type StructuredChildWorkReconcileOutcome | ||
| } from '../../../shared/agent-status-child-work-reconciliation' | ||
| import { | ||
| STRUCTURED_SUPPORTS_STOP_ALL_FACT, | ||
| STRUCTURED_SUPPORTS_TASK_STOP_FACT | ||
| } from '../../../shared/agent-status-child-work-structured-egress' | ||
| import type { StructuredChildWorkEvidence } from '../../../shared/agent-status-child-work-structured-evidence' | ||
| import { | ||
| parseAgentStatusSubject, | ||
| type AgentStatusStructuredSessionSubject | ||
| } from '../../../shared/agent-status-subject' | ||
| import { AgentHookServerIngestStructured } from './server-ingest-structured' | ||
|
|
||
| export abstract class AgentHookServerIngestStructuredChildren extends AgentHookServerIngestStructured { | ||
| /** | ||
| * Admit one structured session's full child-work roster. The parent publication owns | ||
| * the subject and lands first; this refuses to act on a subject the store does not | ||
| * already hold, so a child can never conjure a parent row. | ||
| */ | ||
| ingestStructuredChildWork( | ||
| subject: AgentStatusStructuredSessionSubject, | ||
| evidence: StructuredChildWorkEvidence, | ||
| provider: string | ||
| ): StructuredChildWorkReconcileOutcome | null { | ||
| const parent = parseAgentStatusSubject(subject) | ||
| if (!parent || parent.kind !== 'structured-session') { | ||
| throw new Error('Structured child work requires its exact owner subject') | ||
| } | ||
| const store = this.canonicalStatusStore | ||
| if (!store.getParent(parent)) { | ||
| return null | ||
| } | ||
| // Provider stop capability is a fact about the session, not about any one child. | ||
| store.applyMutation({ | ||
| facts: [ | ||
| { | ||
| subject: parent, | ||
| key: STRUCTURED_SUPPORTS_TASK_STOP_FACT, | ||
| value: evidence.supportsTaskStop | ||
| }, | ||
| { subject: parent, key: STRUCTURED_SUPPORTS_STOP_ALL_FACT, value: evidence.supportsStopAll } | ||
| ] | ||
| }) | ||
| const outcome = reconcileStructuredChildWork({ | ||
| store, | ||
| admission: createAgentChildWorkAdmission(store, { mintChildWorkId: () => randomUUID() }), | ||
| parent, | ||
| provider, | ||
| evidence, | ||
| observedAt: Date.now() | ||
| }) | ||
| if (outcome.rejected.length > 0) { | ||
| console.warn( | ||
| '[agent-status-child-work] refused structured child admissions', | ||
| outcome.rejected.map( | ||
| (entry) => `${entry.providerTaskId || entry.childWorkId}:${entry.reason}` | ||
| ) | ||
| ) | ||
| } | ||
| return outcome | ||
| } | ||
|
|
||
| /** Every canonical child this host holds for one structured session. */ | ||
| getStructuredChildWork(subject: AgentStatusStructuredSessionSubject) { | ||
| const parent = parseAgentStatusSubject(subject) | ||
| return parent ? this.canonicalStatusStore.getChildren(parent) : [] | ||
| } | ||
| } | ||
188 changes: 188 additions & 0 deletions
188
src/main/agent-hooks/structured-child-work-bridge-parity.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,188 @@ | ||
| // Parity gate: the host producer's legacy subagent egress against the renderer bridge's | ||
| // own projection code, not a restatement of it. `subagentSnapshotsFromTasks` is imported | ||
| // from the bridge's module, so a change there reddens this file. | ||
| // | ||
| // The bridge is still the only live producer of native-chat child rows. This gate is what | ||
| // the atomic switch needs before that writer can be turned off. | ||
|
|
||
| import { describe, expect, it } from 'vitest' | ||
| import type { | ||
| AgentSessionBackgroundTask, | ||
| AgentSessionBackgroundTaskState | ||
| } from '../../shared/agent-session-background-task-wire' | ||
| import { | ||
| AGENT_STATUS_MAX_SUBAGENTS, | ||
| type AgentSubagentSnapshot | ||
| } from '../../shared/agent-status-types' | ||
| import { createStructuredChildWorkProducerFixture } from '../../shared/agent-status-child-work-structured-producer.test-fixture' | ||
| import { subagentSnapshotsFromTasks } from '../../shared/structured-session-legacy-subagent-projection' | ||
|
|
||
| function roster(tasks: AgentSessionBackgroundTask[]): AgentSessionBackgroundTaskState { | ||
| return { state: 'monitoring', tasks, supportsTaskStop: true } | ||
| } | ||
|
|
||
| function agent( | ||
| id: string, | ||
| overrides: Partial<AgentSessionBackgroundTask> = {} | ||
| ): AgentSessionBackgroundTask { | ||
| return { id, kind: 'agent', startedAt: 5, ...overrides } | ||
| } | ||
|
|
||
| /** Compare the two projections field by field, holding `startedAt` aside: the bridge | ||
| * publishes the provider's own stamp (or `0`), the host publishes its `firstObservedAt`, | ||
| * and that difference is asserted on its own below. */ | ||
| function withoutStartedAt( | ||
| snapshots: AgentSubagentSnapshot[] | undefined | ||
| ): Omit<AgentSubagentSnapshot, 'startedAt'>[] | undefined { | ||
| return snapshots?.map(({ startedAt: _startedAt, ...rest }) => rest) | ||
| } | ||
|
|
||
| function compare(tasks: AgentSessionBackgroundTask[]): { | ||
| bridge: AgentSubagentSnapshot[] | undefined | ||
| host: AgentSubagentSnapshot[] | undefined | ||
| } { | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| producer.publish(roster(tasks)) | ||
| return { bridge: subagentSnapshotsFromTasks(tasks), host: producer.subagents() } | ||
| } | ||
|
|
||
| function expectParity(tasks: AgentSessionBackgroundTask[]): AgentSubagentSnapshot[] | undefined { | ||
| const { bridge, host } = compare(tasks) | ||
| expect(withoutStartedAt(host)).toEqual(withoutStartedAt(bridge)) | ||
| return host | ||
| } | ||
|
|
||
| describe('legacy subagent egress matches the renderer bridge', () => { | ||
| it('admits only agent-kind rows while the other kinds stay canonical', () => { | ||
| const tasks: AgentSessionBackgroundTask[] = [ | ||
| agent('task-agent', { description: 'review' }), | ||
| { id: 'task-shell', kind: 'command', startedAt: 6 }, | ||
| { id: 'task-watch', kind: 'monitor', startedAt: 7 }, | ||
| { id: 'task-flow', kind: 'workflow', startedAt: 8 }, | ||
| { id: 'task-other', kind: 'unknown', startedAt: 9 } | ||
| ] | ||
| const host = expectParity(tasks) | ||
| expect(host?.map((row) => row.id)).toEqual(['task-agent']) | ||
|
|
||
| // Positive control: the rows the closed vocabulary drops are still in the collection, | ||
| // so an empty subagent projection is a statement about the projection, not the store. | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| producer.publish(roster(tasks)) | ||
| expect(producer.children()).toHaveLength(5) | ||
| // Ordered by first observation then provider id; these share a publication clock. | ||
| expect(producer.backgroundTaskState()?.tasks?.map((task) => [task.id, task.kind])).toEqual([ | ||
| ['task-agent', 'agent'], | ||
| ['task-flow', 'workflow'], | ||
| ['task-other', 'unknown'], | ||
| ['task-shell', 'command'], | ||
| ['task-watch', 'monitor'] | ||
| ]) | ||
| }) | ||
|
|
||
| it.each([ | ||
| ['working', 'working'], | ||
| ['monitoring', 'working'], | ||
| [undefined, 'working'], | ||
| ['done', 'idle'], | ||
| ['idle', 'idle'], | ||
| ['waiting', 'waiting'], | ||
| ['blocked', 'blocked'], | ||
| ['unverifiable', 'unverifiable'] | ||
| ] as const)('maps provider state %s to %s exactly as the bridge does', (state, expected) => { | ||
| const host = expectParity([agent('task-1', state === undefined ? {} : { state })]) | ||
| expect(host?.[0]?.state).toBe(expected) | ||
| }) | ||
|
|
||
| it('trims provider ids, and admits only non-blank ids of at most 64 characters', () => { | ||
| const host = expectParity([ | ||
| agent(' task-padded '), | ||
| agent(' '), | ||
| agent('x'.repeat(64)), | ||
| agent('y'.repeat(65)) | ||
| ]) | ||
| expect(host?.map((row) => row.id)).toEqual(['task-padded', 'x'.repeat(64)]) | ||
| }) | ||
|
|
||
| it('caps accepted rows at AGENT_STATUS_MAX_SUBAGENTS with invalid rows interleaved', () => { | ||
| for (const count of [31, 32, 33]) { | ||
| const tasks: AgentSessionBackgroundTask[] = [] | ||
| for (let index = 0; index < count; index += 1) { | ||
| // Interleaved rejects must not consume the cap. | ||
| tasks.push({ id: `command-${index}`, kind: 'command', startedAt: index }) | ||
| tasks.push(agent(`agent-${String(index).padStart(2, '0')}`, { startedAt: index })) | ||
| } | ||
| const host = expectParity(tasks) | ||
| expect(host).toHaveLength(Math.min(count, AGENT_STATUS_MAX_SUBAGENTS)) | ||
| } | ||
| }) | ||
|
|
||
| it('never lets the bounded projection evict a canonical record', () => { | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| const tasks = Array.from({ length: 40 }, (_unused, index) => | ||
| agent(`agent-${String(index).padStart(2, '0')}`, { startedAt: index }) | ||
| ) | ||
| producer.publish(roster(tasks)) | ||
| expect(producer.subagents()).toHaveLength(AGENT_STATUS_MAX_SUBAGENTS) | ||
| expect(producer.children()).toHaveLength(40) | ||
| }) | ||
|
|
||
| it('settles rows the authoritative roster stopped listing, as the bridge drops them', () => { | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| producer.publish(roster([agent('task-1'), agent('task-2')])) | ||
| producer.publish(roster([agent('task-1')])) | ||
|
|
||
| expect(producer.subagents()?.map((row) => row.id)).toEqual(['task-1']) | ||
| expect(subagentSnapshotsFromTasks([agent('task-1')])?.map((row) => row.id)).toEqual(['task-1']) | ||
| // Membership moved; the record and its history did not disappear. | ||
| const settled = producer.childFor('task-2') | ||
| expect(settled).toMatchObject({ membership: 'settled', state: 'idle', outcome: 'unknown' }) | ||
| }) | ||
| }) | ||
|
|
||
| describe('the intentional timestamp improvement', () => { | ||
| it('publishes host firstObservedAt where the bridge published the provider stamp', () => { | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| producer.publish(roster([agent('task-1', { startedAt: 5 })]), 4_242) | ||
| expect(producer.subagents()?.[0]?.startedAt).toBe(4_242) | ||
| expect(subagentSnapshotsFromTasks([agent('task-1', { startedAt: 5 })])?.[0]?.startedAt).toBe(5) | ||
| // Provider timing is retained separately rather than overwritten. | ||
| expect(producer.childFor('task-1')?.providerTiming).toEqual({ startedAt: 5 }) | ||
| }) | ||
|
|
||
| it('replaces the bridge’s `?? 0` with a real host observation', () => { | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| const task: AgentSessionBackgroundTask = { id: 'task-1', kind: 'agent' } | ||
| producer.publish(roster([task]), 7_777) | ||
| expect(subagentSnapshotsFromTasks([task])?.[0]?.startedAt).toBe(0) | ||
| expect(producer.subagents()?.[0]?.startedAt).toBe(7_777) | ||
| expect(producer.childFor('task-1')?.providerTiming).toBeUndefined() | ||
| }) | ||
|
|
||
| it('keeps first-observation time stable across later updates', () => { | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| producer.publish(roster([agent('task-1')]), 1_000) | ||
| producer.publish(roster([agent('task-1', { totalTokens: 12 })]), 9_000) | ||
| expect(producer.subagents()?.[0]?.startedAt).toBe(1_000) | ||
| }) | ||
| }) | ||
|
|
||
| describe('declared deviations from the bridge', () => { | ||
| it('refuses a provider id carrying control characters that the bridge admits', () => { | ||
| const task = agent('taskone') | ||
| // The bridge only trims and length-checks, so it publishes the row. | ||
| expect(subagentSnapshotsFromTasks([task])?.[0]?.id).toBe('taskone') | ||
| // A canonical alias is a serialized key; a control character cannot be one. | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| producer.publish(roster([task])) | ||
| expect(producer.children()).toHaveLength(0) | ||
| expect(producer.subagents()).toBeUndefined() | ||
| }) | ||
|
|
||
| it('orders the bounded projection by first observation, not by current roster order', () => { | ||
| const producer = createStructuredChildWorkProducerFixture() | ||
| producer.publish(roster([agent('b-task'), agent('a-task')]), 1_000) | ||
| // The provider reorders; the host keeps the order it first observed them in. | ||
| producer.publish(roster([agent('a-task'), agent('b-task')]), 2_000) | ||
| expect(producer.subagents()?.map((row) => row.id)).toEqual(['a-task', 'b-task']) | ||
| }) | ||
| }) |
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
These facts are written on every projection, and
sinkChildrenruns on every journal publication (both the equal-summary and changed-summary branches), so an unchanged roster still advances the canonical store revision.applyMutationhas no no-op short-circuit, so each publication clones every map and validates the whole store, once for the facts plus once per re-committed child.Technical details