Moving SMS Realtime Off Postgres Polling and Onto Redis Streams
A Django app whose SSE endpoints polled Postgres from sleeping request threads, and the outbox, relay and Redis Streams setup that replaced them in March 2026.
In late March 2026 I was doing client work for a YC-backed insurance startup: a Django backend on Railway, Postgres on Supabase, and an SMS inbox in the web app that updated live over server-sent events. Logins were failing with 503s, and the database and the platform with it went down so often that the team had stopped being surprised by it. I worked through it with Codex on GPT-5.4. It traced code paths, read logs and drafted plans quickly. I decided which plan to build, pushed back when a plan was too small, and reviewed every change before it shipped.
What came out of it, over about a week, was a transactional outbox in Postgres, a relay process, per-user Redis Streams and one SSE connection per browser tab. Along the way there were two separate root causes, a couple of wrong turns on staging, and an audit that found a duplicate-delivery bug in the first version.
The login 503s were a different problem
The logs mixed two failures together, and the first useful step was separating them. The 503s on /api/auth/login/ lined up with a synchronous httpx POST to Supabase Auth that had a 10 second read timeout. When Supabase was slow, the request thread waited out the timeout and returned a 503, with local Postgres never involved.
That got its own fix: a 3 second connect timeout, a 5 second read timeout, one retry with a short backoff, and a small circuit breaker that opens for 60 seconds after three failures, so a slow auth provider makes logins fail fast. That cut the outage noise without explaining the database.
Fifteen threads and three streams per tab
Production ran Gunicorn with three Uvicorn workers and ASGI_THREADS=5. Django views written synchronously run in that thread pool, so the whole service had about 15 threads for anything sync. CONN_MAX_AGE was None, which means each of those threads kept its own Postgres connection open forever.
The SMS realtime endpoints were sync generators. Simplified, the unread-count stream looked like this:
SSE_POLL_INTERVAL = 5
SSE_MAX_DURATION = 120
def _event_stream(self, user_id, agency_id, is_admin, view_mode):
start = time.time()
last_count = None
while time.time() - start < SSE_MAX_DURATION:
count = self._get_unread_count(user_id, agency_id, is_admin, view_mode)
if count != last_count:
yield f"event: count_update\ndata: {json.dumps({'unread_count': count})}\n\n"
last_count = count
time.sleep(SSE_POLL_INTERVAL)Each open stream held a thread (and that thread's database connection) for up to two minutes, ran a few COUNT queries every five seconds, then timed out and the browser reconnected. The access logs showed the unread-count stream reconnecting roughly every two minutes per user, which matched the 120 second limit.
The frontend opened the unread-count stream on every authenticated page for the nav badge. The SMS page added a conversation-list stream and an active-conversation stream. So one SMS tab held three of the 15 threads, five tabs could take all of them, and every one of those streams was polling Postgres. The JWT also went in the SSE query string, because EventSource can't set headers, so access tokens were sitting in the access logs.
I didn't have connection-count metrics from before the change, so this diagnosis came from the arithmetic and the access logs. I can't show a graph of the pool filling up.
Outbox, relay, Redis Streams
The first plan Codex wrote was tactical: turn off two of the streams, set CONN_MAX_AGE=0, add indexes. I asked for something that would stop me ever thinking about this again, and the second plan was an architecture change. Postgres stays the source of truth and stops acting as a per-user event bus.
Every SMS write path (inbound webhook, send, mark read, draft edits, opt-in and opt-out) now writes a row to a realtime_outbox_events table in the same transaction as the change itself:
@transaction.atomic
def mark_message_as_read(user, message_id):
conversation_id = ... # looked up with the same filter
updated = Message.objects.filter(
id=message_id, read_at__isnull=True, conversation__deal__agency_id=user.agency_id,
).update(read_at=Now())
if updated:
queue_sms_sync_event(
agency_id=user.agency_id,
conversation_id=conversation_id,
reason="message.read",
invalidations=["conversations", "messages", "unread_count"],
)The event carries no message content. It tells the frontend which query caches to invalidate, and the frontend refetches through the normal API.
A separate relay process, deployed from the same repo with SERVICE_ROLE=relay, claims unpublished rows and appends them to one Redis Stream per user:
with transaction.atomic():
batch = list(
RealtimeOutboxEvent.objects.select_for_update(skip_locked=True)
.filter(published_at__isnull=True, dead_lettered_at__isnull=True,
available_at__lte=timezone.now())
.order_by("id")[:batch_size]
)
# lease the rows for five minutes so a second relay can't grab them
RealtimeOutboxEvent.objects.filter(id__in=[e.id for e in batch]).update(
available_at=timezone.now() + timedelta(minutes=5)
)Failed publishes back off exponentially and dead-letter after five attempts. The relay still polls Postgres, once a second when idle, but that's one process reading a partial index, where before every open tab ran its own count queries.
The browser opens one stream per tab, through a shared provider, and the nav badge and SMS page both read from it. The stream endpoint does XREAD BLOCK on the user's stream and sends a heartbeat every 20 seconds when nothing arrives. Each event goes out with its stream entry ID as the SSE id, so a reconnecting EventSource sends Last-Event-ID and picks up where it left off.
That resume behaviour is why I used Streams over Pub/Sub. Pub/Sub drops anything published while the client is disconnected, and SSE clients disconnect all the time: tab sleeps, network blips, deploys. Streams keep the last thousand entries per user (MAXLEN ~ 1000), so a reconnect reads what it missed from Redis and never touches Postgres.
For auth, the browser calls /api/realtime/session/, which sets a signed cookie scoped to the stream path that lives for 90 seconds and carries a version number, so logging out bumps the version and kills open streams. CONN_MAX_AGE went to 0 so the Supabase pooler owns connection reuse.
Getting bytes through on staging
The code passed its tests on March 26. On staging, the stream returned 200 with correct headers and then sent nothing.
Several things were wrong at once. The async Redis client had a 5 second socket timeout, which killed every 15 second XREAD BLOCK. Django's global GZipMiddleware was wrapping text/event-stream responses, so I replaced it with a subclass that skips streaming SSE responses. Then I made a wrong turn: still seeing no body, I switched the stream to a sync generator, copying another streaming endpoint in the codebase that worked, and added a padding comment at the start to push past any proxy buffer.
That made it worse in a way that looked like progress. Under ASGI, Django serves a sync streaming iterator by consuming it first, and staging logs had the warning saying so. Headers went out, the body arrived only when the stream ended, and a padding prelude can't fix that. Switching back to an async generator, with the longer Redis timeout and the gzip exemption kept, fixed it. The staging canary opened a stream, queued an outbox row, and got the matching event back through relay, Redis and SSE.
Why two services
On March 29 I asked for another audit, because a relay plus Redis looked odd to me when most people seem to do realtime with Redis Pub/Sub and stop there. The answer was that the relay exists because of the outbox. If you're willing to lose the occasional event, you can XADD from transaction.on_commit and drop the relay. I wasn't.
The audit found a P1. The relay published to each recipient's stream in a loop and then marked the row published. If it crashed between XADD and saving published_at, or failed halfway through the recipients, the retry resent the whole row, and recipients who already had the event got it twice. The frontend deduplicated by outbox ID for five minutes, which hid most of it.
The publish now runs as a Lua script that records a receipt key per recipient and outbox row, and returns the existing entry ID if the receipt is already there:
local existing = redis.call('GET', KEYS[2]) -- receipt for (user, outbox_id)
if existing then
return existing
end
local entry_id = redis.call('XADD', KEYS[1], 'MAXLEN', '~', ARGV[3], '*',
'event', ARGV[1], 'payload', ARGV[2])
redis.call('SET', KEYS[2], entry_id, 'EX', tonumber(ARGV[5]))
return entry_idReceipts are also stored on the outbox row, so a retry skips recipients that already succeeded. The audit also found that the stream endpoint refused connections whenever the relay heartbeat was missing, which turned a relay restart into a user-visible outage even though Redis still had readable data. The stream now stays open and reports delivery health on its ready and heartbeat events, and the frontend falls back to polling while delivery is degraded. Per-user stream keys got a seven day TTL so inactive users stop taking memory. 30 backend tests and 7 frontend tests passed, and the PR merged on March 31.
Railway watchPatterns
For a while every merge to the staging branch had to be deployed twice, because Railway's first attempt reported "No changes to watched files" and skipped the deploy without failing. The checked-in railway.toml had watchPatterns = ["/**"], which looks like it matches everything. The live service config disagreed with the repo, though: the dashboard said the builder was Railpack while the deployments were reading railway.toml and building from the Dockerfile, and the skipped deploys carried skippedReason: "No changes to watched files" in their metadata. Watch patterns only exist to gate deploys conditionally, and a standalone backend repo has no use for that, so I removed the line and left railway.toml as the one place that decides how the service builds and starts. That matters more with two services (web and relay) deploying from the same repo, since a relay that silently stays on an old commit is hard to notice.
The same pattern for Discord
In May 2026 the same app was losing Discord deal notifications. They were posted from fire-and-forget daemon threads, and when Discord returned 429 with a retry_after, the thread logged it and dropped the message. Deploys killed pending posts too. I fixed it the same way as the realtime work. Deal creation writes Discord delivery rows in the same transaction, and a worker claims them with select_for_update(skip_locked=True), honours retry_after, retries with backoff and dead-letters what it can't deliver. The daemon-thread path was deleted.
What I'd measure first next time
I'd get Postgres connection counts and per-endpoint thread occupancy onto a dashboard before changing anything. Access logs and thread arithmetic were enough to act on, but they can't tell me how much the database load dropped. What I can say is that after the rollout the database stopped going down.
I'd also check the ASGI serving path for any streaming endpoint before blaming the proxy. Most of the staging time went to buffering theories, and the warning that explained it had been in the application logs from the first deploy. More bugs from the same habit of trusting the happy path are in war stories from production.