Merge branch 'claude/web-llm-pi-ai-config-385e24' into claude/pi-ai-model-discovery

# Conflicts:
#	docs/cordis-catalog/events.md
#	docs/core-data-structures/core.i18n.yaml
#	docs/event-producer-consumer.md
#	packages/host/apiproxy/README.i18n.yaml
#	packages/llm/llm/README.i18n.yaml
This commit is contained in:
Yichen Jiang
2026-08-06 10:50:20 +08:00
822 changed files with 18341 additions and 14753 deletions

View File

@@ -34,7 +34,7 @@ async function harness(): Promise<{ ctx: Context; api: ApiProxy }> {
/** A minimal agent stand-in inside an open turn (the service only reaches `.session`). */
function agentOf(ctx: Context): Agent {
const session = ctx.sessions.create()
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
return { session } as unknown as Agent
}
@@ -185,7 +185,7 @@ describe('approval pending registry', () => {
const abort = new AbortController()
const mux = openMux(api, abort)
const session = ctx.sessions.create()
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
session.append('approval/asked', { id: 'pre-aborted' as ApprovalRequestId, toolName: 'bash' })
const agent = { session } as unknown as Agent
const cancelled = new AbortController()
@@ -308,7 +308,7 @@ describe('approval pending registry', () => {
// Bypass ApprovalService: a log whose sole asked event already has its
// decided partner must not be re-claimed — the answerer delegates.
const session = ctx.sessions.create()
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
session.append('approval/asked', { id: 'stale-ask' as ApprovalRequestId, toolName: 'bash' })
session.append('approval/decided', { id: 'stale-ask' as ApprovalRequestId, outcome: 'rejected' })
const agent = { session } as unknown as Agent
@@ -322,7 +322,7 @@ describe('approval pending registry', () => {
// Bypass ApprovalService: dispatch the waterfall directly with a session
// that has no approval/asked event — the proxy answerer must call next().
const session = ctx.sessions.create()
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
const agent = { session } as unknown as Agent
const outcome = await ctx.waterfall('approval/request', { agent, toolName: 'x' }, () => Promise.resolve('unavailable' as const))
expect(outcome).toBe('unavailable')

View File

@@ -79,7 +79,7 @@ describe('summary blank = conversation not started', () => {
const session = ctx.sessions.create()
attach(session)
appendStandalone(session)
session.append('turn/start', { turn: 0, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 0 })
expect(await listBlank(api, session.id)).toBe(false)
})
})

View File

@@ -10,7 +10,8 @@ import { join } from 'node:path'
import { describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import SessionStore from '@deepseek-ai/dsh-session'
import AgentRegistry, { InboxItemId } from '@deepseek-ai/dsh-agent'
import AgentRegistry from '@deepseek-ai/dsh-agent'
import { MessageId } from '@deepseek-ai/dsh-llm'
import type { Agent } from '@deepseek-ai/dsh-agent'
import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
import type { SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
@@ -89,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 },
@@ -105,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')
@@ -216,7 +217,7 @@ describe('subagent ownership fence', () => {
const queued = await api.sessions.updateQueue(request({
sessionId: originChild.id,
itemId: InboxItemId('queued-item'),
itemId: MessageId('queued-item'),
action: { kind: 'remove' },
}))
expect(queued.result.ok).toBe(false)

View File

@@ -11,8 +11,8 @@ import { MessageId, freezeMessage } from '@deepseek-ai/dsh-llm'
import { describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import AgentRegistry, { InboxItemId } from '@deepseek-ai/dsh-agent'
import type { Agent, InboxItem, InboxPlacement } from '@deepseek-ai/dsh-agent'
import AgentRegistry, { Inbox } from '@deepseek-ai/dsh-agent'
import type { Agent } from '@deepseek-ai/dsh-agent'
import SessionStore from '@deepseek-ai/dsh-session'
import type { SessionId, UserMessage } from '@deepseek-ai/dsh-session'
import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
@@ -20,7 +20,7 @@ import ToolRegistry from '@deepseek-ai/dsh-tools'
import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
import CommandService from '@deepseek-ai/dsh-commands'
import SkillService from '@deepseek-ai/dsh-skill'
import type { HostFrame, MuxFrame } from '../src/api/index.ts'
import type { HostFrame } from '../src/api/index.ts'
import type { RpcRequest, RpcResponse } from '../src/api/rpc.ts'
import { RpcId } from '../src/api/rpc.ts'
import { createApiProxy } from '../src/api-proxy.ts'
@@ -63,7 +63,14 @@ async function harness(options: { commands?: boolean; skills?: boolean } = {}):
/** Register a live structural agent stub (api-proxy-view precedent: only id/session/status/ctx are read). */
function stubAgent(ctx: Context, sessionId?: SessionId): Agent {
const session = ctx.sessions.create(sessionId)
const agent = { id: session.id, session, status: 'idle', ctx } as Agent
const inbox = new Inbox(session, { inserted: () => {}, discarded: () => {}, claimed: () => {} })
const agent = {
id: session.id,
session,
inbox,
status: 'idle',
ctx,
} as Agent
ctx.agents.register(agent)
return agent
}
@@ -78,6 +85,13 @@ async function collect<F>(iterable: AsyncIterable<RpcRequest<F>>, count: number,
return frames
}
/** Read the next payload from an open stream. */
async function nextFrame<F>(iterator: AsyncIterator<RpcRequest<F>>): Promise<F> {
const result = await iterator.next()
if (result.done) throw new Error('stream ended')
return result.value.payload
}
describe('command.list', () => {
it('serves the addressed agent\'s name-sorted catalog', async () => {
const ctx = await harness()
@@ -99,17 +113,6 @@ describe('command.list', () => {
expect(error.code).toBe('internal')
expect(error.message).toContain('command registry')
})
it('does not route a live subagent through the generic command domain', async () => {
const ctx = await harness()
const session = ctx.sessions.create(undefined, { meta: { cwd: '/proj', origin: 'subagent' } })
const agent = { id: session.id, session, status: 'idle', ctx } as Agent
ctx.agents.register(agent)
const api = createApiProxy(ctx, DEFAULTS)
const error = expectErr(await api.commands.list(request({ sessionId: agent.id })))
expect(error).toMatchObject({ code: 'agent-busy' })
})
})
describe('command.execute', () => {
@@ -153,7 +156,7 @@ describe('command.execute', () => {
const api = createApiProxy(ctx, DEFAULTS)
const missing = expectErr(await api.commands.execute(
request({ sessionId: 'session-nope' as SessionId, line: '/x' }), new AbortController().signal))
expect(missing.code).toBe('internal') // Cold Agent-bound access fails loud when persistence is absent.
expect(missing.code).toBe('internal') // no persistence configured: resume fails loud past the gate
const bare = await harness({ commands: false })
const bareApi = createApiProxy(bare, DEFAULTS)
@@ -275,7 +278,7 @@ describe('host/commands-changed frame', () => {
})
})
/** Build one frozen inbox message for the live `agent/inbox/*` events. */
/** Build one frozen inbox message. */
function inboxMessage(id: string, text: string, rpcId?: string): UserMessage {
return freezeMessage({
id: MessageId(id),
@@ -285,28 +288,19 @@ function inboxMessage(id: string, text: string, rpcId?: string): UserMessage {
})
}
/** Build one addressable inbox occurrence around a frozen message. */
function inboxItem(id: string, message: UserMessage, placement: InboxPlacement): InboxItem {
return { id: InboxItemId(id), message, placement }
}
describe('session.updateQueue', () => {
it('routes addressable actions and reports strict steer races', async () => {
it('splices a queued message and reports a lost claim race', async () => {
const ctx = await harness()
const agent = stubAgent(ctx)
const seen: unknown[] = []
agent.updateInbox = (id, action) => {
seen.push({ id, action })
if (id === InboxItemId('present')) return 'applied'
return id === InboxItemId('closed') ? 'steer-unavailable' : 'not-found'
}
const present = inboxMessage('present', 'before')
agent.inbox.splice('next-turn', 0, 0, [present])
const api = createApiProxy(ctx, DEFAULTS)
const applied = await api.sessions.updateQueue({
rpcId: RpcId('q-apply'),
payload: {
sessionId: agent.id,
itemId: InboxItemId('present'),
itemId: MessageId('present'),
action: { kind: 'edit', content: [{ type: 'text', text: 'edited' }] },
},
})
@@ -315,28 +309,15 @@ describe('session.updateQueue', () => {
rpcId: RpcId('q-missing'),
payload: {
sessionId: agent.id,
itemId: InboxItemId('claimed'),
itemId: MessageId('claimed'),
action: { kind: 'remove' },
},
})
expect(expectErr(missing)).toMatchObject({ code: 'queue-item-not-found' })
const closed = await api.sessions.updateQueue({
rpcId: RpcId('q-closed'),
payload: {
sessionId: agent.id,
itemId: InboxItemId('closed'),
action: { kind: 'steer' },
},
expect(agent.inbox.nextTurn[0]).toMatchObject({
id: 'present',
content: [{ type: 'text', text: 'edited' }],
})
expect(expectErr(closed)).toMatchObject({
code: 'steer-unavailable',
details: { itemId: 'closed' },
})
expect(seen).toEqual([
{ id: 'present', action: { kind: 'edit', content: [{ type: 'text', text: 'edited' }] } },
{ id: 'claimed', action: { kind: 'remove' } },
{ id: 'closed', action: { kind: 'steer' } },
])
})
it('rejects a stale occurrence without resuming a cold agent', async () => {
@@ -347,7 +328,7 @@ describe('session.updateQueue', () => {
rpcId: RpcId('q-cold'),
payload: {
sessionId: 'cold-session' as SessionId,
itemId: InboxItemId('stale-item'),
itemId: MessageId('stale-item'),
action: { kind: 'remove' },
},
})
@@ -358,222 +339,64 @@ describe('session.updateQueue', () => {
})
describe('session/queue frames', () => {
it('folds nested mutations observed before their outer enqueue', async () => {
it('publishes authoritative inbox snapshots without duplicating message identity', async () => {
const ctx = await harness()
const api = createApiProxy(ctx, DEFAULTS)
const agent = stubAgent(ctx)
const original = inboxItem('i-edit', inboxMessage('m-edit', 'before'), 'queued')
const edited = inboxItem('i-edit', inboxMessage('m-edit', 'after'), 'queued')
const removed = inboxItem('i-remove', inboxMessage('m-remove', 'remove me'), 'queued')
ctx.on('agent/inbox/enqueue', (subject, item) => {
if (subject !== agent) return
if (item.id === original.id) ctx.emit('agent/inbox/update', agent, edited)
if (item.id === removed.id) ctx.emit('agent/inbox/discard', agent, [removed])
const queued = inboxMessage('m-1', 'queued prompt')
const edited = inboxMessage('m-1', 'edited prompt')
const steering = inboxMessage('m-2', 'steering prompt')
agent.inbox.splice('next-turn', 0, 0, [queued])
agent.inbox.splice('next-step', 0, 0, [steering])
const abort = new AbortController()
const iterator = api.events.mux({
rpcId: RpcId('t-mux-baseline'),
payload: {},
}, abort.signal)[Symbol.asyncIterator]()
const frames = [
await nextFrame(iterator),
await nextFrame(iterator),
]
agent.inbox.splice('next-turn', 0, 1, [edited])
frames.push(await nextFrame(iterator), await nextFrame(iterator))
const injected = freezeMessage({
id: MessageId('m-3'),
role: 'user',
content: [{ type: 'text' as const, text: 'injected context' }],
source: { kind: 'plugin' as const, plugin: 'approval' },
})
const api = createApiProxy(ctx, DEFAULTS)
const live = new AbortController()
const collected = collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-mux-reentrant'), payload: {} }, live.signal), 2, live)
agent.inbox.splice('next-step', 0, 0, [injected])
frames.push(await nextFrame(iterator), await nextFrame(iterator))
abort.abort()
await iterator.return?.()
ctx.emit('agent/inbox/enqueue', agent, original)
ctx.emit('agent/inbox/enqueue', agent, removed)
const liveFrames = (await collected).filter(frame => frame.type === 'session/queue')
expect(liveFrames.map(frame => frame.items)).toEqual([
[{ id: edited.id, placement: edited.placement, message: edited.message }],
])
const replay = new AbortController()
const replayFrames = await collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-mux-reentrant-replay'), payload: {} }, replay.signal), 2, replay)
expect(replayFrames.filter(frame => frame.type === 'session/queue')).toEqual(liveFrames)
})
it('expires unmatched mutations after the synchronous re-entry window', async () => {
const ctx = await harness()
const agent = stubAgent(ctx)
const api = createApiProxy(ctx, DEFAULTS)
const original = inboxItem('i-stale-edit', inboxMessage('m-stale-edit', 'original'), 'queued')
const staleEdit = inboxItem('i-stale-edit', inboxMessage('m-stale-edit', 'stale edit'), 'queued')
const staleTerminal = inboxItem('i-stale-terminal', inboxMessage('m-stale-terminal', 'keep me'), 'queued')
ctx.emit('agent/inbox/update', agent, staleEdit)
ctx.emit('agent/inbox/discard', agent, [staleTerminal])
await Promise.resolve()
ctx.emit('agent/inbox/enqueue', agent, original)
ctx.emit('agent/inbox/enqueue', agent, staleTerminal)
const replay = new AbortController()
const frames = await collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-mux-expired-unseen'), payload: {} }, replay.signal), 2, replay)
expect(frames.filter(frame => frame.type === 'session/queue')).toEqual([{
type: 'session/queue',
sessionId: agent.id,
items: [
{ id: original.id, placement: original.placement, message: original.message },
{ id: staleTerminal.id, placement: staleTerminal.placement, message: staleTerminal.message },
],
}])
})
it('publishes complete live snapshots and replays the latest snapshot on reconnect', async () => {
const ctx = await harness()
const api = createApiProxy(ctx, DEFAULTS)
const agent = stubAgent(ctx)
const live = new AbortController()
const liveStream = api.events.mux({ rpcId: RpcId('t-mux-live'), payload: {} }, live.signal)
// subscribed baseline + one snapshot per accepted inbox occurrence.
const liveCollected = collect<MuxFrame>(liveStream, 3, live)
const queued = inboxItem('i-1', inboxMessage('m-1', 'queued prompt'), 'queued')
const steering = inboxItem('i-2', inboxMessage('m-2', 'steering prompt'), 'steering')
ctx.emit('agent/inbox/enqueue', agent, queued)
ctx.emit('agent/inbox/enqueue', agent, steering)
const liveFrames = (await liveCollected).filter(f => f.type === 'session/queue')
expect(liveFrames).toEqual([
expect(frames.filter(frame => frame.type === 'session/queue')).toEqual([
{
type: 'session/queue',
sessionId: agent.id,
items: [{ id: queued.id, placement: 'queued', message: queued.message }],
items: [
{ id: queued.id, placement: 'queued', message: queued },
{ id: steering.id, placement: 'steering', message: steering },
],
},
{
type: 'session/queue',
sessionId: agent.id,
items: [
{ id: queued.id, placement: 'queued', message: queued.message },
{ id: steering.id, placement: 'steering', message: steering.message },
{ id: edited.id, placement: 'queued', message: edited },
{ id: steering.id, placement: 'steering', message: steering },
],
},
{
type: 'session/queue',
sessionId: agent.id,
items: [
{ id: edited.id, placement: 'queued', message: edited },
{ id: injected.id, placement: 'context', message: injected },
{ id: steering.id, placement: 'steering', message: steering },
],
},
])
// A fresh mux connection replays only the current authoritative snapshot.
const replay = new AbortController()
const replayFrames = await collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-mux-replay'), payload: {} }, replay.signal), 2, replay)
expect(replayFrames.filter(f => f.type === 'session/queue')).toEqual([liveFrames[1]])
})
it('publishes the durable steering event before retiring its transient row', async () => {
const ctx = await harness()
const api = createApiProxy(ctx, DEFAULTS)
const agent = stubAgent(ctx)
const abort = new AbortController()
const collected = collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-steering-order'), payload: {} }, abort.signal), 5, abort)
const steering = inboxItem('i-steering', inboxMessage('m-steering', 'interrupt now'), 'steering')
agent.session.append('turn/start', {
turn: 1,
trigger: { kind: 'message', source: { kind: 'user' } },
})
ctx.emit('agent/inbox/enqueue', agent, steering)
ctx.emit('agent/inbox/dequeue', agent, steering)
agent.session.append('steering/message', {
turn: 1,
message: steering.message,
}, { surfaceOp: 'append' })
const frames = await collected
expect(frames.map(frame => frame.type)).toEqual([
'session/subscribed',
'session/event',
'session/queue',
'session/event',
'session/queue',
])
expect(frames[2]).toMatchObject({
type: 'session/queue',
items: [{ id: steering.id, placement: 'steering' }],
})
expect(frames[3]).toMatchObject({
type: 'session/event',
event: { type: 'steering/message', data: { message: { id: steering.message.id } } },
})
expect(frames[4]).toMatchObject({ type: 'session/queue', items: [] })
})
it('retains claimed steering in re-entrant snapshots until its durable event', async () => {
const ctx = await harness()
const api = createApiProxy(ctx, DEFAULTS)
const agent = stubAgent(ctx)
const steering = inboxItem('i-steering', inboxMessage('m-steering', 'interrupt now'), 'steering')
const queued = inboxItem('i-reentrant', inboxMessage('m-reentrant', 'later'), 'queued')
ctx.on('agent/inbox/dequeue', (subject, item) => {
if (subject === agent && item.id === steering.id) ctx.emit('agent/inbox/enqueue', agent, queued)
})
const abort = new AbortController()
const collected = collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-steering-reentrant-order'), payload: {} }, abort.signal), 6, abort)
agent.session.append('turn/start', {
turn: 1,
trigger: { kind: 'message', source: { kind: 'user' } },
})
ctx.emit('agent/inbox/enqueue', agent, steering)
ctx.emit('agent/inbox/dequeue', agent, steering)
agent.session.append('steering/message', {
turn: 1,
message: steering.message,
}, { surfaceOp: 'append' })
const frames = await collected
expect(frames[3]).toMatchObject({
type: 'session/queue',
items: [
{ id: steering.id, placement: 'steering' },
{ id: queued.id, placement: 'queued' },
],
})
expect(frames[4]).toMatchObject({
type: 'session/event',
event: { type: 'steering/message', data: { message: { id: steering.message.id } } },
})
expect(frames[5]).toMatchObject({
type: 'session/queue',
items: [{ id: queued.id, placement: 'queued' }],
})
})
it('publishes edits in place in the authoritative order', async () => {
const ctx = await harness()
const api = createApiProxy(ctx, DEFAULTS)
const agent = stubAgent(ctx)
const abort = new AbortController()
const collected = collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-mux-updates'), payload: {} }, abort.signal), 5, abort)
const first = inboxItem('i-a', inboxMessage('m-a', 'a'), 'queued')
const second = inboxItem('i-b', inboxMessage('m-b', 'b'), 'queued')
const edited = inboxItem('i-b', inboxMessage('m-b', 'b edited'), 'queued')
ctx.emit('agent/inbox/enqueue', agent, first)
ctx.emit('agent/inbox/enqueue', agent, second)
ctx.emit('agent/inbox/update', agent, edited)
ctx.emit('agent/inbox/dequeue', agent, edited)
const frames = (await collected).filter(frame => frame.type === 'session/queue')
expect(frames.map(frame => frame.items)).toEqual([
[{ id: first.id, placement: first.placement, message: first.message }],
[
{ id: first.id, placement: first.placement, message: first.message },
{ id: second.id, placement: second.placement, message: second.message },
],
[
{ id: first.id, placement: first.placement, message: first.message },
{ id: edited.id, placement: edited.placement, message: edited.message },
],
[{ id: first.id, placement: first.placement, message: first.message }],
])
})
it('publishes an empty snapshot after terminal discard', async () => {
const ctx = await harness()
const api = createApiProxy(ctx, DEFAULTS)
const agent = stubAgent(ctx)
const doomed = inboxItem('i-doomed', inboxMessage('m-5', 'doomed'), 'queued')
ctx.emit('agent/inbox/enqueue', agent, doomed)
ctx.emit('agent/inbox/discard', agent, [doomed])
const abort = new AbortController()
const frames = await collect<MuxFrame>(
api.events.mux({ rpcId: RpcId('t-mux-swept'), payload: {} }, abort.signal), 1, abort)
expect(frames.filter(frame => frame.type === 'session/queue')).toHaveLength(0)
})
})

View File

@@ -59,7 +59,7 @@ function liveAgent(
): Session {
const session = ctx.sessions.create(sid(id), { meta: { cwd: '/proj', ...lineage } })
for (let turn = 1; turn <= turns; turn++) {
session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn })
session.append('user/message', createUserMessage({
content: [{ type: 'text', text: `prompt ${String(turn)}` }],
source: { kind: 'user' },
@@ -67,12 +67,15 @@ function liveAgent(
session.append('turn/end', { turn, reason: { kind: 'completed' } })
}
if (tail !== 'none') {
session.append('turn/start', { turn: turns + 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: turns + 1 })
session.append('user/message', createUserMessage({
content: [{ type: 'text', text: 'open prompt' }],
source: { kind: 'user' },
}), { surfaceOp: 'append' })
if (tail === 'aborted') session.append('turn/end', { turn: turns + 1, reason: { kind: 'aborted' } })
if (tail === 'aborted') session.append('turn/end', {
turn: turns + 1,
reason: { kind: 'aborted', reason: { kind: 'user' } },
})
}
ctx.agents.register({ id: session.id, session, status: 'idle', ctx } as Agent)
return session

View File

@@ -1,16 +1,20 @@
/**
* Projection carrier paths of the host ApiProxy: history tail pages snapshot
* attached state or fold one cold inspected prefix, loadOlder omits the block,
* and live unit changes push session/projection frames.
* Projection carrier paths of the host ApiProxy: the history tail page's
* projections block reads the registry's watermark snapshot (asOfSeq = last
* event seq, one consistent cut); loadOlder pages never carry the block; a
* composition without the registry serves histories without it; a disposed
* registration's key leaves subsequent responses; and every unit change is
* pushed to mux consumers as a session/projection frame minted here.
*/
import { describe, expect, it } from 'vitest'
import { Context } from 'cordis'
import { z } from 'zod'
import AgentRegistry from '@deepseek-ai/dsh-agent'
import AgentRegistry, { Inbox } from '@deepseek-ai/dsh-agent'
import type { Agent } from '@deepseek-ai/dsh-agent'
import { createUserMessage } from '@deepseek-ai/dsh-llm'
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
import type { Session, SessionEvent, SessionHeader } from '@deepseek-ai/dsh-session'
import type { Session } from '@deepseek-ai/dsh-session'
import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
@@ -49,6 +53,8 @@ async function harness(withRegistry: boolean): Promise<{ ctx: Context; session:
await ctx.plugin(AgentRegistry)
if (withRegistry) await ctx.plugin(SessionProjectionRegistry)
const session = ctx.sessions.create()
// The gateway reads both the session and durable inbox baseline.
ctx.agents.register({ id: session.id, session, inbox: new Inbox(session, { inserted: () => {}, discarded: () => {}, claimed: () => {} }), status: 'idle', ctx } as Agent)
return { ctx, session }
}
@@ -80,40 +86,6 @@ describe('session.history projections block', () => {
expect(events.at(-1)?.event.seq).toBe(projections?.asOfSeq)
})
it('folds a cold inspected prefix without publishing an Agent', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
await ctx.plugin(UserInteractionService)
await ctx.plugin(AgentRegistry)
await ctx.plugin(SessionProjectionRegistry)
ctx.sessionProjections.register(lastUserUnit())
const sessionId = SessionId('session-cold-history')
const meta: SessionHeader = { version: 0, id: sessionId, createdAt: 1, cwd: '/tmp' }
const events = [{
type: 'user/message',
seq: 0,
time: 2,
data: createUserMessage({
content: [{ type: 'text', text: 'persisted' }],
source: { kind: 'user' },
}),
surfaceOp: 'append',
}] as SessionEvent[]
ctx.provide('sessionPersistence', {
list: () => Promise.resolve([meta]),
inspect: () => Promise.resolve({ meta, events }),
} as never)
const response = await api(ctx).sessions.history(request({ sessionId }))
expect(response.result.ok).toBe(true)
if (!response.result.ok) throw new Error('unreachable')
expect(response.result.value.projections).toEqual({
asOfSeq: 0,
values: { 'test/last-user': { text: 'persisted' } },
})
expect(ctx.agents.get(sessionId)).toBeUndefined()
})
it('never carries the block on loadOlder pages (beforeSeq present)', async () => {
const { ctx, session } = await harness(true)
ctx.sessionProjections.register(lastUserUnit())
@@ -252,7 +224,7 @@ describe('session/projection push frame', () => {
seedMessages(session, 1)
// Same-reference apply: turn/start does not concern the unit — no frame.
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
seedMessages(session, 1)
const frames = await collected

View File

@@ -57,7 +57,7 @@ async function composed(withTitles = true): Promise<Context> {
function liveAgent(ctx: Context, id: string, turns: number): Session {
const session = ctx.sessions.create(sid(id), { meta: { cwd: '/proj' } })
for (let turn = 1; turn <= turns; turn++) {
session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn })
session.append('user/message', createUserMessage({
content: [{ type: 'text', text: `prompt ${String(turn)}` }],
source: { kind: 'user' },

View File

@@ -138,7 +138,7 @@ describe('session.search', () => {
eventFilters: [
{
kind: 'type',
values: ['user/message', 'assistant/message', 'steering/message'],
values: ['user/message', 'assistant/message'],
},
{ kind: 'surface', values: ['current'] },
],
@@ -182,7 +182,7 @@ describe('session.search', () => {
withBestMatch(0, { sessionId: sid('hidden') }),
withBestMatch(1, { surface: 'shadowed' }),
withBestMatch(2, { type: 'tool/result' }),
withBestMatch(3, { type: 'steering/message', snippet: 'allowed snippet' }),
withBestMatch(3, { type: 'user/message', snippet: 'allowed snippet' }),
],
}),
} as never)

View File

@@ -112,7 +112,7 @@ describe('mux live view computation', () => {
const rawResult = `RAW_RESULT:${'x'.repeat(64 * 1024)}`
const session = ctx.sessions.create()
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
session.append('tool/call', { turn: 1, step: 1, callId: CallId('c-gen'), name: 'gen', arguments: '{}' })
session.append('tool/call', { turn: 1, step: 1, callId: CallId('c-term'), name: 'term', arguments: '{"cmd":"echo hi"}' })
session.append('tool/call', { turn: 1, step: 1, callId: CallId('c-diff'), name: 'diffy', arguments: '{}' })
@@ -175,7 +175,7 @@ describe('mux live view computation', () => {
// history resolves the agent first; a live structural stub is enough (only
// .session is read on this path).
ctx.agents.register({ id: session.id, session, status: 'idle', ctx } as Agent)
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
session.append('tool/call', { turn: 1, step: 1, callId: CallId('h-term'), name: 'term', arguments: '{"cmd":"ls"}' })
// meta rides through to presentResult's ToolResult (the spread arm).
session.append('tool/result', {
@@ -241,7 +241,7 @@ describe('mux live view computation', () => {
const api = createApiProxy(ctx, { provider: 'p', model: 'm', cwd: '/tmp', workspaceRoot: '/tmp' })
const session = ctx.sessions.create()
ctx.agents.register({ id: session.id, session, status: 'idle', ctx } as Agent)
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
const first = appendUserText(session, 'first prompt')
appendAssistantText(session, 'first reply', 1)
const third = appendUserText(session, 'second prompt')
@@ -295,7 +295,7 @@ describe('mux live view computation', () => {
const fiber = await ctx.plugin(Object.assign((inner: Context) => {
session = inner.sessions.create('session-doomed' as SessionId)
}, { inject: ['sessions'] }))
session?.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session?.append('turn/start', { turn: 1 })
session?.append('tool/call', { turn: 1, step: 1, callId: CallId('c-doomed'), name: 'term', arguments: '{"cmd":"x"}' })
// Disposing the owning fiber detaches the session mid-stream; the
// session/disposed listener must clear its open-call table entry.
@@ -314,7 +314,7 @@ describe('mux live view computation', () => {
const collected = collect(stream, 4, abort)
const session = ctx.sessions.create()
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
session.append('turn/start', { turn: 1 })
session.append('tool/call', { turn: 1, step: 1, callId: CallId('c-late'), name: 'term', arguments: '{"cmd":"tail"}' })
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
// The turn/end above cleared the live table; pairing must fall back to

View File

@@ -3,7 +3,7 @@ import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import AgentRegistry from '@deepseek-ai/dsh-agent'
import AgentRegistry, { Inbox } from '@deepseek-ai/dsh-agent'
import type { Agent, AgentFactory } from '@deepseek-ai/dsh-agent'
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
import type { Session } from '@deepseek-ai/dsh-session'
@@ -44,16 +44,15 @@ function stubAgent(session: Session): Agent {
id: session.id,
options: {},
session,
inbox: new Inbox(session, { inserted: () => {}, discarded: () => {}, claimed: () => {} }),
status: 'idle',
acceptsNextStep: false,
ctx: new Context(),
send: () => {},
followup: () => {},
steer: () => ({ outcome: Promise.resolve({ status: 'rejected' as const }) }),
inject: () => {},
send: () => {},
updateInbox: () => 'not-found',
reserveTurnAdmission: () => undefined,
cancel() {},
runMaintenance: task => task(new AbortController().signal),
whenIdle: () => Promise.resolve(),
}
}

View File

@@ -36,11 +36,6 @@ import { hostFrameSchema, muxFrameSchema, askUserQuestionItemSchema } from '../s
import { approvalRequestIdSchema, approvalResponsePayloadSchema } from '../src/api/approvals.schema.ts'
import { askUserQuestionAnswerSchema, questionResponsePayloadSchema } from '../src/api/questions.schema.ts'
import { goalEditRequestSchema } from '../src/api/goals.schema.ts'
import {
subagentHistoryRequestSchema, subagentHistoryValueSchema, subagentListEntrySchema,
subagentListRequestSchema, subagentListValueSchema, subagentPromptRequestSchema,
subagentPromptValueSchema,
} from '../src/api/subagents.schema.ts'
describe('RpcId', () => {
it('brands a raw string at zero runtime cost', () => {
@@ -77,16 +72,9 @@ describe('rpcErrorSchema', () => {
}).code).toBe('model-unavailable')
expect(rpcErrorSchema.parse({ code: 'agent-busy', message: 'm', details: { reason: 'r' } }).code).toBe('agent-busy')
expect(rpcErrorSchema.parse({ code: 'queue-item-not-found', message: 'm', details: { itemId: 'i' } }).code).toBe('queue-item-not-found')
expect(rpcErrorSchema.parse({ code: 'steer-unavailable', message: 'm', details: { itemId: 'i' } }).code).toBe('steer-unavailable')
expect(rpcErrorSchema.parse({ code: 'command-error', message: 'm', details: {} }).code).toBe('command-error')
expect(rpcErrorSchema.parse({ code: 'unknown-command', message: 'm', details: {} }).code).toBe('unknown-command')
expect(rpcErrorSchema.parse({ code: 'title-invalid', message: 'm', details: { sessionId: 's' } }).code).toBe('title-invalid')
expect(rpcErrorSchema.parse({ code: 'subagent-parent-unavailable', message: 'm', details: { parentSessionId: 'p' } }).code).toBe('subagent-parent-unavailable')
expect(rpcErrorSchema.parse({ code: 'subagent-not-found', message: 'm', details: { parentSessionId: 'p', childSessionId: 'c' } }).code).toBe('subagent-not-found')
expect(rpcErrorSchema.parse({ code: 'subagent-catalog-diagnostic', message: 'm', details: { parentSessionId: 'p', childSessionId: 'c', reason: 'corrupt' } }).code).toBe('subagent-catalog-diagnostic')
expect(rpcErrorSchema.parse({ code: 'subagent-not-resumable', message: 'm', details: { childSessionId: 'c' } }).code).toBe('subagent-not-resumable')
expect(rpcErrorSchema.parse({ code: 'subagent-unauthorized', message: 'm', details: { childSessionId: 'c' } }).code).toBe('subagent-unauthorized')
expect(rpcErrorSchema.parse({ code: 'subagent-delivery-unavailable', message: 'm', details: { childSessionId: 'c' } }).code).toBe('subagent-delivery-unavailable')
expect(rpcErrorSchema.parse({ code: 'internal', message: 'm', details: {} }).code).toBe('internal')
})
@@ -142,13 +130,7 @@ describe('sessions domain schemas', () => {
expect(sessionIdSchema.parse('s1')).toBe('s1')
expect(() => sessionIdSchema.parse('')).toThrow()
expect(sessionSummarySchema.parse({ sessionId: 's1', updatedAt: 1, running: false, blank: true })).toMatchObject({ sessionId: 's1', blank: true })
expect(sessionSummarySchema.parse({
sessionId: 's1', updatedAt: 1, running: true, blank: false,
parentSessionId: 'p', origin: 'subagent', cwd: '/x',
})).toMatchObject({ origin: 'subagent', cwd: '/x' })
expect(() => sessionSummarySchema.parse({
sessionId: 's1', updatedAt: 1, running: false, blank: false, origin: 'fork',
})).toThrow()
expect(sessionSummarySchema.parse({ sessionId: 's1', updatedAt: 1, running: true, blank: false, parentSessionId: 'p', cwd: '/x' }).cwd).toBe('/x')
// blank is mandatory: a summary without it fails the parse.
expect(() => sessionSummarySchema.parse({ sessionId: 's1', updatedAt: 1, running: false })).toThrow()
const event = sessionEventSchema.parse({
@@ -280,9 +262,6 @@ describe('sessions domain schemas', () => {
expect(sessionUpdateQueueRequestSchema.parse({
sessionId: 's1', itemId: 'i1', action: { kind: 'remove' },
}).action.kind).toBe('remove')
expect(sessionUpdateQueueRequestSchema.parse({
sessionId: 's1', itemId: 'i1', action: { kind: 'steer' },
}).action.kind).toBe('steer')
expect(() => sessionUpdateQueueRequestSchema.parse({
sessionId: 's1', itemId: 'i1', action: { kind: 'promote' },
})).toThrow()
@@ -292,45 +271,6 @@ describe('sessions domain schemas', () => {
})
})
describe('subagent domain schemas', () => {
it('validates the direct catalog and addressed history pair', () => {
const child = {
kind: 'child', id: 'c', mode: 'continuable', label: 'worker',
activity: 'running', hasChildren: true,
}
const oneShot = {
kind: 'child', id: 'o', mode: 'one-shot', activity: 'inactive', hasChildren: false,
}
const diagnostic = { kind: 'diagnostic', id: 'bad', reason: 'unsupported' }
expect(subagentListEntrySchema.parse(child)).toEqual(child)
expect(subagentListEntrySchema.parse(oneShot)).toEqual(oneShot)
expect(subagentListEntrySchema.parse(diagnostic)).toEqual(diagnostic)
expect(() => subagentListEntrySchema.parse({
kind: 'child', id: 'missing', mode: 'one-shot', activity: 'inactive',
})).toThrow()
expect(subagentListRequestSchema.parse({ parentSessionId: 'p' })).toEqual({ parentSessionId: 'p' })
expect(subagentListValueSchema.parse({
entries: [child, oneShot, diagnostic], parentAvailable: true,
}).entries).toHaveLength(3)
expect(subagentHistoryRequestSchema.parse({
parentSessionId: 'p', childSessionId: 'c', mode: 'continuable', beforeSeq: 4, maxMessages: 2,
}).beforeSeq).toBe(4)
expect(() => subagentHistoryRequestSchema.parse({
parentSessionId: 'p', childSessionId: 'c', mode: 'continuable', maxMessages: 0,
})).toThrow()
expect(subagentHistoryValueSchema.parse({ events: [], hasMore: false }).hasMore).toBe(false)
})
it('validates continuable prompt content and the accepted inbox identity', () => {
expect(subagentPromptRequestSchema.parse({
parentSessionId: 'p', childSessionId: 'c', mode: 'continuable',
content: [{ type: 'text', text: '继续' }],
}).childSessionId).toBe('c')
expect(subagentPromptValueSchema.parse({ messageId: 'm1' }).messageId).toBe('m1')
expect(() => subagentPromptValueSchema.parse({ route: 'started', taskId: 't2' })).toThrow()
})
})
describe('host domain schemas', () => {
it('validates describe request/value', () => {
expect(hostDescribeRequestSchema.parse({})).toEqual({})
@@ -480,7 +420,11 @@ describe('events frame schemas', () => {
{ type: 'question/requested', sessionId: 's', questions: [{ id: 'q', question: 'Q?', options: [{ label: 'L' }], multiSelect: true }] },
{ type: 'question/resolved', sessionId: 's', questionRpcId: 'r', outcome: 'answered' },
{ type: 'session/queue', sessionId: 's', items: [
{ id: 'i1', placement: 'steering', message: { id: 'm1', role: 'user', content: [{ type: 'text', text: 'queued prompt' }], source: { kind: 'user', rpcId: 'r9' } } },
{
id: 'm1',
placement: 'queued',
message: { id: 'm1', role: 'user', content: [{ type: 'text', text: 'queued prompt' }], source: { kind: 'user', rpcId: 'r9' } },
},
] },
{ type: 'session/projection', sessionId: 's', key: 'todos', value: [{ content: 'x', status: 'pending' }], seq: 7 },
{ type: 'stream/error', error: { code: 'internal', message: 'm', details: {} } },
@@ -510,15 +454,26 @@ describe('events frame schemas', () => {
}
})
it('accepts every queue placement and rejects unknown placements', () => {
const item = (placement: string) => ({ type: 'session/queue', sessionId: 's', items: [{
id: 'm', placement,
message: { id: 'm', role: 'user', content: [], source: { kind: 'user' } },
}] })
for (const placement of ['queued', 'steering', 'context']) {
expect(() => muxFrameSchema.parse(item(placement))).not.toThrow()
}
expect(() => muxFrameSchema.parse(item('bogus'))).toThrow()
})
it('rejects a queue snapshot with malformed items', () => {
expect(() => muxFrameSchema.parse({ type: 'session/queue', sessionId: 's', items: 'x' })).toThrow()
expect(() => muxFrameSchema.parse({ type: 'session/queue', sessionId: 's', items: [{ id: '', message: {} }] })).toThrow()
expect(() => muxFrameSchema.parse({ type: 'session/queue', sessionId: 's', items: [{ id: 'i', message: { id: 'm', role: 'user', content: [], source: {} } }] })).toThrow()
expect(() => muxFrameSchema.parse({ type: 'session/queue', sessionId: 's', items: [{ id: '', role: 'user', content: [], source: { kind: 'user' } }] })).toThrow()
expect(() => muxFrameSchema.parse({ type: 'session/queue', sessionId: 's', items: [{ id: 'm', role: 'assistant', content: [], source: { kind: 'user' } }] })).toThrow()
})
it('accepts every host frame branch', () => {
const frames = [
{ type: 'host/session-added', sessionId: 's', blank: true, parentSessionId: 'p', origin: 'subagent' },
{ type: 'host/session-added', sessionId: 's', blank: true, parentSessionId: 'p' },
{ type: 'host/session-added', sessionId: 's', blank: true },
{ type: 'host/session-removed', sessionId: 's' },
{ type: 'host/session-status', sessionId: 's', running: true },
@@ -532,9 +487,6 @@ describe('events frame schemas', () => {
{ type: 'stream/error', error: { code: 'internal', message: 'm', details: {} } },
]
for (const frame of frames) expect(hostFrameSchema.parse(frame)).toMatchObject({ type: frame.type })
expect(() => hostFrameSchema.parse({
type: 'host/session-added', sessionId: 's', blank: true, origin: 'fork',
})).toThrow()
})
})