A production-grade, highly scalable HTTP webhook ingestion service written in Go. Designed to process high-throughput telephony call-completion events, guarantee durable exact-once processing semantics (EOPS), and maintain zero-drift per-account call analytics under at-least-once provider delivery.
Important
Production Incident Solved: This repository addresses critical defects including duplicate event storage, aggregate statistics drift, and unhandled background worker goroutine crashes during deployments. All fixes are verified with 100% test suite passing.
- Incident & Defect Resolution Summary
- Key Features
- Architecture & Idempotency Design
- Database Schema & ER Diagram
- Sequence Diagrams
- π System Engineering & Domain Terminology
- Quick Start & Local Setup
- Configuration Reference
- API Reference & Examples
- Test Suite & Verification
- Repository Structure
- 10,000 Webhooks/Sec Scaling Blueprint
- SOLUTION.md Link
| Incident Symptom | Root Cause Term | Architectural Fix & Terminology |
|---|---|---|
| β Duplicate Call Records | Missing B-Tree Unique Constraint on events(event_id). Non-atomic sequential SELECT-then-INSERT. |
Applied Storage-Tier Uniqueness (idx_events_event_id_unique) and UPSERT (INSERT ... ON CONFLICT DO NOTHING). |
| β Account Call-Count Drifting | Unbounded Redelivery Side-Effects & non-transactional statistic updates. | Encapsulated operations in a single ACID Transaction (IngestEventTx). Statistics increment only if the event is newly inserted. |
| β Recordings Not Marked/Processed | Unmonitored Asynchronous Goroutines lacking exception propagation. | Integrated structured error logging (slog) and worker synchronization via Goroutine Lifecycle Tracking. |
| β In-Flight Data Disappearing on Deploy | Abrupt Process Termination without worker drain signal. | Implemented Graceful Worker Drain (Service.Shutdown) using sync.WaitGroup and context timeouts. |
| β Cache Desynchronization on Restart | Cold-Start Empty State in volatility memory without storage hydration. | Executed Cache Hydration / Warming (InitCache()) on boot with Fallback Lookups on cache miss. |
- π‘οΈ Database-Backed Idempotency: Single-transaction processing (
IngestEventTx) using PostgreSQLON CONFLICT DO NOTHINGguarantees Exact-Once Processing Semantics (EOPS) across horizontal app clusters. - β‘ Concurrency & TOCTOU Prevention: Resolves parallel duplicate webhooks atomically at the database engine level, eliminating Time-of-Check to Time-of-Use (TOCTOU) race conditions.
- π Durable Cache & Fallback: Thread-safe in-memory cache populated from PostgreSQL at boot, featuring Optimistic Read-Side CQRS and fallback querying.
- β Graceful Lifecycle Drain: WaitGroup worker drain blocks SIGTERM process shutdown until all in-flight asynchronous operations (e.g. call recording transcodes) complete.
- π‘οΈ Strict Contract Validation: Enforces RFC-compliant HTTP API contracts, rejecting malformed JSON schemas and illegal statuses with standard 400 Bad Request.
Telephony providers guarantee At-Least-Once Delivery. A single event (event_id) can arrive:
- Sequentially: Due to network retry loops after temporary HTTP 5xx errors or socket timeouts.
- Concurrently: Due to multi-threaded worker dispatch on the provider side.
Note
Application-level deduplication (e.g., sync.Mutex or map[string]bool) fails across multi-node deployments or restarts due to lack of shared memory state.
Request A βββ
βββ same event_id βββΊ [ HTTP Ingest ] βββΊ [ IngestEventTx (PostgreSQL) ]
Request B βββ β
ON CONFLICT(event_id) DO NOTHING
β
βββββββββββββββββββ΄ββββββββββββββββββ
βΌ βΌ
Inserted = True Inserted = False (Duplicate)
(Stats +1, Upsert Call) (Rollback Tx, Return 200 OK)
BEGIN;
-- 1. Insert Event (Idempotency Anchor)
INSERT INTO events (event_id, call_id, account_id, payload)
VALUES ($1, $2, $3, $4)
ON CONFLICT (event_id) DO NOTHING;
-- If rows_affected == 0, duplicate detected -> ROLLBACK & EXIT (inserted = false)
-- 2. Upsert Call Record (State Sync)
INSERT INTO calls (call_id, account_id, status, duration_sec, recording_url, updated_at)
VALUES ($1, $2, $3, $4, $5, now())
ON CONFLICT (call_id) DO UPDATE SET ...;
-- 3. Atomically Increment Per-Account Aggregate Stats
INSERT INTO account_stats (account_id, call_count, total_duration_sec)
VALUES ($1, 1, $2)
ON CONFLICT (account_id) DO UPDATE SET
call_count = account_stats.call_count + 1,
total_duration_sec = account_stats.total_duration_sec + EXCLUDED.total_duration_sec;
COMMIT;erDiagram
events {
text event_id PK "UNIQUE INDEX idx_events_event_id_unique"
text call_id FK "Foreign Key to calls.call_id"
text account_id "Partitioning / Tenant Identifier"
jsonb payload "Raw Webhook JSON Delivery"
timestamptz received_at "Ingestion Timestamp"
}
calls {
text call_id PK "Primary Call Entity Key"
text account_id "Tenant Identifier"
text status "Call Status (completed|failed|no_answer)"
integer duration_sec "Call Duration in Seconds"
text recording_url "Media Recording Location"
boolean recording_processed "Async Job Status Flag"
timestamptz updated_at "State Mutation Timestamp"
}
account_stats {
text account_id PK "Tenant Identifier"
bigint call_count "Aggregate Completed Call Count"
bigint total_duration_sec "Aggregate Total Duration Seconds"
}
events }|--|| calls : "references call_id"
calls }|--|| account_stats : "aggregates under account_id"
sequenceDiagram
autonumber
participant Provider as Webhook Provider
participant Router as HTTP Router
participant Ingest as Ingest Service
participant Store as Store (PostgreSQL)
participant Cache as In-Memory Cache
Provider->>Router: POST /webhooks/calls (JSON Payload)
Router->>Router: Contract Validation (JSON & Event Status)
alt Invalid Payload / Status
Router-->>Provider: HTTP 400 Bad Request
else Valid Payload
Router->>Ingest: Ingest(ctx, Event)
Ingest->>Store: IngestEventTx(ctx, Event)
Note over Store: BEGIN TRANSACTION (ACID Boundary)
Store->>Store: INSERT INTO events ... ON CONFLICT DO NOTHING
alt Event Already Exists (Duplicate Delivery)
Note over Store: ROLLBACK TRANSACTION
Store-->>Ingest: inserted = false
Ingest-->>Router: nil
Router-->>Provider: HTTP 200 OK (Idempotent Ack)
else Event Inserted (New Delivery)
Store->>Store: Upsert Call Entity State
Store->>Store: Increment Account Aggregates (+1 count, +duration)
Note over Store: COMMIT TRANSACTION
Store-->>Ingest: inserted = true
Ingest->>Cache: Record(accountID, durationSec)
opt Recording URL Present
Ingest->>Ingest: Spawn Worker Goroutine (wg.Add)
end
Ingest-->>Router: nil
Router-->>Provider: HTTP 200 OK
end
end
| Technical Term | Definition & Context in System |
|---|---|
| At-Least-Once Delivery | Provider delivery semantic where webhooks are retried until an HTTP 2xx acknowledgement is received, potentially delivering duplicate events. |
| Exact-Once Processing Semantics (EOPS) | Guarantee that regardless of duplicate deliveries, state mutations (database inserts & aggregate stats increments) execute exactly once. |
Idempotency Key (event_id) |
Unique payload identifier used by the service to recognize and discard duplicate deliveries safely. |
| Time-of-Check to Time-of-Use (TOCTOU) | Concurrency bug where checking state (SELECT EXISTS) and updating state (INSERT) are separate operations, causing race conditions. Solved via atomic SQL transactions. |
| ACID Transaction Boundary | Database isolation (BEGIN ... COMMIT) wrapping event insertion, call upsert, and stats updates to guarantee atomicity and rollbacks on failure. |
UPSERT (ON CONFLICT DO NOTHING) |
Atomic database primitive combining insertion and collision detection in a single query execution step. |
| CQRS (Command Query Responsibility Segregation) | Separation of high-throughput write paths (PostgreSQL transaction) and low-latency read paths (In-memory stats cache lookups). |
| Cache Hydration / Warming | Bootstrapping the in-memory cache from PostgreSQL on server startup (InitCache()) to eliminate cold-start cache misses. |
| Graceful Worker Drain | Orderly process shutdown mechanism (sync.WaitGroup) waiting for active asynchronous goroutines to complete before SIGTERM process termination. |
| Micro-Batching | Optimization technique aggregating thousands of stream events into multi-row SQL transactions (pgx.Batch / COPY) to minimize Write-Ahead Logging (WAL) I/O overhead. |
| Transaction-Level Pooling (PgBouncer) | Database proxy mechanism multiplexing thousands of short-lived client connections over a bounded pool of persistent PostgreSQL connections. |
- Docker & Docker Compose (Recommended)
- Go 1.22+ (For local testing without Docker)
To launch PostgreSQL, Redis, and the Webhook service:
docker compose up -d --buildTip
Docker Compose automatically handles PostgreSQL and Redis health checks (pg_isready & redis-cli ping) before booting the application container.
curl -i http://localhost:8080/healthz
# Output: HTTP/1.1 200 OK -> okgo test -v ./...Tears down containers, purges data volumes, and re-applies migrations:
make resetConfiguration is managed via environment variables (loaded via internal/config):
| Variable | Default Value | Description |
|---|---|---|
APP_PORT |
8080 |
HTTP Server port |
POSTGRES_PORT |
5432 |
Host port for PostgreSQL container |
REDIS_PORT |
6379 |
Host port for Redis container |
DATABASE_URL |
postgres://webhook:webhook@localhost:5432/webhook?sslmode=disable |
PostgreSQL DSN string |
REDIS_ADDR |
localhost:6379 |
Redis host:port connection string |
DB_MAX_CONNS |
25 |
Maximum PostgreSQL connection pool size |
Ingests telephony call-completion events.
Content-Type: application/json
{
"event_id": "evt_01H8XK2M9P",
"call_id": "call_9f2ab31c",
"account_id": "acc_123",
"status": "completed",
"duration_sec": 143,
"recording_url": "https://recordings.example.com/9f2ab31c.wav",
"occurred_at": "2026-08-13T09:12:00Z"
}200 OK: Successfully processed or duplicate safely ignored.400 Bad Request: Invalid JSON schema, missing required keys, or unknown status (valid:completed,failed,no_answer).500 Internal Server Error: Storage transaction unexpected failure.
Retrieves cumulative statistics for a given account.
curl -s http://localhost:8080/accounts/acc_123/stats{
"call_count": 1,
"total_duration_sec": 143
}Service liveness check endpoint.
HTTP 200 OK β ok
The repository includes a comprehensive test harness in internal/testutil that isolates test data per account, enabling parallel database execution (go test ./...).
internal/
βββ httpapi/
β βββ handler_test.go # HTTP router, payload validation, status code tests
βββ ingest/
β βββ service_test.go # Idempotency, 10x duplicates, concurrency, redelivery tests
βββ stats/
β βββ cache_test.go # In-memory thread safety & stats aggregation tests
βββ store/
βββ store_test.go # PostgreSQL transaction, conflict resolution, stats tests
- β First Webhook Delivery: Verified event stored, call created, stats = 1.
- β Duplicate Delivery (2x): Verified second identical delivery returns 200 OK, event count = 1, stats = 1.
- β Heavy Redelivery (10x): 10 identical POST calls result in exactly 1 record and +1 stats increment.
- β Redelivery After HTTP 200: Late redeliveries after successful ack are safely ignored.
- β
Concurrent Redeliveries: 15 concurrent goroutines posting the same
event_idsimultaneously result in exactly 1 DB record and +1 stats update. - β
Independent Event IDs: Different
event_idpayloads for the same account accumulate stats correctly. - β Transaction Rollback: Invalid SQL steps trigger full rollback without partial updates.
- β
Payload Validation: Rejects malformed JSON and unknown statuses (
unknown_status) with 400 Bad Request.
webhook-ingest/
βββ .github/
β βββ workflows/
β βββ ci.yml # GitHub Actions CI automated build & test pipeline
βββ cmd/
β βββ server/
β βββ main.go # Entrypoint, DB wiring, cache hydration, graceful shutdown
βββ internal/
β βββ config/ # Environment variable parser
β βββ httpapi/ # HTTP routing, payload validation, response helpers
β βββ ingest/ # Core ingestion service & worker pool management
β βββ redisclient/ # Redis connection manager
β βββ stats/ # Thread-safe in-memory account aggregate cache
β βββ store/ # PostgreSQL repository & transactional query logic
β βββ testutil/ # Shared integration test harness & isolated fixtures
βββ migrations/
β βββ 001_init.sql # Initial PostgreSQL tables (events, calls, account_stats)
β βββ 002_unique_event_id.sql # Unique constraint migration on events(event_id)
βββ Dockerfile # Multi-stage optimized Go build
βββ docker-compose.yml # Environment compose configuration
βββ LICENSE # MIT Open Source License
βββ Makefile # Helper developer targets
βββ README.md # System documentation
βββ SOLUTION.md # Detailed solution writeup & 10k/sec scaling blueprint
To scale the architecture from local execution to 10,000 webhooks/sec, the synchronous HTTP-to-PostgreSQL path must be transformed into an asynchronous event-streaming pipeline:
[ 10,000 Webhooks/sec ]
β
βΌ
βββββββββββββββββββββββββββββββββββ
β Edge API Gateway / Ingestion β (Ultra-fast validation & 202 Accepted)
ββββββββββββββββββ¬βββββββββββββββββ
β
βΌ Pushes to Message Stream
βββββββββββββββββββββββββββββββββββ
β Apache Kafka / AWS Kinesis / β (Partitioned by account_id)
β Redis Streams β
ββββββββββββββββββ¬βββββββββββββββββ
β
βΌ Parallel Worker Consumption
βββββββββββββββββββββββββββββββββββ
β Ingest Consumer Worker Pool β
ββββββββββββββββββ¬βββββββββββββββββ
β (Micro-batching pgx.Batch / COPY)
βΌ
βββββββββββββββββββββββββββββββββββ
β PgBouncer (Connection Pooler) β
ββββββββββββββββββ¬βββββββββββββββββ
β
βΌ
βββββββββββββββββββββββββββββββββββ
β PostgreSQL Cluster / Distributedβ (Primary-Replica / Citus)
βββββββββββββββββββββββββββββββββββ
Tip
See SOLUTION.md for full architectural design specs, micro-batching configurations, and Redis write-behind caching strategies.
For full details on defect analysis, deduplication rationale, and evaluation answers, check out SOLUTION.md.