How to Trigger Automated Actions from Analytics Events
Build automated actions from analytics events in Node.js and PostgreSQL: rules, an outbox, webhooks, idempotent retries and rate limits that stop alert fatigue.
An Acme Shop visitor adds a camera to the cart, gets distracted and closes the tab. Your analytics table records that exact moment, and then nobody does anything with it. The data is a few milliseconds away from an action, yet it sits in a dashboard waiting for a human to notice.
Automated actions from analytics events close that gap. A rule watches the event stream, decides something matters, and calls a webhook: a recovery email, a Slack message to the sales team, a flag on a user. The idea is simple. The engineering is not, because actions have side effects, and a duplicated side effect is worse than a duplicated event.
In this article you will build a small action engine on the stack you already have: Node.js, pg and the events table. It supports rules that fire on an event and rules that fire when an event does not happen. It records every action in an outbox table, sends signed webhooks, retries safely and limits how often one user or one team gets pinged.
It assumes the collector from the Node.js analytics API and the buffering from async event tracking. It does not cover brokers; for that, see event-driven analytics architecture.
The Problem: Events Are Facts, Actions Are Promises
An event says “this happened”. An action says “do something about it”. The first can be duplicated, delayed or replayed without harm, because a duplicate row is mostly a counting error.
The second cannot. Sending a customer two discount emails costs trust and money.
That asymmetry drives every design decision below. You need a place to record “we decided to act” separately from “we acted”, you need a key that prevents duplicate decisions, and you need retries that cannot double-send. Here are the three actions Acme Shop will implement:
- Welcome webhook: on
signup_completed, call the email service once per user. - High-value purchase alert: on
purchase_completedabove a threshold, post to the sales channel. - Cart recovery: when a session has
add_to_cartbut nopurchase_completedafter 60 minutes, call the email service once per session.
The first two react to an event. The third reacts to the absence of one, which is why a pure per-event trigger is not enough.
Architecture for Automated Actions from Analytics Events
collector (POST /v1/events)
|
| INSERT INTO events ...
v
evaluateEventRules(event) ---- decides ----+
|
sweepAbsenceRules() (every minute) ---------+
v
actions (outbox table)
UNIQUE (dedupe_key)
|
NOTIFY wake-up | poll fallback
v
dispatcher (SKIP LOCKED)
|
signed webhook (HTTPS)
status: pending, sent, failed, dead, skipped
Every component has one job. Rules decide and never call the network. The outbox stores the decision durably in PostgreSQL, so a crash cannot lose it.
The dispatcher performs the side effect and records the outcome. This is the transactional outbox pattern, and it keeps the decision and the event write in the same transaction.
The trade-off is latency and polling. A database outbox is not a message broker, and you will hit its limits at high volume. For an online store sending thousands of actions an hour, it is more than enough, and it removes an entire infrastructure component.
The Outbox Table
This table extends your schema. It does not change events; it adds an action log beside it.
CREATE TABLE actions (
action_id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
rule_name text NOT NULL,
dedupe_key text NOT NULL UNIQUE,
subject_key text NOT NULL,
payload jsonb NOT NULL,
status text NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'sent', 'failed', 'dead', 'skipped')),
attempts int NOT NULL DEFAULT 0,
run_after timestamptz NOT NULL DEFAULT now(),
created_at timestamptz NOT NULL DEFAULT now(),
sent_at timestamptz
);
CREATE INDEX actions_ready_idx ON actions (run_after)
WHERE status IN ('pending', 'failed');
CREATE INDEX actions_subject_idx ON actions (rule_name, subject_key, created_at);
Three columns carry the design. dedupe_key is unique, so inserting the same decision twice does nothing. subject_key names who the action is about (a user or a session), which the cooldown check uses. The partial index keeps the dispatcher’s lookup small even after millions of finished rows; see PostgreSQL performance tuning for analytics for why partial indexes fit queues.
Rules That Fire on an Event
Define rules as plain objects. A rule has a name, a match function, a dedupe key builder and a payload builder. Keep them pure so you can unit test each one without a database.
// rules.js
export const eventRules = [
{
name: 'welcome_webhook',
matches: (e) => e.event_name === 'signup_completed' && e.user_id,
subject: (e) => `user:${e.user_id}`,
dedupeKey: (e) => `welcome_webhook:${e.user_id}`,
payload: (e) => ({ user_id: e.user_id, occurred_at: e.occurred_at }),
cooldownSeconds: 0,
},
{
name: 'high_value_purchase',
matches: (e) =>
e.event_name === 'purchase_completed' &&
Number(e.properties?.amount_cents) >= 50000,
subject: () => 'team:sales',
dedupeKey: (e) => `high_value_purchase:${e.event_id}`,
payload: (e) => ({
event_id: e.event_id,
amount_cents: Number(e.properties.amount_cents),
page_url: e.page_url,
}),
cooldownSeconds: 900,
},
];
Notice the dedupe keys. The welcome rule keys on the user, so ten signup_completed events for one user produce one webhook. The purchase rule keys on event_id, so a client retry of the same purchase produces one alert, and two different purchases produce two.
Choosing the key is choosing what “the same thing” means, and it is the most important line in each rule. The purchase rule also sets a 15 minute cooldown, which the cooldown section explains.
We once double-counted purchases because a retry reached the collector twice with new event IDs. That was a client bug, and the fix was in the client. The same bug here would double-send alerts.
Dedupe keys built on stable business identifiers, such as a user ID or an order ID from properties, protect you even when event IDs misbehave. For a purchase, prefer properties.order_id if you track one.
Evaluating rules inside the insert transaction
// evaluate.js
import { eventRules } from './rules.js';
import { inCooldown } from './cooldown.js';
export async function enqueueForEvent(client, event) {
for (const rule of eventRules) {
if (!rule.matches(event)) continue;
const subject = rule.subject(event);
if (await inCooldown(client, rule.name, subject, rule.cooldownSeconds)) continue;
const { rowCount } = await client.query(
`INSERT INTO actions (rule_name, dedupe_key, subject_key, payload)
VALUES ($1, $2, $3, $4)
ON CONFLICT (dedupe_key) DO NOTHING`,
[rule.name, rule.dedupeKey(event), subject, JSON.stringify(rule.payload(event))]
);
if (rowCount > 0) {
await client.query(`SELECT pg_notify('action_ready', $1)`, [rule.name]);
}
}
}
Call enqueueForEvent in the same transaction that inserts the event, right after the INSERT INTO events. If the transaction rolls back, no action exists. If it commits, both exist. That removes the classic failure where an event is stored but its action is lost, or an action is sent for an event that never saved.
Here is that wiring, written against the collector’s insert. Only a newly stored event reaches the rules, so a client retry that hits ON CONFLICT DO NOTHING never evaluates twice.
// store.js
import { enqueueForEvent } from './evaluate.js';
export async function storeEvent(pool, event) {
const client = await pool.connect();
try {
await client.query('BEGIN');
const { rowCount } = await client.query(
`INSERT INTO events
(event_id, event_name, occurred_at, anonymous_id, user_id,
session_id, page_url, referrer, user_agent, properties)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
ON CONFLICT (event_id) DO NOTHING`,
[event.event_id, event.event_name, event.occurred_at, event.anonymous_id,
event.user_id, event.session_id, event.page_url, event.referrer,
event.user_agent, JSON.stringify(event.properties)]
);
if (rowCount === 1) await enqueueForEvent(client, event);
await client.query('COMMIT');
} catch (err) {
await client.query('ROLLBACK');
throw err;
} finally {
client.release();
}
}
The cost is a slightly longer insert transaction. If you batch inserts as in the async article, loop over the batch inside one transaction. Keep rule functions fast and free of network calls, or your ingestion latency will carry their latency.
Rules That Fire on Absence
Cart recovery depends on something not happening. No event arrives to trigger it, so a scheduled query has to look for it. Run this sweep every minute.
// sweep.js
export async function sweepCartAbandonment(pool) {
await pool.query(
`INSERT INTO actions (rule_name, dedupe_key, subject_key, payload, run_after)
SELECT 'cart_recovery',
'cart_recovery:' || c.session_id,
'session:' || c.session_id,
jsonb_build_object('session_id', c.session_id,
'anonymous_id', c.anonymous_id,
'last_cart_at', c.last_cart_at),
now() + interval '5 minutes'
FROM (
SELECT session_id, min(anonymous_id) AS anonymous_id,
max(occurred_at) AS last_cart_at
FROM events
WHERE event_name = 'add_to_cart'
AND occurred_at >= now() - interval '6 hours'
GROUP BY session_id
HAVING max(occurred_at) < now() - interval '60 minutes'
) c
WHERE NOT EXISTS (
SELECT 1 FROM events p
WHERE p.session_id = c.session_id
AND p.event_name = 'purchase_completed'
AND p.occurred_at >= c.last_cart_at
)
ON CONFLICT (dedupe_key) DO NOTHING`
);
}
The time window has two edges for a reason. The HAVING clause defines “abandoned”: the latest cart event in the session is more than 60 minutes old. The 6 hour filter bounds the scan so the query stays cheap and can use the (event_name, occurred_at) index from the tuning article. The dedupe key on session_id means the same session produces one recovery action, no matter how many sweeps see it.
My first draft put the 60 minute cutoff in the WHERE clause instead. A shopper who added an item at minute 0 and another at minute 50 then looked abandoned at minute 60, because the minute 50 event was filtered out before max() ran. Moving the cutoff into HAVING fixes it, since the maximum now covers every cart event in the window.
There is a subtle bug to avoid. A late-arriving purchase_completed event, delayed by a mobile client retrying offline, can land after the sweep already queued the recovery email. Therefore the sweep sets a five minute run_after delay, and the dispatcher re-checks the condition just before sending (see stillValid below). A short grace period costs little and removes the most embarrassing false positive: emailing someone who just bought.
Note also that this sweep identifies a session, not a person. To email someone, you need a user_id and a consent record, which this series does not handle. Do not send marketing email to an anonymous visitor. In this example the webhook receiver does the lookup and consent check.
The Dispatcher: Claim, Send, Record
The dispatcher claims a batch of ready rows, sends each one, and records the result. The key tool is row locking with SKIP LOCKED, documented in the PostgreSQL SELECT reference. It lets several workers pull from the same table without blocking each other or taking the same row.
// dispatcher.js
import crypto from 'node:crypto';
import pg from 'pg';
const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL, max: 4 });
const WEBHOOK_URL = process.env.ACTION_WEBHOOK_URL; // set by you, never by users
const WEBHOOK_SECRET = process.env.ACTION_WEBHOOK_SECRET;
const MAX_ATTEMPTS = 6;
function sign(body) {
return crypto.createHmac('sha256', WEBHOOK_SECRET).update(body).digest('hex');
}
// Conditions that must still hold at send time. Return false to skip the action.
const stillValid = {
async cart_recovery(client, action) {
const { rowCount } = await client.query(
`SELECT 1 FROM events
WHERE session_id = $1 AND event_name = 'purchase_completed'
AND occurred_at >= $2
LIMIT 1`,
[action.payload.session_id, action.payload.last_cart_at]
);
return rowCount === 0;
},
};
async function claimBatch(client, limit) {
const { rows } = await client.query(
`SELECT action_id, rule_name, dedupe_key, payload, attempts
FROM actions
WHERE status IN ('pending', 'failed') AND run_after <= now()
ORDER BY run_after
LIMIT $1
FOR UPDATE SKIP LOCKED`,
[limit]
);
return rows;
}
async function deliver(action) {
const body = JSON.stringify({ rule: action.rule_name, data: action.payload });
const res = await fetch(WEBHOOK_URL, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Idempotency-Key': action.dedupe_key,
'X-Signature': sign(body),
},
body,
signal: AbortSignal.timeout(5000),
});
if (!res.ok) throw new Error(`webhook returned ${res.status}`);
}
export async function runOnce() {
const client = await pool.connect();
try {
await client.query('BEGIN');
const batch = await claimBatch(client, 20);
for (const action of batch) {
try {
const check = stillValid[action.rule_name];
if (check && !(await check(client, action))) {
await client.query(
`UPDATE actions SET status = 'skipped' WHERE action_id = $1`,
[action.action_id]);
continue;
}
await deliver(action);
await client.query(
`UPDATE actions SET status = 'sent', sent_at = now(), attempts = attempts + 1
WHERE action_id = $1`, [action.action_id]);
} catch (err) {
const attempts = action.attempts + 1;
const dead = attempts >= MAX_ATTEMPTS;
const delaySeconds = Math.min(3600, 2 ** attempts * 15);
await client.query(
`UPDATE actions
SET status = $2, attempts = $3,
run_after = now() + make_interval(secs => $4)
WHERE action_id = $1`,
[action.action_id, dead ? 'dead' : 'failed', attempts, delaySeconds]);
}
}
await client.query('COMMIT');
return batch.length;
} catch (err) {
await client.query('ROLLBACK');
throw err;
} finally {
client.release();
}
}
while (true) {
const n = await runOnce().catch((e) => { console.error(e); return 0; });
if (n === 0) await new Promise((r) => setTimeout(r, 2000));
}
Read the failure paths carefully. An action whose precondition no longer holds becomes skipped, which keeps an audit trail without sending anything. The webhook gets a five second timeout, so a hung receiver cannot hold database locks forever.
Failures back off exponentially up to an hour. After six attempts the row becomes dead, which stops the infinite retry loop and leaves a row for a human to inspect.
The remaining hazard is a crash between a successful HTTP call and the UPDATE. On restart, the row is still pending, so the dispatcher sends again. This is at-least-once delivery, and you cannot make it exactly-once across a network boundary.
The defense is the Idempotency-Key header: the receiver stores keys it has processed and ignores repeats. Document that requirement for whoever owns the receiver.
One design note. This dispatcher holds a transaction open while it calls the network. With a five second timeout and a batch of 20, the worst case is long, but it keeps the code short.
If you need higher throughput, claim rows by setting a locked_until column, commit, send, then update. The simpler form is a reasonable start.
For signing, the receiver should recompute the HMAC with the shared secret and compare it using a constant-time check. The Node.js crypto module provides createHmac and timingSafeEqual for this. Never take the webhook URL from event data; an attacker who controls an event property could otherwise aim your server at internal addresses.
Waking the Dispatcher with LISTEN/NOTIFY
Polling every two seconds is fine, but you can cut latency with LISTEN/NOTIFY. The PostgreSQL NOTIFY documentation says notifications are delivered only after the sending transaction commits, and the payload must stay under 8000 bytes in the default configuration.
// enqueueForEvent already calls pg_notify('action_ready', ...) after each new row.
// in dispatcher.js, on a dedicated connection:
const listener = await pool.connect();
await listener.query('LISTEN action_ready');
listener.on('notification', () => { runOnce().catch(console.error); });
Treat the notification as a hint, never as the work. A listener that is disconnected misses notifications, and PostgreSQL does not replay them. The polling loop is the safety net that finds anything missed. Send only a short label in the payload and read the real data from the table.
Use a dedicated connection for LISTEN. A pooled connection that returns to the pool loses its listening state, and a pooler in transaction mode, such as PgBouncer, does not support it at all.
Rate Limits, Cooldowns and Alert Fatigue
Rules multiply. A threshold that fires ten times a day in a quiet week fires four hundred times on Black Friday. A team that gets four hundred messages mutes the channel, and the next real alarm goes unseen. That is alert fatigue, and it is an engineering failure, not a human one.
Add a cooldown check before inserting into the outbox. This one allows at most one action per rule and subject within a window:
// cooldown.js
export async function inCooldown(client, ruleName, subjectKey, seconds) {
if (!seconds) return false;
const { rowCount } = await client.query(
`SELECT 1 FROM actions
WHERE rule_name = $1 AND subject_key = $2
AND created_at > now() - make_interval(secs => $3)
LIMIT 1`,
[ruleName, subjectKey, seconds]
);
return rowCount > 0;
}
enqueueForEvent above already calls it before each insert and skips the rule when it returns true. The actions_subject_idx index serves this lookup. Tune the window per rule: a welcome email gets none because the dedupe key already handles it, the sales alert gets 15 minutes, and a “conversion rate dropped” alert might get thirty. The cost is that suppressed events leave no action row, so count them in your own logs if you need to know how much was muted.
Cooldowns differ from dedupe keys. A dedupe key says “this exact cause is already handled”. A cooldown says “this kind of thing happened recently, so stay quiet”. You usually need both.
| Control | Stops | Example | Cost |
|---|---|---|---|
| Dedupe key | Same cause acted on twice | One cart email per session | Key choice must be right |
| Cooldown | Bursts from the same rule and subject | One sales alert per 15 minutes | Suppressed events lose detail |
| Per-user daily cap | Too many messages to one person | Max 2 emails per user per day | Needs a cross-rule count |
| Digest | Many small alerts | Hourly summary of purchases | Adds delay |
| Severity routing | Noise in the urgent channel | Page only on checkout failures | Needs agreed severity levels |
A mistake I have seen in production is alerting on a raw count threshold with no baseline. A rule that fired when errors exceeded 50 an hour sent constant noise in the daytime and nothing at night when a real outage dropped traffic to zero. Prefer rates against a baseline, and pair every alert rule with an owner and a documented response. If nobody knows what to do when it fires, delete it.
Testing and Observing the Engine
Unit test rules as pure functions with sample events. Then test the pipeline end to end by inserting an event and asserting on the outbox row. Before launch, run the dispatcher against a local mock receiver that fails randomly, and check that every action ends as sent exactly once from the receiver’s point of view.
Monitor the outbox with queries, not a new tool. These three answer most questions:
-- backlog and age of the oldest ready action
SELECT count(*), min(run_after) FROM actions
WHERE status IN ('pending', 'failed') AND run_after <= now();
-- dead actions needing a human
SELECT rule_name, count(*) FROM actions WHERE status = 'dead' GROUP BY 1;
-- actions per rule in the last day
SELECT rule_name, status, count(*) FROM actions
WHERE created_at > now() - interval '1 day' GROUP BY 1, 2 ORDER BY 1, 2;
Archive or delete old sent and skipped rows on a schedule, for example after 30 days, so the table does not grow without bound.
How Real Systems Do This
Product analytics tools ship this feature as a first-class concept. Segment calls it destinations and functions, Mixpanel and Amplitude offer behavioral cohorts that sync to messaging tools, and PostHog provides webhooks and a pipeline of destinations driven by events. Check each product’s current documentation for the exact feature names. They all solve the same problems you just handled: matching, deduplication, retries and rate limits.
Webhook providers set the conventions you copied. Signed payloads with an HMAC header and an idempotency key are the common contract, and receivers are expected to tolerate repeats. When you operate a receiver, treat every webhook as possibly duplicated, because it is.
At larger scale, teams move the matching out of the database and into a stream processor reading from a broker, which is the territory of the architecture article. The outbox and the idempotent action contract stay the same, so what you build here carries over.
Decision Framework
- Is the action reversible and cheap, such as a Slack message? If not, such as email, billing or account changes, require a dedupe key and a human-reviewed rule.
- Does the rule react to an event or to the absence of one? Use an insert-time rule for the first and a scheduled sweep for the second.
- What does “the same thing” mean? Pick a dedupe key from a stable business identifier.
- Can the decision share a transaction with the event write? If yes, use the outbox in the same transaction.
- How many actions per hour? As a rule of thumb, a few thousand an hour is comfortable in PostgreSQL. Beyond that, load test the dispatcher and plan for a broker.
- Who owns each alert, and what do they do when it fires? If the answer is unclear, do not ship the rule.
When NOT to Use This
- Standard lifecycle messaging: if you need onboarding emails, push campaigns and segmentation, buy a customer engagement platform. Building consent management, unsubscribe handling and deliverability yourself costs more than the subscription.
- Actions with legal or financial weight: refunds, account bans and credit decisions should not hang off raw analytics events. Analytics data has duplicates, gaps and bot traffic, so route those through your transactional system.
- Sub-second reactions at high volume: fraud blocking in the request path needs a synchronous check or a stream processor, not a polling outbox.
Common Mistakes
- Calling the webhook inside the collector request, so a slow receiver slows ingestion and a failed call loses the action.
- Skipping the dedupe key, which sends a customer the same email on every retry or replay.
- Using
LISTEN/NOTIFYas the queue, so a restart drops notifications and actions never run. - Retrying forever with no dead state, which hides a broken receiver behind an endless loop.
- Alerting on raw counts without cooldowns, which trains the team to ignore the channel.
- Accepting the webhook URL from event data, which lets an attacker aim your server at internal hosts.
Key Takeaways
- Separate deciding from doing: rules write to an outbox, and a dispatcher performs the side effect.
- Write the outbox row in the same transaction as the event, with a unique dedupe key.
- Choose dedupe keys from stable business identifiers, not just
event_id. - Use scheduled sweeps for absence rules, and re-check the condition before sending.
- Claim work with
FOR UPDATE SKIP LOCKED, back off on failure, and mark exhausted actions as dead. - Treat
LISTEN/NOTIFYas a wake-up hint backed by polling. - Add cooldowns, caps and digests before the first noisy day, not after.
FAQ
How do I trigger an action when an analytics event happens?
Evaluate rules when the event is stored, write a row to an outbox table in the same transaction, and let a separate dispatcher send the webhook. A unique dedupe key on the outbox row prevents duplicate actions on retries.
Is LISTEN/NOTIFY reliable enough for a job queue?
No, not by itself. Notifications reach only connected listeners and are not replayed. Use it to wake a worker quickly, and keep a table with a polling query as the source of truth.
How do I avoid sending the same webhook twice?
You cannot fully prevent it, because a crash can happen after the call succeeds. Send an idempotency key with every request, and make the receiver ignore keys it has already processed. That gives effectively-once behavior over at-least-once delivery.
How do I trigger an action when something does not happen, like an abandoned cart?
Run a scheduled query that finds sessions with an add_to_cart event and no later purchase_completed after a set delay. Insert a deduplicated action for each, and re-check the condition just before sending.
Conclusion
Automated actions from analytics events turn analytics from a report into a reflex. The engine is small: rules, an outbox with dedupe keys, and a careful dispatcher. The hard part is restraint, so decide what must never happen twice and what is not worth a notification.
Rule of thumb: record the decision before you perform the action, and never perform the same action twice for the same cause. Next, move from rules to questions with text-to-SQL on a local LLM.
Last updated on 9 October 2026.
