跳转到内容
搜索文档

Workflow

最后更新 查看 MarkdownAgent 设置

ThinkWorkflow 在持久作业需要单个模型驱动推理步骤时将 Think 连接到 Cloudflare Workflows。

当 Workflow 拥有流程时使用:

  • 持久的多步骤编排
  • 审批门控或长等待
  • 可重试的确定性副作用
  • 应产生带类型结构化输出的 Think 轮次

循环提示词保持为定时任务,简单一次性后台轮次使用 submitMessages()。步骤重要的作业使用 Workflows。

API

@cloudflare/think/workflows 导入:

import { ThinkWorkflow } from "@cloudflare/think/workflows";

扩展 ThinkWorkflow 并在 run() 内调用 step.prompt()

import { z } from "zod";
import { ThinkWorkflow } from "@cloudflare/think/workflows";

const draftSchema = z.object({
	title: z.string(),
	summary: z.string(),
	labels: z.array(z.string()),
});

export class TriageWorkflow extends ThinkWorkflow {
	async run(event, step) {
		const draft = await step.prompt("triage-issue", {
			prompt: `Triage issue #${event.payload.issueNumber}`,
			output: draftSchema,
			timeout: "3 days",
		});

		await step.do("apply-labels", async () => {
			await this.agent.applyLabels(draft.labels);
		});
	}
}
import { z } from "zod";
import { ThinkWorkflow } from "@cloudflare/think/workflows";
import type { ThinkWorkflowStep } from "@cloudflare/think/workflows";
import type { AgentWorkflowEvent } from "agents/workflows";

const draftSchema = z.object({
	title: z.string(),
	summary: z.string(),
	labels: z.array(z.string()),
});

export class TriageWorkflow extends ThinkWorkflow<TriageAgent, Params> {
	async run(event: AgentWorkflowEvent<Params>, step: ThinkWorkflowStep) {
		const draft = await step.prompt("triage-issue", {
			prompt: `Triage issue #${event.payload.issueNumber}`,
			output: draftSchema,
			timeout: "3 days",
		});

		await step.do("apply-labels", async () => {
			await this.agent.applyLabels(draft.labels);
		});
	}
}

在 Think Agent 内用 runWorkflow() 启动 Workflow:

export class TriageAgent extends Think {
	async triageIssue(issueNumber) {
		return this.runWorkflow(
			"TRIAGE_WORKFLOW",
			{ issueNumber },
			{ metadata: { issueNumber } },
		);
	}
}
export class TriageAgent extends Think<Env> {
	async triageIssue(issueNumber: number): Promise<string> {
		return this.runWorkflow(
			"TRIAGE_WORKFLOW",
			{ issueNumber },
			{ metadata: { issueNumber } },
		);
	}
}

runWorkflow() 创建 Workflow 实例并注入 ThinkWorkflowrun() 内重连 this.agent 所需的 Agent 身份。优先于直接调用 Workflows 绑定:

// Avoid this for Agent workflows. It does not include Agent context.
await this.env.TRIAGE_WORKFLOW.create({ params: { issueNumber } });

当等待中的 Workflow 需要外部信号(如人工审批)时,从 Agent 使用 sendWorkflowEvent()

await this.sendWorkflowEvent("TRIAGE_WORKFLOW", workflowId, {
	type: "approval",
	payload: { approved: true },
});

step.prompt() 接受提示词字符串与 Zod 对象 schema。Workflow 调用 Agent 前 schema 转为 JSON Schema。Think 然后运行完整智能体轮次:Agent 可跨多步使用工具,并通过调用参数匹配 schema 的内部 final_answer 工具返回结构化结果。这使用普通工具调用而非流式 response_format,因此在 Think 支持的所有提供商(含拒绝流式请求上 JSON Schema 响应的 Workers AI)上工作。Workflow 恢复时,载荷在返回带类型的值前再次用原始 Zod schema 验证。

无法表示为 JSON Schema 的不支持 Zod 特性在创建提示词步骤时失败。Think 不会静默修复无效模型输出。若模型未产生有效 final_answer 调用,提交进入终态错误且 step.prompt() 抛出。

