☰
Operit A2A Server 任务生命周期解析:从 task ID 生成到流式状态映射的完整实现
2026/9/26 3:16:56 网站建设 项目流程
  • AI Agent
  • 人工智能
  • 大模型
  • AI 应用
  • 工具调用
  • 本地部署
  • MCP Clients
  • Agent 记忆

【免费下载链接】Operit

The most powerful AI agent and AI chat software on Android/Operit是一款Android上能力最为强大、发展最久的AI Agent

项目地址:https://gitcode.com/gh_mirrors/op/Operit
点击查看免费下载

Operit 通过既有外部 HTTP 服务对外提供 A2A 1.0 Server 能力,其中最关键的一环是把 A2A 协议的"任务"(Task)安全地映射到 Operit 内部独立的流式聊天执行上。本文基于仓库中的 A2A Server 任务生命周期设计文档 及其配套的 协议边界文档、文档与验证记录,结合 A2aTaskManager.kt 与 A2aHttpHandler.kt 的完整实现,逐层拆解服务端 task ID 生成、上下文与独立聊天的绑定、任务状态机、returnImmediately语义、SSE 流式映射以及查询与取消机制。读完本文,你将掌握 Operit A2A 任务从提交到终态的完整生命周期,并可直接对照源码与 A2A Server 协议文档 进行二次开发或联调。

一、任务生命周期在 A2A Server 中的定位

A2A(Agent2Agent)1.0 协议要求 Server 对每次调用暴露一个可寻址的 Task 资源:客户端提交消息、查询任务状态、订阅事件流,并可随时取消。Operit 的设计约束(见 index.md)要求:

  • 复用既有聊天执行器,但将每个 A2A 请求隔离为独立聊天与独立任务;
  • A2A 调用不共享、不篡改用户正在进行的普通聊天;
  • 任务只存活于外部 HTTP 服务运行期间,客户端可以查询和取消服务仍存活期间的任务。

围绕这三个约束,任务生命周期层(integrations/a2a/下的A2aTaskManager)承担了全部执行管理职责,而协议序列化(JSON-RPC 封装、SSE 编码)由A2aHttpHandler负责。二者分工清晰:A2aTaskManager只管理执行,不关心线协议(源码注释原文:"The A2A protocol layer owns serialization; this class only manages execution.",见 A2aTaskManager.kt)。

二、服务端 task ID 生成与任务注册

每个 A2A 请求进入系统后,第一步是生成服务端唯一的 task ID。A2aTaskManager.submit()的实现如下(A2aTaskManager.kt):

fun submit(message: A2aIncomingMessage): A2aTaskSnapshot { val contextId = resolveContextId(message) val chatId = contextChats[contextId] val taskId = UUID.randomUUID().toString() val record = TaskRecord(taskId, contextId) tasks[taskId] = record val job = serviceScope.launch(Dispatchers.IO) { executeTask(record, message.text, chatId) } record.attachJob(job) return record.snapshot() }

要点拆解:

  • task ID 采用UUID.randomUUID()生成,与服务端既有聊天 ID、上下文 ID 完全解耦,客户端拿到的result.task.id就是这个值;
  • 任务注册在ConcurrentHashMap<String, TaskRecord>(tasks)中,天然支持并发提交与并发查询;
  • 每个任务持有一个Job(record.attachJob(job)),使取消任务可以直接取消底层协程;
  • submit立即返回record.snapshot(),此时状态为TASK_STATE_SUBMITTED,后续真正的执行在Dispatchers.IO上异步进行。

注意:任务记录是纯内存态的,文档明确"Tasks are held in memory for the lifetime of the external HTTP service"(external_a2a_server.md),服务重启或停用后任务记录即被清空,且 A2A 任务状态不做持久化。

三、上下文 ID 与独立 Operit 聊天的映射

3.1 上下文解析:没有 contextId 则自动生成

private fun resolveContextId(message: A2aIncomingMessage): String { return message.contextId?.trim()?.takeIf { it.isNotBlank() } ?: UUID.randomUUID().toString() }

(A2aTaskManager.kt)

请求若携带非空contextId则原样采用(会做 trim 与空白校验);否则为本次任务生成全新上下文。这正是 A2A 1.0 的语义:contextId表示"同一会话"的延续标识。

3.2 一个上下文对应一个 Operit 聊天

映射表是contextChats: ConcurrentHashMap<String, String>(A2aTaskManager.kt),键为 A2AcontextId,值为 Operit 内部聊天 ID。映射策略在executeTask中落地:

val request = ExternalChatRequest( requestId = record.id, message = message, createNewChat = existingChatId == null, chatId = existingChatId, createIfNone = false, returnToolStatus = false )

(A2aTaskManager.kt)

  • 首个请求:contextChats中没有该上下文对应的聊天,existingChatId == null,于是createNewChat = true,为这个 A2A 上下文创建一条全新的 Operit 聊天;
  • 后续请求:同一上下文的新消息携带相同contextId,从contextChats命中既有聊天 ID,chatId = existingChatId、createNewChat = false,在同一条聊天里续写;
  • 流式会话启动成功后,contextChats[record.contextId] = session.chatId完成绑定(A2aTaskManager.kt)。

这套映射就是"每个新的 A2A 上下文创建独立 Operit 聊天,后续关联到同一上下文的新请求复用对应聊天"这一设计目标的源码级实现(对应 02_task_lifecycle.md 修改项第 2 条)。

3.3 与普通聊天的隔离

由于 A2A 任务只会操作自己上下文专属的新聊天(createIfNone = false且绝不向既有普通聊天传chatId),因此A2A 调用不会共享或篡改用户正在进行的普通聊天——这是 02_task_lifecycle.md 预期结果的第一条,也是"任务隔离"约束(index.md 意图部分)的直接体现。此外,从 01_protocol_boundary.md 可以确认,A2A 与既有 REST 路由只共享 HTTP Server、Token 与聊天执行器,不共享协议对象。

四、任务状态机:从提交到终态

任务状态常量定义在 A2aTaskManager.kt:

状态常量值是否终态
TASK_STATE_SUBMITTEDTASK_STATE_SUBMITTED否
TASK_STATE_WORKINGTASK_STATE_WORKING否
TASK_STATE_INPUT_REQUIREDTASK_STATE_INPUT_REQUIRED否
TASK_STATE_AUTH_REQUIREDTASK_STATE_AUTH_REQUIRED否
TASK_STATE_COMPLETEDTASK_STATE_COMPLETED是
TASK_STATE_CANCELEDTASK_STATE_CANCELED是
TASK_STATE_FAILEDTASK_STATE_FAILED是
TASK_STATE_REJECTEDTASK_STATE_REJECTED是

isTerminalState()判定COMPLETED、CANCELED、FAILED、REJECTED四种终态。实际执行路径中(A2aTaskManager.kt)用到的主要迁移为:

  1. SUBMITTED → WORKING:executor.startStreaming(request)返回Started后,record.start(session)将状态置为WORKING并向订阅者发布非 final 的Status事件(L199-L214);
  2. WORKING → COMPLETED:响应流收集完毕,session.responseStreamSession.currentState()不是InputProcessingState.Error时record.complete();
  3. 任意非终态 → FAILED:启动失败(ExternalChatStreamingStartResult.Failed)或最终状态为Error、或执行过程抛异常时record.fail(message);
  4. 任意非终态 → CANCELED:客户端调用CancelTask、或服务关闭时cancelForShutdown()触发。

TaskRecord用一把synchronized(lock)保护所有状态迁移与输出追加,保证并发场景(一边流式写、一边查询/取消)下状态一致。output以字符串累加方式保存全部文本输出,A2aTaskSnapshot每次快照都携带id、contextId、state、output与可选的error(L312-L320)。

值得注意的边界:Operit 不实现TASK_STATE_INPUT_REQUIRED(即不接受带taskId的消息续写),协议文档与parseIncomingMessage都会拒绝这类消息(A2aHttpHandler.kt)。

五、SendMessage 与 returnImmediately 语义

SendMessage的行为由configuration.returnImmediately决定(A2aHttpHandler.kt):

METHOD_SEND_MESSAGE -> { val sendRequest = parseSendMessageRequest(request.params) val submittedTask = taskManager.submit(sendRequest.message) val resultTask = if (sendRequest.returnImmediately) { submittedTask } else { runBlocking { taskManager.awaitTerminalTask(submittedTask.id) } } jsonRpcResultResponse(request.id, JSONObject().put("task", taskToJson(resultTask))) }
  • returnImmediately = false(默认):JSON-RPC 请求会阻塞等待任务进入终态再返回完整结果,客户端一次请求即拿到最终task(含artifacts文本);
  • returnImmediately = true:SendMessage立即返回提交瞬间的任务快照(通常是TASK_STATE_WORKING),客户端随后用GetTask轮询推进。

