•21 min read

Temporalワークフローオーケストレーション: TypeScriptで堅牢な分散ステートマシンを構築する

Temporalワークフローオーケストレーション: TypeScriptで堅牢な分散ステートマシンを構築する

堅牢で長時間稼働する分散システムを構築する際には、状態管理、耐障害性、運用上の複雑さにおいて重大な課題が生じます。従来のアプローチでは、手動のリトライロジック、補償メカニズム、サービス境界を越えた永続的な状態管理によって複雑でエラーが発生しやすいコードベースになりがちでした。Temporal.io は、分散ステートマシンをオーケストレーションするための耐久性のある実行エンジンを提供することで、これらの複雑さを抽象化する重要なインフラストラクチャコンポーネントとして登場しました。このガイドでは、TypeScript を使用して Temporal を活用し、回復力の高いアプリケーションを構築するためのアーキテクチャ原則、実装パターン、および運用上の考慮事項について詳しく説明します。

Audio Briefing
0:00 / 0:00

分散ステートマシンの問題

分散システムは、本質的にネットワークパーティション、サービス障害、部分的な障害、予測不可能なレイテンシなどの課題に直面します。それぞれが独自の障害モードを持つ異なるマイクロサービス間で多段階のビジネスプロセスをオーケストレーションするには、原子性、一貫性、分離性、耐久性(ACID特性、またはその分散システム版であるBASE)を確保するための洗練されたメカニズムが必要です。

一般的な問題は次のとおりです。

  • 状態の喪失: サービスが処理中にクラッシュし、インメモリの状態が失われ、手動での回復または複雑な永続化レイヤーが必要になります。
  • 不完全なトランザクション: 多段階の操作が、一部のステップが完了した後、他のステップが完了する前に失敗し、システムの状態が不整合になります。
  • 手動のリトライとタイムアウト: 指数バックオフ、ジッター、サーキットブレーカーを備えた堅牢なリトライポリシーをサービス呼び出し全体に実装することは容易ではなく、しばしば重複します。
  • 補償ロジック: 部分的に完了した操作を元に戻す(例:在庫引き落としが失敗した後に支払いを払い戻す)には、明示的で複雑な「Saga」パターンが必要です。
  • 可観測性: 複数のサービスにわたる長時間実行プロセスの正確な状態と進行状況を追跡することは困難です。

メッセージキュー(例:Kafka、RabbitMQ)は、非同期通信とメッセージに対するある程度の耐久性を提供しますが、多段階プロセスの状態を本質的に管理したり、オーケストレーションロジック自体の耐久性のある実行を提供したりするわけではありません。これらの上に構築されたカスタムステートマシンは、Temporal がすぐに提供する機能の多くを再実装することがよくあります。

Advertisement

Temporalのアーキテクチャパラダイム

Temporal は、ワークフローの状態と実行ロジックを専用のフォールトトレラントなサービスに外部化することで、これらの課題に対処します。これにより、開発者は複雑で長時間実行されるビジネスプロセスを通常のコードとして記述でき、基盤となる分散システムのプリミティブを抽象化できます。

