- Conversion `async def` → `def` de TOUTES les routes dont le corps ne contient ni `await`, ni `async with`, ni `async for`, ni `asyncio` (scan automatique corps par corps sur app/ : 352 converties, 0 dangereuses, vérifié `asyncio`/`run_coroutine`/`.result()` absents). FastAPI exécute ces handlers dans son threadpool → tout leur SQLite (`get_conn()` + `conn.execute`) quitte l'event loop, sans changer une ligne de logique. - Répartition : api_v2 60, dashboard 40, collections 25, board 23, workspace 19, wiki 17, permissions 14, api 14, main.py 6, + 35 fichiers. - Les 4 routers prioritaires de l'audit sont couverts par ce lot : api_v2 60 + dashboard 40 + collections 25 + board 23 = 148 conversions (le reste de leurs routes attend la phase 2 : elles ont de vrais `await`). - Reste (phase 2) : les 311 routes avec de vrais `await` → enrouler les blocs DB dans `await anyio.to_thread.run_sync(...)` ; pas de wrapper partagé livré (rien ne l'appellerait — YAGNI jusqu'au premier usage). suite **1037/1037** (233 s) · `ruff check app tests` OK · docs à jour
109 lines
4.3 KiB
Python
109 lines
4.3 KiB
Python
"""FlowDeck — /api/v2/sync endpoints (v6.0.0 PWA offline sync, Bearer v6.4.0).
|
|
|
|
Auth: ``Authorization: Bearer <token>`` (scopes ``read`` for delta/status,
|
|
``write`` for batch). The legacy ``flowdeck_session`` cookie is still accepted
|
|
as a fallback so the installed PWA/service worker keeps syncing.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
|
|
from fastapi import APIRouter, Header, HTTPException, Query, Request
|
|
from fastapi.responses import JSONResponse
|
|
|
|
from app.auth.session import SessionManager
|
|
from app.services.api_v2_helpers import get_bearer_user, has_scope
|
|
from app.services.sync_engine import SyncEngine
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
router = APIRouter(prefix="/api/v2/sync", tags=["sync"])
|
|
|
|
_engine = SyncEngine()
|
|
|
|
|
|
def _user(request: Request, authorization: str | None = None,
|
|
*, required_scope: str = "read") -> dict:
|
|
"""Bearer-first auth with session-cookie fallback (offline.js compat)."""
|
|
auth = authorization or request.headers.get("authorization") or ""
|
|
if auth and auth.lower().startswith("bearer "):
|
|
try:
|
|
user = get_bearer_user(request, authorization)
|
|
except HTTPException:
|
|
raise HTTPException(
|
|
status_code=401, detail="Invalid or expired API token"
|
|
) from None
|
|
if not has_scope(user.get("_token_scopes"), required_scope):
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail=f"Insufficient scope. Required: {required_scope}",
|
|
)
|
|
return user
|
|
user = SessionManager.decode_session(request.cookies.get("flowdeck_session", ""))
|
|
if not user:
|
|
raise HTTPException(status_code=401, detail="Authentication required")
|
|
return user
|
|
|
|
|
|
@router.get("/delta")
|
|
async def sync_delta(
|
|
request: Request,
|
|
since: float = Query(default=0, description="Epoch seconds (ou ms) du dernier sync"),
|
|
workspace_id: int = Query(default=None),
|
|
authorization: str | None = Header(default=None),
|
|
):
|
|
"""Pull server-side changes since `since` (for the given workspace)."""
|
|
user = _user(request, authorization, required_scope="read")
|
|
if workspace_id is None:
|
|
raise HTTPException(status_code=400, detail="workspace_id is required")
|
|
result = await _engine.get_delta(user["id"], since, workspace_id)
|
|
if result.get("error") == "forbidden":
|
|
return JSONResponse({"detail": "Forbidden"}, status_code=403)
|
|
return result
|
|
|
|
|
|
@router.post("/batch")
|
|
async def sync_batch(request: Request, authorization: str | None = Header(default=None)):
|
|
"""Apply a batch of offline mutations and return per-mutation results."""
|
|
user = _user(request, authorization, required_scope="write")
|
|
try:
|
|
body = await request.json()
|
|
except Exception:
|
|
raise HTTPException(status_code=400, detail="Invalid JSON body") from None
|
|
|
|
mutations = body.get("mutations") or []
|
|
device_id = body.get("device_id") or "unknown"
|
|
if not isinstance(mutations, list) or not mutations:
|
|
return {"results": [], "conflicts": [], "server_time": SyncEngine._now_epoch()}
|
|
|
|
result = await _engine.apply_batch(user["id"], mutations, device_id)
|
|
result["server_time"] = SyncEngine._now_epoch()
|
|
return result
|
|
|
|
|
|
@router.get("/status")
|
|
def sync_status(request: Request, workspace_id: int = Query(default=None),
|
|
authorization: str | None = Header(default=None)):
|
|
"""Synchronization status for the workspace (pending server queue, last sync)."""
|
|
user = _user(request, authorization, required_scope="read")
|
|
from app.db import get_conn
|
|
with get_conn() as conn:
|
|
if not SyncEngine._can_access(conn, user["id"], workspace_id):
|
|
return JSONResponse({"detail": "Forbidden"}, status_code=403)
|
|
pending = conn.execute(
|
|
"SELECT COUNT(*) AS n FROM offline_sync_queue WHERE user_id=? AND status='pending'",
|
|
(user["id"],),
|
|
).fetchone()["n"]
|
|
last = conn.execute(
|
|
"SELECT MAX(created_at) AS last FROM offline_sync_queue "
|
|
"WHERE user_id=? AND status='synced'",
|
|
(user["id"],),
|
|
).fetchone()["last"]
|
|
return {
|
|
"pending_count": pending,
|
|
"last_sync": last,
|
|
"is_syncing": False,
|
|
"server_time": SyncEngine._now_epoch(),
|
|
"workspace_id": workspace_id,
|
|
}
|