From b56628ced78e6332bcf45986b84e9082d2754205 Mon Sep 17 00:00:00 2001 From: _Kerman Date: Fri, 24 Jul 2026 17:55:08 +0800 Subject: [PATCH] refactor(agent-loop): clarify pending message flow --- .../cordis/tool-cordis/src/api-catalog.ts | 4 +-- packages/core/agent-loop/src/agent.ts | 33 +++++++++---------- packages/core/agent/src/types.ts | 6 ++-- 3 files changed, 20 insertions(+), 23 deletions(-) diff --git a/packages/cordis/tool-cordis/src/api-catalog.ts b/packages/cordis/tool-cordis/src/api-catalog.ts index 00d0c5f39d..afd6f06539 100644 --- a/packages/cordis/tool-cordis/src/api-catalog.ts +++ b/packages/cordis/tool-cordis/src/api-catalog.ts @@ -906,8 +906,8 @@ export const EVENT_API: readonly EventApiEntry[] = [ name: 'agent/inbox/enqueue', mode: 'emit', signature: '\'agent/inbox/enqueue\'(this: Scoped, agent: Agent, message: AgentMessage): void', - jsDoc: '/**\n * A frozen item entered the queued or steering inbox.\n * @param agent - the owning agent.\n * @param message - accepted content, source, and correlation identity.\n * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.\n * @mode emit\n */', - summary: 'A frozen item entered the queued or steering inbox.', + jsDoc: '/**\n * An item entered the queued or steering inbox.\n * @param agent - the owning agent.\n * @param message - accepted content, source, and correlation identity.\n * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.\n * @mode emit\n */', + summary: 'An item entered the queued or steering inbox.', }, { name: 'agent/prompt-submit', diff --git a/packages/core/agent-loop/src/agent.ts b/packages/core/agent-loop/src/agent.ts index ab6a4fea1c..ec0e1c9910 100644 --- a/packages/core/agent-loop/src/agent.ts +++ b/packages/core/agent-loop/src/agent.ts @@ -13,7 +13,6 @@ import { createScope } from '@deepseek-ai/dsh-scope' import type { Scope } from '@deepseek-ai/dsh-scope' import type { AgentMessage, - AgentMessageId as AgentMessageIdType, CancelOptions, AgentInterruptReason, AgentOptions, @@ -92,7 +91,7 @@ export class ReactLoopAgent extends Agent { send( content: ContentBlock[], options: SendOptions = { target: 'next-turn', wakeup: true, source: { kind: 'user' } }, - ): AgentMessageIdType { + ): AgentMessageId { const id = AgentMessageId(randomUUID()) const { target, wakeup, source } = options if (target === 'next-step' && !wakeup) { @@ -135,10 +134,10 @@ export class ReactLoopAgent extends Agent { if (cause.kind !== 'disposed') emitAgentEvent(this.loopCtx, this, 'agent/cancel-requested', cause) } if (!options.keepInbox) { - const discarded = [ - ...this.queued, - ...this.outbox.filter((item): item is PendingMessage => 'id' in item), - ] + const discarded: AgentMessage[] = [...this.queued] + for (const message of this.outbox) { + if ('id' in message) discarded.push(message) + } // Clear before abort observers run: replacement work belongs to the next turn. this.queued.length = 0 this.outbox.length = 0 @@ -429,18 +428,18 @@ export class ReactLoopAgent extends Agent { /** Commit the outbox and report whether it contained steering. */ private drainOutbox(turn: number): boolean { let steered = false - for (const item of this.outbox.splice(0)) { - if (!('id' in item)) { - this.session.append('user/message', item, { surfaceOp: 'append' }) - continue + for (const message of this.outbox.splice(0)) { + if ('id' in message) { + steered = true + emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', message) + this.session.append( + 'steering/message', + { turn, content: message.content, source: message.source }, + { surfaceOp: 'append' }, + ) + } else { + this.session.append('user/message', message, { surfaceOp: 'append' }) } - steered = true - emitAgentEvent(this.loopCtx, this, 'agent/inbox/dequeue', item) - this.session.append( - 'steering/message', - { turn, content: item.content, source: item.source }, - { surfaceOp: 'append' }, - ) } return steered } diff --git a/packages/core/agent/src/types.ts b/packages/core/agent/src/types.ts index dae911b678..f3e4106339 100644 --- a/packages/core/agent/src/types.ts +++ b/packages/core/agent/src/types.ts @@ -158,7 +158,7 @@ export abstract class Agent { /** * The unified delivery primitive over the (`target` × `wakeup`) matrix. - * Detaches, validates, and freezes one lossless-JSON item, then routes it: + * It routes the caller's typed content and source as follows: * * - `next-turn` (default) queues an item that becomes the sole ordinary * message of its own FIFO-ordered turn; `wakeup` (default `true`) wakes a @@ -169,8 +169,6 @@ export abstract class Agent { * without running the model: an open turn stages it for the next safe log * position, while an idle injection appends it immediately without opening * a turn. - * - * Invalid input throws synchronously before any notification, enqueue, or append. * @param content - the model-facing content blocks to deliver. * @param options - target queue, wakeup decision, and source. * @returns the accepted message's {@link AgentMessageId}, stable across its `agent/inbox/*` events. @@ -289,7 +287,7 @@ declare module 'cordis' { */ 'agent/status'(this: Scoped, agent: Agent, status: AgentStatus): void /** - * A frozen item entered the queued or steering inbox. + * An item entered the queued or steering inbox. * @param agent - the owning agent. * @param message - accepted content, source, and correlation identity. * Scope-filtered dispatch (`@deepseek-ai/dsh-scope`): agent-scoped listeners receive only that agent.