mirror of
https://github.com/stablyai/orca.git
synced 2026-10-09 08:02:35 +00:00
feat(native-chat): say how a withdrawn Codex send joined the turn it names
The answered turn becomes { turnItemId, via }: 'start' when Codex
answered the send's turn/start with that turn, 'steer' when Orca steered
it into that running turn. A client can then tell the turn's opener from
a send steered into it even when the opener is on an older page.
This commit is contained in:
@@ -133,9 +133,9 @@ describe('the turn Codex answered a send into but has not opened', () => {
|
||||
it('is the turn the latest armed send was answered into', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
|
||||
echoes.arm('client-2')
|
||||
echoes.bindTurn('client-2', 'thread-1', 'turn-2')
|
||||
echoes.bindTurn('client-2', 'thread-1', 'turn-2', 'start')
|
||||
|
||||
expect(echoes.answeredUnopenedTurn('thread-1', NONE_OPEN)).toBe('turn-2')
|
||||
})
|
||||
@@ -143,7 +143,7 @@ describe('the turn Codex answered a send into but has not opened', () => {
|
||||
it('is none once Codex opened that turn', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
|
||||
|
||||
expect(echoes.answeredUnopenedTurn('thread-1', new Set(['turn-1']))).toBeNull()
|
||||
})
|
||||
@@ -151,7 +151,7 @@ describe('the turn Codex answered a send into but has not opened', () => {
|
||||
it('is none once that turn ended, even with its send still armed for an echo', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
|
||||
// A completed end leaves its unechoed send armed.
|
||||
expect(echoes.endTurn('thread-1', 'turn-1', { status: 'completed' })).toEqual([])
|
||||
|
||||
@@ -161,9 +161,9 @@ describe('the turn Codex answered a send into but has not opened', () => {
|
||||
it('skips a turn a wait left unopened, and still names an earlier one', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
|
||||
echoes.arm('client-2')
|
||||
echoes.bindTurn('client-2', 'thread-1', 'turn-2')
|
||||
echoes.bindTurn('client-2', 'thread-1', 'turn-2', 'start')
|
||||
|
||||
echoes.leftUnopened('thread-1', 'turn-2')
|
||||
|
||||
@@ -174,7 +174,7 @@ describe('the turn Codex answered a send into but has not opened', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.arm('client-2')
|
||||
echoes.bindTurn('client-2', 'thread-2', 'turn-2')
|
||||
echoes.bindTurn('client-2', 'thread-2', 'turn-2', 'start')
|
||||
|
||||
expect(echoes.answeredUnopenedTurn('thread-1', NONE_OPEN)).toBeNull()
|
||||
})
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
import type { ProviderDiagnostic } from '../../shared/agent-session-failure'
|
||||
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
|
||||
import type {
|
||||
AgentJournalItemIdentity,
|
||||
AgentJournalTurnJoin
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
|
||||
/** Sends awaiting their echo. One bound to a turn that ended without taking it settles from that
|
||||
* end; any other whose echo never arrives is retired by the journal's recovery on exit. */
|
||||
@@ -35,10 +38,15 @@ export type CodexDispatchEchoes = {
|
||||
/** Drops an armed send whose write never reached the provider. */
|
||||
disarm: (clientMessageId: string) => void
|
||||
/**
|
||||
* Binds a send to the turn Codex answered it into. Returns that turn's end when the answer is
|
||||
* read after it; a send that end settles is no longer armed.
|
||||
* Binds a send to the turn Codex answered it into, and how it joined that turn. Returns that
|
||||
* turn's end when the answer is read after it; a send that end settles is no longer armed.
|
||||
*/
|
||||
bindTurn: (clientMessageId: string, threadId: string, turnId: string) => CodexTurnEnd | null
|
||||
bindTurn: (
|
||||
clientMessageId: string,
|
||||
threadId: string,
|
||||
turnId: string,
|
||||
via: AgentJournalTurnJoin
|
||||
) => CodexTurnEnd | null
|
||||
/** The turn the latest armed send was answered into that is neither in `openTurnIds`, ended, nor
|
||||
* left unopened through a wait: one Codex has picked for the send but not opened. */
|
||||
answeredUnopenedTurn: (threadId: string, openTurnIds: ReadonlySet<string>) => string | null
|
||||
@@ -49,7 +57,11 @@ export type CodexDispatchEchoes = {
|
||||
* completed, which echoes its pending input first, so one it never echoed waits for recovery. An
|
||||
* interrupt withdraws an un-echoed send, steered or the turn's own input: neither reached history.
|
||||
*/
|
||||
endTurn: (threadId: string, turnId: string, end: CodexTurnEnd) => string[]
|
||||
endTurn: (
|
||||
threadId: string,
|
||||
turnId: string,
|
||||
end: CodexTurnEnd
|
||||
) => { clientMessageId: string; via: AgentJournalTurnJoin }[]
|
||||
/** Submission origin for this exact send, retained until its echo settles it. */
|
||||
requestOrigin: (clientMessageId: string) => CodexDispatchRequestOrigin | null
|
||||
/** Highest causal sequence assigned to a dispatch in this session. */
|
||||
@@ -61,7 +73,11 @@ export type CodexDispatchEchoes = {
|
||||
export function createCodexDispatchEchoes(): CodexDispatchEchoes {
|
||||
const armed = new Map<
|
||||
string,
|
||||
{ requestedAt: number | null; sequence: number; turn?: { threadId: string; turnId: string } }
|
||||
{
|
||||
requestedAt: number | null
|
||||
sequence: number
|
||||
turn?: { threadId: string; turnId: string; via: AgentJournalTurnJoin }
|
||||
}
|
||||
>()
|
||||
const endedTurns = new Map<string, CodexTurnEnd>()
|
||||
const unopenedTurns = new Set<string>()
|
||||
@@ -85,12 +101,12 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes {
|
||||
},
|
||||
settle: (clientMessageId) => armed.delete(clientMessageId),
|
||||
disarm: (clientMessageId) => void armed.delete(clientMessageId),
|
||||
bindTurn: (clientMessageId, threadId, turnId) => {
|
||||
bindTurn: (clientMessageId, threadId, turnId, via) => {
|
||||
const entry = armed.get(clientMessageId)
|
||||
if (!entry) {
|
||||
return null
|
||||
}
|
||||
entry.turn = { threadId, turnId }
|
||||
entry.turn = { threadId, turnId, via }
|
||||
const end = endedTurns.get(turnKey(threadId, turnId)) ?? null
|
||||
if (end && settles(end)) {
|
||||
armed.delete(clientMessageId)
|
||||
@@ -132,10 +148,10 @@ export function createCodexDispatchEchoes(): CodexDispatchEchoes {
|
||||
}
|
||||
const settled = [...armed].flatMap(([clientMessageId, entry]) =>
|
||||
entry.turn && turnKey(entry.turn.threadId, entry.turn.turnId) === turn
|
||||
? [clientMessageId]
|
||||
? [{ clientMessageId, via: entry.turn.via }]
|
||||
: []
|
||||
)
|
||||
for (const clientMessageId of settled) {
|
||||
for (const { clientMessageId } of settled) {
|
||||
armed.delete(clientMessageId)
|
||||
}
|
||||
return settled
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import type {
|
||||
AgentJournalAnsweredTurnIdentity,
|
||||
AgentJournalItemIdentity,
|
||||
AgentSessionJournalIdentity
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
@@ -86,7 +87,7 @@ export type CodexStructuredSessionAdapterDeps = {
|
||||
| { providerIdentity: AgentJournalItemIdentity }
|
||||
| ({
|
||||
state: 'rejected'
|
||||
answeredInTurn: AgentJournalItemIdentity
|
||||
answeredInTurn: AgentJournalAnsweredTurnIdentity
|
||||
} & AgentJournalDispatchRejection)
|
||||
)
|
||||
) => void
|
||||
|
||||
@@ -317,7 +317,10 @@ describe('the turn a withdrawn Codex send names', () => {
|
||||
expect(rig.turnRecords.length).toBeGreaterThan(0)
|
||||
expect(new Set(rig.turnRecords.map(agentJournalItemKey)).size).toBe(1)
|
||||
expect(rig.settlements).toEqual([
|
||||
expect.objectContaining({ clientMessageId: 'client-1', answeredInTurn: rig.turnRecords[0] })
|
||||
expect.objectContaining({
|
||||
clientMessageId: 'client-1',
|
||||
answeredInTurn: { turn: rig.turnRecords[0], via: 'start' }
|
||||
})
|
||||
])
|
||||
})
|
||||
|
||||
@@ -329,8 +332,14 @@ describe('the turn a withdrawn Codex send names', () => {
|
||||
rig.turns.end('interrupted')
|
||||
|
||||
expect(rig.settlements.map((settlement) => [settlement.clientMessageId, settlement])).toEqual([
|
||||
['client-1', expect.objectContaining({ answeredInTurn: rig.turnRecords[0] })],
|
||||
['client-2', expect.objectContaining({ answeredInTurn: rig.turnRecords[0] })]
|
||||
[
|
||||
'client-1',
|
||||
expect.objectContaining({ answeredInTurn: { turn: rig.turnRecords[0], via: 'start' } })
|
||||
],
|
||||
[
|
||||
'client-2',
|
||||
expect.objectContaining({ answeredInTurn: { turn: rig.turnRecords[0], via: 'steer' } })
|
||||
]
|
||||
])
|
||||
})
|
||||
|
||||
@@ -345,7 +354,7 @@ describe('the turn a withdrawn Codex send names', () => {
|
||||
|
||||
await expect(sending).resolves.toMatchObject({
|
||||
state: 'rejected',
|
||||
answeredInTurn: rig.turnRecords[0]
|
||||
answeredInTurn: { turn: rig.turnRecords[0], via: 'start' }
|
||||
})
|
||||
})
|
||||
|
||||
@@ -369,9 +378,11 @@ describe('a send bound to a turn', () => {
|
||||
it('dies with the settlement its turn end makes', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
|
||||
|
||||
expect(echoes.endTurn('thread-1', 'turn-1', { status: 'interrupted' })).toEqual(['client-1'])
|
||||
expect(echoes.endTurn('thread-1', 'turn-1', { status: 'interrupted' })).toEqual([
|
||||
{ clientMessageId: 'client-1', via: 'start' }
|
||||
])
|
||||
expect(echoes.size).toBe(0)
|
||||
expect(echoes.settle('client-1')).toBe(false)
|
||||
})
|
||||
@@ -379,20 +390,20 @@ describe('a send bound to a turn', () => {
|
||||
it('dies with its child, which forgets recorded turn ends too', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
|
||||
echoes.endTurn('thread-2', 'turn-2', { status: 'interrupted' })
|
||||
|
||||
echoes.clear()
|
||||
|
||||
expect(echoes.size).toBe(0)
|
||||
echoes.arm('client-2')
|
||||
expect(echoes.bindTurn('client-2', 'thread-2', 'turn-2')).toBeNull()
|
||||
expect(echoes.bindTurn('client-2', 'thread-2', 'turn-2', 'start')).toBeNull()
|
||||
})
|
||||
|
||||
it('is matched by thread as well as turn id', () => {
|
||||
const echoes = createCodexDispatchEchoes()
|
||||
echoes.arm('client-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1')
|
||||
echoes.bindTurn('client-1', 'thread-1', 'turn-1', 'start')
|
||||
|
||||
expect(echoes.endTurn('thread-2', 'turn-1', { status: 'interrupted' })).toEqual([])
|
||||
expect(echoes.size).toBe(1)
|
||||
|
||||
@@ -17,7 +17,7 @@ import {
|
||||
agentSessionFailureWords,
|
||||
type AgentJournalDispatchRejection
|
||||
} from '../../shared/agent-session-failure-words'
|
||||
import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types'
|
||||
import type { AgentJournalAnsweredTurnIdentity } from '../../shared/agent-session-journal-types'
|
||||
import type { CodexTurnEnd } from './codex-structured-dispatch-echo'
|
||||
import { codexTurnLifecycleIdentity } from './codex-structured-journal-translation-turns'
|
||||
import type { CodexSession } from './codex-structured-session-state'
|
||||
@@ -43,8 +43,8 @@ export function codexDispatchRejection(
|
||||
export type CodexTurnEndSettlement = {
|
||||
clientMessageId: string
|
||||
state: 'rejected'
|
||||
/** The turn Codex answered the send into, whose end settled it. */
|
||||
answeredInTurn: AgentJournalItemIdentity
|
||||
/** The turn Codex answered the send into, whose end settled it, and how the send joined it. */
|
||||
answeredInTurn: AgentJournalAnsweredTurnIdentity
|
||||
} & AgentJournalDispatchRejection
|
||||
|
||||
function errorDetail(params: unknown): ProviderDiagnostic | undefined {
|
||||
@@ -97,10 +97,14 @@ export function settleCodexSendsInEndedTurn(
|
||||
return
|
||||
}
|
||||
const rejection = codexTurnEndRejection(end)
|
||||
const answeredInTurn = codexTurnLifecycleIdentity(frame.sessionId, turnId)
|
||||
for (const clientMessageId of session.dispatchEchoes.endTurn(session.threadId, turnId, end)) {
|
||||
const turn = codexTurnLifecycleIdentity(frame.sessionId, turnId)
|
||||
for (const { clientMessageId, via } of session.dispatchEchoes.endTurn(
|
||||
session.threadId,
|
||||
turnId,
|
||||
end
|
||||
)) {
|
||||
if (rejection) {
|
||||
settle({ clientMessageId, state: 'rejected', answeredInTurn, ...rejection })
|
||||
settle({ clientMessageId, state: 'rejected', answeredInTurn: { turn, via }, ...rejection })
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,8 @@
|
||||
import { agentSessionFailureFact, providerDiagnosticOf } from '../../shared/agent-session-failure'
|
||||
import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
|
||||
import type {
|
||||
AgentJournalMessageItem,
|
||||
AgentJournalTurnJoin
|
||||
} from '../../shared/agent-session-journal-types'
|
||||
import type { NativeChatBlock } from '../../shared/native-chat-types'
|
||||
import type { AgentSessionDispatchOutcome } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
|
||||
import {
|
||||
@@ -108,7 +111,7 @@ async function steerCodexTurn(
|
||||
host: CodexTurnHost,
|
||||
expectedTurnId: string,
|
||||
input: { clientMessageId: string; body: AgentJournalMessageItem; timeoutMs?: number }
|
||||
): Promise<{ turnId: string } | null> {
|
||||
): Promise<{ turnId: string; via: 'steer' } | null> {
|
||||
try {
|
||||
const answer = await host.connection.request(
|
||||
'turn/steer',
|
||||
@@ -120,7 +123,7 @@ async function steerCodexTurn(
|
||||
},
|
||||
{ timeoutMs: input.timeoutMs }
|
||||
)
|
||||
return { turnId: readCodexTurnId(answer) ?? expectedTurnId }
|
||||
return { turnId: readCodexTurnId(answer) ?? expectedTurnId, via: 'steer' }
|
||||
} catch (error) {
|
||||
if (isCodexAppServerRequestError(error) || isCodexAppServerUnsupportedError(error)) {
|
||||
return null
|
||||
@@ -142,7 +145,7 @@ export async function startCodexTurn(
|
||||
requestedAt?: number
|
||||
timeoutMs?: number
|
||||
}
|
||||
): Promise<{ turnId: string | null } | false> {
|
||||
): Promise<{ turnId: string | null; via: AgentJournalTurnJoin } | false> {
|
||||
// Armed before the write: the echo and `turn/started` can both land while the
|
||||
// response is in flight, and the start must snapshot this send in its frontier.
|
||||
if (!host.dispatchEchoes.arm(input.clientMessageId, input.requestedAt)) {
|
||||
@@ -168,7 +171,7 @@ export async function startCodexTurn(
|
||||
},
|
||||
{ timeoutMs: input.timeoutMs }
|
||||
)
|
||||
return { turnId: readCodexTurnId(answer) }
|
||||
return { turnId: readCodexTurnId(answer), via: 'start' }
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -187,7 +190,7 @@ export async function dispatchCodexTurn(
|
||||
},
|
||||
timeoutMs: number | undefined
|
||||
): Promise<AgentSessionDispatchOutcome> {
|
||||
let answer: { turnId: string | null } | false
|
||||
let answer: { turnId: string | null; via: AgentJournalTurnJoin } | false
|
||||
try {
|
||||
answer = await startCodexTurn(session, { ...input, timeoutMs })
|
||||
} catch (error) {
|
||||
@@ -211,13 +214,21 @@ export async function dispatchCodexTurn(
|
||||
}
|
||||
// An answer read after the turn it names already ended is settled by that end.
|
||||
const endedFirst = answer.turnId
|
||||
? session.dispatchEchoes.bindTurn(input.clientMessageId, session.threadId, answer.turnId)
|
||||
? session.dispatchEchoes.bindTurn(
|
||||
input.clientMessageId,
|
||||
session.threadId,
|
||||
answer.turnId,
|
||||
answer.via
|
||||
)
|
||||
: null
|
||||
const rejection = endedFirst ? codexTurnEndRejection(endedFirst) : null
|
||||
return rejection && answer.turnId
|
||||
? {
|
||||
state: 'rejected',
|
||||
answeredInTurn: codexTurnLifecycleIdentity(input.sessionId, answer.turnId),
|
||||
answeredInTurn: {
|
||||
turn: codexTurnLifecycleIdentity(input.sessionId, answer.turnId),
|
||||
via: answer.via
|
||||
},
|
||||
...rejection
|
||||
}
|
||||
: { state: 'admitted' }
|
||||
|
||||
@@ -6,6 +6,7 @@ import {
|
||||
type UnreadAgentSessionFailureFact
|
||||
} from '../../../shared/agent-session-failure'
|
||||
import { agentJournalSubmissionKey } from '../../../shared/agent-session-journal-item-key'
|
||||
import type { AgentJournalAnsweredTurn } from '../../../shared/agent-session-journal-types'
|
||||
import { journalDispatchRowApplies } from './journal-dispatch-settlement'
|
||||
import type { JournalReducerState } from './journal-reducer'
|
||||
import {
|
||||
@@ -34,14 +35,11 @@ export function applyJournalDispatchRow(
|
||||
} else {
|
||||
delete submission.rejection
|
||||
}
|
||||
if (
|
||||
row.state === 'rejected' &&
|
||||
typeof row.answeredInTurnItemId === 'string' &&
|
||||
row.answeredInTurnItemId
|
||||
) {
|
||||
submission.answeredInTurnItemId = row.answeredInTurnItemId
|
||||
const answeredInTurn = row.state === 'rejected' ? readAnsweredTurn(row.answeredInTurn) : undefined
|
||||
if (answeredInTurn) {
|
||||
submission.answeredInTurn = answeredInTurn
|
||||
} else {
|
||||
delete submission.answeredInTurnItemId
|
||||
delete submission.answeredInTurn
|
||||
}
|
||||
submission.resolvedAt = row.state === 'pending' ? null : row.ts
|
||||
if (row.state === 'pending') {
|
||||
@@ -70,6 +68,19 @@ export function applyJournalDispatchRow(
|
||||
})
|
||||
}
|
||||
|
||||
/** A stored answered turn; one malformed, or naming a way of joining this build does not know, is
|
||||
* dropped, never the row. */
|
||||
function readAnsweredTurn(value: unknown): AgentJournalAnsweredTurn | undefined {
|
||||
if (typeof value !== 'object' || value === null) {
|
||||
return undefined
|
||||
}
|
||||
const turnItemId = 'turnItemId' in value ? value.turnItemId : undefined
|
||||
const via = 'via' in value ? value.via : undefined
|
||||
return typeof turnItemId === 'string' && turnItemId && (via === 'start' || via === 'steer')
|
||||
? { turnItemId, via }
|
||||
: undefined
|
||||
}
|
||||
|
||||
/** A stored rejection fact, read where it can be placed; a kind it cannot place is kept as
|
||||
* written, so the classifier still knows a fact was there without this build claiming what it
|
||||
* says. Shared with the queued-draft table, whose returned card mirrors its submission. */
|
||||
|
||||
@@ -111,7 +111,12 @@ export function journalDispatchRowBuilder(
|
||||
reason: boundedDispatchReason(input),
|
||||
...(input.state === 'rejected' ? { rejection: input.rejection } : {}),
|
||||
...(input.state === 'rejected' && input.answeredInTurn
|
||||
? { answeredInTurnItemId: agentJournalItemKey(input.answeredInTurn) }
|
||||
? {
|
||||
answeredInTurn: {
|
||||
turnItemId: agentJournalItemKey(input.answeredInTurn.turn),
|
||||
via: input.answeredInTurn.via
|
||||
}
|
||||
}
|
||||
: {}),
|
||||
...journalRowBase(state().epoch, seq, input.fence, ts),
|
||||
...(input.recovered ? { recovered: input.recovered } : {}),
|
||||
|
||||
@@ -16,6 +16,7 @@
|
||||
import type { AgentSessionFailureFact } from '../../../shared/agent-session-failure'
|
||||
import {
|
||||
AGENT_SESSION_JOURNAL_SCHEMA_VERSION,
|
||||
type AgentJournalAnsweredTurn,
|
||||
type AgentJournalDispatchState,
|
||||
type AgentJournalItemBody,
|
||||
type AgentJournalMessageItem,
|
||||
@@ -150,10 +151,10 @@ export type JournalDispatchRow = JournalRowBase & {
|
||||
/** On `rejected`: why, typed. Older readers keep the key and ignore it; a malformed one is
|
||||
* dropped when read, never the row. */
|
||||
rejection?: AgentSessionFailureFact
|
||||
/** On `rejected`: the turn record a Codex send was answered into (`turn/start` or `turn/steer`)
|
||||
* when that turn's end settled the send. Absent on every other row. Older readers keep the key
|
||||
* and ignore it. */
|
||||
answeredInTurnItemId?: string
|
||||
/** On `rejected`: the turn a Codex send was answered into, and how it joined it, when that
|
||||
* turn's end settled the send. Absent on every other row. Older readers keep the key and ignore
|
||||
* it; a malformed one is dropped when read, never the row. */
|
||||
answeredInTurn?: AgentJournalAnsweredTurn
|
||||
}
|
||||
|
||||
/** An item mutation may name its own producer, because one batch can CREATE
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import type { AgentJournalDispatchRejection } from '../../../shared/agent-session-failure-words'
|
||||
import type {
|
||||
AgentJournalAnsweredTurnIdentity,
|
||||
AgentJournalCursor,
|
||||
AgentJournalItemBody,
|
||||
AgentJournalItemIdentity,
|
||||
@@ -42,7 +43,7 @@ export type ResolveDispatchInput = {
|
||||
* `agentSessionFailureWords`, never written by hand. */
|
||||
| ({
|
||||
state: 'rejected'
|
||||
answeredInTurn?: AgentJournalItemIdentity
|
||||
answeredInTurn?: AgentJournalAnsweredTurnIdentity
|
||||
} & AgentJournalDispatchRejection)
|
||||
| { state: 'unknown'; reason?: string | null }
|
||||
)
|
||||
|
||||
@@ -94,7 +94,7 @@ async function sendWithdrawnByItsTurnEnd() {
|
||||
clientMessageId: 'send-1',
|
||||
state: 'rejected',
|
||||
...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }),
|
||||
answeredInTurn: TURN,
|
||||
answeredInTurn: { turn: TURN, via: 'start' },
|
||||
fence: 1
|
||||
})
|
||||
return journal
|
||||
@@ -107,23 +107,25 @@ describe('the turn a rejected submission was answered into', () => {
|
||||
|
||||
expect(journal.submission('send-1')).toMatchObject({
|
||||
dispatchState: 'rejected',
|
||||
answeredInTurnItemId: turnItemId
|
||||
answeredInTurn: { turnItemId, via: 'start' }
|
||||
})
|
||||
expect(readAgentSessionHydrationPage(journal).submissions).toEqual([
|
||||
expect.objectContaining({ answeredInTurnItemId: turnItemId })
|
||||
expect.objectContaining({ answeredInTurn: { turnItemId, via: 'start' } })
|
||||
])
|
||||
await journals.closeAll()
|
||||
const replayed = await journals.open({ identity: IDENTITY, stateDirectory: root! })
|
||||
expect(replayed.submission('send-1')).toMatchObject({ answeredInTurnItemId: turnItemId })
|
||||
expect(replayed.submission('send-1')).toMatchObject({
|
||||
answeredInTurn: { turnItemId, via: 'start' }
|
||||
})
|
||||
})
|
||||
|
||||
it('is absent on a take-back that names no turn, as on rows from older hosts', async () => {
|
||||
const { journal } = await sendHandedOverThenWithdrawn()
|
||||
|
||||
expect(journal.submission('send-1')).toMatchObject({ dispatchState: 'rejected' })
|
||||
expect(journal.submission('send-1')).not.toHaveProperty('answeredInTurnItemId')
|
||||
expect(journal.submission('send-1')).not.toHaveProperty('answeredInTurn')
|
||||
expect(readAgentSessionHydrationPage(journal).submissions[0]).not.toHaveProperty(
|
||||
'answeredInTurnItemId'
|
||||
'answeredInTurn'
|
||||
)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -12,6 +12,7 @@ import type {
|
||||
|
||||
import type { AgentSessionBackgroundTaskStops } from '../../../shared/agent-child-work-stop-targets'
|
||||
import type {
|
||||
AgentJournalAnsweredTurnIdentity,
|
||||
AgentJournalItemIdentity,
|
||||
AgentJournalItemBody,
|
||||
AgentJournalMessageItem,
|
||||
@@ -180,7 +181,7 @@ export type AgentSessionDispatchOutcome =
|
||||
* provider answered the send into, which ended before the answer was read. */
|
||||
| ({
|
||||
state: 'rejected'
|
||||
answeredInTurn?: AgentJournalItemIdentity
|
||||
answeredInTurn?: AgentJournalAnsweredTurnIdentity
|
||||
} & AgentJournalDispatchRejection)
|
||||
/** The call did not settle. Never re-send on the user's behalf. */
|
||||
| { state: 'unknown'; reason: string }
|
||||
|
||||
@@ -1,4 +1,7 @@
|
||||
import type { AgentJournalItemIdentity } from '../../../shared/agent-session-journal-types'
|
||||
import type {
|
||||
AgentJournalAnsweredTurnIdentity,
|
||||
AgentJournalItemIdentity
|
||||
} from '../../../shared/agent-session-journal-types'
|
||||
import type { AgentJournalDispatchRejection } from '../../../shared/agent-session-failure-words'
|
||||
import type { StructuredAgentSessionMutationContext } from './structured-agent-session-host-mutations'
|
||||
import { structuredAgentSessionConversationFence } from './structured-agent-session-provider-child'
|
||||
@@ -14,7 +17,7 @@ export async function settleStructuredAgentSessionLateDispatch(
|
||||
| { providerIdentity: AgentJournalItemIdentity }
|
||||
| ({
|
||||
state: 'rejected'
|
||||
answeredInTurn?: AgentJournalItemIdentity
|
||||
answeredInTurn?: AgentJournalAnsweredTurnIdentity
|
||||
} & AgentJournalDispatchRejection)
|
||||
| { state: 'unknown'; reason: string }
|
||||
)
|
||||
|
||||
@@ -317,12 +317,12 @@ describe('the turn a withdrawn Codex send was answered into', () => {
|
||||
item.body.kind === 'turn' ? [item.itemId] : []
|
||||
),
|
||||
named: snapshot.submissions.find((entry) => entry.clientMessageId === clientMessageId)
|
||||
?.answeredInTurnItemId,
|
||||
onPage: onPage?.answeredInTurnItemId
|
||||
?.answeredInTurn,
|
||||
onPage: onPage?.answeredInTurn
|
||||
}
|
||||
}
|
||||
|
||||
it("is that turn's record, for the send that opened it and for one steered into it", async () => {
|
||||
it("is that turn's record, started by the send that opened it and steered by a later one", async () => {
|
||||
const opening = await send('look around')
|
||||
await vi.waitFor(() => expect(answers).toBe(1))
|
||||
turns.start()
|
||||
@@ -336,13 +336,13 @@ describe('the turn a withdrawn Codex send was answered into', () => {
|
||||
|
||||
const { turnRecords } = await answeredInto(opening)
|
||||
expect(turnRecords).toHaveLength(1)
|
||||
expect(await answeredInto(opening)).toMatchObject({
|
||||
named: turnRecords[0],
|
||||
onPage: turnRecords[0]
|
||||
})
|
||||
expect(await answeredInto(steered)).toMatchObject({
|
||||
named: turnRecords[0],
|
||||
onPage: turnRecords[0]
|
||||
const started = { turnItemId: turnRecords[0], via: 'start' }
|
||||
const steeredIn = { turnItemId: turnRecords[0], via: 'steer' }
|
||||
expect(await answeredInto(opening)).toEqual({ turnRecords, named: started, onPage: started })
|
||||
expect(await answeredInto(steered)).toEqual({
|
||||
turnRecords,
|
||||
named: steeredIn,
|
||||
onPage: steeredIn
|
||||
})
|
||||
})
|
||||
|
||||
@@ -373,7 +373,7 @@ describe('the turn a withdrawn Codex send was answered into', () => {
|
||||
)
|
||||
const { turnRecords, named } = await answeredInto(sent)
|
||||
expect(turnRecords).toHaveLength(1)
|
||||
expect(named).toBe(turnRecords[0])
|
||||
expect(named).toEqual({ turnItemId: turnRecords[0], via: 'start' })
|
||||
})
|
||||
})
|
||||
|
||||
|
||||
@@ -319,7 +319,8 @@ export const AgentJournalSubmissionSchema = z.object({
|
||||
submittedAt: z.number(),
|
||||
resolvedAt: z.number().nullable(),
|
||||
submittedSequence: z.number().int().optional(),
|
||||
answeredInTurnItemId: z.string().min(1).optional(),
|
||||
// Open like `dispatchState`: a way of joining a newer host names must not drop the submission.
|
||||
answeredInTurn: z.object({ turnItemId: z.string().min(1), via: z.string().min(1) }).optional(),
|
||||
recovered: z.literal(true).optional(),
|
||||
handoverRecorded: z.literal(true).optional(),
|
||||
handedOverAt: z.number().optional(),
|
||||
|
||||
@@ -391,6 +391,18 @@ export type AgentJournalRenderItem = AgentJournalProducerLinkage & {
|
||||
|
||||
// ─── Submissions ────────────────────────────────────────────────────────────
|
||||
|
||||
/** The turn a send was answered into: its record's item id, and how the send joined it. `start`:
|
||||
* the provider answered the send's start request with that turn; `steer`: Orca steered it into
|
||||
* that running turn. Known limit: a start the provider silently folds into a running turn reads
|
||||
* as `start`. Open: a newer host may name another way, which a reader leaves unclaimed. */
|
||||
export type AgentJournalAnsweredTurn = { turnItemId: string; via: AgentJournalTurnJoin }
|
||||
export type AgentJournalTurnJoin = 'start' | 'steer'
|
||||
/** The same, as a writer names it: the turn record's identity, keyed when the row is written. */
|
||||
export type AgentJournalAnsweredTurnIdentity = {
|
||||
turn: AgentJournalItemIdentity
|
||||
via: AgentJournalTurnJoin
|
||||
}
|
||||
|
||||
export const AGENT_JOURNAL_DISPATCH_STATES = ['pending', 'accepted', 'rejected', 'unknown'] as const
|
||||
export type AgentJournalDispatchState = (typeof AGENT_JOURNAL_DISPATCH_STATES)[number]
|
||||
|
||||
@@ -415,10 +427,10 @@ export type AgentJournalSubmission = {
|
||||
* never stored. Needed because a rejected send's own row moves to its rejection, which erases
|
||||
* where it was sent. Absent from hosts that predate it. */
|
||||
submittedSequence?: number
|
||||
/** On `rejected`: the item id of the turn record a Codex send was answered into, when that turn
|
||||
* ended without taking it. Absent on every other send: accepted ones (the echo places them),
|
||||
* Claude, queued take-backs, restart recovery, and rows from hosts that predate it. */
|
||||
answeredInTurnItemId?: string
|
||||
/** On `rejected`: the turn a Codex send was answered into, when that turn ended without taking
|
||||
* it. Absent on every other send: accepted ones (the echo places them), Claude, queued
|
||||
* take-backs, restart recovery, and rows from hosts that predate it. */
|
||||
answeredInTurn?: AgentJournalAnsweredTurn
|
||||
/** Set when crash reconciliation resolved the dispatch, not the provider. A live
|
||||
* `unknown` is a send still outstanding; a recovered one outlived its writer. */
|
||||
recovered?: true
|
||||
|
||||
@@ -40,7 +40,7 @@ function withoutPosition(page: AgentSessionHistoryPage): AgentSessionHistoryPage
|
||||
return {
|
||||
...page,
|
||||
submissions: page.submissions.map(
|
||||
({ submittedSequence: _submitted, answeredInTurnItemId: _turn, ...rest }) => rest
|
||||
({ submittedSequence: _submitted, answeredInTurn: _turn, ...rest }) => rest
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -78,10 +78,13 @@ test('a released client folds and draws a page whose submissions carry their jou
|
||||
state: 'rejected',
|
||||
...agentSessionFailureWords(agentSessionFailureFact('cancelled'), { surface: 'rejection' }),
|
||||
answeredInTurn: {
|
||||
provider: 'legacy',
|
||||
agent: 'codex',
|
||||
sessionId: IDENTITY.sessionId,
|
||||
recordId: 'turn-lifecycle:turn-1'
|
||||
turn: {
|
||||
provider: 'legacy',
|
||||
agent: 'codex',
|
||||
sessionId: IDENTITY.sessionId,
|
||||
recordId: 'turn-lifecycle:turn-1'
|
||||
},
|
||||
via: 'start'
|
||||
},
|
||||
fence: 1
|
||||
})
|
||||
@@ -93,7 +96,9 @@ test('a released client folds and draws a page whose submissions carry their jou
|
||||
})
|
||||
)
|
||||
expect(page.submissions.find((entry) => entry.clientMessageId === 'turn-ended')).toEqual(
|
||||
expect.objectContaining({ answeredInTurnItemId: expect.any(String) })
|
||||
expect.objectContaining({
|
||||
answeredInTurn: { turnItemId: expect.any(String), via: 'start' }
|
||||
})
|
||||
)
|
||||
|
||||
const checkout = await materializeReleaseCheckout(BASELINE_REF)
|
||||
@@ -117,7 +122,7 @@ test('a released client folds and draws a page whose submissions carry their jou
|
||||
const state = reduce(empty, { type: 'history-page', page: from })
|
||||
return {
|
||||
submissions: state.submissions.map(
|
||||
({ submittedSequence: _s, answeredInTurnItemId: _t, ...rest }) => rest
|
||||
({ submittedSequence: _s, answeredInTurn: _t, ...rest }) => rest
|
||||
),
|
||||
messages: project(state.items, [], state.submissions)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user