コアコンセプト

  1. コードとしてのワークフロー: Temporal ワークフローは、アクティビティをオーケストレーションする耐久性のあるフォールトトレラントな関数です。標準のアプリケーションコード(例:TypeScript)として記述され、数日、数週間、あるいは数年間一時停止しても、あたかも順次実行されているかのように見えます。
  2. 決定論的実行: これは Temporal の基礎です。ワークフローコードは決定論的でなければなりません。つまり、同じ入力と(その履歴からの)同じイベントシーケンスが与えられた場合、常に同じ出力と(アクティビティ、タイマーのスケジューリングなどの)同じコマンドシーケンスを生成する必要があります。これにより、Temporal はワーカーの障害後やバージョンアップグレード中にワークフローの実行履歴をリプレイして状態を回復できます。非決定論的な操作(例:Date.now()、Math.random()、直接I/O)はワークフローコード内では厳しく禁止されており、アクティビティ内にカプセル化する必要があります。
  3. イベントソーシングと履歴リプレイ: ワークフローによって発行されたすべての状態変更とコマンドは、その実行履歴にイベントとして記録されます。ワークフローを処理しているワーカーが失敗した場合、別のワーカーがワークフローを引き継ぎ、最初から履歴をリプレイしてその正確な状態を再構築し、失敗した時点から実行を続行できます。このメカニズムは、透過的なフォールトトレランスと耐久性を提供します。
  4. アクティビティ: アクティビティは、Temporal における非決定論的で副作用のある作業の単位です。これらは、外部システム(データベース、API、メッセージキュー、ファイルシステム)との相互作用をカプセル化します。アクティビティは冪等に設計されており、設定可能なポリシーで Temporal によって自動的にリトライできます。
  5. ワーカー: ワーカーは、ワークフローとアクティビティの実装をホストするアプリケーションプロセスです。Temporal クラスター上のタスクキューをポーリングし、タスクを実行し、結果を報告します。ワーカーはステートレスであり、水平方向にスケーリングできます。
  6. タスクキュー: ワークフローとアクティビティは、タスクキューを介して Temporal クラスターと通信します。ワークフローがアクティビティをスケジュールすると、特定の Activity Task Queue にタスクが配置されます。ワークフローが実行する必要がある場合、Workflow Task Queue にタスクが配置されます。ワーカーはこれらのキューを購読します。これにより、タスクプロデューサーとコンシューマーが分離されます。
  7. 可視性: Temporal クラスターは、実行中および完了したすべてのワークフローの状態、履歴、進行状況を検査するための API と UI(Temporal Web UI)を提供し、デバッグと運用監視に役立ちます。

TypeScript SDK

Temporal TypeScript SDK は、ワークフローとアクティビティを定義し、クライアント接続を管理し、ワーカーを実行するためのデコレータとユーティリティ関数を提供します。開発者エクスペリエンスの向上とコンパイル時の安全性のため、TypeScript の強力な型付けを活用しています。

回復力の高いワークフローの構築:実践的な例(注文処理)

Eコマースの注文処理プロセスを考えてみましょう。

  1. 支払い処理。
  2. 在庫引き落とし。
  3. 注文発送。
  4. 確認メール送信。

このプロセスは、どのステップでの障害にも耐え、一貫性を確保する必要があります。

プロジェクトのセットアップ

mkdir temporal-order-processing
cd temporal-order-processing
npm init -y
npm install @temporalio/client @temporalio/worker @temporalio/workflow @temporalio/testing typescript ts-node
npm install -D @types/node
npx tsc --init

tsconfig.json を次のように変更します。

{
  "compilerOptions": {
    "target": "es2020",
    "module": "commonjs",
    "rootDir": "./",
    "outDir": "./dist",
    "esModuleInterop": true,
    "forceConsistentCasingInFileNames": true,
    "strict": true,
    "skipLibCheck": true,
    "experimentalDecorators": true,
    "emitDecoratorMetadata": true
  }
}

アクティビティの定義 (activities.ts)

アクティビティは、外部の非決定論的な操作が行われる場所です。

// activities.ts
import { ApplicationFailure } from '@temporalio/workflow';

export async function processPayment(orderId: string, amount: number): Promise<string> {
  console.log(`Processing payment for order ${orderId}, amount ${amount}...`);
  // Simulate external payment gateway call
  await new Promise(resolve => setTimeout(resolve, 1000));
  if (Math.random() < 0.1) { // 10% chance of failure
    console.error(`Payment failed for order ${orderId}`);
    throw ApplicationFailure.create({
      message: `Payment gateway error for order ${orderId}`,
      nonRetryable: false, // Allow retries
    });
  }
  console.log(`Payment successful for order ${orderId}. Transaction ID: tx-${orderId}`);
  return `tx-${orderId}`;
}

export async function deductInventory(orderId: string, itemId: string, quantity: number): Promise<boolean> {
  console.log(`Deducting inventory for order ${orderId}, item ${itemId}, quantity ${quantity}...`);
  // Simulate external inventory service call
  await new Promise(resolve => setTimeout(resolve, 800));
  if (Math.random() < 0.05) { // 5% chance of failure
    console.error(`Inventory deduction failed for order ${orderId}`);
    throw ApplicationFailure.create({
      message: `Inventory service error for order ${orderId}`,
      nonRetryable: false,
    });
  }
  console.log(`Inventory deducted for order ${orderId}.`);
  return true;
}

export async function shipOrder(orderId: string, address: string): Promise<string> {
  console.log(`Shipping order ${orderId} to ${address}...`);
  // Simulate external shipping service call
  await new Promise(resolve => setTimeout(resolve, 1500));
  if (Math.random() < 0.02) { // 2% chance of failure
    console.error(`Shipping failed for order ${orderId}`);
    throw ApplicationFailure.create({
      message: `Shipping service error for order ${orderId}`,
      nonRetryable: false,
    });
  }
  console.log(`Order ${orderId} shipped. Tracking ID: trk-${orderId}`);
  return `trk-${orderId}`;
}

export async function sendConfirmationEmail(orderId: string, email: string): Promise<boolean> {
  console.log(`Sending confirmation email for order ${orderId} to ${email}...`);
  // Simulate email service call
  await new Promise(resolve => setTimeout(resolve, 500));
  console.log(`Confirmation email sent for order ${orderId}.`);
  return true;
}

export async function refundPayment(orderId: string, transactionId: string): Promise<boolean> {
  console.log(`Refunding payment for order ${orderId}, transaction ${transactionId}...`);
  await new Promise(resolve => setTimeout(resolve, 700));
  console.log(`Payment refunded for order ${orderId}.`);
  return true;
}

export async function restoreInventory(orderId: string, itemId: string, quantity: number): Promise<boolean> {
  console.log(`Restoring inventory for order ${orderId}, item ${itemId}, quantity ${quantity}...`);
  await new Promise(resolve => setTimeout(resolve, 600));
  console.log(`Inventory restored for order ${orderId}.`);
  return true;
}

ワークフローの定義 (workflows.ts)

ワークフローはアクティビティをオーケストレーションし、プロセス全体の状態を管理します。

// workflows.ts
import { proxyActivities, ApplicationFailure, workflowInfo } from '@temporalio/workflow';
import * as activities from './activities';

const {
  processPayment,
  deductInventory,
  shipOrder,
  sendConfirmationEmail,
  refundPayment,
  restoreInventory,
} = proxyActivities<typeof activities>({
  startToCloseTimeout: '1 minute', // Max time for an activity to complete
  retry: {
    initialInterval: '1 second',
    backoffCoefficient: 2,
    maximumInterval: '10 seconds',
    maximumAttempts: 5, // Retry up to 5 times
    nonRetryableErrorTypes: ['NonRetryableError'], // Custom non-retryable error
  },
});

export interface OrderDetails {
  orderId: string;
  amount: number;
  itemId: string;
  quantity: number;
  customerEmail: string;
  shippingAddress: string;
}

export async function orderFulfillmentWorkflow(orderDetails: OrderDetails): Promise<string> {
  const { orderId, amount, itemId, quantity, customerEmail, shippingAddress } = orderDetails;
  let paymentTransactionId: string | undefined;
  let inventoryDeducted = false;

  try {
    // Step 1: Process Payment
    console.log(`Workflow ${workflowInfo().workflowId}: Starting payment processing.`);
    paymentTransactionId = await processPayment(orderId, amount);
    console.log(`Workflow ${workflowInfo().workflowId}: Payment processed. Transaction ID: ${paymentTransactionId}`);

    // Step 2: Deduct Inventory
    console.log(`Workflow ${workflowInfo().workflowId}: Starting inventory deduction.`);
    await deductInventory(orderId, itemId, quantity);
    inventoryDeducted = true;
    console.log(`Workflow ${workflowInfo().workflowId}: Inventory deducted.`);

    // Step 3: Ship Order
    console.log(`Workflow ${workflowInfo().workflowId}: Starting order shipping.`);
    const trackingId = await shipOrder(orderId, shippingAddress);
    console.log(`Workflow ${workflowInfo().workflowId}: Order shipped. Tracking ID: ${trackingId}`);

    // Step 4: Send Confirmation Email (best effort, no compensation needed)
    console.log(`Workflow ${workflowInfo().workflowId}: Sending confirmation email.`);
    await sendConfirmationEmail(orderId, customerEmail);
    console.log(`Workflow ${workflowInfo().workflowId}: Confirmation email sent.`);

    return `Order ${orderId} fulfilled successfully. Tracking ID: ${trackingId}`;

  } catch (error) {
    console.error(`Workflow ${workflowInfo().workflowId}: Order fulfillment failed: ${error.message}`);

    // Compensation logic (Saga pattern)
    if (inventoryDeducted) {
      console.log(`Workflow ${workflowInfo().workflowId}: Compensating: Restoring inventory.`);
      await restoreInventory(orderId, itemId, quantity);
    }
    if (paymentTransactionId) {
      console.log(`Workflow ${workflowInfo().workflowId}: Compensating: Refunding payment.`);
      await refundPayment(orderId, paymentTransactionId);
    }

    throw ApplicationFailure.create({
      message: `Order fulfillment failed for ${orderId}: ${error.message}`,
      nonRetryable: true, // Mark workflow failure as non-retryable for the client
    });
  }
}

ワーカーの実装 (worker.ts)

ワーカープロセスは、ワークフローとアクティビティのコードをホストし、実行します。

// worker.ts
import { Worker } from '@temporalio/worker';
import * as activities from './activities';
import * as workflows from './workflows';

async function run() {
  const worker = await Worker.create({
    workflowsPath: require.resolve('./workflows'),
    activities,
    taskQueue: 'order-processing-task-queue',
  });

  console.log('Temporal Worker started...');
  await worker.run();
  console.log('Temporal Worker stopped.');
}

run().catch((err) => {
  console.error(err);
  process.exit(1);
});

クライアントとの対話 (client.ts)

クライアントはワークフローの実行を開始します。

// client.ts
import { Connection, Client } from '@temporalio/client';
import { orderFulfillmentWorkflow, OrderDetails } from './workflows';

async function startOrderWorkflow() {
  const connection = await Connection.connect({ address: 'localhost:7233' }); // Connect to Temporal Cluster
  const client = new Client({ connection });

  const orderDetails: OrderDetails = {
    orderId: `order-${Date.now()}`,
    amount: 100.00,
    itemId: 'SKU-123',
    quantity: 1,
    customerEmail: 'customer@example.com',
    shippingAddress: '123 Main St, Anytown, USA',
  };

  console.log(`Starting order fulfillment workflow for order ${orderDetails.orderId}...`);

  const handle = await client.workflow.start(orderFulfillmentWorkflow, {
    args: [orderDetails],
    taskQueue: 'order-processing-task-queue',
    workflowId: `order-fulfillment-${orderDetails.orderId}`,
  });

  console.log(`Workflow started, ID: ${handle.workflowId}, Run ID: ${handle.runId}`);

  // Wait for the workflow to complete
  try {
    const result = await handle.result();
    console.log(`Workflow completed: ${result}`);
  } catch (error) {
    console.error(`Workflow failed: ${error.message}`);
  }
}

startOrderWorkflow().catch((err) => {
  console.error(err);
  process.exit(1);
});

例の実行

  1. Temporal クラスターの起動: ローカル開発には Docker Compose を使用します。
    docker run --rm --name temporal-dev -p 7233:7233 -p 8080:8080 temporalio/auto-setup:1.20.0
    
  2. ワーカーの起動:
    npx ts-node worker.ts
    
  3. クライアントの起動(別のターミナルで):
    npx ts-node client.ts
    

