mirror of
https://github.com/stablyai/orca.git
synced 2026-10-04 16:02:08 +00:00
* refactor(native-chat): give the structured chat host one required logger The structured chat runtime took an optional onError callback that the desktop never passed, so a late dispatch settlement, an unanswered-dispatch release, a journal event-sink write and a provider lifecycle delivery that failed were dropped with no trace. Other host failures went to scattered console.warn calls, which reach nothing in a packaged desktop build. The runtime and host now take one required logger (warn/error with a scope and fields). The production logger writes each entry as a failed span to <userData>/logs/main.trace.ndjson, which the diagnostic bundle collects, and to the console (stderr under a supervised headless host). The runtime and the host wrap it so a logger that throws never fails what it reports, and the install refuses without one. Sites that deliberately kept a recovery-capsule error out of the log still log no error object. * refactor(native-chat): hand the chat host's collaborators the logger, and give orcad its trace file The delivery loop, idle sweep, queued-message drain, lease renewer, event sink, conversation map and provider start/exit settlement each took an internal error callback that the host mapped onto the logger. They now take the logger itself and log under their own scope. The event sink keeps one onFailed hook, which decides whether to stop the provider, not whether to report. The dead-generation settlement returns its failure so each caller logs it under its own scope. orcad now installs the desktop's local trace sink under its own data root, so a headless host's chat failures reach <data-root>/logs/main.trace.ndjson as well as stderr. Also passes the logger in the test fixtures the first commit missed, which tc:node caught. * fix(native-chat): keep repeated chat failures from flooding the trace file, and record their causes - The production structured-chat logger writes a repeated failure (same level, scope, session, message and error text) once per 5 minutes, carrying how many repeats it swallowed; the tracked set is capped at 256. - Trace entries now carry the error's code (and SQLite errcode) and up to three causes by name and message. - A chat read whose conversation will not open is logged through the host's logger (open-for-read), and so are the runtime's chat-tab bookkeeping failures that already hold the host. - orcad writes its own orcad.trace.ndjson, closes it after every quit handler, and flushes it on process exit; a trace file that cannot be opened leaves tracing off instead of stopping the app or orcad. - Tests: the desktop wiring test proves the logger reaches the trace sink, and the privacy tests read every level the logger received. * fix(native-chat): log a created chat's tab-publication and launch-prompt failures through the host's logger * fix(native-chat): key a repeated chat failure on everything its entry writes The repeat suppression keyed on the message and the error's text, so two refusals with the same code but different causes, a plain error and a refusal of one code, or two object-valued errors shared a key and the second was swallowed for five minutes. The key is now the entry's whole written content (fields, code, errcode, refusal reason, cause chain, a stable rendering of a non-error value) plus the error's name and message; a refusal's reason is also written. * test(native-chat): pin that an error's name keeps two repeated failures apart * test(native-chat): build the refusal in the repeat-key test as the wire does * fix(native-chat): read an error's code and a refusal's reason by narrowing, not Reflect.get
254 lines
10 KiB
TypeScript
254 lines
10 KiB
TypeScript
import { mkdtemp, rm } from 'node:fs/promises'
|
|
import { tmpdir } from 'node:os'
|
|
import { join } from 'node:path'
|
|
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
|
import type {
|
|
AgentSessionHistoryRequest,
|
|
AgentSessionHistoryResult,
|
|
AgentSessionStatusSummary,
|
|
AgentSessionSubscribeEvent
|
|
} from '../../src/shared/agent-session-wire'
|
|
import type { AgentJournalCursor } from '../../src/shared/agent-session-journal-types'
|
|
import {
|
|
EMPTY_STRUCTURED_AGENT_SESSION,
|
|
reduceStructuredAgentSession
|
|
} from '../../src/shared/structured-agent-session-reducer'
|
|
import {
|
|
hasUnansweredStructuredAgentSessionDispatch,
|
|
projectStructuredAgentSessionStatus
|
|
} from '../../src/shared/structured-agent-session-projection'
|
|
import { createTrackedJournalOpener } from '../../src/main/native-chat/agent-session-journal/journal-host-database-test-support'
|
|
import { readAgentSessionHistory } from '../../src/main/native-chat/agent-session-wire/agent-session-history-page'
|
|
import { AgentSessionSubscribers } from '../../src/main/native-chat/agent-session-wire/structured-agent-session-subscribers'
|
|
import { StructuredAgentSessionStatusFeed } from '../../src/main/native-chat/agent-session-wire/structured-agent-session-status-feed'
|
|
|
|
const mocks = vi.hoisted(() => ({ call: vi.fn(), subscribe: vi.fn() }))
|
|
vi.mock('@/runtime/structured-agent-session-client', () => ({
|
|
callStructuredAgentSession: mocks.call,
|
|
subscribeStructuredAgentSession: mocks.subscribe
|
|
}))
|
|
|
|
import {
|
|
getStructuredAgentSessionReadOwner,
|
|
resetStructuredAgentSessionReadOwnersForTests
|
|
} from '../../src/renderer/src/components/native-chat/structured-agent-session-read-owner'
|
|
import { createStructuredAgentSessionLogger } from '../../src/main/native-chat/agent-session-wire/structured-agent-session-logger'
|
|
|
|
const SESSION = 'cursor-body-regression'
|
|
const target = { kind: 'local' } as const
|
|
const journals = createTrackedJournalOpener()
|
|
let root: string
|
|
|
|
beforeEach(async () => {
|
|
vi.resetAllMocks()
|
|
root = await mkdtemp(join(tmpdir(), 'orca-cursor-body-'))
|
|
})
|
|
afterEach(async () => {
|
|
resetStructuredAgentSessionReadOwnersForTests()
|
|
await journals.closeAll()
|
|
await rm(root, { recursive: true, force: true })
|
|
})
|
|
|
|
async function fixture() {
|
|
const journal = await journals.open({
|
|
identity: {
|
|
sessionId: SESSION,
|
|
workspaceId: 'folder-workspace',
|
|
hostId: 'local',
|
|
agent: 'codex',
|
|
providerHandle: { kind: 'codex', threadId: 'thread-1' }
|
|
},
|
|
stateDirectory: join(root, 'journal')
|
|
})
|
|
async function appendOutput(index: number) {
|
|
await journal.appendItem(
|
|
{ provider: 'orca', clientMessageId: `output-${index}` },
|
|
{ kind: 'status', text: `Tool output ${index}` },
|
|
{ fence: 1 }
|
|
)
|
|
}
|
|
for (let index = 1; index < 99; index += 1) {
|
|
await appendOutput(index)
|
|
}
|
|
await journal.appendSubmission({
|
|
clientMessageId: 'pending-send',
|
|
payloadFingerprint: 'prompt',
|
|
body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'Run tools' }] },
|
|
fence: 1
|
|
})
|
|
expect(journal.cursor().sequence).toBe(100)
|
|
const initial = structuredClone(
|
|
readAgentSessionHistory(journal, { sessionId: SESSION, direction: 'tail' })
|
|
)
|
|
const accept = () =>
|
|
journal.resolveDispatch({
|
|
clientMessageId: 'pending-send',
|
|
fence: 1,
|
|
state: 'accepted',
|
|
providerIdentity: { provider: 'codex', threadId: 'thread-1', turnId: 'turn-1', ordinal: 0 }
|
|
})
|
|
return { journal, initial, appendOutput, accept }
|
|
}
|
|
|
|
describe('structured session cursor/body regression', () => {
|
|
it('replaces retained pending submissions together with a real bounded snapshot at 140', async () => {
|
|
const { journal, initial, appendOutput, accept } = await fixture()
|
|
const retained = reduceStructuredAgentSession(EMPTY_STRUCTURED_AGENT_SESSION, {
|
|
type: 'event',
|
|
event: { type: 'snapshot', sessionId: SESSION, page: initial.page, fence: 1 }
|
|
})
|
|
expect(retained.submissions[0]?.dispatchState).toBe('pending')
|
|
await accept()
|
|
for (let index = 102; index <= 140; index += 1) {
|
|
await appendOutput(index)
|
|
}
|
|
const bounded = readAgentSessionHistory(journal, {
|
|
sessionId: SESSION,
|
|
direction: 'tail',
|
|
limit: 1
|
|
}).page
|
|
expect(bounded.liveCursor?.sequence).toBe(140)
|
|
expect(bounded.items).not.toContainEqual(retained.items.at(-1))
|
|
expect(bounded.submissions).toEqual([])
|
|
|
|
const replaced = reduceStructuredAgentSession(retained, {
|
|
type: 'event',
|
|
event: { type: 'snapshot', sessionId: SESSION, page: bounded, fence: 1 }
|
|
})
|
|
expect(replaced.cursor).toEqual(bounded.liveCursor)
|
|
expect(replaced.submissions).toEqual(bounded.submissions)
|
|
expect(hasUnansweredStructuredAgentSessionDispatch(replaced.submissions, 1)).toBe(false)
|
|
})
|
|
|
|
it.each([40, 401])(
|
|
'replays an off-page dispatch after %i missed rows without stranding pending state',
|
|
async (missedRows) => {
|
|
const { journal, appendOutput, accept } = await fixture()
|
|
let hostSummary: AgentSessionStatusSummary | undefined
|
|
const feed = new StructuredAgentSessionStatusFeed({
|
|
logger: createStructuredAgentSessionLogger(),
|
|
sessions: new Map([
|
|
[
|
|
SESSION,
|
|
{
|
|
journal,
|
|
fence: 1,
|
|
params: { location: { workspaceId: 'folder-workspace' }, provider: 'codex' }
|
|
}
|
|
]
|
|
]),
|
|
getRecord: () => null,
|
|
now: () => 1_000,
|
|
onStatusChanged: (summary) => {
|
|
hostSummary = summary
|
|
}
|
|
})
|
|
const subscribers = new AgentSessionSubscribers({
|
|
onJournalPublished: (sessionId, published) => feed.publish(sessionId, published)
|
|
})
|
|
const delayedOlder = Promise.withResolvers<AgentSessionHistoryResult>()
|
|
let warm = false
|
|
mocks.call.mockImplementation((_target, _method, request: AgentSessionHistoryRequest) => {
|
|
// Hold the measured bounded page before its asynchronous older-page fill can mask it.
|
|
if (warm && missedRows === 40 && request.direction === 'before') {
|
|
return delayedOlder.promise
|
|
}
|
|
const result = readAgentSessionHistory(journal, {
|
|
...request,
|
|
...(warm && missedRows === 40 ? { limit: 1 } : {})
|
|
})
|
|
return Promise.resolve(
|
|
structuredClone({
|
|
...result,
|
|
page: { ...result.page, fence: 1, hostNow: 1234 },
|
|
providerSession: { key: 'session_id', id: 'provider-1' }
|
|
})
|
|
)
|
|
})
|
|
mocks.subscribe.mockImplementation(
|
|
(
|
|
_target,
|
|
request: { cursor?: AgentJournalCursor },
|
|
onEvent: (event: AgentSessionSubscribeEvent) => void
|
|
) =>
|
|
Promise.resolve({
|
|
unsubscribe: subscribers.open({
|
|
id: 'pane',
|
|
sessionId: SESSION,
|
|
journal,
|
|
fence: 1,
|
|
cursor: request.cursor,
|
|
emit: (event) => onEvent(structuredClone(event))
|
|
})
|
|
})
|
|
)
|
|
const owner = getStructuredAgentSessionReadOwner(SESSION, target)
|
|
const unlisten = owner.subscribe(() => {})
|
|
const deactivate = owner.activate()
|
|
await vi.waitFor(() => expect(mocks.subscribe).toHaveBeenCalledTimes(1))
|
|
expect(owner.getSnapshot().state.cursor?.sequence).toBe(100)
|
|
expect(owner.getSnapshot().state.items.at(-1)?.body).toMatchObject({ role: 'user' })
|
|
expect(owner.getSnapshot().state.submissions[0]?.dispatchState).toBe('pending')
|
|
expect(owner.getSnapshot().providerSession).toEqual({ key: 'session_id', id: 'provider-1' })
|
|
expect(owner.getSnapshot().state.hostClock?.hostNow).toBe(1234)
|
|
expect(mocks.call).toHaveBeenCalledTimes(1)
|
|
deactivate()
|
|
|
|
await accept()
|
|
for (let index = 102; index <= 100 + missedRows; index += 1) {
|
|
await appendOutput(index)
|
|
}
|
|
const tail = readAgentSessionHistory(journal, {
|
|
sessionId: SESSION,
|
|
direction: 'tail',
|
|
limit: missedRows === 40 ? 1 : 200
|
|
}).page
|
|
expect(tail.liveCursor?.sequence).toBe(100 + missedRows)
|
|
expect(tail.submissions).toEqual([])
|
|
feed.publish(SESSION, journal)
|
|
// IPC/RPC copies values; the journal mutates its own submission records in place.
|
|
expect(owner.getSnapshot().state.submissions[0]?.dispatchState).toBe('pending')
|
|
warm = true
|
|
const stop = owner.activate()
|
|
|
|
await vi.waitFor(() => expect(owner.getSnapshot().state.cursor).toEqual(journal.cursor()))
|
|
const caughtUp = owner.getSnapshot().state
|
|
if (missedRows === 40) {
|
|
expect({
|
|
cursor: caughtUp.cursor?.sequence,
|
|
dispatch: caughtUp.submissions[0]?.dispatchState,
|
|
unansweredDispatch: hasUnansweredStructuredAgentSessionDispatch(caughtUp.submissions, 1)
|
|
}).toEqual({ cursor: 140, dispatch: 'accepted', unansweredDispatch: false })
|
|
}
|
|
|
|
await journal.appendItem(
|
|
{ provider: 'orca', clientMessageId: 'completed-turn' },
|
|
{ kind: 'turn', turnId: 'turn-1', state: 'completed' },
|
|
{ fence: 1 }
|
|
)
|
|
subscribers.publish(SESSION, journal)
|
|
await vi.waitFor(() => expect(owner.getSnapshot().state.cursor).toEqual(journal.cursor()))
|
|
const settled = owner.getSnapshot().state
|
|
expect(hostSummary?.status).toBe('idle')
|
|
expect(
|
|
projectStructuredAgentSessionStatus(settled.items, settled.submissions, settled.fence)
|
|
).toBe(hostSummary?.status)
|
|
expect(settled.submissions).toEqual(journal.snapshot().submissions)
|
|
expect({
|
|
cursor: caughtUp.cursor?.sequence,
|
|
dispatch: caughtUp.submissions[0]?.dispatchState,
|
|
unansweredDispatch: hasUnansweredStructuredAgentSessionDispatch(caughtUp.submissions, 1)
|
|
}).toEqual({ cursor: 100 + missedRows, dispatch: 'accepted', unansweredDispatch: false })
|
|
expect(mocks.call).toHaveBeenCalledTimes(1)
|
|
expect(mocks.subscribe).toHaveBeenCalledTimes(2)
|
|
expect(mocks.subscribe.mock.calls[1]?.[1]).toEqual({
|
|
sessionId: SESSION,
|
|
cursor: { epoch: journal.cursor().epoch, sequence: 100 }
|
|
})
|
|
|
|
stop()
|
|
unlisten()
|
|
}
|
|
)
|
|
})
|