•15 min read

Temporal Workflow Orchestration: Building Resilient Distributed State Machines in TypeScript

Temporal Workflow Orchestration: Building Resilient Distributed State Machines in TypeScript

The construction of robust, long-running distributed systems presents significant challenges in state management, fault tolerance, and operational complexity. Traditional approaches often lead to intricate, error-prone codebases burdened by manual retry logic, compensation mechanisms, and persistent state management across service boundaries. Temporal.io emerges as a critical infrastructure component, abstracting these complexities by providing a durable execution engine for orchestrating distributed state machines. This guide details the architectural principles, implementation patterns, and operational considerations for leveraging Temporal with TypeScript to build highly resilient applications.

Audio Briefing
0:00 / 0:00

The Distributed State Machine Problem

Distributed systems inherently face challenges such as network partitions, service failures, partial failures, and unpredictable latency. Orchestrating multi-step business processes across disparate microservices, each with its own failure modes, requires sophisticated mechanisms to ensure atomicity, consistency, isolation, and durability (ACID properties, or their distributed equivalent, BASE).

Common issues include:

  • Lost State: A service crashes mid-process, losing in-memory state and requiring manual recovery or complex persistence layers.
  • Incomplete Transactions: A multi-step operation fails after some steps complete but before others, leading to an inconsistent system state.
  • Manual Retries and Timeouts: Implementing robust retry policies with exponential backoff, jitter, and circuit breakers across service calls is non-trivial and often duplicated.
  • Compensation Logic: Reversing partially completed operations (e.g., refunding a payment after inventory deduction fails) requires explicit, often complex, "saga" patterns.
  • Observability: Tracing the exact state and progress of a long-running process across multiple services is difficult.

Message queues (e.g., Kafka, RabbitMQ) provide asynchronous communication and some level of durability for messages, but they do not inherently manage the state of a multi-step process or provide durable execution of the orchestration logic itself. Custom state machines built atop these often re-implement much of what Temporal offers out-of-the-box.

Advertisement

Temporal's Architectural Paradigm

Temporal addresses these challenges by externalizing workflow state and execution logic to a dedicated, fault-tolerant service. It allows developers to write complex, long-running business processes as ordinary code, abstracting away the underlying distributed systems primitives.

Core Concepts

  1. Workflow as Code: A Temporal Workflow is a durable, fault-tolerant function that orchestrates Activities. It is written as standard application code (e.g., TypeScript) and appears to execute sequentially, even if it pauses for days, weeks, or years.
  2. Deterministic Execution: This is the cornerstone of Temporal. Workflow code must be deterministic, meaning that given the same input and the same sequence of events (from its history), it must always produce the same output and the same sequence of commands (e.g., scheduling activities, timers). This enables Temporal to replay workflow execution history to recover state after a worker failure or during version upgrades. Non-deterministic operations (e.g., Date.now(), Math.random(), direct I/O) are strictly prohibited within workflow code; they must be encapsulated within Activities.
  3. Event Sourcing & History Replay: Every state change and command issued by a Workflow is recorded as an event in its execution history. If a worker processing a Workflow fails, another worker can pick up the Workflow, replay its history from the beginning, and reconstruct its exact state, then continue execution from the point of failure. This mechanism provides transparent fault tolerance and durability.
  4. Activities: Activities are the units of non-deterministic, side-effecting work in Temporal. They encapsulate interactions with external systems (databases, APIs, message queues, file systems). Activities are designed to be idempotent and can be retried automatically by Temporal with configurable policies.
  5. Workers: Workers are application processes that host Workflow and Activity implementations. They poll Task Queues on the Temporal Cluster, execute tasks, and report results. Workers are stateless and can be scaled horizontally.
  6. Task Queues: Workflows and Activities communicate with the Temporal Cluster via Task Queues. When a Workflow schedules an Activity, a task is placed on a specific Activity Task Queue. When a Workflow needs to execute, a task is placed on a Workflow Task Queue. Workers subscribe to these queues. This decouples task producers from consumers.
  7. Visibility: The Temporal Cluster provides APIs and a UI (Temporal Web UI) to inspect the state, history, and progress of all running and completed workflows, aiding in debugging and operational monitoring.

