Performer
ジョブの処理内容はPerformerを継承したクラスとして記述します。
基本形
import { Performer, type JobContext } from 'tsumugi/performer';
export class SendMail extends Performer<{ to: string; subject: string }, void, {}, Env> {
async perform(payload: { to: string; subject: string }, ctx: JobContext): Promise<void> {
const remaining = ctx.deadlineAt - Date.now();
await fetch('https://api.example.com/mail', {
method: 'POST',
body: JSON.stringify(payload),
// AbortSignal.timeoutは負の値を受け付けないため、期限を過ぎていれば即座に中断する
signal: remaining > 0 ? AbortSignal.timeout(remaining) : AbortSignal.abort(),
});
}
}型引数は順に、ペイロード、戻り値、必須キーの宣言、Envです。
bindingはWorkerEntrypointと同様にコンストラクタで受け取るため、this.envから参照可能です。
binding名
binding名はWorkerのエントリからexportした名前で解決され、export class SendMailと書けばbinding名はSendMailになり、別途の登録は不要です。
// src/index.ts
export { SendMail } from './performers/send-mail.js';defineTsumugiのperformersには、performerをまとめたモジュールをそのまま渡します。
これはペイロードと必須キーの型を導出するためのものであり、実行時の解決には利用されません。
import * as performers from './performers/index.js';
export * from './performers/index.js';
const tsumugi = defineTsumugi({ performers, /* ... */ });別名を付ける場合はバレルのexportで変更してください。
実行時の解決先とperformersのキーが同じ1箇所から決定されるため、両者が一致します。
// src/performers/index.ts
export { SendMail as MAIL } from './send-mail.js';Workerのエントリでのみ別名を付けると、実行時はMAILで解決される一方でperformersのキーはSendMailのまま残り、型の上のbinding名と一致しなくなります。
// src/index.ts
// 実行時のexportだけが変わるので、これだけでは足りない
export { SendMail as MAIL } from './performers/send-mail.js';実行文脈
performの第2引数にJobContextが渡されます。
| フィールド | 内容 |
|---|---|
jobId | <binding>#<shard>:<localId>形式のジョブID |
attempt | 1始まりの試行回数 |
idempotencyKey | ジョブ単位で一定の値、再実行でも同じ値 |
deadlineAt | タイムアウトが切れる時刻、epochミリ秒 |
heartbeat | 実行中であることを報告する関数 |
log | 処理途中のメッセージを保存する非同期関数 |
traceparent | 投入時のトレース情報。未指定の場合はnull |
spawn | Flowのノードとして実行中に子ノードを追加する関数 |
at-least-onceでは同じジョブが2回実行される場合があるため、外部への副作用はidempotencyKeyを使って冪等にしてください。
中断が必要なケースではdeadlineAtからAbortSignalを組み立てて渡します。AbortSignalはRPCの引数として渡せない制約があるため、Tsumugiは時刻のみを渡します。
AbortSignal.timeoutは負の値を受け付けないため、期限を過ぎている場合はAbortSignal.abort()を利用してください。
const remaining = ctx.deadlineAt - Date.now();
const signal = remaining > 0 ? AbortSignal.timeout(remaining) : AbortSignal.abort();中断は協調的な要求であり、応じないperformerは期限を過ぎても実行を継続する可能性があります。
spawnの使用方法はFlowを参照してください。
heartbeat
所要時間が入力によって大きく変わる処理では、timeoutMsを最長の場合に合わせる必要があります。
ctx.heartbeat()を実行すると、無応答の判定が最後の報告時刻を起点として行われます。timeoutMsは1回の報告間隔に対して設定してください。
class Import extends Performer<{ rows: string[] }, void, {}, Env> {
async perform(payload: { rows: string[] }, ctx: JobContext): Promise<void> {
for (const [index, row] of payload.rows.entries()) {
await store(row);
await ctx.heartbeat((index + 1) / payload.rows.length);
}
}
}引数の進捗は0以上1以下です。範囲外の値は0以上1以下へ丸められ、数値以外は進捗なしの報告として扱われます。
報告した進捗はジョブの詳細画面に表示されます。
実行間隔には5秒の下限があります。これより短い間隔で実行しても、報告は5秒に1回までに制限されます。
報告に失敗しても例外にはならず、報告が無いジョブと同様に扱われます。
ログ
ctx.log(message)で処理途中のメッセージを保存できます。呼び出しにはawaitが必要です。
await ctx.log('Import started');
await importRows(payload.rows);
await ctx.log('Import completed');ログは呼び出し時に保存され、ジョブの詳細画面に試行回数、記録時刻、本文が取得できます。
保存に成功した記録は、performerが異常終了した場合も残ります。
1ジョブあたり最新の20件まで保存され、超過した古い記録は削除されます。
本文はJavaScriptの文字列長で2,000文字まで制限され、超過分は切り捨てられます。
リトライした場合においても、保存済みのログは引き継がれ、件数はジョブ全体で制限されます。
保存間隔の下限/送信待機の上限は1,000msです。
保存に失敗した場合、console.errorへ出力され、ジョブの処理は継続します。
古い試行や終了したジョブからの呼び出しは保存されません。
トレース情報
直接のジョブ投入で指定したtraceparentはctx.traceparentから取得可能です。未指定の場合はnullになります。
リトライ時も投入時と同値が参照されます。
トレースSDKを使用する場合は、この値を親のコンテキストとしてspanを作成してください。
投入時の形式はジョブの投入を参照してください。
失敗の通知
例外をthrowすると失敗として扱われ、試行回数が残っている場合に限りリトライされます。
8,192文字を超える戻り値と直列化できない値は保存されません。Flowのノードでの扱いはFlowを参照してください。
class ChargeCard extends Performer<{ customerId: string; amountJpy: number }, void, { concurrencyKey: true }, Env> {
async perform(payload: { customerId: string; amountJpy: number }): Promise<void> {
const res = await this.env.PAYMENT.charge(payload);
if (!res.ok) throw new Error(`payment failed: ${res.status}`);
}
}発生した例外のメッセージは試行履歴に保存され、ダッシュボードの詳細画面に表示されます。
本文は2,000文字で打ち切られ、1ジョブあたり20件まで保持されます。
キーを必須にする
第3型引数に{ concurrencyKey: true }または{ uniqueKey: true }を指定すると、そのperformerへの投入時にキーの指定が必須になります。
class ChargeCard extends Performer<Payload, void, { concurrencyKey: true }, Env> {}キーは投入時に文字列として渡します。performer側の関数で導出することはできません。
必須化はtsumugi.enqueueとtsumugi.jobs(env)で適用され、キーの渡し忘れはコンパイルエラーになります。
await tsumugi.enqueue(env, { binding: 'CHARGE', payload, concurrencyKey: 'customer:c1' });
await tsumugi.jobs(env).enqueue('CHARGE', payload, { concurrencyKey: 'customer:c1' });WARNING
トップレベルのenqueue(env, input)とcreateClient()では必須化が適用されません。
投入経路を参照してください。
別Workerへの配置
performerはservice binding越しに別のWorkerへの配置が可能です。
// 相手側のWorker
import { Performer, type JobContext } from 'tsumugi/performer';
export class SendMail extends Performer<{ to: string; subject: string }, void, {}, Env> {
async perform(payload: { to: string; subject: string }, ctx: JobContext): Promise<void> {
// ...
}
}
// WorkerEntrypointの名前付きexportに加えてdefaultも必要
export default {
async fetch(): Promise<Response> {
return new Response('performer only', { status: 404 });
},
} satisfies ExportedHandler<Env>;wrangler.jsoncのservice bindingで、binding名とentrypointを対応させます。
"services": [
{ "binding": "MAIL", "service": "my-mailer", "entrypoint": "SendMail" },
],呼び出し側のperformersには、クラスの代わりにremote()を置きます。
import { remote } from 'tsumugi';
import type { SendMail } from 'my-mailer';
const performers = { ...local, MAIL: remote<SendMail>() };同一Workerのperformerと別Workerのperformerは混在可能です。
別Worker時の制約
spawnで要求した子ノードはperformが完了してから作成されます。呼び出した時点では実行されません。
別Workerではctx.spawnが非同期の呼び出しになるため、awaitが必要になります。awaitせずにperformが終わると、その要求は失われます。
await ctx.spawn('child', 'MAIL', payload);テスト
tsumugi/testingはDurable ObjectとQueuesを起動せずにperformerを呼び出します。
import { createTestContext, runPerformer } from 'tsumugi/testing';
declare const ctx: ExecutionContext;
declare const env: Env;
const result = await runPerformer(new SendWelcome(ctx, env), { userId: 'u_1' });
if (result.ok) console.log(result.value);
else console.error(result.error);PerformerのコンストラクタはWorkerEntrypointと同じ(ctx, env)の2引数です。runPerformerが要求するのはperformのみのため、this.envを使わないperformerはperformを持つ通常のオブジェクトでも検証できます。
runPerformerは例外を再送出せず、成功と失敗を同じ形式で返します。
実行文脈を差し替える
const ctx = createTestContext({ attempt: 3 });
await runPerformer(performer, payload, ctx);deadlineAtを過去や近い将来に置くと、期限に対する振る舞いを検証できます。
const ctx = createTestContext({ deadlineAt: Date.now() + 50 });
await runPerformer(performer, payload, ctx);スケジューラとバックオフ
時刻と乱数は引数で渡すため、fake timersは不要です。
import { fixedClock, nextAttempt, schedule } from 'tsumugi/testing';
// 3回目の再試行の時刻
nextAttempt({ attempts: 3, maxAttempts: 5, backoff, now: Date.now() });
// このポリシーで投入される件数
schedule({ now, jobs, policy, bucket });Durable Objectを経由するテスト
Durable ObjectとQueuesを経由した動作を検証する場合は@cloudflare/vitest-pool-workersが必要です。tsumugi/testingはこの範囲を扱いません。