mirror of
https://github.com/stablyai/orca.git
synced 2026-10-03 08:02:12 +00:00
* fix(journal): a batch revision restates the producer of each row it revises The reducer rebuilds a row's producer linkage from its NEWEST revision, and absence is a positive claim: no agent id means the session's own agent wrote the row. So any revision written without the stamp hands a subagent's row back to its parent, permanently. Three host paths revise rows they did not write, from the render item they already hold, and all three dropped the stamp: - answering a prompt re-appended the asker's row with the fence only; - dead-generation settlement failed running tool calls and cancelled pending prompts through a lifecycle batch; - stale-session settlement on acquire cancelled lost prompts the same way. The batch path could not carry a producer at all: linkage was removed from the batch row because one row covers N mutations, with a note that a mixed batch would have to stamp per mutation. Dead-generation settlement is such a batch already, and Codex settlement is about to become one. So each item mutation now names its own producer, inline like the row base. A mutation that names none falls back to the row's linkage, which is what a batch read before. Parse sanitizes a bad per-mutation id the same way it does a row's: the field is dropped and the mutation kept. No schema version bump. An older host's mutation validator ignores unknown keys, so it reads a stamped mutation as the session's own, which is exactly what it shows today. Old journals carry no stamp and read as before. Turn revisions still carry nothing: a turn is the session's unit of work, and the live-turn scans rely on a turn row never carrying linkage. The note recording that invariant is updated to the new write sites. * fix(native-chat): attribute a Codex subagent's journal rows to the subagent Codex journals every thread on its app-server connection into the session's journal, and a spawned subagent's items arrive on the child's own thread. None of those rows carried producer linkage, so under the journal's rule that absence means the session's own agent wrote a row, every child's command, message, reasoning, prompt and status row read as the PARENT's: the parent could show its child's running command, its child's reasoning as "thinking", and its child's prose as its own latest line. The Claude lane's model is reused, not reinvented: the same fields and the same absence rule. What differs is how the producer is known. Orca opens exactly one thread per app-server, so any other thread is one Codex spawned. That decides WHETHER a row is a child's from its first frame, announced or not, and the thread id is final at once: it is never re-minted the way a tool-call reference is, so no correction ledger is needed for identity. - agentId: the child thread id, the same id the status side keys a Codex child on. - parentAgentId: the thread whose stream carried the child's `started` activity. Codex emits that item on the spawning agent's own session, so a child that spawned a grandchild is named; the session's own thread is not. Other activity kinds ride whichever agent acted and are not used. - producerKind: 'agent'. - attempt: which run of the child the row's own turn was, counted from the child turns the roster already observes; absent on the first run. Taken from the row's turn rather than the child's latest, so a persistent shell that outlives its turn keeps its run across revisions. - providerParentRef is omitted: a Codex child's frames carry no parent reference of their own beyond the thread id, which is already agentId. One resolver, owned by the roster (which already owns what is known about each child thread), is handed to every writer: items, streams, generic and summary rows, prompts, compactions, goals, and the three settlement batches. The session-end settlement mixes every thread's rows in one batch, so each mutation names its own producer. Turn rows stay unstamped: Codex writes them only for the primary thread. The spawn-group roster row stays unstamped on purpose: a child's frame can trigger its write, but it is the parent's list of its children. The translator's construction moves to a parts module so the translator stays a router under the line cap, and the item streams reuse one append-and-publish helper instead of two copies. Children are never swept at turn end; nothing here changes that. * test(native-chat): pin Codex subagent attribution at every writer and every parent reader Two layers, so a stamp that is correct in the store and never read, or read and never persisted, cannot pass. The readers, through the real path: translator, deferred sink, on-disk journal, snapshot. Each is a defect on main: the parent named its child's running command as its own tool, read its child's reasoning as itself thinking, showed its child's compaction as its activity line, and quoted its child's prose as its latest line (checked after closing and reopening the journal, so the stamp is read back from disk). The transcript still renders the child's rows. The writers, through a sink that records the linkage of every plain append, batch mutation and lifecycle transition: start, streamed checkpoint and completion of one command all restate the child; a row that beats the spawn announcement is still the child's; a grandchild names the child that announced it, while an `interacted` activity names no parent; a follow-up turn is the child's second run, and a shell that outlives its turn keeps its own; the exit batch settles each thread's rows under its own producer and the turn row under none; a child's provider frames, approval and goal rows are its own; nothing is stamped while the session thread is still opening; and the spawn-group row stays the parent's. * test(native-chat): pin linkage forwarding on the sink's lifecycle-transition path A Codex child's goal row is written through a lifecycle transition, so a sink that forwarded only the fence there would file the child's goal as the session's own. * test(native-chat): type the Codex item fixtures as thread items * refactor(journal): keep a row's producer across revisions that name none The reducer took a row's producer linkage from its newest revision, so every writer that revised a row it did not write - a prompt answer, a dead-generation or stale-session settlement, the reopen sweep of stale subagent rosters - had to restate the producer or silently hand a subagent's row to the session's own agent. Three of those writers had been patched to restate it; the next one to forget would reintroduce the bug. Attribution is now fixed by a row's first write. A revision that names no producer keeps the row's existing linkage; one that names any replaces the whole bundle, which is how a provisional stamp is still corrected in place. A row re-created after a tombstone starts with nothing. The reducer runs the same fold on replay, so the kept producer survives a reopen. The three restatements are removed. Per-mutation linkage on lifecycle batches stays: a batch can create a row (a Codex child's prompt, or a child's item settled before any checkpoint landed) and one batch can mix producers. * test(journal): pin producer inheritance in the reducer and across a reopen A revision naming no producer keeps the row's, on the plain item path and in a batch settling a child's row beside the session's own; one naming any replaces the bundle wholesale; a tombstone clears it; a stale revision cannot touch it; and a reopened journal replays it exactly as it was folded live. * refactor(codex): name the translator's writer factory for what it builds * docs(codex): say why a settled row names its producer * test(journal): drop a producer test the stale-revision guards make unreachable The stale revision is dropped whole by two independent guards before the inheritance rule runs, so its producer assertion could never fail; the reducer's own stale-revision tests already cover the drop. Also say what the batch sink does forward: each mutation's own producer.
189 lines
5.8 KiB
TypeScript
189 lines
5.8 KiB
TypeScript
import type { NativeChatSubagentState } from '../../shared/native-chat-types'
|
|
import { MAX_SUBAGENT_FIELD_CHARS } from '../../shared/native-chat-subagent-summary'
|
|
|
|
const MAX_CHILDREN = 128
|
|
const MAX_SETTLED_TURNS = 256
|
|
/** Turn ordinals remembered per child; a row from an older run reads as its first. */
|
|
const MAX_TURN_ORDINALS_PER_CHILD = 64
|
|
|
|
export type CodexChildExecution = {
|
|
turnId: string
|
|
state: NativeChatSubagentState
|
|
}
|
|
|
|
export type CodexExecutionChild = {
|
|
agentThreadId: string
|
|
registered: boolean
|
|
label: string | null
|
|
parentTurnId: string | null
|
|
/** The thread whose stream carried this child's `started` activity. Codex
|
|
* emits that item on the spawning agent's own session, so it names the parent. */
|
|
spawnerThreadId: string | null
|
|
execution: CodexChildExecution | null
|
|
/** Which run each observed turn was, in the order the child's turns began. */
|
|
turnOrdinals: Map<string, number>
|
|
turnCount: number
|
|
}
|
|
|
|
/** Child turn events own execution; activity items only identify the child. */
|
|
export class CodexSubagentExecutions {
|
|
private readonly children = new Map<string, CodexExecutionChild>()
|
|
private readonly settledTurns = new Map<string, NativeChatSubagentState>()
|
|
|
|
register(
|
|
agentThreadId: string,
|
|
label: string | null,
|
|
parentTurnId: string | null | undefined,
|
|
spawnerThreadId?: string
|
|
): CodexExecutionChild | undefined {
|
|
const child = this.child(agentThreadId)
|
|
if (!child) {
|
|
return undefined
|
|
}
|
|
if (!child.registered || parentTurnId !== undefined) {
|
|
child.parentTurnId = parentTurnId ?? null
|
|
}
|
|
// A child is spawned once; its announcement is delivered twice, never by another thread.
|
|
child.spawnerThreadId ??= spawnerThreadId ?? null
|
|
child.registered = true
|
|
// Retain one overflow unit so the journal can append its per-row truncation marker.
|
|
child.label ??=
|
|
label
|
|
?.trim()
|
|
.replace(/\s+/g, ' ')
|
|
.slice(0, MAX_SUBAGENT_FIELD_CHARS + 1) || null
|
|
return child
|
|
}
|
|
|
|
observeTurn(
|
|
agentThreadId: string,
|
|
turnId: string,
|
|
state: NativeChatSubagentState
|
|
): { child: CodexExecutionChild; execution: CodexChildExecution } | null {
|
|
const key = JSON.stringify([agentThreadId, turnId])
|
|
const settled = this.settledTurns.get(key)
|
|
if (state === 'working' && settled !== undefined) {
|
|
return null
|
|
}
|
|
const child = this.child(agentThreadId)
|
|
if (!child) {
|
|
return null
|
|
}
|
|
this.numberTurn(child, turnId)
|
|
if (
|
|
state === 'working' &&
|
|
child.execution?.turnId === turnId &&
|
|
child.execution.state !== 'working'
|
|
) {
|
|
return null
|
|
}
|
|
const execution = { turnId, state: settled ?? state }
|
|
if (state !== 'working') {
|
|
this.settledTurns.set(key, execution.state)
|
|
while (this.settledTurns.size > MAX_SETTLED_TURNS) {
|
|
const oldest = this.settledTurns.keys().next().value
|
|
if (oldest === undefined) {
|
|
break
|
|
}
|
|
this.settledTurns.delete(oldest)
|
|
}
|
|
}
|
|
if (state === 'working' || !child.execution || child.execution.turnId === turnId) {
|
|
child.execution = execution
|
|
}
|
|
return { child, execution }
|
|
}
|
|
|
|
/** Survives the child's turn, so a row outliving that turn can still name it. */
|
|
label(agentThreadId: string): string | null {
|
|
return this.children.get(agentThreadId)?.label ?? null
|
|
}
|
|
|
|
spawnerOf(agentThreadId: string): string | null {
|
|
return this.children.get(agentThreadId)?.spawnerThreadId ?? null
|
|
}
|
|
|
|
/** Which run of the child a turn was: 1 for the turn it was spawned into, then
|
|
* one more per follow-up turn. Null when the turn was never observed. */
|
|
turnOrdinal(agentThreadId: string, turnId: string): number | null {
|
|
return this.children.get(agentThreadId)?.turnOrdinals.get(turnId) ?? null
|
|
}
|
|
|
|
workingChildren(): CodexExecutionChild[] {
|
|
return [...this.children.values()].filter(
|
|
(child) => child.registered && child.execution?.state === 'working'
|
|
)
|
|
}
|
|
|
|
settleSession(): void {
|
|
for (const child of this.children.values()) {
|
|
if (child.execution?.state === 'working') {
|
|
child.execution = { ...child.execution, state: 'unverifiable' }
|
|
}
|
|
}
|
|
}
|
|
|
|
clear(): void {
|
|
this.children.clear()
|
|
this.settledTurns.clear()
|
|
}
|
|
|
|
/** Retention bounds are not observable through the child/turn API, so expose the two counts. */
|
|
retentionSizes(): { children: number; settledTurns: number } {
|
|
return { children: this.children.size, settledTurns: this.settledTurns.size }
|
|
}
|
|
|
|
private child(agentThreadId: string): CodexExecutionChild | undefined {
|
|
const existing = this.children.get(agentThreadId)
|
|
if (existing) {
|
|
return existing
|
|
}
|
|
if (this.children.size >= MAX_CHILDREN) {
|
|
const settled = [...this.children].find(([, child]) => child.execution?.state !== 'working')
|
|
if (!settled) {
|
|
return undefined
|
|
}
|
|
this.children.delete(settled[0])
|
|
}
|
|
const child: CodexExecutionChild = {
|
|
agentThreadId,
|
|
registered: false,
|
|
label: null,
|
|
parentTurnId: null,
|
|
spawnerThreadId: null,
|
|
execution: null,
|
|
turnOrdinals: new Map(),
|
|
turnCount: 0
|
|
}
|
|
this.children.set(agentThreadId, child)
|
|
return child
|
|
}
|
|
|
|
private numberTurn(child: CodexExecutionChild, turnId: string): void {
|
|
if (child.turnOrdinals.has(turnId)) {
|
|
return
|
|
}
|
|
child.turnCount += 1
|
|
child.turnOrdinals.set(turnId, child.turnCount)
|
|
if (child.turnOrdinals.size > MAX_TURN_ORDINALS_PER_CHILD) {
|
|
const oldest = child.turnOrdinals.keys().next().value
|
|
if (oldest !== undefined) {
|
|
child.turnOrdinals.delete(oldest)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
export function codexChildTurnState(status: unknown): NativeChatSubagentState {
|
|
if (status === 'completed') {
|
|
return 'completed'
|
|
}
|
|
if (status === 'interrupted') {
|
|
return 'stopped'
|
|
}
|
|
if (status === 'failed') {
|
|
return 'failed'
|
|
}
|
|
return 'unverifiable'
|
|
}
|