行为说明

  • Agent 可先使用工具。 step.prompt() 轮次是完整智能体轮次:Agent 可跨多步调用自有工具再调用最终答案工具。若期望 Agent 在回答前使用工具,至少允许 maxSteps: 2maxSteps: 1 强制第一步回答且不能调用任何其他工具。
  • 结构化轮次期间强制工具使用。 为保证 Agent 以结构化答案终止(而非纯文本回复),Think 为该轮设置 toolChoice。不要在 step.prompt() 轮次从 beforeTurn 重写 toolChoice — 否则可能阻止 Agent 调用最终答案工具,导致提示词失败。
  • think_final_answer 保留。 Think 注入内部 think_final_answer 工具承载结构化结果。此名称(及任何 think_final_answer_* 变体)保留;其调用与结果从持久对话剥离,后续轮次看不到 Think 内部管道。
  • 模型须支持流式工具调用。 Think 流式传输每轮,因此 step.prompt() 仅适用于流式时可靠发出强制工具调用的模型。强工具调用方(如 OpenAI gpt-4o-mini、Anthropic claude-haiku-4-5、Workers AI @cf/moonshotai/kimi-k2.6)已验证可用。部分模型仅在非流式请求上遵守强制 toolChoice,流式时会纯文本回复并停止 — 例如 Workers AI @cf/meta/llama-3.3-70b-instruct-fp8-fast。这些模型轮次结束无 think_final_answer 调用且 step.prompt() 失败(Model ended the turn without calling the think_final_answer tool);请改用流式工具调用可用的模型。

运行方式

调用读起来像阻塞步骤,但不保持长生命周期 Durable Object RPC 打开。

  1. step.do("<name>:submit", ...) 创建或查找幂等 Think 提交。
  2. Think 通过正常提交队列运行提交的轮次。
  3. 提交达到 completederrorabortedskipped 时,Think 记录待处理 Workflow 通知。
  4. Think 用 sendWorkflowEvent() 与 Durable Object alarm 排空通知发件箱直至投递成功。
  5. step.waitForEvent("<name>:wait", ...) 恢复 Workflow。
  6. step.prompt() 验证结构化输出或抛出带类型的错误。

机器可读输出携带在待处理通知与 Workflow 事件载荷中。Think 不在提交账本上存储单独 output_json 列,投递后清除通知载荷。投递后 Workflow 拥有持久结果。

幂等性

默认 step.prompt() 从 Workflow 身份与步骤名称推断幂等键:

think-workflow:<workflowName>:<workflowId>:<stepName>

循环中传递字符串 key 以区分同 step 名的重复使用:

await step.prompt("summarize-file", {
	key: file.path,
	prompt: `Summarize ${file.path}`,
	output: summarySchema,
});

提示词文本不是推断键的一部分,但 Think 存储 Workflow 元数据与提示词/配置指纹供诊断。

超时

传递 timeout 控制 Workflow 等待终态事件的时长。等待超时时,step.prompt() 默认取消 Think 提交并抛出 ThinkPromptTimeoutError

当有意让 Think 提交在 Workflow 停止等待后继续时,设置 cancelOnTimeout: false

与其他原语的边界

周期性提示词提交或确定性调度处理函数使用 getScheduledTasks()

getScheduledTasks() {
	return {
		dailySummary: {
			schedule: "every day at 09:00",
			timezone: "UTC",
			prompt: "Generate the daily report."
		},
		dailyWorkflow: {
			schedule: "every day at 09:00",
			timezone: "UTC",
			retry: { maxAttempts: 3 },
			handler: async ({ idempotencyKey, scheduledFor, timezone }) => {
				await this.env.REPORT_WORKFLOW.create({
					id: idempotencyKey,
					params: { scheduledFor, timezone }
				});
			}
		}
	};
}

调用方稍后检查提交状态的一次性持久轮次使用 submitMessages()

需要在 Agent 内恢复的应用自有幂等 Agent 作业使用 startFiber()。Think 的 Workflow 通知投递不使用 fiber;使用私有发件箱,因需存储事件直至投递成功。

流程有多步确定性步骤、长等待或人工审批时使用 Workflows。

这篇文档对您有帮助吗?