跳到主要内容

Workflows API

Workflows API 提供用于与 Mastra 中的自动化 Workflow 交互并执行它们的方法。

获取所有 Workflow
获取所有 Workflow的直接链接

检索所有可用 Workflow 的列表:

const workflows = await mastraClient.listWorkflows()

获取 Workflow run 数量
获取 Workflow run 数量的直接链接

通过单个请求检索每个 Workflow 的 runningsuspended 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 的实例:

src/mastra/workflows/test-workflow.ts
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:

string
此 Workflow run 实例的唯一标识符

eventTimestamp:

Date
事件的时间戳

payload:

object
包含 currentStep(id、status、output、payload)和 workflowState(status、steps 记录)

动态 Workflow
动态 Workflow的直接链接

beta

动态 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,
})