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:

  1. Create — Business logic calls create_notification with a NotificationData contract.
  2. Persist — A NotificationEvent row is inserted into PostgreSQL via products/notifications/backend/models.py.
  3. Publish — The event JSON is published to the notifications Kafka topic.
  4. Consume — The Go consumer reads the message and publishes to Redis channel notifications:{organization_id}.
  5. Stream — The livestream service streams Redis messages to connected browsers over SSE.
  6. Display — The frontend parses JSON into InAppNotification objects and updates the Kea store, triggering toast renders.
  7. Sync — User actions invoke REST endpoints to mark items read, creating NotificationReadState records and invalidating the unread count cache.

Summary

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →