LobeHub Upstash Workflow 实战指南:基于 QStash 的三层异步工作流架构与三大核心模式
LobeHub 使用 Upstash Workflow + QStash 在后台执行大量异步批量任务(评测运行、记忆提取、onboarding 处理等),并沉淀出一套标准化的"三层架构 + 三大模式"实现规范。读完本文,你将掌握 dry-run 预演、fan-out 分片扇出、单任务执行这三种模式的设计动机与完整代码模板,并能按 LobeHub 仓库中 agent-eval-run 等真实工作流的写法,独立开发一个具备幂等性、限流控制和可观测性的新工作流。
一、为什么是"三大模式":平台约束决定了设计
LobeHub 的每个 Upstash Workflow 都组合了三个核心模式。它们之所以存在,是因为平台在三个维度上约束了你:速率限制(rate limit)让盲目 fan-out 很危险,步骤上限(step limit)限制了单个工作流的规模,而幂等性(idempotency)要求重试不能造成重复处理。
- Dry-Run Mode(预演模式)——不触发实际执行,先拿到统计信息;
- Fan-Out Pattern(扇出模式)——把大批量拆成小分片并行处理;
- Single Task Execution(单任务执行)——每次工作流执行只处理一条数据。
完整的实现规范定义在仓库的 SKILL.md 中,配套的四个参考文档分别为 实现模板、最佳实践、真实案例 和 云项目部署。
二、三层架构总览
所有工作流遵循同一套 3 层架构:
Layer 1: Entry Point (process-*)
├─ 校验前置条件
├─ 计算待处理数据总量
├─ 过滤已处理数据
├─ 支持 dry-run 模式(仅返回统计)
└─ 有工作要处理时触发 Layer 2
Layer 2: Pagination (paginate-*)
├─ 处理基于游标(cursor)的分页
├─ 对大批量实现 fan-out
├─ 递归处理所有页
└─ 为每条数据触发 Layer 3
Layer 3: Single Task Execution (execute-* / generate-*)
└─ 执行单条数据的实际业务逻辑
仓库中的实际落地位置
从源码结构看,当前仓库的落地方式与规范文档略有演进:规范文档描述的是 Next.js App Router 的逐目录路由形态(src/app/(backend)/api/workflows/{name}/process-*/route.ts),而当前主服务已迁移到 Hono 服务端的集中式路由。对应关系如下:
- 统一入口:工作流捕获路由/api/workflows/%5B%5B...route%5D%5D/route.ts) 把所有
POST /api/workflows/*请求委托给 Hono 应用; - 路由注册:workflows Hono 根路由 用
new Hono().basePath('/api/workflows')挂载了agent-eval-run、memory-user-memory、agent-signal、task、goal、topic-auto-summary、verify等九个业务工作流; - Workflow 类:位于
apps/server/src/workflows/下,如 agentEvalRun; - 各层 Handler:位于
apps/server/src/router-hono/workflows/{name}/workflows/下,如 runBenchmark、paginateTestCases。
规范中给出的标准文件结构如下(新工作流在开源仓中按此组织):
src/
├── app/(backend)/api/workflows/
│ └── {workflow-name}/
│ ├── process-{entities}/route.ts # Layer 1
│ ├── paginate-{entities}/route.ts # Layer 2
│ └── execute-{entity}/route.ts # Layer 3
└── server/workflows/
└── {workflowName}/
└── index.ts # Workflow 类
规范文档推荐的延伸阅读导航(均已转换为本仓库路径):
| 你想做的事 | 参考文档 |
|---|---|
| 从零写 Workflow 类 + 三层路由 | references/implementation.md |
| 调优 flowControl、错误处理、日志与测试 | references/best-practices.md |
| 看两个端到端的真实工作流 | references/examples.md |
| 在 lobehub-cloud 上部署(re-export、云专属操作) | references/cloud.md |
三、三大模式详解(60 秒版)
1. Dry-Run Mode:在任何副作用之前短路
Layer 1 在产生任何副作用之前短路返回,让调用方能"预演"将要发生什么:
if (dryRun) {
return {
...result,
dryRun: true,
message: `[DryRun] Would process ${itemsNeedingProcessing.length} items`,
};
}
典型用途:在真正提交批量任务之前,先确认将要处理多少条数据。仓库中的 runBenchmark 就是这个写法的直接实例:
// If dryRun mode, return statistics only
if (dryRun) {
console.info('[run-benchmark] Dry run: %d test cases would execute', testCaseIds.length);
return {
...result,
dryRun: true,
message: `[DryRun] Would execute ${testCaseIds.length} test cases`,
};
}
2. Fan-Out Pattern:拆片扇出,绕开步骤上限
Layer 2 把超大批量切分成块(chunk),并用每个分片递归重新触发自身。这样做是为了在一页数据条数过多时避免触碰工作流的步骤上限:
const CHUNK_SIZE = 20;
if (itemIds.length > CHUNK_SIZE) {
const chunks = chunk(itemIds, CHUNK_SIZE);
await Promise.all(
chunks.map((ids, idx) =>
context.run(`workflow:fanout:${idx + 1}/${chunks.length}`, () =>
WorkflowClass.triggerPaginateItems({ itemIds: ids }),
),
),
);
}
默认常量:PAGE_SIZE = 50(每页条数)、CHUNK_SIZE = 20(每个扇出分片的条数)。这两个常量在当前仓库的 paginateTestCases 中一字不差地保留:
const CHUNK_SIZE = 20; // Max items to process directly
const PAGE_SIZE = 50; // Items per page
当 fan-out 分片直接带着 itemIds 回到 Layer 2 时,会跳过分页逻辑直接逐条触发 Layer 3(见 paginateTestCases L42-L57)。
3. Single Task Execution:一次执行只处理一条
Layer 3 每次调用恒定处理一条数据。并行度由 Layer 2 扇出到大量 Layer 3 调用来实现,并由 flowControl 控制节奏:
export const { POST } = serve<ExecutePayload>(
async (context) => {
const { itemId } = context.requestPayload ?? {};
if (!itemId) return { success: false, error: 'Missing itemId' };
const item = await context.run('workflow:get-item', () => getItem(itemId));
const result = await context.run('workflow:execute', () => processItem(item));
await context.run('workflow:save', () => saveResult(itemId, result));
return { success: true, itemId, result };
},
{
flowControl: { key: 'workflow.execute', parallelism: 10, ratePerSecond: 5 },
},
);
四、Workflow 类:触发器与过滤器的统一封装
规范模板(详见 implementation.md)要求每个工作流配一个 Workflow 类,包含三部分:WORKFLOW_PATHS 常量、各层 payload 类型、静态 trigger 方法,外加一个用于保持管道幂等的静态过滤方法:
import { Client } from '@upstash/workflow';
import debug from 'debug';
const log = debug('lobe-server:workflows:{workflow-name}');
// Workflow paths
const WORKFLOW_PATHS = {
processItems: '/api/workflows/{workflow-name}/process-items',
paginateItems: '/api/workflows/{workflow-name}/paginate-items',
executeItem: '/api/workflows/{workflow-name}/execute-item',
} as const;
// Payload types
export interface ProcessItemsPayload {
dryRun?: boolean;
force?: boolean;
}
export interface PaginateItemsPayload {
cursor?: string;
itemIds?: string[]; // For fanout chunks
}
export interface ExecuteItemPayload {
itemId: string;
}
const getWorkflowUrl = (path: string): string => {
const baseUrl = process.env.APP_URL;
if (!baseUrl) throw new Error('APP_URL is required to trigger workflows');
return new URL(path, baseUrl).toString();
};
const getWorkflowClient = (): Client => {
const token = process.env.QSTASH_TOKEN;
if (!token) throw new Error('QSTASH_TOKEN is required to trigger workflows');
const config: ConstructorParameters<typeof Client>[0] = { token };
if (process.env.QSTASH_URL) {
(config as Record<string, unknown>).url = process.env.QSTASH_URL;
}
return new Client(config);
};
export class {WorkflowName}Workflow {
private static client: Client;
private static getClient(): Client {
if (!this.client) this.client = getWorkflowClient();
return this.client;
}
static triggerProcessItems(payload: ProcessItemsPayload) {
const url = getWorkflowUrl(WORKFLOW_PATHS.processItems);
log('Triggering process-items workflow');
return this.getClient().trigger({ body: payload, url });
}
static triggerPaginateItems(payload: PaginateItemsPayload) {
const url = getWorkflowUrl(WORKFLOW_PATHS.paginateItems);
log('Triggering paginate-items workflow');
return this.getClient().trigger({ body: payload, url });
}
static triggerExecuteItem(payload: ExecuteItemPayload) {
const url = getWorkflowUrl(WORKFLOW_PATHS.executeItem);
log('Triggering execute-item workflow: %s', payload.itemId);
return this.getClient().trigger({ body: payload, url });
}
/**
* 过滤出仍需要处理的数据(如检查 Redis 缓存、数据库状态)。
* 只返回真正需要处理的那些 —— 让整条管道保持幂等。
*/
static async filterItemsNeedingProcessing(itemIds: string[]): Promise<string[]> {
if (itemIds.length === 0) return [];
// 检查已有状态,返回需要处理的数据
return itemIds;
}
}
对照真实实现 AgentEvalRunWorkflow,可以看到两点演进:
- 触发客户端不再在类内部惰性构造,而是统一复用全局的
workflowClient(来自 qstash 客户端封装,详见第七节); filterTestCasesNeedingExecution给出幂等过滤的完整示范:查询该 run 下所有 RunTopic,只保留status === 'pending'的用例 ID,其余一律跳过——这就是"重试不会重复处理"的落点。
五、三层路由完整模板
Layer 1:入口点(process-*)
职责:校验前置条件、计算统计信息、支持 dry-run 模式。
import { serve } from '@upstash/workflow/nextjs';
import { getServerDB } from '@/database/server';
import { WorkflowClass, type ProcessPayload } from '@/server/workflows/{workflowName}';
export const { POST } = serve<ProcessPayload>(
async (context) => {
const { dryRun, force } = context.requestPayload ?? {};
console.log('[{workflow}:process] Starting with payload:', { dryRun, force });
const allItemIds = await context.run('{workflow}:get-all-items', async () => {
const db = await getServerDB();
// 查询数据库中符合条件的项目
return items.map((item) => item.id);
});
console.log('[{workflow}:process] Total eligible items:', allItemIds.length);
if (allItemIds.length === 0) {
return { success: true, totalEligible: 0, message: 'No eligible items found' };
}
const itemsNeedingProcessing = await context.run('{workflow}:filter-existing', () =>
WorkflowClass.filterItemsNeedingProcessing(allItemIds),
);
const result = {
success: true,
totalEligible: allItemIds.length,
toProcess: itemsNeedingProcessing.length,
alreadyProcessed: allItemIds.length - itemsNeedingProcessing.length,
};
// Dry-run 在任何副作用前短路
if (dryRun) {
console.log('[{workflow}:process] Dry run mode, returning statistics only');
return {
...result,
dryRun: true,
message: `[DryRun] Would process ${itemsNeedingProcessing.length} items`,
};
}
if (itemsNeedingProcessing.length === 0) {
return { ...result, message: 'All items already processed' };
}
await context.run('{workflow}:trigger-paginate', () => WorkflowClass.triggerPaginateItems({}));
return {
...result,
message: `Triggered pagination for ${itemsNeedingProcessing.length} items`,
};
},
{
flowControl: {
key: '{workflow}.process',
parallelism: 1, // 单实例 —— 避免重复处理
ratePerSecond: 1,
},
},
);
当前仓库中 runBenchmarkHandler 的执行顺序与模板完全一致:校验 runId/userId → 检查 run 状态(running 且未 force 时拒绝)→ 拉取全部测试用例 → filter-existing 过滤 → dry-run 短路 → 无待执行项提前返回 → 更新 run 状态为 running → 触发 paginate 层。
Layer 2:分页(paginate-*)
职责:处理基于游标的分页,对大批量实现 fan-out。
import { serve } from '@upstash/workflow/nextjs';
import { chunk } from 'es-toolkit/compat';
import { getServerDB } from '@/database/server';
import { WorkflowClass, type PaginatePayload } from '@/server/workflows/{workflowName}';
const PAGE_SIZE = 50;
const CHUNK_SIZE = 20;
export const { POST } = serve<PaginatePayload>(
async (context) => {
const { cursor, itemIds: payloadItemIds } = context.requestPayload ?? {};
console.log('[{workflow}:paginate] Starting:', {
cursor,
itemIdsCount: payloadItemIds?.length ?? 0,
});
// 若传入了特定 itemIds(来自 fanout 分片),直接处理
if (payloadItemIds && payloadItemIds.length > 0) {
await Promise.all(
payloadItemIds.map((itemId) =>
context.run(`{workflow}:execute:${itemId}`, () =>
WorkflowClass.triggerExecuteItem({ itemId }),
),
),
);
return { success: true, processedItems: payloadItemIds.length };
}
// 逐页遍历所有数据
const itemBatch = await context.run('{workflow}:get-batch', async () => {
const db = await getServerDB();
const items = await db.query(...);
if (!items.length) return { ids: [] };
const last = items.at(-1);
return {
ids: items.map((item) => item.id),
cursor: last ? last.id : undefined,
};
});
const batchItemIds = itemBatch.ids;
const nextCursor = 'cursor' in itemBatch ? itemBatch.cursor : undefined;
if (batchItemIds.length === 0) {
return { success: true, message: 'Pagination complete' };
}
const itemIds = await context.run('{workflow}:filter-existing', () =>
WorkflowClass.filterItemsNeedingProcessing(batchItemIds),
);
if (itemIds.length > 0) {
if (itemIds.length > CHUNK_SIZE) {
// Fan out —— 用每个分片递归重入分页层
const chunks = chunk(itemIds, CHUNK_SIZE);
console.log('[{workflow}:paginate] Fanout mode:', {
chunks: chunks.length,
chunkSize: CHUNK_SIZE,
});
await Promise.all(
chunks.map((ids, idx) =>
context.run(`{workflow}:fanout:${idx + 1}/${chunks.length}`, () =>
WorkflowClass.triggerPaginateItems({ itemIds: ids }),
),
),
);
} else {
// 直接处理这一页
await Promise.all(
itemIds.map((itemId) =>
context.run(`{workflow}:execute:${itemId}`, () =>
WorkflowClass.triggerExecuteItem({ itemId }),
),
),
);
}
}
// 尾调用进入下一页
if (nextCursor) {
await context.run('{workflow}:next-page', () =>
WorkflowClass.triggerPaginateItems({ cursor: nextCursor }),
);
}
return {
success: true,
processedItems: itemIds.length,
skippedItems: batchItemIds.length - itemIds.length,
nextCursor: nextCursor ?? null,
};
},
{
flowControl: {
key: '{workflow}.paginate',
parallelism: 20,
ratePerSecond: 5,
},
},
);
实际的 paginateTestCasesHandler 还示范了一个规范未展开但生产中必要的细节:分页前先检查 run 是否被中止(runStatus === 'aborted' 时返回 { cancelled: true }),避免用户取消后分页层仍在盲目触发下游任务。
Layer 3:执行(execute-* / generate-*)
职责:对恰好一条数据执行实际业务逻辑。
import { serve } from '@upstash/workflow/nextjs';
import { getServerDB } from '@/database/server';
import { WorkflowClass, type ExecutePayload } from '@/server/workflows/{workflowName}';
export const { POST } = serve<ExecutePayload>(
async (context) => {
const { itemId } = context.requestPayload ?? {};
if (!itemId) {
return { success: false, error: 'Missing itemId' };
}
const db = await getServerDB();
const item = await context.run('{workflow}:get-item', async () => {
// 查询数据库获取该条目
return item;
});
if (!item) {
return { success: false, error: 'Item not found' };
}
const result = await context.run('{workflow}:process-item', async () => {
const workflow = new WorkflowClass(db, itemId);
return workflow.generate(); // 或 process()、execute() 等
});
await context.run('{workflow}:save-result', async () => {
const workflow = new WorkflowClass(db, itemId);
return workflow.saveToRedis(result); // 或 saveToDatabase() 等
});
return { success: true, itemId, result };
},
{
flowControl: {
key: '{workflow}.execute',
parallelism: 10,
ratePerSecond: 5,
},
},
);
一个实现细节值得注意:当前仓库用 runStep 辅助函数 统一包装 context.run() 调用(runStep(context, 'step-name', fn)),所有层的步骤名依然保持"工作流前缀 + 语义名"的约定,与规范中的 context.run() 用法等价。
六、最佳实践与常见坑
以下内容来自 best-practices.md,是脚手架落地后必须逐条对照的清单。
1. 错误处理:显式失败,而不是异常炸掉工作流
export const { POST } = serve<Payload>(
async (context) => {
const { itemId } = context.requestPayload ?? {};
if (!itemId) {
return { success: false, error: 'Missing itemId in payload' };
}
try {
const result = await context.run('step-name', () => doWork(itemId));
return { success: true, itemId, result };
} catch (error) {
console.error('[workflow:error]', error);
return {
success: false,
error: error instanceof Error ? error.message : 'Unknown error',
};
}
},
{ flowControl: { ... } },
);
2. 日志:统一前缀方便跨 QStash 仪表盘与 grep 排查
console.log('[{workflow}:{layer}] Starting with payload:', payload);
console.log('[{workflow}:{layer}] Processing items:', { count: items.length });
console.log('[{workflow}:{layer}] Completed:', result);
console.error('[{workflow}:{layer}:error]', error);
仓库实现中则使用 debug 包(如 debug('lobe-server:workflows:run-benchmark'))作为结构化替代,前缀约定不变。
3. 返回值:形状要匹配层的职责
入口层返回统计信息,执行层返回单条结果:
// 成功
return { success: true, itemId, result, message: 'Optional success message' };
// 错误
return { success: false, error: 'Error description', itemId };
// 统计(入口点)
return {
success: true,
totalEligible: 100,
toProcess: 80,
alreadyProcessed: 20,
dryRun: true, // 如适用
message: 'Summary message',
};
4. flowControl:按层调并发
// Layer 1: 入口 —— 单实例,避免重复处理
flowControl: { key: '{workflow}.process', parallelism: 1, ratePerSecond: 1 }
// Layer 2: 分页 —— 中等并发
flowControl: { key: '{workflow}.paginate', parallelism: 20, ratePerSecond: 5 }
// Layer 3: 执行 —— 更高并发以并行处理条目
flowControl: { key: '{workflow}.execute', parallelism: 10, ratePerSecond: 5 }
为什么是这些默认值:
- Layer 1 恒用
parallelism: 1,保证并发触发不会让两个实例启动同一批任务; - Layer 2 可以较宽地扇出(10–20),因为分页本身开销小;
- Layer 3 默认封顶 5–10,依据外部 API 的速率限制上下调。
5. context.run() 的黄金法则:步骤名必须唯一
Upstash 按步骤名(step name)做去重(dedupe)——这是"重试不重复执行"的机制本身,但也意味着重复步骤名会静默丢数据:
// ✅ 唯一步骤名
await Promise.all(
items.map((item) => context.run(`{workflow}:execute:${item.id}`, () => processItem(item))),
);
// ❌ 同名步骤 —— Upstash 按步骤名去重,你会丢失数据
await Promise.all(items.map((item) => context.run(`{workflow}:execute`, () => processItem(item))));
其他要求:步骤名带前缀且描述性强({workflow}:step-name);每个步骤保持幂等(可安全重试);不要嵌套 context.run(),保持扁平。
6. Payload 校验:快速失败
在函数顶部校验,让失败显式化,而不是让 undefined 级联成后期难以定位的故障:
const { itemId, configId } = context.requestPayload ?? {};
if (!itemId) return { success: false, error: 'Missing itemId in payload' };
if (!configId) return { success: false, error: 'Missing configId in payload' };
7. 数据库连接:每个工作流只取一次
getServerDB() 是异步的,在每个步骤里重复获取会白白增加延迟:
const db = await getServerDB();
const item = await context.run('get-item', () => itemModel.findById(db, itemId));
const result = await context.run('save-result', () => resultModel.create(db, result));
8. 测试:dry-run 路径与完整路径都要覆盖
describe('WorkflowName', () => {
it('should process items successfully', async () => {
const items = await createTestItems();
await WorkflowClass.triggerProcessItems({ dryRun: false });
await waitForCompletion();
const results = await getResults();
expect(results).toHaveLength(items.length);
});
it('should support dryRun mode', async () => {
const result = await WorkflowClass.triggerProcessItems({ dryRun: true });
expect(result).toMatchObject({
success: true,
dryRun: true,
totalEligible: expect.any(Number),
toProcess: expect.any(Number),
});
});
});
常见坑速查
| 坑 | 后果 | 正确做法 |
|---|---|---|
复用 context.run() 步骤名 |
Upstash 按名去重,批量处理只剩一条 | 步骤名带上条目 ID:`process:${item.id}` |
| 跳过 payload 校验 | undefined 级联成后期难懂的失败 |
顶部 fail fast,返回 { success: false, error } |
| 跳过 filter 步骤 | 已处理的数据被重复处理 | 先 filterExisting(allItems) 再触发执行 |
| 日志前缀不一致 | 无法按 workflow+layer grep 定位 | 统一 [workflow:layer] 前缀 |
七、环境配置与 QStash 客户端
环境变量
所有工作流依赖以下环境变量(SKILL.md 的"Environment Variables"一节):
# 所有工作流必需
APP_URL=https://your-app.com # 工作流端点的基础 URL
QSTASH_TOKEN=qstash_xxx # QStash 认证令牌
# 可选(自定义 QStash URL)
QSTASH_URL=https://custom-qstash.com
APP_URL 与 QSTASH_TOKEN 的强校验体现在 Workflow 类的 getWorkflowUrl / 客户端构造中:缺少任一项会直接抛出 'APP_URL is required to trigger workflows' / 'QSTASH_TOKEN is required to trigger workflows' 错误(见 agentEvalRun/index.ts L131-L135)。
仓库中统一封装的客户端
当前仓库不鼓励每个 Workflow 类各自 new Client,而是复用 src/libs/qstash 中的单例:
qstashClient:OtelQstashClient,作为 Upstash Workflowserve()的选项传入,携带 Vercel Deployment Protection 的 bypass 头(VERCEL_AUTOMATION_BYPASS_SECRET),并对每次publishJSON记录 OTEL 指标;workflowClient:OtelWorkflowClient,用于workflowClient.trigger()触发工作流,成功/失败都会经recordUpstashWorkflowEvent上报(含workflowRunId、status、重试参数等);verifyQStashSignature:用QSTASH_CURRENT_SIGNING_KEY/QSTASH_NEXT_SIGNING_KEY校验 QStash 回调签名(未配置签名密钥时跳过校验),配套中间件见 qstashAuth。
Hono 路由的注册方式(serve 来自 @upstash/workflow/hono,并叠加 OTEL 包装)可参考 agent-eval-run 路由:
app.post(
'/run-benchmark',
serve(
withOtelMetricsForUpstashWorkflows(runBenchmarkHandler, {
url: '/api/workflows/agent-eval-run/run-benchmark',
}),
{ ...runBenchmarkWorkflowOptions, qstashClient },
),
);
八、仓库内的端到端实例
examples.md 给出了两个"逐字遵循本模式"的工作流(位于 lobehub-cloud 私有仓,作为模式参照):
实例 1:welcome-placeholder(为用户生成 AI 欢迎占位内容)
- Layer 1
process-users:入口,检查符合条件的用户; - Layer 2
paginate-users:分页遍历活跃用户; - Layer 3
generate-user:为单个用户生成占位内容。
关键特性:过滤掉 Redis 中已有缓存占位内容的用户;paidOnly 标记限定订阅用户;dryRun 统计模式;大批量用户走 fan-out(CHUNK_SIZE=20)。Layer 3 形态:
export const { POST } = serve<GenerateUserPlaceholderPayload>(async (context) => {
const { userId } = context.requestPayload ?? {};
const workflow = new WelcomePlaceholderWorkflow(db, userId);
const placeholders = await context.run('generate', () => workflow.generate());
return { success: true, userId, placeholdersCount: placeholders.length };
});
实例 2:agent-welcome(为 AI Agent 生成欢迎语与开场问题)
三层结构同构(process-agents / paginate-agents / generate-agent),同样具备 Redis 缓存过滤、paidOnly、dryRun 与 fan-out。
两个工作流的差异仅在于:实体类型(users vs agents)、业务逻辑(占位内容生成 vs 欢迎语生成)、数据源(不同的数据库查询)。其余一切——三层拆分、dry-run 处理、fan-out、filter-existing、flowControl 调参——完全相同。这正是该模式的意义:一旦内化该模式,新增工作流基本就是"实体替换"。
在开源仓内,agent-eval-run(评测基准运行)与 memory-user-memory(用户记忆提取)是可直接阅读的端到端实现,目录分别位于 agent-eval-run 与 memory-user-memory。
九、云项目(lobehub-cloud)部署模式
如果工作流要在 lobehub-cloud 私有仓中部署,遵循 cloud.md 的约定:
两种归属模式
- Cloud-Only 工作流:功能专属云用户(AI 生成、付费特性),直接实现在
lobehub-cloud/src/app/(backend)/api/workflows/下,无需 re-export; - Re-export 工作流:核心实现在开源 lobehub,云端只需转发。原因是云部署必须 serve 这些端点,而 lobehub 子模块代码不能被云路由直接访问;同时为未来的云专属覆写留了口子。
Re-export 的正确写法
// lobehub-cloud/src/app/(backend)/api/workflows/feature/layer/route.ts
export { POST } from 'lobehub/src/app/(backend)/api/workflows/feature/layer/route';
关键点:必须使用 lobehub/src/... 路径,不能用 @/...,否则会触发 Circular definition of import alias 'POST' 循环导入错误。
TypeScript 路径映射
云项目通过 tsconfig paths 实现"云代码优先、开源回退"的解析顺序:
// lobehub-cloud/tsconfig.json
{
"compilerOptions": {
"paths": {
"@/*": ["./src/*", "./lobehub/src/*"]
}
}
}
解析顺序:先查 ./src/*(云代码),再回退 ./lobehub/src/*(开源),允许云端覆写特定模块而其余部分沿用 lobehub 默认。
归属决策与迁移
- 放到 lobehub(开源):功能对所有用户有用、无专有业务逻辑、可以开源;
- 放到 cloud(私有):付费/高级特性、依赖云专属服务、包含专有算法;
- 迁移工作流从云到开源的五步:拷贝路由 → 去除云专属依赖(换成通用接口)→ 云端建 re-export → Workflow 类移入
apps/server/src/workflows/→ 更新云端 import 为lobehub/apps/server/src/workflows/feature。
十、新工作流开发 Checklist
以下清单完整继承自 SKILL.md 的"Checklist for New Workflows"一节,建议在新建工作流时逐项打勾。
规划
- [ ] 明确要处理的实体(users、agents、items、…)
- [ ] 定义单条数据的业务逻辑
- [ ] 确定过滤逻辑(Redis 缓存、数据库状态、…)
实现
- [ ] 用 TypeScript interface 定义各层 payload 类型
- [ ] 创建带静态 trigger 方法的 Workflow 类
- [ ] Layer 1:入口路由,支持 dry-run
- [ ] Layer 1:过滤逻辑,避免重复工作
- [ ] Layer 2:带 fan-out 的分页
- [ ] Layer 3:单任务执行(每次运行一条)
- [ ] 为每一层配置合适的
flowControl - [ ] 使用带工作流前缀的一致化日志
- [ ] 校验所有必需的 payload 参数
- [ ] 保证
context.run()步骤名唯一
质量与部署
- [ ] 返回一致的响应形状
- [ ] 配置云部署(如在 lobehub-cloud,参见 references/cloud.md)
- [ ] 编写集成测试(
dryRun路径 + 完整执行路径) - [ ] 先用 dry-run 做冒烟测试
- [ ] 全量上线前用小批量试跑验证
结语
LobeHub 的 Upstash Workflow 规范本质上是对"平台三大约束"(限流、步骤上限、幂等)的一组固定应答:入口层单实例 + dry-run 预演守住幂等与可观测性,分页层用游标 + fan-out 把任意规模的批量拆到步骤上限以内,执行层恒定单条 + flowControl 把并发压进外部 API 的速率预算。三层各自独立、可单独重试与限流,组合起来就是一条可水平扩展的后台处理管道。仓库中的 agent-eval-run、memory-user-memory 等实现都是这套模式的直接产物,新工作流只需按本文的模板与 Checklist 完成"实体替换"即可上线。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00