跳转到内容
搜索文档

动态工作流

最后更新 查看 MarkdownAgent 设置

在 Dynamic Worker 内运行工作流 (Workflow),以为在运行时加载的代码提供持久执行。工作流中的每个步骤都可以在故障中存活,可以休眠数小时或数天,可以等待外部事件,并且从它停止的地方精确恢复 — 即使隔离区 (isolate) 在步骤之间被回收。

因为 Dynamic Workers 是按需创建的,所以您不必预先注册每个工作流或单独管理它们。在需要时加载代码,工作流引擎在后台处理持久性和重试。这对于一次性执行以及长时间运行的多步流程同样有效。

例如,您可能正在构建:

  • SaaS 平台:每个租户定义其自己的自动化 — 引导序列、审批链或计费重试逻辑 — 并且您需要每个流程都能持久运行,而无需为每个客户部署单独的工作流。
  • AI 代理框架:代理在运行时生成并执行多步骤计划,并且每个计划都需要在重启中存活,在工具调用之间休眠,并等待人工批准。
  • 多租户作业系统:每个客户提交其自己的处理逻辑 — 数据转换、webhook 链、计划任务 — 并且您希望每个步骤都能持久保存进度并在失败时重试,而无需构建您自己的编排器。

@cloudflare/dynamic-workflows 库将您的 Worker Loader 连接到工作流引擎,以便每个 Dynamic Worker 获得持久步骤(step.do(), step.sleep(), step.waitForEvent())的全部功能,而无需您自己构建管道。

在本指南中,您将使用 @cloudflare/dynamic-workflows 库来设置 Worker Loader,编写具有持久步骤的 Dynamic Worker,并触发工作流实例。

理解模型

此设置由三部分组成:

  • Worker Loader:您部署的主 Worker。它接收请求,决定加载哪个 Dynamic Worker,并创建工作流实例。您编写此代码。
  • Dynamic Worker:定义工作流实际执行内容的每租户代码 — 其步骤、休眠和事件等待。每个 Dynamic Worker 都在运行时按需加载。
  • DynamicWorkflow 类:库创建的工作流入口点。当工作流引擎需要执行步骤时,此类会为该实例加载正确的 Dynamic Worker,并在其中运行该步骤。
Architecture

它们是这样协同工作的:

  • Worker Loader 接收请求,加载租户的 Dynamic Worker,并赋予其一个带有租户 ID 标记的工作流绑定。
  • Dynamic Worker 调用 env.WORKFLOWS.create() 启动新的工作流实例。租户 ID 自动与实例一起保存。
  • 工作流引擎运行在 Dynamic Worker 中定义的步骤 — step.do(), step.waitForEvent(), step.sleep()。每个步骤都是持久的:其结果被持久化,并在成功后不会重新运行。
  • 如果隔离区在步骤之间被回收(例如,在休眠期间或等待事件时),引擎会从实例中读取租户 ID,通过 Worker Loader 重新加载相同的 Dynamic Worker,并从停止的地方恢复。

该库提供了两个函数来处理 Worker Loader 和工作流引擎之间的连接,因此您不必手动标记请求、解析有效负载或编写自己的 WorkflowEntrypoint 子类。

  • wrapWorkflowBinding:创建一个标记有元数据(例如 { tenantId })的工作流绑定,您可以将其传递给 Dynamic Worker。该库将该元数据附加到 Dynamic Worker 创建的每个实例上,因此引擎可以将每个实例追溯到正确的租户。
  • createDynamicWorkflowEntrypoint:创建 DynamicWorkflow 类,当引擎恢复时,该类会重新加载正确的 Dynamic Worker。您向其提供一个获取元数据并返回租户的工作流类的回调,当需要运行步骤时,库会调用该回调。

安装库

该库处理 Worker Loader 和工作流引擎之间的连接,因此您无需手动标记请求、解析有效负载或编写自己的 WorkflowEntrypoint 子类。

npm i @cloudflare/dynamic-workflows