awaitTerminalTask底层是TaskRecord.awaitTerminal()(A2aTaskManager.kt):若任务已终态则直接返回快照,否则挂起在CompletableDeferred<A2aTaskSnapshot>上,由complete/fail/cancel等终态迁移来补全。这正对应 02_task_lifecycle.md 修改项第 3 条:"SendMessage按 A2A 1.0 的returnImmediately语义返回任务,后台将流式聊天执行映射为工作、完成、失败或取消状态"。

请求参数校验

parseSendMessageRequest(A2aHttpHandler.kt)还会拒绝taskPushNotificationConfig(抛出A2aPushNotificationNotSupportedException,错误码-32003)、校验acceptedOutputModes必须包含text/plain(否则-32005),并接受可选的historyLength。消息本身要求ROLE_USER、非空messageId、至少一个纯文本 Part({"text": "..."}),多 Part 会以换行拼接(L353-L387)。

六、流式任务:SendStreamingMessage 与 SubscribeToTask

6.1 SSE 管道建立

SendStreamingMessage提交任务后立即返回text/event-stream响应;SubscribeToTask对已存在的活动任务建立同样的流。两者的共同实现是streamingResponse()(A2aHttpHandler.kt):

  • 用PipedInputStream(64 * 1024)/PipedOutputStream连接协程事件流与 HTTP 响应体;
  • 通过Channel<A2aTaskEvent>(Channel.UNLIMITED)转发任务事件;
  • 订阅时requireActive控制:SendStreamingMessage允许终态任务(任务可能已经瞬间完成),SubscribeToTask则要求任务处于活动状态,否则抛出A2aUnsupportedOperationException(-32004,"cannot be subscribed after it is terminal");
  • 响应头附带A2A-Version: 1.0、Cache-Control: no-cache、Connection: keep-alive、X-Accel-Buffering: no,并在ExternalChatHttpServer.useGzipWhenAccepted中对 SSE 响应禁用 gzip(ExternalChatHttpServer.kt),保证逐帧推送不被缓冲。

6.2 事件序列:initial task → artifactUpdate → statusUpdate(final:true)

每个data:行都是一个带 JSON-RPC envelope 的 A2A 1.0 Stream Response(每条数据只有一个成员),写入时对多行 JSON 逐行加data:前缀(writeSseEvent,L542-L551)。典型事件流:

data: {"jsonrpc":"2.0","id":2,"result":{"task":{"id":"task-uuid","contextId":"context-uuid","status":{"state":"TASK_STATE_WORKING"}}}} data: {"jsonrpc":"2.0","id":2,"result":{"artifactUpdate":{"taskId":"task-uuid","contextId":"context-uuid","artifact":{"artifactId":"task-uuid-result","parts":[{"text":"这是回答的第一部分。"}]},"append":true,"lastChunk":false}}} data: {"jsonrpc":"2.0","id":2,"result":{"statusUpdate":{"taskId":"task-uuid","contextId":"context-uuid","status":{"state":"TASK_STATE_COMPLETED"},"final":true}}}

事件映射规则(见 A2aHttpHandler.kt):

  • 首个事件携带result.task(当前快照);
  • 输出增量映射为result.artifactUpdate:artifactId固定为"$taskId-result",append: true、lastChunk: false;
  • 状态变化映射为result.statusUpdate:携带taskId、contextId、status.state(终态时附带含message的Message,role 为ROLE_AGENT)与final标记;
  • 收到Status事件且final == true后流关闭(L232-L240)。

流式收集侧,A2aTaskManager.executeTask通过ExternalChatResponseSanitizer.sanitizeStream(...)清洗流式响应后逐块record.appendOutput(chunk),每块追加都会发布一条Artifact事件(A2aTaskManager.kt)——这就是"后台将流式聊天执行映射为工作、完成、失败或取消状态"的完整链路。

七、快照查询、分页列表与取消

7.1 GetTask:即时快照

GetTask从tasks表取出TaskRecord并返回snapshot()(A2aTaskManager.kt)。任务不存在时抛A2aTaskNotFoundException(JSON-RPC 错误码-32001)。

7.2 ListTasks:过滤与分页

listTasks(contextId, state)(A2aTaskManager.kt)支持按上下文与状态过滤、按任务 ID 排序,然后由 listTasksToJson 做游标分页:

  • pageSize默认为 50,上限 100(DEFAULT_TASK_LIST_PAGE_SIZE = 50、MAX_TASK_LIST_PAGE_SIZE = 100);
  • pageToken是上一页最后一个任务的 ID:indexOfFirst { task -> task.id == token } + 1作为下一页起始;token 非法时抛-32602;
  • 还有剩余任务时返回nextPageToken(取当前页最后一个任务 ID);
  • status参数必须落在VALID_TASK_STATES集合内。

