A production-grade, microservices-based personalization platform that tracks user behavior in real-time, processes events asynchronously through Redis Streams, and dynamically updates your frontend — all within milliseconds.
┌─────────────────────────────────────────────────────────────────┐
│ Frontend (React + SDK) │
│ Port 80 | Framer Motion | SCSS Modules │
└─────────────────────────────┬───────────────────────────────────┘
│ HTTP / JWT
▼
┌─────────────────────────────────────────────────────────────────┐
│ API Gateway (Node.js) │
│ Port 3000 | JWT Auth | Rate Limiting | Proxy │
└──────────────────────┬───────────────────┬──────────────────────┘
│ │
┌────────────▼──────┐ ┌────────▼──────────────┐
│ Event Service │ │ Decision Engine │
│ (FastAPI :8001) │ │ (FastAPI :8002) │
│ │ │ + Stream Worker │
└────────┬──────────┘ └──────────┬─────────────┘
│ │
│ XADD │ XREADGROUP
▼ ▼
┌────────────────────────────────────────────────┐
│ Redis │
│ │
│ Streams: │
│ • user_events → consumed by Decision Eng │
│ • decision_events → audit trail │
│ • email_stream → consumed by Notif Service │
│ │
│ Cache: │
│ • user_decision:{id} TTL 5min │
│ • user_behavior:{id} TTL 1hr │
└───────────────────────┬────────────────────────┘
│ XREADGROUP
▼
┌─────────────────────────┐
│ Notification Service │
│ (FastAPI :8003) │
│ Sends emails via SMTP │
└─────────────────────────┘
┌─────────────────────────┐
│ PostgreSQL │
│ users | event_logs │
│ decision_logs │
└─────────────────────────┘
| Service | Port | Technology | Purpose |
|---|---|---|---|
| API Gateway | 3000 | Node.js + Express | Auth, rate limiting, reverse proxy |
| Event Service | 8001 | FastAPI | Ingest events → Redis Stream |
| Decision Engine | 8002 | FastAPI + Worker | Rule evaluation + caching |
| Notification Service | 8003 | FastAPI + Worker | Email sending via SMTP |
| Frontend | 80 | React + Vite | SaaS UI |
| Redis | 6379 | Redis 7 | Streams + Cache |
| Redis Commander | 8081 | Web UI | Redis debugging |
| PostgreSQL | 5432 | Postgres 16 | Persistent storage |
- Docker 24+
- Docker Compose v2+
git clone https://github.com/yourname/personaflux.git
cd personaflux
# Copy env files
cp api-gateway/.env.example api-gateway/.env
cp frontend/.env.example frontend/.envdocker-compose up --buildWait ~60 seconds for all services to initialise.
| URL | Description |
|---|---|
| http://localhost | Frontend UI |
| http://localhost:3000/health | API Gateway health |
| http://localhost:8081 | Redis Commander |
| http://localhost:8001/docs | Event Service Swagger |
| http://localhost:8002/docs | Decision Engine Swagger |
| http://localhost:8003/docs | Notification Service Swagger |
Email: demo@personaflux.ai
Password: demo123
<script src="http://localhost/sdk/persona.js"></script>
<script>
Persona.init("user_123", {
apiBase: "http://localhost:3000",
sdkKey: "sdk_personaflux_key_dev",
pollIntervalMs: 2000,
onDecision: (decision) => {
console.log("Decision received:", decision);
},
onDiscount: (percent) => {
console.log(`Show ${percent}% discount!`);
},
});
// Manual event tracking
Persona.track("pricing_view", { source: "header_nav" });
</script>The SDK automatically tracks:
- Page views — fires on load + SPA navigation
- Pricing views — detected by URL containing
/pricing - Clicks — elements with
.premium,.upgrade,data-track="premium" - Idle time — after 30s of no interaction
- Feature hovers — elements with
data-track="feature"
<!-- Show only when discount is active -->
<div data-persona-show="discount">
🎉 You have an exclusive discount!
</div>
<!-- Highlighted automatically by SDK -->
<button class="btn-premium" data-track="premium">
Upgrade to Pro
</button>All protected endpoints require a JWT token:
Authorization: Bearer <token>
SDK calls use a lightweight API key:
X-SDK-Key: sdk_personaflux_key_dev
X-User-Id: user_123
{
"name": "Jane Smith",
"email": "jane@example.com",
"password": "securepassword"
}Response: { "token": "...", "user": { "id": "...", "email": "...", "name": "..." } }
{ "email": "jane@example.com", "password": "securepassword" }{
"user_id": "user_123",
"event_type": "pricing_view",
"metadata": { "source": "navbar" },
"timestamp": 1700000000000,
"page_url": "https://yoursite.com/pricing"
}Allowed event_types: page_view, pricing_view, premium_click, feature_hover, cta_click, idle, click, scroll, form_submit, exit_intent
{
"source": "cache",
"decision": {
"user_id": "user_123",
"show_discount": true,
"discount_percent": 20,
"show_idle_popup": false,
"trigger_email": true,
"email_template": "premium_interest",
"show_premium_cta": true,
"highlight_features": ["analytics"],
"personalization_score": 0.65,
"rules_triggered": ["pricing_engagement_discount", "premium_click_trigger"],
"evaluated_at": 1700000000000
}
}{
"to_email": "user@example.com",
"template": "premium_interest",
"name": "Jane",
"discount_percent": "20"
}| Rule | Trigger | Action |
|---|---|---|
pricing_engagement_discount |
Pricing page viewed ≥ 2× | Show 15% discount |
premium_click_trigger |
Premium button clicked ≥ 1× | Trigger email + highlight CTA |
discount_boost_premium |
Discount + premium click | Boost discount to 20% |
idle_popup |
No activity for 30+ seconds | Show retention popup |
feature_highlight_analytics |
Total page views ≥ 3 | Highlight analytics feature |
feature_highlight_integrations |
Feature hovers ≥ 2 | Highlight integrations feature |
user_events
└── Published by: Event Service (on every /track call)
└── Consumed by: Decision Engine Worker (consumer group: decision-engine-group)
└── Format: { user_id, event_type, metadata, timestamp, page_url, session_id }
decision_events
└── Published by: Decision Engine Worker (after every evaluation)
└── Consumed by: (future: analytics, A/B testing)
└── Format: { user_id, decision: JSON, ts }
email_stream
└── Published by: Decision Engine Worker (when trigger_email = true)
└── Consumed by: Notification Service Worker (consumer group: notification-service-group)
└── Format: { user_id, template, discount_percent, ts }
Consumer groups allow:
- Multiple workers to consume the same stream without duplicate processing
- Pending message tracking — messages that were delivered but not ACKed
- Retry on crash — unacknowledged messages are re-delivered
XGROUP CREATE user_events decision-engine-group $ MKSTREAM
XREADGROUP GROUP decision-engine-group worker-1 COUNT 20 BLOCK 2000 STREAMS user_events >
XACK user_events decision-engine-group <msg-id>
If a message fails processing, it stays in the Pending Entries List (PEL). On startup, the worker calls XPENDING_RANGE to reprocess any stuck messages. After MAX_RETRIES, the message is acknowledged to prevent infinite loops (in production, move to a dead-letter queue).
curl -X POST http://localhost:3000/auth/register \
-H "Content-Type: application/json" \
-d '{"name":"Test User","email":"test@test.com","password":"password123"}'# Track 2 pricing views → triggers discount rule
for i in 1 2; do
curl -X POST http://localhost:3000/sdk/track \
-H "Content-Type: application/json" \
-H "X-SDK-Key: sdk_personaflux_key_dev" \
-H "X-User-Id: test_user_1" \
-d '{"user_id":"test_user_1","event_type":"pricing_view","metadata":{}}'
done
# Track premium click → triggers email rule
curl -X POST http://localhost:3000/sdk/track \
-H "Content-Type: application/json" \
-H "X-SDK-Key: sdk_personaflux_key_dev" \
-H "X-User-Id: test_user_1" \
-d '{"user_id":"test_user_1","event_type":"premium_click","metadata":{}}'# Wait 1-2 seconds for the worker to process
sleep 2
curl http://localhost:3000/sdk/decision/test_user_1 \
-H "X-SDK-Key: sdk_personaflux_key_dev"# View notification service logs
docker logs pf-notification-servicecurl -X POST http://localhost:8003/test-email \
-H "Content-Type: application/json" \
-d '{"to_email":"you@example.com","template":"premium_interest","name":"Test","discount_percent":"20"}'# Using Redis CLI
docker exec -it pf-redis redis-cli
# View stream length
XLEN user_events
XLEN email_stream
XLEN decision_events
# View messages
XRANGE user_events - + COUNT 5
# View consumer group info
XINFO GROUPS user_events
# View pending messages
XPENDING user_events decision-engine-group - + 10# 1. Build the frontend
cd frontend
cp .env.example .env
# Edit .env: set VITE_API_URL=https://your-render-backend.onrender.com
npm install
npm run build
# 2. Deploy to Netlify
npm install -g netlify-cli
netlify login
netlify deploy --prod --dir=distOr push to GitHub and connect at netlify.com. Set:
- Build command:
cd frontend && npm run build - Publish directory:
frontend/dist - Environment variable:
VITE_API_URL=https://your-backend.onrender.com
Add frontend/_redirects:
/* /index.html 200
- Push your repo to GitHub
- Go to render.com → New → Web Service
- Create one service per microservice:
API Gateway:
- Root directory:
api-gateway - Build:
npm ci - Start:
node server.js - Environment: copy all vars from
.env.example
Event Service / Decision Engine / Notification Service:
- Runtime: Python 3.11
- Build:
pip install -r requirements.txt - Start:
uvicorn main:app --host 0.0.0.0 --port $PORT
- upstash.com → Create Redis database (free)
- Copy the
REDIS_URL(format:rediss://...) - Set
REDIS_URLin all backend services on Render
- supabase.com → New project
- Go to Settings → Database → Connection string
- Run
postgres/init/01_schema.sqlin the SQL editor - Set
DATABASE_URLin all backend services
| Service | Provider | Free Tier |
|---|---|---|
| Frontend | Netlify | 100GB bandwidth |
| API Gateway | Render | 750 hrs/month |
| Event Service | Render | 750 hrs/month |
| Decision Engine | Render | 750 hrs/month |
| Redis | Upstash | 10,000 req/day |
| PostgreSQL | Supabase | 500MB storage |
Each microservice is stateless and can be scaled independently:
# Scale decision engine workers
decision-engine:
deploy:
replicas: 4 # 4 workers consuming from the same consumer groupRedis Streams with consumer groups handle parallel processing automatically — each message is delivered to exactly one worker.
# Trim stream to prevent unbounded growth
XTRIM user_events MAXLEN ~ 100000
# Set in docker-compose via redis.conf:
# stream-node-max-bytes 4096
# stream-node-max-entries 128Event logs can be partitioned by occurred_at using PostgreSQL table partitioning:
PARTITION BY RANGE (occurred_at);
CREATE TABLE event_logs_2025_01 PARTITION OF event_logs
FOR VALUES FROM ('2025-01-01') TO ('2025-02-01');- Decision cache TTL: 5 minutes (configurable) — prevents re-evaluating on every poll
- Behavior cache TTL: 1 hour — resets session after inactivity
- Rate limit keys TTL: 61 seconds — sliding window
In production, add nginx upstream to round-robin across API Gateway instances:
upstream gateway {
least_conn;
server gateway-1:3000;
server gateway-2:3000;
server gateway-3:3000;
}personaflux/
├── api-gateway/
│ ├── server.js # Express app + proxy setup
│ ├── middleware/
│ │ ├── auth.js # JWT verification
│ │ └── rateLimit.js # Redis sliding window rate limiter
│ ├── routes/
│ │ ├── auth.js # Register / login
│ │ ├── events.js # Proxy → event-service
│ │ ├── decision.js # Proxy → decision-engine
│ │ └── notify.js # Proxy → notification-service
│ ├── Dockerfile
│ └── package.json
│
├── event-service/
│ ├── main.py # FastAPI app + /track endpoint
│ ├── Dockerfile
│ └── requirements.txt
│
├── decision-engine/
│ ├── main.py # FastAPI app + /decision/:id endpoint
│ ├── worker.py # Redis Stream consumer + rule evaluation
│ ├── rules.py # RuleEngine + UserBehavior + Decision models
│ ├── entrypoint.py # Starts API + worker in same container
│ ├── Dockerfile
│ └── requirements.txt
│
├── notification-service/
│ ├── main.py # FastAPI app + email worker
│ ├── templates.py # Email templates
│ ├── Dockerfile
│ └── requirements.txt
│
├── frontend/
│ ├── src/
│ │ ├── pages/
│ │ │ ├── Home.tsx # Landing page
│ │ │ ├── Pricing.tsx # Pricing with live discount
│ │ │ ├── Demo.tsx # Interactive demo sandbox
│ │ │ └── Auth.tsx # Login + Register
│ │ ├── components/
│ │ │ ├── Navbar.tsx # Animated navbar
│ │ │ └── PipelineStatus.tsx # Real-time pipeline visualizer
│ │ ├── hooks/
│ │ │ └── useDecision.ts # Polling hook
│ │ └── utils/
│ │ ├── api.ts # Axios instance + interceptors
│ │ └── AuthContext.tsx
│ ├── Dockerfile
│ └── package.json
│
├── sdk/
│ └── persona.js # Drop-in browser SDK
│
├── postgres/
│ └── init/
│ └── 01_schema.sql # Database schema + seed
│
└── docker-compose.yml
| Variable | Default | Description |
|---|---|---|
JWT_SECRET |
...dev_secret |
JWT signing key |
SDK_API_KEY |
sdk_personaflux_key_dev |
SDK authentication key |
REDIS_HOST |
redis |
Redis hostname |
RATE_LIMIT_MAX |
200 |
Requests per minute per IP |
EVENT_SERVICE_URL |
http://event-service:8001 |
Event service base URL |
DECISION_ENGINE_URL |
http://decision-engine:8002 |
Decision engine base URL |
| Variable | Default | Description |
|---|---|---|
REDIS_URL |
redis://redis:6379 |
Redis connection string |
SERVICE_MODE |
all |
all | api | worker |
| Variable | Default | Description |
|---|---|---|
SMTP_USER |
`` | Gmail address (leave empty for dev mode) |
SMTP_PASS |
`` | Gmail app password |
Why Redis Streams over Kafka? Redis Streams provide Kafka-like semantics (consumer groups, offsets, replay) without the operational overhead. For sub-100k events/day, Redis is simpler, faster, and free.
Why rule-based over ML? Rules are deterministic, debuggable, and latency-free. The same Redis behavior hash can later feed a lightweight ML model without architecture changes.
Why polling over WebSockets? Polling every 2s is sufficient for personalization (sub-second latency isn't required). It avoids persistent connection management and works across all deployment platforms including serverless.
Why a separate Decision cache? Decouples the worker (async) from the API (sync). The frontend gets fast responses even if the worker is briefly behind. TTL prevents stale decisions.
Built with ❤️ using FastAPI, Redis Streams, React, and Node.js