How to Implement Real-Time In-App Notifications in PostHog
PostHog's notification system combines Django backend events, Kafka messaging, Redis pub/sub, and Server-Sent Events (SSE) to deliver instant in-app alerts to users without polling.
PostHog's open-source platform includes a production-ready real-time notification system that enables instant in-app alerts through a distributed event pipeline. Understanding how to implement real-time in-app notifications within this architecture allows you to extend the platform with custom alert types while leveraging existing Kafka, Redis, and SSE infrastructure. This guide walks through the complete implementation using actual source code from the PostHog/posthog repository.
Architectural Overview
The real-time notification system spans five distinct layers working in concert to deliver sub-second alerts.
Notification Creation — Business logic in products/notifications/backend/logic.py creates NotificationEvent records and publishes them to Kafka via the create_notification function.
Message Transport — A Go-based consumer in livestream/events/notification_consumer.go reads Kafka events and republishes them to organization-specific Redis channels using the NotificationKafkaConsumer.processMessage method.
Realtime Delivery — The frontend opens an SSE connection to /api/livestream/notifications/, handled by the livestream service, which streams messages from Redis directly to the browser.
UI State Management — The React frontend uses Kea logic in frontend/src/lib/components/NotificationsMenu/notificationsMenuLogic.tsx to manage the notification list, unread counts, and dropdown state.
API Surface — REST endpoints in products/notifications/backend/presentation/views.py expose NotificationsViewSet for listing notifications, retrieving unread counts, and marking items read or unread.
Caching Layer — Unread counts are cached per user and organization in products/notifications/backend/cache.py to avoid expensive database queries on every poll.
Backend Implementation
Creating Notification Events
To trigger a notification, services call the facade API which validates data and persists events. In products/notifications/backend/facade/api.py, the create_notification function accepts a NotificationData contract and handles persistence:
# products/notifications/backend/facade/api.py
from products.notifications.backend.facade.contracts import NotificationData
from products.notifications.backend.logic import create_notification
def send_user_mention_notification(user_id: int, org_id: int, text: str) -> None:
data = NotificationData(
organization_id=org_id,
notification_type="user_mention",
title="You were mentioned",
body=text,
target_type="user",
target_id=str(user_id),
resource_type=None,
resource_id=None,
)
create_notification(data)
According to the PostHog source code, the create_notification function in products/notifications/backend/logic.py performs three critical operations: it saves a NotificationEvent record to the database (defined in products/notifications/backend/models.py), builds a JSON payload, and calls _publish_to_kafka to push the event to the notifications Kafka topic.
Processing Events with the Kafka Consumer
The livestream service runs a Go consumer that bridges Kafka and Redis. In livestream/events/notification_consumer.go, the NotificationKafkaConsumer.processMessage method extracts the organization ID and publishes to a scoped Redis channel:
func (c *NotificationKafkaConsumer) processMessage(ctx context.Context, value []byte) {
var data struct {
OrganizationID string `json:"organization_id"`
}
if err := json.Unmarshal(value, &data); err != nil {
// handle error
}
channel := fmt.Sprintf("notifications:%s", data.OrganizationID)
cmd := c.redisClient.B().Spublish().Channel(channel).Message(string(value)).Build()
_ = c.redisClient.Do(ctx, cmd).Error()
}
This pattern ensures that only users belonging to the specific organization receive the notification, maintaining strict data isolation while enabling horizontal scaling of the consumer pool.
REST API and Caching Layer
The backend exposes REST endpoints through NotificationsViewSet in products/notifications/backend/presentation/views.py. The unread_count action checks a cached value before querying the database:
# Conceptual usage of the cache layer
from products.notifications.backend.cache import get_unread_count, set_unread_count
def get_cached_unread(user_id, org_id):
count = get_unread_count(user_id, org_id)
if count is None:
count = calculate_unread_from_db(user_id, org_id)
set_unread_count(user_id, org_id, count, ttl=60)
return count
When users mark notifications as read via the mark_read endpoint, the backend creates a NotificationReadState record and calls invalidate_unread_count to bust the cache immediately.
Frontend Implementation
Establishing the SSE Connection
The browser connects to the livestream service using the connectToNotificationsSSE function exported from frontend/src/layout/navigation-3000/sidepanel/panels/activity/notificationsSSE.ts. This utility manages the EventSource lifecycle, authentication headers, and automatic reconnection with exponential backoff:
import { connectToNotificationsSSE } from 'layout/navigation-3000/sidepanel/panels/activity/notificationsSSE'
const abort = new AbortController()
connectToNotificationsSSE(
'/api/livestream/notifications/',
user.apiToken,
abort.signal,
(notification) => {
// notification conforms to InAppNotification type
notificationsMenuLogic.actions.addNotification(notification)
},
{
onFirstMessage: () => console.log('SSE connected'),
onError: (err) => console.error('SSE error', err),
}
)
The SSE connection streams JSON payloads that match the InAppNotification TypeScript interface defined in the frontend types system.
State Management with Kea
Notification state is managed via Kea logic in frontend/src/lib/components/NotificationsMenu/notificationsMenuLogic.tsx. This logic handles the notification list, active tabs, and unread counts:
// Conceptual usage within a React component
import { useActions, useValues } from 'kea'
import { notificationsMenuLogic } from 'lib/components/NotificationsMenu/notificationsMenuLogic'
function NotificationDropdown() {
const { notifications, unreadCount } = useValues(notificationsMenuLogic)
const { addNotification, markRead } = useActions(notificationsMenuLogic)
// Connect to SSE on mount to populate notifications
}
The logic also coordinates with the API to persist read states, calling endpoints like api/environments/{team_id}/notifications/{id}/mark_read/ when users interact with specific items.
Rendering Toast Notifications
Real-time toasts are rendered by notificationToasts.tsx in the same directory as the menu logic. When the SSE client receives a new message, the Kea store updates trigger React re-renders, displaying ephemeral toast alerts while simultaneously updating the persistent notification bell count.
Data Flow Walkthrough
Understanding the complete path of a notification ensures you can debug or extend the system effectively:
- Create — Business logic calls
create_notificationwith aNotificationDatacontract. - Persist — A
NotificationEventrow is inserted into PostgreSQL viaproducts/notifications/backend/models.py. - Publish — The event JSON is published to the
notificationsKafka topic. - Consume — The Go consumer reads the message and publishes to Redis channel
notifications:{organization_id}. - Stream — The livestream service streams Redis messages to connected browsers over SSE.
- Display — The frontend parses JSON into
InAppNotificationobjects and updates the Kea store, triggering toast renders. - Sync — User actions invoke REST endpoints to mark items read, creating
NotificationReadStaterecords and invalidating the unread count cache.
Summary
- Backend creation happens through
create_notificationinproducts/notifications/backend/logic.py, which persists to PostgreSQL and publishes to Kafka. - Real-time relay is handled by
NotificationKafkaConsumerinlivestream/events/notification_consumer.go, pushing messages to organization-scoped Redis channels. - Frontend streaming uses
connectToNotificationsSSEto establish Server-Sent Events connections that feed data into Kea logic. - State persistence relies on
NotificationsViewSetinproducts/notifications/backend/presentation/views.pyfor CRUD operations and read-state management. - Performance optimization comes from
products/notifications/backend/cache.py, which caches unread counts with a 60-second TTL to reduce database load.
Frequently Asked Questions
How does PostHog handle SSE reconnection errors?
The connectToNotificationsSSE function in frontend/src/layout/navigation-3000/sidepanel/panels/activity/notificationsSSE.ts implements automatic reconnection using a retryWithBackoff mechanism from the generic livestream client. If the EventSource connection drops due to network errors or server restarts, the client waits an exponentially increasing interval before attempting to reconnect, ensuring users resume receiving real-time notifications without manual page refreshes.
What is the role of Redis in the notification pipeline?
Redis serves as the low-latency pub/sub broker between the Go consumer and the livestream service. According to the source code in livestream/events/notification_consumer.go, the NotificationKafkaConsumer publishes messages to Redis channels named notifications:{organization_id}, enabling the livestream service to subscribe to organization-specific streams and push updates only to relevant users via SSE.
How are unread counts cached and invalidated?
Unread counts are cached per user and organization using Django's cache framework, as implemented in products/notifications/backend/cache.py with a default UNREAD_COUNT_TTL_SECONDS of 60 seconds. When a user marks a notification as read or unread through the mark_read or mark_unread endpoints in views.py, the backend calls invalidate_unread_count to immediately clear the cached value, ensuring the next API request returns an accurate count fresh from the database.
Can I extend this system to support push notifications?
Yes, the architecture supports extension by modifying the create_notification function in products/notifications/backend/logic.py to publish to additional Kafka topics or integrating external services alongside the existing Redis publication. You would add new consumers (similar to notification_consumer.go) to process these events for push delivery services like Firebase Cloud Messaging or OneSignal, while the existing SSE pipeline continues to handle in-app alerts without modification.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →