Skip to content

Repository files navigation

Namma Push Engine

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.


⚡ Features

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.

🔄 Full System Flow

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
Loading

🔐 Redis Session Authentication

Namma Push Engine uses a shared Redis session store for authentication — no extra auth service required.

flowchart LR
    A([🖥️ Backend]) -- "SET NS:&lbrace;token&rbrace; &lbrace;driverId&rbrace; EX 86400" --> R[(🗄️ Redis)]
    U([📱 Mobile App]) -- "gRPC connect\ntoken=&lbrace;token&rbrace;" --> E([🚀 Push Engine])
    E -- "GET NS:&lbrace;token&rbrace;" --> R
    R -- "&lbrace;driverId&rbrace; ✅" --> E
    E -- "Stream LIVE" --> U
Loading

Why Redis Auth Is Faster

Method Latency Extra Service
❌ HTTP Auth API 30 – 100ms per connect Requires auth microservice
✅ Redis Session < 1ms per connect None — Redis is already there

Session Lifecycle

stateDiagram-v2
    [*] --> Active : Backend sets NS:&lbrace;token&rbrace;
    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:&lbrace;token&rbrace; (ban/logout)
    Expired --> [*]
    Revoked --> [*]
Loading

Key Format Reference

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

⚙️ Configuration (notification_service.dhall)

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 },
    ...
}

Full Parameter Reference

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

🚀 Running the Engine

Docker Compose (Recommended)

docker-compose up --build -d

This will automatically spin up the required Redis instance and compile/start the Namma Push Engine gateway.

Manual Build

cd crates/notification_service
cargo run --release

🧪 Test with Mock Sender

cd mock-sender
npm install
npm run start

The mock sender automatically:

  1. Injects Redis sessions (NS:{token}) for all configured test clients
  2. Pushes a test notification to every shard of each client's stream (N{clientId}{shard})
  3. Your connected node-client receives it instantly via the live gRPC stream

📊 Metrics & Observability

Prometheus metrics are auto-exported at:

http://<your-server-ip>:9091/metrics

Point Grafana to this endpoint to monitor the following exported Prometheus metrics:

Counters

  • total_notifications: Total number of notifications processed (labeled by category)
  • delivered_notifications: Total successfully delivered to clients over gRPC
  • retried_notifications: Total notifications queued for retry
  • expired_notifications: Total notifications that expired (TTL exceeded)
  • cleanup_push_skipped: Stream cleanup executions that were skipped

Gauges

  • connected_clients: Current number of active gRPC streams (live clients)

Histograms

  • notification_latency: Latency of notification delivery (P50/P95/P99)
  • notification_client_connection_duration: Duration of client connections
  • measure_duration: Core loop and critical section processing duration
  • call_external_api: External FCM / Auth API call latencies
  • termination: Graceful shutdown durations
  • incoming_api: gRPC request processing latencies
  • channel_delay: Internal mpsc channel delays and PubSub message latencies

About

Namma Push Engine is a lightning-fast, ultra-optimized, gRPC-based push-notification service

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages