•18 min read

Điều phối Workflow Temporal: Xây dựng máy trạng thái phân tán bền vững trong TypeScript

Điều phối Workflow Temporal: Xây dựng máy trạng thái phân tán bền vững trong TypeScript

Việc xây dựng các hệ thống phân tán mạnh mẽ, hoạt động lâu dài đặt ra những thách thức đáng kể trong quản lý trạng thái, khả năng chịu lỗi và độ phức tạp trong vận hành. Các phương pháp truyền thống thường dẫn đến các codebase phức tạp, dễ xảy ra lỗi, bị gánh nặng bởi logic thử lại thủ công, cơ chế bù trừ và quản lý trạng thái bền vững trên các ranh giới dịch vụ. Temporal.io nổi lên như một thành phần cơ sở hạ tầng quan trọng, trừu tượng hóa những phức tạp này bằng cách cung cấp một công cụ thực thi bền vững để điều phối các máy trạng thái phân tán. Hướng dẫn này trình bày chi tiết các nguyên tắc kiến trúc, các mẫu triển khai và các cân nhắc vận hành để tận dụng Temporal với TypeScript nhằm xây dựng các ứng dụng có khả năng phục hồi cao.

Audio Briefing
0:00 / 0:00

Vấn đề máy trạng thái phân tán

Các hệ thống phân tán vốn phải đối mặt với các thách thức như phân vùng mạng, lỗi dịch vụ, lỗi một phần và độ trễ không thể đoán trước. Việc điều phối các quy trình kinh doanh nhiều bước trên các microservice khác nhau, mỗi microservice có các chế độ lỗi riêng, đòi hỏi các cơ chế tinh vi để đảm bảo tính nguyên tử, tính nhất quán, tính cô lập và tính bền vững (các thuộc tính ACID, hoặc tương đương phân tán của chúng, BASE).

Các vấn đề thường gặp bao gồm:

  • Mất trạng thái: Một dịch vụ gặp sự cố giữa chừng, làm mất trạng thái trong bộ nhớ và yêu cầu khôi phục thủ công hoặc các lớp bền vững phức tạp.
  • Giao dịch không hoàn chỉnh: Một hoạt động nhiều bước thất bại sau khi một số bước hoàn thành nhưng trước khi các bước khác hoàn thành, dẫn đến trạng thái hệ thống không nhất quán.
  • Thử lại và hết thời gian thủ công: Việc triển khai các chính sách thử lại mạnh mẽ với exponential backoff, jitter và circuit breaker trên các cuộc gọi dịch vụ là không hề đơn giản và thường bị trùng lặp.
  • Logic bù trừ: Hoàn tác các hoạt động đã hoàn thành một phần (ví dụ: hoàn tiền thanh toán sau khi trừ kho thất bại) yêu cầu các mẫu "saga" rõ ràng, thường phức tạp.
  • Khả năng quan sát: Việc theo dõi trạng thái chính xác và tiến trình của một quy trình chạy dài trên nhiều dịch vụ là khó khăn.

Các hàng đợi tin nhắn (ví dụ: Kafka, RabbitMQ) cung cấp giao tiếp không đồng bộ và một mức độ bền vững cho tin nhắn, nhưng chúng không tự quản lý trạng thái của một quy trình nhiều bước hoặc cung cấp khả năng thực thi bền vững của chính logic điều phối. Các máy trạng thái tùy chỉnh được xây dựng trên các hệ thống này thường triển khai lại nhiều thứ mà Temporal cung cấp sẵn.

Advertisement

Mô hình kiến trúc của Temporal

Temporal giải quyết những thách thức này bằng cách đưa trạng thái quy trình làm việc và logic thực thi ra bên ngoài một dịch vụ chuyên dụng, có khả năng chịu lỗi. Nó cho phép các nhà phát triển viết các quy trình kinh doanh phức tạp, chạy dài dưới dạng mã thông thường, trừu tượng hóa các nguyên thủy hệ thống phân tán cơ bản.

