Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions backend/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
from routes.submit_room_line import submit_room_line_bp
from routes.admin import admin_bp
from routes.frontend import frontend_bp
from routes.analytics import analytics_bp
from services.db import redis_client
from services.canvas_counter import get_canvas_draw_count
from services.graphql_service import commit_transaction_via_graphql
Expand Down Expand Up @@ -143,6 +144,7 @@ def handle_all_exceptions(e):

# Frontend serving must be last to avoid route conflicts
app.register_blueprint(frontend_bp)
app.register_blueprint(analytics_bp)

if __name__ == '__main__':
if not redis_client.exists('res-canvas-draw-count'):
Expand Down
6 changes: 6 additions & 0 deletions backend/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,12 @@

LOG_FILE = "backend_graphql.log"

# Analytics / LLM configuration
ANALYTICS_ENABLED = os.getenv("ANALYTICS_ENABLED", "True") == "True"
OPENAI_API_KEY = os.getenv("OPENAI_API_KEY")
ANALYTICS_COLLECTION_NAME = os.getenv("ANALYTICS_COLLECTION_NAME", "analytics_events")
ANALYTICS_AGGREGATES_COLLECTION = os.getenv("ANALYTICS_AGGREGATES_COLLECTION", "analytics_aggregates")

JWT_SECRET = os.getenv("JWT_SECRET", "dev-insecure-change-me")
JWT_ISSUER = "rescanvas"

Expand Down
49 changes: 49 additions & 0 deletions backend/routes/analytics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
from flask import Blueprint, request, jsonify
from services.db import analytics_aggregates_coll
from services.analytics_service import query_recent
from services.insights_generator import generate_insights
import logging

analytics_bp = Blueprint('analytics_bp', __name__)
logger = logging.getLogger(__name__)


@analytics_bp.route('/api/analytics/recent', methods=['GET'])
def recent_events():
room = request.args.get('roomId')
limit = int(request.args.get('limit', '100'))
docs = query_recent(room, limit=limit)
# convert ObjectId to str if needed
def clean(d):
d = dict(d)
d.pop('_id', None)
return d
return jsonify([clean(d) for d in docs])


@analytics_bp.route('/api/analytics/overview', methods=['GET'])
def overview():
room = request.args.get('roomId')
q = {} if not room else {"roomId": str(room)}
ag = analytics_aggregates_coll.find_one(q) or {}
ag.pop('_id', None)
return jsonify(ag)


@analytics_bp.route('/api/analytics/insights', methods=['POST'])
def insights():
data = request.get_json() or {}
room = data.get('roomId')
q = {} if not room else {"roomId": str(room)}
ag = analytics_aggregates_coll.find_one(q) or {}
try:
res = generate_insights(ag)
return jsonify(res)
except Exception:
logger.exception('Failed to generate insights')
return jsonify({"error": "failed"}), 500


@analytics_bp.route('/api/analytics/health', methods=['GET'])
def health():
return jsonify({"ok": True})
20 changes: 20 additions & 0 deletions backend/routes/socketio_handlers.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from flask_socketio import join_room, leave_room, emit
from services.socketio import socketio
from services.db import rooms_coll, shares_coll, users_coll
from services.analytics_service import ingest_event
import logging
from config import JWT_SECRET
import jwt
Expand Down Expand Up @@ -122,6 +123,16 @@ def on_join_room(data):
payload = {'roomId': room_id, 'userId': user_id, 'username': username_to_emit, 'members': members}
logging.getLogger(__name__).info('socket: emitting user_joined to room %s payload=%s sid=%s had_cached_claims=%s', room_id, payload, sid, had_cached)
emit('user_joined', payload, room=f"room:{room_id}")
try:
# record join event for analytics (anonymized inside service)
ingest_event({
'roomId': room_id,
'userId': user_id,
'eventType': 'join',
'payload': {'username': username_to_emit}
})
except Exception:
pass
try:
emit('server_debug', {'action': 'emitted_user_joined', 'sid': sid, 'had_cached_claims': had_cached, 'payload': payload}, room=None)
except Exception:
Expand Down Expand Up @@ -150,5 +161,14 @@ def on_leave_room(data):
payload = {'roomId': room_id, 'username': username_to_emit, 'members': members}
logging.getLogger(__name__).info('socket: emitting user_left to room %s payload=%s', room_id, payload)
emit('user_left', payload, room=f"room:{room_id}")
try:
ingest_event({
'roomId': room_id,
'userId': None, # leave events can be anonymous
'eventType': 'leave',
'payload': {'username': username_to_emit}
})
except Exception:
pass
except Exception:
pass
29 changes: 29 additions & 0 deletions backend/routes/submit_room_line.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
from services.graphql_service import commit_transaction_via_graphql
from services.db import redis_client, strokes_coll, rooms_coll, shares_coll
from services.socketio_service import push_to_room
from services.analytics_service import ingest_event
from services.canvas_counter import get_canvas_draw_count, increment_canvas_draw_count
from services.crypto_service import unwrap_room_key, encrypt_for_room, wrap_room_key
import nacl.signing, nacl.encoding
Expand Down Expand Up @@ -157,6 +158,20 @@ def submit_room_line():
'type': room_type
})

