refactor(agent): unify sourced message delivery

This commit is contained in:
_Kerman
2026-07-24 22:38:50 +08:00
parent 009d113e0e
commit 992cf894af
197 changed files with 1890 additions and 2132 deletions

View File

@@ -2,12 +2,11 @@ import { describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import Loader from '@cordisjs/plugin-loader'
import AgentRegistry, { AgentMessageId } from '@deepseek-ai/dsh-agent'
import type { Agent, AgentStatus, AliasSendOptions } from '@deepseek-ai/dsh-agent'
import type { Agent, AgentStatus } from '@deepseek-ai/dsh-agent'
import CommandService from '@deepseek-ai/dsh-commands'
import GoalService from '@deepseek-ai/dsh-goal'
import type { GoalRef } from '@deepseek-ai/dsh-goal'
import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
import { Session, SessionId } from '@deepseek-ai/dsh-session'
import { Session, SessionId, type UserMessageData } from '@deepseek-ai/dsh-session'
import * as commandGoal from '@deepseek-ai/dsh-command-goal'
interface Harness {
@@ -17,24 +16,9 @@ interface Harness {
readonly plugin: Awaited<ReturnType<Context['plugin']>>
}
/** Number the next balanced injection or message turn. */
function nextTurn(session: Session): number {
return session.events.reduce(
(maximum, event) => event.type === 'turn/start' ? Math.max(maximum, event.data.turn) : maximum,
0,
) + 1
}
/** Append one idle injection using the public Agent contract's balanced shape. */
function appendInjection(session: Session, content: ContentBlock[], options?: AliasSendOptions): void {
const source: MessageSource = options?.source ?? { kind: 'plugin', plugin: '' }
const turn = nextTurn(session)
session.append('turn/start', { turn, trigger: { kind: 'injection', source } })
session.append('user/message', {
content,
source,
}, { surfaceOp: 'append' })
session.append('turn/end', { turn, reason: { kind: 'completed' } })
/** Append one idle injection using the public Agent contract. */
function appendInjection(session: Session, input: UserMessageData): void {
session.append('user/message', input, { surfaceOp: 'append' })
}
/** Build a live idle agent accepted by the exact-identity goal service. */
@@ -50,7 +34,7 @@ function stubAgent(id: string): { agent: Agent; session: Session } {
send: () => AgentMessageId('stub'),
followup: () => AgentMessageId('stub'),
steer: () => AgentMessageId('stub'),
inject(content, options) { appendInjection(session, content, options); return AgentMessageId('stub') },
inject(input) { appendInjection(session, input); return AgentMessageId('stub') },
cancel() { status = 'idle' },
retry() {},
whenIdle() { return Promise.resolve() },
@@ -126,7 +110,7 @@ describe('/goal human command', () => {
expect(created.text).toContain('Rounds: 0/256')
expect(created.text).toContain('Activation: armed')
expect(test.ctx.goals.get(test.agent)?.objective).toBe('finish the release')
expect(test.session.events.map(event => event.type)).toEqual(['turn/start', 'user/message', 'turn/end'])
expect(test.session.events.map(event => event.type)).toEqual(['user/message'])
const count = test.session.events.length
await expect(run(test, ' replacement')).resolves.toEqual({

View File

@@ -215,9 +215,7 @@ export function apply(ctx: Context): void {
}
state.attempt = reservation
try {
agent.followup(content, {
source: { kind: 'goal', goalId: goal.id, revision: goal.revision, round },
})
agent.followup({ content: content, source: { kind: 'goal', goalId: goal.id, revision: goal.revision, round } })
} catch (error: unknown) {
state.attempt = undefined
ctx.logger.warn(`goal-session: could not queue round ${round} for agent "${agent.id}": ${renderThrown(error)}`)
@@ -308,7 +306,6 @@ export function apply(ctx: Context): void {
ctx.on('agent/cancel-requested', (agent, cause) => {
const state = stateFor(agent)
const attempt = state.attempt
state.attempt = undefined
state.competingQueued = false
const goal = currentGoal(state)
if (goal?.phase === 'active' && goal.activation === 'armed') {
@@ -316,6 +313,12 @@ export function apply(ctx: Context): void {
disarm(state)
return
}
// An admitted round closes durably as aborted; retain it so the normal
// turn outcome path appends pause after cancellation reaches idle.
// Pausing here would stage context into the active outbox only for this
// same cancel() call to discard it.
if (attempt.turn !== undefined || attempt.phase === 'admitted') return
state.attempt = undefined
try {
applyOutcome(state, goal, { kind: 'pause', reason: cause.kind })
} catch (error: unknown) {

View File

@@ -4,7 +4,7 @@ import type { Agent } from '@deepseek-ai/dsh-agent'
import { agentEvents } from '@deepseek-ai/dsh-agent'
import AgentLoop from '@deepseek-ai/dsh-agent-loop'
import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
import GoalService, { GoalId } from '@deepseek-ai/dsh-goal'
import GoalService, { foldGoal, GoalId } from '@deepseek-ai/dsh-goal'
import type { GoalView } from '@deepseek-ai/dsh-goal'
import { LlmAdapter, LlmError } from '@deepseek-ai/dsh-llm'
import type { GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
@@ -273,7 +273,7 @@ describe('same-session goal driving', () => {
? Promise.resolve({ kind: 'block', reason: 'stop this round' })
: next())
test.ctx.on('goal/changed', (agent, change) => {
if (change.operation === 'block') agent.followup([{ type: 'text', text: 'inspect the blocker' }])
if (change.operation === 'block') agent.followup({ content: [{ type: 'text', text: 'inspect the blocker' }], source: { kind: 'user' } })
})
test.ctx.goals.create(test.agent, { objective: 'stop and inspect' })
@@ -314,13 +314,17 @@ describe('same-session goal driving', () => {
const goal = await waitForGoal(test.ctx, test.agent, current => current?.phase === 'paused')
expect(goal).toMatchObject({ roundsStarted: 1, activation: 'disarmed' })
expect(foldGoal(test.agent.session.events)).toMatchObject({
goal: { phase: 'paused', revision: 2 },
roundsStarted: 1,
})
expect(test.adapter.requests).toHaveLength(1)
})
it('lets already-queued human work finish before reserving the next round', async () => {
const test = await harness([textResponse('human answer'), textResponse('goal answer')])
test.ctx.goals.create(test.agent, { objective: 'continue after the human', maxGoalRounds: 1 })
test.agent.followup([{ type: 'text', text: 'human goes first' }])
test.agent.followup({ content: [{ type: 'text', text: 'human goes first' }], source: { kind: 'user' } })
await waitForGoal(test.ctx, test.agent, goal => goal?.phase === 'blocked')
@@ -361,7 +365,7 @@ describe('same-session goal driving', () => {
test.ctx.on('agent/inbox/enqueue', (agent, info) => {
if (agent !== test.agent || info.source.kind !== 'goal' || inserted) return
inserted = true
agent.followup([{ type: 'text', text: 'human joined the pending batch' }])
agent.followup({ content: [{ type: 'text', text: 'human joined the pending batch' }], source: { kind: 'user' } })
})
test.ctx.goals.create(test.agent, { objective: 'yield to nested human input', maxGoalRounds: 1 })
@@ -444,11 +448,11 @@ describe('same-session goal driving', () => {
// Reject only the goal-sourced round follow-up, not the state-change injection
// that precedes it.
const realFollowup = test.agent.followup.bind(test.agent)
vi.spyOn(test.agent, 'followup').mockImplementation((content, options) => {
if (options?.source?.kind === 'goal') {
vi.spyOn(test.agent, 'followup').mockImplementation((input) => {
if (input.source.kind === 'goal') {
throw new Error('queue rejected')
}
return realFollowup(content, options)
return realFollowup(input)
})
test.ctx.goals.create(test.agent, { objective: 'handle queue failure' })
@@ -465,12 +469,12 @@ describe('same-session goal driving', () => {
it('preserves a custom agent side effect when followup disarms before throwing', async () => {
const test = await harness([])
const realFollowup = test.agent.followup.bind(test.agent)
vi.spyOn(test.agent, 'followup').mockImplementation((content, options) => {
if (options?.source?.kind === 'goal') {
vi.spyOn(test.agent, 'followup').mockImplementation((input) => {
if (input.source.kind === 'goal') {
test.ctx.goals.disarm(test.agent)
throw new Error('queue rejected after disarm')
}
return realFollowup(content, options)
return realFollowup(input)
})
test.ctx.goals.create(test.agent, { objective: 'preserve the newer activation state' })
@@ -548,9 +552,7 @@ describe('same-session goal driving', () => {
it('blocks forged goal attribution without touching an absent reservation', async () => {
const test = await harness([])
test.agent.followup([{ type: 'text', text: 'forged automatic work' }], {
source: { kind: 'goal', goalId: GoalId('forged-goal'), revision: 1, round: 1 },
})
test.agent.followup({ content: [{ type: 'text', text: 'forged automatic work' }], source: { kind: 'goal', goalId: GoalId('forged-goal'), revision: 1, round: 1 } })
await test.agent.whenIdle()
expect(test.adapter.requests).toHaveLength(0)
@@ -559,7 +561,7 @@ describe('same-session goal driving', () => {
it('does not invent goal state when ordinary queued work is cancelled', async () => {
const test = await harness([])
test.agent.followup([{ type: 'text', text: 'cancel ordinary work' }])
test.agent.followup({ content: [{ type: 'text', text: 'cancel ordinary work' }], source: { kind: 'user' } })
test.agent.cancel({ kind: 'user' })
await test.agent.whenIdle()
@@ -569,7 +571,7 @@ describe('same-session goal driving', () => {
it('disarms without durably pausing when cancellation belongs to unrelated human work', async () => {
const test = await harness(['hang'])
test.agent.followup([{ type: 'text', text: 'inspect something first' }])
test.agent.followup({ content: [{ type: 'text', text: 'inspect something first' }], source: { kind: 'user' } })
await waitForRequests(test.adapter, 1)
const created = test.ctx.goals.create(test.agent, { objective: 'continue after inspection' })

View File

@@ -490,7 +490,8 @@ export class GoalService extends Service {
const pending: PendingGoalChange = { change, activation, applied: false }
cache.pending.push(pending)
try {
agent.inject(renderGoalChange(change), {
agent.inject({
content: renderGoalChange(change),
source: { kind: 'goal', goalId: ref.id, revision: ref.revision, round: 0, change },
})
} catch (error: unknown) {

View File

@@ -1,9 +1,9 @@
import { describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import AgentRegistry, { agentEvents, AgentMessageId } from '@deepseek-ai/dsh-agent'
import type { Agent, AgentStatus, AliasSendOptions } from '@deepseek-ai/dsh-agent'
import type { Agent, AgentStatus } from '@deepseek-ai/dsh-agent'
import { HarnessError, type ContentBlock, type MessageSource } from '@deepseek-ai/dsh-llm'
import SessionStore, { Session, SessionId } from '@deepseek-ai/dsh-session'
import SessionStore, { Session, SessionId, type UserMessageData } from '@deepseek-ai/dsh-session'
import GoalService, {
GoalError,
GoalId,
@@ -13,10 +13,7 @@ import GoalService, {
} from '@deepseek-ai/dsh-goal'
import type { GoalChangeMeta, GoalRef, GoalSnapshotChangeMeta } from '@deepseek-ai/dsh-goal'
interface DeferredInjection {
content: ContentBlock[]
options: AliasSendOptions | undefined
}
type DeferredInjection = UserMessageData
interface StubAgent {
agent: Agent
@@ -32,23 +29,9 @@ function nextTurn(session: Session): number {
return session.events.reduce((max, event) => event.type === 'turn/start' ? Math.max(max, event.data.turn) : max, 0) + 1
}
/** Mirror the public Agent.inject idle/open-turn contract for domain tests. */
function appendInjection(session: Session, content: ContentBlock[], options?: AliasSendOptions): void {
const source: MessageSource = options?.source ?? { kind: 'plugin', plugin: '' }
const context = {
content,
source,
}
const last = session.events.at(-1)
const open = last !== undefined && last.type !== 'turn/end'
if (open) {
session.append('user/message', context, { surfaceOp: 'append' })
return
}
const turn = nextTurn(session)
session.append('turn/start', { turn, trigger: { kind: 'injection', source } })
session.append('user/message', context, { surfaceOp: 'append' })
session.append('turn/end', { turn, reason: { kind: 'completed' } })
/** Mirror the public Agent.inject contract for domain tests. */
function appendInjection(session: Session, input: UserMessageData): void {
session.append('user/message', input, { surfaceOp: 'append' })
}
/** Build a registry-compatible agent around one concrete session. */
@@ -66,9 +49,9 @@ function stubAgentForSession(session: Session): StubAgent {
send: () => AgentMessageId('stub'),
followup: () => AgentMessageId('stub'),
steer: () => AgentMessageId('stub'),
inject(content, options) {
if (shouldDefer) deferred.push({ content, options })
else appendInjection(session, content, options)
inject(input) {
if (shouldDefer) deferred.push(input)
else appendInjection(session, input)
return AgentMessageId('stub')
},
cancel() {},
@@ -83,7 +66,7 @@ function stubAgentForSession(session: Session): StubAgent {
setStatus(value) { status = value },
drain() {
shouldDefer = false
for (const injection of deferred.splice(0)) appendInjection(session, injection.content, injection.options)
for (const injection of deferred.splice(0)) appendInjection(session, injection)
},
}
}
@@ -112,7 +95,7 @@ function appendRound(session: Session, ref: GoalRef, round: number): void {
}
describe('GoalService creation and replay', () => {
it('applies the configured default and writes one balanced verbatim context snapshot', async () => {
it('applies the configured default and writes one verbatim context snapshot', async () => {
vi.useFakeTimers()
vi.setSystemTime(1_700_000_000_000)
const { ctx, agent, session } = await harness({ defaultMaxGoalRounds: 17 })
@@ -133,8 +116,8 @@ describe('GoalService creation and replay', () => {
})
expect(goal.id).toMatch(/^goal-/)
expect(seen).toEqual(['create'])
expect(session.events.map(event => event.type)).toEqual(['turn/start', 'user/message', 'turn/end'])
const context = session.events[1]
expect(session.events.map(event => event.type)).toEqual(['user/message'])
const context = session.events[0]
expect(context?.type).toBe('user/message')
if (context?.type !== 'user/message') throw new Error('expected goal context')
expect(context.data.source).toMatchObject({ kind: 'goal', goalId: goal.id, revision: 1, round: 0 })
@@ -438,7 +421,7 @@ describe('GoalService mutations', () => {
expect(deferred).toHaveLength(3)
expect(session.events).toHaveLength(0)
appendInjection(session, [{ type: 'text', text: 'unrelated' }], { source: { kind: 'plugin', plugin: 'test' } })
appendInjection(session, { content: [{ type: 'text', text: 'unrelated' }], source: { kind: 'plugin', plugin: 'test' } })
expect(ctx.goals.get(agent)).toMatchObject({ revision: 3, phase: 'paused' })
test.drain()
expect(deferred).toHaveLength(0)
@@ -472,9 +455,9 @@ describe('GoalService mutations', () => {
const stub = stubAgent('goal-rejected-injection')
const append = stub.agent.inject.bind(stub.agent)
let reject = true
stub.agent.inject = (content, options) => {
stub.agent.inject = (input) => {
if (reject) throw new Error('injection rejected')
return append(content, options)
return append(input)
}
ctx.agents.register(stub.agent)
@@ -493,7 +476,7 @@ describe('GoalService mutations', () => {
test.ctx.goals.edit(test.agent, created, { objective: 'ordered edit' })
const second = test.deferred[1]
if (second === undefined) throw new Error('expected a second deferred goal mutation')
appendInjection(test.session, second.content, second.options)
appendInjection(test.session, second)
expect(() => test.ctx.goals.get(test.agent)).toThrow('advance the current goal')
})
@@ -548,10 +531,10 @@ describe('GoalService mutations', () => {
createdAt: 12,
updatedAt: 12,
}
appendInjection(session, renderGoalChange(change), {
appendInjection(session, { content: renderGoalChange(change),
source: { kind: 'goal', goalId: change.goal.id, revision: 1, round: 0, change },
})
appendInjection(session, [{ type: 'text', text: 'corrupt' }], {
appendInjection(session, { content: [{ type: 'text', text: 'corrupt' }],
source: {
kind: 'goal', goalId: change.goal.id, revision: 2, round: 0,
change: { ...change, operation: 'edit', extra: true } as never,
@@ -645,7 +628,7 @@ describe('goal replay validation', () => {
expect(decodeGoalChange(undefined)).toBeUndefined()
expect(decodeGoalChange({ kind: 'other' })).toBeUndefined()
const session = new Session(SessionId('unrelated'))
appendInjection(session, [{ type: 'text', text: 'other' }], {
appendInjection(session, { content: [{ type: 'text', text: 'other' }],
source: { kind: 'plugin', plugin: 'test' },
})
expect(foldGoal(session.events)).toEqual({ roundsStarted: 0 })

View File

@@ -12,7 +12,7 @@ All calls are exclusive, so a model-ordered batch observes earlier mutations and
All three canonical values match the compact JSON already rendered to Native callers: `{ goal: null }` or `{ goal: { id, revision, objective, phase, roundsStarted, maxGoalRounds, blockedReason? }, activation }`. Programmatic consumers therefore receive the same domain structure without parsing the rendered JSON.
An autonomous goal round that successfully reports `complete` or `blocked` contributes the existing terminal `agent/turn-stop` decision for that physical turn. Direct-human mutations never contribute this stop: the assistant may acknowledge the change and concurrent human steering remains available to the loop.
An autonomous goal round that successfully reports `complete` or `blocked` marks that tool execution with `concludeTurn()` so the physical turn stops after the step. Direct-human mutations never contribute this stop: the assistant may acknowledge the change and concurrent human steering remains available to the loop.
## Authority

View File

@@ -2,11 +2,11 @@ import { describe, expect, it } from 'vitest'
import { Context } from 'cordis'
import Loader from '@cordisjs/plugin-loader'
import AgentRegistry, { agentEvents, AgentMessageId } from '@deepseek-ai/dsh-agent'
import type { Agent, AgentStatus, AliasSendOptions } from '@deepseek-ai/dsh-agent'
import type { Agent, AgentStatus } from '@deepseek-ai/dsh-agent'
import GoalService, { GoalId } from '@deepseek-ai/dsh-goal'
import type { GoalRef } from '@deepseek-ai/dsh-goal'
import { CallId } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
import type { MessageSource } from '@deepseek-ai/dsh-llm'
import { SESSION_FORMAT_VERSION, Session, SessionId } from '@deepseek-ai/dsh-session'
import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
import ToolRegistry from '@deepseek-ai/dsh-tools'
@@ -34,12 +34,8 @@ function stubAgent(rawId: string, supplied?: Session): StubAgent {
send: () => AgentMessageId('stub'),
followup: () => AgentMessageId('stub'),
steer: () => AgentMessageId('stub'),
inject(content: ContentBlock[], options?: AliasSendOptions) {
const source = options?.source ?? { kind: 'plugin', plugin: '' }
session.append('user/message', {
content,
source,
}, { surfaceOp: 'append' })
inject(input) {
session.append('user/message', input, { surfaceOp: 'append' })
return AgentMessageId('stub')
},
cancel() {},