"""FlowDeck — Public API v2 : Agent & Skill marketplace (v6.6.0, phase 5). Thin Bearer+scopes wrappers over the existing agent logic (AgentEngine, `agent_skills`, the gallery service) so third-party integrations can drive FlowDeck Agent without a browser session: * ``/api/v2/agents`` — agents CRUD, conversations, synchronous runs (JSON, the SSE stream stays an internal/UI concern), audit journal & rollback. * ``/api/v2/skills`` — the skill marketplace: CRUD, portable export/import and the built-in gallery of installable presets. Rules honoured (see docs/API_GUIDE_V6.md): one code path (the engine and the gallery service are reused, never re-implemented), JSON only, no secrets or internal columns, rate limit + audit + idempotency on every mutation. """ from __future__ import annotations import json import time from fastapi import APIRouter, Header, HTTPException, Request from fastapi.responses import JSONResponse from app.db import get_conn from app.routers.agent import _default_agent from app.services import skill_gallery from app.services.agent_engine import AgentEngine, undo_action from app.services.api_v2_helpers import ( audit_log, check_idempotency, check_v2_rate_limit, get_bearer_user, has_scope, paginate_headers, parse_pagination, row_to_dict, store_idempotency, ) from app.services.llm_client import LLMClient from app.services.llm_config import get_user_llm_key router = APIRouter(prefix="/api/v2", tags=["api-v2-agent"]) # ── Shared guards ────────────────────────────────────────────────────────── def _guard(request: Request, authorization: str | None, *, write: bool = False) -> dict: """Bearer auth + per-token rate limit (+ write scope when required).""" user = get_bearer_user(request, authorization) ip = request.client.host if request.client else "unknown" if not check_v2_rate_limit(user.get("_token_hash"), ip): raise HTTPException(429, "Rate limit exceeded: 300 req/min per token") if write and not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") return user async def _json_body(request: Request) -> dict: try: body = await request.json() except Exception: # noqa: BLE001 return {} return body if isinstance(body, dict) else {} def _workspace_of(request: Request, body: dict | None = None) -> int | None: """Workspace resolution mirrors the internal agent router: explicit param wins, then the token's own workspace, else NULL (shared/global scope).""" body = body or {} raw = body.get("workspace_id") or request.query_params.get("workspace_id") if raw is None: return None try: return int(raw) except (TypeError, ValueError): return None def _owned_conversation(conn, conversation_id: int, user_id: int): """Conversation visible to this token's user (ownership is enforced here, unlike the session router where the browser is already authenticated).""" return conn.execute( "SELECT * FROM agent_conversations WHERE id=? AND user_id=?", (conversation_id, user_id), ).fetchone() def _engine_for(user_id: int, workspace_id: int | None, provider: str | None) -> AgentEngine: engine = AgentEngine(user_id, workspace_id=workspace_id) if provider: user_key = get_user_llm_key(user_id, provider) if user_key and user_key.get("api_key"): engine.llm = LLMClient( provider=provider, api_key=user_key["api_key"], api_base=user_key.get("api_base") or None, ) else: engine.llm = LLMClient(provider=provider) return engine # ── Agents ───────────────────────────────────────────────────────────────── @router.get("/agents") async def list_agents_v2(request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization) limit, offset = parse_pagination(request) ws = _workspace_of(request) with get_conn() as conn: _default_agent(conn, user["id"]) clause = "WHERE workspace_id IS ? OR workspace_id=?" total = conn.execute(f"SELECT COUNT(*) FROM agents {clause}", (ws, ws)).fetchone()[0] rows = conn.execute( f"SELECT * FROM agents {clause} ORDER BY agent_type, name LIMIT ? OFFSET ?", (ws, ws, limit, offset), ).fetchall() return JSONResponse( content={"agents": [row_to_dict(r) for r in rows], "total": total, "limit": limit, "offset": offset}, headers=paginate_headers(total), ) @router.post("/agents") async def create_agent_v2(request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) body = await _json_body(request) name = (body.get("name") or "").strip() or "Custom Agent" ws = _workspace_of(request, body) with get_conn() as conn: try: cur = conn.execute( """INSERT INTO agents (workspace_id, name, icon, agent_type, description, system_instructions, model, scope_json, trigger_json, approval_mode, created_by) VALUES (?,?,?,?,?,?,?,?,?,?,?)""", (ws, name, body.get("icon", "🤖"), body.get("agent_type", "custom"), body.get("description", ""), body.get("system_instructions", ""), body.get("model", "gpt-4o"), json.dumps(body.get("scope", {})), json.dumps(body.get("trigger", {})), body.get("approval_mode", "auto"), user["id"]), ) conn.commit() agent_id = cur.lastrowid row = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() except Exception as exc: # noqa: BLE001 raise HTTPException(409, f"Cannot create agent: {exc}") from exc audit_log(user, "agent.create", "agent", agent_id, name, request) data = {"id": agent_id, "name": name, "status": "created", "agent": row_to_dict(row)} key = (request.headers.get("Idempotency-Key") or "").strip() if key: store_idempotency(key, user["id"], data, 201) return JSONResponse(content=data, status_code=201) @router.get("/agents/{agent_id}") async def get_agent_v2(agent_id: int, request: Request, authorization: str | None = Header(default=None)): _guard(request, authorization) with get_conn() as conn: row = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() if not row: raise HTTPException(404, "Agent not found") return row_to_dict(row) @router.put("/agents/{agent_id}") async def update_agent_v2(agent_id: int, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) body = await _json_body(request) with get_conn() as conn: existing = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() if not existing: raise HTTPException(404, "Agent not found") sets, params = [], [] for col in ("name", "icon", "description", "system_instructions", "model", "approval_mode", "is_active"): if col in body: sets.append(f"{col}=?") params.append(body[col]) if "scope" in body: sets.append("scope_json=?") params.append(json.dumps(body["scope"])) if "trigger" in body: sets.append("trigger_json=?") params.append(json.dumps(body["trigger"])) if sets: params.append(agent_id) conn.execute(f"UPDATE agents SET {', '.join(sets)} WHERE id=?", params) conn.commit() audit_log(user, "agent.update", "agent", agent_id, "", request) return {"id": agent_id, "status": "updated"} @router.delete("/agents/{agent_id}") async def delete_agent_v2(agent_id: int, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) with get_conn() as conn: if not conn.execute("SELECT id FROM agents WHERE id=?", (agent_id,)).fetchone(): raise HTTPException(404, "Agent not found") conn.execute("DELETE FROM agents WHERE id=?", (agent_id,)) conn.commit() audit_log(user, "agent.delete", "agent", agent_id, "", request) return {"id": agent_id, "status": "deleted"} # ── Conversations (static paths declared before /agents/{agent_id}) ──────── @router.get("/agents/conversations") async def list_conversations_v2(request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization) limit, offset = parse_pagination(request) with get_conn() as conn: total = conn.execute( "SELECT COUNT(*) FROM agent_conversations WHERE user_id=?", (user["id"],) ).fetchone()[0] rows = conn.execute( """SELECT * FROM agent_conversations WHERE user_id=? ORDER BY updated_at DESC LIMIT ? OFFSET ?""", (user["id"], limit, offset), ).fetchall() return JSONResponse( content={"conversations": [row_to_dict(r) for r in rows], "total": total, "limit": limit, "offset": offset}, headers=paginate_headers(total), ) @router.post("/agents/conversations") async def create_conversation_v2(request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) body = await _json_body(request) ws = _workspace_of(request, body) agent_id = body.get("agent_id") with get_conn() as conn: if agent_id is not None: agent = conn.execute("SELECT id FROM agents WHERE id=?", (agent_id,)).fetchone() if not agent: raise HTTPException(404, "Agent not found") agent_id = agent["id"] else: agent_id = _default_agent(conn, user["id"])["id"] cur = conn.execute( """INSERT INTO agent_conversations (agent_id, user_id, title, context_json, provider, model) VALUES (?,?,?,?,?,?)""", (agent_id, user["id"], body.get("title") or "New conversation", json.dumps({"workspace_id": ws}), body.get("provider") or "", body.get("model") or ""), ) conv_id = cur.lastrowid conn.commit() row = conn.execute("SELECT * FROM agent_conversations WHERE id=?", (conv_id,)).fetchone() audit_log(user, "agent.conversation.create", "agent_conversation", conv_id, "", request) data = {"id": conv_id, "status": "created", "conversation": row_to_dict(row)} key = (request.headers.get("Idempotency-Key") or "").strip() if key: store_idempotency(key, user["id"], data, 201) return JSONResponse(content=data, status_code=201) @router.get("/agents/conversations/{conversation_id}") async def get_conversation_v2(conversation_id: int, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization) with get_conn() as conn: conv = _owned_conversation(conn, conversation_id, user["id"]) if not conv: raise HTTPException(404, "Conversation not found") messages = conn.execute( "SELECT * FROM agent_messages WHERE conversation_id=? ORDER BY created_at, id", (conversation_id,), ).fetchall() return {"conversation": row_to_dict(conv), "messages": [row_to_dict(m) for m in messages]} @router.delete("/agents/conversations/{conversation_id}") async def delete_conversation_v2(conversation_id: int, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) with get_conn() as conn: if not _owned_conversation(conn, conversation_id, user["id"]): raise HTTPException(404, "Conversation not found") conn.execute("DELETE FROM agent_conversations WHERE id=?", (conversation_id,)) conn.commit() audit_log(user, "agent.conversation.delete", "agent_conversation", conversation_id, "", request) return {"id": conversation_id, "status": "deleted"} @router.get("/agents/conversations/{conversation_id}/actions") async def list_actions_v2(conversation_id: int, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization) with get_conn() as conn: if not _owned_conversation(conn, conversation_id, user["id"]): raise HTTPException(404, "Conversation not found") rows = conn.execute( "SELECT * FROM agent_actions WHERE conversation_id=? ORDER BY created_at, id", (conversation_id,), ).fetchall() return {"actions": [row_to_dict(r) for r in rows]} @router.post("/agents/actions/{action_id}/undo") async def undo_action_v2(action_id: int, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) with get_conn() as conn: row = conn.execute( """SELECT a.id FROM agent_actions a JOIN agent_conversations c ON c.id = a.conversation_id WHERE a.id=? AND c.user_id=?""", (action_id, user["id"]), ).fetchone() if not row: raise HTTPException(404, "Action not found") try: undo_action(action_id) except ValueError as exc: raise HTTPException(400, str(exc)) from exc except Exception as exc: # noqa: BLE001 raise HTTPException(500, f"Rollback failed: {exc}") from exc audit_log(user, "agent.action.undo", "agent_action", action_id, "", request) return {"id": action_id, "status": "reverted"} # ── Runs (JSON — the SSE stream stays internal) ──────────────────────────── def _collect_run_events(events: list[dict]) -> dict: """Aggregate an engine event stream into a JSON run result. Engine events are flat (``{"type": "final", "content": ...}``), the same shape the SSE panel consumes. """ final = None reasoning = [] actions = [] error = None for ev in events: etype = ev.get("type") if etype == "final": final = ev.get("content") or final elif etype == "reasoning": reasoning.append(ev.get("content") or "") elif etype == "action": actions.append({k: v for k, v in ev.items() if k != "type"}) elif etype == "error": error = ev.get("message") or "run failed" return { "status": "failed" if error else "completed", "final": final, "error": error, "reasoning": reasoning, "actions": actions, } @router.post("/agents/conversations/{conversation_id}/run") async def run_conversation_v2(conversation_id: int, request: Request, authorization: str | None = Header(default=None)): """Synchronous agent run: buffers the engine stream and returns JSON. Third parties get one HTTP round-trip instead of an SSE subscription; the same AgentEngine, permissions, journal and webhooks are used as the UI. """ user = _guard(request, authorization, write=True) idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) body = await _json_body(request) objective = (body.get("message") or body.get("objective") or "").strip() if not objective: raise HTTPException(400, "message is required") with get_conn() as conn: conv = _owned_conversation(conn, conversation_id, user["id"]) if not conv: raise HTTPException(404, "Conversation not found") eff_provider = body.get("provider") or conv["provider"] or None eff_model = body.get("model") or conv["model"] or None if body.get("provider") is not None or body.get("model") is not None: conn.execute( "UPDATE agent_conversations SET provider=?, model=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (body.get("provider", conv["provider"] or ""), body.get("model", conv["model"] or ""), conversation_id), ) conn.commit() conv_context = {} try: conv_context = json.loads(conv["context_json"] or "{}") or {} except (TypeError, ValueError): conv_context = {} ws = _workspace_of(request, body) if ws is None: ws = conv_context.get("workspace_id") engine = _engine_for(user["id"], ws, eff_provider) started = time.time() events = [ ev async for ev in engine.run( conversation_id, objective, model=eff_model, mentions=body.get("mentions"), files=body.get("files"), skill_id=body.get("skill_id"), skill_ids=body.get("skill_ids"), extra_context=body.get("context"), ) ] result = _collect_run_events(events) with get_conn() as conn: actions = conn.execute( "SELECT * FROM agent_actions WHERE conversation_id=? ORDER BY created_at, id", (conversation_id,), ).fetchall() payload = { "conversation_id": conversation_id, "status": result["status"], "final": result["final"], "error": result["error"], "reasoning": result["reasoning"], "actions": [row_to_dict(a) for a in actions], "events": events, "duration_ms": int((time.time() - started) * 1000), } audit_log(user, "agent.run", "agent_conversation", conversation_id, objective[:200], request) status_code = 200 if result["status"] == "completed" else 500 data = payload key = (request.headers.get("Idempotency-Key") or "").strip() if key: store_idempotency(key, user["id"], data, status_code) return JSONResponse(content=data, status_code=status_code) @router.post("/agents/{agent_id}/trigger") async def trigger_agent_v2(agent_id: int, request: Request, authorization: str | None = Header(default=None)): """Fire a custom agent from an external integration (JSON, synchronous).""" user = _guard(request, authorization, write=True) body = await _json_body(request) ws = _workspace_of(request, body) with get_conn() as conn: agent = conn.execute("SELECT * FROM agents WHERE id=?", (agent_id,)).fetchone() if not agent: raise HTTPException(404, "Agent not found") cur = conn.execute( """INSERT INTO agent_conversations (agent_id, user_id, title, context_json) VALUES (?,?,?,?)""", (agent_id, user["id"], f"Run: {agent['name']}", json.dumps({"workspace_id": ws})), ) conv_id = cur.lastrowid conn.commit() objective = (agent["system_instructions"] or "").strip() or f"Exécute l'agent « {agent['name']} »." if body.get("message"): objective = f"{objective}\n\n{body['message']}" engine = _engine_for(user["id"], ws, agent["model"] or None) started = time.time() events = [ev async for ev in engine.run(conv_id, objective, model=agent["model"])] result = _collect_run_events(events) payload = { "conversation_id": conv_id, "agent_id": agent_id, "status": result["status"], "final": result["final"], "error": result["error"], "reasoning": result["reasoning"], "actions": result["actions"], "duration_ms": int((time.time() - started) * 1000), } audit_log(user, "agent.trigger", "agent", agent_id, objective[:200], request) return JSONResponse(content=payload, status_code=200 if result["status"] == "completed" else 500) # ── Skill marketplace ────────────────────────────────────────────────────── @router.get("/skills") async def list_skills_v2(request: Request, authorization: str | None = Header(default=None)): _guard(request, authorization) limit, offset = parse_pagination(request) ws = _workspace_of(request) with get_conn() as conn: total = conn.execute( "SELECT COUNT(*) FROM agent_skills WHERE workspace_id IS ? OR workspace_id=?", (ws, ws), ).fetchone()[0] rows = conn.execute( """SELECT * FROM agent_skills WHERE workspace_id IS ? OR workspace_id=? ORDER BY name LIMIT ? OFFSET ?""", (ws, ws, limit, offset), ).fetchall() return JSONResponse( content={"skills": [row_to_dict(r) for r in rows], "total": total, "limit": limit, "offset": offset}, headers=paginate_headers(total), ) @router.post("/skills") async def create_skill_v2(request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) body = await _json_body(request) ws = _workspace_of(request, body) try: fields = skill_gallery.parse_payload( {k: body[k] for k in ("name", "description", "prompt_template", "allowed_tools") if k in body} | {"format": skill_gallery.EXPORT_FORMAT} ) except ValueError as exc: raise HTTPException(400, str(exc)) from exc try: row, created = skill_gallery.upsert_skill( fields, workspace_id=ws, created_by=user["id"], overwrite=bool(body.get("overwrite")), ) except ValueError as exc: raise HTTPException(409, str(exc)) from exc audit_log(user, "skill.create", "skill", row.get("id"), fields["name"], request) data = {"id": row.get("id"), "name": fields["name"], "status": "created" if created else "updated", "skill": row_to_dict(row) if row else {}} key = (request.headers.get("Idempotency-Key") or "").strip() if key: store_idempotency(key, user["id"], data, 201 if created else 200) return JSONResponse(content=data, status_code=201 if created else 200) # Gallery & import are static segments: declared before /skills/{skill_id} so # FastAPI never tries to coerce "gallery" into an int path parameter. @router.get("/skills/gallery") async def skills_gallery_v2(request: Request, authorization: str | None = Header(default=None)): _guard(request, authorization) presets = skill_gallery.list_gallery() return {"gallery": presets, "total": len(presets), "install": "POST /api/v2/skills/gallery/{slug}/install"} @router.post("/skills/gallery/{slug}/install") async def install_gallery_skill_v2(slug: str, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) preset = skill_gallery.get_gallery(slug) if not preset: raise HTTPException(404, f"Unknown gallery skill: {slug}") body = await _json_body(request) ws = _workspace_of(request, body) try: row, created = skill_gallery.upsert_skill( skill_gallery.parse_payload(preset), workspace_id=ws, created_by=user["id"], overwrite=bool(body.get("overwrite", True)), ) except ValueError as exc: raise HTTPException(409, str(exc)) from exc audit_log(user, "skill.gallery.install", "skill", row.get("id"), slug, request) data = {"slug": slug, "id": row.get("id"), "name": row.get("name"), "status": "installed" if created else "updated", "skill": row_to_dict(row)} return JSONResponse(content=data, status_code=201 if created else 200) @router.post("/skills/import") async def import_skill_v2(request: Request, authorization: str | None = Header(default=None)): """Import a portable skill document (from another FlowDeck instance).""" user = _guard(request, authorization, write=True) idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) body = await _json_body(request) payload = body.get("payload") if isinstance(body.get("payload"), dict) else body try: fields = skill_gallery.parse_payload(payload) except ValueError as exc: raise HTTPException(400, str(exc)) from exc ws = _workspace_of(request, body) try: row, created = skill_gallery.upsert_skill( fields, workspace_id=ws, created_by=user["id"], overwrite=bool(body.get("overwrite")), ) except ValueError as exc: raise HTTPException(409, str(exc)) from exc audit_log(user, "skill.import", "skill", row.get("id"), fields["name"], request) data = {"id": row.get("id"), "name": fields["name"], "status": "imported" if created else "updated", "skill": row_to_dict(row)} key = (request.headers.get("Idempotency-Key") or "").strip() if key: store_idempotency(key, user["id"], data, 201 if created else 200) return JSONResponse(content=data, status_code=201 if created else 200) @router.get("/skills/{skill_id}") async def get_skill_v2(skill_id: int, request: Request, authorization: str | None = Header(default=None)): _guard(request, authorization) with get_conn() as conn: row = conn.execute("SELECT * FROM agent_skills WHERE id=?", (skill_id,)).fetchone() if not row: raise HTTPException(404, "Skill not found") return row_to_dict(row) @router.get("/skills/{skill_id}/export") async def export_skill_v2(skill_id: int, request: Request, authorization: str | None = Header(default=None)): """Portable JSON document — POST it to /api/v2/skills/import elsewhere.""" _guard(request, authorization) with get_conn() as conn: row = conn.execute("SELECT * FROM agent_skills WHERE id=?", (skill_id,)).fetchone() if not row: raise HTTPException(404, "Skill not found") return skill_gallery.export_skill(row) @router.delete("/skills/{skill_id}") async def delete_skill_v2(skill_id: int, request: Request, authorization: str | None = Header(default=None)): user = _guard(request, authorization, write=True) with get_conn() as conn: row = conn.execute("SELECT name FROM agent_skills WHERE id=?", (skill_id,)).fetchone() if not row: raise HTTPException(404, "Skill not found") conn.execute("DELETE FROM agent_skills WHERE id=?", (skill_id,)) conn.commit() audit_log(user, "skill.delete", "skill", skill_id, row["name"] or "", request) return {"id": skill_id, "status": "deleted"} @router.post("/skills/{skill_id}/apply") async def apply_skill_v2(skill_id: int, request: Request, authorization: str | None = Header(default=None)): """Open a conversation pre-loaded with the skill (ready to run).""" user = _guard(request, authorization, write=True) body = await _json_body(request) ws = _workspace_of(request, body) with get_conn() as conn: skill = conn.execute("SELECT * FROM agent_skills WHERE id=?", (skill_id,)).fetchone() if not skill: raise HTTPException(404, "Skill not found") agent_id = _default_agent(conn, user["id"])["id"] cur = conn.execute( """INSERT INTO agent_conversations (agent_id, user_id, title, context_json) VALUES (?,?,?,?)""", (agent_id, user["id"], skill["name"], json.dumps({"workspace_id": ws if ws is not None else skill["workspace_id"], "skill_id": skill_id})), ) conv_id = cur.lastrowid conn.commit() audit_log(user, "skill.apply", "skill", skill_id, skill["name"], request) return JSONResponse(content={"conversation_id": conv_id, "skill": skill["name"], "status": "ready"}, status_code=201)