Clause AI 平台用 Mastra 协调四个专用 Agent 实现租约文档分析,完整处理了中途崩溃、部分写入、速率限制等失败场景,提供可重试的状态持久化方案。
构建一个 AI 功能很容易。构建一个可靠的多 Agent 管道,协调四个专业 AI Agent、持久化中间状态、在重试时跳过已完成的步骤、在模型思考时保持 API 响应——这才是困难的部分。
本文深入解析了 Clause AI 背后的架构。Clause AI 是一个分析租赁和 lease 协议的平台,它提取关键条款、标记风险条款,并允许用户使用 RAG(检索增强生成)技术与自己的合同对话——这一切都由一个协调的多 Agent 管道驱动。
构建 AI 驱动的文档分析工具,最 naive 的做法是一个函数调一次 LLM、解析输出、存入数据库。在它崩掉之前看起来一切正常。
一旦引入多个步骤——解析、摘要、嵌入、风险分析——事情就开始出问题:
Clause AI 从设计之初就考虑了这些失败模式,而不是事后补救。
系统没有用一个试图包揽一切的巨大 prompt,而是使用了四个专用 Agent——每个 Agent 都有明确的职责、调优的模型参数和结构化输出 schema。
每个 Agent 都配置了不同的推理层级和温度。Parser Agent 以低温度(0.2)运行,因为提取需要精确性——你希望对文档内容做确定性的忠实复现。Query Agent 温度更高(0.7),因为对话式响应受益于更自然的措辞。
每个 Agent 都使用 Zod schema 生成经过验证的结构化输出。以 Parser Agent 为例,它返回一个带 nullable 字段的 typed 对象——如果文档中缺少某个信息,Agent 返回 null 而不是虚构数据。
const ResponseSchema = z.object({
title: z.string().nullable(),
type: z.enum(AGREEMENT_TYPES).nullable(),
metadata: z
.object({
effectiveDate: z.string().nullable(),
expiryDate: z.string().nullable(),
autoRenewal: z.boolean().nullable(),
governingLaw: z.string().nullable(),
})
.nullable(),
parties: z
.array(
z.object({
name: z.string(),
role: z.enum(AGREEMENT_PARTY_ROLES),
address: z.string().nullable(),
}),
)
.nullable(),
sections: z
.array(
z.object({
ref: z.string(),
type: z.enum(SECTION_CLAUSE_TYPES),
heading: z.string(),
content: z.string(),
}),
)
.nullable(),
error: z.string().nullable(),
});
这种 schema 优先的方法意味着下游 Agent 和数据库写入可以信任所收到数据的形状。不需要防御性解析,不需要"希望 LLM 返回了正确的格式"——要么通过验证,要么直接失败。
编排层是这个架构复杂度预算的真正消耗点。项目使用 Mastra——一个 TypeScript 原生的 Agent 编排框架——来定义工作流为可组合的顺序管道,支持分支、迭代和共享状态。
以下是主工作流定义,精简到核心:
export const agentWorkflow = createWorkflow({
id: "agent-workflow",
inputSchema: z.any(),
outputSchema: z.object({ status: z.string() }),
stateSchema: WorkflowStateSchema,
})
.then(initiateStateHydration)
.branch([[async ({ state }) => !state.skipParserAgent, parsingWorkflow]])
.branch([[async ({ state }) => !state.skipSummaryAgent, summaryWorkflow]])
.then(embeddingWorkflow)
.branch([[async ({ state }) => !state.skipRiskAgent, riskWorkflow]])
.then(finishStep)
.commit();
.then() 调用将步骤串联执行。.branch() 调用根据运行时状态条件性地执行子工作流。容错性就在这里体现——下一节详细展开。
每个子工作流本身都是两步模式:LLM 调用 → 数据库持久化。
export const summaryWorkflow = createWorkflow({
id: "summary-workflow",
inputSchema: z.any(),
outputSchema: z.any(),
stateSchema: WorkflowStateSchema,
})
.then(summaryAgentStep) // LLM call: generate summary
.then(storeSummaryStep) // DB call: persist to Postgres
.commit();
这种分离是刻意为之的。LLM 步骤只写入工作流状态——一个内存中的临时对象。数据库步骤是独立的、可重试的操作。如果数据库写入失败,LLM 结果不会丢失——它存在于状态中,可以重试而无需重新运行昂贵的模型调用。
整个系统最关键的设计模式是状态水合(state hydration)——每次管道运行的第一步。
在任何 Agent 执行之前,工作流从数据库读取协议的当前状态。如果前一次运行已经完成了解析步骤(数据库中已存在 title、type、parties 和 sections),水合步骤在工作流状态中设置 skipParserAgent: true。主管道的 .branch() 看到这个标志就会完全跳过解析子工作流。
const hydrateWorkflowState = async (agreementId, userId, state) => {
const agreement = await AgreementsService.fetchAgreement(
agreementId,
userId,
);
const [dbSections, dbRisks] = await Promise.all([
AgreementsService.fetchSectionsByAgreement(agreementId, userId),
AgreementsService.fetchRisksByAgreement(agreementId, userId),
]);
return {
...state,
title: agreement.title,
sections: dbSections.length > 0 ? dbSections : state.sections,
risks: dbRisks.length > 0 ? dbRisks : state.risks,
// Skip flags based on what already exists
skipParserAgent: Boolean(
agreement.title &&
agreement.metadata &&
agreement.parties &&
dbSections.length > 0,
),
skipSummaryAgent: Boolean(agreement.summary),
skipRiskAgent: Boolean(dbRisks.length > 0),
};
};
崩溃恢复是零成本的。如果 Worker 在解析之后、摘要之前死亡,重试会从中断处继续。
不会浪费 API 调用。已完成的 LLM 步骤不会重新运行。
本质上是幂等的。对同一协议运行两次工作流产生相同结果,不会产生副作用。
还有一个 forceRestart 标志可以完全绕过水合步骤,当用户明确希望从头重新处理文档时很有用。
管道中的步骤并非都是简单的 A→B 链。嵌入和风险工作流使用 Mastra 的 .foreach() 原语将工作扇出到多个条目。
解析完成后,每个 section 需要一个向量嵌入用于语义搜索。嵌入工作流:
.foreach() 扇出,为每个 section 运行一个 per-section 子工作流export const embeddingWorkflow = createWorkflow({
id: "embedding-workflow",
inputSchema: z.any(),
outputSchema: z.any(),
stateSchema: WorkflowStateSchema,
})
.then(prepareEmbeddingSectionsStep)
.foreach(embeddingPerSectionWorkflow)
.commit();
风险工作流遵循类似的扇出模式,但有一个变化:分析前按条款类型对 section 分组。不是分析 15+ 个独立 section,而是将它们分组到逻辑类别(租金、终止、维护等)中,然后在一次 LLM 调用中分析每个组。这减少了 API 调用次数,同时保持每个 prompt 聚焦。
export const riskWorkflow = createWorkflow({
id: "risk-workflow",
inputSchema: z.any(),
outputSchema: z.any(),
stateSchema: WorkflowStateSchema,
})
.then(prepareRiskSectionsStep) // Group sections by type
.foreach(riskAnalyzeStep) // Analyze each group
.then(storeRiskResultStep) // Persist all results
.commit();
风险评分本身是刻意保守的。Risk Agent 的系统 prompt 明确说明"不存在风险是一个有效且预期的结果",并设置了很高的门槛:只标记"在法律备忘录中值得提及"的风险。每个被识别出的风险都有一个数字评分(0–100)映射到严重程度级别。
处理管道完成后,协议就准备好进行交互式问答了。Query Agent 在架构上与其他三个 Agent 不同——它是按需为每个用户问题运行的,而不是批量管道的一部分,并且使用 tool calling 来决定它需要什么信息。
该 Agent 可以访问两个工具:
fetchSectionsTool — 生成查询嵌入,对 pgvector 中协议 section 执行向量相似度搜索,在 token 预算内返回最相关的 sectionfetchRisksTool — 从数据库返回预先识别出的风险(无需嵌入步骤)这里关键的设计选择是:Agent 自己决定是否使用工具。对于可以用协议元数据(已注入系统 prompt)或对话历史回答的简单问题,不会调用工具。这保持了简单查询的快速响应。
// Token-budgeted retrieval instead of fixed top-K
const selectedSections = [];
for (const s of sections) {
const tokens = estimateTokens(s.content) + estimateTokens(s.heading);
if (tokens + usedTokens > maxTokens) break;
usedTokens += tokens;
selectedSections.push({
section: s.ref,
heading: s.heading,
content: s.content,
similarity: s.similarity,
});
}
sections 工具使用基于 token 预算的检索,而不是固定的 top-K 数量。由于 Query Agent 已经在上下文窗口中携带了对话历史和协议元数据,盲目返回 10 个 section 可能会撑爆上下文并降低响应质量。相反,section 会被添加直到达到 token 上限,不管最终是几个。
问答流程是全异步的。当用户发送问题时:
这使得 API 即使在模型需要几秒钟推理复杂问题时也能保持响应。
整个处理管道通过 BullMQ(Redis 后端的任务队列)与 API 层解耦。当用户上传文档时,API 不会inline 启动 AI 处理——而是将任务入队。
这解决了两类问题:
消息 worker 同时处理文件处理任务和邮件通知任务,根据 job name 路由:
if (job.name === PROCESS_FILE_JOB) {
await WorkflowService.startAgreementProcessing(agreementId, fileId, userId);
} else if (job.name === EMAIL_NOTIFICATION_JOB) {
await NotificationService.sendEmailNotification(email, type, payload);
}
每个 Agent 都配置了多个模型 fallback。如果主模型返回 429(限流),系统将其标记为在 retry-after header 指定的时长内不可用,并自动 fallback 到下一个可用模型。
const parserAgent = new Agent({
id: "parser-agent",
name: "Parser Agent",
instructions: Instructions,
model: getAvailableModels().map((model) => ({
id: model,
model: model,
modelSettings: {
reasoning: "low",
temperature: 0.2,
},
})),
});
这意味着管道不会因为临时限流而失败——它会优雅地降级到其他模型并继续处理。
从上传到交互式问答的完整管道:
每一步都是独立的、可重试的、幂等的。中间状态在步骤之间持久化。工作流可以中断并恢复,不会丢失进度或浪费 API 调用。
构建多 Agent 系统不是调用多个 LLM 那么简单——而是对它们的编排。真正的工程工作在于脚手架:
.foreach() 扇出处理可变长度的工作(section、风险组),无需硬编码批量大小多 Agent 方式不仅仅是一个架构选择——更是一个可靠性策略。每个 Agent 都有聚焦的职责、清晰的契约,以及一个不会污染管道其余部分的故障边界。