Các khái niệm cốt lõi

  1. Workflow as Code: Một Temporal Workflow là một hàm bền vững, có khả năng chịu lỗi, điều phối các Activity. Nó được viết dưới dạng mã ứng dụng tiêu chuẩn (ví dụ: TypeScript) và dường như thực thi tuần tự, ngay cả khi nó tạm dừng trong nhiều ngày, nhiều tuần hoặc nhiều năm.
  2. Thực thi xác định: Đây là nền tảng của Temporal. Mã Workflow phải xác định, nghĩa là với cùng một đầu vào và cùng một chuỗi sự kiện (từ lịch sử của nó), nó phải luôn tạo ra cùng một đầu ra và cùng một chuỗi lệnh (ví dụ: lên lịch hoạt động, bộ hẹn giờ). Điều này cho phép Temporal phát lại lịch sử thực thi quy trình làm việc để khôi phục trạng thái sau khi worker gặp lỗi hoặc trong quá trình nâng cấp phiên bản. Các hoạt động không xác định (ví dụ: Date.now(), Math.random(), I/O trực tiếp) bị nghiêm cấm trong mã quy trình làm việc; chúng phải được đóng gói trong các Activity.
  3. Event Sourcing & History Replay: Mọi thay đổi trạng thái và lệnh được phát ra bởi một Workflow đều được ghi lại dưới dạng một sự kiện trong lịch sử thực thi của nó. Nếu một worker xử lý một Workflow gặp lỗi, một worker khác có thể tiếp tục Workflow, phát lại lịch sử của nó từ đầu và xây dựng lại trạng thái chính xác của nó, sau đó tiếp tục thực thi từ điểm lỗi. Cơ chế này cung cấp khả năng chịu lỗi và độ bền minh bạch.
  4. Activities: Activities là các đơn vị công việc không xác định, có tác dụng phụ trong Temporal. Chúng đóng gói các tương tác với các hệ thống bên ngoài (cơ sở dữ liệu, API, hàng đợi tin nhắn, hệ thống tệp). Activities được thiết kế để có tính bất biến và có thể được Temporal tự động thử lại với các chính sách có thể cấu hình.
  5. Workers: Workers là các quy trình ứng dụng lưu trữ các triển khai Workflow và Activity. Chúng thăm dò Task Queues trên Temporal Cluster, thực thi các tác vụ và báo cáo kết quả. Workers không trạng thái và có thể được mở rộng theo chiều ngang.
  6. Task Queues: Workflows và Activities giao tiếp với Temporal Cluster thông qua Task Queues. Khi một Workflow lên lịch một Activity, một tác vụ được đặt vào một Activity Task Queue cụ thể. Khi một Workflow cần thực thi, một tác vụ được đặt vào một Workflow Task Queue. Workers đăng ký các hàng đợi này. Điều này tách rời các nhà sản xuất tác vụ khỏi người tiêu dùng.
  7. Visibility: Temporal Cluster cung cấp các API và giao diện người dùng (Temporal Web UI) để kiểm tra trạng thái, lịch sử và tiến trình của tất cả các quy trình làm việc đang chạy và đã hoàn thành, hỗ trợ gỡ lỗi và giám sát hoạt động.

TypeScript SDK

Temporal TypeScript SDK cung cấp các decorator và hàm tiện ích để định nghĩa các quy trình làm việc và hoạt động, quản lý kết nối máy khách và chạy các worker. Nó tận dụng tính năng gõ mạnh của TypeScript để cải thiện trải nghiệm nhà phát triển và an toàn thời gian biên dịch.

Xây dựng một Workflow có khả năng phục hồi: Một ví dụ thực tế (Xử lý đơn hàng)

Hãy xem xét quy trình thực hiện đơn hàng thương mại điện tử:

  1. Xử lý thanh toán.
  2. Trừ kho.
  3. Giao hàng.
  4. Gửi xác nhận.

Quy trình này phải có khả năng phục hồi trước các lỗi ở bất kỳ bước nào và đảm bảo tính nhất quán.

Thiết lập dự án

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

Sửa đổi tsconfig.json để bao gồm:

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

Định nghĩa Activity (activities.ts)

Activities là nơi xảy ra các hoạt động bên ngoài, không xác định.

// 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;
}

Định nghĩa Workflow (workflows.ts)

Workflow điều phối các hoạt động và quản lý trạng thái quy trình tổng thể.

// 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
    });
  }
}

Triển khai Worker (worker.ts)

Quy trình worker lưu trữ và thực thi mã workflow và activity.

// 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);
});

Tương tác với Client (client.ts)

Client khởi tạo việc thực thi workflow.

// 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);
});

Chạy ví dụ

  1. Khởi động Temporal Cluster: Sử dụng Docker Compose để phát triển cục bộ:
    docker run --rm --name temporal-dev -p 7233:7233 -p 8080:8080 temporalio/auto-setup:1.20.0
    
  2. Khởi động Worker:
    npx ts-node worker.ts
    
  3. Khởi động Client (trong một terminal riêng):
    npx ts-node client.ts
    

