Skip to content

About

PersonaFlux AI is a persona-driven AI platform that leverages LLMs and NLP to generate structured user personas, enabling personalized experiences, intelligent agents, and context-aware decision systems.

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Latest commit

 

History

2 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 

Repository files navigation

PersonaFlux AI — Real-Time Personalization Engine

Docker FastAPI Redis React

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.


Architecture Diagram

┌─────────────────────────────────────────────────────────────────┐
│                     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           │
                    └─────────────────────────┘

Services Overview

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

Quick Start

Prerequisites

  • Docker 24+
  • Docker Compose v2+

1. Clone & configure

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/.env

2. Start everything

docker-compose up --build

Wait ~60 seconds for all services to initialise.

3. Open the app

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

4. Demo credentials

Email:    demo@personaflux.ai
Password: demo123

JavaScript SDK Usage

Drop-in integration (any website)

<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>

Automatic tracking

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"

Decision-driven UI

<!-- 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>

API Documentation

Authentication

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

Endpoints

POST /auth/register

{
  "name": "Jane Smith",
  "email": "jane@example.com",
  "password": "securepassword"
}

Response: { "token": "...", "user": { "id": "...", "email": "...", "name": "..." } }

POST /auth/login

{ "email": "jane@example.com", "password": "securepassword" }

POST /events/track (or POST /sdk/track — public)

{
  "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

GET /decision/decision/:userId (or GET /sdk/decision/:userId — public)

{
  "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
  }
}

POST /notify/test-email

{
  "to_email": "user@example.com",
  "template": "premium_interest",
  "name": "Jane",
  "discount_percent": "20"
}

Personalization Rules

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

Redis Stream Architecture

Streams

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

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>

Retry Logic

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).


Testing Guide

1. Register and get a token

curl -X POST http://localhost:3000/auth/register \
  -H "Content-Type: application/json" \
  -d '{"name":"Test User","email":"test@test.com","password":"password123"}'

2. Track events (test the pipeline)

# 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":{}}'

3. Poll the decision

# 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"

4. Check the email was triggered

# View notification service logs
docker logs pf-notification-service

5. Send a test email

curl -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"}'

6. Inspect Redis streams

# 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

Free Deployment Guide (Step by Step)

Frontend → Netlify

# 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=dist

Or 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

Backend → Render.com (free tier)

  1. Push your repo to GitHub
  2. Go to render.com → New → Web Service
  3. 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

Redis → Upstash (free tier)

  1. upstash.com → Create Redis database (free)
  2. Copy the REDIS_URL (format: rediss://...)
  3. Set REDIS_URL in all backend services on Render

PostgreSQL → Supabase (free tier)

  1. supabase.com → New project
  2. Go to Settings → Database → Connection string
  3. Run postgres/init/01_schema.sql in the SQL editor
  4. Set DATABASE_URL in all backend services

Cost: $0/month

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

Scaling Strategy

Horizontal Scaling

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 group

Redis Streams with consumer groups handle parallel processing automatically — each message is delivered to exactly one worker.

Redis Stream Tuning

# 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 128

Database Sharding

Event 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');

Caching Strategy

  • 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

Load Balancing

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;
}

Project Structure

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

Environment Variables Reference

API Gateway

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

Decision Engine

Variable Default Description
REDIS_URL redis://redis:6379 Redis connection string
SERVICE_MODE all all | api | worker

Notification Service

Variable Default Description
SMTP_USER `` Gmail address (leave empty for dev mode)
SMTP_PASS `` Gmail app password

System Design Decisions

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

About

PersonaFlux AI is a persona-driven AI platform that leverages LLMs and NLP to generate structured user personas, enabling personalized experiences, intelligent agents, and context-aware decision systems.

Topics

Resources

Stars

1 star

Watchers

0 watching

Forks

Releases

Packages

Contributors