# Receive replies once and fan them out

A webhook receiver that verifies each prospect_response delivery, stores it, answers at once, then runs any number of handlers, with replay for failures.

Works with: Slack.

Endpoints used:

- [`POST /v1/campaigns/{campaign_id}/webhooks`](https://docs.versionseven.ai/api-reference/campaigns/create-webhook) Create a webhook
- [`GET /v1/campaigns/{campaign_id}/webhooks`](https://docs.versionseven.ai/api-reference/campaigns/list-webhooks) List a campaign's webhooks
- [`POST /v1/campaigns/{campaign_id}/webhooks/{webhook_id}/secret`](https://docs.versionseven.ai/api-reference/campaigns/rotate-webhook-secret) Set or rotate a webhook secret
- [`GET /v1/webhooks/examples`](https://docs.versionseven.ai/api-reference/reference/list-webhook-examples) List webhook examples

A campaign has one webhook, so one receiver gets every reply. This recipe is that receiver: it verifies the delivery, stores it, answers Victoria AI straight away, then hands the reply to each handler you've switched on. A handler is a small file that does one thing, such as posting to Slack or appending a row to a sheet, and the other webhook recipes in this Cookbook are written as handlers for it. When a handler fails, replay the stored delivery instead of waiting for a prospect to write again.

## Before you start

- A signing secret of your own, at least 16 characters, in `VICTORIA_WEBHOOK_SECRET`. One way to make one is `openssl rand -hex 32`.
- A public HTTPS address for the receiver. Victoria AI doesn't deliver to `localhost` or private addresses.
- A token of your own in `RELAY_TOKEN`, which protects the replay endpoint. Leave it unset to disable replay.
- For the Slack handler, a Slack [incoming webhook](https://api.slack.com/messaging/webhooks) URL in `SLACK_WEBHOOK_URL`.
- An API key with `campaigns:write` in `VICTORIA_API_KEY`, to register the webhook. See [Authentication](https://docs.versionseven.ai/guides/authentication).
- Node.js 18 or later with `express`, or Python 3.10 or later with `flask` and `requests`.

## How it works

1. Victoria AI posts a [`prospect_response`](https://docs.versionseven.ai/api-reference/webhooks/prospect-response) event to `/hooks/<campaign_id>`. The payload names the campaign but carries no id, so the id lives in the URL you register.
2. The receiver checks `X-Signature-256` against the raw body, skips an `idempotency_key` it has seen, and appends the delivery to a file.
3. It answers `200` at once. Victoria AI then considers the delivery done and won't retry it, so the receiver owns what happens next.
4. Each handler in `HANDLERS` runs in turn. A failure is logged and doesn't stop the others.
5. `POST /replay/<idempotency_key>` runs the handlers again for a stored delivery.

## 1. Run the receiver

**Node.js (Express)**

```javascript
// receiver.mjs
import crypto from "node:crypto";
import { appendFile, readFile } from "node:fs/promises";
import express from "express";

const secret = process.env.VICTORIA_WEBHOOK_SECRET;
const relayToken = process.env.RELAY_TOKEN; // unset: replay is disabled
const deliveriesFile = process.env.DELIVERIES_FILE ?? "deliveries.ndjson";
const handlerNames = (process.env.HANDLERS ?? "console").split(",").map((name) => name.trim()).filter(Boolean);
const port = Number(process.env.PORT ?? 3000);

if (!secret || secret.length < 16) throw new Error("VICTORIA_WEBHOOK_SECRET must be at least 16 characters");

// A handler is ./handlers/<name>.mjs whose default export takes a delivery.
const handlers = await Promise.all(
  handlerNames.map(async (name) => ({ name, handle: (await import(`./handlers/${name}.mjs`)).default }))
);

function isValidSignature(rawBody, header) {
  if (typeof header !== "string") return false;
  const digest = crypto.createHmac("sha256", secret).update(rawBody).digest("hex");
  const expected = Buffer.from(`sha256=${digest}`);
  const received = Buffer.from(header);
  return expected.length === received.length && crypto.timingSafeEqual(expected, received);
}

// Every stored delivery, one JSON object per line. A database works the same way.
async function readDeliveries() {
  try {
    const text = await readFile(deliveriesFile, "utf8");
    return text.split("\n").filter(Boolean).map((line) => JSON.parse(line));
  } catch (error) {
    if (error.code === "ENOENT") return [];
    throw error;
  }
}

// Keys already received, so a restart doesn't repeat deliveries.
const processed = new Set((await readDeliveries()).map((delivery) => delivery.idempotency_key));

async function dispatch(delivery) {
  const results = [];
  for (const { name, handle } of handlers) {
    try {
      await handle(delivery);
      results.push({ handler: name, ok: true });
    } catch (error) {
      console.error(`${name} failed for ${delivery.idempotency_key}: ${error.message}`);
      results.push({ handler: name, ok: false, error: error.message });
    }
  }
  return results;
}

const app = express();

// express.raw keeps the body as the exact bytes that were signed.
app.post("/hooks/:campaignId", express.raw({ type: "application/json" }), async (req, res) => {
  if (!isValidSignature(req.body, req.get("X-Signature-256"))) return res.sendStatus(401);

  const event = JSON.parse(req.body.toString("utf8"));
  if (event.event !== "prospect_response" || processed.has(event.idempotency_key)) {
    return res.sendStatus(200);
  }

  const delivery = {
    received_at: new Date().toISOString(),
    campaign_id: req.params.campaignId,
    idempotency_key: event.idempotency_key,
    event,
  };
  try {
    await appendFile(deliveriesFile, `${JSON.stringify(delivery)}\n`);
  } catch (error) {
    console.error(`Couldn't store ${event.idempotency_key}: ${error.message}`);
    return res.sendStatus(500); // not stored, so let Victoria AI retry
  }
  processed.add(event.idempotency_key);
  res.sendStatus(200); // stored: acknowledged before the handlers run

  dispatch(delivery);
});

app.post("/replay/:idempotencyKey", async (req, res) => {
  if (!relayToken || req.get("Authorization") !== `Bearer ${relayToken}`) return res.sendStatus(403);
  const delivery = (await readDeliveries()).find((candidate) => candidate.idempotency_key === req.params.idempotencyKey);
  if (!delivery) return res.sendStatus(404);
  res.json({ idempotency_key: delivery.idempotency_key, results: await dispatch(delivery) });
});

app.listen(port, () => console.log(`Receiver listening on port ${port} with handlers: ${handlerNames.join(", ")}`));
```

**Python (Flask)**

```python
# receiver.py
import hashlib
import hmac
import importlib
import json
import os
import threading
from datetime import datetime, timezone

from flask import Flask, abort, jsonify, request

SECRET = os.environ["VICTORIA_WEBHOOK_SECRET"].encode()
RELAY_TOKEN = os.environ.get("RELAY_TOKEN")  # unset: replay is disabled
DELIVERIES_FILE = os.environ.get("DELIVERIES_FILE", "deliveries.ndjson")
HANDLER_NAMES = [name.strip() for name in os.environ.get("HANDLERS", "console").split(",") if name.strip()]

if len(SECRET) < 16:
    raise SystemExit("VICTORIA_WEBHOOK_SECRET must be at least 16 characters")

# A handler is handlers/<name>.py with a handle(delivery) function.
HANDLERS = [(name, importlib.import_module(f"handlers.{name}").handle) for name in HANDLER_NAMES]

app = Flask(__name__)
write_lock = threading.Lock()


def is_valid_signature(raw_body: bytes, header: str | None) -> bool:
    if not header:
        return False
    expected = "sha256=" + hmac.new(SECRET, raw_body, hashlib.sha256).hexdigest()
    return hmac.compare_digest(expected, header)


def read_deliveries() -> list[dict]:
    """Every stored delivery, one JSON object per line. A database works the same way."""
    try:
        with open(DELIVERIES_FILE, encoding="utf-8") as file:
            return [json.loads(line) for line in file if line.strip()]
    except FileNotFoundError:
        return []


# Keys already received, so a restart doesn't repeat deliveries.
processed = {delivery["idempotency_key"] for delivery in read_deliveries()}


def dispatch(delivery: dict) -> list[dict]:
    results = []
    for name, handle in HANDLERS:
        try:
            handle(delivery)
            results.append({"handler": name, "ok": True})
        except Exception as error:  # noqa: BLE001 - one handler's failure must not stop the others
            app.logger.error("%s failed for %s: %s", name, delivery["idempotency_key"], error)
            results.append({"handler": name, "ok": False, "error": str(error)})
    return results


@app.post("/hooks/<campaign_id>")
def receive(campaign_id: str):
    raw_body = request.get_data()  # the exact bytes that were signed
    if not is_valid_signature(raw_body, request.headers.get("X-Signature-256")):
        abort(401)

    event = request.get_json()
    key = event.get("idempotency_key")
    if event.get("event") != "prospect_response" or key in processed:
        return "", 200

    delivery = {
        "received_at": datetime.now(timezone.utc).isoformat(),
        "campaign_id": campaign_id,
        "idempotency_key": key,
        "event": event,
    }
    try:
        with write_lock, open(DELIVERIES_FILE, "a", encoding="utf-8") as file:
            file.write(json.dumps(delivery) + "\n")
    except OSError as error:
        app.logger.error("Couldn't store %s: %s", key, error)
        abort(500)  # not stored, so let Victoria AI retry
    processed.add(key)

    # Stored: acknowledge now, run the handlers after the response.
    threading.Thread(target=dispatch, args=(delivery,), daemon=True).start()
    return "", 200


@app.post("/replay/<idempotency_key>")
def replay(idempotency_key: str):
    if not RELAY_TOKEN or request.headers.get("Authorization") != f"Bearer {RELAY_TOKEN}":
        abort(403)
    delivery = next((d for d in read_deliveries() if d["idempotency_key"] == idempotency_key), None)
    if delivery is None:
        abort(404)
    return jsonify({"idempotency_key": idempotency_key, "results": dispatch(delivery)})
```

Start it with `HANDLERS=console,slack node receiver.mjs`, or `HANDLERS=console,slack flask --app receiver run --port 3000`. In production, run the Flask app under a WSGI server such as Gunicorn.

## 2. Add handlers

A handler receives the stored delivery: `campaign_id` from the URL, `idempotency_key`, `received_at`, and the full `event`. These two ship with the receiver; the Google Sheets handler is its own recipe, [Log prospect replies to Google Sheets](https://docs.versionseven.ai/cookbook/log-replies-to-google-sheets).

**Node.js console**

```javascript
// handlers/console.mjs
export default async function handle({ campaign_id, event }) {
  const lead = event.lead ?? {};
  const name = [lead.first_name, lead.last_name].filter(Boolean).join(" ") || "A prospect";
  const sentiment = event.ai_response?.sentiment ?? "unread";
  console.log(`[${campaign_id}] ${sentiment} reply from ${name} to "${event.campaign}"`);
}
```

**Node.js Slack**

```javascript
// handlers/slack.mjs
const slackWebhookUrl = process.env.SLACK_WEBHOOK_URL;

// Slack treats &, < and > as control characters in message text.
function escapeForSlack(text) {
  return String(text).replace(/&/g, "&amp;").replace(/</g, "&lt;").replace(/>/g, "&gt;");
}

function slackMessage(event) {
  const lead = event.lead ?? {};
  const name = [lead.first_name, lead.last_name].filter(Boolean).join(" ") || "A prospect";
  const who = lead.company ? `${name} (${lead.company})` : name;
  const channel = event.channel === "linkedin" ? "on LinkedIn" : "by email";
  const reply = event.prospect_message ?? "";
  const excerpt = reply.length > 1000 ? `${reply.slice(0, 1000)}…` : reply;
  const sentiment = event.ai_response?.sentiment;

  const lines = [`*${escapeForSlack(who)}* replied ${channel} to *${escapeForSlack(event.campaign ?? "a campaign")}*`];
  if (sentiment) lines.push(`Sentiment: ${escapeForSlack(sentiment)}`);
  lines.push(...escapeForSlack(excerpt).split("\n").map((line) => `> ${line}`));
  return { text: lines.join("\n") };
}

export default async function handle({ event }) {
  const response = await fetch(slackWebhookUrl, {
    method: "POST",
    headers: { "Content-Type": "application/json" },
    body: JSON.stringify(slackMessage(event)),
    signal: AbortSignal.timeout(10_000),
  });
  if (!response.ok) throw new Error(`Slack answered ${response.status}`);
}
```

**Python console**

```python
# handlers/console.py
def handle(delivery: dict) -> None:
    event = delivery["event"]
    lead = event.get("lead") or {}
    name = " ".join(part for part in (lead.get("first_name"), lead.get("last_name")) if part) or "A prospect"
    sentiment = (event.get("ai_response") or {}).get("sentiment") or "unread"
    print(f'[{delivery["campaign_id"]}] {sentiment} reply from {name} to "{event.get("campaign")}"', flush=True)
```

**Python Slack**

```python
# handlers/slack.py
import os

import requests

SLACK_WEBHOOK_URL = os.environ["SLACK_WEBHOOK_URL"]


def escape_for_slack(text) -> str:
    # Slack treats &, < and > as control characters in message text.
    return str(text).replace("&", "&amp;").replace("<", "&lt;").replace(">", "&gt;")


def slack_message(event: dict) -> dict:
    lead = event.get("lead") or {}
    name = " ".join(part for part in (lead.get("first_name"), lead.get("last_name")) if part) or "A prospect"
    who = f"{name} ({lead['company']})" if lead.get("company") else name
    channel = "on LinkedIn" if event.get("channel") == "linkedin" else "by email"
    reply = event.get("prospect_message") or ""
    excerpt = reply[:1000] + "…" if len(reply) > 1000 else reply
    sentiment = (event.get("ai_response") or {}).get("sentiment")

    campaign = event.get("campaign") or "a campaign"
    lines = [f"*{escape_for_slack(who)}* replied {channel} to *{escape_for_slack(campaign)}*"]
    if sentiment:
        lines.append(f"Sentiment: {escape_for_slack(sentiment)}")
    lines.extend(f"> {line}" for line in escape_for_slack(excerpt).split("\n"))
    return {"text": "\n".join(lines)}


def handle(delivery: dict) -> None:
    response = requests.post(SLACK_WEBHOOK_URL, json=slack_message(delivery["event"]), timeout=10)
    response.raise_for_status()
```

A handler that only wants some replies decides for itself: check `event.ai_response.sentiment`, `agent_action` or `out_of_office` at the top and return.

## 3. Register the webhook

Register the receiver on each campaign, with the campaign's id in the path, using [`POST /v1/campaigns/{campaign_id}/webhooks`](https://docs.versionseven.ai/api-reference/campaigns/create-webhook):

```bash
CAMPAIGN_ID=550e8400-e29b-41d4-a716-446655440000
curl -X POST https://api.versionseven.ai/v1/campaigns/$CAMPAIGN_ID/webhooks \
  -H "Authorization: Bearer $VICTORIA_API_KEY" \
  -H "Content-Type: application/json" \
  -d "{\"webhook_url\": \"https://replies.example.com/hooks/$CAMPAIGN_ID\", \"secret\": \"$VICTORIA_WEBHOOK_SECRET\", \"replace\": true}"
```

`replace: true` matters here. A campaign has one webhook: if a different URL already holds it, the request without `replace` answers `409 WEBHOOK_SLOT_TAKEN` and names that host, so nothing moves by accident. With `replace`, the webhook is repointed at the receiver (`200`, with `replaced: true`) and the previous URL stops getting the campaign's events. That's the point of this recipe: whatever used to receive replies becomes a handler instead.

[`GET /v1/campaigns/{campaign_id}/webhooks`](https://docs.versionseven.ai/api-reference/campaigns/list-webhooks) shows what a campaign points at today. To change the secret later, [`POST /v1/campaigns/{campaign_id}/webhooks/{webhook_id}/secret`](https://docs.versionseven.ai/api-reference/campaigns/rotate-webhook-secret) sets a new one; restart the receiver with the new `VICTORIA_WEBHOOK_SECRET` first, then rotate, so no delivery is signed with a key the receiver doesn't hold.

## 4. Send a test delivery, then replay it

Sign one of the examples from [`GET /v1/webhooks/examples`](https://docs.versionseven.ai/api-reference/reference/list-webhook-examples) yourself and post it to the receiver:

```bash
BODY=$(curl -s https://api.versionseven.ai/v1/webhooks/examples -H "Authorization: Bearer $VICTORIA_API_KEY" | python3 -c 'import json,sys; print(json.dumps(json.load(sys.stdin)["examples"][0]))')
SIGNATURE=$(printf '%s' "$BODY" | openssl dgst -sha256 -hmac "$VICTORIA_WEBHOOK_SECRET" | sed 's/^.* //')
KEY=$(printf '%s' "$BODY" | python3 -c 'import json,sys; print(json.load(sys.stdin)["idempotency_key"])')

curl -X POST http://localhost:3000/hooks/$CAMPAIGN_ID \
  -H "Content-Type: application/json" \
  -H "X-Signature-256: sha256=$SIGNATURE" \
  --data "$BODY"

curl -X POST http://localhost:3000/replay/$KEY -H "Authorization: Bearer $RELAY_TOKEN"
```

The first request answers `200`, the console handler prints a line and the Slack message appears. Sending it again answers `200` and does nothing, because the key is stored. The replay answers with each handler's result, for example `{"results": [{"handler": "console", "ok": true}, {"handler": "slack", "ok": false, "error": "Slack answered 500"}]}`, which is how you re-run a handler after fixing whatever it talks to.

## What the receiver guarantees, and what it doesn't

- A delivery is stored before it's acknowledged. If the file can't be written, the receiver answers `500` and Victoria AI retries every 15 minutes, up to 5 attempts within 48 hours, with the same `idempotency_key`.
- Once acknowledged, a delivery is the receiver's. Victoria AI doesn't know whether a handler failed, so watch the log and use replay. For a receiver that should fail the whole delivery when one destination fails, do the work before answering, as [Post prospect replies to Slack](https://docs.versionseven.ai/cookbook/post-replies-to-slack) does.
- A lead's reply is delivered once per campaign, and again with a new key when a later reply is positive after an earlier one wasn't. A handler that counts replies should count keys, not leads.
- Several campaigns can share one receiver: register each campaign with its own `/hooks/<campaign_id>` path.

## Next steps

- [Log prospect replies to Google Sheets](https://docs.versionseven.ai/cookbook/log-replies-to-google-sheets), a handler for this receiver.
- [Post prospect replies to Slack](https://docs.versionseven.ai/cookbook/post-replies-to-slack), the smallest receiver, for when Slack is the only destination.
- The same hub without a server: [n8n](https://docs.versionseven.ai/cookbook/n8n-reply-hub-and-leads), [Zapier](https://docs.versionseven.ai/cookbook/zapier-reply-hub-and-leads) or [Make](https://docs.versionseven.ai/cookbook/make-reply-hub-and-leads).
- [Receiving webhooks](https://docs.versionseven.ai/guides/webhooks) for delivery, retries and the one-webhook rule.