ワーカーとクライアントからのログを観察してください。障害発生時にアクティビティがリトライされ、後のステップが失敗した場合に補償ロジックが実行されるのがわかります。

高度な Temporal パターン

分散トランザクションのためのSaga

注文処理の例は、基本的な Saga パターンを示しています。Saga は、各トランザクションが単一サービス内のデータを更新し、次のステップをトリガーするイベントを発行する一連のローカルトランザクションです。ステップが失敗した場合、先行する完了したステップを元に戻すために補償トランザクションが実行されます。Temporal ワークフローは、ワークフローの catch ブロック内で補償ロジックを直接定義できるため、Saga を本質的にサポートしています。

orderFulfillmentWorkflow は、どの補償アクティビティ(restoreInventory、refundPayment)を呼び出す必要があるかを判断するために、inventoryDeducted と paymentTransactionId フラグを明示的にチェックします。これにより、分散サービス全体での原子性が保証されます。

シグナル

シグナルは、実行中のワークフローに送信される非同期メッセージです。これにより、外部システムはワークフローの完了を待たずに、ワークフローと対話し、その状態を更新できます。

例:注文の承認

// workflows.ts (add to existing workflow file)
import { defineSignal, set  } from '@temporalio/workflow';

export const orderApprovalSignal = defineSignal<[boolean]>('orderApproval');

export async function orderFulfillmentWorkflow(orderDetails: OrderDetails): Promise<string> {
  // ... (existing code) ...

  let approved = false;
  set('approved', approved); // Set initial state for query

  // Wait for approval signal
  console.log(`Workflow ${workflowInfo().workflowId}: Waiting for order approval.`);
  await workflow.wait(async () => approved);
  console.log(`Workflow ${workflowInfo().workflowId}: Order approved, proceeding.`);

  // ... (rest of the workflow logic) ...
}
// client.ts (add to existing client file)
// ...
const handle = await client.workflow.start(orderFulfillmentWorkflow, { /* ... */ });

// Later, from another service or UI:
await handle.signal(orderApprovalSignal, true);
console.log(`Signal sent to approve workflow ${handle.workflowId}`);
// ...

クエリ

クエリは、実行中のワークフローに対して、その現在の状態を変更せずに取得するための同期呼び出しです。これらは読み取り専用操作です。

例:注文ステータスの取得

// workflows.ts (add to existing workflow file)
import { defineQuery, set, get } from '@temporalio/workflow';

export const getOrderStatusQuery = defineQuery<string>('getOrderStatus');

export async function orderFulfillmentWorkflow(orderDetails: OrderDetails): Promise<string> {
  // ... (existing code) ...

  let currentStatus = 'Initialized';
  set('currentStatus', currentStatus); // Set initial state

  workflow.set  (getOrderStatusQuery, () => get<string>('currentStatus'));

  try {
    currentStatus = 'Payment Processing';
    set('currentStatus', currentStatus);
    paymentTransactionId = await processPayment(orderId, amount);

    currentStatus = 'Inventory Deduction';
    set('currentStatus', currentStatus);
    await deductInventory(orderId, itemId, quantity);

    currentStatus = 'Shipping';
    set('currentStatus', currentStatus);
    const trackingId = await shipOrder(orderId, shippingAddress);

    currentStatus = 'Confirmation Sent';
    set('currentStatus', currentStatus);
    await sendConfirmationEmail(orderId, customerEmail);

    currentStatus = 'Completed';
    set('currentStatus', currentStatus);
    return `Order ${orderId} fulfilled successfully. Tracking ID: ${trackingId}`;

  } catch (error) {
    currentStatus = `Failed: ${error.message}`;
    set('currentStatus', currentStatus);
    // ... compensation logic ...
    throw error;
  }
}
// client.ts (add to existing client file)
// ...
const handle = await client.workflow.start(orderFulfillmentWorkflow, { /* ... */ });

