> ## Documentation Index
> Fetch the complete documentation index at: https://polos.dev/docs/llms.txt
> Use this file to discover all available pages before exploring further.

# Event-triggered workflows

Event-triggered workflows execute automatically when events are published to specific topics. Perfect for building reactive systems, webhooks, and event-driven architectures.

## Defining event-triggered workflows

Use `trigger_on_event` to specify which topic should trigger the workflow:

<CodeGroup>
  ```python Python theme={null}
  from polos import workflow, WorkflowContext, EventPayload

  @workflow(
      id="user-signup-handler",
      trigger_on_event="user/signup"
  )
  async def user_signup_handler(ctx: WorkflowContext, payload: EventPayload):
      # Triggered when event published to "user.signup" topic
      user_id = payload.data["user_id"]
      email = payload.data["email"]

      # Send welcome email
      await ctx.step.run("send_welcome", send_welcome_email, user_id, email)

      # Create sample data
      await ctx.step.run("setup_account", create_sample_data, user_id)

      return {"status": "onboarded", "user_id": user_id}
  ```

  ```typescript TypeScript theme={null}
  import { defineWorkflow, WorkflowContext, EventPayload } from '@polos/sdk';

  const userSignupHandler = defineWorkflow<EventPayload, void, { status: string; userId: string }>(
    {
      id: 'user-signup-handler',
      triggerOnEvent: 'user/signup',
    },
    async (ctx, payload) => {
      // Triggered when event published to "user.signup" topic
      const userId = payload.data['user_id'];
      const email = payload.data['email'];

      // Send welcome email
      await ctx.step.run('send_welcome', () => sendWelcomeEmail(userId, email));

      // Create sample data
      await ctx.step.run('setup_account', () => createSampleData(userId));

      return { status: 'onboarded', userId };
    }
  );
  ```
</CodeGroup>

## Event payload structure

Event-triggered workflows receive the event in their payload:

<CodeGroup>
  ```python Python theme={null}
  @workflow(
      id="notification-handler",
      trigger_on_event="notifications/new"
  )
  async def notification_handler(ctx: WorkflowContext, payload: EventPayload):
      # Single event
      print(f"Event ID: {payload.id}")
      print(f"Sequence ID: {payload.sequence_id}")
      print(f"Topic: {payload.topic}")
      print(f"Event Type: {payload.event_type}")
      print(f"Data: {payload.data}")  # This is a dict
      print(f"Created At: {payload.created_at}")

      # Process the event
      notification_data = payload.data
      await ctx.step.run("send", send_notification, notification_data)
  ```

  ```typescript TypeScript theme={null}
  import { defineWorkflow, WorkflowContext, EventPayload } from '@polos/sdk';

  const notificationHandler = defineWorkflow<EventPayload, void, void>(
    {
      id: 'notification-handler',
      triggerOnEvent: 'notifications/new',
    },
    async (ctx, payload) => {
      // Single event
      console.log(`Event ID: ${payload.id}`);
      console.log(`Sequence ID: ${payload.sequenceId}`);
      console.log(`Topic: ${payload.topic}`);
      console.log(`Event Type: ${payload.eventType}`);
      console.log(`Data: ${JSON.stringify(payload.data)}`);
      console.log(`Created At: ${payload.createdAt}`);

      // Process the event
      const notificationData = payload.data;
      await ctx.step.run('send', () => sendNotification(notificationData));
    }
  );
  ```
</CodeGroup>

**Event structure:**

```json theme={null}
{
  "id": "evt_123abc",
  "sequence_id": 456
  "topic": "notifications.new",
  "event_type": "notification_created",
  "data": {
    "user_id": "user_456",
    "message": "You have a new message"
  },
  "created_at": "2025-01-28T10:30:00Z"
}
```

## Publishing events

Publish events from workflows or external systems:

### From within a workflow

<CodeGroup>
  ```python Python theme={null}
  @workflow
  async def create_user(ctx: WorkflowContext, input: CreateUserInput):
      # Create user
      user = await ctx.step.run("create", create_user_record, input)

      # Publish event (triggers event-triggered workflows)
      await ctx.step.publish_event(
          "publish_signup",
          topic="user/signup",
          data={
              "user_id": user.id,
              "email": user.email,
              "name": user.name
          },
          event_type="user_created"
      )

      return {"user_id": user.id}
  ```

  ```typescript TypeScript theme={null}
  import { defineWorkflow, WorkflowContext } from '@polos/sdk';

  interface CreateUserInput {
    email: string;
    name: string;
  }

  const createUser = defineWorkflow<CreateUserInput, void, { userId: string }>(
    { id: 'create-user' },
    async (ctx, input) => {
      // Create user
      const user = await ctx.step.run('create', () => createUserRecord(input));

      // Publish event (triggers event-triggered workflows)
      await ctx.step.publishEvent('publish_signup', {
        topic: 'user/signup',
        data: {
          user_id: user.id,
          email: user.email,
          name: user.name,
        },
        type: 'user_created',
      });

      return { userId: user.id };
    }
  );
  ```
