"""FlowDeck — /api/v2/sync endpoints (v6.0.0 PWA offline sync, Bearer v6.4.0). Auth: ``Authorization: Bearer `` (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") 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") async 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, }