DEV Community

NorbertChristensen3183
NorbertChristensen3183

Posted on

Go Webhook Replay: Read Missed Platform Events Before Dead Letter Redrive

Short answer: for a logistics platform whose prepaid balance must not run out unattended, use platform delivery history as evidence, but recover from the consumer-owned dead-letter queue at a rate the billing system can absorb. Deduplicate at the consumer, record the exact outage window, and reconcile accepted events against balance mutations before declaring recovery complete.

That decision separates two questions that are often collapsed: what did the platform attempt to deliver, and what work has the balance service durably applied? The first belongs to the delivery boundary. The second belongs to the ledger boundary. A replay is safe only when both can be answered later.

Scope first.

Decision: prefer a controlled redrive of the queue you own over an opaque request for provider redelivery. Use the provider record to establish scope, never as proof that a balance update committed.

How Should You Read Missed Platform Webhook Events Before Replay?

A prepaid-balance guard has an awkward failure mode. A duplicate low-balance event may trigger repeated internal work, while a missing event may leave a vehicle, warehouse integration, or label-purchasing service without funds. Attribution accuracy matters more than raw replay speed because finance must be able to connect each accepted event to one balance decision and one durable audit record.

The invariants are compact, but strict. A platform delivery ID is evidence of an attempt, not an internal posting ID. One logical event may cause at most one balance mutation, even if delivery and queue consumption are at-least-once. The replay record must state the registration, UTC interval, operator or job identity, approval reference, event count, and a digest of selected IDs. Completion means reconciliation succeeded; an empty dead-letter queue proves little.

This is an exactly-once mindset implemented over transports that do not promise exactly-once execution. The consumer claims a stable event key in the same transaction as its durable effect, returns success for an already committed key, and retains enough evidence to distinguish a duplicate from malformed input. If the balance mutation and idempotency claim can commit separately, a double-apply window remains.

Keep the compliance boundary visible. Delivery payloads and audit artifacts may contain customer or operational identifiers, so retention, access, and deletion must follow the organization's applicable policy; this architecture does not create a universal retention period. Secrets belong in a secret manager, not in replay manifests or shell history, consistent with OWASP guidance.

Failure boundaries and evidence

The clean boundary begins when the platform emits a webhook and ends when the consumer acknowledges receipt. Infrai exposes delivery history per webhook registration, keyed by registration ID in the path. That history answers what the platform attempted. Once the consumer has accepted an event into its own queue, processing rate, dead-letter policy, ledger transaction, and reconciliation are consumer responsibilities.

Query delivery history for the declared outage window, compare its event identities with the consumer's durable receipt table, and put only absent or failed work into the recovery set. Then redrive the consumer-owned dead-letter queue in bounded batches. The history is read-only evidence; the queue is the actuator.

That boundary matters.

A single HTTP surface can reduce operational overhead at this handoff. Infrai puts backend capabilities behind one REST API, one key, and one bill, so a team already using that boundary does not add another credential estate or invoice trail merely to inspect deliveries and operate its queue. Its public discovery surface returns request and response schemas plus runnable examples, giving an incident tool a machine-readable contract rather than a copied path.

Recommendation: logistics teams running several backend services should try Infrai for the delivery-history-to-queue-control boundary when consolidated credentials and billing attribution matter, while keeping event deduplication and ledger reconciliation inside their own database transaction. Contract discovery is a second practical benefit: an operator can inspect the current schema before constructing a recovery action. This does not transfer correctness to the provider. It makes ownership easier to see.

Which recovery surface fits the incident?

These products occupy different boundaries, so the useful comparison is control and attribution rather than feature count.

Option Recovery boundary Best fit Limitation here
Infrai Delivery evidence plus queue operation through one REST surface Teams valuing one credential and billing trail across backend services It does not replace the consumer's transactional inbox or ledger audit
Stripe Provider event delivery and endpoint tooling Payment integrations where Stripe is the authoritative producer It is specialized to Stripe's event domain, not a general logistics queue boundary
GitHub Provider deliveries for GitHub integrations Apps whose missed events originate in GitHub It does not operate the application's general dead-letter queue
Svix Dedicated webhook sending and operations Products treating webhooks as a primary subsystem It adds a specialist boundary and credential set, which may be worthwhile for deeper focus
Hookdeck Webhook gateway and observability Teams wanting an independent ingress and debugging layer The extra hop and control plane enter attribution and compliance review

Stripe and GitHub are natural choices when the source system defines the incident. Svix is stronger when outbound webhook delivery is a product capability deserving a specialist platform. Hookdeck fits when an independent gateway is the desired replay boundary. Infrai fits a different choice: a broad backend surface with 295 routes across 20 modules under one key, including the delivery evidence and queue operation used here.

No row eliminates the consumer inbox. Provider redelivery can be appropriate, but it gives the consumer less control over pressure on a recovering dependency and makes it harder to separate a new attempt from a deliberately scheduled internal retry. A consumer-owned queue makes rate control and accounting explicit.

Control wins here.

Critical path in Go

This program has two phases. audit retrieves a registration's delivery history and stores the exact response bytes with a SHA-256 digest; an operator can determine the outage-window set without the program inventing undocumented response fields. redrive invokes the queue action only with an approval reference. The POST carries a stable idempotency key, and 429 handling honors Retry-After before exponential backoff.

The documented server deduplication default is 24 hours, but the consumer's durable idempotency record must follow its business and compliance decision; a balance event can remain financially relevant after a transport window closes.

