Claude
Skills
Sign in
Back

event-driven-architect

Included with Lifetime
$97 forever

Designs event-driven architectures with event sourcing, CQRS, pub/sub patterns, and domain events for decoupled systems. Use when users request "event sourcing", "CQRS", "domain events", "pub/sub", or "event-driven".

General

What this skill does


# Event-Driven Architect

Build decoupled, scalable systems with event-driven patterns.

## Core Workflow

1. **Identify domain events**: Define what happened
2. **Design event schema**: Structure event payloads
3. **Implement event bus**: Publish and subscribe
4. **Add event handlers**: React to events
5. **Consider CQRS**: Separate reads and writes
6. **Enable event sourcing**: Store event history

## Event Fundamentals

### Event Structure

```typescript
// events/base.ts
export interface DomainEvent<T = unknown> {
  id: string;
  type: string;
  aggregateId: string;
  aggregateType: string;
  payload: T;
  metadata: {
    timestamp: Date;
    version: number;
    correlationId?: string;
    causationId?: string;
    userId?: string;
  };
}

// Type-safe event creator
export function createEvent<T>(
  type: string,
  aggregateType: string,
  aggregateId: string,
  payload: T,
  metadata?: Partial<DomainEvent['metadata']>
): DomainEvent<T> {
  return {
    id: crypto.randomUUID(),
    type,
    aggregateType,
    aggregateId,
    payload,
    metadata: {
      timestamp: new Date(),
      version: 1,
      ...metadata,
    },
  };
}
```

### Define Domain Events

```typescript
// events/order.events.ts
export interface OrderCreatedPayload {
  customerId: string;
  items: Array<{
    productId: string;
    quantity: number;
    price: number;
  }>;
  totalAmount: number;
  shippingAddress: Address;
}

export interface OrderPaidPayload {
  paymentId: string;
  amount: number;
  method: 'card' | 'bank' | 'wallet';
}

export interface OrderShippedPayload {
  trackingNumber: string;
  carrier: string;
  estimatedDelivery: string;
}

export interface OrderCancelledPayload {
  reason: string;
  cancelledBy: string;
  refundAmount?: number;
}

// Event types
export type OrderEvent =
  | DomainEvent<OrderCreatedPayload> & { type: 'OrderCreated' }
  | DomainEvent<OrderPaidPayload> & { type: 'OrderPaid' }
  | DomainEvent<OrderShippedPayload> & { type: 'OrderShipped' }
  | DomainEvent<OrderCancelledPayload> & { type: 'OrderCancelled' };

// Event creators
export const OrderEvents = {
  created: (orderId: string, payload: OrderCreatedPayload) =>
    createEvent('OrderCreated', 'Order', orderId, payload),

  paid: (orderId: string, payload: OrderPaidPayload) =>
    createEvent('OrderPaid', 'Order', orderId, payload),

  shipped: (orderId: string, payload: OrderShippedPayload) =>
    createEvent('OrderShipped', 'Order', orderId, payload),

  cancelled: (orderId: string, payload: OrderCancelledPayload) =>
    createEvent('OrderCancelled', 'Order', orderId, payload),
};
```

## Event Bus

### In-Memory Event Bus

```typescript
// events/event-bus.ts
import { EventEmitter } from 'events';
import { DomainEvent } from './base';

type EventHandler<T = unknown> = (event: DomainEvent<T>) => Promise<void>;

class EventBus {
  private emitter = new EventEmitter();
  private handlers = new Map<string, EventHandler[]>();

  async publish<T>(event: DomainEvent<T>): Promise<void> {
    console.log(`Publishing event: ${event.type}`, event);

    // Store event (for event sourcing)
    await this.storeEvent(event);

    // Emit to handlers
    this.emitter.emit(event.type, event);
    this.emitter.emit('*', event); // Wildcard for all events
  }

  async publishAll(events: DomainEvent[]): Promise<void> {
    for (const event of events) {
      await this.publish(event);
    }
  }

  subscribe<T>(eventType: string, handler: EventHandler<T>): () => void {
    const wrappedHandler = async (event: DomainEvent<T>) => {
      try {
        await handler(event);
      } catch (error) {
        console.error(`Error handling ${eventType}:`, error);
        // Could emit to dead letter queue here
      }
    };

    this.emitter.on(eventType, wrappedHandler);

    // Return unsubscribe function
    return () => {
      this.emitter.off(eventType, wrappedHandler);
    };
  }

  subscribeAll(handler: EventHandler): () => void {
    return this.subscribe('*', handler);
  }

  private async storeEvent(event: DomainEvent): Promise<void> {
    await db.event.create({
      data: {
        id: event.id,
        type: event.type,
        aggregateId: event.aggregateId,
        aggregateType: event.aggregateType,
        payload: event.payload as any,
        metadata: event.metadata as any,
        createdAt: event.metadata.timestamp,
      },
    });
  }
}

export const eventBus = new EventBus();
```