try:
ingest_event({
'roomId': roomId,
'userId': actor_id or drawing.get('user'),
'eventType': 'stroke_created',
'payload': {
'color': drawing.get('color'),
'lineWidth': drawing.get('lineWidth')
},
'ts': drawing['timestamp']
})
except Exception:
pass

# Update room's updatedAt so the Dashboard's "Last edited" reflects drawing activity
try:
rooms_coll.update_one({'_id': room['_id']}, {'$set': {'updatedAt': datetime.utcnow()}})
Expand All @@ -180,6 +195,20 @@ def submit_room_line():
'type': 'public'
})

try:
ingest_event({
'roomId': roomId,
'userId': actor_id or drawing.get('user'),
'eventType': 'stroke_created',
'payload': {
'color': drawing.get('color'),
'lineWidth': drawing.get('lineWidth')
},
'ts': drawing['timestamp']
})
except Exception:
pass

try:
rooms_coll.update_one({'_id': room['_id']}, {'$set': {'updatedAt': datetime.utcnow()}})
except Exception:
Expand Down
65 changes: 65 additions & 0 deletions backend/services/analytics_service.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
"""
Simple analytics ingestion and query helpers.
This module stores anonymized events into MongoDB and provides small helpers
used by socket handlers and the aggregation worker.
"""
import time
import logging
from services.db import analytics_coll
from bson.objectid import ObjectId
from config import ANALYTICS_ENABLED

logger = logging.getLogger(__name__)


def _anonymize_user(user_id):
"""Return a deterministic anonymized user id for privacy-aware aggregation."""
if not user_id:
return None
# Keep short hash-like form
try:
import hashlib
return hashlib.sha1(str(user_id).encode('utf-8')).hexdigest()[:16]
except Exception:
return str(user_id)


def ingest_event(event: dict):
"""Ingest an analytics event into the analytics collection.

Event shape (examples):
{
"roomId": "...",
"userId": "...",
"eventType": "stroke_created" | "join" | "leave" | "heartbeat",
"payload": { ... },
"ts": 1234567890
}
"""
if not ANALYTICS_ENABLED:
return None
try:
ev = dict(event)
ev.setdefault('ts', int(time.time() * 1000))
# anonymize userId for privacy
if 'userId' in ev and ev['userId']:
ev['anonUserId'] = _anonymize_user(ev.get('userId'))
ev.pop('userId', None)
# ensure roomId string
if 'roomId' in ev and isinstance(ev['roomId'], ObjectId):
ev['roomId'] = str(ev['roomId'])
analytics_coll.insert_one(ev)
return ev
except Exception:
logger.exception('Failed to ingest analytics event')
return None


def query_recent(roomId=None, limit=100):
q = {} if not roomId else {"roomId": str(roomId)}
try:
docs = list(analytics_coll.find(q).sort([('ts', -1)]).limit(limit))
return docs
except Exception:
logger.exception('Failed to query recent analytics events')
return []
18 changes: 17 additions & 1 deletion backend/services/db.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,16 @@
invites_coll = mongo_client[DB_NAME]["room_invites"]
notifications_coll = mongo_client[DB_NAME]["notifications"]

# Analytics collections
try:
from config import ANALYTICS_COLLECTION_NAME, ANALYTICS_AGGREGATES_COLLECTION
except Exception:
ANALYTICS_COLLECTION_NAME = "analytics_events"
ANALYTICS_AGGREGATES_COLLECTION = "analytics_aggregates"

analytics_coll = mongo_client[DB_NAME][ANALYTICS_COLLECTION_NAME]
analytics_aggregates_coll = mongo_client[DB_NAME][ANALYTICS_AGGREGATES_COLLECTION]

# TTL index on refresh token expiresAt so expired refresh tokens are removed automatically
try:
refresh_tokens_coll.create_index("expiresAt", expireAfterSeconds=0)
Expand All @@ -61,4 +71,10 @@
users_coll.create_index("username", unique=True)
rooms_coll.create_index([("ownerId", 1), ("type", 1)])
shares_coll.create_index([("roomId", 1), ("userId", 1)], unique=True)
strokes_coll.create_index([("roomId", 1), ("ts", 1)])
strokes_coll.create_index([("roomId", 1), ("ts", 1)])
try:
analytics_coll.create_index([("roomId", 1), ("ts", 1)])
analytics_coll.create_index([("eventType", 1)])
analytics_aggregates_coll.create_index([("roomId", 1)])
except Exception:
pass
74 changes: 74 additions & 0 deletions backend/services/insights_generator.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
"""
LLM-based insights generator. Uses OpenAI if API key is available; otherwise
returns simple rule-based summaries.
"""
import logging
from config import OPENAI_API_KEY
import json