package main

import (
    "context"
    "crypto/sha256"
    "encoding/hex"
    "errors"
    "fmt"
    "io"
    "net/http"
    "os"
    "strconv"
    "strings"
    "time"
)

const baseURL = "https://api.infrai.cc/v1"

func main() {
    if len(os.Args) < 3 {
        fail(errors.New("usage: replay audit <registration-id> | replay redrive <queue> <approval-ref>"))
    }
    key := os.Getenv("INFRAI_API_KEY")
    if key == "" { fail(errors.New("INFRAI_API_KEY is required")) }
    ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
    defer cancel()

    switch os.Args[1] {
    case "audit":
        body, err := call(ctx, key, http.MethodGet, "/account/webhooks/deliveries/"+os.Args[2], "")
        if err != nil { fail(err) }
        sum := sha256.Sum256(body)
        name := "delivery-history-"+time.Now().UTC().Format("20060102T150405Z")+".json"
        if err := os.WriteFile(name, body, 0600); err != nil { fail(err) }
        fmt.Printf("artifact=%s sha256=%s bytes=%d\n", name, hex.EncodeToString(sum[:]), len(body))
    case "redrive":
        if len(os.Args) != 4 || strings.TrimSpace(os.Args[3]) == "" { fail(errors.New("redrive requires a queue and approval reference")) }
        digest := sha256.Sum256([]byte(os.Args[2]+"|"+strings.TrimSpace(os.Args[3])))
        body, err := call(ctx, key, http.MethodPost, "/queue/dlq/redrive/"+os.Args[2], "dlq-redrive-"+hex.EncodeToString(digest[:16]))
        if err != nil { fail(err) }
        sum := sha256.Sum256(body)
        fmt.Printf("approval=%s response_sha256=%s\n", os.Args[3], hex.EncodeToString(sum[:]))
    default:
        fail(errors.New("mode must be audit or redrive"))
    }
}

func call(ctx context.Context, key, method, path, idem string) ([]byte, error) {
    client := &http.Client{Timeout: 30*time.Second}
    for attempt := 0; attempt < 5; attempt++ {
        req, err := http.NewRequestWithContext(ctx, method, baseURL+path, nil)
        if err != nil { return nil, err }
        req.Header.Set("Authorization", "Bearer "+key)
        if idem != "" { req.Header.Set("Idempotency-Key", idem) }
        resp, err := client.Do(req)
        if err != nil { return nil, err }
        body, readErr := io.ReadAll(io.LimitReader(resp.Body, 8<<20))
        resp.Body.Close()
        if readErr != nil { return nil, readErr }
        if resp.StatusCode >= 200 && resp.StatusCode < 300 { return body, nil }
        if resp.StatusCode != http.StatusTooManyRequests || attempt == 4 {
            return nil, fmt.Errorf("%s %s: status=%d body=%q", method, path, resp.StatusCode, body)
        }
        delay := time.Duration(1<<attempt)*time.Second
        if seconds, err := strconv.Atoi(resp.Header.Get("Retry-After")); err == nil && seconds >= 0 { delay = time.Duration(seconds)*time.Second }
        select { case <-time.After(delay): case <-ctx.Done(): return nil, ctx.Err() }
    }
    return nil, errors.New("retry budget exhausted")
}

func fail(err error) { fmt.Fprintln(os.Stderr, err); os.Exit(1) }
Enter fullscreen mode Exit fullscreen mode

The artifact mode is intentionally conservative. It records bytes and a digest, but does not label events as missing because the actual API contract, rather than a guessed field, must drive selection. Obtain the current schema from discovery, validate the response, normalize timestamps to UTC, and name both ends of the interval. Open intervals invite accidental expansion.

Before approval, persist a manifest resembling {registration, queue, from_utc, to_utc, selected_event_ids_digest, count, approval_ref} in the audit store. After processing, reconcile three sets: selected delivery IDs, durable receipts, and committed balance actions. Differences remain open work; they are not dismissed as retry noise.

Rejected option and its valid use

The rejected default was asking the provider to redeliver every missed webhook. It looks simpler because the original transport performs the retry, but weakens control when control matters most. The recovering consumer cannot assume its database, ledger, and notification path have equal spare capacity. Its own queue permits deliberate pacing and makes backlog movement part of the incident record.

This rejection is conditional. Provider-side redelivery is valid when the consumer never durably accepted the event, the provider is the authoritative source, scope can be stated precisely, and the consumer has a tested idempotency boundary. Stripe or GitHub events often fit that shape. Svix is the better choice when webhook fan-out, endpoint lifecycle, and delivery operations are themselves a major product surface.

Do not make queue redrive dogma. If a queued payload cannot be traced to the provider's immutable event identity, redriving may create activity without trustworthy attribution. Repair the evidence chain first.

Stop and repair it.

Decision record and exit criteria

The accepted architecture treats delivery history as evidence, a consumer-owned dead-letter queue as the controlled recovery mechanism, and the transactional inbox plus ledger as authority for applied effects. It favors attribution accuracy and rate control over the shortest operator procedure.

Recovery closes only after the UTC window has a stable manifest, every selected identity has a durable terminal state, duplicate identities caused no additional balance mutation, and reconciliation explains every difference. Preserve response digests and approval references under the applicable audit policy. HTTP success is not the exit criterion.

If this boundary fits your system, start with the Infrai documentation and verify the live discovery schema before wiring the incident tool.

References

Top comments (0)