import type { AgentEngine } from '../agent/engine.js'; import type { RunOptions } from '../agent/types.js'; import type { AgentEvent } from '../types.js'; import type { A2AAgentCard, A2AAgentSkill, A2AJsonRpcRequest, A2AJsonRpcResponse, A2ASendTaskParams, A2AGetTaskParams, A2ACancelTaskParams, A2ATask, A2ATaskEvent, A2AMessage } from './types.js'; export interface A2AServerAgentConfig { name: string; description: string; url: string; version?: string; skills?: A2AAgentSkill[]; provider?: A2AAgentCard['provider']; streaming?: boolean; } export interface A2ATaskStore { get(id: string): A2ATask | undefined | Promise; set(task: A2ATask): void | Promise; delete?(id: string): void | Promise; } /** * 引擎事件 → A2A 事件映射器。 * * 返回 `null` 表示该事件不转发到 A2A 流(静默丢弃)。 * `acc` 在单次运行内共享,可用于累积最终文本等跨事件状态。 */ export type A2AEventMapper = (event: AgentEvent, task: A2ATask, acc: A2AEventAccumulator) => A2ATaskEvent | null; /** 单次运行内共享的累积状态 */ export interface A2AEventAccumulator { /** 累积的最终文本输出 */ finalText: string; /** 供自定义映射器存放任意中间状态 */ [key: string]: unknown; } /** 单次 A2A 任务运行的上下文 */ export interface A2ARunContext { task: A2ATask; params: A2ASendTaskParams; /** 从 params.message 提取的纯文本 prompt */ message: string; /** 取消信号 —— 由 tasks/cancel 或客户端断连触发 */ abortSignal: AbortSignal; } /** * 任务执行器:产出引擎事件流。 * * 相较 `getEngine`,本接口不假设事件来自 `AgentEngine`, * 从而允许接入方在自有服务层(多租户解析、线程落库、可观测性等) * 之上驱动运行,仍复用内核的协议编排、心跳与取消能力。 */ export type A2ATaskRunner = (ctx: A2ARunContext) => AsyncIterable; export interface A2AServerHandlerOptions { /** * 引擎工厂(便捷路径):内核负责调用 `engine.run(...)`。 * 与 `runTask` 二选一;同时提供时 `runTask` 优先。 */ getEngine?: (params: A2ASendTaskParams) => AgentEngine | Promise; /** * 自定义任务执行器(高级路径):接管事件来源。 * * 适用于事件由应用服务层产生(而非直接操作 AgentEngine)的场景, * 内核依旧提供任务状态机、心跳、取消与事件映射。 */ runTask?: A2ATaskRunner; buildRunOptions?: (params: A2ASendTaskParams, task: A2ATask) => Partial | Promise>; taskStore?: A2ATaskStore; /** * 自定义事件映射表(按 `AgentEvent.type` 索引),覆盖同名的内置映射。 * * 用于接入方按需改写事件形态,新增事件类型只需加一行映射, * 无需修改流式主循环。返回 `null` 的事件不会下发。 */ eventMappers?: Record; /** * 未被内置或自定义映射覆盖的事件的兜底处理。 * - `'drop'`(默认):丢弃,保持 A2A 流干净 * - `'passthrough'`:以 `metadata` 形式透传,便于调试或富客户端消费 */ unknownEventStrategy?: 'drop' | 'passthrough'; /** * SSE 心跳间隔(毫秒)。设为 `0` 或不设则禁用。 * * 引擎长时间无事件输出时(如长工具调用),定期产出心跳事件, * 避免中间代理/网关因空闲超时切断连接。 */ heartbeatIntervalMs?: number; /** 任务生命周期钩子(落库、指标、会话映射等副作用挂载点) */ onTaskCreated?: (task: A2ATask, params: A2ASendTaskParams) => void | Promise; onTaskUpdated?: (task: A2ATask) => void | Promise; } export interface A2AStreamingJsonRpcResponse { jsonrpc: '2.0'; id: string; stream: AsyncGenerator; } export type A2AHandlerResponse = A2AJsonRpcResponse | A2AStreamingJsonRpcResponse; export declare class InMemoryA2ATaskStore implements A2ATaskStore { private tasks; get(id: string): A2ATask | undefined; set(task: A2ATask): void; delete(id: string): void; } export declare function createA2AAgentCard(config: A2AServerAgentConfig): A2AAgentCard; export declare function createA2AMessage(text: string, role?: 'user' | 'agent', taskId?: string, contextId?: string): A2AMessage; /** withHeartbeat 产出项:引擎事件或保活心跳 */ export type HeartbeatItem = { kind: 'event'; value: T; } | { kind: 'heartbeat'; }; /** * 为任意异步迭代器叠加心跳保活。 * * 当上游在 `intervalMs` 内没有产出事件时,产出一个 `heartbeat` 项, * 供调用方下发 keep-alive,避免反向代理因空闲超时切断长连接。 * * 实现要点(手写 race 极易出错之处): * 1. `pending` 跨轮复用 —— 心跳胜出时**不能**丢弃已发起的 `next()`, * 否则会漏事件或重复消费迭代器; * 2. 定时器在每轮 `finally` 中**必定清理** —— 否则长任务会累积大量 * 游离 timer,持续占用事件循环; * 3. 结束时调用上游 `return()` —— 保证下游提前 break 时上游能及时释放。 */ export declare function withHeartbeat(source: AsyncIterable, intervalMs: number, abortSignal?: AbortSignal): AsyncGenerator>; /** * 内置引擎事件映射表。 * * 以表驱动替代 switch,接入方可通过 `eventMappers` 覆盖任意条目 * 或新增事件类型,无需改动流式主循环。 */ export declare const DEFAULT_A2A_EVENT_MAPPERS: Record; /** * 引擎事件 → A2A 事件的默认映射(保留具名导出以兼容既有调用方)。 */ export declare function mapAgentEventToA2AEvent(event: AgentEvent, task: A2ATask): A2ATaskEvent | null; export declare class A2AServerHandler { private readonly options; private store; private readonly mappers; /** 运行中任务的中止控制器 —— 支撑 tasks/cancel 真正停止引擎 */ private readonly running; constructor(options: A2AServerHandlerOptions); /** 持久化任务并触发更新钩子(钩子异常不影响协议主流程) */ private persist; /** * 解析本次运行的事件流来源。 * * 优先使用 `runTask`(应用自行驱动),否则回落到 `getEngine` 便捷路径。 * `abortSignal` 一律置于 overrides 展开之后,确保取消链路不被覆盖。 */ private resolveEventStream; /** 应用映射表,未命中时按 unknownEventStrategy 兜底 */ private mapEvent; handleJsonRpc(request: A2AJsonRpcRequest): Promise; sendTask(params: A2ASendTaskParams): Promise; sendTaskStreaming(params: A2ASendTaskParams): AsyncGenerator; /** 创建任务、落库并触发创建钩子 */ private beginTask; getTask(params: A2AGetTaskParams): Promise; cancelTask(params: A2ACancelTaskParams): Promise; private createSubmittedTask; } //# sourceMappingURL=server.d.ts.map