Made with ❤️ by blacklovertech
Namma Push Engine is a lightning-fast, ultra-optimized, gRPC-based push-notification service written in Rust. It acts as a real-time bridge between your backend services and client apps like driver/rider apps, dashboards, and live tracking views — with zero-HTTP authentication via Redis sessions.
| Feature | Description |
|---|---|
| 🦀 Pure Rust Performance | Built on tokio + tonic. Handles 100,000+ persistent bi-directional gRPC streams on a single server. |
| ⏱️ Ultra-Low Latency | P50 < 10ms · P95 < 25ms · P99 < 50ms at 50,000 TPS |
| 🎯 Dirty-Bit Polling | Epoll-style adaptive polling — only scans streams with unread messages. Near-zero Redis I/O at idle. |
| 🔐 Redis Session Auth | Sub-millisecond auth via shared Redis session store. Zero external HTTP calls on connect. |
| 🔔 FCM Fallback | Automatically fires Firebase Cloud Messaging when a client is offline or the app is killed. |
| ♻️ Zero GC Pauses | No garbage collector — eliminates the tail-latency spikes that plague Node.js and Go at scale. |
sequenceDiagram
autonumber
actor User as 👤 User (Mobile App)
participant Backend as 🖥️ Your Backend
participant Redis as 🗄️ Redis
participant Engine as 🚀 Namma Push Engine
Note over Backend,Redis: Phase 1 — Authentication Setup
Backend->>Redis: SET NS:{token} {driverId} EX 86400
Redis-->>Backend: ✅ OK (Session stored, TTL 24h)
Note over User,Engine: Phase 2 — Client Connection
User->>Engine: gRPC stream open (header: token={token})
Engine->>Redis: GET NS:{token}
Redis-->>Engine: {driverId}
Engine-->>User: ✅ Auth OK — stream is LIVE
Note over Backend,Engine: Phase 3 — Notification Delivery
Backend->>Redis: XADD N{driverId}{shard} * [id, title, body, entity...]
Redis-->>Backend: ✅ Stream entry ID
loop Every poll cycle (dirty-bit tracking)
Engine->>Redis: XREAD N{driverId}{0..N_SHARDS}
Redis-->>Engine: Notification entries
Engine-->>User: 📨 Push via open gRPC stream
end
Note over User,Engine: Phase 4 — Acknowledgement or Fallback
alt User ACKs notification
User->>Engine: ACK (notification_id)
Engine->>Redis: XDEL N{driverId}{shard} {entry_id}
else Client offline / timeout
Engine->>Backend: 🔔 FCM Push via Firebase API
end
Namma Push Engine uses a shared Redis session store for authentication — no extra auth service required.
flowchart LR
A([🖥️ Backend]) -- "SET NS:{token} {driverId} EX 86400" --> R[(🗄️ Redis)]
U([📱 Mobile App]) -- "gRPC connect\ntoken={token}" --> E([🚀 Push Engine])
E -- "GET NS:{token}" --> R
R -- "{driverId} ✅" --> E
E -- "Stream LIVE" --> U
| Method | Latency | Extra Service |
|---|---|---|
| ❌ HTTP Auth API | 30 – 100ms per connect | Requires auth microservice |
| ✅ Redis Session | < 1ms per connect | None — Redis is already there |
stateDiagram-v2
[*] --> Active : Backend sets NS:{token}
Active --> Connected : App connects (GET succeeds)
Connected --> Active : App disconnects (key still valid)
Active --> Expired : TTL reaches 0 (auto-cleanup)
Active --> Revoked : Backend DEL NS:{token} (ban/logout)
Expired --> [*]
Revoked --> [*]
| Redis Key | Value | TTL | Purpose |
|---|---|---|---|
NS:{token} |
{driverId} |
24h (backend-set) | Session auth lookup |
N{driverId}{shard} |
Stream entries | Notification TTL | Delivery stream |
active-notification |
PubSub messages | — | Dirty-bit signal channel |
let redis_cfg = {
host = "0.0.0.0", -- Redis host
port = 30001, -- Redis port
partition = 0, -- DB index (sessions use NS: prefix on same DB)
pool_size = 50, -- Connection pool size
...
}
in {
grpc_port = 50051, -- Mobile apps connect here
max_shards = +5, -- Stream shards per client (parallelism)
redis_cfg = redis_cfg,
fcm_cfg = { api_key = "...", enabled = True },
...
}| Parameter | Type | Default | Description |
|---|---|---|---|
grpc_port |
Integer | 50051 |
gRPC server port (mobile apps connect here) |
http_server_port |
Integer | 9091 |
Health check + Prometheus /metrics endpoint |
redis_cfg |
Object | — | host, port, pool_size, partition, cluster support |
fcm_cfg |
Object | — | Firebase key + enabled flag for offline fallback |
logger_cfg |
Object | INFO |
Log verbosity: TRACE / DEBUG / INFO / WARN / ERROR |
max_shards |
Integer | 5 |
Redis stream shards per client — higher = more parallelism |
channel_buffer |
Integer | 100000 |
Internal mpsc channel capacity for PubSub events |
request_timeout_seconds |
Integer | 60 |
Drop unresponsive clients after N seconds |
retry_delay_millis |
Integer | 1000 |
Re-send interval for unacknowledged notifications |
expired_cleanup_delay_millis |
Integer | 500 |
GC interval to clear TTL-expired stream entries |
read_all_connected_client_notifications |
Boolean | False |
Backfill all active clients on startup |
docker-compose up --build -dThis will automatically spin up the required Redis instance and compile/start the Namma Push Engine gateway.
cd crates/notification_service
cargo run --releasecd mock-sender
npm install
npm run startThe mock sender automatically:
- Injects Redis sessions (
NS:{token}) for all configured test clients - Pushes a test notification to every shard of each client's stream (
N{clientId}{shard}) - Your connected
node-clientreceives it instantly via the live gRPC stream
Prometheus metrics are auto-exported at:
http://<your-server-ip>:9091/metrics
Point Grafana to this endpoint to monitor the following exported Prometheus metrics:
total_notifications: Total number of notifications processed (labeled by category)delivered_notifications: Total successfully delivered to clients over gRPCretried_notifications: Total notifications queued for retryexpired_notifications: Total notifications that expired (TTL exceeded)cleanup_push_skipped: Stream cleanup executions that were skipped
connected_clients: Current number of active gRPC streams (live clients)
notification_latency: Latency of notification delivery (P50/P95/P99)notification_client_connection_duration: Duration of client connectionsmeasure_duration: Core loop and critical section processing durationcall_external_api: External FCM / Auth API call latenciestermination: Graceful shutdown durationsincoming_api: gRPC request processing latencieschannel_delay: Internal mpsc channel delays and PubSub message latencies