TypeScript SDK

The Temporal TypeScript SDK provides decorators and utility functions to define workflows and activities, manage client connections, and run workers. It leverages TypeScript's strong typing for improved developer experience and compile-time safety.

Building a Resilient Workflow: A Practical Example (Order Processing)

Consider an e-commerce order fulfillment process:

  1. Process Payment.
  2. Deduct Inventory.
  3. Ship Order.
  4. Send Confirmation.

This process must be resilient to failures at any step and ensure consistency.

Project Setup

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

Modify tsconfig.json to include:

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

Activity Definitions (activities.ts)

Activities are where external, non-deterministic operations occur.

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

Workflow Definition (workflows.ts)

The workflow orchestrates activities and manages the overall process state.

// 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 Implementation (worker.ts)

The worker process hosts and executes the workflow and activity code.

// 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 Interaction (client.ts)

The client initiates the workflow execution.

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

Running the Example

  1. Start Temporal Cluster: Use Docker Compose for local development:
    docker run --rm --name temporal-dev -p 7233:7233 -p 8080:8080 temporalio/auto-setup:1.20.0
    
  2. Start Worker:
    npx ts-node worker.ts
    
  3. Start Client (in a separate terminal):
    npx ts-node client.ts
    

Observe the logs from the worker and client. You will see activities retrying on failure and compensation logic executing if a later step fails.

Advanced Temporal Patterns

Sagas for Distributed Transactions

The order fulfillment example demonstrates a basic Saga pattern. A Saga is a sequence of local transactions where each transaction updates data within a single service and publishes an event to trigger the next step. If a step fails, compensating transactions are executed to undo the preceding completed steps. Temporal workflows inherently support Sagas by allowing you to define compensation logic directly within the catch block of your workflow.

The orderFulfillmentWorkflow explicitly checks inventoryDeducted and paymentTransactionId flags to determine which compensation activities (restoreInventory, refundPayment) need to be called. This ensures atomicity across distributed services.

Signals

Signals are asynchronous messages sent to a running workflow. They allow external systems to interact with and update the state of a workflow without waiting for its completion.

Example: Approving an Order

// 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 are synchronous calls to a running workflow to retrieve its current state without altering it. They are read-only operations.

Example: Getting Order Status

// 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 allow breaking down complex workflows into smaller, manageable, and independently executable units. They provide better modularity, isolation, and can have their own retry policies and timeouts.

Example: Processing individual order items as 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 and Retries

Temporal provides robust, configurable timeout and retry policies for both Workflows and Activities.

  • Activity Options:

    • startToCloseTimeout: Maximum time an Activity Execution is allowed to run after it has been picked up by a Worker.
    • scheduleToCloseTimeout: Maximum time from when an Activity Task is scheduled to when it completes. Includes queue time and retries.
    • scheduleToStartTimeout: Maximum time from when an Activity Task is scheduled to when a Worker starts processing it.
    • heartbeatTimeout: Interval at which an Activity must report its progress. If a heartbeat is missed, the Activity is retried. Useful for long-running activities.
    • retry: Configurable policy including initialInterval, backoffCoefficient, maximumInterval, maximumAttempts, and nonRetryableErrorTypes.
  • Workflow Options:

    • workflowExecutionTimeout: Maximum time a Workflow Execution is allowed to run.
    • workflowRunTimeout: Maximum time a single Workflow Run is allowed to run (a Workflow Execution can have multiple runs due to retries or continuations).
    • workflowTaskTimeout: Maximum time a Workflow Task is allowed to run.

These options are crucial for building resilient systems that can gracefully handle transient failures and prevent indefinite blocking.

Advertisement

Testing Temporal Workflows

The @temporalio/testing package provides a local test environment for unit and integration testing workflows and activities without needing a running Temporal Cluster.

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