Examples/Unexpected Worker Shutdown Handling

Unexpected Worker Shutdown Handling

Demonstrates how to handle unexpected worker crashes with configurable retry strategies, per-message-type overrides, and detailed error context.

Unexpected Worker Shutdown Handling

This example demonstrates how to handle unexpected worker crashes with configurable strategies and detailed error context.

Overview

When a worker process dies unexpectedly (OOM crash, SIGKILL, or hard crash), pending requests may be left in limbo. The shutdown handling feature provides:

  • Reliable crash detection via exit events (works even for OOM kills)
  • Configurable strategies: Reject immediately or retry automatically
  • Per-message-type overrides: Different strategies for different operations
  • Rich error context: WorkerCrashedError with crash details and attempt counts

When to Use Each Strategy

Reject Strategy (strategy: 'reject')

Use for non-idempotent operations where retrying could cause side effects:

  • Payment processing
  • Database writes
  • External API calls that aren't safe to replay
unexpectedShutdown: {
  strategy: 'reject',  // default behavior
  processPayment: { strategy: 'reject' },
}

Retry Strategy (strategy: 'retry')

Use for idempotent operations that can be safely retried:

  • Read-only queries
  • Computations
  • Data processing tasks
unexpectedShutdown: {
  strategy: 'reject',
  compute: { strategy: 'retry', attempts: 2 },  // retry up to 2 times
  fetchData: { strategy: 'retry' },  // uses default 1 attempt
}

Files

Shared Message Definitions

Define message types that both host and worker import:

messages.ts
/**
 * Shared message definitions for the shutdown-handling example
 *
 * This file is imported by both the host and worker to ensure
 * type safety and avoid duplication.
 */

import { DefineMessages } from 'isolated-workers';

/**
 * Message types demonstrating different shutdown scenarios
 */
export type Messages = DefineMessages<{
  // Idempotent operation - safe to retry
  compute: {
    payload: { value: number; shouldCrash?: boolean };
    result: { result: number };
  };

  // Non-idempotent operation - should reject on crash
  processPayment: {
    payload: { paymentId: string; amount: number; shouldCrash?: boolean };
    result: { success: boolean };
  };

  // Operation that may fail - demonstrate retry limits
  processBatch: {
    payload: { items: number[]; crashAfterItems?: number };
    result: { processed: number };
  };
}>;

Host (Client)

Demonstrates different shutdown strategies and error handling:

host.ts
/**
 * Shutdown Handling Example - Host (Client) Side
 *
 * Demonstrates different strategies for handling unexpected worker crashes:
 * - 'reject' strategy: Immediately reject pending requests
 * - 'retry' strategy: Automatically retry failed requests
 * - Per-message-type overrides: Different strategies for different operations
 */

import { createWorker, WorkerCrashedError } from 'isolated-workers';
import { fileURLToPath } from 'url';
import { dirname, join } from 'path';
import type { Messages } from './messages.js';

const __dirname = dirname(fileURLToPath(import.meta.url));

async function createTestWorker() {
  return createWorker<Messages>({
    script: join(__dirname, 'worker.ts'),
    timeout: 10000,
    unexpectedShutdown: {
      strategy: 'reject', // default for all messages
      compute: { strategy: 'retry', attempts: 3 }, // retry idempotent ops up to 3 times
      processBatch: { strategy: 'retry', attempts: 3 }, // retry with limit
      // processPayment uses default 'reject' - non-idempotent
    },
  });
}

