← Back to all posts

Engineering

Exactly-once is a WHERE clause: surviving duplicate carrier webhooks

Carriers deliver webhooks at least once, out of order, and sometimes twice under different ids. Here is how Handset turns that into one message status, one billed minute, and one event to your endpoint, with a unique index, a handful of guarded UPDATEs, and no dedupe service.

Jose Zamudio · Founder, Handset

·7 min read

A carrier webhook is a promise with the quality of a rumor. Telnyx tells us a message was delivered, and it means it, but it also reserves the right to tell us again in four seconds, to tell us after it told us the message failed, and to tell us under a brand-new event id because the first delivery timed out on our side and the retry got re-signed. None of that is a bug on their end. At-least-once is the only delivery guarantee a webhook can honestly make, and every carrier makes it.

What a customer sees on the other side of Handset should be none of that. A message is delivered once. A call minute is billed once. Your webhook endpoint gets one message.delivered event, not three. "Exactly once" is not something the carrier gives us, so it has to be something we build, and the interesting part is how little it takes: one unique index, a few UPDATE statements that refuse to run twice, and a transactional outbox. There is no dedupe service, no Redis set of seen ids, and no state machine table. The database already knows how to say no.

The ingest server does three things and nothing else

Carrier webhooks land on a process whose entire job fits in a package comment:

// The ingest server receives carrier webhooks. It does exactly three things,
// fast: verify the signature, persist the raw event (idempotently, keyed by
// the carrier's event id), and enqueue a processing job. All routing and
// business logic happens in workers — this process must answer in tens of
// milliseconds because call-control latency is dead air in the caller's ear.

Verification is Ed25519 over the timestamp and raw body, with a five-minute tolerance in both directions so a replayed or future-dated payload is rejected before it touches anything. After that, the handler does not look at what the event means. It writes it down and hands it off. If the write fails, it returns a 5xx, which makes the carrier retry, which is exactly the right outcome for a transient database problem.

The write is where the first layer of exactly-once lives:

CREATE TABLE carrier_events (
  id               text PRIMARY KEY,
  carrier          text NOT NULL DEFAULT 'telnyx',
  carrier_event_id text NOT NULL,
  event_type       text NOT NULL,
  payload          jsonb NOT NULL,
  processed_at     timestamptz,
  error            text,
  created_at       timestamptz NOT NULL DEFAULT now(),
  UNIQUE (carrier, carrier_event_id)
);

Every carrier event has an id, and the pair (carrier, carrier_event_id) is unique. Ingest does not check whether it has seen the event before. It inserts and lets the constraint decide:

return pgx.BeginFunc(ctx, s.Store.Pool, func(tx pgx.Tx) error {
	tag, err := tx.Exec(ctx, `
		INSERT INTO carrier_events (id, carrier, carrier_event_id, event_type, payload)
		VALUES ($1, $2, $3, $4, $5)
		ON CONFLICT (carrier, carrier_event_id) DO NOTHING`,
		ids.New(ids.CarrierEvent), carrierName, ev.ProviderEventID, string(ev.Type), payload)
	if err != nil {
		return err
	}
	if tag.RowsAffected() == 0 {
		return nil // duplicate delivery
	}
	_, err = s.River.InsertTx(ctx, tx, ProcessArgs{
		Carrier: carrierName, ProviderEventID: ev.ProviderEventID,
	}, opts)
	return err
})

The detail that matters is that the job enqueue is in the same transaction as the insert. A duplicate delivery gets zero rows affected, so it gets no job. There is no window where the event row exists but the job does not, and no window where two jobs exist for one row. A check-then-insert in application code has that window. ON CONFLICT DO NOTHING does not.

The event id is the wrong key for everything after this

If the story ended there it would be a short post. It does not, because the carrier's event id only identifies a delivery, not a fact. Telnyx can tell us "message X was delivered" under event id A, and then, after our side timed out and it retried, tell us the same thing under event id B. Both pass the unique index. Both get a job. And our own job runner retries too: a worker can finish the database work, crash before it marks the event processed, and run the whole handler again.

So the second layer keys on the resource, not the event. A delivery receipt does not ask "have I seen this event?" It asks "is there anything left to do to this message?" and lets the row answer:

tag, err := tx.Exec(ctx, `
	UPDATE messages SET status = $2, error_code = $3
	 WHERE id = $1 AND status NOT IN ('delivered', 'failed')`,
	messageID, newStatus, errCode)
if err != nil {
	return err
}
if tag.RowsAffected() == 0 {
	return nil // duplicate receipt
}

delivered and failed are terminal. Once a message is in one of them, no receipt moves it, including a receipt that says the opposite. The first terminal word from the carrier wins, and everything after the guard, recording the status history, emitting the customer event, happens only when the UPDATE actually changed a row.

The same shape shows up everywhere a carrier event changes state. Dispatching an outbound message moves it from queued to sent with WHERE status IN ('queued', 'scheduled'), and bills the segments only if tag.RowsAffected() == 1. A hangup is guarded by WHERE ended_at IS NULL, and the billable minutes are recorded inside the same transaction, so a second hangup finds nothing to end and nothing to charge:

err := tx.QueryRow(ctx, `
	UPDATE calls SET
	       status = CASE status
	                  WHEN 'ringing'     THEN 'missed'
	                  WHEN 'in_progress' THEN 'completed'
	                  ELSE status
	                END,
	       duration_seconds = $2,
	       ended_at = $3
	 WHERE id = $1 AND ended_at IS NULL
	 RETURNING status, transcribe`, callID, dur, ev.OccurredAt.UTC()).Scan(&finalStatus, &transcribed)
if errors.Is(err, pgx.ErrNoRows) {
	return nil // duplicate hangup
}

This is the whole trick. The guard is in the WHERE clause, which means it is evaluated under the same row lock as the write. Two workers racing on the same hangup cannot both see ended_at IS NULL, because the second one blocks on the first one's lock and then sees the row the first one committed. There is no separate "claims" table to keep consistent with the real one, because the real one is the claims table.

Out of order is handled by retrying, not by sorting

The one case a status guard cannot cover is an event that arrives too early. A delivery receipt can reach ingest before the dispatcher has stamped the message with its carrier reference, because the carrier accepted the message and fired the receipt in the time it took our transaction to commit. The receipt handler looks up the message by that reference and finds nothing.

We do not buffer it or sort events by their timestamps. The handler returns an error, the job runner retries with backoff, and by the second attempt the row is there. The event was already persisted at ingest, so nothing is lost; it just waits its turn. A state machine that reorders events by occurred_at would be more elegant and would have to be right about clocks on two companies' servers. A retry only has to be right about the fact that the dispatcher will eventually commit.

The one place we do use a compare-and-swap

Sequential ring is the feature where a call rings one phone, then the next, then the next, until someone answers. Each leg's hangup webhook is what triggers dialing the following leg. A duplicate hangup here is not harmless: it would dial the same person twice, or dial the next target while the previous one is still ringing.

The ring plan lives in a jsonb column with a cursor, and advancing the cursor is a compare-and-swap:

// Exactly-once claim of this position; losing it means a duplicate
// delivery's worker (or an answer) got there first.
tag, err := s.Store.Pool.Exec(ctx, `
	UPDATE calls SET ring_plan = jsonb_set(ring_plan, '{next}', to_jsonb($3::int))
	 WHERE id = $1 AND (ring_plan->>'next')::int = $2 AND status = 'ringing'`,
	callID, next, next+1)
if tag.RowsAffected() == 0 {
	return true, nil // someone else took this turn
}

The worker says "I believe the next target is index 2, and I am claiming it by setting it to 3." If another worker already did that, the predicate fails, zero rows change, and this worker walks away without dialing. The status = 'ringing' condition also stops a late hangup from dialing another leg after someone has already picked up. It is still a WHERE clause doing the work; it just happens to compare a value inside a JSON document instead of a column.

One transition, one event to you

Everything above would be wasted if the customer-facing side re-introduced duplicates. It does not, because the outbound event is written inside the same transaction as the state change it describes:

// Emit must be called inside the same transaction as the state change that
// produced the event — event and change commit or roll back together.
func (e *Emitter) Emit(ctx context.Context, tx pgx.Tx, accountID string, tenantID *string, eventType string, data any) error {
	eventID := ids.New(ids.Event)
	if _, err := tx.Exec(ctx, `
		INSERT INTO outbound_events (id, account_id, tenant_id, type, event_version, payload)
		VALUES ($1, $2, $3, $4, $5, $6)`,
		eventID, accountID, tenantID, eventType, Version, payload); err != nil {
		return err
	}
	// ... one webhook_deliveries row per subscribed endpoint, plus a delivery job

If the guarded UPDATE touched zero rows, Emit is never called, so a duplicate carrier receipt produces no customer event at all. If it touched one row and the transaction rolled back, the event rolls back with it. The id on the event you receive is the id of that row, and it stays the same across every delivery attempt.

Which brings up the honest caveat. Delivery to your endpoint is at-least-once, for the same reason the carrier's delivery to us is: we can get a 2xx from you and crash before recording it, and then we will send that event again with the same id. We retry up to eight times with backoff, sign every attempt with a fresh timestamp, and stop calling an endpoint after a hundred consecutive failures. If your handler is not idempotent on the event id, it should be. Exactly-once on our side is the property that the same thing never happens twice; the same event can still be said twice.

What we learned

The phrase "exactly-once" hides a lot of sins when it is implemented as "we remember what we have seen." Memory has to be kept somewhere, kept consistent with the thing it is guarding, and kept forever. Every layer here avoids that by making the guard and the write the same statement: a unique index on the raw event, a WHERE status NOT IN (...) on the message, a WHERE ended_at IS NULL on the call, a jsonb_set behind a predicate on the ring cursor, and an outbox row in the transaction that earned it. The billing ledger does not even have its own dedupe; it writes a usage row in the same transaction as the transition that made the minute billable, and inherits exactly-once from the guard.

The tell that it works is what the tests look like. They do not mock a carrier that behaves. They replay the same inbound SMS through the whole pipeline and count one message. They deliver the same hangup twice to a sequential ring and count the same number of legs. They re-run a failure marker and count one message.failed event. The carrier gets to be a rumor. The row is the fact.

Share:XLinkedIn

More from Handset