Workflows API
Workflows API 提供用于与 Mastra 中的自动化 Workflow 交互并执行它们的方法。
获取所有 Workflow获取所有 Workflow的直接链接
检索所有可用 Workflow 的列表:
const workflows = await mastraClient.listWorkflows()
获取 Workflow run 数量获取 Workflow run 数量的直接链接
通过单个请求检索每个 Workflow 的 running 和 suspended run 数量。数量在服务器上计算,并以 Workflow 的注册表键为键。该键是指在 Mastra 配置中注册 Workflow 时使用的键,可能与 Workflow 自身的 id 不同:
const runCounts = await mastraClient.listWorkflowRunCounts()
// { "cityWorkflow": { running: 2, suspended: 1 }, ... }
返回:Record<string, { running: number; suspended: number }>
服务器可能会在两次请求之间缓存数量数秒。早于此端点的服务器会响应 404 Not Found;当客户端可能与较旧部署通信时,请处理此错误。
使用特定 Workflow使用特定 Workflow的直接链接
通过 ID 获取特定 Workflow 的实例:
export const testWorkflow = createWorkflow({
id: 'city-workflow',
})
const workflow = mastraClient.getWorkflow('city-workflow')
Workflow 方法Workflow 方法的直接链接
details()details的直接链接
检索 Workflow 的详细信息:
const details = await workflow.details()
createRun()createrun的直接链接
创建新的 Workflow run 实例:
const run = await workflow.createRun()
// Or with an existing runId
const run = await workflow.createRun({ runId: 'existing-run-id' })
// Or with a resourceId to associate the run with a specific resource
const run = await workflow.createRun({
runId: 'my-run-id',
resourceId: 'user-123',
})
resourceId 参数将 Workflow run 与特定资源(例如用户 ID、租户 ID)关联。此值会随 run 一起持久化,之后可用于筛选和查询 run。
startAsync()startasync的直接链接
启动 Workflow run 并等待其完成,将完整结果作为 Workflow 输出返回。
const run = await workflow.createRun()
const result = await run.startAsync({
inputData: {
city: 'New York',
},
})
你也可以传入 initialState 来设置 Workflow state 的初始值:
const result = await run.startAsync({
inputData: {
city: 'New York',
},
initialState: {
count: 0,
items: [],
},
})
initialState 对象应与 Workflow stateSchema 中定义的结构匹配。更多详情请参阅 Workflow State。
要将 run 与特定资源关联,请向 createRun() 传入 resourceId:
const run = await workflow.createRun({ resourceId: 'user-123' })
const result = await run.startAsync({
inputData: {
city: 'New York',
},
})
start()start的直接链接
启动 Workflow run 而不等待其完成(触发后即不再等待)。它会立即返回成功消息。之后可在 Workflow 实例上使用 runById() 检查结果:
const run = await workflow.createRun()
await run.start({
inputData: {
city: 'New York',
},
})
// Poll for results later
const result = await workflow.runById(run.runId)
这适用于需要启动执行并在之后检查结果的长时间运行 Workflow。
resumeAsync()resumeasync的直接链接
恢复暂停的 Workflow 步骤并等待完整结果:
const run = await workflow.createRun({ runId: prevRunId })
const result = await run.resumeAsync({
step: 'step-id',
resumeData: { key: 'value' },
})
resume()resume的直接链接
恢复暂停的 Workflow 步骤而不等待其完成:
const run = await workflow.createRun({ runId: prevRunId })
await run.resume({
step: 'step-id',
resumeData: { key: 'value' },
})
当 .foreach() 步骤跨多个迭代暂停时,传入 forEachIndex(从零开始;0 对应第一次迭代),一次恢复一个迭代。未指定的迭代会保持暂停状态。
await run.resume({
step: 'approve',
resumeData: { ok: true },
forEachIndex: 1, // resumes the second iteration
})
resumeAsync() 和 resumeStream() 也支持 forEachIndex。
cancel()cancel的直接链接
取消正在运行的 Workflow:
const run = await workflow.createRun({ runId: existingRunId })
const result = await run.cancel()
// Returns: { message: 'Workflow run canceled' }
此方法会停止所有正在运行的步骤,并阻止后续步骤执行。检查 abortSignal 参数的步骤可以通过清理资源(超时、网络请求等)来响应取消操作。
有关取消的工作原理以及如何编写响应取消操作的步骤,请参阅 Run.cancel() 参考。
stream()stream的直接链接
流式传输 Workflow 执行以获取实时更新:
const run = await workflow.createRun()
const stream = await run.stream({
inputData: {
city: 'New York',
},
})
for await (const chunk of stream) {
console.log(JSON.stringify(chunk, null, 2))
}
runById()runbyid的直接链接
获取 Workflow run 的执行结果:
const result = await workflow.runById(runId)
// Or with options for performance optimization:
const result = await workflow.runById(runId, {
fields: ['status', 'result'], // Only fetch specific fields
withNestedWorkflows: false, // Skip expensive nested workflow data
requestContext: { userId: 'user-123' }, // Optional request context
})
Run 结果格式
Workflow run 结果包含以下内容:
runId:
eventTimestamp:
payload:
动态 Workflow动态 Workflow的直接链接
动态 Workflow 处于 beta 阶段。在 API 稳定之前,可能会发生不伴随主版本升级的破坏性变更。
动态 Workflow 是以 JSON 表示的 Workflow 定义。服务器会持久化每个定义,并将其注册为可运行的 Workflow。有关定义格式,请参阅动态 Workflow。
listDynamicWorkflows()listdynamicworkflows的直接链接
列出动态 Workflow 定义,并可选择按 status('active' | 'archived')和 authorId 筛选:
const { definitions, total } = await mastraClient.listDynamicWorkflows({
status: 'active',
})
upsertDynamicWorkflow()upsertdynamicworkflow的直接链接
创建或替换动态 Workflow 定义。服务器会验证并持久化定义,然后实时注册它以供执行:
const stored = await mastraClient.upsertDynamicWorkflow({
id: 'greeting-workflow',
description: 'Returns a greeting for the supplied name',
inputSchema: {
type: 'object',
properties: { name: { type: 'string' } },
required: ['name'],
},
outputSchema: {
type: 'object',
properties: { message: { type: 'string' } },
required: ['message'],
},
graph: [
{
type: 'mapping',
id: 'create-greeting',
mapConfig: JSON.stringify({
message: { template: 'Hello, ${initData.name}!' },
}),
},
],
})
当根定义嵌套了尚不存在的辅助 Workflow 时,请通过 dependencies 在同一请求中传入这些 Workflow。服务器会将整个 bundle 作为一个单元进行验证和注册,并通过 dependencyIds 回显辅助 Workflow 的 ID:
const stored = await mastraClient.upsertDynamicWorkflow({
id: 'root-workflow',
// ...schemas and graph referencing 'helper-workflow'...
dependencies: [helperDefinition],
})
console.log(stored.dependencyIds) // ['helper-workflow']
getDynamicWorkflow()getdynamicworkflow的直接链接
获取用于定义管理的动态 Workflow 实例。要执行动态 Workflow,请像处理其他 Workflow 一样使用 getWorkflow(id).createRun():
const dynamicWorkflow = mastraClient.getDynamicWorkflow('greeting-workflow')
dynamicWorkflow.details()dynamicworkflowdetails的直接链接
检索持久化的定义,包括 schema、graph、状态和时间戳:
const definition = await dynamicWorkflow.details()
dynamicWorkflow.delete()dynamicworkflowdelete的直接链接
删除存储的定义并注销实时 Workflow:
await dynamicWorkflow.delete()
执行动态 Workflow执行动态 Workflow的直接链接
注册后,动态 Workflow 将通过常规 Workflow API 运行:
const workflow = mastraClient.getWorkflow('greeting-workflow')
const run = await workflow.createRun()
const result = await run.startAsync({ inputData: { name: 'Ada' } })
调度调度的直接链接
调度通过 createWorkflow 上的 schedule 字段在代码中声明。client SDK 提供读取和操作方法,用于在运行时管理 Workflow 调度。请参阅调度的 Workflow。
createSchedule()createschedule的直接链接
通过传入 workflowId 创建 Workflow 调度。
const schedule = await mastraClient.createSchedule({
workflowId: 'daily-report',
cron: '0 9 * * *',
inputData: { reportType: 'summary' },
})
listSchedules()listschedules的直接链接
列出 Workflow 调度,并可选择按 Workflow ID 或状态筛选。
const schedules = await mastraClient.listSchedules({
workflowId: 'daily-report',
status: 'active',
})
getSchedule()getschedule的直接链接
按 ID 获取单个 Workflow 调度。
const schedule = await mastraClient.getSchedule('daily-report')
updateSchedule()updateschedule的直接链接
更新 Workflow 调度。
const updated = await mastraClient.updateSchedule('daily-report', {
cron: '0 10 * * *',
inputData: { reportType: 'summary' },
})
deleteSchedule()deleteschedule的直接链接
删除 Workflow 调度。
await mastraClient.deleteSchedule('daily-report')
runSchedule()runschedule的直接链接
立即触发一次 Workflow 调度,而不更改其 cron 周期。
const run = await mastraClient.runSchedule('daily-report')
pauseSchedule()pauseschedule的直接链接
暂停调度,使调度器停止触发它。返回更新后的调度。
await mastraClient.pauseSchedule('daily-report')
resumeSchedule()resumeschedule的直接链接
恢复暂停的调度。下一次触发时间将从当前时间重新计算,因此暂停很久的调度不会触发积压任务。返回更新后的调度。
await mastraClient.resumeSchedule('daily-report')
listScheduleTriggers()listscheduletriggers的直接链接
列出 Workflow 调度的触发历史记录,包括每次触发所关联的 run 摘要。
const { triggers } = await mastraClient.listScheduleTriggers('daily-report', {
limit: 50,
})