async function main() {
  console.log('=== Shutdown Handling Example ===\n');

  let worker = await createTestWorker();
  console.log(`Worker spawned with PID: ${worker.pid}\n`);

  // Test 1: Normal operation (no crash)
  console.log('Test 1: Normal operation');
  try {
    const result = await worker.send('compute', { value: 5 });
    console.log('✓ Result:', result.result, '\n');
  } catch (err) {
    console.error('✗ Unexpected error:', (err as Error).message, '\n');
  }

  // Test 2: Idempotent operation with crash (retry will also crash - demonstrates exhaustion)
  // Note: When a worker crashes, the retry spawns a NEW worker process.
  // If the crash condition is in the payload (shouldCrash: true), retries will also crash.
  // This demonstrates retry exhaustion for idempotent operations.
  console.log('Test 2: Idempotent operation with crash - retry exhaustion');
  try {
    const result = await worker.send('compute', {
      value: 7,
      shouldCrash: true, // Will crash on every attempt
    });
    console.log('✗ Should have crashed, got:', result, '\n');
  } catch (err) {
    if (err instanceof WorkerCrashedError) {
      console.log('✓ WorkerCrashedError after retries:', err.message);
      console.log(`  - Reason: ${err.reason.type}`);
      console.log(`  - Attempt: ${err.attempt}/${err.maxAttempts}\n`);
    } else {
      console.error('✗ Unexpected error:', (err as Error).message, '\n');
    }
  }

  // Need a fresh worker since Test 2 exhausted retries
  await worker.close();
  worker = await createTestWorker();
  console.log(`Fresh worker spawned with PID: ${worker.pid}\n`);

  // Test 3: Non-idempotent operation (should reject immediately, no retry)
  console.log('Test 3: Non-idempotent operation (should reject, not retry)');
  try {
    // This payment will crash, and since processPayment uses 'reject' strategy,
    // it should fail with WorkerCrashedError immediately (no retries)
    const result = await worker.send('processPayment', {
      paymentId: 'pay-123',
      amount: 100,
      shouldCrash: true,
    });
    console.log('✗ Payment should have been rejected, got:', result, '\n');
  } catch (err) {
    if (err instanceof WorkerCrashedError) {
      console.log('✓ WorkerCrashedError (as expected):', err.message);
      console.log(`  - Attempt: ${err.attempt}/${err.maxAttempts}`);
      console.log(`  - No retries (non-idempotent operation)\n`);
    } else {
      console.error('✗ Unexpected error:', (err as Error).message, '\n');
    }
  }

  // Need a fresh worker since Test 3 crashed
  await worker.close();
  worker = await createTestWorker();
  console.log(`Fresh worker spawned with PID: ${worker.pid}\n`);

  // Test 4: Batch processing with retry exhaustion
  console.log('Test 4: Batch processing with retry exhaustion');
  try {
    // This batch will crash immediately (crashAfterItems: 0), and will keep crashing
    // until we exhaust all 3 retry attempts
    const result = await worker.send('processBatch', {
      items: [1, 2, 3, 4, 5],
      crashAfterItems: 0, // Crash immediately every time
    });
    console.log('✗ Should have exhausted retries, got:', result, '\n');
  } catch (err) {
    if (err instanceof WorkerCrashedError) {
      console.log('✓ WorkerCrashedError after retries:', err.message);
      console.log(`  - Attempt: ${err.attempt}/${err.maxAttempts}\n`);
    } else {
      console.error('✗ Unexpected error:', (err as Error).message, '\n');
    }
  }

  // Need another fresh worker since Test 4 exhausted retries
  await worker.close();
  worker = await createTestWorker();
  console.log(`Fresh worker spawned with PID: ${worker.pid}\n`);

  // Test 5: Successful batch processing (no crash)
  console.log('Test 5: Successful batch processing');
  try {
    const result = await worker.send('processBatch', {
      items: [1, 2, 3],
      // No crashAfterItems = no crash
    });
    console.log('✓ Processed batch of', result.processed, 'items\n');
  } catch (err) {
    console.error('✗ Unexpected error:', (err as Error).message, '\n');
  }

  // Test 6: Cleanup
  console.log('Test 6: Cleanup');
  await worker.close();
  console.log('✓ Worker closed successfully\n');

  console.log('=== All tests completed ===');
}

main().catch((err) => {
  console.error('Host error:', err);
  process.exit(1);
});

Worker

Worker that processes requests and can crash on specific commands:

worker.ts
/**
 * Shutdown Handling Example - Worker (Server) Side
 *
 * This worker demonstrates various scenarios including:
 * - Normal message processing
 * - Crashing during processing (for retry testing)
 * - Processing that takes time
 */

import { startWorkerServer, Handlers } from 'isolated-workers';
import type { Messages } from './messages.js';

const handlers: Handlers<Messages> = {
  compute: ({ value, shouldCrash = false }) => {
    console.log(`Worker: computing ${value}^2 (shouldCrash=${shouldCrash})`);

    if (shouldCrash) {
      console.log('Worker: crashing as requested');
      process.exit(1);
    }

    return { result: value * value };
  },

  processPayment: ({ paymentId, amount, shouldCrash = false }) => {
    console.log(`Worker: processing payment ${paymentId} for $${amount}`);

    if (shouldCrash) {
      console.log('Worker: crashing during payment processing');
      process.exit(1);
    }

    return { success: true };
  },

  processBatch: ({ items, crashAfterItems }) => {
    console.log(`Worker: processing batch of ${items.length} items`);

    if (crashAfterItems !== undefined && crashAfterItems >= 0) {
      console.log(
        `Worker: will crash after processing ${crashAfterItems} items`
      );
      // Process up to the crash point, then crash
      for (let i = 0; i < Math.min(crashAfterItems, items.length); i++) {
        // Simulate processing
        items[i] * 2;
      }
      console.log('Worker: crashing as requested');
      process.exit(1);
    }

    return { processed: items.length };
  },
};

async function main() {
  console.log('Worker starting...');
  console.log(`Worker PID: ${process.pid}`);

  const server = await startWorkerServer(handlers);

  console.log('Worker ready and accepting requests');

  process.on('SIGTERM', async () => {
    console.log('Worker received SIGTERM');
    await server.stop();
    process.exit(0);
  });
}

main().catch((err) => {
  console.error('Worker error:', err);
  process.exit(1);
});

Running the Example

cd examples && pnpm run:shutdown-handling

Key Concepts

Worker-Level Default Strategy

The top-level unexpectedShutdown.strategy applies to all message types by default:

const worker = await createWorker<Messages>({
  script: './worker.js',
  unexpectedShutdown: {
    strategy: 'reject', // default for all messages
  },
});

Per-Message-Type Overrides

Override the default for specific message types:

const worker = await createWorker<Messages>({
  script: './worker.js',
  unexpectedShutdown: {
    strategy: 'reject', // default for most
    compute: { strategy: 'retry', attempts: 2 }, // override for compute
    processBatch: { strategy: 'retry' }, // override with default 1
  },
});

WorkerCrashedError

When a crash occurs (or retries are exhausted), requests reject with WorkerCrashedError:

try {
  await worker.send('processPayment', { paymentId: '123', amount: 100 });
} catch (err) {
  if (err instanceof WorkerCrashedError) {
    console.error('Worker crashed:', err.message);
    console.error('Reason:', err.reason); // exit code, signal, or error
    console.error('Attempt:', err.attempt, '/', err.maxAttempts);
  }
}

Retry Flow

  1. Worker crashes (exit, error, or close event)
  2. handleUnexpectedShutdown() called for each pending request
  3. For each request:
    • Look up strategy for its message type
    • If reject OR attempt >= maxAttempts → reject with WorkerCrashedError
    • If retry AND attempt < maxAttempts → queue for retry
  4. If any requests queued → spawn new worker, re-send each with attempt++
  5. Original Promise stays the same—user awaits the same resolve/reject

Crash Detection

The library detects crashes via exit events, which fire reliably even for:

  • OOM kills (SIGKILL)
  • Hard crashes (uncaught exceptions without handlers)
  • Manual termination (kill -9 )
// Detection happens automatically:
child.on('exit', (code, signal) => {
  // Triggers shutdown handling for all pending requests
});

Running the Example

Run the example

bash
pnpm run:shutdown-handling