diff --git a/.agents/notes/implemented/simplification/2026-07-17-one-send-one-turn.i18n.yaml b/.agents/notes/implemented/simplification/2026-07-17-one-send-one-turn.i18n.yaml index 7d841e0bf9..a0fd1ed4ec 100644 --- a/.agents/notes/implemented/simplification/2026-07-17-one-send-one-turn.i18n.yaml +++ b/.agents/notes/implemented/simplification/2026-07-17-one-send-one-turn.i18n.yaml @@ -3,4 +3,4 @@ # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write .agents/notes/implemented/simplification/2026-07-17-one-send-one-turn.md 2026-07-17-one-send-one-turn.md: 3ae43f137206f25bdbc563875c17e24211f17d6b -2026-07-17-one-send-one-turn.zh.md: 5ccdb2192048ecf795415bcd427f967df6a609fb +2026-07-17-one-send-one-turn.zh.md: 097090073f44194a8f8d4578a8cc39ffc723eb79 diff --git a/docs/core-data-structures/llm-streaming.i18n.yaml b/docs/core-data-structures/llm-streaming.i18n.yaml index c3924f184b..eca9de899d 100644 --- a/docs/core-data-structures/llm-streaming.i18n.yaml +++ b/docs/core-data-structures/llm-streaming.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write docs/core-data-structures/llm-streaming.md -llm-streaming.md: 6811611768a0ec577360a8ff82792899cf3b8fec -llm-streaming.zh.md: 35374af6a20086f15384840689dba5aa24351750 +llm-streaming.md: cdb2f76d9a47192cacedf545ae3d1ebbd985251e +llm-streaming.zh.md: 1cc7b06014224637d391d0a605315bff79b4a9fe diff --git a/docs/core-data-structures/llm-streaming.md b/docs/core-data-structures/llm-streaming.md index 7fa02b6e3b..cdb2f76d9a 100644 --- a/docs/core-data-structures/llm-streaming.md +++ b/docs/core-data-structures/llm-streaming.md @@ -144,7 +144,7 @@ declare class BlockAssembler { * Assemble all blocks seen so far, in stream order. * @returns one block per seen index, except that max-token truncation drops * tool calls that cannot be executed safely; an open block assembles from - * accumulated deltas (an unknown block type never closed by `block-end` throws). + * its accumulated deltas (an unknown block type never closed by `block-end` throws). */ blocks(): ContentBlock[]; /** Usage from the `usage` chunk; undefined until one arrives. */ diff --git a/docs/core-data-structures/llm-streaming.zh.md b/docs/core-data-structures/llm-streaming.zh.md index c25d140895..1cc7b06014 100644 --- a/docs/core-data-structures/llm-streaming.zh.md +++ b/docs/core-data-structures/llm-streaming.zh.md @@ -144,7 +144,7 @@ declare class BlockAssembler { * Assemble all blocks seen so far, in stream order. * @returns one block per seen index, except that max-token truncation drops * tool calls that cannot be executed safely; an open block assembles from - * accumulated deltas (an unknown block type never closed by `block-end` throws). + * its accumulated deltas (an unknown block type never closed by `block-end` throws). */ blocks(): ContentBlock[]; /** Usage from the `usage` chunk; undefined until one arrives. */ diff --git a/docs/module-graph.md b/docs/module-graph.md index 0a93add1ee..762c56443e 100644 --- a/docs/module-graph.md +++ b/docs/module-graph.md @@ -383,7 +383,6 @@ flowchart TD pkg_token_meter --> pkg_invariants pkg_token_meter --> pkg_llm pkg_token_meter --> pkg_session - pkg_agent --> pkg_brand pkg_agent --> pkg_invariants pkg_agent --> pkg_llm pkg_agent --> pkg_scope @@ -1059,7 +1058,7 @@ flowchart TD | [`lsp`](../packages/lsp/lsp) | `lsp` | [`brand`](../packages/util/brand), [`invariants`](../packages/support/invariants), [`llm`](../packages/llm/llm) | | [`sandbox`](../packages/sandbox/sandbox) | `sandbox` | [`invariants`](../packages/support/invariants), [`llm`](../packages/llm/llm) | | [`token-meter`](../packages/llm/token-meter) | `llm` | [`invariants`](../packages/support/invariants), [`llm`](../packages/llm/llm), [`session`](../packages/core/session) | -| [`agent`](../packages/core/agent) | `core` | [`brand`](../packages/util/brand), [`invariants`](../packages/support/invariants), [`llm`](../packages/llm/llm), [`scope`](../packages/core/scope), [`session`](../packages/core/session), [`system-prompt`](../packages/core/system-prompt) | +| [`agent`](../packages/core/agent) | `core` | [`invariants`](../packages/support/invariants), [`llm`](../packages/llm/llm), [`scope`](../packages/core/scope), [`session`](../packages/core/session), [`system-prompt`](../packages/core/system-prompt) | | [`bash`](../packages/bash/bash) | `bash` | [`invariants`](../packages/support/invariants), [`sandbox`](../packages/sandbox/sandbox), [`subprocess`](../packages/subprocess/subprocess) | | [`fs`](../packages/fs/fs) | `fs` | [`brand`](../packages/util/brand), [`invariants`](../packages/support/invariants), [`llm`](../packages/llm/llm), [`sandbox`](../packages/sandbox/sandbox) | | [`compact`](../packages/compact/compact) | `compact` | [`invariants`](../packages/support/invariants), [`llm`](../packages/llm/llm), [`session`](../packages/core/session) | diff --git a/packages/acp/acp/tests/turns.spec.ts b/packages/acp/acp/tests/turns.spec.ts index a4ae22080e..65305ba672 100644 --- a/packages/acp/acp/tests/turns.spec.ts +++ b/packages/acp/acp/tests/turns.spec.ts @@ -31,28 +31,28 @@ describe('ACP prompt lifecycle', () => { harness = undefined }) - it('maps a max-token turn without losing its committed text', async () => { + it('settles after a max-token turn without losing its committed text', async () => { harness = await makeBridgeHarness({ script: [maxTokensResponse('cut off')] }) const sessionId = await newSession(harness) const result = await harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] }) - expect(result.stopReason).toBe('max_tokens') + expect(result.stopReason).toBe('end_turn') await vi.waitFor(() => { expect(messageText(harness!)).toBe('cut off') }) }) - it('rejects a failed turn and never publishes its partial chunks', async () => { + it('settles after a failed turn and never publishes its partial chunks', async () => { harness = await makeBridgeHarness({ script: [errorResponse('provider boom')] }) const sessionId = await newSession(harness) await expect(harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })) - .rejects.toThrow(/turn failed: provider boom/) + .resolves.toEqual({ stopReason: 'end_turn' }) expect(messageText(harness)).toBe('') }) - it('rejects an ordinary plugin failure through the same prompt boundary', async () => { + it('settles after an ordinary plugin failure', async () => { harness = await makeBridgeHarness({ script: [textResponse('must not run')] }) harness.ctx.on('agent/step', () => { throw new Error('plugin pre-step failed') }) const sessionId = await newSession(harness) await expect(harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })) - .rejects.toThrow(/turn failed: plugin pre-step failed/) + .resolves.toEqual({ stopReason: 'end_turn' }) }) it('settles even when an earlier turn observer throws', async () => { @@ -202,27 +202,27 @@ describe('ACP prompt lifecycle', () => { await vi.waitFor(() => { expect(messageText(harness!)).toBe('recovered') }) }) - it('a failed turn with no retry still rejects, at quiescence', async () => { + it('a failed turn with no retry settles at quiescence', async () => { harness = await makeBridgeHarness({ script: [errorResponse('terminal boom')] }) let offered = 0 harness.ctx.on('agent/request-error', async () => { offered += 1 }) const sessionId = await newSession(harness) await expect(harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })) - .rejects.toThrow(/turn failed: terminal boom/) + .resolves.toEqual({ stopReason: 'end_turn' }) expect(offered).toBe(1) }) - it('an admission-blocked prompt settles cancelled instead of hanging', async () => { + it('an admission-blocked prompt settles instead of hanging', async () => { harness = await makeBridgeHarness({ script: [] }) harness.ctx.on('agent/prompt-submit', async () => ({ kind: 'block' as const, reason: 'policy said no' })) const sessionId = await newSession(harness) await expect(harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })) - .resolves.toEqual({ stopReason: 'cancelled' }) + .resolves.toEqual({ stopReason: 'end_turn' }) // The blocked prompt opened no turn and streamed nothing. expect(messageText(harness)).toBe('') }) - it('discards and settles a turnless prompt retained by its admission policy', async () => { + it('settles a turnless prompt retained by its admission policy', async () => { harness = await makeBridgeHarness({ script: [] }) harness.ctx.on('agent/prompt-submit', async () => ({ kind: 'block' as const, @@ -233,7 +233,7 @@ describe('ACP prompt lifecycle', () => { const agent = harness.ctx.agents.get(SessionId(sessionId))! await expect(harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })) - .resolves.toEqual({ stopReason: 'cancelled' }) + .resolves.toEqual({ stopReason: 'end_turn' }) expect(agent.status).toBe('idle') expect(agent.session.events.some(event => event.type === 'turn/start')).toBe(false) }) @@ -244,6 +244,6 @@ describe('ACP prompt lifecycle', () => { const sessionId = await newSession(harness) await expect(harness.client.prompt({ sessionId, prompt: [{ type: 'text', text: 'go' }] })) - .resolves.toEqual({ stopReason: 'cancelled' }) + .resolves.toEqual({ stopReason: 'end_turn' }) }) }) diff --git a/packages/core/agent-loop/README.i18n.yaml b/packages/core/agent-loop/README.i18n.yaml index 3a82c693ca..6b6b160e1f 100644 --- a/packages/core/agent-loop/README.i18n.yaml +++ b/packages/core/agent-loop/README.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write packages/core/agent-loop/README.md -README.md: a1617a1ef871f61157e0d70a06d055168170dced -README.zh.md: 6ba945a41e700331929dabb557802c14256921fb +README.md: c64af9bcea176f4d70b402ad6b84391ae15759d2 +README.zh.md: 6a878a1ffc31380466df99698a905ee0e38e24a4 diff --git a/packages/core/agent-loop/tests/loop.spec.ts b/packages/core/agent-loop/tests/loop.spec.ts index 9ccde685d0..2d856dd3ee 100644 --- a/packages/core/agent-loop/tests/loop.spec.ts +++ b/packages/core/agent-loop/tests/loop.spec.ts @@ -308,7 +308,6 @@ describe('agent loop', () => { send(agent, 'start') await waitForIdle(ctx, agent) - const types = agent.session.events.map(e => e.type) const steering = agent.session.events.find(e => e.type === 'user/message' && JSON.stringify(e.data.content).includes('change of plans')) expect(steering).toBeDefined() diff --git a/packages/core/agent-loop/tests/properties.spec.ts b/packages/core/agent-loop/tests/properties.spec.ts index a75baf9203..c5dec8c142 100644 --- a/packages/core/agent-loop/tests/properties.spec.ts +++ b/packages/core/agent-loop/tests/properties.spec.ts @@ -78,7 +78,7 @@ function userMessageTexts(agent: Agent): string[] { function turnNumbers(agent: Agent): number[] { return agent.session.events .filter(e => e.type === 'turn/start') - .map(e => (e.data as { turn: number }).turn) + .map(e => e.data.turn) } function turnEndNumbers(agent: Agent): number[] { diff --git a/packages/core/agent-loop/tests/request-error.spec.ts b/packages/core/agent-loop/tests/request-error.spec.ts index a7b35bfdbf..96b6bfc045 100644 --- a/packages/core/agent-loop/tests/request-error.spec.ts +++ b/packages/core/agent-loop/tests/request-error.spec.ts @@ -67,10 +67,6 @@ describe('agent/request-error', () => { }) ctx.on('agent/request-error', async (subject, context) => { expect(subject).toBe(agent) - expect(agent.session.events.at(-1)).toMatchObject({ - type: 'step/end', - data: { turn: context.turn, step: context.step }, - }) seen.push(context) return { kind: 'retry' } }) @@ -89,7 +85,7 @@ describe('agent/request-error', () => { code: 'RATE_LIMIT', }, { - turn: 2, + turn: 1, step: 1, code: 'SERVICE_UNAVAILABLE', }, diff --git a/packages/core/agent/package.json b/packages/core/agent/package.json index db0637ecc2..58495e1a63 100644 --- a/packages/core/agent/package.json +++ b/packages/core/agent/package.json @@ -21,7 +21,6 @@ "files": [ "lib/index.js", "lib/invariant.js", - "lib/types/**/*.js", "lib/types/**/*.d.ts", "lib/types/**/*.d.ts.map", "src" diff --git a/packages/core/session/tests/fork.spec.ts b/packages/core/session/tests/fork.spec.ts index 94e2bd74ce..8a3a568762 100644 --- a/packages/core/session/tests/fork.spec.ts +++ b/packages/core/session/tests/fork.spec.ts @@ -139,17 +139,17 @@ describe('SessionStore.fork', () => { const reasons: TurnEndReason[] = [ { kind: 'completed' }, { kind: 'aborted', reason: { kind: 'user' } }, - { kind: 'error', error: new Error('model failed') }, + { kind: 'error', error: 'model failed' }, { kind: 'aborted', reason: { kind: 'disposed' } }, { kind: 'max-tokens' }, { kind: 'interrupted' }, ] - for (const reason of reasons) { - const source = ctx.sessions.create(SessionId(`parent-${reason.kind}`)) + for (const [index, reason] of reasons.entries()) { + const source = ctx.sessions.create(SessionId(`parent-${index}`)) appendClosedTurn(source, 1, reason.kind, reason) - const child = sessions.fork(source, lastSeq(source), SessionId(`child-${reason.kind}`)) + const child = sessions.fork(source, lastSeq(source), SessionId(`child-${index}`)) expect(inherited(child).at(-1)?.type).toBe('turn/end') expect(child.header.seedLength).toBe(source.events.length) diff --git a/packages/core/session/tests/invariant.spec.ts b/packages/core/session/tests/invariant.spec.ts index bd7440eaa2..3cf089fe5d 100644 --- a/packages/core/session/tests/invariant.spec.ts +++ b/packages/core/session/tests/invariant.spec.ts @@ -326,7 +326,7 @@ describe('session-log invariants', () => { unresolved.append('step/start', { turn: 1, step: 1 }) unresolved.append('tool/call', { turn: 1, step: 1, callId: CallId('c1'), name: 'echo', arguments: '{}' }) unresolved.append('step/end', { turn: 1, step: 1 }) - unresolved.append('turn/end', { turn: 1, reason: { kind: 'error', error: new Error('boom') } }) + unresolved.append('turn/end', { turn: 1, reason: { kind: 'error', error: 'boom' } }) }).not.toThrow() }) @@ -384,16 +384,16 @@ describe('session-log invariants', () => { const { ctx } = await setup() // Balanced seed: between turns. expect(() => ctx.sessions.create(SessionId('inherited-between-turns'), { seed: [ - { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } }, + { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }, { type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'completed' } } }, ] })).not.toThrow() // Unbalanced seed: inside the open turn, which the relation permits. const open = ctx.sessions.create(SessionId('inherited-inside-open-turn'), { seed: [ - { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } }, + { type: 'turn/start', seq: 0, time: 1, data: { turn: 1 } }, ] }) expect(open.events.map(event => event.type)).toEqual(['turn/start', 'session/end-seed']) // Still open afterwards: the boundary moves no cursor. - expect(() => open.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } })) + expect(() => open.append('turn/start', { turn: 2 })) .toThrow(/turn 1 is still open/) expect(() => open.append('turn/end', { turn: 1, reason: { kind: 'completed' } })).not.toThrow() }) diff --git a/packages/examples/cli-demo/src/cli.ts b/packages/examples/cli-demo/src/cli.ts index f40cd03f9a..68eeadae0a 100644 --- a/packages/examples/cli-demo/src/cli.ts +++ b/packages/examples/cli-demo/src/cli.ts @@ -226,7 +226,9 @@ export async function runOneShot(ctx: Context, options: OneShotOptions): Promise options.onEvent(sessionId, event) } catch (error: unknown) { outputError = toError(error) - agent.cancel({ kind: 'user' }) + queueMicrotask(() => { + agent.cancel({ kind: 'user' }) + }) } } diff --git a/packages/examples/cli-demo/tests/cli.spec.ts b/packages/examples/cli-demo/tests/cli.spec.ts index 3a7c29fadb..ff1916fe42 100644 --- a/packages/examples/cli-demo/tests/cli.spec.ts +++ b/packages/examples/cli-demo/tests/cli.spec.ts @@ -354,7 +354,7 @@ describe('runOneShot and executeCli', () => { }) }) - it('counts a failed retry attempt once even though it has no assistant message', async () => { + it('reports usage committed by the recovered assistant message', async () => { const failed = { inputTokens: 11, outputTokens: 2, cacheReadTokens: 3 } const recovered = { inputTokens: 7, outputTokens: 5, reasoningTokens: 4 } const { ctx } = await harness([failedResponse(failed), textResponse('done', recovered)]) @@ -362,9 +362,8 @@ describe('runOneShot and executeCli', () => { const result = await runOneShot(ctx, { task: 'task' }) expect(result.usage).toEqual({ - inputTokens: 18, - outputTokens: 7, - cacheReadTokens: 3, + inputTokens: 7, + outputTokens: 5, reasoningTokens: 4, }) }) @@ -422,7 +421,8 @@ describe('runOneShot and executeCli', () => { const outcome = await result expect(outcome).toMatchObject({ type: 'result', output: 'streamed' }) const events = streamed.map(item => item.event) - expect(events[0]).toMatchObject({ type: 'turn/start', data: { turn: 3 } }) + expect(events.find(event => event.type === 'turn/start')) + .toMatchObject({ type: 'turn/start', data: { turn: 3 } }) expect(events.at(-1)).toMatchObject({ type: 'turn/end', data: { turn: 3 } }) expect(streamed.every(item => item.sessionId === agent.session.id)).toBe(true) expect(events.some(event => event.type === 'user/message' @@ -462,7 +462,7 @@ describe('runOneShot and executeCli', () => { const failed = await harness([]) failed.ctx.on('agent/prompt-submit', async () => { throw new Error('admission exploded') }) - await expect(runOneShot(failed.ctx, { task: 'task' })).rejects.toThrow('not admitted') + await expect(runOneShot(failed.ctx, { task: 'task' })).resolves.toMatchObject({ output: '' }) }) it('emits partial data without attributing a turn outcome', async () => { diff --git a/packages/host/apiproxy/tests/api-proxy-cold.spec.ts b/packages/host/apiproxy/tests/api-proxy-cold.spec.ts index b9d5018975..0e82805ab3 100644 --- a/packages/host/apiproxy/tests/api-proxy-cold.spec.ts +++ b/packages/host/apiproxy/tests/api-proxy-cold.spec.ts @@ -90,7 +90,7 @@ describe('attached updatedAt excludes end-seed', () => { const worked = 1_000_000 const resumed = ctx.sessions.create(sid('resumed-untouched'), { seed: [ - { type: 'turn/start', seq: 0, time: worked, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } }, + { type: 'turn/start', seq: 0, time: worked, data: { turn: 1 } }, { type: 'turn/end', seq: 1, time: worked, data: { turn: 1, reason: { kind: 'completed' } } }, ], meta: { cwd: '/proj', createdAt: 500 }, @@ -106,7 +106,7 @@ describe('attached updatedAt excludes end-seed', () => { expect(summary?.updatedAt).toBe(worked) // Real work appended after end-seed does move it. - resumed.append('turn/start', { turn: 2, trigger: { kind: 'message', source: { kind: 'user' } } }) + resumed.append('turn/start', { turn: 2 }) const after = await api.sessions.list(request({})) if (!after.result.ok) throw new Error('list failed') const moved = after.result.value.items.find(item => item.sessionId === 'resumed-untouched') diff --git a/packages/llm/llm-retry/tests/retry.spec.ts b/packages/llm/llm-retry/tests/retry.spec.ts index f011fb3e5d..27f936d936 100644 --- a/packages/llm/llm-retry/tests/retry.spec.ts +++ b/packages/llm/llm-retry/tests/retry.spec.ts @@ -144,15 +144,8 @@ function alwaysConfig(backoff: BackoffConfig = {}): AlwaysRetryPolicyConfig { } } -function waitForIdle(ctx: Context, agent: Agent): Promise { - return new Promise((resolve) => { - const dispose = ctx.on('agent/status', (subject, status) => { - if (subject === agent && status === 'idle') { - dispose() - resolve() - } - }) - }) +function waitForIdle(_ctx: Context, agent: Agent): Promise { + return agent.whenIdle() } function waitForRetry(ctx: Context, agent: Agent, retryNumber: number): Promise> { @@ -175,7 +168,7 @@ afterEach(async () => { }) describe('provider-routed retry policy', () => { - it('records the scheduled delay before opening a fresh request attempt', async () => { + it('records the scheduled delay before retrying the request', async () => { vi.useFakeTimers() const adapter = new ScriptedAdapter([ new LlmError('busy', 'RATE_LIMIT', { status: 429 }), @@ -214,7 +207,7 @@ describe('provider-routed retry policy', () => { expect(adapter.requests).toHaveLength(2) expect(agent.session.events.filter(item => item.type === 'step/start').map(item => item.data)) - .toEqual([{ turn: 1, step: 1 }, { turn: 2, step: 1 }]) + .toEqual([{ turn: 1, step: 1 }]) expect(agent.session.deriveMessages().at(-1)).toEqual({ id: expect.any(String) as unknown, role: 'assistant', @@ -251,7 +244,7 @@ describe('provider-routed retry policy', () => { expect(agent.session.events.filter(event => event.type === 'assistant/message').map(event => ({ turn: event.data.turn, step: event.data.step, - }))).toEqual([{ turn: 2, step: 1 }]) + }))).toEqual([{ turn: 1, step: 1 }]) expect(agent.session.deriveMessages().at(-1)).toMatchObject({ role: 'assistant', content: [{ type: 'text', text: 'recovered' }], @@ -291,7 +284,7 @@ describe('provider-routed retry policy', () => { expect(agent.session.events.filter(event => event.type === 'assistant/message').map(event => ({ turn: event.data.turn, step: event.data.step, - }))).toEqual([{ turn: 2, step: 1 }]) + }))).toEqual([{ turn: 1, step: 1 }]) expect(agent.session.events.some(event => event.type === 'tool/call')).toBe(false) expect(toolExecutions).toBe(0) expect(agent.session.deriveMessages().at(-1)).toMatchObject({ @@ -534,9 +527,9 @@ describe('provider-routed retry policy', () => { backoff: { initialDelayMs: 1, maxDelayMs: 1 }, }), }, (ctx) => { - ctx.on('agent/request', async (_agent, turn, _step, _signal, next) => ({ + ctx.on('agent/request', async (_agent, _turn, _step, _signal, next) => ({ ...await next(), - provider: turn === 1 ? 'mock' : 'other', + provider: adapter.requests.length === 0 ? 'mock' : 'other', })) })) const agent = context.agentLoop.create(SessionId('retry-provider-budgets'), { diff --git a/packages/llm/llm-retry/tests/transport-recovery.spec.ts b/packages/llm/llm-retry/tests/transport-recovery.spec.ts index f9de22a120..90c210f7bb 100644 --- a/packages/llm/llm-retry/tests/transport-recovery.spec.ts +++ b/packages/llm/llm-retry/tests/transport-recovery.spec.ts @@ -55,14 +55,8 @@ async function harness( return ctx } -function waitForIdle(ctx: Context, agent: Agent): Promise { - return new Promise((resolve) => { - const dispose = ctx.on('agent/status', (subject, status) => { - if (subject !== agent || status !== 'idle') return - dispose() - resolve() - }) - }) +function waitForIdle(_ctx: Context, agent: Agent): Promise { + return agent.whenIdle() } function sendAndWait(ctx: Context, agent: Agent): Promise { @@ -109,7 +103,7 @@ describe('bounded retry through the real DeepSeek HTTP/SSE adapter', () => { expect(server?.requests).toHaveLength(1) expect(agent.session.events.filter(event => event.type === 'step/start') .map(event => [event.data.turn, event.data.step])) - .toEqual([[1, 1], [2, 1]]) + .toEqual([[1, 1]]) expect(agent.session.events.filter(event => event.type === 'llm/retry').map(event => event.data.failure.code)) .toEqual(['TRANSPORT']) expect(finalAssistantText(agent)).toBe('connected after retry') @@ -141,7 +135,7 @@ describe('bounded retry through the real DeepSeek HTTP/SSE adapter', () => { )).toHaveLength(failedChunkCount) expect(agent.session.events.filter(event => event.type === 'assistant/message') .map(event => [event.data.turn, event.data.step])) - .toEqual([[2, 1]]) + .toEqual([[1, 1]]) expect(agent.session.events.filter(event => event.type === 'llm/retry').map(event => event.data.failure.code)) .toEqual(['TRANSPORT']) expect(finalAssistantText(agent)).toBe('recovered response') @@ -166,7 +160,7 @@ describe('bounded retry through the real DeepSeek HTTP/SSE adapter', () => { .toEqual(['EMPTY_RESPONSE']) expect(agent.session.events.filter(event => event.type === 'assistant/message') .map(event => [event.data.turn, event.data.step])) - .toEqual([[2, 1]]) + .toEqual([[1, 1]]) expect(agent.session.events.at(-1)).toMatchObject({ type: 'turn/end', data: { reason: { kind: 'completed' } }, @@ -234,7 +228,7 @@ describe('bounded retry through the real DeepSeek HTTP/SSE adapter', () => { await sendAndWait(context, agent) expect(server.requests).toHaveLength(3) - expect(agent.session.events.filter(event => event.type === 'step/start')).toHaveLength(3) + expect(agent.session.events.filter(event => event.type === 'step/start')).toHaveLength(1) expect(agent.session.events.filter(event => event.type === 'llm/retry')).toHaveLength(2) expect(agent.session.events.at(-1)).toMatchObject({ type: 'turn/end', diff --git a/packages/llm/llm/README.i18n.yaml b/packages/llm/llm/README.i18n.yaml index 49ff6d7c48..15e79631a9 100644 --- a/packages/llm/llm/README.i18n.yaml +++ b/packages/llm/llm/README.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write packages/llm/llm/README.md -README.md: d343449d1530bf70a3a8c57f883894e29c42d18f -README.zh.md: 4dc4a0ca06378116d05fdb4b9b048738930511fd +README.md: dc7499a6854fe9a45c1297aa2a1a67aea92eaf6f +README.zh.md: 28f281908c2fde708f491194b441a32a8dcba8a5 diff --git a/packages/llm/llm/tests/service.spec.ts b/packages/llm/llm/tests/service.spec.ts index af87442744..95142aade2 100644 --- a/packages/llm/llm/tests/service.spec.ts +++ b/packages/llm/llm/tests/service.spec.ts @@ -262,13 +262,16 @@ describe('LlmService', () => { messages: [], })) - expect(chunks.at(-1)).toMatchObject({ + const finish = chunks.at(-1) + expect(finish).toMatchObject({ type: 'finish', reason: { kind: 'error', - failure: { code: 'NO_ADAPTER', message: expect.stringContaining('no adapter registered') }, + failure: { code: 'NO_ADAPTER' }, }, }) + if (finish?.type !== 'finish' || finish.reason.kind !== 'error') throw new Error('expected error finish') + expect(finish.reason.failure.message).toContain('no adapter registered') }) it.each(['done', 'value'] as const)('normalizes a throwing IteratorResult.%s getter', async (field) => { diff --git a/packages/sdk/sdk-client/tests/sdk-client.spec.ts b/packages/sdk/sdk-client/tests/sdk-client.spec.ts index ca2d30b3cc..3717e60357 100644 --- a/packages/sdk/sdk-client/tests/sdk-client.spec.ts +++ b/packages/sdk/sdk-client/tests/sdk-client.spec.ts @@ -175,7 +175,7 @@ describe('DeepSeekHarness', () => { await using harness = new DeepSeekHarness({ launch: fakeLaunch() }) captured = harness const result = await harness.run('scoped') - expect(result.finalResponse).toBe('scoped') + expect(result.finalResponse).toBe('hello from fake runtime') } // After scope exit the runtime is closed: reuse fails loudly. await expect(captured.run('after')).rejects.toThrow(TransportClosedError) diff --git a/packages/telemetry/session-telemetry-otel/tests/otel.spec.ts b/packages/telemetry/session-telemetry-otel/tests/otel.spec.ts index fa2b67db2b..1b720c194d 100644 --- a/packages/telemetry/session-telemetry-otel/tests/otel.spec.ts +++ b/packages/telemetry/session-telemetry-otel/tests/otel.spec.ts @@ -94,7 +94,7 @@ describe('TelemetryOtel wire', () => { const { ctx, fiber } = await boot(url) const session = ctx.sessions.create(SessionId('wire'), { meta: { cwd: '/tmp/w' } }) session.append('turn/start', { turn: 1 }) - session.append('turn/end', { turn: 1, reason: { kind: 'error', error: new Error('boom') } }) + session.append('turn/end', { turn: 1, reason: { kind: 'error', error: 'boom' } }) await fiber.dispose() expect(captures.length).toBeGreaterThan(0) diff --git a/packages/telemetry/session-telemetry/tests/telemetry.spec.ts b/packages/telemetry/session-telemetry/tests/telemetry.spec.ts index 88a42a5191..410ea2d379 100644 --- a/packages/telemetry/session-telemetry/tests/telemetry.spec.ts +++ b/packages/telemetry/session-telemetry/tests/telemetry.spec.ts @@ -126,7 +126,7 @@ describe('TelemetryCoordinator capture', () => { }), }, { surfaceOp: 'append' }) session.append('telemetry-test/opaque', { payload: { nested: [] } }) - session.append('turn/end', { turn: 1, reason: { kind: 'error', error: new Error('boom') } }) + session.append('turn/end', { turn: 1, reason: { kind: 'error', error: 'boom' } }) const severities = backend.ledger().map(r => [r.attributes['event.type'], r.severity]) expect(severities).toEqual([ ['turn/start', 'info'], diff --git a/packages/ui/tui/tests/tui.spec.ts b/packages/ui/tui/tests/tui.spec.ts index 458f3d57e7..25f149a980 100644 --- a/packages/ui/tui/tests/tui.spec.ts +++ b/packages/ui/tui/tests/tui.spec.ts @@ -452,7 +452,7 @@ describe('goodbye message and /resume', () => { it.each([ [{ kind: 'aborted', reason: { kind: 'user' } }, 'cancelled'], - [{ kind: 'error', error: new Error('failed') }, 'error'], + [{ kind: 'error', error: 'failed' }, 'error'], [{ kind: 'aborted', reason: { kind: 'disposed' } }, 'cancelled'], [{ kind: 'max-tokens' }, 'max tokens'], [{ kind: 'interrupted' }, 'interrupted'], @@ -3687,6 +3687,9 @@ describe('pi-tui chat lifecycle and transcript', () => { result.terminal.send('\r') await tick() expect(result.agent.cancelled).toContainEqual({ kind: 'user' }) + result.agent.status = 'idle' + agentEvents(result.ctx, result.agent).emit('agent/status', 'idle') + await tick() expect(result.exit).toHaveBeenCalledWith(0) const events = await setup() @@ -3714,7 +3717,7 @@ describe('pi-tui chat lifecycle and transcript', () => { events.session.append('turn/start', { turn: 6 }) events.session.append('turn/end', { turn: 6, - reason: { kind: 'error', error: { message: 'structured provider failure', code: 'SERVER' } }, + reason: { kind: 'error', error: 'structured provider failure' }, }) events.session.append('turn/start', { turn: 8 }) events.session.append('turn/end', { @@ -3732,7 +3735,6 @@ describe('pi-tui chat lifecycle and transcript', () => { expect(events.terminal.output).toContain('structured provider failure') expect(events.terminal.output).toContain('output-token limit') expect(events.terminal.output).toContain('previous process ended') - expect(events.terminal.output).toContain('Turn stopped: the agent was disposed') expect(events.terminal.output).toContain('Turn ended: plugin-policy') expect(events.terminal.output).toContain('was disposed') await dispose(events) diff --git a/python/sdk/tests/manual_sdk_agent_smoke.py b/python/sdk/tests/manual_sdk_agent_smoke.py index 751b7fc0bf..ae94f1bb07 100644 --- a/python/sdk/tests/manual_sdk_agent_smoke.py +++ b/python/sdk/tests/manual_sdk_agent_smoke.py @@ -73,9 +73,7 @@ def run_smoke(repo_root: Path, keep_sessions: bool) -> None: "Please reply with a short confirmation and do not call tools.", session_id="sdk-smoke-main", ) - print(f"turn_status={result.status}") print(f"final_response={result.final_response}") - assert result.status == "ok", result assert "configured HTTP model endpoint" in result.final_response assert len(MockCompletionHandler.requests) == 1 request = MockCompletionHandler.requests[0] diff --git a/python/sdk/tests/test_client.py b/python/sdk/tests/test_client.py index 52ceac9f4d..f16f5f27e9 100644 --- a/python/sdk/tests/test_client.py +++ b/python/sdk/tests/test_client.py @@ -39,6 +39,9 @@ for line in sys.stdin: print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True) elif method == "session/prompt": params = msg.get("params") or {} + print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": params["sessionId"], "event": {"type": "agent/inbox/spliced", "data": {"target": "next-turn", "start": 0, "inserted": [{"id": "message-1"}]}}}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": params["sessionId"], "status": "running"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) print(json.dumps({ "jsonrpc": "2.0", "method": "session.event", @@ -57,10 +60,9 @@ for line in sys.stdin: }), flush=True) print(json.dumps({ "jsonrpc": "2.0", - "method": "session.finished", - "params": {"sessionId": params["sessionId"], "status": "ok"}, + "method": "session.status", + "params": {"sessionId": params["sessionId"], "status": "idle"}, }), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -83,9 +85,8 @@ for line in sys.stdin: ) as harness: result = harness.run("say hello", session_id="main") - assert result.status == "ok" assert result.final_response == "hello from runtime" - assert result.events[0]["type"] == "assistant/message" + assert result.events[-1]["type"] == "assistant/message" dumped_env = json.loads(env_dump.read_text()) assert dumped_env["DEEPSEEK_API_KEY"] == "env-key" assert dumped_env["DEEPSEEK_BASE_URL"] == "http://127.0.0.1:4321" @@ -113,9 +114,11 @@ for line in sys.stdin: if method == "initialize": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True) elif method == "session/prompt": + print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": "main", "event": {"type": "agent/inbox/spliced", "data": {"target": "next-turn", "start": 0, "inserted": [{"id": "message-1"}]}}}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": "main", "status": "running"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": "main", "childSessionId": "child"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": "main", "status": "ok"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": "main", "status": "idle"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -133,8 +136,7 @@ for line in sys.stdin: on_notification=lambda notification: seen.append(notification.method), ) - assert result.status == "ok" - assert seen == ["subagent.started", "session.finished"] + assert seen == ["session.event", "session.status", "subagent.started", "session.status"] def test_relative_cwd_is_absolute_in_process_environment_and_wire( @@ -189,10 +191,12 @@ for line in sys.stdin: if method == "initialize": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True) elif method == "session/prompt": + print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": "main", "event": {"type": "agent/inbox/spliced", "data": {"target": "next-turn", "start": 0, "inserted": [{"id": "message-1"}]}}}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": "main", "status": "running"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": "main", "childSessionId": "child"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "subagent.finished", "params": {"parentSessionId": "main", "childSessionId": "child", "status": "ok", "stopReason": "completed"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": "main", "status": "ok"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": "main", "status": "idle"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -205,11 +209,12 @@ for line in sys.stdin: ) as harness: result = harness.run("spawn a helper", session_id="main") - assert result.status == "ok" assert [notification.method for notification in result.notifications] == [ + "session.event", + "session.status", "subagent.started", "subagent.finished", - "session.finished", + "session.status", ] @@ -229,6 +234,9 @@ for line in sys.stdin: print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True) elif method == "session/prompt": root = (msg.get("params") or {})["sessionId"] + print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": root, "event": {"type": "agent/inbox/spliced", "data": {"target": "next-turn", "start": 0, "inserted": [{"id": "message-1"}]}}}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": root, "status": "running"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": root, "childSessionId": "child"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": "child", "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "child response"}]}}}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "subagent.started", "params": {"parentSessionId": "child", "childSessionId": "grandchild"}}), flush=True) @@ -236,8 +244,7 @@ for line in sys.stdin: print(json.dumps({"jsonrpc": "2.0", "method": "subagent.finished", "params": {"parentSessionId": "child", "childSessionId": "grandchild", "status": "ok"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "subagent.finished", "params": {"parentSessionId": root, "childSessionId": "child", "status": "ok"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": root, "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "root response"}]}}}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": root, "status": "ok"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": root, "status": "idle"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -256,10 +263,11 @@ for line in sys.stdin: ) assert harness.client._notifications.qsize() == 0 - assert result.status == "ok" assert result.final_response == "root response" - assert [event["data"]["content"][0]["text"] for event in result.events] == ["root response"] + assert [event["data"]["content"][0]["text"] for event in result.events if event["type"] == "assistant/message"] == ["root response"] assert [notification.method for notification in result.notifications] == [ + "session.event", + "session.status", "subagent.started", "session.event", "subagent.started", @@ -267,7 +275,7 @@ for line in sys.stdin: "subagent.finished", "subagent.finished", "session.event", - "session.finished", + "session.status", ] assert seen == [notification.method for notification in result.notifications] @@ -287,10 +295,12 @@ for line in sys.stdin: elif method == "session/prompt": params = msg.get("params") or {} print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": "other", "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "wrong session"}]}}}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": "other", "status": "ok"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": "other", "status": "idle"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": params["sessionId"], "event": {"type": "agent/inbox/spliced", "data": {"target": "next-turn", "start": 0, "inserted": [{"id": "message-1"}]}}}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": params["sessionId"], "status": "running"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": params["sessionId"], "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "right session"}]}}}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": params["sessionId"], "status": "ok"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": params["sessionId"], "status": "idle"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -303,9 +313,8 @@ for line in sys.stdin: ) as harness: result = harness.run("stay in your lane", session_id="main") - assert result.status == "ok" assert result.final_response == "right session" - assert [notification.payload.get("sessionId") for notification in result.notifications] == ["main", "main"] + assert [notification.payload.get("sessionId") for notification in result.notifications] == ["main"] * 4 def test_high_level_session_run_does_not_accumulate_global_notifications(tmp_path: Path) -> None: @@ -322,9 +331,11 @@ for line in sys.stdin: print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"serverInfo": {"name": "fake-runtime"}}}), flush=True) elif method == "session/prompt": params = msg.get("params") or {} + print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": params["sessionId"], "event": {"type": "agent/inbox/spliced", "data": {"target": "next-turn", "start": 0, "inserted": [{"id": "message-1"}]}}}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": params["sessionId"], "status": "running"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": params["sessionId"], "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "ok"}]}}}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": params["sessionId"], "status": "ok"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": params["sessionId"], "status": "idle"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -333,11 +344,10 @@ for line in sys.stdin: with DeepSeekHarness(launch_args_override=(sys.executable, str(script)), cwd=str(tmp_path)) as harness: result = harness.run("one turn", session_id="main") - assert result.status == "ok" assert harness.client._notifications.qsize() == 0 -def test_session_run_waits_for_late_finished_without_replaying_stale_notifications(tmp_path: Path) -> None: +def test_session_run_waits_for_late_idle_without_replaying_stale_notifications(tmp_path: Path) -> None: script = tmp_path / "fake_runtime.py" script.write_text( """ @@ -355,15 +365,17 @@ for line in sys.stdin: turn += 1 params = msg.get("params") or {} session_id = params["sessionId"] + message_id = f"message-{turn}" + print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": session_id, "event": {"type": "agent/inbox/spliced", "data": {"target": "next-turn", "start": 0, "inserted": [{"id": message_id}]}}}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": session_id, "status": "running"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": message_id}}), flush=True) if turn == 1: print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": session_id, "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "first"}]}}}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": session_id, "status": "ok"}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": session_id, "status": "idle"}}), flush=True) else: - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) time.sleep(0.05) print(json.dumps({"jsonrpc": "2.0", "method": "session.event", "params": {"sessionId": session_id, "event": {"type": "assistant/message", "data": {"content": [{"type": "text", "text": "second"}]}}}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "method": "session.finished", "params": {"sessionId": session_id, "status": "ok"}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "method": "session.status", "params": {"sessionId": session_id, "status": "idle"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -376,7 +388,7 @@ for line in sys.stdin: assert first.final_response == "first" assert second.final_response == "second" - assert [notification.payload.get("sessionId") for notification in second.notifications] == ["main", "main"] + assert [notification.payload.get("sessionId") for notification in second.notifications] == ["main"] * 4 def test_client_starts_subprocess_sends_requests_and_routes_notifications(tmp_path: Path) -> None: @@ -394,7 +406,7 @@ for line in sys.stdin: elif method == "session/prompt": params = msg.get("params") or {} print(json.dumps({"jsonrpc": "2.0", "method": "llm/request", "params": {"requestId": "req-1", "sessionId": params["sessionId"], "model": "dsagent", "messages": []}}), flush=True) - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break @@ -531,7 +543,7 @@ for line in sys.stdin: elif method in {"emit-first", "emit-second"}: print(json.dumps({"jsonrpc": "2.0", "method": "tick", "params": {"source": method}}), flush=True) elif method == "session/prompt": - print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"accepted": True}}), flush=True) + print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {"messageId": "message-1"}}), flush=True) elif method == "shutdown": print(json.dumps({"jsonrpc": "2.0", "id": msg["id"], "result": {}}), flush=True) break