</CodeGroup>

### From external systems (API)

<CodeGroup>
  ```python Python theme={null}
  import httpx

  async def publish_event(topic: str, data: dict):
      """Publish an event via API."""
      async with httpx.AsyncClient() as client:
          response = await client.post(
              "https://api.polos.ai/api/v1/events/publish",
              headers={
                  "Authorization": "Bearer YOUR_API_KEY",
                  "Content-Type": "application/json"
              },
              json={
                  "topic": topic,
                  "events": [{
                      "data": data,
                      "event_type": "custom_event"
                  }]
              }
          )
          response.raise_for_status()

  # Trigger workflows listening to "order.completed"
  await publish_event(
      "order.completed",
      {"order_id": "ord_123", "total": 99.99}
  )
  ```

  ```typescript TypeScript theme={null}
  async function publishEvent(topic: string, data: Record<string, unknown>) {
    /** Publish an event via API. */
    const response = await fetch('https://api.polos.ai/api/v1/events/publish', {
      method: 'POST',
      headers: {
        Authorization: 'Bearer YOUR_API_KEY',
        'Content-Type': 'application/json',
      },
      body: JSON.stringify({
        topic,
        events: [
          {
            data,
            event_type: 'custom_event',
          },
        ],
      }),
    });

    if (!response.ok) {
      throw new Error(`Request failed: ${response.status}`);
    }
  }

  // Trigger workflows listening to "order.completed"
  await publishEvent('order.completed', { order_id: 'ord_123', total: 99.99 });
  ```
</CodeGroup>

## Batch events for processing

Process multiple events together for efficiency:

<CodeGroup>
  ```python Python theme={null}
  @workflow(
      id="batch-processor",
      trigger_on_event="analytics/events",
      batch_size=10,              # Process up to 10 events
      batch_timeout_seconds=30    # Or wait max 30 seconds
  )
  async def batch_processor(ctx: WorkflowContext, payload: BatchEventPayload):
      # Multiple events in the payload
      events = payload.events

      print(f"Processing batch of {len(events)} events")

      # Extract all event data
      analytics_data = [event.data for event in events]

      # Process batch
      await ctx.step.run("batch_insert", insert_analytics, analytics_data)

      return {"processed": len(events)}
  ```

  ```typescript TypeScript theme={null}
  import { defineWorkflow, WorkflowContext, BatchEventPayload } from '@polos/sdk';

  const batchProcessor = defineWorkflow<BatchEventPayload, void, { processed: number }>(
    {
      id: 'batch-processor',
      triggerOnEvent: 'analytics/events',
      batchSize: 10,              // Process up to 10 events
      batchTimeoutSeconds: 30,    // Or wait max 30 seconds
    },
    async (ctx, payload) => {
      // Multiple events in the payload
      const events = payload.events;

      console.log(`Processing batch of ${events.length} events`);

      // Extract all event data
      const analyticsData = events.map((event) => event.data);

      // Process batch
      await ctx.step.run('batch_insert', () => insertAnalytics(analyticsData));

      return { processed: events.length };
    }
  );
  ```
</CodeGroup>

**Batching behavior:**

* Workflow triggers when either `batch_size` is reached or `batch_timeout_seconds` elapses
* If 10 events arrive in 5 seconds → triggers immediately with 10 events
* If only 3 events arrive in 30 seconds → triggers with 3 events after timeout

**Batch payload structure:**

```json theme={null}
{
  "events": [
    {
      "id": "evt_1",
      "sequence_id": 1001,
      "topic": "analytics.events",
      "data": {"action": "click"},
      "created_at": "2025-01-28T10:30:00Z"
    },
    {
      "id": "evt_2",
      "sequence_id": 1002,
      "topic": "analytics.events",
      "data": {"action": "view"},
      "created_at": "2025-01-28T10:30:05Z"
    }
  ]
}
```

## Multiple handlers for one topic

Multiple workflows can listen to the same topic:

<CodeGroup>
  ```python Python theme={null}
  @workflow(
      id="immediate-handler",
      trigger_on_event="order/created"
  )
  async def immediate_handler(ctx: WorkflowContext, payload: EventPayload):
      """Process each order immediately."""
      order_id = payload.data["order_id"]

      await ctx.step.run("send_confirmation", send_order_confirmation, order_id)
      return {"handler": "immediate"}

  @workflow(
      id="batched-handler",
      trigger_on_event="order/created",
      batch_size=5,
      batch_timeout_seconds=60
  )
  async def batched_handler(ctx: WorkflowContext, payload: BatchEventPayload):
      """Process orders in batches for analytics."""
      events = payload.events

      order_ids = [e.data["order_id"] for e in events]
      await ctx.step.run("batch_analytics", update_analytics, order_ids)

      return {"handler": "batched", "count": len(events)}
  ```

  ```typescript TypeScript theme={null}
  import {
    defineWorkflow,
    WorkflowContext,
    EventPayload,
    BatchEventPayload,
  } from '@polos/sdk';

  const immediateHandler = defineWorkflow<EventPayload, void, { handler: string }>(
    {
      id: 'immediate-handler',
      triggerOnEvent: 'order/created',
    },
    async (ctx, payload) => {
      /** Process each order immediately. */
      const orderId = payload.data['order_id'];

      await ctx.step.run('send_confirmation', () => sendOrderConfirmation(orderId));
      return { handler: 'immediate' };
    }
  );

  const batchedHandler = defineWorkflow<BatchEventPayload, void, { handler: string; count: number }>(
    {
      id: 'batched-handler',
      triggerOnEvent: 'order/created',
      batchSize: 5,
      batchTimeoutSeconds: 60,
    },
    async (ctx, payload) => {
      /** Process orders in batches for analytics. */
      const events = payload.events;

      const orderIds = events.map((e) => e.data['order_id']);
      await ctx.step.run('batch_analytics', () => updateAnalytics(orderIds));

      return { handler: 'batched', count: events.length };
    }
  );
  ```
</CodeGroup>

**When an event is published:**

* `immediate-handler` triggers once per event
* `batched-handler` triggers once per batch (up to 5 events or 60 seconds)

## Event filtering

Filter events by `event_type`:

<CodeGroup>
  ```python Python theme={null}
  @workflow(
      id="high-priority-handler",
      trigger_on_event="notifications/system"
  )
  async def high_priority_handler(ctx: WorkflowContext, payload: PayloadEvent):
      # Only process high priority notifications
      if payload.event_type == "high_priority":
          await ctx.step.run("alert", send_alert, payload.data)
      else:
          print(f"Skipping event type: {payload.event_type}")
  ```

  ```typescript TypeScript theme={null}
  import { defineWorkflow, WorkflowContext, EventPayload } from '@polos/sdk';

  const highPriorityHandler = defineWorkflow<EventPayload, void, void>(
    {
      id: 'high-priority-handler',
      triggerOnEvent: 'notifications/system',
    },
    async (ctx, payload) => {
      // Only process high priority notifications
      if (payload.eventType === 'high_priority') {
        await ctx.step.run('alert', () => sendAlert(payload.data));
      } else {
        console.log(`Skipping event type: ${payload.eventType}`);
      }
    }
  );
  ```
</CodeGroup>

**Alternative approach: Use multiple topics**

<CodeGroup>
  ```python Python theme={null}
  # Publish to specific topics
  await ctx.step.publish_event(
      "publish",
      topic="notifications/high_priority",  # Specific topic
      data=notification_data
  )

  @workflow(
      id="high-priority-handler",
      trigger_on_event="notifications/high_priority"  # Only high priority
  )
  async def high_priority_handler(ctx: WorkflowContext, payload: EventPayload):
      # All events on this topic are high priority
      await ctx.step.run("alert", send_alert, payload.data)
  ```

  ```typescript TypeScript theme={null}
  import { defineWorkflow, WorkflowContext, EventPayload } from '@polos/sdk';

  // Publish to specific topics
  await ctx.step.publishEvent('publish', {
    topic: 'notifications/high_priority', // Specific topic
    data: notificationData,
  });

  const highPriorityHandler = defineWorkflow<EventPayload, void, void>(
    {
      id: 'high-priority-handler',
      triggerOnEvent: 'notifications/high_priority', // Only high priority
    },
    async (ctx, payload) => {
      // All events on this topic are high priority
      await ctx.step.run('alert', () => sendAlert(payload.data));
    }
  );
  ```
</CodeGroup>

## Event topic patterns

Use topic hierarchies for organization:

```python theme={null}
# User events
"user/signup"
"user/login"
"user/deleted"

# Order events
"order/created"
"order/completed"
"order/cancelled"

# System events
"system/error"
"system/warning"
"system/maintenance"
```

## Key takeaways

* **Event-triggered workflows** execute automatically when events are published
* **Use `trigger_on_event`** to specify the topic
* **Publish events** from workflows or external systems via API
* **Batch events** for efficiency with `batch_size` and `batch_timeout_seconds`
* **Multiple handlers** can listen to the same topic
