メインコンテンツへ移動

中断と再開

Workflow は任意の Step で一時停止し、追加データの収集、API コールバックの待機、高コストな処理の抑制、human-in-the-loop 入力の要求を行えます。Workflow が中断されると、現在の実行状態が Snapshot として保存されます。後から特定の Step ID から Workflow を再開し、その Snapshot に記録された状態を正確に復元できます。Snapshot は設定済みのストレージプロバイダーに保存され、デプロイやアプリケーションの再起動後も維持されます。

suspend() で Workflow を一時停止する
pausing-a-workflow-with-suspendへの直接リンク

特定の Step で Workflow の実行を一時停止するには、suspend() を使用します。Step の execute ブロックでは、resumeData の値を使用して中断条件を定義できます。

  • 条件を満たさない場合、Workflow は一時停止して suspend() を返します。
  • 条件を満たす場合、Workflow は Step 内の残りのロジックを続行します。

suspend() で Workflow を一時停止する

src/mastra/workflows/test-workflow.ts
const step1 = createStep({
id: 'step-1',
inputSchema: z.object({
userEmail: z.string(),
}),
outputSchema: z.object({
output: z.string(),
}),
resumeSchema: z.object({
approved: z.boolean(),
}),
execute: async ({ inputData, resumeData, suspend }) => {
const { userEmail } = inputData
const { approved } = resumeData ?? {}

if (!approved) {
return await suspend({})
}

return {
output: `Email sent to ${userEmail}`,
}
},
})

export const testWorkflow = createWorkflow({
id: 'test-workflow',
inputSchema: z.object({
userEmail: z.string(),
}),
outputSchema: z.object({
output: z.string(),
}),
})
.then(step1)
.commit()

resume() で Workflow を再開する
restarting-a-workflow-with-resumeへの直接リンク

一時停止した Step から中断中の Workflow を再開するには、resume() を使用します。Step の resumeSchema と一致する resumeData を渡して中断条件を満たし、実行を続けます。

resume() で Workflow を再開する

import { step1 } from './workflows/test-workflow'

const workflow = mastra.getWorkflow('testWorkflow')
const run = await workflow.createRun()

await run.start({
inputData: {
userEmail: 'alex@example.com',
},
})

const handleResume = async () => {
const result = await run.resume({
step: step1,
resumeData: { approved: true },
})
}

step オブジェクトを渡すと、resumeData に完全な型安全性が与えられます。一方、ID がユーザー入力やデータベースから得られる場合は、柔軟性を高めるために Step ID を渡すこともできます。

const result = await run.resume({
step: 'step-1',
resumeData: { approved: true },
})

中断中の Step が1つだけの場合、step 引数を完全に省略できます。Mastra は Workflow で最後に中断した Step を再開します。

runId だけを使用して再開する場合は、まず createRun() で実行インスタンスを作成します。

const workflow = mastra.getWorkflow('testWorkflow')
const run = await workflow.createRun({ runId: '123' })

const stream = run.resume({
resumeData: { approved: true },
})

resume() は、HTTP エンドポイント、イベントハンドラー、人による入力への応答、タイマーなど、アプリケーション内のどこからでも呼び出せます。

const midnight = new Date()
midnight.setUTCHours(24, 0, 0, 0)

setTimeout(async () => {
await run.resume({
step: 'step-1',
resumeData: { approved: true },
})
}, midnight.getTime() - Date.now())

suspendData で中断データにアクセスする
accessing-suspend-data-with-suspenddataへの直接リンク

Step が中断された後、再開時に suspend() へ渡したデータへアクセスしたい場合があります。このデータへアクセスするには、Step の execute 関数で suspendData パラメーターを使用します。

src/mastra/workflows/user-approval.ts
const approvalStep = createStep({
id: 'user-approval',
inputSchema: z.object({
requestId: z.string(),
}),
resumeSchema: z.object({
approved: z.boolean(),
}),
suspendSchema: z.object({
reason: z.string(),
requestDetails: z.string(),
}),
outputSchema: z.object({
result: z.string(),
}),
execute: async ({ inputData, resumeData, suspend, suspendData }) => {
const { requestId } = inputData
const { approved } = resumeData ?? {}

// On first execution, suspend with context
if (!approved) {
return await suspend({
reason: 'User approval required',
requestDetails: `Request ${requestId} pending review`,
})
}

// On resume, access the original suspend data
const suspendReason = suspendData?.reason || 'Unknown'
const details = suspendData?.requestDetails || 'No details'

return {
result: `${details} - ${suspendReason} - Decision: ${approved ? 'Approved' : 'Rejected'}`,
}
},
})

Step の再開時、suspendData パラメーターには最初の中断時に suspend() 関数へ渡したデータがそのまま自動設定されます。Workflow を中断した理由のコンテキストを維持し、再開処理でその情報を使用できます。

中断中の実行を識別する
中断中の実行を識別するへの直接リンク

Workflow が中断されると、一時停止した Step から再開します。Workflow の status で中断中であることを確認し、suspended を使用して一時停止した Step またはネストされた Workflow を識別できます。

const workflow = mastra.getWorkflow('testWorkflow')
const run = await workflow.createRun()

const result = await run.start({
inputData: {
userEmail: 'alex@example.com',
},
})

if (result.status === 'suspended') {
console.log(result.suspended[0])
await run.resume({
step: result.suspended[0],
resumeData: { approved: true },
})
}

出力例
出力例への直接リンク

suspended 配列には、実行中に中断した Workflow と Step の ID が含まれます。resume() の呼び出し時にこれらを step パラメーターへ渡すと、中断した実行経路を指定して再開できます。

['nested-workflow', 'step-1']

中断した実行を復元する
中断した実行を復元するへの直接リンク

アプリケーションでストレージから中断中の実行を復元する必要がある場合は、workflow.getWorkflowRunById()createWorkflowStateReader() を使用します。Reader は、生の Snapshot 構造を読み取ることなく、中断した Step、再開ラベル、Step のペイロード、Step の出力を公開します。

src/mastra/workflows/recover-run.ts
import { createWorkflowStateReader } from '@mastra/core/workflows'

const workflow = mastra.getWorkflow('testWorkflow')
const state = await workflow.getWorkflowRunById('run-123')

if (state?.status === 'suspended') {
const reader = createWorkflowStateReader(state)
const suspendedStep = reader.getSuspendedStep()
const approvalLabel = reader.getResumeLabel('approve')
const run = await workflow.createRun({ runId: state.runId })

await run.resume({
step: approvalLabel?.stepId ?? suspendedStep?.path,
resumeData: { approved: true },
forEachIndex: approvalLabel?.foreachIndex,
})
}

ネストされた Workflow では、suspendedStep.path に再開パスが含まれます。foreach の中断では、ラベルが特定の反復を指す場合、一致する再開ラベルに foreachIndex が含まれます。

Sleep
Sleepへの直接リンク

Sleep メソッドを使用すると、Workflow レベルで実行を一時停止でき、ステータスは waiting になります。一方、suspend() は特定の Step 内で実行を一時停止し、ステータスを suspended にします。

利用可能なメソッド:

  • .sleep():指定したミリ秒の間、一時停止します
  • .sleepUntil():指定した日時まで一時停止します