配置您的 Worker Loader

您的 Worker Loader 需要两个绑定

  • 一个 Worker Loader 绑定 (LOADER),以在运行时加载 Dynamic Workers。
  • 一个 工作流绑定 (WORKFLOWS),它指向 DynamicWorkflow 类。这是工作流引擎用来将每个实例路由到正确的 Dynamic Worker 的入口点。
{
  "$schema": "./node_modules/wrangler/config-schema.json",
  "name": "my-worker-loader",
  "main": "src/index.ts",
  // Set this to today's date
  "compatibility_date": "2026-08-17",
  "worker_loaders": [
    {
      "binding": "LOADER"
    }
  ],
  "workflows": [
    {
      "name": "dynamic-workflow",
      "binding": "WORKFLOWS",
      "class_name": "DynamicWorkflow"
    }
  ]
}
name = "my-worker-loader"
main = "src/index.ts"
# Set this to today's date
compatibility_date = "2026-08-17"

[[worker_loaders]]
binding = "LOADER"

[[workflows]]
name = "dynamic-workflow"
binding = "WORKFLOWS"
class_name = "DynamicWorkflow"

创建 Worker Loader

Worker Loader 是您将 Dynamic Workers 连接到工作流引擎的地方。在此文件中,您可以定义:

  • 如何加载租户的代码:一个接收租户 ID、获取其代码并赋予它们工作流绑定的函数。绑定是使用 wrapWorkflowBinding 创建的,它使用租户 ID 标记每个工作流实例,以便引擎稍后可以路由回正确的代码。

  • 引擎如何恢复工作流:使用 createDynamicWorkflowEntrypoint,您可以定义一个回调,引擎在需要运行步骤时调用该回调。回调从实例元数据接收租户 ID,并返回租户的工作流类。这就是持久执行跨隔离区重启起作用的原因 — 引擎知道如何重新加载正确的代码。

import {
	createDynamicWorkflowEntrypoint,
	DynamicWorkflowBinding,
	wrapWorkflowBinding,
} from "@cloudflare/dynamic-workflows";

// Required: re-exporting puts the class on cloudflare:workers exports,
// which is how wrapWorkflowBinding builds per-tenant RPC stubs.
export { DynamicWorkflowBinding };

function loadTenant(env, tenantId) {
	return env.LOADER.get(tenantId, async () => ({
		compatibilityDate: "2026-01-01",
		mainModule: "index.js",
		modules: { "index.js": await fetchTenantCode(tenantId) },
		// The Dynamic Worker uses this exactly like a real Workflow binding;
		// every create() is tagged with { tenantId } automatically.
		env: { WORKFLOWS: wrapWorkflowBinding({ tenantId }) },
	}));
}

// The entrypoint name must match `class_name` in the workflows binding of your Wrangler config file.
export const DynamicWorkflow = createDynamicWorkflowEntrypoint(
	async ({ env, metadata }) => {
		const stub = loadTenant(env, metadata.tenantId);
		return stub.getEntrypoint("TenantWorkflow");
	},
);

export default {
	fetch(request, env) {
		const tenantId = request.headers.get("x-tenant-id");
		return loadTenant(env, tenantId).getEntrypoint().fetch(request);
	},
};
import {
	createDynamicWorkflowEntrypoint,
	DynamicWorkflowBinding,
	wrapWorkflowBinding,
	type WorkflowRunner,
} from "@cloudflare/dynamic-workflows";

// Required: re-exporting puts the class on cloudflare:workers exports,
// which is how wrapWorkflowBinding builds per-tenant RPC stubs.
export { DynamicWorkflowBinding };

interface Env {
	WORKFLOWS: Workflow;
	LOADER: WorkerLoader;
}

function loadTenant(env: Env, tenantId: string) {
	return env.LOADER.get(tenantId, async () => ({
		compatibilityDate: "2026-01-01",
		mainModule: "index.js",
		modules: { "index.js": await fetchTenantCode(tenantId) },
		// The Dynamic Worker uses this exactly like a real Workflow binding;
		// every create() is tagged with { tenantId } automatically.
		env: { WORKFLOWS: wrapWorkflowBinding({ tenantId }) },
	}));
}