Quan sát nhật ký từ worker và client. Bạn sẽ thấy các hoạt động thử lại khi gặp lỗi và logic bù trừ được thực thi nếu một bước sau đó thất bại.

Các mẫu Temporal nâng cao

Sagas cho các giao dịch phân tán

Ví dụ về thực hiện đơn hàng minh họa một mẫu Saga cơ bản. Saga là một chuỗi các giao dịch cục bộ trong đó mỗi giao dịch cập nhật dữ liệu trong một dịch vụ duy nhất và xuất bản một sự kiện để kích hoạt bước tiếp theo. Nếu một bước thất bại, các giao dịch bù trừ sẽ được thực thi để hoàn tác các bước đã hoàn thành trước đó. Các workflow của Temporal vốn hỗ trợ Sagas bằng cách cho phép bạn định nghĩa logic bù trừ trực tiếp trong khối catch của workflow.

orderFulfillmentWorkflow kiểm tra rõ ràng các cờ inventoryDeducted và paymentTransactionId để xác định các hoạt động bù trừ nào (restoreInventory, refundPayment) cần được gọi. Điều này đảm bảo tính nguyên tử trên các dịch vụ phân tán.

Signals

Signals là các tin nhắn không đồng bộ được gửi đến một workflow đang chạy. Chúng cho phép các hệ thống bên ngoài tương tác và cập nhật trạng thái của một workflow mà không cần chờ hoàn thành.

Ví dụ: Phê duyệt đơn hàng

// 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}`);
// ...

Queries

Queries là các cuộc gọi đồng bộ đến một workflow đang chạy để truy xuất trạng thái hiện tại của nó mà không làm thay đổi nó. Chúng là các hoạt động chỉ đọc.

Ví dụ: Lấy trạng thái đơn hàng

// 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}`);
// ...

Child Workflows

Child Workflows cho phép chia nhỏ các workflow phức tạp thành các đơn vị nhỏ hơn, dễ quản lý và có thể thực thi độc lập. Chúng cung cấp tính mô đun, cô lập tốt hơn và có thể có các chính sách thử lại và thời gian chờ riêng.

Ví dụ: Xử lý các mặt hàng đơn hàng riêng lẻ dưới dạng child workflows

// 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,
    });
  }
}

Timeouts và Retries

Temporal cung cấp các chính sách thời gian chờ và thử lại mạnh mẽ, có thể cấu hình cho cả Workflows và Activities.

  • Tùy chọn Activity:

    • startToCloseTimeout: Thời gian tối đa cho phép một Activity Execution chạy sau khi nó đã được một Worker nhận.
    • scheduleToCloseTimeout: Thời gian tối đa từ khi một Activity Task được lên lịch cho đến khi nó hoàn thành. Bao gồm thời gian trong hàng đợi và các lần thử lại.
    • scheduleToStartTimeout: Thời gian tối đa từ khi một Activity Task được lên lịch cho đến khi một Worker bắt đầu xử lý nó.
    • heartbeatTimeout: Khoảng thời gian mà một Activity phải báo cáo tiến độ của nó. Nếu một heartbeat bị bỏ lỡ, Activity sẽ được thử lại. Hữu ích cho các hoạt động chạy dài.
    • retry: Chính sách có thể cấu hình bao gồm initialInterval, backoffCoefficient, maximumInterval, maximumAttempts và nonRetryableErrorTypes.
  • Tùy chọn Workflow:

    • workflowExecutionTimeout: Thời gian tối đa cho phép một Workflow Execution chạy.
    • workflowRunTimeout: Thời gian tối đa cho phép một Workflow Run duy nhất chạy (một Workflow Execution có thể có nhiều lần chạy do thử lại hoặc tiếp tục).
    • workflowTaskTimeout: Thời gian tối đa cho phép một Workflow Task chạy.

Các tùy chọn này rất quan trọng để xây dựng các hệ thống có khả năng phục hồi có thể xử lý các lỗi tạm thời một cách duyên dáng và ngăn chặn việc chặn vô thời hạn.

Advertisement

Kiểm thử Temporal Workflows

Gói @temporalio/testing cung cấp một môi trường kiểm thử cục bộ để kiểm thử đơn vị và tích hợp các workflow và activity mà không cần một Temporal Cluster đang chạy.

// 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