logger = logging.getLogger(__name__)


def _summarize_aggregates(aggregates: dict):
# Build a concise prompt from aggregates
try:
total_strokes = aggregates.get('total_strokes', 0)
active_users = aggregates.get('active_users', 0)
top_colors = aggregates.get('top_colors', [])
collab_pairs = aggregates.get('collaboration_pairs', [])

summary = (
f"Total strokes: {total_strokes}. Active users: {active_users}. "
f"Top colors: {', '.join(top_colors[:5]) if top_colors else 'N/A'}. "
f"Top collaboration pairs: {', '.join([f'{p[0]}-{p[1]}' for p in collab_pairs[:5]]) if collab_pairs else 'N/A'}."
)
return summary
except Exception:
return "No summary available"


def generate_insights(aggregates: dict):
"""Generate human-readable insights and recommendations.

If OpenAI API key is present, attempt a short chat completion. If not,
return a simple deterministic summary and a few heuristic recommendations.
"""
try:
prompt_summary = _summarize_aggregates(aggregates)
if OPENAI_API_KEY:
try:
import openai
openai.api_key = OPENAI_API_KEY
system = "You are an analytics assistant for a collaborative drawing app. Provide a short summary and 3 actionable recommendations to improve collaboration and room health."
response = openai.ChatCompletion.create(
model="gpt-4o-mini",
messages=[
{"role": "system", "content": system},
{"role": "user", "content": f"Here are aggregates: {json.dumps(aggregates)}. Produce a short summary and 3 actionable recommendations."}
],
max_tokens=400,
temperature=0.6,
)
text = response.choices[0].message.content
return {"summary": text, "source": "openai"}
except Exception:
logger.exception('OpenAI call failed, falling back to heuristic summary')

# Fallback heuristic
summary = _summarize_aggregates(aggregates)
recommendations = []
if aggregates.get('active_users', 0) < 2:
recommendations.append('Encourage users to invite collaborators or schedule group drawing sessions to improve engagement.')
if aggregates.get('avg_stroke_rate', 0) < 1:
recommendations.append('Introduce prompts or templates to encourage drawing activity.')
if aggregates.get('anomaly_score', 0) > 0.7:
recommendations.append('Investigate sudden surges in activity — could be bot behavior or an event-driven spike.')
if not recommendations:
recommendations = [
'Highlight active users in the room to promote collaboration.',
'Surface weekly summary emails with top contributors and trending palettes.'
]
return {"summary": summary, "recommendations": recommendations, "source": "heuristic"}
except Exception:
logger.exception('generate_insights error')
return {"summary": "Unable to generate insights.", "recommendations": []}
53 changes: 53 additions & 0 deletions backend/tests/test_analytics_service.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
import sys
import types
import time


class FakeColl:
def __init__(self):
self.storage = []

def insert_one(self, doc):
# emulate pymongo returning an InsertOneResult-like object
self.storage.append(dict(doc))
class R: pass
r = R()
r.inserted_id = len(self.storage) - 1
return r

def find(self, q=None):
# return all docs matching roomId if provided
q = q or {}
room = q.get('roomId')
if room:
return [d for d in self.storage if d.get('roomId') == room]
return list(self.storage)


def test_ingest_and_query_recent(monkeypatch):
# Inject a fake services.db module to avoid importing real dependencies (redis/mongo)
fake_db = types.ModuleType('services.db')
fake_analytics = FakeColl()
fake_db.analytics_coll = fake_analytics
sys.modules['services.db'] = fake_db

# Now import the analytics service (it will import services.db from sys.modules)
import services.analytics_service as analytics_service

ev = analytics_service.ingest_event({
'roomId': 'room-test-1',
'userId': 'user-123',
'eventType': 'stroke_created',
'payload': {'color': '#ff0000'},
'ts': int(time.time() * 1000)
})

assert ev is not None
# userId should be removed and anonUserId present
assert 'anonUserId' in ev and 'userId' not in ev

recent = analytics_service.query_recent('room-test-1', limit=10)
assert isinstance(recent, list)
assert len(recent) >= 1
# Ensure stored doc contains our eventType
assert any(r.get('eventType') == 'stroke_created' for r in recent)
Loading
Loading