Manager Dashboard and Log Management System
Using the manager-worker pattern in production made it clear that a dedicated manager dashboard was necessary. I built a system to monitor worker status, task progress, and logs all from a single screen.
Dashboard API Design
from fastapi import APIRouter, Depends
from typing import List
router = APIRouter(prefix="/api/manager", tags=["manager"])
@router.get("/overview")
async def get_overview(
manager: ManagerSession = Depends(get_manager),
):
workers = await manager.get_worker_status()
tasks = await manager.get_task_summary()
return {
"total_workers": len(workers),
"active_workers": sum(1 for w in workers if w["status"] == "active"),
"total_tasks": tasks["total"],
"completed_tasks": tasks["completed"],
"failed_tasks": tasks["failed"],
"pending_tasks": tasks["pending"],
"uptime_seconds": manager.get_uptime(),
}
@router.get("/workers")
async def list_workers(
manager: ManagerSession = Depends(get_manager),
):
workers = await manager.get_worker_details()
return {"data": workers}
@router.get("/workers/{worker_id}/logs")
async def get_worker_logs(
worker_id: str,
limit: int = 100,
level: str = None,
manager: ManagerSession = Depends(get_manager),
):
logs = await manager.get_worker_logs(worker_id, limit=limit, level=level)
return {"data": logs}
Log Collection and Aggregation
I implemented a centralized log management system that collects logs from each worker and supports search and filtering.
class LogAggregator:
def __init__(self, redis_store):
self.store = redis_store
async def collect_logs(self, session_ids: List[str]) -> List[dict]:
all_logs = []
for sid in session_ids:
key = f"claude:session:{sid}:logs"
logs = await self.store.redis.lrange(key, 0, -1)
for log_str in logs:
log = json.loads(log_str)
log["session_id"] = sid
all_logs.append(log)
all_logs.sort(key=lambda x: x["timestamp"], reverse=True)
return all_logs
async def search_logs(self, query: str, session_ids: List[str] = None):
logs = await self.collect_logs(session_ids or [])
return [
log for log in logs
if query.lower() in log.get("content", "").lower()
]
async def get_error_summary(self, session_ids: List[str]) -> dict:
logs = await self.collect_logs(session_ids)
errors = [l for l in logs if l["level"] == "error"]
return {
"total_errors": len(errors),
"recent_errors": errors[:10],
"error_rate": len(errors) / max(len(logs), 1),
}
Real-Time Monitoring
SSE (Server-Sent Events) delivers real-time updates to the dashboard.
from fastapi.responses import StreamingResponse
@router.get("/stream")
async def stream_events(manager: ManagerSession = Depends(get_manager)):
async def event_generator():
pubsub = manager.store.redis.pubsub()
await pubsub.subscribe("claude:events")
async for message in pubsub.listen():
if message["type"] == "message":
data = json.loads(message["data"])
yield f"data: {json.dumps(data)}\n\n"
return StreamingResponse(
event_generator(),
media_type="text/event-stream",
)
Log Retention Policy
As logs accumulate they put pressure on Redis memory, so I applied a retention policy.
async def apply_retention_policy(self, max_age_hours: int = 24):
cutoff = datetime.now() - timedelta(hours=max_age_hours)
sessions = await self.store.redis.smembers("claude:sessions:active")
for sid in sessions:
key = f"claude:session:{sid}:logs"
logs = await self.store.redis.lrange(key, 0, -1)
for i, log_str in enumerate(reversed(logs)):
log = json.loads(log_str)
if datetime.fromisoformat(log["timestamp"]) < cutoff:
await self.store.redis.ltrim(key, 0, len(logs) - i - 2)
break
Wrap-Up
The manager dashboard made operations significantly easier. The error summary and real-time monitoring in particular enabled fast response when issues arose.