Skip to content

Performer

ジョブの処理内容はPerformerを継承したクラスとして記述します。

基本形

ts
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になり、別途の登録は不要です。

ts
// src/index.ts
export { SendMail } from './performers/send-mail.js';

defineTsumugiperformersには、performerをまとめたモジュールをそのまま渡します。
これはペイロードと必須キーの型を導出するためのものであり、実行時の解決には利用されません。

ts
import * as performers from './performers/index.js';

export * from './performers/index.js';

const tsumugi = defineTsumugi({ performers, /* ... */ });

別名を付ける場合はバレルのexportで変更してください。
実行時の解決先とperformersのキーが同じ1箇所から決定されるため、両者が一致します。

ts
// src/performers/index.ts
export { SendMail as MAIL } from './send-mail.js';

Workerのエントリでのみ別名を付けると、実行時はMAILで解決される一方でperformersのキーはSendMailのまま残り、型の上のbinding名と一致しなくなります。

ts
// src/index.ts
// 実行時のexportだけが変わるので、これだけでは足りない
export { SendMail as MAIL } from './performers/send-mail.js';

実行文脈

performの第2引数にJobContextが渡されます。

フィールド内容
jobId<binding>#<shard>:<localId>形式のジョブID
attempt1始まりの試行回数
idempotencyKeyジョブ単位で一定の値、再実行でも同じ値
deadlineAtタイムアウトが切れる時刻、epochミリ秒
heartbeat実行中であることを報告する関数
log処理途中のメッセージを保存する非同期関数
traceparent投入時のトレース情報。未指定の場合はnull
spawnFlowのノードとして実行中に子ノードを追加する関数

at-least-onceでは同じジョブが2回実行される場合があるため、外部への副作用はidempotencyKeyを使って冪等にしてください。

中断が必要なケースではdeadlineAtからAbortSignalを組み立てて渡します。
AbortSignalはRPCの引数として渡せない制約があるため、Tsumugiは時刻のみを渡します。

AbortSignal.timeoutは負の値を受け付けないため、期限を過ぎている場合はAbortSignal.abort()を利用してください。

ts
const remaining = ctx.deadlineAt - Date.now();
const signal = remaining > 0 ? AbortSignal.timeout(remaining) : AbortSignal.abort();

中断は協調的な要求であり、応じないperformerは期限を過ぎても実行を継続する可能性があります。

spawnの使用方法はFlowを参照してください。

heartbeat

所要時間が入力によって大きく変わる処理では、timeoutMsを最長の場合に合わせる必要があります。

ctx.heartbeat()を実行すると、無応答の判定が最後の報告時刻を起点として行われます。
timeoutMsは1回の報告間隔に対して設定してください。

ts
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が必要です。

ts
await ctx.log('Import started');
await importRows(payload.rows);
await ctx.log('Import completed');

ログは呼び出し時に保存され、ジョブの詳細画面に試行回数、記録時刻、本文が取得できます。
保存に成功した記録は、performerが異常終了した場合も残ります。

1ジョブあたり最新の20件まで保存され、超過した古い記録は削除されます。
本文はJavaScriptの文字列長で2,000文字まで制限され、超過分は切り捨てられます。
リトライした場合においても、保存済みのログは引き継がれ、件数はジョブ全体で制限されます。

保存間隔の下限/送信待機の上限は1,000msです。
保存に失敗した場合、console.errorへ出力され、ジョブの処理は継続します。
古い試行や終了したジョブからの呼び出しは保存されません。

トレース情報

直接のジョブ投入で指定したtraceparentctx.traceparentから取得可能です。未指定の場合はnullになります。
リトライ時も投入時と同値が参照されます。

トレースSDKを使用する場合は、この値を親のコンテキストとしてspanを作成してください。
投入時の形式はジョブの投入を参照してください。

失敗の通知

例外をthrowすると失敗として扱われ、試行回数が残っている場合に限りリトライされます。
8,192文字を超える戻り値と直列化できない値は保存されません。Flowのノードでの扱いはFlowを参照してください。

ts
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への投入時にキーの指定が必須になります。

ts
class ChargeCard extends Performer<Payload, void, { concurrencyKey: true }, Env> {}

キーは投入時に文字列として渡します。performer側の関数で導出することはできません。

必須化はtsumugi.enqueuetsumugi.jobs(env)で適用され、キーの渡し忘れはコンパイルエラーになります。

ts
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への配置が可能です。

ts
// 相手側の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を対応させます。

jsonc
"services": [
  { "binding": "MAIL", "service": "my-mailer", "entrypoint": "SendMail" },
],

呼び出し側のperformersには、クラスの代わりにremote()を置きます。

ts
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が終わると、その要求は失われます。

ts
await ctx.spawn('child', 'MAIL', payload);

テスト

tsumugi/testingはDurable ObjectとQueuesを起動せずにperformerを呼び出します。

ts
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は例外を再送出せず、成功と失敗を同じ形式で返します。

実行文脈を差し替える

ts
const ctx = createTestContext({ attempt: 3 });
await runPerformer(performer, payload, ctx);

deadlineAtを過去や近い将来に置くと、期限に対する振る舞いを検証できます。

ts
const ctx = createTestContext({ deadlineAt: Date.now() + 50 });
await runPerformer(performer, payload, ctx);

スケジューラとバックオフ

時刻と乱数は引数で渡すため、fake timersは不要です。

ts
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はこの範囲を扱いません。