// The entrypoint name must match `class_name` in the workflows binding of your Wrangler config file.
export const DynamicWorkflow = createDynamicWorkflowEntrypoint<Env>(
	async ({ env, metadata }) => {
		const stub = loadTenant(env, metadata.tenantId as string);
		return stub.getEntrypoint("TenantWorkflow") as unknown as WorkflowRunner;
	},
);

export default {
	fetch(request: Request, env: Env) {
		const tenantId = request.headers.get("x-tenant-id")!;
		return loadTenant(env, tenantId).getEntrypoint().fetch(request);
	},
};

当请求到达时,会发生以下情况:

  1. fetch 处理程序从请求头中读取租户 ID。
  2. loadTenant 调用 env.LOADER.get() 来加载(或重用)该租户的 Dynamic Worker。Dynamic Worker 收到 WORKFLOWS: wrapWorkflowBinding({ tenantId }) 作为绑定,其外观和行为类似于普通的工作流绑定。
  3. 请求转发给 Dynamic Worker 的 fetch 处理程序,现在它可以调用 env.WORKFLOWS.create() 来启动一个工作流实例。

当该工作流实例稍后需要运行步骤时 — 例如,在 step.sleep() 之后或当新的隔离区接管它时 — 工作流引擎会调用 DynamicWorkflow 类上的 run()。该库从存储在实例上的元数据中读回 tenantId,并调用您传递给 createDynamicWorkflowEntrypoint 的回调。该回调加载该租户的 Dynamic Worker,并返回其 TenantWorkflow 类,以便引擎可以在原始代码中执行下一步。

编写 Dynamic Worker

Dynamic Worker 是您的用户编写的代码,它不需要了解任何关于路由层的信息。它是一个标准的工作流,照常使用 step.do(), step.sleep()step.waitForEvent() — 从它的角度来看,env.WORKFLOWS 是一个常规的工作流绑定。

import { WorkflowEntrypoint } from "cloudflare:workers";

export class TenantWorkflow extends WorkflowEntrypoint {
	async run(event, step) {
		return step.do("greet", async () => `Hello, ${event.payload.name}!`);
	}
}

export default {
	async fetch(request, env) {
		const instance = await env.WORKFLOWS.create({
			params: await request.json(),
		});
		// instance is an RPC stub — .id is an RpcPromise, so await it.
		return Response.json({ id: await instance.id });
	},
};
import { WorkflowEntrypoint } from "cloudflare:workers";

export class TenantWorkflow extends WorkflowEntrypoint {
	async run(event, step) {
		return step.do("greet", async () => `Hello, ${event.payload.name}!`);
	}
}

export default {
	async fetch(request, env) {
		const instance = await env.WORKFLOWS.create({
			params: await request.json(),
		});
		// instance is an RPC stub — .id is an RpcPromise, so await it.
		return Response.json({ id: await instance.id });
	},
};

普通工作流行为仍然适用。工作流 ID,.status().pause(),重试,休眠和持久步骤不受此架构影响。该库仅添加了 Worker Loader 和 Dynamic Worker 之间的路由。

触发动态工作流

向带有租户 ID 标头和 JSON 有效负载的 Worker Loader 发送 POST 请求。Worker Loader 加载匹配的 Dynamic Worker,它调用 env.WORKFLOWS.create() 并返回新的实例 ID。

curl -X POST http://localhost:8787/ \
  -H "x-tenant-id: tenant-42" \
  -H "Content-Type: application/json" \
  -d '{"name": "Alice"}'

检查工作流状态

使用上一个请求返回的实例 ID 来检查工作流状态。有关状态 API 的更多信息,请参考 Workers API 参考

curl "http://localhost:8787/api/status?instanceId=YOUR_INSTANCE_ID"

相关资源

这篇文档对您有帮助吗?