@cloudflare/think 编程式提交(Programmatic Submissions):Webhook 与 RPC 场景下的持久化回合提交、幂等重试与状态检查
【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents
导读
在 Cloudflare Workers 上构建 Think Agent 时,Webhook 处理器、RPC 调用方和父级 Worker 通常有严格的超时限制:它们需要"立即接受一个回合、快速返回、安全重试、稍后查看结果"。submitMessages()正是为此设计的——它不等待推理完成,而是先把回合持久化到 Durable Object 的 SQLite 中,返回一个提交记录,再异步执行模型回合。本文以@cloudflare/think的编程式提交(Programmatic Submissions)为核心,讲解submitMessages()的完整 API、五种状态与幂等语义、inspectSubmission/listSubmissions/cancelSubmission/deleteSubmissions的管理操作,并结合仓库源码与examples/think-submissions示例工程,说明它如何与saveMessages()、startFiber()、Workflows、定时任务(getScheduledTasks())正确分层使用。
为什么需要编程式提交
如果你的代码调用saveMessages()并等待回合结束,而调用方(比如一个 Webhook 处理器)在等待中超时了,你无法判断该回合究竟是"从未被接受"、"排队中"、"正在运行"还是"已经完成"。此时盲目重试可能重复插入一条用户消息并启动第二个回合——这在支付回调、订单同步等外部事件处理中是灾难性的。
submitMessages()创建了一个持久化的接受边界(durable acceptance boundary):它在推理运行之前就可靠地接受该回合并返回提交记录,让调用方可以先安全返回,稍后再查询状态:
const submission = await this.submitMessages( [ { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: "Process webhook event 123" }] } ], { idempotencyKey: "webhook-event-123" } );调用方可以立即返回submission.submissionId。之后调用inspectSubmission(submissionId)或listSubmissions()检查状态。
从源码结构看,submitMessages()的提交记录持久化在cf_think_submissions表中(见 packages/think/src/think.ts),提交时即写入status = 'pending'、created_at、序列化的messages_json与metadata_json,随后通过schedule(0, "_drainThinkSubmissions", ...)调度一个幂等的队列排空任务。也就是说,"接受"与"执行"被刻意拆开:接受是同步持久化的,执行是异步排队的。
API 概览
const submission = await this.submitMessages(messages, { submissionId: "optional-stable-id", // 可选:自定义稳定 ID idempotencyKey: "external-job-id", // 可选:外部系统的幂等键 metadata: { source: "webhook" }, // 可选:附加元数据(JSON 可序列化) channel: "web", // 可选:该提交所属的频道 ID });submitMessages()接受可序列化的UIMessage[]。它不接受saveMessages((messages) => ...)支持的函数形式——因为持久化提交在执行前就要保存工作,无法存储闭包。消息数组必须至少包含一条消息,否则实现会抛出"submitMessages requires at least one message"(见 packages/think/src/think.ts 中_admitTurn内的校验)。
参数详解(源码确认)
对照 packages/think/src/think.ts 中的SubmitMessagesOptions类型:
| 参数 | 类型 | 说明 |
|---|---|---|
submissionId | string | 可选的稳定提交 ID,缺省时使用crypto.randomUUID()。提交记录同时也作为request_id使用 |
idempotencyKey | string | 外部系统的工作 ID,用于幂等重试去重 |
metadata | Record<string, unknown> | 任意 JSON 可序列化元数据,与提交一起持久化,可通过inspectSubmission读回 |
channel | string | 提交所属的频道 ID,会被盖章(stamp)到用户消息上,使排空执行时的回合能从历史中重新解析频道 |
提交状态机
ThinkSubmissionStatus共有六种状态(类型定义见 packages/think/src/think.ts):
| 状态 | 含义 |
|---|---|
pending | 已接受,正在等待其回合 |
running | 已被 Agent 认领并正在执行 |
completed | Think 回合成功完成 |
aborted | 提交被取消 |
skipped | 提交执行前回合状态被重置 |
error | 执行失败或恢复不安全 |
inspectSubmission()返回的ThinkSubmissionInspection对象还包含submissionId、idempotencyKey、requestId、error、metadata、createdAt、startedAt、completedAt等字段,方便调用方完整还原一次提交的生命周期。前四种状态之外的aborted/skipped/error与completed一起构成"终态",终态记录会一直保留,直到你显式删除。
幂等重试:从外部系统传入幂等键
从你的外部系统传入idempotencyKey。用相同键重试时,submitMessages()不会插入重复消息,而是返回已存在的提交记录,并标记accepted: false:
const first = await this.submitMessages(messages, { idempotencyKey: payload.id }); const retry = await this.submitMessages(messages, { idempotencyKey: payload.id }); // retry.submissionId === first.submissionId // retry.accepted === false如果同时传入submissionId和idempotencyKey,它们必须指向同一个提交。实现中依次读取existingById与existingByKey,一旦发现二者指向不同行,会抛出"submissionId and idempotencyKey refer to different submissions",而不是在两个身份中随意二选一(源码见 packages/think/src/think.ts 的submitMessages实现)。
accepted: true表示本次调用是首次接受;accepted: false表示命中了已存在的记录。仓库测试packages/think/src/tests/submissions.test.ts中有专门用例验证:
- 顺序重试去重:"deduplicates retries by idempotency key without appending duplicate messages",断言
first.accepted === true、retry.accepted === false且不会重复追加消息; - 并发首次提交去重:"deduplicates concurrent first submissions with the same idempotency key";
- 冲突检测:"rejects conflicting submission id and idempotency key pairs",断言抛出同一错误消息。
检查、列出、取消与删除
// 检查单个提交的当前状态 const current = await this.inspectSubmission(submission.submissionId); // 列出活跃中的提交(按状态过滤) const active = await this.listSubmissions({ status: ["pending", "running"] }); // 持久化取消:跨 Worker / Durable Object RPC 边界有效 await this.cancelSubmission(submission.submissionId, "No longer needed");取消是持久的(durable cancellation)。与仅作用于单次请求的AbortSignal不同,cancelSubmission()通过AbortController中止进行中的请求,并把状态更新为aborted、写入error_message与completed_at。源码中cancelSubmission()还会调用abortRequest(request_id, reason)并发出submission:status事件。它只对pending/running状态生效;终态记录调用后直接返回。
终态记录会一直保留,直到你删除它们。批量清理可配合completedBefore时间戳做定期 GC:
await this.deleteSubmissions({ status: ["completed", "error", "aborted"], completedBefore: new Date(Date.now() - 7 * 24 * 60 * 60 * 1000) });deleteSubmissions()默认只处理四种终态(completed、aborted、skipped、error),limit默认 100、上限 500,并把删除 SQL 按 SQLite 100 个绑定参数上限分批执行以减少往返(源码见 packages/think/src/think.ts)。另有单条删除deleteSubmission(submissionId),同样只允许删除终态记录。建议在 cron 或定时任务中周期性调用,避免提交表无限增长。
Session 行为:先记账、后入话,严格 FIFO
Think 的执行语义很关键:先把已接受的提交写入提交账本(submission ledger),只有当该提交开始执行时才把消息追加到对话Session。这保证了严格 FIFO 回合语义——后接受的提交在其自身回合开始之前,对模型不可见。源码中_drainSubmissions()按created_at ASC, submission_id ASC顺序每次取一条pending记录执行,执行时通过_admitTurn({ admission: "execute-submission" })进入回合队列,实现先进先出。
由此派生的行为:
- 如果某提交在消息被应用之前被取消(包括已认领但仍在等待回合的情况),这些消息不会持久化到对话中;
- 如果聊天被清空或在某个 pending 提交运行前回合状态被重置,该提交会被标记为
skipped——_markPendingSubmissionsSkipped()会把所有未执行的 pending 提交置为skipped。测试 "marks pending submissions as skipped on turn reset" 与 "does not let stream errors override aborted or skipped submission results" 验证了这一点; - 因此,取消、重置等管理操作可以安全进行,不会污染对话历史。
与saveMessages()的取舍
saveMessages()与submitMessages()只是同一个统一入口runTurn()的两种便捷模式:runTurn({ mode: "wait" })对应saveMessages(),runTurn({ mode: "submit" })对应submitMessages()(见 docs/think/index.md)。
调用方可以等待完整回合结束时用saveMessages():
const result = await this.saveMessages(messages); // result.status is final超时歧义会让重试不安全时用submitMessages():
const result = await this.submitMessages(messages, { idempotencyKey }); // result.status is the accepted submission state补充说明:
submit模式不接受函数输入(submitMessages((m) => ...)不可用),因为提交无法持久化闭包;- 阻塞模式(
wait/stream/continuation)不能嵌套:在活跃回合内部(例如工具execute中)调用会因死锁回合队列而抛错,此时应改用submitMessages()(持久化、当前回合释放队列后运行)或addMessages()(仅写转录、不跑推理); waitUntilStable()在调用方需要避免当前聊天 UI 处于回合中途时继续接受新工作时仍然有用,但它不是持久化准入所必需的——已接受的提交由 Think 序列化,且在自己的回合开始前不会把消息追加进 Session。
与 Workflows、startFiber()的分层
这三者解决不同粒度的"持久化"问题:
| API | 持久化单元 | 适用场景 |
|---|---|---|
submitMessages() | 一次 Think 对话回合 | Webhook / RPC 调用方提交单次回合,要求快速 ACK + 安全重试 |
startFiber() | Agent 拥有的周边副作用工作 | 围绕回合的应用级作业:一次接受 webhook、恢复序列化的聊天线程、发布可见回复 |
| Workflows | 多步骤编排 | 每步重试、长等待、外部事件、人工审批、或作为更大流程的一部分触发 Think |
关键的分工原则:submitMessages()负责 Think 的消息/Session 准入;startFiber()负责围绕它的应用作业;Workflows 负责多步编排。不要用托管 fiber 替换 Think 内部的chatRecovery恢复 fiber,除非该工作确实是调用方需要检查或去重的应用级作业。
组合示例:外层 fiber + 内层提交
当外部作业大于 Think 回合本身时,两个 API 配合使用:
await this.startFiber( "reply-to-webhook", async (ctx) => { ctx.stash({ webhookId, threadId }); const submission = await this.submitMessages(messages, { idempotencyKey: `turn:${webhookId}`, metadata: { threadId } }); await postVisibleReply(threadId, submission.submissionId); }, { idempotencyKey: `webhook:${webhookId}`, waitForCompletion: true } );外层托管 fiber 回答的是:"这个 Webhook 作业是否被接受、恢复、取消或解决?";内层提交回答的是:"这次 Think 回合是否被准入到会话并完成?"——保持两个边界分离。
Workflows 也可以组合使用本 API:
const submission = await this.agent.submitMessages(messages, { idempotencyKey: event.payload.jobId });定时任务与编程式提交的关系
声明式定时提示任务(getScheduledTasks())在底层走的是同一条持久化提交路径:Think 在启动时对声明做对账(reconcile)、持久化下一次一次性调度、每次运行后重新武装下一次。prompt任务创建的就是一个带稳定幂等键的submitMessages()提交,因此重试不会重复工作(详见 docs/think/index.md 与 examples/think-submissions/README.md)。
选择原则:
- 触发器是周期性且代码声明的 → 用
getScheduledTasks(); - 外部调用方或 Webhook 创建一次性工作→ 直接用
submitMessages()。
完整实战:examples/think-submissions示例工程
仓库的 examples/think-submissions 提供了一个完整可运行的控制台示例,专门演示持久化提交的全生命周期:即时 ACK、幂等重试、队列状态、取消和定时任务。
运行
npm install npm start打开开发 URL,在仪表盘中依次体验:
- 提交一个提示词;
- 观察返回的
{ submissionId, accepted, status }即时回执; - 观看提交依次经过
pending、running到达终态; - 用相同
idempotencyKey重试,确认不会创建重复回合; - 取消一个
pending或running的提交; - 检查代码声明的每小时定时任务(本地想快速观察触发,可临时改成
every 5 minutes)。
服务端封装
examples/think-submissions/src/server.ts 定义了一个TaskAgent extends Think,通过@callable()把五个 API 暴露给客户端:
export class TaskAgent extends Think<Env> { getModel() { return "@cf/moonshotai/kimi-k2.7-code"; } getScheduledTasks(): ThinkScheduledTasks { return { hourlyQueueDigest: { schedule: "every 1 hour", prompt: "Write a concise hourly reminder that durable background task queues should be checked for stuck or failed work." } }; } @callable() async submitTask(prompt: string, idempotencyKey?: string) { const key = idempotencyKey?.trim() || undefined; return this.submitMessages( [ { id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: prompt }] } ], { idempotencyKey: key, metadata: { source: "example", promptPreview: prompt.slice(0, 120) } } ); } @callable() async inspectTask(submissionId: string): Promise<ThinkSubmissionInspection | null> { return this.inspectSubmission(submissionId); } @callable() async listTasks(status?: ThinkSubmissionStatus) { return this.listSubmissions({ status, limit: 25 }); } @callable() async cancelTask(submissionId: string) { await this.cancelSubmission(submissionId, "Cancelled from dashboard"); } }对应的 examples/think-submissions/wrangler.jsonc 配置了TaskAgent的 Durable Object 绑定与new_sqlite_classes迁移(tag: "v1")、nodejs_compat兼容标志和 AI binding:
{ "$schema": "./node_modules/wrangler/config-schema.json", "ai": { "binding": "AI", "remote": true }, "compatibility_date": "2026-06-11", "compatibility_flags": ["nodejs_compat"], "durable_objects": { "bindings": [{ "class_name": "TaskAgent", "name": "TaskAgent" }] }, "main": "src/server.ts", "migrations": [{ "new_sqlite_classes": ["TaskAgent"], "tag": "v1" }] }消息格式
提交消息使用 AI SDK 的UIMessage结构,每条消息需要id、role和parts:
{ id: crypto.randomUUID(), role: "user", parts: [{ type: "text", text: prompt }] }metadata中可以附带source、promptPreview等上下文,方便在inspectSubmission()时识别作业来源。
何时使用本 API:决策速查
- 调用方能等待Think 回合结束 →
saveMessages(); - 调用方需要快速持久化回执 + 安全重试(Webhook、RPC、父 Worker)→
submitMessages(); - 调用方还拥有外部副作用(接受一次 webhook、恢复 provider 状态、发布可见回复)→ 在外层再包一个
startFiber(); - 相同的持久化回合应按周期创建 →
getScheduledTasks(); - 作业是多步骤流程(重试、审批、长等待)→ Workflows。
编程式提交是整个 Think "工作必须比请求活得更久" 理念的关键一环:它把"接受工作"与"执行工作"解耦,让 Webhook 和 RPC 调用方在严格超时下也能安全地提交、重试与追踪 Think 回合。
【免费下载链接】agentsBuild and deploy AI Agents on Cloudflare项目地址: https://gitcode.com/GitHub_Trending/agents1/agents
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考