// Later, to query status:
const status = await handle.query(getOrderStatusQuery);
console.log(`Current order status: ${status}`);
// ...

子ワークフロー

子ワークフローを使用すると、複雑なワークフローを、より小さく、管理しやすく、独立して実行可能な単位に分割できます。これにより、モジュール性、分離性が向上し、独自の再試行ポリシーとタイムアウトを設定できます。

例:個々の注文品目を子ワークフローとして処理する

// workflows.ts (add to existing workflow file)
import { proxyActivities, ApplicationFailure, workflowInfo, startChild } from '@temporalio/workflow';
import * as activities from './activities';

// Define a new child workflow
export async function processOrderItemWorkflow(orderId: string, itemId: string, quantity: number): Promise<boolean> {
  console.log(`Child Workflow ${workflowInfo().workflowId}: Processing item ${itemId} for order ${orderId}`);
  await activities.deductInventory(orderId, itemId, quantity);
  // Potentially more item-specific logic
  return true;
}

export async function orderFulfillmentWorkflow(orderDetails: OrderDetails): Promise<string> {
  const { orderId, amount, itemId, quantity, customerEmail, shippingAddress } = orderDetails;
  let paymentTransactionId: string | undefined;
  let itemProcessed = false; // Track if child workflow completed

  try {
    // ... process payment ...

    // Step 2: Deduct Inventory using a child workflow
    console.log(`Workflow ${workflowInfo().workflowId}: Starting child workflow for item processing.`);
    await startChild(processOrderItemWorkflow, {
      args: [orderId, itemId, quantity],
      workflowId: `${orderId}-item-${itemId}-processing`,
      taskQueue: 'order-processing-task-queue',
    });
    itemProcessed = true;
    console.log(`Workflow ${workflowInfo().workflowId}: Child workflow for item processing completed.`);

    // ... ship order, send email ...

    return `Order ${orderId} fulfilled successfully.`;

  } catch (error) {
    console.error(`Workflow ${workflowInfo().workflowId}: Order fulfillment failed: ${error.message}`);

    // Compensation logic
    if (itemProcessed) {
      console.log(`Workflow ${workflowInfo().workflowId}: Compensating: Restoring inventory via activity.`);
      await activities.restoreInventory(orderId, itemId, quantity);
    }
    if (paymentTransactionId) {
      console.log(`Workflow ${workflowInfo().workflowId}: Compensating: Refunding payment.`);
      await activities.refundPayment(orderId, paymentTransactionId);
    }

    throw ApplicationFailure.create({
      message: `Order fulfillment failed for ${orderId}: ${error.message}`,
      nonRetryable: true,
    });
  }
}

タイムアウトとリトライ

Temporal は、ワークフローとアクティビティの両方に対して、堅牢で設定可能なタイムアウトおよびリトライポリシーを提供します。

  • アクティビティオプション:

    • startToCloseTimeout: アクティビティ実行がワーカーによってピックアップされてから実行が許可される最大時間。
    • scheduleToCloseTimeout: アクティビティタスクがスケジュールされてから完了するまでの最大時間。キュー時間とリトライを含む。
    • scheduleToStartTimeout: アクティビティタスクがスケジュールされてからワーカーが処理を開始するまでの最大時間。
    • heartbeatTimeout: アクティビティが進行状況を報告しなければならない間隔。ハートビートが失われた場合、アクティビティはリトライされる。長時間実行されるアクティビティに役立つ。
    • retry: initialInterval、backoffCoefficient、maximumInterval、maximumAttempts、nonRetryableErrorTypes を含む設定可能なポリシー。
  • ワークフローオプション:

    • workflowExecutionTimeout: ワークフロー実行が実行を許可される最大時間。
    • workflowRunTimeout: 単一のワークフロー実行が実行を許可される最大時間(ワークフロー実行は、リトライまたは継続により複数の実行を持つことができる)。
    • workflowTaskTimeout: ワークフロータスクが実行を許可される最大時間。

