Files
flowdeck/app/routers/api_v2.py
T
bruno 7be96f0618
FlowDeck CI / lint (push) Canceled after 0s
FlowDeck CI / test (push) Canceled after 0s
FlowDeck CI / docker (push) Canceled after 0s
fix: A29 + A42(partiel) — publish partagé, fuite password_hash, data_dir (v7.5.0)
- A29 — `app/services/publish.py` : slugify titré unique (fallback aléatoire),
  404 si la page n'existe pas, événements centralisés. Les 3 paires
  publish/unpublish déléguent (sharing = front, board, v2) :
  · board : mise à jour aveugle → 404 + contrôle de session ajouté
  · board : perd `share_mode='anyone'` en bonus, v2 : perd `is_shared=1` —
    le share dialog reste l'unique propriétaire de ces drapeaux
  · v2 : slug fourni conservé, slug vidé aussi à la dépublication (avant : laissé)
  · `/users/me` ×2 et listings collections ×3 = contrats versionnés distincts,
    décision documentée (on garde)
- Byproduct sécurité — `GET /api/users/me` (v1) et le contexte de `/accounts`
  faisaient `SELECT *` sur users → password_hash / login_attempts / locked_until
  exposés → colonnes whitelistées (liste v2)
- A42 (partiel) — 9 copies de `Path(os.environ.get("FLOWDECK_DATA_DIR", "/data"))`
  → `settings.data_dir` (property : lecture à chaque accès, les tests
  monkeypatchent l'env) ; cache Gitea : évacuation des entrées expirées à chaque
  écriture. Reste : client httpx partagé (52 créations, cache par event loop)

tests : test_publish_service_shared_and_safe, test_users_me_no_secret_columns,
test_gitea_cache_evicts_expired

suite **1034/1034** · `ruff check app tests` OK · OpenAPI 511 chemins / 7.5.0
docs (ROADMAP/CHANGELOG/WORKLOAD/VERSION) à jour
2026-10-01 10:01:39 -04:00

2298 lines
118 KiB
Python

"""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 ( # noqa: F401 — require_scope est utilisé par les handlers
audit_log,
check_idempotency,
check_v2_rate_limit,
get_bearer_user,
has_scope,
paginate_headers,
parse_pagination,
require_scope,
row_to_dict,
store_idempotency,
to_iso8601,
validate_scopes_input,
)
from app.services.automations import fire_event as _fire_event
from app.services.publish import fire_published, fire_unpublished, publish, unpublish
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 <token>"}
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("get_me")
# never expose secrets
return d
@router.patch("/users/me")
async def patch_me(request: Request, authorization: str | None = Header(default=None)):
user = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_workspace")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(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")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 — A25 : PAS de try ici,
# une exception doit interrompre la transaction (sinon la collection est
# commitée sans son schéma et l'erreur disparaît).
from app.services.db_templates import materialize_properties
materialize_properties(conn, cid, schema)
# default view
conn.execute("INSERT INTO collection_views (collection_id, name, view_type, config_json) VALUES (?,?,?,?)", (cid, "Default View", "table", json.dumps({"visible_properties": ["Title"]})))
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(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")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_linked_db")
# 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:
logger.exception("create_linked_db")
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:
logger.exception("create_linked_db")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_collection_page_v2")
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:
logger.exception("create_collection_page_v2")
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:
logger.exception("create_collection_page_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("patch_page_v2")
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:
logger.exception("patch_page_v2")
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 = require_scope("write")(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:
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:
logger.exception("delete_page_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
# 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:
logger.exception("restore_page_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_relation_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_view_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("patch_view_v2")
# 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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_comment_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("patch_comment_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_mention_v2")
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:
logger.exception("create_mention_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("add_favorite_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("remove_favorite_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_share_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
try:
body = await request.json()
except Exception:
body = {}
slug_in = (body.get("slug") or body.get("publish_slug") or "").strip() or None
slug, _title = publish(page_id, explicit_slug=slug_in)
audit_log(user, "page.publish", "page", page_id, slug, request)
await fire_published(page_id, slug)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
unpublish(page_id)
await fire_unpublished(page_id)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("create_sprint_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("patch_sprint_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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:
logger.exception("burndown_v2")
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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
# A25 : pas de try — un échec de matérialisation doit interrompre la
# transaction plutôt que de commiter une collection sans schéma.
from app.services.db_templates import materialize_properties
materialize_properties(conn, cid, schema)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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, '<mark>', '</mark>', '...', 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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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=<hex>)" 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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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 = require_scope("write")(request, authorization)
_v2_rate_check(request, user)
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}