Files
flowdeck/app/routers/api_v2_agent.py
T
bruno 224bda74d5
FlowDeck CI / lint (push) Canceled after 0s
FlowDeck CI / test (push) Canceled after 0s
FlowDeck CI / docker (push) Canceled after 0s
fix: A21 phase 1 — 352 routes async sans await → threadpool (v7.8.0)
- Conversion `async def` → `def` de TOUTES les routes dont le corps ne contient
  ni `await`, ni `async with`, ni `async for`, ni `asyncio` (scan automatique
  corps par corps sur app/ : 352 converties, 0 dangereuses, vérifié
  `asyncio`/`run_coroutine`/`.result()` absents). FastAPI exécute ces handlers
  dans son threadpool → tout leur SQLite (`get_conn()` + `conn.execute`) quitte
  l'event loop, sans changer une ligne de logique.
- Répartition : api_v2 60, dashboard 40, collections 25, board 23,
  workspace 19, wiki 17, permissions 14, api 14, main.py 6, + 35 fichiers.
- Les 4 routers prioritaires de l'audit sont couverts par ce lot :
  api_v2 60 + dashboard 40 + collections 25 + board 23 = 148 conversions
  (le reste de leurs routes attend la phase 2 : elles ont de vrais `await`).
- Reste (phase 2) : les 311 routes avec de vrais `await` → enrouler les blocs
  DB dans `await anyio.to_thread.run_sync(...)` ; pas de wrapper partagé livré
  (rien ne l'appellerait — YAGNI jusqu'au premier usage).

suite **1037/1037** (233 s) · `ruff check app tests` OK · docs à jour
2026-10-01 10:53:26 -04:00

658 lines
28 KiB
Python

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