7.3 CancelTask:协同取消

cancelTask(taskId)调TaskRecord.cancel()(A2aTaskManager.kt):

  1. transitionToCancelled()在锁内把非终态任务置为CANCELED并记录错误信息;
  2. 对活动中的流式会话调用responseStreamSession.cancel()中断响应流;
  3. 取消底层协程job.cancel();
  4. 补全terminalTask(让阻塞等待SendMessage的客户端解除挂起);
  5. 发布final: true的Status事件(SSE 订阅端据此关闭流)。

已处终态的任务调用CancelTask会抛A2aTaskNotCancelableException(-32002)。executeTask捕获到协程CancellationException时也会markCancelled()并向上抛出(L155-L157),保证取消路径状态一致。取消后的任务快照即为返回结果,状态为TASK_STATE_CANCELED。

八、生命周期边界与清理

8.1 服务存活期内的任务

A2aTaskManager与外部 HTTP 服务同生命周期:ExternalChatHttpServer构造时创建A2aHttpHandler(进而创建A2aTaskManager),stopServer()时调用a2aHandler.close()(ExternalChatHttpServer.kt)。A2aTaskManager.close()(L109-L113)会:

  • 对所有未终态任务执行cancelForShutdown()(与主动取消相同:中断会话、取消协程、发布 final 事件);
  • 清空tasks与contextChats两张表。

因此"客户端可以查询和取消服务仍存活期间的任务"(02_task_lifecycle.md 预期结果第 2 条)成立;服务停止后任务记录随内存释放而消失,不做持久化。

8.2 路由与协议边界

A2A 入口挂在既有serve()分发链上,且优先级先于健康检查与 REST 路由(ExternalChatHttpServer.kt):

session.uri == A2aHttpHandler.AGENT_CARD_PATH -> a2aHandler.handleAgentCard(session).withCors() session.uri == A2aHttpHandler.A2A_PATH -> a2aHandler.handleJsonRpc(session).withCors()
  • GET /.well-known/agent-card.json(AGENT_CARD_PATH)无需 Bearer Token,从请求Host头动态拼出本次可访问的/a2aJSON-RPC URL(buildJsonRpcEndpoint,L483-L492),并声明protocolVersion: "1.0"、JSONRPCbinding、流式能力(pushNotifications: false)、text/plain输入输出模式与 Bearer 安全方案;
  • POST /a2a除OPTIONS外全部走既有requireBearerToken鉴权,并校验A2A-Version请求头:非1.0返回VersionNotSupportedError(-32009,L471-L481)。

这印证了 01_protocol_boundary.md 的设计:协议解析与任务管理收敛在integrations/a2a/,既有/api/health、/api/external-chat、/api/web/*与 Web Chat 路由保持原样。

九、能力边界速查

根据 external_a2a_server.md 与源码,A2A Server 的能力边界为:

  • 输入/输出模式:仅text/plain,不接受文件、结构化数据 Part;
  • 推送通知:不支持(Agent Card 声明pushNotifications: false,请求带taskPushNotificationConfig即报-32003);
  • 续写:不支持带taskId的消息(无TASK_STATE_INPUT_REQUIRED);
  • 历史与持久化:不返回任务历史,任务状态不跨服务重启持久化;
  • 客户端:本模块不提供 A2A Client,且 A2A 路由不并入 MCP。

联调时请以 A2A Server 协议文档 中的 curl/JSON 示例为准,它同时涵盖了 Agent Card 发现、Bearer 鉴权、六种 JSON-RPC 方法及 SSE 事件格式的完整范例。

  • AI Agent
  • 人工智能
  • 大模型
  • AI 应用
  • 工具调用
  • 本地部署
  • MCP Clients
  • Agent 记忆

【免费下载链接】Operit

The most powerful AI agent and AI chat software on Android/Operit是一款Android上能力最为强大、发展最久的AI Agent

项目地址:https://gitcode.com/gh_mirrors/op/Operit
点击查看免费下载

相关推荐

上一篇:Ingress NGINX 注解风险等级与作用域全解:annotations-risk 治理指南
下一篇:Flipper Zero 手电筒插件(Flashlight)源码剖析与 GPIO 控制实战

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询