Spaces:
Running
Running
| """AP Commander β FastAPI server for AP Clerk, Oversight, and Curriculum endpoints.""" | |
| from __future__ import annotations | |
| import uuid | |
| import logging | |
| from datetime import datetime, timedelta | |
| from typing import Any, Dict, List, Optional | |
| from fastapi import FastAPI, HTTPException, Request | |
| from fastapi.responses import JSONResponse, HTMLResponse | |
| from fastapi.middleware.cors import CORSMiddleware | |
| from .models import ( | |
| ResetRequest, ResetResponse, | |
| StepRequest, StepResponse, | |
| StateResponse, TaskInfo, | |
| OversightResetRequest, OversightResetResponse, | |
| OversightStepRequest, OversightStepResponse, | |
| CurriculumRequest, CurriculumResponse, | |
| ) | |
| from .environment import APClerkEnvironment | |
| from .tasks import TASKS | |
| from .ui import HTML_PAGE | |
| from oversight_environment import OversightEnvironment | |
| logging.basicConfig(level=logging.INFO) | |
| logger = logging.getLogger("ap-clerk-env") | |
| app = FastAPI( | |
| title="AP Commander β Multi-Agent Enterprise Financial Environment", | |
| description=( | |
| "OpenEnv-compatible multi-agent environment for enterprise AP workflows. " | |
| "AP Clerk agent (27 tasks) + Fleet AI Oversight Agent + Adaptive Curriculum. " | |
| "Themes: Multi-Agent (#1), Long-Horizon Planning (#2), World Modeling (#3), " | |
| "Self-Improvement (#4)." | |
| ), | |
| version="4.0.0", | |
| ) | |
| app.add_middleware( | |
| CORSMiddleware, | |
| allow_origins=["*"], | |
| allow_methods=["*"], | |
| allow_headers=["*"], | |
| ) | |
| _sessions: Dict[str, APClerkEnvironment] = {} | |
| _oversight_sessions: Dict[str, OversightEnvironment] = {} | |
| _session_times: Dict[str, datetime] = {} | |
| _session_run_ids: Dict[str, str] = {} # session_id β run_id | |
| _SESSION_TTL = timedelta(hours=1) | |
| # Server-side curriculum history β populated by /step when done=True. | |
| # Keyed by run_id (X-Run-Id request header, or "default"). | |
| # Rolling window of 200 entries per run_id to prevent unbounded growth. | |
| _curriculum_history: Dict[str, List[dict]] = {} | |
| _CURRICULUM_WINDOW = 200 | |
| # In-memory stats counters | |
| _stats: Dict[str, Any] = { | |
| "total_episodes": 0, | |
| "completed_episodes": 0, | |
| "total_score": 0.0, | |
| "task_counts": {}, # task_id β {"count": int, "score_sum": float} | |
| } | |
| def _prune_sessions() -> None: | |
| """Remove sessions older than _SESSION_TTL.""" | |
| cutoff = datetime.utcnow() - _SESSION_TTL | |
| stale = [sid for sid, t in _session_times.items() if t < cutoff] | |
| for sid in stale: | |
| _sessions.pop(sid, None) | |
| _oversight_sessions.pop(sid, None) | |
| _session_times.pop(sid, None) | |
| _session_run_ids.pop(sid, None) | |
| def _get_session(session_id: str) -> APClerkEnvironment: | |
| if session_id not in _sessions: | |
| raise HTTPException( | |
| status_code=404, | |
| detail=f"Session {session_id!r} not found. Call /reset first." | |
| ) | |
| return _sessions[session_id] | |
| def _get_oversight_session(session_id: str) -> OversightEnvironment: | |
| if session_id not in _oversight_sessions: | |
| raise HTTPException( | |
| status_code=404, | |
| detail=f"Oversight session {session_id!r} not found. Call /oversight/reset first." | |
| ) | |
| return _oversight_sessions[session_id] | |
| async def root(): | |
| return HTML_PAGE | |
| async def health(): | |
| return { | |
| "status": "ok", | |
| "environment": "ap-commander", | |
| "version": "4.0.0", | |
| "agents": ["ap-clerk", "oversight"], | |
| "total_tasks": len(TASKS), | |
| "themes": ["multi-agent", "long-horizon", "world-modeling", "self-improvement"], | |
| } | |
| async def list_tasks(request: Request = None): | |
| tasks = [ | |
| TaskInfo(task_id=tid, name=spec.name, difficulty=spec.difficulty, description=spec.description) | |
| for tid, spec in TASKS.items() | |
| ] | |
| accept = (request.headers.get("accept", "") if request else "") | |
| if "text/html" in accept: | |
| return HTMLResponse(_render_tasks_page(tasks)) | |
| return tasks | |
| def _render_tasks_page(tasks: list) -> str: | |
| diff_color = {"easy": "#22c55e", "medium": "#eab308", "hard": "#ef4444", "long-horizon": "#8b5cf6", "oversight": "#0ea5e9"} | |
| diff_order = ["easy", "medium", "hard", "long-horizon", "oversight"] | |
| diff_label = {"easy": "Easy", "medium": "Medium", "hard": "Hard", "long-horizon": "Long-Horizon", "oversight": "Oversight"} | |
| groups: dict = {d: [] for d in diff_order} | |
| for t in tasks: | |
| groups.setdefault(t.difficulty, []).append(t) | |
| sections_html = "" | |
| for diff in diff_order: | |
| bucket = groups.get(diff, []) | |
| if not bucket: | |
| continue | |
| col = diff_color.get(diff, "#94a3b8") | |
| label = diff_label.get(diff, diff) | |
| cards = "" | |
| for t in bucket: | |
| spec = TASKS.get(t.task_id) | |
| steps = getattr(spec, "max_steps", 1) if spec else 1 | |
| cards += f""" | |
| <div class="task-card"> | |
| <div class="card-top"> | |
| <span class="badge" style="background:{col}22;color:{col};border:1px solid {col}44">{label}</span> | |
| <span class="steps-pill">{steps} step{"s" if steps > 1 else ""}</span> | |
| </div> | |
| <div class="task-name">{t.name}</div> | |
| <div class="task-id">{t.task_id}</div> | |
| <div class="task-desc">{t.description}</div> | |
| </div>""" | |
| sections_html += f""" | |
| <div class="diff-section"> | |
| <div class="diff-header"> | |
| <span class="diff-dot" style="background:{col}"></span> | |
| <h2>{label}</h2> | |
| <span class="diff-count">{len(bucket)} task{"s" if len(bucket) != 1 else ""}</span> | |
| </div> | |
| <div class="cards-grid">{cards}</div> | |
| </div>""" | |
| total = len(tasks) | |
| return f"""<!DOCTYPE html> | |
| <html lang="en"> | |
| <head> | |
| <meta charset="UTF-8"><meta name="viewport" content="width=device-width,initial-scale=1"> | |
| <title>AP Commander β Task Library</title> | |
| <style> | |
| *,*::before,*::after{{box-sizing:border-box;margin:0;padding:0}} | |
| :root{{--bg:#020817;--bg2:#0a1628;--bg3:#0f1f3d;--glass:rgba(10,22,56,0.7);--border:rgba(14,165,233,0.18);--border2:rgba(14,165,233,0.35);--accent:#0ea5e9;--teal:#14b8a6;--text:#f0f9ff;--dim:#94a3b8;--dimmer:#475569;--grad:linear-gradient(135deg,#0ea5e9,#14b8a6)}} | |
| body{{background:var(--bg);color:var(--text);font-family:-apple-system,BlinkMacSystemFont,'Segoe UI',system-ui,sans-serif;line-height:1.6;min-height:100vh}} | |
| body::before{{content:'';position:fixed;inset:0;background-image:linear-gradient(rgba(14,165,233,0.04) 1px,transparent 1px),linear-gradient(90deg,rgba(14,165,233,0.04) 1px,transparent 1px);background-size:40px 40px;pointer-events:none;z-index:0}} | |
| nav{{position:sticky;top:0;z-index:100;display:flex;align-items:center;justify-content:space-between;padding:0 40px;height:60px;background:rgba(2,8,23,0.9);backdrop-filter:blur(16px);border-bottom:1px solid var(--border)}} | |
| .nav-brand{{display:flex;align-items:center;gap:10px;text-decoration:none;color:var(--text);font-size:16px;font-weight:800;background:var(--grad);-webkit-background-clip:text;-webkit-text-fill-color:transparent;background-clip:text}} | |
| .nav-links{{display:flex;gap:8px}} | |
| .nav-links a{{color:var(--dim);text-decoration:none;font-size:13px;padding:5px 12px;border-radius:7px;transition:all .2s}} | |
| .nav-links a:hover{{color:var(--text);background:rgba(14,165,233,0.08)}} | |
| .nav-links a.cta{{background:var(--accent);color:#020817;font-weight:700}} | |
| .nav-links a.cta:hover{{background:#38bdf8}} | |
| main{{max-width:1100px;margin:0 auto;padding:60px 40px;position:relative;z-index:1}} | |
| .page-header{{margin-bottom:56px}} | |
| .page-label{{font-size:11px;font-weight:700;text-transform:uppercase;letter-spacing:1.5px;color:var(--accent);margin-bottom:10px}} | |
| .page-title{{font-size:clamp(28px,4vw,44px);font-weight:900;letter-spacing:-1.5px;line-height:1.1;margin-bottom:12px;background:var(--grad);-webkit-background-clip:text;-webkit-text-fill-color:transparent;background-clip:text}} | |
| .page-sub{{font-size:15px;color:var(--dim);max-width:500px;line-height:1.7}} | |
| .total-pill{{display:inline-flex;align-items:center;gap:6px;background:rgba(14,165,233,0.1);border:1px solid var(--border2);border-radius:20px;padding:5px 14px;font-size:12px;color:var(--accent);font-weight:700;margin-top:16px}} | |
| .diff-section{{margin-bottom:48px}} | |
| .diff-header{{display:flex;align-items:center;gap:10px;margin-bottom:18px;padding-bottom:14px;border-bottom:1px solid var(--border)}} | |
| .diff-dot{{width:10px;height:10px;border-radius:50%;flex-shrink:0}} | |
| .diff-header h2{{font-size:18px;font-weight:800;letter-spacing:-.3px}} | |
| .diff-count{{font-size:12px;color:var(--dim);background:rgba(255,255,255,0.05);border:1px solid var(--border);border-radius:20px;padding:2px 10px;font-weight:600;margin-left:auto}} | |
| .cards-grid{{display:grid;grid-template-columns:repeat(auto-fill,minmax(300px,1fr));gap:14px}} | |
| .task-card{{background:var(--glass);backdrop-filter:blur(12px);border:1px solid var(--border);border-radius:12px;padding:20px;transition:border-color .2s,transform .2s,box-shadow .2s}} | |
| .task-card:hover{{border-color:var(--border2);transform:translateY(-2px);box-shadow:0 12px 32px rgba(14,165,233,0.08)}} | |
| .card-top{{display:flex;align-items:center;gap:8px;margin-bottom:10px}} | |
| .badge{{display:inline-flex;align-items:center;padding:3px 10px;border-radius:20px;font-size:10px;font-weight:700;text-transform:uppercase;letter-spacing:.4px}} | |
| .steps-pill{{margin-left:auto;font-size:10px;color:var(--dimmer);background:rgba(255,255,255,0.04);border:1px solid var(--border);border-radius:20px;padding:2px 8px;font-weight:600;font-family:monospace}} | |
| .task-name{{font-size:15px;font-weight:800;letter-spacing:-.2px;margin-bottom:4px}} | |
| .task-id{{font-size:10px;color:var(--dimmer);font-family:monospace;margin-bottom:10px}} | |
| .task-desc{{font-size:12px;color:var(--dim);line-height:1.6}} | |
| footer{{text-align:center;padding:32px;font-size:12px;color:var(--dimmer);border-top:1px solid var(--border);margin-top:24px;position:relative;z-index:1}} | |
| @media(max-width:700px){{main{{padding:40px 20px}}nav{{padding:0 20px}}.cards-grid{{grid-template-columns:1fr}}}} | |
| </style> | |
| </head> | |
| <body> | |
| <nav> | |
| <a class="nav-brand" href="/">⬑ AP Commander</a> | |
| <div class="nav-links"> | |
| <a href="/">Home</a> | |
| <a href="/docs">API Docs</a> | |
| <a href="https://github.com/Vayuputra2401/RL-Agent" target="_blank">GitHub</a> | |
| <a href="https://huggingface.co/spaces/Pathikreet/ap-commander-training" target="_blank" class="cta">Training Space β</a> | |
| </div> | |
| </nav> | |
| <main> | |
| <div class="page-header"> | |
| <div class="page-label">AP Commander</div> | |
| <div class="page-title">Task Library</div> | |
| <div class="page-sub">Every task is generated fresh from a seeded RNG β same seed, same invoice. No static dataset. The agent must reason, not memorise.</div> | |
| <div class="total-pill">⬑ {total} tasks across 5 difficulty tiers</div> | |
| </div> | |
| {sections_html} | |
| </main> | |
| <footer>AP Commander Β· Pathikreet Chowdhury Β· Anubhav Bhattacharya Β· Radhika Ravi Β· Meta PyTorch OpenEnv Γ Scaler 2026</footer> | |
| </body> | |
| </html>""" | |
| async def reset(body: Optional[ResetRequest] = None, request: Request = None): | |
| if body is None: | |
| body = ResetRequest() | |
| _prune_sessions() | |
| session_id = body.session_id or str(uuid.uuid4()) | |
| run_id = (request.headers.get("X-Run-Id", "default") if request else "default") | |
| env = APClerkEnvironment() | |
| try: | |
| obs = env.reset(body.task_id, seed=body.seed) | |
| except ValueError as e: | |
| raise HTTPException(status_code=400, detail=str(e)) | |
| _sessions[session_id] = env | |
| _session_times[session_id] = datetime.utcnow() | |
| _session_run_ids[session_id] = run_id | |
| # Track episode start in stats | |
| _stats["total_episodes"] += 1 | |
| tid = body.task_id | |
| if tid not in _stats["task_counts"]: | |
| _stats["task_counts"][tid] = {"count": 0, "score_sum": 0.0} | |
| logger.info("reset session=%s task=%s seed=%s", session_id, body.task_id, body.seed) | |
| return ResetResponse( | |
| observation=obs, | |
| session_id=session_id, | |
| info={"message": f"Episode started for task '{body.task_id}'", | |
| "seed": body.seed}, | |
| ) | |
| async def step(body: StepRequest): | |
| env = _get_session(body.session_id) | |
| try: | |
| obs, reward, done, info = env.step(body.action) | |
| except RuntimeError as e: | |
| raise HTTPException(status_code=400, detail=str(e)) | |
| # Update stats on episode completion | |
| if done: | |
| _stats["completed_episodes"] += 1 | |
| _stats["total_score"] += reward.score | |
| tid = obs.task_id | |
| if tid not in _stats["task_counts"]: | |
| _stats["task_counts"][tid] = {"count": 0, "score_sum": 0.0} | |
| _stats["task_counts"][tid]["count"] += 1 | |
| _stats["task_counts"][tid]["score_sum"] += reward.score | |
| # Record in server-side curriculum history (keyed by run_id from /reset) | |
| run_id = _session_run_ids.get(body.session_id, "default") | |
| bucket = _curriculum_history.setdefault(run_id, []) | |
| bucket.append({"task_id": tid, "score": reward.score}) | |
| if len(bucket) > _CURRICULUM_WINDOW: | |
| del bucket[:-_CURRICULUM_WINDOW] | |
| logger.info( | |
| "step session=%s decision=%s score=%.3f", | |
| body.session_id, body.action.decision, reward.score, | |
| ) | |
| return StepResponse(observation=obs, reward=reward, done=done, info=info) | |
| async def state(session_id: str): | |
| env = _get_session(session_id) | |
| s = env.state | |
| return StateResponse( | |
| session_id=session_id, | |
| task_id=s["task_id"], | |
| step_count=s["step_count"], | |
| episode_score=s["episode_score"], | |
| done=s["done"], | |
| current_observation=s["current_observation"], | |
| ) | |
| async def stats(): | |
| completed = _stats["completed_episodes"] | |
| mean_score = round(_stats["total_score"] / completed, 4) if completed > 0 else 0.0 | |
| task_breakdown = { | |
| tid: { | |
| "count": v["count"], | |
| "mean_score": round(v["score_sum"] / v["count"], 4) if v["count"] > 0 else 0.0, | |
| } | |
| for tid, v in _stats["task_counts"].items() | |
| } | |
| return { | |
| "total_episodes": _stats["total_episodes"], | |
| "completed_episodes": completed, | |
| "mean_score": mean_score, | |
| "task_breakdown": task_breakdown, | |
| } | |
| # ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| # Oversight Agent endpoints (Theme #1 Fleet AI) | |
| # ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| async def oversight_reset(body: Optional[OversightResetRequest] = None): | |
| """ | |
| Start a new Oversight Agent session. | |
| The agent receives a batch of completed AP Clerk episodes and must | |
| identify which ones contain suspicious/fraudulent decisions. | |
| """ | |
| if body is None: | |
| body = OversightResetRequest() | |
| _prune_sessions() | |
| session_id = str(uuid.uuid4()) | |
| env = OversightEnvironment() | |
| obs = env.reset(seed=body.seed, num_episodes=body.num_episodes) | |
| obs.session_id = session_id | |
| _oversight_sessions[session_id] = env | |
| _session_times[session_id] = datetime.utcnow() | |
| logger.info("oversight_reset session=%s num_episodes=%d", session_id, body.num_episodes) | |
| return OversightResetResponse( | |
| observation=obs, | |
| session_id=session_id, | |
| info={ | |
| "message": f"Oversight session started with {body.num_episodes} episodes to review.", | |
| "seed": body.seed, | |
| "instructions": ( | |
| "Review each episode_summary and submit OversightAction for each. " | |
| "Verdicts: CLEAR | FLAG_FOR_REVIEW | ESCALATE_TO_AUDIT. " | |
| "Include specific numeric signal in your reasoning." | |
| ), | |
| }, | |
| ) | |
| async def oversight_step(body: OversightStepRequest): | |
| """ | |
| Submit an Oversight verdict for one episode in the batch. | |
| Repeat until all episodes are reviewed (done=True). | |
| """ | |
| env = _get_oversight_session(body.session_id) | |
| try: | |
| obs, reward, done, info = env.step(body.action) | |
| except RuntimeError as e: | |
| raise HTTPException(status_code=400, detail=str(e)) | |
| logger.info( | |
| "oversight_step session=%s episode=%s verdict=%s score=%.3f", | |
| body.session_id, body.action.episode_id, body.action.verdict, reward.score, | |
| ) | |
| return OversightStepResponse(observation=obs, reward=reward, done=done, info=info) | |
| async def oversight_state(session_id: str): | |
| """Get current state of an Oversight session.""" | |
| env = _get_oversight_session(session_id) | |
| return {"session_id": session_id, **env.state} | |
| # ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| # Adaptive Curriculum endpoint (Theme #4 Self-Improvement) | |
| # ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | |
| # Task difficulty ordering for curriculum | |
| _DIFFICULTY_ORDER = ["easy", "medium", "hard", "long-horizon", "oversight"] | |
| _CURRICULUM_THRESHOLDS = { | |
| "easy": 0.70, # unlock medium | |
| "medium": 0.65, # unlock hard | |
| "hard": 0.68, # unlock long-horizon | |
| "long-horizon": 0.72, # unlock oversight | |
| } | |
| _TASKS_BY_DIFFICULTY: Dict[str, List[str]] = {} | |
| for _tid, _spec in TASKS.items(): | |
| _d = _spec.difficulty | |
| _TASKS_BY_DIFFICULTY.setdefault(_d, []).append(_tid) | |
| async def curriculum_next_task(body: Optional[CurriculumRequest] = None, request: Request = None): | |
| """ | |
| Recommend the next task based on server-verified episode history. | |
| The curriculum uses server-side records populated by /step completions. | |
| Client-provided session_history is accepted only as a cold-start fallback | |
| when no server-side history exists for this run_id. | |
| Difficulty ladder: easy β medium β hard β long-horizon β oversight | |
| """ | |
| if body is None: | |
| body = CurriculumRequest() | |
| run_id = (request.headers.get("X-Run-Id", "default") if request else "default") | |
| server_history = _curriculum_history.get(run_id, []) | |
| # Use server-side history; fall back to client only when server has no data | |
| if server_history: | |
| raw_history = server_history | |
| else: | |
| raw_history = [{"task_id": e.task_id, "score": e.score} | |
| for e in body.session_history] | |
| history = raw_history # list of dicts with task_id + score | |
| # Compute mean score per difficulty level from recent history | |
| # history entries are dicts: {"task_id": str, "score": float} | |
| scores_by_diff: Dict[str, List[float]] = {} | |
| for entry in history: | |
| tid = entry["task_id"] if isinstance(entry, dict) else entry.task_id | |
| sc = entry["score"] if isinstance(entry, dict) else entry.score | |
| spec = TASKS.get(tid) | |
| if spec: | |
| diff = spec.difficulty | |
| scores_by_diff.setdefault(diff, []).append(sc) | |
| mean_by_diff = { | |
| d: sum(s) / len(s) for d, s in scores_by_diff.items() if s | |
| } | |
| # Determine highest unlocked difficulty | |
| current_diff = "easy" | |
| unlocked = ["easy"] | |
| for diff in _DIFFICULTY_ORDER[1:]: | |
| prev_diff = _DIFFICULTY_ORDER[_DIFFICULTY_ORDER.index(diff) - 1] | |
| thresh = _CURRICULUM_THRESHOLDS.get(prev_diff, 0.70) | |
| if mean_by_diff.get(prev_diff, 0.0) >= thresh: | |
| current_diff = diff | |
| unlocked.append(diff) | |
| else: | |
| break | |
| # Choose least-practiced task at current difficulty | |
| available = _TASKS_BY_DIFFICULTY.get(current_diff, []) | |
| if not available: | |
| available = _TASKS_BY_DIFFICULTY.get("easy", []) | |
| current_diff = "easy" | |
| practiced_counts: Dict[str, int] = {} | |
| for entry in history: | |
| tid = entry["task_id"] if isinstance(entry, dict) else entry.task_id | |
| practiced_counts[tid] = practiced_counts.get(tid, 0) + 1 | |
| recommended = min(available, key=lambda tid: practiced_counts.get(tid, 0)) | |
| # Build reason string | |
| if current_diff == "easy": | |
| reason = "Starting with easy tasks to build foundational skills." | |
| else: | |
| prev = _DIFFICULTY_ORDER[_DIFFICULTY_ORDER.index(current_diff) - 1] | |
| mean = mean_by_diff.get(prev, 0.0) | |
| thresh = _CURRICULUM_THRESHOLDS.get(prev, 0.70) | |
| reason = ( | |
| f"Mean score on {prev} tasks is {mean:.2f} " | |
| f"(threshold: {thresh:.2f}) β unlocked {current_diff}." | |
| ) | |
| return CurriculumResponse( | |
| recommended_task_id=recommended, | |
| difficulty=current_diff, | |
| reason=reason, | |
| unlocked_tasks=[t for d in unlocked for t in _TASKS_BY_DIFFICULTY.get(d, [])], | |
| ) | |
| async def generic_handler(request: Request, exc: Exception): | |
| logger.exception("Unhandled error: %s", exc) | |
| return JSONResponse(status_code=500, content={"detail": str(exc)}) | |
| def start(): | |
| import uvicorn | |
| uvicorn.run("app.main:app", host="0.0.0.0", port=7860, workers=1) | |