"""FlowDeck — Public API v2 (v6.3.0). Full Notion-clone REST wrapper around existing services. Thin wrappers: one code path, JSON only, Bearer+scopes, pagination, RFC7807 errors delegated to global handler, audit + idempotency. """ from __future__ import annotations import hashlib import json import logging import secrets from datetime import datetime from fastapi import APIRouter, Header, HTTPException, Request from fastapi.responses import JSONResponse, Response from app.db import get_conn 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, to_iso8601, validate_scopes_input, ) from app.services.automations import fire_event as _fire_event logger = logging.getLogger(__name__) router = APIRouter(prefix="/api/v2", tags=["api-v2"]) def _hash(token: str) -> str: return hashlib.sha256(token.encode()).hexdigest() def _v2_rate_check(request: Request, user: dict) -> None: ip = request.client.host if request.client else "unknown" th = user.get("_token_hash") if not check_v2_rate_limit(th, ip): raise HTTPException(status_code=429, detail="Rate limit exceeded: 300 req/min per token") # ── Tokens ──────────────────────────────────────────────────────────────── @router.post("/tokens") async def create_token(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "API token").strip()[:100] scopes = validate_scopes_input(body.get("scopes") or "read,write") expires_at = body.get("expires_at") # Idempotency idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) token = f"fd_{secrets.token_urlsafe(32)}" prefix = token[:12] th = _hash(token) exp_val = None if expires_at: try: # accept ISO string exp_val = str(expires_at) # validate parse datetime.fromisoformat(exp_val.replace("Z", "+00:00")) except Exception as err: raise HTTPException(400, "Invalid expires_at, use ISO-8601") from err with get_conn() as conn: try: cur = conn.execute( "INSERT INTO api_tokens (user_id, name, token_hash, token_prefix, scopes, expires_at) VALUES (?, ?, ?, ?, ?, ?)", (user["id"], name, th, prefix, scopes, exp_val), ) conn.commit() tid = cur.lastrowid except Exception as e: raise HTTPException(409, f"Token creation failed: {e}") from None row = conn.execute("SELECT id, name, token_prefix, scopes, expires_at, created_at FROM api_tokens WHERE id=?", (tid,)).fetchone() audit_log(user, "token.create", "api_token", tid, f"scopes={scopes}", request) data = {"id": tid, "name": row["name"], "token": token, "prefix": prefix, "scopes": scopes, "expires_at": to_iso8601(row["expires_at"]) if row["expires_at"] else None, "note": "Copy token now — shown once. Use as Authorization: Bearer "} key = (request.headers.get("Idempotency-Key") or request.headers.get("idempotency-key") or "").strip() if key: store_idempotency(key, user["id"], data, 200) return data @router.get("/tokens") async def list_tokens(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT id, name, token_prefix, scopes, expires_at, last_used_at, created_at, revoked FROM api_tokens WHERE user_id=? ORDER BY created_at DESC", (user["id"],)).fetchall() out = [] for r in rows: d = dict(r) d["created_at"] = to_iso8601(d.get("created_at")) d["expires_at"] = to_iso8601(d.get("expires_at")) if d.get("expires_at") else None d["last_used_at"] = to_iso8601(d.get("last_used_at")) if d.get("last_used_at") else None # never expose hash out.append({k: v for k, v in d.items() if k != "token_hash"}) return {"tokens": out} @router.delete("/tokens/{token_id}") async def revoke_token(token_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT id, user_id FROM api_tokens WHERE id=?", (token_id,)).fetchone() if not row: raise HTTPException(404, "Token not found") if row["user_id"] != user["id"] and not user.get("is_admin"): raise HTTPException(403, "Not your token") conn.execute("UPDATE api_tokens SET revoked=1 WHERE id=?", (token_id,)) conn.commit() audit_log(user, "token.revoke", "api_token", token_id, "", request) return {"id": token_id, "status": "revoked"} @router.post("/tokens/{token_id}/rotate") async def rotate_token(token_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT id, user_id, name, scopes FROM api_tokens WHERE id=?", (token_id,)).fetchone() if not row: raise HTTPException(404, "Token not found") if row["user_id"] != user["id"] and not user.get("is_admin"): raise HTTPException(403, "Not your token") # revoke old conn.execute("UPDATE api_tokens SET revoked=1 WHERE id=?", (token_id,)) new_token = f"fd_{secrets.token_urlsafe(32)}" th = _hash(new_token) prefix = new_token[:12] cur = conn.execute("INSERT INTO api_tokens (user_id, name, token_hash, token_prefix, scopes) VALUES (?, ?, ?, ?, ?)", (row["user_id"], row["name"], th, prefix, row["scopes"] or "read,write")) conn.commit() nid = cur.lastrowid audit_log(user, "token.rotate", "api_token", token_id, f"new_id={nid}", request) return {"id": nid, "token": new_token, "prefix": prefix, "scopes": row["scopes"], "note": "Copy token now — shown once"} # ── Users ───────────────────────────────────────────────────────────────── @router.get("/users/me") async def get_me(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: row = conn.execute("SELECT id, login, full_name, email, avatar_url, avatar_color, is_admin, is_active, auth_method, sidebar_config, notification_prefs, timezone, created_at FROM users WHERE id=?", (user["id"],)).fetchone() if not row: raise HTTPException(404, "User not found") d = row_to_dict(row) # parse json prefs for k in ("notification_prefs", "sidebar_config"): if isinstance(d.get(k), str): try: d[k] = json.loads(d[k] or "{}") except Exception: pass # never expose secrets return d @router.patch("/users/me") async def patch_me(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} allowed = {"full_name", "email", "avatar_color", "notification_prefs", "sidebar_config", "timezone"} updates = {} for k in allowed: if k in body: updates[k] = body[k] if not updates: raise HTTPException(400, "No updatable fields") # validation if "email" in updates and updates["email"] and "@" not in str(updates["email"]): raise HTTPException(400, "Invalid email") with get_conn() as conn: sets = [] params = [] for k, v in updates.items(): if k in ("notification_prefs", "sidebar_config"): v = json.dumps(v) if isinstance(v, (dict, list)) else str(v) sets.append(f"{k}=?") params.append(v) params.append(user["id"]) conn.execute(f"UPDATE users SET {', '.join(sets)} WHERE id=?", params) conn.commit() row = conn.execute("SELECT id, login, full_name, email, avatar_url, avatar_color, is_admin, timezone, notification_prefs FROM users WHERE id=?", (user["id"],)).fetchone() audit_log(user, "user.update", "user", user["id"], "", request) return row_to_dict(row) @router.get("/users/search") async def search_users(request: Request, q: str = "", authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) q = (q or request.query_params.get("q") or "").strip() if not q: return {"users": []} like = f"%{q}%" with get_conn() as conn: rows = conn.execute("SELECT id, login, full_name, email, avatar_url, avatar_color FROM users WHERE login LIKE ? OR email LIKE ? OR full_name LIKE ? LIMIT 20", (like, like, like)).fetchall() return {"users": [dict(r) for r in rows]} # ── Workspaces ──────────────────────────────────────────────────────────── @router.get("/workspaces") async def list_workspaces(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: total = conn.execute("SELECT COUNT(*) FROM workspaces WHERE owner_id=? OR id IN (SELECT workspace_id FROM workspace_members WHERE user_id=?)", (user["id"], user["id"])).fetchone()[0] rows = conn.execute("SELECT w.*, wm.role FROM workspaces w LEFT JOIN workspace_members wm ON wm.workspace_id=w.id AND wm.user_id=? WHERE w.owner_id=? OR w.id IN (SELECT workspace_id FROM workspace_members WHERE user_id=?) ORDER BY w.created_at DESC LIMIT ? OFFSET ?", (user["id"], user["id"], user["id"], limit, offset)).fetchall() out = [] for r in rows: d = dict(r) d["created_at"] = to_iso8601(d.get("created_at")) try: d["settings"] = json.loads(d.get("settings_json") or "{}") except Exception: d["settings"] = {} out.append(d) return {"workspaces": out, "total": total, "limit": limit, "offset": offset} @router.post("/workspaces") async def create_workspace(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) try: body = await request.json() except Exception: body = {} name = (body.get("name") or "").strip() if not name: raise HTTPException(400, "name is required") settings_json = json.dumps(body.get("settings") or body.get("settings_json") or {}) with get_conn() as conn: cur = conn.execute("INSERT INTO workspaces (name, owner_id, settings_json) VALUES (?, ?, ?)", (name, user["id"], settings_json)) wid = cur.lastrowid # owner is implicitly admin member try: conn.execute("INSERT OR IGNORE INTO workspace_members (workspace_id, user_id, role) VALUES (?, ?, 'owner')", (wid, user["id"])) except Exception: pass conn.commit() row = conn.execute("SELECT * FROM workspaces WHERE id=?", (wid,)).fetchone() audit_log(user, "workspace.create", "workspace", wid, name, request) data = {"id": wid, "name": name, "owner_id": user["id"], "status": "created", "workspace": row_to_dict(row)} key = (request.headers.get("Idempotency-Key") or "").strip() if key: store_idempotency(key, user["id"], data, 200) return JSONResponse(content=data, status_code=201) @router.get("/workspaces/{workspace_id}") async def get_workspace(workspace_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: row = conn.execute("SELECT * FROM workspaces WHERE id=?", (workspace_id,)).fetchone() if not row: raise HTTPException(404, "Workspace not found") # ACL: must be member or owner member = conn.execute("SELECT role FROM workspace_members WHERE workspace_id=? AND user_id=?", (workspace_id, user["id"])).fetchone() is_owner = row["owner_id"] == user["id"] if not is_owner and not member and not user.get("is_admin"): raise HTTPException(404, "Workspace not found") members = conn.execute("SELECT u.id, u.login, u.full_name, u.avatar_url, wm.role FROM workspace_members wm JOIN users u ON u.id=wm.user_id WHERE wm.workspace_id=? ORDER BY wm.joined_at", (workspace_id,)).fetchall() d = row_to_dict(row) d["members"] = [dict(m) for m in members] return d @router.patch("/workspaces/{workspace_id}") async def patch_workspace(workspace_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM workspaces WHERE id=?", (workspace_id,)).fetchone() if not row: raise HTTPException(404, "Workspace not found") if row["owner_id"] != user["id"] and not user.get("is_admin"): # check admin member mem = conn.execute("SELECT role FROM workspace_members WHERE workspace_id=? AND user_id=?", (workspace_id, user["id"])).fetchone() if not mem or mem["role"] not in ("owner", "admin"): raise HTTPException(403, "Only owner/admin can edit workspace") name = body.get("name", row["name"]) sj = body.get("settings_json") or body.get("settings") if sj is not None: sj = json.dumps(sj) if isinstance(sj, (dict, list)) else str(sj) else: sj = row["settings_json"] conn.execute("UPDATE workspaces SET name=?, settings_json=? WHERE id=?", (name, sj, workspace_id)) conn.commit() audit_log(user, "workspace.update", "workspace", workspace_id, "", request) return {"id": workspace_id, "status": "updated"} @router.delete("/workspaces/{workspace_id}") async def delete_workspace(workspace_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT * FROM workspaces WHERE id=?", (workspace_id,)).fetchone() if not row: raise HTTPException(404, "Workspace not found") if row["owner_id"] != user["id"] and not user.get("is_admin"): raise HTTPException(403, "Only owner can delete workspace") conn.execute("DELETE FROM workspaces WHERE id=?", (workspace_id,)) conn.commit() audit_log(user, "workspace.delete", "workspace", workspace_id, "", request) return {"id": workspace_id, "status": "deleted"} @router.get("/workspaces/{workspace_id}/members") async def list_workspace_members(workspace_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: ws = conn.execute("SELECT id FROM workspaces WHERE id=?", (workspace_id,)).fetchone() if not ws: raise HTTPException(404, "Workspace not found") rows = conn.execute("SELECT u.id, u.login, u.full_name, u.avatar_url, wm.role, wm.joined_at FROM workspace_members wm JOIN users u ON u.id=wm.user_id WHERE wm.workspace_id=? ORDER BY wm.joined_at", (workspace_id,)).fetchall() return {"members": [row_to_dict(r) for r in rows]} @router.post("/workspaces/{workspace_id}/members") async def invite_member(workspace_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} target_id = body.get("user_id") or body.get("uid") email = (body.get("email") or "").strip() role = (body.get("role") or "editor").strip().lower() if role not in ("owner", "admin", "editor", "viewer", "commenter"): raise HTTPException(400, "Invalid role") with get_conn() as conn: ws = conn.execute("SELECT owner_id FROM workspaces WHERE id=?", (workspace_id,)).fetchone() if not ws: raise HTTPException(404, "Workspace not found") # only owner/admin can invite if ws["owner_id"] != user["id"] and not user.get("is_admin"): mem = conn.execute("SELECT role FROM workspace_members WHERE workspace_id=? AND user_id=?", (workspace_id, user["id"])).fetchone() if not mem or mem["role"] not in ("owner", "admin"): raise HTTPException(403, "Only owner/admin can invite") uid = target_id if not uid and email: u = conn.execute("SELECT id FROM users WHERE email=?", (email,)).fetchone() if not u: raise HTTPException(404, f"User with email {email} not found") uid = u["id"] if not uid: raise HTTPException(400, "user_id or email required") try: conn.execute("INSERT INTO workspace_members (workspace_id, user_id, role) VALUES (?, ?, ?)", (workspace_id, uid, role)) except Exception: conn.execute("UPDATE workspace_members SET role=? WHERE workspace_id=? AND user_id=?", (role, workspace_id, uid)) conn.commit() audit_log(user, "workspace.invite", "workspace", workspace_id, f"uid={uid} role={role}", request) return {"workspace_id": workspace_id, "user_id": uid, "role": role, "status": "added"} @router.patch("/workspaces/{workspace_id}/members/{uid}") async def update_member_role(workspace_id: int, uid: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} role = (body.get("role") or "").strip().lower() if role not in ("owner", "admin", "editor", "viewer", "commenter"): raise HTTPException(400, "Invalid role") with get_conn() as conn: ws = conn.execute("SELECT owner_id FROM workspaces WHERE id=?", (workspace_id,)).fetchone() if not ws: raise HTTPException(404, "Workspace not found") if ws["owner_id"] != user["id"] and not user.get("is_admin"): mem = conn.execute("SELECT role FROM workspace_members WHERE workspace_id=? AND user_id=?", (workspace_id, user["id"])).fetchone() if not mem or mem["role"] not in ("owner", "admin"): raise HTTPException(403, "Only owner/admin can change roles") conn.execute("UPDATE workspace_members SET role=? WHERE workspace_id=? AND user_id=?", (role, workspace_id, uid)) if conn.total_changes == 0: raise HTTPException(404, "Member not found") conn.commit() audit_log(user, "workspace.role_change", "workspace", workspace_id, f"uid={uid} role={role}", request) return {"workspace_id": workspace_id, "user_id": uid, "role": role, "status": "updated"} @router.delete("/workspaces/{workspace_id}/members/{uid}") async def remove_member(workspace_id: int, uid: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: ws = conn.execute("SELECT owner_id FROM workspaces WHERE id=?", (workspace_id,)).fetchone() if not ws: raise HTTPException(404, "Workspace not found") if ws["owner_id"] != user["id"] and not user.get("is_admin"): mem = conn.execute("SELECT role FROM workspace_members WHERE workspace_id=? AND user_id=?", (workspace_id, user["id"])).fetchone() if not mem or mem["role"] not in ("owner", "admin"): raise HTTPException(403, "Only owner/admin can remove members") conn.execute("DELETE FROM workspace_members WHERE workspace_id=? AND user_id=?", (workspace_id, uid)) conn.commit() audit_log(user, "workspace.remove_member", "workspace", workspace_id, f"uid={uid}", request) return {"workspace_id": workspace_id, "user_id": uid, "status": "removed"} # ── Collections ─────────────────────────────────────────────────────────── @router.get("/collections") async def list_collections_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) ws_filter = request.query_params.get("workspace_id") q = (request.query_params.get("query") or "").strip() with get_conn() as conn: where = [] params: list = [] if ws_filter: try: wid = int(ws_filter) where.append("c.workspace_id=?") params.append(wid) except ValueError: pass if q: where.append("(c.name LIKE ? OR c.description LIKE ?)") like = f"%{q}%" params.extend([like, like]) clause = ("WHERE " + " AND ".join(where)) if where else "" total = conn.execute(f"SELECT COUNT(*) FROM collections c {clause}", params).fetchone()[0] rows = conn.execute(f"SELECT c.* FROM collections c {clause} ORDER BY c.name LIMIT ? OFFSET ?", (*params, limit, offset)).fetchall() cols = [] for r in rows: d = row_to_dict(r) # filter by visibility: skip private not visible (best-effort) cols.append(d) resp = {"collections": cols, "total": total, "limit": limit, "offset": offset} return JSONResponse(content=resp, headers=paginate_headers(total)) @router.post("/collections") async def create_collection_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) try: body = await request.json() except Exception: body = {} name = (body.get("name") or "").strip() if not name: raise HTTPException(400, "name is required") description = body.get("description", "") icon = body.get("icon", "📋") workspace_id = body.get("workspace_id") schema = body.get("schema") or body.get("schema_json") or [] if isinstance(schema, str): try: schema = json.loads(schema) except Exception: schema = [] schema_json = json.dumps(schema) with get_conn() as conn: cur = conn.execute("INSERT INTO collections (name, description, icon, schema_json, workspace_id, created_by) VALUES (?, ?, ?, ?, ?, ?)", (name, description, icon, schema_json, workspace_id, user["id"])) cid = cur.lastrowid # materialize properties if schema provided try: from app.services.db_templates import materialize_properties materialize_properties(conn, cid, schema) except Exception: pass # default view try: conn.execute("INSERT INTO collection_views (collection_id, name, view_type, config_json) VALUES (?, ?, ?, ?)", (cid, "Default View", "table", json.dumps({"visible_properties": ["Title"]}))) except Exception: pass conn.commit() row = conn.execute("SELECT * FROM collections WHERE id=?", (cid,)).fetchone() audit_log(user, "collection.create", "collection", cid, name, request) data = {"id": cid, "name": name, "status": "created", "collection": 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("/collections/{collection_id}") async def get_collection_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: row = conn.execute("SELECT * FROM collections WHERE id=?", (collection_id,)).fetchone() if not row: raise HTTPException(404, "Collection not found") pages = conn.execute("SELECT id, title, icon, position, property_values_json, created_at FROM collection_pages WHERE collection_id=? ORDER BY position LIMIT 50", (collection_id,)).fetchall() d = row_to_dict(row) d["pages"] = [row_to_dict(p) for p in pages] return d @router.patch("/collections/{collection_id}") async def patch_collection_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM collections WHERE id=?", (collection_id,)).fetchone() if not row: raise HTTPException(404, "Collection not found") name = body.get("name", row["name"]) description = body.get("description", row["description"]) icon = body.get("icon", row["icon"]) schema = body.get("schema") or body.get("schema_json") if schema is not None: sj = json.dumps(schema) if isinstance(schema, (list, dict)) else str(schema) else: sj = row["schema_json"] conn.execute("UPDATE collections SET name=?, description=?, icon=?, schema_json=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (name, description, icon, sj, collection_id)) conn.commit() audit_log(user, "collection.update", "collection", collection_id, "", request) return {"id": collection_id, "status": "updated"} @router.delete("/collections/{collection_id}") async def delete_collection_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT * FROM collections WHERE id=?", (collection_id,)).fetchone() if not row: raise HTTPException(404, "Collection not found") conn.execute("DELETE FROM collections WHERE id=?", (collection_id,)) conn.commit() audit_log(user, "collection.delete", "collection", collection_id, "", request) return {"id": collection_id, "status": "deleted"} @router.post("/collections/{collection_id}/linked") async def create_linked_db(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "").strip() or f"Linked DB {collection_id}" with get_conn() as conn: src = conn.execute("SELECT * FROM collections WHERE id=?", (collection_id,)).fetchone() if not src: raise HTTPException(404, "Collection not found") cur = conn.execute("INSERT INTO collections (name, description, icon, schema_json, workspace_id, created_by) VALUES (?, ?, ?, ?, ?, ?)", (name, src["description"], src["icon"], src["schema_json"], src["workspace_id"] if "workspace_id" in src.keys() else None, user["id"])) nid = cur.lastrowid # copy data source as linked try: conn.execute("INSERT INTO collection_data_sources (collection_id, source_collection_id, is_linked) VALUES (?, ?, 1)", (nid, collection_id)) except Exception: pass # copy views + properties (light) rows = conn.execute("SELECT * FROM collection_properties WHERE collection_id=?", (collection_id,)).fetchall() for p in rows: try: conn.execute("INSERT INTO collection_properties (collection_id, name, prop_type, options_json, position) VALUES (?, ?, ?, ?, ?)", (nid, p["name"], p["prop_type"], p["options_json"], p["position"])) except Exception: pass vrows = conn.execute("SELECT * FROM collection_views WHERE collection_id=?", (collection_id,)).fetchall() for v in vrows: try: conn.execute("INSERT INTO collection_views (collection_id, name, view_type, config_json, position) VALUES (?, ?, ?, ?, ?)", (nid, v["name"], v["view_type"], v["config_json"], v["position"])) except Exception: pass conn.commit() audit_log(user, "collection.linked", "collection", nid, f"src={collection_id}", request) return {"id": nid, "name": name, "status": "created"} @router.post("/collections/{collection_id}/task") async def toggle_task(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT is_task FROM collections WHERE id=?", (collection_id,)).fetchone() if not row: raise HTTPException(404, "Collection not found") cur_val = row["is_task"] if "is_task" in row.keys() else 0 new_val = 0 if cur_val else 1 conn.execute("UPDATE collections SET is_task=? WHERE id=?", (new_val, collection_id)) conn.commit() audit_log(user, "collection.toggle_task", "collection", collection_id, str(new_val), request) return {"id": collection_id, "is_task": bool(new_val)} @router.get("/collections/{collection_id}/sources") async def list_sources(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") rows = conn.execute("SELECT * FROM collection_data_sources WHERE collection_id=? ORDER BY position", (collection_id,)).fetchall() return {"sources": [dict(r) for r in rows]} @router.post("/collections/{collection_id}/sources") async def add_source(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} src_id = body.get("source_collection_id") or body.get("source_id") if not src_id: raise HTTPException(400, "source_collection_id required") with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") if not conn.execute("SELECT id FROM collections WHERE id=?", (src_id,)).fetchone(): raise HTTPException(404, "Source collection not found") try: conn.execute("INSERT INTO collection_data_sources (collection_id, source_collection_id) VALUES (?, ?)", (collection_id, src_id)) conn.commit() except Exception as e: raise HTTPException(409, str(e)) from None audit_log(user, "collection.add_source", "collection", collection_id, str(src_id), request) return {"collection_id": collection_id, "source_collection_id": src_id, "status": "added"} @router.delete("/collections/{collection_id}/sources/{source_id}") async def remove_source(collection_id: int, source_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM collection_data_sources WHERE collection_id=? AND (id=? OR source_collection_id=?)", (collection_id, source_id, source_id)) conn.commit() audit_log(user, "collection.remove_source", "collection", collection_id, str(source_id), request) return {"status": "removed"} # ── Collection pages ────────────────────────────────────────────────────── @router.get("/collections/{collection_id}/pages") async def list_collection_pages_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") total = conn.execute("SELECT COUNT(*) FROM collection_pages WHERE collection_id=?", (collection_id,)).fetchone()[0] # filters: filter[status]=Done etc., sort, fields # Simple: filter by property name via property_values_json LIKE (best-effort), sort by position or title sort = request.query_params.get("sort") or "" order = "position" desc = False if sort: if sort.startswith("-"): desc = True sort = sort[1:] # allow sorting by title/position/created_at if sort in ("title", "position", "created_at", "updated_at"): order = sort direction = "DESC" if desc else "ASC" rows = conn.execute(f"SELECT * FROM collection_pages WHERE collection_id=? ORDER BY {order} {direction} LIMIT ? OFFSET ?", (collection_id, limit, offset)).fetchall() # apply filter[xxx] in-memory (small) filters = {k[7:-1]: v for k, v in request.query_params.items() if k.startswith("filter[") and k.endswith("]")} fields = request.query_params.get("fields") fields_set = set(fields.split(",")) if fields else None out = [] for r in rows: d = row_to_dict(r) # property filter (AND) if filters: try: pv = json.loads(r["property_values_json"] or "{}") if isinstance(r["property_values_json"], str) else r["property_values_json"] except Exception: pv = {} ok = True for fk, fv in filters.items(): # lookup by prop id or name found = False for kk, vv in (pv or {}).items(): if str(kk) == str(fk) or str(kk).lower() == fk.lower(): if str(vv) == str(fv): found = True break # also check title if filter field is title if fk == "title" and d.get("title") == fv: found = True if not found: ok = False break if not ok: continue if fields_set: d = {k: v for k, v in d.items() if k in fields_set or k in ("id", "collection_id")} out.append(d) return JSONResponse(content={"pages": out, "total": total, "limit": limit, "offset": offset}, headers=paginate_headers(total)) @router.post("/collections/{collection_id}/pages") async def create_collection_page_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") idem = check_idempotency(request, user["id"]) if idem: return JSONResponse(content=idem["data"], status_code=idem["status"]) try: body = await request.json() except Exception: body = {} title = (body.get("title") or body.get("name") or "Untitled").strip() or "Untitled" icon = body.get("icon", "file") parent_id = body.get("parent_id") prop_vals = body.get("property_values") or body.get("properties") or body.get("property_values_json") or {} if isinstance(prop_vals, str): try: prop_vals = json.loads(prop_vals) except Exception: prop_vals = {} with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") # validate properties if helper exists try: pass # light validation: we rely on existing validators except Exception: pass max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_pages WHERE collection_id=?", (collection_id,)).fetchone()[0] # apply auto props try: props_list = [dict(r) for r in conn.execute("SELECT * FROM collection_properties WHERE collection_id=?", (collection_id,)).fetchall()] from app.services.property_types import apply_auto_properties as _aap _aap(props_list, prop_vals, user, is_create=True) except Exception: pass cur = conn.execute("INSERT INTO collection_pages (collection_id, title, icon, position, parent_id, property_values_json) VALUES (?, ?, ?, ?, ?, ?)", (collection_id, title, icon, max_pos, parent_id, json.dumps(prop_vals))) pid = cur.lastrowid conn.commit() row = conn.execute("SELECT * FROM collection_pages WHERE id=?", (pid,)).fetchone() audit_log(user, "page.create", "collection_page", pid, title, request) try: await _fire_event("collection.page.created", {"page_id": pid, "collection_id": collection_id, "title": title}) except Exception: pass data = {"id": pid, "title": title, "status": "created", "page": 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("/pages/{page_id}") async def get_page_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: row = conn.execute("SELECT * FROM collection_pages WHERE id=?", (page_id,)).fetchone() if not row: # also try pages table (block pages) row2 = conn.execute("SELECT * FROM pages WHERE id=?", (page_id,)).fetchone() if not row2: raise HTTPException(404, "Page not found") d = row_to_dict(row2) # v6.5.0: resolve synced blocks server-side (fresh content). if (d.get("content_format") or "blocks") == "blocks" and d.get("content"): from app.services.synced_blocks import resolve_content_json d["content"] = resolve_content_json(d["content"], d["content_format"]) return d d = row_to_dict(row) # property_values_json already parsed by row_to_dict # v6.5.0: expose the row's content page when it exists (no lazy # creation on a read-only endpoint). content_page_id = conn.execute( "SELECT id FROM pages WHERE collection_row_id=?", (page_id,), ).fetchone() d["content_page_id"] = content_page_id["id"] if content_page_id else None return d @router.patch("/pages/{page_id}") async def patch_page_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM collection_pages WHERE id=?", (page_id,)).fetchone() if not row: raise HTTPException(404, "Page not found") title = body.get("title", row["title"]) icon = body.get("icon", row["icon"]) pos = body.get("position", row["position"]) parent_id = body.get("parent_id", row["parent_id"]) pv_raw = row["property_values_json"] or "{}" try: stored = json.loads(pv_raw) if isinstance(pv_raw, str) else dict(pv_raw) except Exception: stored = {} incoming = body.get("property_values") or body.get("properties") if incoming is not None: if isinstance(incoming, str): try: incoming = json.loads(incoming) except Exception: incoming = {} # merge for k, v in (incoming or {}).items(): stored[str(k)] = v # apply auto props try: props_list = [dict(r) for r in conn.execute("SELECT * FROM collection_properties WHERE collection_id=?", (row["collection_id"],)).fetchall()] from app.services.property_types import apply_auto_properties as _aap _aap(props_list, stored, user, is_create=False) except Exception: pass conn.execute("UPDATE collection_pages SET title=?, icon=?, position=?, parent_id=?, property_values_json=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (title, icon, pos, parent_id, json.dumps(stored), page_id)) conn.commit() audit_log(user, "page.update", "collection_page", page_id, "", request) try: await _fire_event("collection.page.updated", {"page_id": page_id, "collection_id": row["collection_id"], "title": title}) except Exception: pass return {"id": page_id, "status": "updated"} @router.delete("/pages/{page_id}") async def delete_page_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT * FROM collection_pages WHERE id=?", (page_id,)).fetchone() if not row: raise HTTPException(404, "Page not found") conn.execute("DELETE FROM collection_pages WHERE id=?", (page_id,)) conn.commit() audit_log(user, "page.delete", "collection_page", page_id, "", request) try: await _fire_event("collection.page.deleted", {"page_id": page_id, "collection_id": row["collection_id"]}) except Exception: pass return {"id": page_id, "status": "deleted"} @router.post("/pages/{page_id}/restore") async def restore_page_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") # For soft-deleted pages (deleted_at) - but collection_pages has no deleted_at; handle pages table with get_conn() as conn: row = conn.execute("SELECT deleted_at FROM pages WHERE id=?", (page_id,)).fetchone() if row and row["deleted_at"]: conn.execute("UPDATE pages SET deleted_at=NULL WHERE id=?", (page_id,)) conn.commit() try: await _fire_event("page.restored", {"page_id": page_id}) except Exception: pass return {"id": page_id, "status": "restored"} raise HTTPException(404, "Page not found or not deleted") @router.post("/pages/{page_id}/move") async def move_page_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM collection_pages WHERE id=?", (page_id,)).fetchone() if not row: raise HTTPException(404, "Page not found") parent_id = body.get("parent_id", row["parent_id"]) position = body.get("position", row["position"]) conn.execute("UPDATE collection_pages SET parent_id=?, position=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (parent_id, position, page_id)) conn.commit() audit_log(user, "page.move", "collection_page", page_id, f"parent={parent_id} pos={position}", request) return {"id": page_id, "status": "moved"} @router.get("/pages/{page_id}/sub-items") async def list_sub_items_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: if not conn.execute("SELECT id FROM collection_pages WHERE id=?", (page_id,)).fetchone(): raise HTTPException(404, "Page not found") rows = conn.execute("SELECT * FROM collection_pages WHERE parent_id=? ORDER BY position", (page_id,)).fetchall() return {"sub_items": [row_to_dict(r) for r in rows]} @router.post("/pages/{page_id}/sub-items") async def create_sub_item_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: parent = conn.execute("SELECT collection_id FROM collection_pages WHERE id=?", (page_id,)).fetchone() if not parent: raise HTTPException(404, "Page not found") title = (body.get("title") or "Untitled").strip() max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_pages WHERE parent_id=?", (page_id,)).fetchone()[0] pv = json.dumps(body.get("property_values") or {}) cur = conn.execute("INSERT INTO collection_pages (collection_id, title, parent_id, position, property_values_json) VALUES (?, ?, ?, ?, ?)", (parent["collection_id"], title, page_id, max_pos, pv)) nid = cur.lastrowid conn.commit() audit_log(user, "page.create_subitem", "collection_page", nid, title, request) return {"id": nid, "status": "created"} @router.get("/pages/{page_id}/dependencies") async def list_dependencies_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT * FROM page_dependencies WHERE page_id=?", (page_id,)).fetchall() return {"dependencies": [dict(r) for r in rows]} @router.post("/pages/{page_id}/dependencies") async def add_dependency_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} dep_id = body.get("dependency_id") or body.get("depends_on") dtype = body.get("dependency_type") or "blocks" if not dep_id: raise HTTPException(400, "dependency_id required") with get_conn() as conn: if not conn.execute("SELECT id FROM collection_pages WHERE id=?", (page_id,)).fetchone(): raise HTTPException(404, "Page not found") if not conn.execute("SELECT id FROM collection_pages WHERE id=?", (dep_id,)).fetchone(): raise HTTPException(404, "Dependency page not found") try: conn.execute("INSERT INTO page_dependencies (page_id, dependency_id, dependency_type) VALUES (?, ?, ?)", (page_id, dep_id, dtype)) conn.commit() except Exception as e: raise HTTPException(409, str(e)) from None audit_log(user, "page.add_dependency", "collection_page", page_id, str(dep_id), request) return {"status": "added"} @router.delete("/pages/{page_id}/dependencies/{dep_id}") async def remove_dependency_v2(page_id: int, dep_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM page_dependencies WHERE page_id=? AND dependency_id=?", (page_id, dep_id)) conn.commit() audit_log(user, "page.remove_dependency", "collection_page", page_id, str(dep_id), request) return {"status": "removed"} # ── Properties ──────────────────────────────────────────────────────────── @router.get("/collections/{collection_id}/properties") async def list_properties_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") rows = conn.execute("SELECT * FROM collection_properties WHERE collection_id=? ORDER BY position", (collection_id,)).fetchall() return {"properties": [row_to_dict(r) for r in rows]} @router.post("/collections/{collection_id}/properties") async def create_property_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "").strip() if not name: raise HTTPException(400, "name is required") ptype = body.get("prop_type") or body.get("type") or "text" with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_properties WHERE collection_id=?", (collection_id,)).fetchone()[0] try: cur = conn.execute("INSERT INTO collection_properties (collection_id, name, prop_type, options_json, number_format, position, required, visible_in_views, validation_json, group_name) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", (collection_id, name, ptype, json.dumps(body.get("options") or []), body.get("number_format") or "number", max_pos, int(bool(body.get("required"))), int(bool(body.get("visible_in_views", True))), json.dumps(body.get("validation") or {}), (body.get("group_name") or "").strip())) conn.commit() pid = cur.lastrowid except Exception as e: raise HTTPException(409, f"Property exists: {e}") from None audit_log(user, "property.create", "property", pid, name, request) return {"id": pid, "name": name, "prop_type": ptype, "status": "created"} @router.patch("/properties/{prop_id}") async def patch_property_v2(prop_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM collection_properties WHERE id=?", (prop_id,)).fetchone() if not row: raise HTTPException(404, "Property not found") name = body.get("name", row["name"]) opts = json.dumps(body.get("options", json.loads(row["options_json"] or "[]"))) nf = body.get("number_format", row["number_format"]) req = int(bool(body.get("required", row["required"]))) vis = int(bool(body.get("visible_in_views", row["visible_in_views"]))) vj = json.dumps(body.get("validation", json.loads(row["validation_json"] or "{}"))) if "validation" in body else (row["validation_json"] if "validation_json" in row.keys() else "{}") grp = body.get("group_name", row["group_name"] if "group_name" in row.keys() else "") conn.execute("UPDATE collection_properties SET name=?, options_json=?, number_format=?, required=?, visible_in_views=?, validation_json=?, group_name=? WHERE id=?", (name, opts, nf, req, vis, vj, grp, prop_id)) conn.commit() audit_log(user, "property.update", "property", prop_id, "", request) return {"id": prop_id, "status": "updated"} @router.delete("/properties/{prop_id}") async def delete_property_v2(prop_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: if not conn.execute("SELECT id FROM collection_properties WHERE id=?", (prop_id,)).fetchone(): raise HTTPException(404, "Property not found") conn.execute("DELETE FROM collection_properties WHERE id=?", (prop_id,)) conn.commit() audit_log(user, "property.delete", "property", prop_id, "", request) return {"id": prop_id, "status": "deleted"} @router.post("/properties/{prop_id}/relation") async def create_relation_v2(prop_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} related_id = body.get("related_collection_id") reverse = (body.get("reverse_name") or "").strip() if not related_id: raise HTTPException(400, "related_collection_id required") with get_conn() as conn: row = conn.execute("SELECT collection_id FROM collection_properties WHERE id=?", (prop_id,)).fetchone() if not row: raise HTTPException(404, "Property not found") conn.execute("UPDATE collection_properties SET prop_type='relation', related_collection_id=?, reverse_name=? WHERE id=?", (related_id, reverse, prop_id)) if reverse: max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_properties WHERE collection_id=?", (related_id,)).fetchone()[0] try: conn.execute("INSERT INTO collection_properties (collection_id, name, prop_type, related_collection_id, reverse_name, position) VALUES (?, ?, 'relation', ?, ?, ?)", (related_id, reverse, row["collection_id"], "", max_pos)) except Exception: pass conn.commit() return {"id": prop_id, "status": "updated"} @router.post("/properties/evaluate-formula") async def evaluate_formula_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) try: body = await request.json() except Exception: body = {} expr = body.get("expression") or body.get("formula") if not expr: raise HTTPException(400, "expression required") ctx = body.get("context") or {} try: from app.services.formula_engine import FormulaEngine res = FormulaEngine().evaluate(expr, ctx) except Exception as e: raise HTTPException(400, f"Formula error: {e}") from None return {"result": res, "expression": expr} @router.post("/properties/compute-rollup") async def compute_rollup_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) try: body = await request.json() except Exception: body = {} for k in ("collection_id", "relation_property_id", "target_property_id", "page_id"): if k not in body: raise HTTPException(400, f"{k} required") try: from app.services.rollup_engine import RollupEngine res = RollupEngine().compute(body["collection_id"], body["relation_property_id"], body["target_property_id"], body["page_id"], body.get("function", "count")) except Exception as e: raise HTTPException(400, f"Rollup error: {e}") from None return {"result": res} # ── Views & dashboards ──────────────────────────────────────────────────── @router.get("/collections/{collection_id}/views") async def list_views_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT * FROM collection_views WHERE collection_id=? ORDER BY position", (collection_id,)).fetchall() return {"views": [row_to_dict(r) for r in rows]} @router.post("/collections/{collection_id}/views") async def create_view_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "New View").strip() vtype = body.get("view_type") or body.get("type") or "table" config = body.get("config") or body.get("config_json") or {} with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_views WHERE collection_id=?", (collection_id,)).fetchone()[0] cur = conn.execute("INSERT INTO collection_views (collection_id, name, view_type, config_json, position, created_by) VALUES (?, ?, ?, ?, ?, ?)", (collection_id, name, vtype, json.dumps(config), max_pos, user["id"])) vid = cur.lastrowid conn.commit() audit_log(user, "view.create", "view", vid, name, request) try: await _fire_event("collection.view.created", {"view_id": vid, "collection_id": collection_id, "name": name, "view_type": vtype}) except Exception: pass return {"id": vid, "name": name, "view_type": vtype, "status": "created"} @router.patch("/views/{view_id}") async def patch_view_v2(view_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM collection_views WHERE id=?", (view_id,)).fetchone() if not row: raise HTTPException(404, "View not found") cfg = json.loads(row["config_json"] or "{}") if "config" in body: cfg.update(body["config"]) elif "config_json" in body: try: cfg.update(json.loads(body["config_json"]) if isinstance(body["config_json"], str) else body["config_json"]) except Exception: pass # also flat keys for k in ("group_by", "sub_group_by", "wip_limits", "card_size", "cover_property", "cover_mode", "card_properties", "visible_properties", "filters", "sorts", "date_property"): if k in body: cfg[k] = body[k] name = body.get("name", row["name"]) vtype = body.get("view_type") or body.get("type") or row["view_type"] conn.execute("UPDATE collection_views SET name=?, view_type=?, config_json=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (name, vtype, json.dumps(cfg), view_id)) conn.commit() audit_log(user, "view.update", "view", view_id, "", request) return {"id": view_id, "status": "updated"} @router.delete("/views/{view_id}") async def delete_view_v2(view_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: if not conn.execute("SELECT id FROM collection_views WHERE id=?", (view_id,)).fetchone(): raise HTTPException(404, "View not found") conn.execute("DELETE FROM collection_views WHERE id=?", (view_id,)) conn.commit() audit_log(user, "view.delete", "view", view_id, "", request) return {"id": view_id, "status": "deleted"} @router.post("/views/{view_id}/save-as") async def save_as_view_v2(view_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "").strip() or "Copy" with get_conn() as conn: row = conn.execute("SELECT * FROM collection_views WHERE id=?", (view_id,)).fetchone() if not row: raise HTTPException(404, "View not found") max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_views WHERE collection_id=?", (row["collection_id"],)).fetchone()[0] cur = conn.execute("INSERT INTO collection_views (collection_id, name, view_type, config_json, position, created_by) VALUES (?, ?, ?, ?, ?, ?)", (row["collection_id"], name, row["view_type"], row["config_json"], max_pos, user["id"])) nid = cur.lastrowid conn.commit() return {"id": nid, "name": name, "status": "created"} @router.get("/collections/{collection_id}/dashboards") async def list_dashboards_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT * FROM collection_dashboards WHERE collection_id=? ORDER BY created_at", (collection_id,)).fetchall() return {"dashboards": [row_to_dict(r) for r in rows]} @router.post("/collections/{collection_id}/dashboards") async def create_dashboard_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "Dashboard").strip() layout = body.get("layout") or body.get("layout_json") or {"columns": 1, "widgets": []} with get_conn() as conn: cur = conn.execute("INSERT INTO collection_dashboards (collection_id, name, layout_json) VALUES (?, ?, ?)", (collection_id, name, json.dumps(layout))) did = cur.lastrowid conn.commit() return {"id": did, "name": name, "status": "created"} @router.patch("/dashboards/{dashboard_id}") async def patch_dashboard_v2(dashboard_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM collection_dashboards WHERE id=?", (dashboard_id,)).fetchone() if not row: raise HTTPException(404, "Dashboard not found") name = body.get("name", row["name"]) layout = body.get("layout") or body.get("layout_json") if layout is not None: layout_json = json.dumps(layout) else: layout_json = row["layout_json"] conn.execute("UPDATE collection_dashboards SET name=?, layout_json=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (name, layout_json, dashboard_id)) conn.commit() return {"id": dashboard_id, "status": "updated"} @router.delete("/dashboards/{dashboard_id}") async def delete_dashboard_v2(dashboard_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM collection_dashboards WHERE id=?", (dashboard_id,)) conn.commit() return {"id": dashboard_id, "status": "deleted"} # ── Comments & mentions ─────────────────────────────────────────────────── @router.get("/pages/{page_id}/comments") async def list_comments_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: total = conn.execute("SELECT COUNT(*) FROM comments WHERE target_id=? OR page_id=?", (page_id, page_id)).fetchone()[0] rows = conn.execute("SELECT c.*, u.login, u.full_name FROM comments c LEFT JOIN users u ON u.id=c.user_id WHERE c.target_id=? OR c.page_id=? ORDER BY c.created_at LIMIT ? OFFSET ?", (page_id, page_id, limit, offset)).fetchall() return {"comments": [row_to_dict(r) for r in rows], "total": total, "limit": limit, "offset": offset} @router.post("/pages/{page_id}/comments") async def create_comment_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} text = (body.get("body") or body.get("content") or "").strip() if not text: raise HTTPException(400, "body is required") with get_conn() as conn: cur = conn.execute("INSERT INTO comments (page_id, user_id, body, target_type, target_id, anchor_block_id, anchor_start, anchor_end) VALUES (?, ?, ?, 'page', ?, ?, ?, ?)", (page_id, user["id"], text, page_id, body.get("anchor_block_id"), body.get("anchor_start"), body.get("anchor_end"))) nid = cur.lastrowid conn.commit() row = conn.execute("SELECT * FROM comments WHERE id=?", (nid,)).fetchone() audit_log(user, "comment.create", "comment", nid, text[:80], request) try: await _fire_event("comment.added", {"comment_id": nid, "page_id": page_id, "user_id": user["id"]}) except Exception: pass return {"id": nid, "status": "created", "comment": row_to_dict(row)} @router.patch("/comments/{comment_id}") async def patch_comment_v2(comment_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM comments WHERE id=?", (comment_id,)).fetchone() if not row: raise HTTPException(404, "Comment not found") if row["user_id"] != user["id"] and not user.get("is_admin"): raise HTTPException(403, "Not your comment") body_text = body.get("body", row["body"]) resolved = body.get("resolved", row["resolved"]) was_resolved = int(row["resolved"] or 0) conn.execute("UPDATE comments SET body=?, resolved=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (body_text, int(bool(resolved)), comment_id)) conn.commit() if int(bool(resolved)) and not was_resolved: try: await _fire_event("comment.resolved", {"comment_id": comment_id, "page_id": row["page_id"]}) except Exception: pass return {"id": comment_id, "status": "updated"} @router.delete("/comments/{comment_id}") async def delete_comment_v2(comment_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT * FROM comments WHERE id=?", (comment_id,)).fetchone() if not row: raise HTTPException(404, "Comment not found") if row["user_id"] != user["id"] and not user.get("is_admin"): raise HTTPException(403, "Not your comment") conn.execute("DELETE FROM comments WHERE id=?", (comment_id,)) conn.commit() return {"id": comment_id, "status": "deleted"} @router.post("/pages/{page_id}/mentions") async def create_mention_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} targets = body.get("user_ids") or body.get("mentions") or [] if isinstance(targets, int): targets = [targets] if not targets: raise HTTPException(400, "user_ids required") created = 0 with get_conn() as conn: for uid in targets: try: conn.execute("INSERT INTO notifications (user_id, actor_id, ntype, title, message, resource_type, resource_id, url) VALUES (?, ?, 'mention', 'You were mentioned', ?, 'page', ?, ?)", (uid, user["id"], body.get("message") or f"Mentioned in page {page_id}", page_id, f"/pages/{page_id}")) created += 1 except Exception: pass conn.commit() if created: try: await _fire_event("mention.added", {"page_id": page_id, "user_ids": [u for u in targets if isinstance(u, int)], "count": created}) except Exception: pass return {"mentions": created, "status": "created"} # ── Notifications ───────────────────────────────────────────────────────── @router.get("/notifications") async def list_notifications_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) unread = request.query_params.get("unread") with get_conn() as conn: where = "user_id=?" params: list = [user["id"]] if unread == "1": where += " AND is_read=0" total = conn.execute(f"SELECT COUNT(*) FROM notifications WHERE {where}", params).fetchone()[0] rows = conn.execute(f"SELECT * FROM notifications WHERE {where} ORDER BY created_at DESC LIMIT ? OFFSET ?", (*params, limit, offset)).fetchall() return {"notifications": [row_to_dict(r) for r in rows], "total": total, "limit": limit, "offset": offset} @router.get("/notifications/unread-count") async def unread_count(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: cnt = conn.execute("SELECT COUNT(*) FROM notifications WHERE user_id=? AND is_read=0", (user["id"],)).fetchone()[0] return {"unread": cnt} @router.post("/notifications/{notif_id}/read") async def mark_read(notif_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: conn.execute("UPDATE notifications SET is_read=1 WHERE id=? AND user_id=?", (notif_id, user["id"])) conn.commit() return {"id": notif_id, "status": "read"} @router.post("/notifications/read-all") async def mark_all_read(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: conn.execute("UPDATE notifications SET is_read=1 WHERE user_id=?", (user["id"],)) conn.commit() return {"status": "all read"} @router.patch("/users/me/preferences") async def patch_prefs(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT notification_prefs FROM users WHERE id=?", (user["id"],)).fetchone() try: cur = json.loads(row["notification_prefs"] or "{}") except Exception: cur = {} cur.update(body) conn.execute("UPDATE users SET notification_prefs=? WHERE id=?", (json.dumps(cur), user["id"])) conn.commit() return {"preferences": cur} # ── Favorites / Tags / Recents ──────────────────────────────────────────── @router.get("/favorites") async def list_favorites_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT f.*, p.title, p.workspace_id FROM favorites f JOIN pages p ON p.id=f.page_id WHERE f.user_id=? ORDER BY f.position", (user["id"],)).fetchall() return {"favorites": [row_to_dict(r) for r in rows]} @router.post("/favorites") async def add_favorite_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} pid = body.get("page_id") if not pid: raise HTTPException(400, "page_id required") with get_conn() as conn: try: conn.execute("INSERT INTO favorites (user_id, page_id) VALUES (?, ?)", (user["id"], pid)) conn.commit() except Exception as e: raise HTTPException(409, str(e)) from None try: await _fire_event("favorite.added", {"page_id": pid, "user_id": user["id"]}) except Exception: pass return {"page_id": pid, "status": "added"} @router.delete("/favorites/{page_id}") async def remove_favorite_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM favorites WHERE user_id=? AND page_id=?", (user["id"], page_id)) conn.commit() try: await _fire_event("favorite.removed", {"page_id": page_id, "user_id": user["id"]}) except Exception: pass return {"page_id": page_id, "status": "removed"} @router.get("/tags") async def list_tags_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) q = (request.query_params.get("q") or "").strip() with get_conn() as conn: if q: rows = conn.execute("SELECT * FROM tags WHERE user_id=? AND name LIKE ? ORDER BY name", (user["id"], f"%{q}%")).fetchall() else: rows = conn.execute("SELECT * FROM tags WHERE user_id=? ORDER BY name", (user["id"],)).fetchall() return {"tags": [dict(r) for r in rows]} @router.post("/tags") async def create_tag_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "").strip() if not name: raise HTTPException(400, "name required") color = body.get("color", "#787774") with get_conn() as conn: try: cur = conn.execute("INSERT INTO tags (name, color, user_id) VALUES (?, ?, ?)", (name, color, user["id"])) tid = cur.lastrowid conn.commit() except Exception as e: raise HTTPException(409, str(e)) from None return {"id": tid, "name": name, "color": color, "status": "created"} @router.patch("/tags/{tag_id}") async def patch_tag_v2(tag_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM tags WHERE id=? AND user_id=?", (tag_id, user["id"])).fetchone() if not row: raise HTTPException(404, "Tag not found") name = body.get("name", row["name"]) color = body.get("color", row["color"]) conn.execute("UPDATE tags SET name=?, color=? WHERE id=?", (name, color, tag_id)) conn.commit() return {"id": tag_id, "status": "updated"} @router.delete("/tags/{tag_id}") async def delete_tag_v2(tag_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM tags WHERE id=? AND user_id=?", (tag_id, user["id"])) conn.commit() return {"id": tag_id, "status": "deleted"} @router.post("/pages/{page_id}/tags") async def attach_tag_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} tag_id = body.get("tag_id") if not tag_id: raise HTTPException(400, "tag_id required") with get_conn() as conn: try: conn.execute("INSERT INTO page_tags (page_id, tag_id) VALUES (?, ?)", (page_id, tag_id)) conn.commit() except Exception as e: raise HTTPException(409, str(e)) from None return {"page_id": page_id, "tag_id": tag_id, "status": "attached"} @router.delete("/pages/{page_id}/tags/{tag_id}") async def detach_tag_v2(page_id: int, tag_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM page_tags WHERE page_id=? AND tag_id=?", (page_id, tag_id)) conn.commit() return {"status": "detached"} @router.get("/recents") async def list_recents_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit = int(request.query_params.get("limit", "20")) st = request.query_params.get("source_type") with get_conn() as conn: if st: rows = conn.execute("SELECT * FROM recents WHERE user_id=? AND source_type=? ORDER BY accessed_at DESC LIMIT ?", (user["id"], st, limit)).fetchall() else: rows = conn.execute("SELECT * FROM recents WHERE user_id=? ORDER BY accessed_at DESC LIMIT ?", (user["id"], limit)).fetchall() return {"recents": [row_to_dict(r) for r in rows]} # ── Sharing / publish ───────────────────────────────────────────────────── @router.get("/pages/{page_id}/shares") async def list_shares_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT * FROM page_shares WHERE page_id=?", (page_id,)).fetchall() return {"shares": [dict(r) for r in rows]} @router.post("/pages/{page_id}/shares") async def create_share_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} perm = (body.get("permission") or "view").strip().lower() if perm not in ("view", "comment", "edit"): raise HTTPException(400, "Invalid permission. Use view, comment, or edit") email = (body.get("email") or "").strip() uid = body.get("user_id") if not email and not uid: raise HTTPException(400, "email or user_id required") with get_conn() as conn: cur = conn.execute("INSERT INTO page_shares (page_id, shared_with_user_id, shared_with_email, permission, created_by) VALUES (?, ?, ?, ?, ?)", (page_id, uid, email, perm, user["id"])) conn.execute("UPDATE pages SET is_shared=1 WHERE id=?", (page_id,)) conn.commit() nid = cur.lastrowid audit_log(user, "share.create", "share", nid, f"page={page_id}", request) try: await _fire_event("page.shared", {"page_id": page_id, "share_id": nid, "permission": perm}) except Exception: pass return {"id": nid, "page_id": page_id, "status": "shared"} @router.patch("/shares/{share_id}") async def patch_share_v2(share_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} perm = (body.get("permission") or "").strip().lower() if perm not in ("view", "comment", "edit"): raise HTTPException(400, "Invalid permission") with get_conn() as conn: if not conn.execute("SELECT id FROM page_shares WHERE id=?", (share_id,)).fetchone(): raise HTTPException(404, "Share not found") conn.execute("UPDATE page_shares SET permission=? WHERE id=?", (perm, share_id)) conn.commit() return {"id": share_id, "status": "updated"} @router.delete("/shares/{share_id}") async def delete_share_v2(share_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT page_id FROM page_shares WHERE id=?", (share_id,)).fetchone() if not row: raise HTTPException(404, "Share not found") conn.execute("DELETE FROM page_shares WHERE id=?", (share_id,)) # unset is_shared if no shares left cnt = conn.execute("SELECT COUNT(*) FROM page_shares WHERE page_id=?", (row["page_id"],)).fetchone()[0] if cnt == 0: conn.execute("UPDATE pages SET is_shared=0 WHERE id=?", (row["page_id"],)) conn.commit() return {"id": share_id, "status": "revoked"} @router.post("/pages/{page_id}/publish") async def publish_page_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} slug = (body.get("slug") or body.get("publish_slug") or f"p-{page_id}-{secrets.token_urlsafe(6)}").strip() with get_conn() as conn: if not conn.execute("SELECT id FROM pages WHERE id=?", (page_id,)).fetchone(): raise HTTPException(404, "Page not found") conn.execute("UPDATE pages SET is_published=1, publish_slug=?, is_shared=1 WHERE id=?", (slug, page_id)) conn.commit() audit_log(user, "page.publish", "page", page_id, slug, request) try: await _fire_event("page.published", {"page_id": page_id, "slug": slug}) except Exception: pass return {"page_id": page_id, "slug": slug, "url": f"/p/{slug}", "status": "published"} @router.delete("/pages/{page_id}/publish") async def unpublish_page_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("UPDATE pages SET is_published=0 WHERE id=?", (page_id,)) conn.commit() try: await _fire_event("page.unpublished", {"page_id": page_id}) except Exception: pass return {"page_id": page_id, "status": "unpublished"} # ── History ─────────────────────────────────────────────────────────────── @router.get("/pages/{page_id}/history") async def list_history_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT h.*, u.login FROM page_history h LEFT JOIN users u ON u.id=h.user_id WHERE h.page_id=? ORDER BY h.created_at DESC", (page_id,)).fetchall() # also page_versions for block pages vrows = conn.execute("SELECT * FROM page_versions WHERE page_id=? ORDER BY created_at DESC", (page_id,)).fetchall() return {"history": [row_to_dict(r) for r in rows], "versions": [row_to_dict(r) for r in vrows]} @router.post("/pages/{page_id}/history/restore") async def restore_history_v2(page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} hid = body.get("history_id") or body.get("id") or body.get("version_id") if not hid: raise HTTPException(400, "history_id required") with get_conn() as conn: h = conn.execute("SELECT * FROM page_versions WHERE id=? AND page_id=?", (hid, page_id)).fetchone() if h: conn.execute("UPDATE pages SET content=?, title=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", (h["blocks_json"], h["title"], page_id)) conn.commit() return {"page_id": page_id, "restored_version": hid, "status": "restored"} h2 = conn.execute("SELECT * FROM page_history WHERE id=? AND page_id=?", (hid, page_id)).fetchone() if h2: try: snap = json.loads(h2["snapshot_json"] or "{}") except Exception: snap = {} # best-effort restore content if snap.get("content"): conn.execute("UPDATE pages SET content=? WHERE id=?", (snap["content"], page_id)) conn.commit() return {"page_id": page_id, "restored_version": hid, "status": "restored"} raise HTTPException(404, "History not found") # ── Sprints ─────────────────────────────────────────────────────────────── @router.get("/collections/{collection_id}/sprints") async def list_sprints_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: total = conn.execute("SELECT COUNT(*) FROM sprints WHERE collection_id=?", (collection_id,)).fetchone()[0] rows = conn.execute("SELECT * FROM sprints WHERE collection_id=? ORDER BY created_at DESC LIMIT ? OFFSET ?", (collection_id, limit, offset)).fetchall() return {"sprints": [row_to_dict(r) for r in rows], "total": total, "limit": limit, "offset": offset} @router.post("/collections/{collection_id}/sprints") async def create_sprint_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "Sprint").strip() start = body.get("start_date") or body.get("start") or "" end = body.get("end_date") or body.get("end") or "" if not start or not end: raise HTTPException(400, "start_date and end_date required (YYYY-MM-DD)") with get_conn() as conn: cur = conn.execute("INSERT INTO sprints (collection_id, name, start_date, end_date, goal) VALUES (?, ?, ?, ?, ?)", (collection_id, name, start, end, body.get("goal") or "")) sid = cur.lastrowid conn.commit() try: await _fire_event("sprint.created", {"sprint_id": sid, "collection_id": collection_id, "name": name}) except Exception: pass return {"id": sid, "name": name, "status": "created"} @router.patch("/sprints/{sprint_id}") async def patch_sprint_v2(sprint_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM sprints WHERE id=?", (sprint_id,)).fetchone() if not row: raise HTTPException(404, "Sprint not found") name = body.get("name", row["name"]) start = body.get("start_date", row["start_date"]) end = body.get("end_date", row["end_date"]) goal = body.get("goal", row["goal"]) status = body.get("status", row["status"]) conn.execute("UPDATE sprints SET name=?, start_date=?, end_date=?, goal=?, status=? WHERE id=?", (name, start, end, goal, status, sprint_id)) conn.commit() try: await _fire_event("sprint.updated", {"sprint_id": sprint_id, "collection_id": row["collection_id"], "name": name, "status": status}) except Exception: pass return {"id": sprint_id, "status": "updated"} @router.delete("/sprints/{sprint_id}") async def delete_sprint_v2(sprint_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM sprints WHERE id=?", (sprint_id,)) conn.commit() return {"id": sprint_id, "status": "deleted"} @router.post("/sprints/{sprint_id}/assign") async def assign_sprint_v2(sprint_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} pid = body.get("page_id") if not pid: raise HTTPException(400, "page_id required") with get_conn() as conn: try: conn.execute("INSERT INTO sprint_pages (sprint_id, page_id, velocity_points) VALUES (?, ?, ?)", (sprint_id, pid, body.get("velocity_points", 1))) conn.commit() except Exception as e: raise HTTPException(409, str(e)) from None return {"sprint_id": sprint_id, "page_id": pid, "status": "assigned"} @router.delete("/sprints/{sprint_id}/assign/{page_id}") async def unassign_sprint_v2(sprint_id: int, page_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM sprint_pages WHERE sprint_id=? AND page_id=?", (sprint_id, page_id)) conn.commit() return {"status": "removed"} @router.get("/sprints/{sprint_id}/burndown") async def burndown_v2(sprint_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: s = conn.execute("SELECT * FROM sprints WHERE id=?", (sprint_id,)).fetchone() if not s: raise HTTPException(404, "Sprint not found") pages = conn.execute("SELECT sp.*, cp.property_values_json FROM sprint_pages sp JOIN collection_pages cp ON cp.id=sp.page_id WHERE sp.sprint_id=?", (sprint_id,)).fetchall() total = len(pages) # crude: completed where status property == Done (best-effort) completed = 0 for p in pages: try: pv = json.loads(p["property_values_json"] or "{}") for v in pv.values(): if str(v).lower() in ("done", "completed", "terminé"): completed += 1 break except Exception: pass remaining = total - completed # ideal linear ideal = [round(total * (1 - i / 10)) for i in range(11)] return {"sprint_id": sprint_id, "total": total, "completed": completed, "remaining": remaining, "ideal": ideal} # ── Templates ───────────────────────────────────────────────────────────── @router.get("/collections/{collection_id}/templates") async def list_templates_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT * FROM page_templates WHERE collection_id=? ORDER BY created_at", (collection_id,)).fetchall() return {"templates": [row_to_dict(r) for r in rows]} @router.post("/collections/{collection_id}/templates") async def create_template_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} name = (body.get("name") or "Template").strip() pv = json.dumps(body.get("property_values") or body.get("property_values_json") or {}) cj = json.dumps(body.get("content") or body.get("content_json") or []) with get_conn() as conn: cur = conn.execute("INSERT INTO page_templates (collection_id, name, property_values_json, content_json) VALUES (?, ?, ?, ?)", (collection_id, name, pv, cj)) tid = cur.lastrowid conn.commit() return {"id": tid, "name": name, "status": "created"} @router.patch("/templates/{template_id}") async def patch_template_v2(template_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM page_templates WHERE id=?", (template_id,)).fetchone() if not row: raise HTTPException(404, "Template not found") name = body.get("name", row["name"]) pv = json.dumps(body.get("property_values", json.loads(row["property_values_json"] or "{}"))) if "property_values" in body else row["property_values_json"] cj = json.dumps(body.get("content", json.loads(row["content_json"] or "[]"))) if "content" in body else row["content_json"] conn.execute("UPDATE page_templates SET name=?, property_values_json=?, content_json=? WHERE id=?", (name, pv, cj, template_id)) conn.commit() return {"id": template_id, "status": "updated"} @router.delete("/templates/{template_id}") async def delete_template_v2(template_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM page_templates WHERE id=?", (template_id,)) conn.commit() return {"id": template_id, "status": "deleted"} @router.post("/templates/{template_id}/apply") async def apply_template_v2(template_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: tpl = conn.execute("SELECT * FROM page_templates WHERE id=?", (template_id,)).fetchone() if not tpl: raise HTTPException(404, "Template not found") max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_pages WHERE collection_id=?", (tpl["collection_id"],)).fetchone()[0] pv = tpl["property_values_json"] or "{}" cur = conn.execute("INSERT INTO collection_pages (collection_id, title, position, property_values_json) VALUES (?, ?, ?, ?)", (tpl["collection_id"], tpl["name"], max_pos, pv)) pid = cur.lastrowid conn.commit() return {"template_id": template_id, "page_id": pid, "status": "applied"} @router.get("/templates/database") async def list_db_templates_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT * FROM database_templates ORDER BY name").fetchall() return {"templates": [row_to_dict(r) for r in rows]} @router.post("/templates/database/{template_id}/apply") async def apply_db_template_v2(template_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: tpl = conn.execute("SELECT * FROM database_templates WHERE id=?", (template_id,)).fetchone() if not tpl: raise HTTPException(404, "Template not found") name = (body.get("name") or tpl["name"]).strip() schema = json.loads(tpl["schema_json"] or "[]") cur = conn.execute("INSERT INTO collections (name, description, icon, schema_json, workspace_id, created_by) VALUES (?, ?, ?, ?, ?, ?)", (name, tpl["description"], tpl["icon"] if "icon" in tpl.keys() else "📋", json.dumps(schema), body.get("workspace_id"), user["id"])) cid = cur.lastrowid try: from app.services.db_templates import materialize_properties materialize_properties(conn, cid, schema) except Exception: pass conn.commit() return {"collection_id": cid, "name": name, "status": "created"} # ── Export / Import ─────────────────────────────────────────────────────── @router.get("/pages/{page_id}/export") async def export_page_v2(page_id: int, request: Request, format: str = "markdown", authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) fmt = (format or request.query_params.get("format") or "markdown").lower() if fmt not in ("markdown", "html", "pdf"): raise HTTPException(400, "format must be markdown, html or pdf") with get_conn() as conn: row = conn.execute("SELECT * FROM pages WHERE id=?", (page_id,)).fetchone() if not row: # try collection_pages row2 = conn.execute("SELECT * FROM collection_pages WHERE id=?", (page_id,)).fetchone() if not row2: raise HTTPException(404, "Page not found") # collection pages: return JSON return {"page": row_to_dict(row2), "format": fmt} # block pages: delegate to export service from app.services.export import export_page as _export try: data, mime, fname = _export(row, fmt) # type: ignore return Response(content=data, media_type=mime, headers={"Content-Disposition": f'attachment; filename="{fname}"'}) except Exception as e: raise HTTPException(500, f"Export failed: {e}") from None @router.get("/collections/{collection_id}/export/csv") async def export_csv_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) import csv import io with get_conn() as conn: if not conn.execute("SELECT id FROM collections WHERE id=?", (collection_id,)).fetchone(): raise HTTPException(404, "Collection not found") props = [dict(r) for r in conn.execute("SELECT * FROM collection_properties WHERE collection_id=? ORDER BY position", (collection_id,)).fetchall()] rows = conn.execute("SELECT * FROM collection_pages WHERE collection_id=? ORDER BY position", (collection_id,)).fetchall() out = io.StringIO() writer = csv.writer(out) header = ["Title"] + [p["name"] for p in props] writer.writerow(header) for r in rows: try: pv = json.loads(r["property_values_json"] or "{}") except Exception: pv = {} vals = [r["title"]] for p in props: vals.append(str(pv.get(str(p["id"])) or pv.get(p["name"]) or "")) writer.writerow(vals) return Response(content=out.getvalue().encode("utf-8"), media_type="text/csv", headers={"Content-Disposition": f'attachment; filename="collection-{collection_id}.csv"'}) @router.post("/collections/{collection_id}/import/csv") async def import_csv_v2(collection_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: form = await request.form() file = form.get("file") data = await file.read() if file else b"" text = data.decode("utf-8", errors="ignore") except Exception as err: raise HTTPException(400, "file required (multipart)") from err import csv import io reader = csv.DictReader(io.StringIO(text)) created = 0 with get_conn() as conn: for row in reader: title = row.get("Title") or row.get("title") or "Untitled" # map remaining columns to property names pv = {} # resolve prop name -> id props = {p["name"]: p["id"] for p in conn.execute("SELECT id, name FROM collection_properties WHERE collection_id=?", (collection_id,)).fetchall()} for k, v in row.items(): if k in ("Title", "title"): continue pid = props.get(k) if pid: pv[str(pid)] = v max_pos = conn.execute("SELECT COALESCE(MAX(position), -1)+1 FROM collection_pages WHERE collection_id=?", (collection_id,)).fetchone()[0] conn.execute("INSERT INTO collection_pages (collection_id, title, position, property_values_json) VALUES (?, ?, ?, ?)", (collection_id, title, max_pos, json.dumps(pv))) created += 1 conn.commit() return {"imported": created, "status": "ok"} # ── Forges / Projects ───────────────────────────────────────────────────── @router.get("/projects") async def list_projects_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) with get_conn() as conn: rows = conn.execute("SELECT * FROM projects ORDER BY proj_type, owner, name").fetchall() return {"projects": [row_to_dict(r) for r in rows]} @router.get("/projects/{owner}/{repo}/tree") async def project_tree_v2(owner: str, repo: str, request: Request, path: str = "", authorization: str | None = Header(default=None)): get_bearer_user(request, authorization) # proxy to gitea client? Return placeholder listing from projects table with get_conn() as conn: proj = conn.execute("SELECT * FROM projects WHERE owner=? AND name=?", (owner, repo)).fetchone() if not proj: raise HTTPException(404, "Project not found") # delegate to gitea API if available (best-effort) try: from app.services.gitea_client import gitea tree = await gitea.list_repo_files(owner, repo, path or "") return {"owner": owner, "repo": repo, "path": path, "tree": tree} except Exception: return {"owner": owner, "repo": repo, "path": path, "tree": []} # ── Search (FTS5 + LIKE fallback) ───────────────────────────────────────── @router.get("/search") async def search_v2(request: Request, query: str = "", workspace_id: int | None = None, type: str = "all", authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) q = (query or request.query_params.get("query") or "").strip() if not q: return {"results": [], "query": q} like = f"%{q}%" with get_conn() as conn: pages = [] # try FTS5 try: rows = conn.execute("SELECT p.id, p.title, p.content, p.workspace_id, snippet(pages_fts, -1, '', '', '...', 32) as snippet FROM pages_fts f JOIN pages p ON p.id=f.rowid WHERE pages_fts MATCH ? LIMIT 20", (q,)).fetchall() pages = [{"type": "page", "id": r["id"], "title": r["title"], "snippet": r["snippet"]} for r in rows] except Exception: rows = conn.execute("SELECT id, title FROM pages WHERE title LIKE ? OR content LIKE ? LIMIT 20", (like, like)).fetchall() pages = [{"type": "page", "id": r["id"], "title": r["title"]} for r in rows] # collections colls = conn.execute("SELECT id, name FROM collections WHERE name LIKE ? LIMIT 10", (like,)).fetchall() results = pages + [{"type": "collection", "id": r["id"], "title": r["name"]} for r in colls] return {"query": q, "results": results} # ── Admin ───────────────────────────────────────────────────────────────── @router.get("/admin/users") async def admin_list_users_v2(request: Request, limit: int = 30, offset: int = 0, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not user.get("is_admin") and not has_scope(user.get("_token_scopes"), "admin"): raise HTTPException(403, "Admin scope required") limit = max(1, min(limit, 100)) with get_conn() as conn: total = conn.execute("SELECT COUNT(*) FROM users").fetchone()[0] rows = conn.execute("SELECT id, login, full_name, email, is_admin, is_active, created_at FROM users ORDER BY id LIMIT ? OFFSET ?", (limit, offset)).fetchall() return {"users": [row_to_dict(r) for r in rows], "total": total, "limit": limit, "offset": offset} @router.patch("/admin/users/{uid}") async def admin_patch_user_v2(uid: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not user.get("is_admin") and not has_scope(user.get("_token_scopes"), "admin"): raise HTTPException(403, "Admin scope required") try: body = await request.json() except Exception: body = {} with get_conn() as conn: if not conn.execute("SELECT id FROM users WHERE id=?", (uid,)).fetchone(): raise HTTPException(404, "User not found") sets = [] params: list = [] for k in ("is_active", "is_admin", "full_name", "email"): if k in body: sets.append(f"{k}=?") params.append(int(body[k]) if k in ("is_active", "is_admin") else body[k]) if not sets: raise HTTPException(400, "No fields") params.append(uid) conn.execute(f"UPDATE users SET {', '.join(sets)} WHERE id=?", params) conn.commit() audit_log(user, "admin.user_update", "user", uid, "", request) return {"id": uid, "status": "updated"} @router.delete("/admin/users/{uid}") async def admin_delete_user_v2(uid: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not user.get("is_admin") and not has_scope(user.get("_token_scopes"), "admin"): raise HTTPException(403, "Admin scope required") if uid == user["id"]: raise HTTPException(400, "Cannot delete yourself") with get_conn() as conn: conn.execute("DELETE FROM users WHERE id=?", (uid,)) conn.commit() audit_log(user, "admin.user_delete", "user", uid, "", request) return {"id": uid, "status": "deleted"} @router.get("/admin/audit-logs") async def admin_audit_logs_v2(request: Request, limit: int = 50, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not user.get("is_admin") and not has_scope(user.get("_token_scopes"), "admin"): raise HTTPException(403, "Admin scope required") limit = max(1, min(limit, 200)) with get_conn() as conn: rows = conn.execute("SELECT * FROM api_audit_log ORDER BY created_at DESC LIMIT ?", (limit,)).fetchall() return {"logs": [row_to_dict(r) for r in rows]} # ── Webhooks (v6.4.0 prod) ──────────────────────────────────────────────── @router.get("/webhooks/events") async def list_webhook_events_v2(request: Request, authorization: str | None = Header(default=None)): """Catalogue of deliverable events (+ wildcard syntax).""" user = get_bearer_user(request, authorization) _v2_rate_check(request, user) from app.services.webhook_outbound import EVENTS return {"events": EVENTS, "wildcards": ["*", "page.*", "collection.*"]} @router.get("/webhooks") async def list_webhooks_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: total = conn.execute("SELECT COUNT(*) AS n FROM webhook_subscriptions").fetchone()["n"] rows = conn.execute( "SELECT * FROM webhook_subscriptions ORDER BY created_at DESC LIMIT ? OFFSET ?", (limit, offset), ).fetchall() return JSONResponse( content={"webhooks": [dict(r) for r in rows]}, headers=paginate_headers(total), ) @router.post("/webhooks") async def create_webhook_v2(request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} from app.services.webhook_outbound import EVENTS, _event_matches url = (body.get("url") or "").strip() event = (body.get("event") or "page.created").strip() secret = (body.get("secret") or "").strip() if not url or not url.startswith("http"): raise HTTPException(400, "url must start with http") if not event or (event not in EVENTS and not (event.endswith(".*") or event in ("*", "all"))): raise HTTPException(400, f"Unknown event '{event}'. See GET /api/v2/webhooks/events") # make sure the pattern matches at least one known event if not any(_event_matches(event, e) for e in EVENTS): raise HTTPException(400, f"Event pattern '{event}' matches no known event") with get_conn() as conn: cur = conn.execute("INSERT INTO webhook_subscriptions (url, event, secret) VALUES (?, ?, ?)", (url, event, secret)) wid = cur.lastrowid conn.commit() row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (wid,)).fetchone() audit_log(user, "webhook.create", "webhook", wid, url, request) return {"id": wid, "status": "created", "webhook": dict(row) if row else {}, "signature_header": "X-FlowDeck-Signature (HMAC-SHA256, sha256=)" if secret else None} @router.patch("/webhooks/{webhook_id}") async def patch_webhook_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") try: body = await request.json() except Exception: body = {} with get_conn() as conn: row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (webhook_id,)).fetchone() if not row: raise HTTPException(404, "Webhook not found") url = body.get("url", row["url"]) event = body.get("event", row["event"]) secret = body.get("secret", row["secret"]) active = int(bool(body.get("active", row["active"]))) conn.execute("UPDATE webhook_subscriptions SET url=?, event=?, secret=?, active=? WHERE id=?", (url, event, secret, active, webhook_id)) conn.commit() audit_log(user, "webhook.update", "webhook", webhook_id, "", request) return {"id": webhook_id, "status": "updated"} @router.delete("/webhooks/{webhook_id}") async def delete_webhook_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: conn.execute("DELETE FROM webhook_subscriptions WHERE id=?", (webhook_id,)) conn.commit() audit_log(user, "webhook.delete", "webhook", webhook_id, "", request) return {"id": webhook_id, "status": "deleted"} @router.post("/webhooks/{webhook_id}/test") async def test_webhook_v2(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") with get_conn() as conn: row = conn.execute("SELECT * FROM webhook_subscriptions WHERE id=?", (webhook_id,)).fetchone() if not row: raise HTTPException(404, "Webhook not found") # live delivery via the prod dispatcher (HMAC + retry + journal), # direct to this subscription only (no wildcard fan-out) from app.services.webhook_outbound import deliver_to_sub ok = await deliver_to_sub(webhook_id, row["url"], "ping", {"webhook_id": webhook_id, "test": True}, row["secret"] or "") audit_log(user, "webhook.test", "webhook", webhook_id, f"ok={ok}", request) return {"webhook_id": webhook_id, "status": "tested", "delivered": ok} @router.get("/webhooks/{webhook_id}/deliveries") async def list_deliveries_v2(webhook_id: int, request: Request, status: str | None = None, authorization: str | None = Header(default=None)): user = get_bearer_user(request, authorization) _v2_rate_check(request, user) limit, offset = parse_pagination(request) with get_conn() as conn: if status: total = conn.execute( "SELECT COUNT(*) AS n FROM webhook_deliveries WHERE webhook_id=? AND status=?", (webhook_id, status)).fetchone()["n"] rows = conn.execute( "SELECT * FROM webhook_deliveries WHERE webhook_id=? AND status=? " "ORDER BY created_at DESC LIMIT ? OFFSET ?", (webhook_id, status, limit, offset)).fetchall() else: total = conn.execute( "SELECT COUNT(*) AS n FROM webhook_deliveries WHERE webhook_id=?", (webhook_id,)).fetchone()["n"] rows = conn.execute( "SELECT * FROM webhook_deliveries WHERE webhook_id=? " "ORDER BY created_at DESC LIMIT ? OFFSET ?", (webhook_id, limit, offset)).fetchall() return JSONResponse( content={"deliveries": [row_to_dict(r) for r in rows]}, headers=paginate_headers(total), ) @router.post("/webhooks/{webhook_id}/retry") async def retry_webhook_deliveries(webhook_id: int, request: Request, authorization: str | None = Header(default=None)): """Manually retry failed deliveries for a webhook.""" user = get_bearer_user(request, authorization) _v2_rate_check(request, user) if not has_scope(user.get("_token_scopes"), "write"): raise HTTPException(403, "Insufficient scope. Required: write") from app.services.webhook_outbound import retry_due_deliveries with get_conn() as conn: # Force retry by setting next_retry_at to the past conn.execute( """UPDATE webhook_deliveries SET next_retry_at = strftime('%s', 'now', '-1 second') WHERE webhook_id = ? AND status = 'retrying'""", (webhook_id,), ) conn.commit() retried = await retry_due_deliveries() audit_log(user, "webhook.retry", "webhook", webhook_id, f"retried={retried}", request) return {"webhook_id": webhook_id, "status": "retried", "retried_count": retried} @router.post("/webhooks/verify-signature") async def verify_webhook_signature(request: Request, authorization: str | None = Header(default=None)): """Verify a webhook signature (for debugging/testing).""" user = get_bearer_user(request, authorization) _v2_rate_check(request, user) from app.services.webhook_outbound import verify_signature try: body = await request.json() except Exception: body = {} secret = body.get("secret", "") payload = body.get("payload", "{}") signature = body.get("signature", "") is_valid = verify_signature(secret, payload.encode(), signature) audit_log(user, "webhook.signature_verify", "webhook", 0, f"valid={is_valid}", request) return {"valid": is_valid, "secret": secret[:10] + "..." if len(secret) > 10 else secret}