### Redis-Based Event Bus

```typescript
// events/redis-event-bus.ts
import { Redis } from 'ioredis';
import { DomainEvent } from './base';

const publisher = new Redis(process.env.REDIS_URL!);
const subscriber = new Redis(process.env.REDIS_URL!);

class RedisEventBus {
  private handlers = new Map<string, Set<(event: DomainEvent) => Promise<void>>>();

  constructor() {
    subscriber.on('message', async (channel, message) => {
      const event = JSON.parse(message) as DomainEvent;
      const handlers = this.handlers.get(channel) || new Set();

      for (const handler of handlers) {
        try {
          await handler(event);
        } catch (error) {
          console.error(`Error handling ${event.type}:`, error);
        }
      }
    });
  }

  async publish(event: DomainEvent): Promise<void> {
    const channel = `events:${event.type}`;
    await publisher.publish(channel, JSON.stringify(event));

    // Also store in stream for replay
    await publisher.xadd(
      `stream:${event.aggregateType}`,
      '*',
      'event',
      JSON.stringify(event)
    );
  }

  subscribe(eventType: string, handler: (event: DomainEvent) => Promise<void>): () => void {
    const channel = `events:${eventType}`;

    if (!this.handlers.has(channel)) {
      this.handlers.set(channel, new Set());
      subscriber.subscribe(channel);
    }

    this.handlers.get(channel)!.add(handler);

    return () => {
      this.handlers.get(channel)?.delete(handler);
    };
  }
}

export const eventBus = new RedisEventBus();
```

## Event Handlers

### Handler Registration

```typescript
// handlers/order.handlers.ts
import { eventBus } from '../events/event-bus';
import { OrderEvent } from '../events/order.events';

// Email notification on order created
eventBus.subscribe<OrderCreatedPayload>('OrderCreated', async (event) => {
  await emailService.send({
    to: await getUserEmail(event.payload.customerId),
    template: 'order-confirmation',
    data: {
      orderId: event.aggregateId,
      items: event.payload.items,
      total: event.payload.totalAmount,
    },
  });
});

// Update inventory on order created
eventBus.subscribe<OrderCreatedPayload>('OrderCreated', async (event) => {
  for (const item of event.payload.items) {
    await inventoryService.reserve(item.productId, item.quantity);
  }
});

// Analytics tracking
eventBus.subscribe<OrderPaidPayload>('OrderPaid', async (event) => {
  await analytics.track('order_completed', {
    orderId: event.aggregateId,
    amount: event.payload.amount,
    paymentMethod: event.payload.method,
  });
});

// Notify shipping on order paid
eventBus.subscribe<OrderPaidPayload>('OrderPaid', async (event) => {
  await shippingService.createShipment(event.aggregateId);
});

// Handle cancellation
eventBus.subscribe<OrderCancelledPayload>('OrderCancelled', async (event) => {
  // Release inventory
  const order = await orderRepository.findById(event.aggregateId);
  for (const item of order.items) {
    await inventoryService.release(item.productId, item.quantity);
  }

  // Process refund
  if (event.payload.refundAmount) {
    await paymentService.refund(event.aggregateId, event.payload.refundAmount);
  }

  // Send cancellation email
  await emailService.send({
    to: await getUserEmail(order.customerId),
    template: 'order-cancelled',
    data: {
      orderId: event.aggregateId,
      reason: event.payload.reason,
    },
  });
});
```

## Event Sourcing

### Aggregate with Events

```typescript
// aggregates/order.aggregate.ts
import { DomainEvent } from '../eve

Related in General