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.
Last updated
- Slack
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 isopenssl rand -hex 32. - A public HTTPS address for the receiver. Victoria AI doesn't deliver to
localhostor 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 URL in
SLACK_WEBHOOK_URL. - An API key with
campaigns:writeinVICTORIA_API_KEY, to register the webhook. See Authentication. - Node.js 18 or later with
express, or Python 3.10 or later withflaskandrequests.
How it works
- Victoria AI posts a
prospect_responseevent to/hooks/<campaign_id>. The payload names the campaign but carries no id, so the id lives in the URL you register. - The receiver checks
X-Signature-256against the raw body, skips anidempotency_keyit has seen, and appends the delivery to a file. - It answers
200at once. Victoria AI then considers the delivery done and won't retry it, so the receiver owns what happens next. - Each handler in
HANDLERSruns in turn. A failure is logged and doesn't stop the others. POST /replay/<idempotency_key>runs the handlers again for a stored delivery.
1. Run the receiver
// 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(", ")}`));# 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.
// 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}"`);
}// 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, "&").replace(/</g, "<").replace(/>/g, ">");
}
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}`);
}# 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)# 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("&", "&").replace("<", "<").replace(">", ">")
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:
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 shows what a campaign points at today. To change the secret later, POST /v1/campaigns/{campaign_id}/webhooks/{webhook_id}/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 yourself and post it to the receiver:
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
500and Victoria AI retries every 15 minutes, up to 5 attempts within 48 hours, with the sameidempotency_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 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, a handler for this receiver.
- Post prospect replies to Slack, the smallest receiver, for when Slack is the only destination.
- The same hub without a server: n8n, Zapier or Make.
- Receiving webhooks for delivery, retries and the one-webhook rule.