これらのオプションは、一時的な障害を適切に処理し、無期限のブロックを防ぐことができる回復力の高いシステムを構築するために不可欠です。

Advertisement

Temporal ワークフローのテスト

@temporalio/testing パッケージは、実行中の Temporal クラスターを必要とせずに、ワークフローとアクティビティの単体テストと統合テストを行うためのローカルテスト環境を提供します。

// test/workflows.test.ts
import { TestWorkflowEnvironment } from '@temporalio/testing';
import { Worker } from '@temporalio/worker';
import { Connection, Client } from '@temporalio/client';
import { orderFulfillmentWorkflow, OrderDetails } from '../workflows';
import * as activities from '../activities';

let testEnv: TestWorkflowEnvironment;
let worker: Worker;
let client: Client;

beforeAll(async () => {
  testEnv = await TestWorkflowEnvironment.createTimeSkipping();
  client = new Client({ connection: testEnv.nativeConnection });

  worker = await Worker.create({
    connection: testEnv.nativeConnection,
    taskQueue: 'test-task-queue',
    workflowsPath: require.resolve('../workflows'),
    activities,
  });

  await worker.run(); // Start the worker in the background
}, 30000); // Increase timeout for setup

afterAll(async () => {
  await worker?.shutdown();
  await testEnv?.shutdown();
});

describe('orderFulfillmentWorkflow', () => {
  it('should successfully fulfill an order', async () => {
    const orderDetails: OrderDetails = {
      orderId: 'test-order-1',
      amount: 100,
      itemId: 'SKU-TEST-1',
      quantity: 1,
      customerEmail: 'test@example.com',
      shippingAddress: '123 Test St',
    };

    // Mock activities to ensure deterministic behavior and control outcomes
    const mockActivities = {
      processPayment: jest.fn(() => Promise.resolve('tx-test-1')),
      deductInventory: jest.fn(() => Promise.resolve(true)),
      shipOrder: jest.fn(() => Promise.resolve('trk-test-1')),
      sendConfirmationEmail: jest.fn(() => Promise.resolve(true)),
      refundPayment: jest.fn(() => Promise.resolve(true)),
      restoreInventory: jest.fn(() => Promise.resolve(true)),
    };

    // Override activities for this test run
    worker.activities = mockActivities;

    const result = await client.workflow.execute(orderFulfillmentWorkflow, {
      args: [orderDetails],
      workflowId: 'test-order-workflow-1',
      taskQueue: 'test-task-queue',
    });

    expect(result).toContain('fulfilled successfully');
    expect(mockActivities.processPayment).toHaveBeenCalledWith('test-order-1', 100);
    expect(mockActivities.deductInventory).toHaveBeenCalledWith('test-order-1', 'SKU-TEST-1', 1);
    expect(mockActivities.shipOrder).toHaveBeenCalledWith('test-order-1', '123 Test St');
    expect(mockActivities.sendConfirmationEmail).toHaveBeenCalledWith('test@example.com', 'test-order-1');
    expect(mockActivities.refundPayment).not.toHaveBeenCalled();
    expect(mockActivities.restoreInventory).not.toHaveBeenCalled();
  });

  it('should compensate if inventory deduction fails', async () => {
    const orderDetails: OrderDetails = {
      orderId: 'test-order-2',
      amount: 200,
      itemId: 'SKU-TEST-2',
      quantity: 2,
      customerEmail: 'test2@example.com',
      shippingAddress: '456 Test Ave',
    };

    const mockActivities = {
      processPayment: jest.fn(() => Promise.resolve('tx-test-2')),
      deductInventory: jest.fn(() => Promise.reject(new Error('Inventory unavailable'))),
      shipOrder: jest.fn(), // Should not be called
      sendConfirmationEmail: jest.fn(), // Should not be called
      refundPayment: jest.fn(() => Promise.resolve(true)),
      restoreInventory: jest.fn(() => Promise.resolve(true)),
    };

    worker.activities = mockActivities;
Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement