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

# Webhook to Iceberg

> Accept webhook events, persist them to a Stream, and publish normalized rows to Iceberg.

Use this pattern for payment, product, telemetry, or partner webhooks that must
be acknowledged quickly but retained for downstream analysis.

```mermaid theme={null}
flowchart LR
    X["Webhook sender"] --> W["Ingress Worker\nauth + envelope"]
    W --> S["Stream\ndurable events"]
    S --> P["Pipeline\nfilter + normalize"]
    P --> K["Sink + Catalog"]
    K --> I["Iceberg table"]
```

## 1. Bind the Stream

```jsonc theme={null}
{
  "name": "webhook-ingress",
  "main": "worker.js",
  "compatibility_date": "2026-08-27",
  "pipelines": [
    { "binding": "WEBHOOKS", "stream": "partner_webhooks" }
  ],
  "vars": {
    "WEBHOOK_ISSUER": "partner"
  }
}
```

Store the verification secret as an encrypted deployment secret, not in
`vars`.

Create `partner_webhooks` with an immutable structured schema. The binding
references this resource; it does not define its schema:

```json theme={null}
{
  "fields": [
    { "name": "event_id", "type": "string", "required": true },
    { "name": "event_type", "type": "string", "required": true },
    { "name": "issuer", "type": "string", "required": true },
    { "name": "received_at", "type": "timestamp", "required": true },
    { "name": "payload", "type": "json", "required": true }
  ]
}
```

## 2. Verify and ingest

```js theme={null}
export default {
  async fetch(request, env) {
    if (request.method !== "POST") {
      return new Response("Method not allowed", { status: 405 });
    }

    const raw = await request.text();
    const valid = await verifySignature(
      raw,
      request.headers.get("x-webhook-signature"),
      env.WEBHOOK_SECRET,
    );
    if (!valid) return new Response("Invalid signature", { status: 401 });

    let payload;
    try {
      payload = JSON.parse(raw);
    } catch {
      return new Response("Invalid JSON", { status: 400 });
    }
    await env.WEBHOOKS.send([{
      event_id: payload.id,
      event_type: payload.type,
      issuer: env.WEBHOOK_ISSUER,
      received_at: new Date().toISOString(),
      payload,
    }]);

    return new Response(null, { status: 202 });
  },
};
```

Return `202` only after `send` resolves. At that point the Stream has durably
acknowledged the record; lakehouse publication can happen independently.

## 3. Normalize in a Pipeline

```sql theme={null}
INSERT INTO webhook_events
SELECT event_id,
       event_type,
       issuer,
       received_at,
       payload.customer.id AS customer_id,
       payload.data AS data
FROM partner_webhooks
WHERE event_id IS NOT NULL
  AND event_type IS NOT NULL;
```

Map `webhook_events` to an Iceberg Sink. The Sink and Catalog deduplicate
Pipeline batch identities, so a crash and retry cannot publish the same batch
as a second logical commit.

## 4. Keep the raw envelope

If audit or replay matters, fan out the same Pipeline input to a second Sink:

```sql theme={null}
INSERT INTO raw_webhooks
SELECT * FROM partner_webhooks;

INSERT INTO webhook_events
SELECT event_id, event_type, issuer, received_at, payload.data AS data
FROM partner_webhooks
WHERE event_id IS NOT NULL;
```

The raw Iceberg table preserves the source envelope. The normalized table is
the stable analytical contract.

## Production checks

* Verify the signature before writing to the Stream.
* Preserve the source event ID even if deduplication happens later.
* Keep payloads below the 5 MiB Stream request limit.
* Treat schema changes as versioned data contracts.
* Alert on Stream validation errors and Pipeline cursor lag.


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.