diff --git a/CHANGELOG.md b/CHANGELOG.md index d66f59b..31c2215 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,39 @@ # Changelog - FlowDeck +## v5.3.0 (2026-09-08) — Realtime : édition collaborative en direct + +> Deux utilisateurs peuvent maintenant éditer une même page **en même temps** +> : les blocs, le titre, les curseurs et la présence se synchronisent via une +> connexion WebSocket, avec fusion last-write-wins et repli sur polling. + +- **Passerelle WebSocket** `WS /ws/pages/{page_id}` (app/services/realtime_server.py + + app/routers/realtime.py) : salons en mémoire par page, authentifiés via le + cookie de session (refus 4401), page introuvable/supprimée → 4404. +- **Protocole de synchronisation** : `hello`/`sync`/`op`/`ack`/`title`/`sel`/ + `ping`/`peer_join`/`peer_leave` ; chaque page porte une version incrémentée à + chaque opération appliquée ; un client périmé reçoit un `sync` complet. +- **Merge last-write-wins par bloc** : insert, update (remplacement de bloc), + delete, move — appliqués dans l'ordre d'arrivée ; une frappe en cours dans un + bloc focalisé n'est jamais écrasée par une mise à jour distante. +- **Client éditeur** (app/templates/_page_editor_realtime.html) : diff local → + ops (update/delete/move/insert), envoi dès que l'utilisateur tape (debounce), + application des ops distantes avec re-render discret + re-focus du bloc actif. +- **Présence & curseurs** : avatars colorés dans la barre du haut et curseurs + de chaque pair positionnés dans le document (calque dédié, mise à jour à + l'édition/au scroll). +- **Titre synchronisé** : diffusion du titre (debounce) sans écraser un titre + en cours d'édition. +- **Fallback polling** : si le WebSocket est indisponible, rafraîchissement + toutes les 10 s de l'état serveur (adopté uniquement si pas de brouillon + local) avec reconnexion automatique. +- **Persistance** : écriture debounce (~1 s) de `content` + `title` en base, + flush immédiat à la déconnexion du dernier client. +- **Sécurité** : `connect-src` CSP étendu à `ws:` ; aucune donnée sensible + transmise (les peers n'exposent que id/login/full_name/couleur). +- **Tests** (tests/test_realtime.py) : auth, page absente, hello→sync, merge + LWW à deux clients + persistance, présence join/leave, curseurs, titre, + resync des clients périmés, `apply_op`/`merge_ops` unitaires. + ## v5.2.0 (2026-09-07) - Automations : un moteur de règles (if-this-then-that) > Nouveau moteur de règles : brancher des actions sur des événements de pages diff --git a/ROADMAP.md b/ROADMAP.md index affb90e..f794feb 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -423,13 +423,19 @@ app/ - [x] **Quick actions** — navigation, création, commandes - [x] **Recherche full-text** — SQLite FTS5 sur pages + propriétés (prérequis technique de la palette) -### v5.1.0 — Automations ⬜ (non commencé) +### v5.1.0 — Automations ✅ (2026-09-07) +> **Objectif** : moteur de règles if-this-then-that + boutons cliquables. **COMPLETED**. -> 💡 Note : les **webhooks sortants** (`webhook_subscriptions`, CRUD `/workspace/webhooks`, -> `services/webhook_outbound.py`) sont déjà implémentés — bonne base pour les actions automation. - -- [ ] **Database automations** — moteur de règles if-this-then-that (trigger + condition + action) -- [ ] **Buttons** — boutons cliquables déclenchant des actions +- [x] **Database automations** — moteur de règles if-this-then-that (trigger + condition + action) + - Déclencheurs : événement (`page.created/updated/deleted/moved`, `collection.*`) / cron (`*/N`, minute fixe, `@hourly`, `@daily`) / bouton + - Conditions combinables : eq, neq, contains, not_contains, is_empty, is_not_empty, changed + - Actions : webhook (X-FlowDeck-Secret), set_property (validé), create_page (interpolation `[[prop]]`/`{{title}}`), notify + - Tables `automations` + `automation_runs` (migration v5), scheduler de fond (60 s), historique des exécutions +- [x] **Buttons** — boutons cliquables déclenchant des actions + - Bloc `button` dans l'éditeur (menu slash) + picker d'automation inline + - Endpoint `/api/automations/{id}/run` (exempt CSRF) + UI Settings → Automations +- [x] Hooks événements dans collections.py / board.py / workspace.py (apply_page_template) +- [x] **11 tests** `tests/test_automations.py` ; suite complète 289 verte ### v5.2.0 — Infrastructure & Polish 🔄 (en cours) > ✅ **Priorité n°1 livrée** : les **migrations versionnées** (voir `app/migrations.py` + @@ -473,6 +479,17 @@ app/ - [x] **Templates de database prédéfinis** — 6 templates seedés (CRM, Project tracker, Task list, Content calendar, Meeting notes, Reading list) + galerie au clic « Database » (Get Started) et `/database` ; `POST /db/api` & `/db/inline/api` acceptent `template` et matérialisent les propriétés - [x] **Validation des propriétés** (required, unique, min/max) — côté serveur (`validate_property_rule`, 400 + messages) + UI (modale propriété + erreurs de cellule) +### v5.13.0 — Collaboration temps réel ✅ (2026-09-08) +> **Objectif** : édition collaborative en direct (parité Notion en équipe). **COMPLETED**. +> Chantier n°1 de la parité Notion ; socle pour v5.14.0 (synced blocks). + +- [x] **WebSocket gateway** — endpoint `WS /ws/pages/{id}`, auth via cookie session (refus 4401), page introuvable/supprimée → 4404 ; rooms par page en mémoire, chargées depuis la base au premier connect, droppées quand vides +- [x] **Présence** — avatars des utilisateurs connectés (topbar), couleur par utilisateur, join/leave broadcasté ; `welcome` = self + peers +- [x] **Curseurs live** — position curseur/sélection des autres éditeurs (block + offset), calque dédié, label avec nom, mise à jour à l'édition/au scroll +- [x] **Merge de modifications** — diffusion des ops de blocs insert/update/delete/move avec **last-write-wins par bloc** + version de page (stale → sync complet) ; titre synchronisé (debounce) ; persistance debounce ~1 s + flush à la déconnexion du dernier client +- [x] **Fallback polling** — si WS indisponible, rafraîchissement diff toutes les 10 s (adopté seulement sans brouillon local) + reconnexion automatique +- [x] CSP `connect-src` étendu à `ws:` ; **12 tests** `tests/test_realtime.py` ; suite complète 298 verte (+3 PDF pré-existants) + ### v5.4.0 — Expérience éditeur (nouveautés, parité Notion) - [ ] **Backlinks** — section « Lié depuis… » en bas de page (scan des liens internes) - [ ] **Duplicates** — « Duplicate » sur page + collection (menu `...`) @@ -557,15 +574,7 @@ app/ - [ ] **Full-width mode** — toggle pour passer la page en pleine largeur (comme Notion) - [ ] **Small text / typo options** — option de page : taille de police réduite, serif/mono -## v5.13.0 — Collaboration temps réel ⬜ (non commencé) -> **Objectif** : avancer depuis v6.0.0 le chantier n°1 de la parité Notion en équipe. -> Sans realtime, FlowDeck reste un outil mono-utilisateur partagé. - -- [ ] **WebSocket gateway** — endpoint `/ws/pages/{id}`, auth via session, rooms par page -- [ ] **Présence** — avatars des utilisateurs connectés sur la page (topbar), couleur par utilisateur -- [ ] **Curseurs live** — position des curseurs/sélections des autres éditeurs, label avec nom -- [ ] **Merge de modifications** — diffusion des opérations de blocs (insert/update/delete/move) avec last-write-wins par bloc + version de page pour détecter les conflits -- [ ] **Fallback polling** — si WS indisponible, rafraîchissement diff toutes les 10 s +## v5.13.0 — Collaboration temps réel ✅ (livré — voir section Completed) ## v5.14.0 — Synced blocks ⬜ (non commencé) > **Objectif** : avancer depuis v6.0.0 un bloc Notion très utilisé (même contenu dans @@ -604,8 +613,8 @@ app/ 1. ~~**v5.2.0 → Migrations versionnées**~~ ✅ livré (`schema_version` + `app/migrations.py`) 2. ~~**v5.0.0 → Command palette + FTS5**~~ ✅ livré (palette Ctrl+K + `GET /api/search`) 3. ~~**v5.3.0 → Inline databases + templates + validation**~~ ✅ livré (slash `/database`, 6 templates, validation propriétés) -4. **v5.13.0 → Realtime (WS + présence)** — chantier de parité Notion n°1, socle aussi pour v5.14 (synced blocks) -5. **v5.10.0 → Interactions de bloc** (drag&drop, undo/redo, duplicate) — gros impact UX, effort modéré +1. ~~**v5.13.0 → Realtime (WS + présence)**~~ ✅ livré (`app/services/realtime_server.py` + `WS /ws/pages/{id}`, présence, curseurs live, merge LWW, 12 tests) +2. **v5.10.0 → Interactions de bloc** (drag&drop, undo/redo, duplicate) — gros impact UX, effort modéré — prochaine priorité --- ## Résumé des phases @@ -617,11 +626,11 @@ Base + Kanban Éditeur + Gitea UX Pro MVP Onboard + UI Notion + Tags + Admin + Sharing COMPLETED + GitHub OAuth + Library -v4.0.2 ✅ v4.1–4.9 ✅ v4.10 ✅ v5.0–5.9 ⬜ v5.10–5.12 ⬜ v5.13–5.14 ⬜ v6.0 ⬜ +v4.0.2 ✅ v4.1–4.9 ✅ v4.10 ✅ v5.0–5.3 ✅ v5.4–5.12 ⬜ v5.13 ✅ · v5.14 ⬜ v6.0 ⬜ Quality DB views, Agent IA Palette → Éditeur bloc Realtime + Pro + Agent & Tests Templates & COMPLETED Automations (drag&drop, Synced blocks - Collaboration Embeds, Import, undo/redo) + - DB avancée, Wiki-links, - Calendrier, AI Templates & lock + Collaboration Realtime, undo/redo) + + DB avancée, Wiki-links, + Calendrier, AI Templates & lock -*Dernière mise à jour: 2026-09-06 — v5.0.0 (palette+FTS5) et v5.2.0 (migrations versionnées) livrés ; v5.3.0 Database Avancée livré (inline `/database`, 6 templates prédéfinis, validation propriétés, 280 tests) ; reste v5.4 → v6.0* +*Dernière mise à jour: 2026-09-08 — v5.13.0 Realtime livré (WebSocket gateway, présence, curseurs live, merge LWW + fallback polling, 12 tests, suite 298 verte) ; prochaine priorité v5.10.0 Interactions de bloc ; reste v5.4 → v6.0* diff --git a/VERSION b/VERSION index 7cbea07..e230c83 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -5.2.0 \ No newline at end of file +5.3.0 \ No newline at end of file diff --git a/app/main.py b/app/main.py index 352cb05..718217e 100644 --- a/app/main.py +++ b/app/main.py @@ -18,6 +18,7 @@ from app.routers import dashboard, board, notes, api, auth, webhooks, collection from app.routers.notifications import router as notifications_router from app.routers.automations import router as automations_router from app.routers.collaboration import router as collaboration_router +from app.routers.realtime import router as realtime_router from app.routers.gitea import router as gitea_router from app.routers.github_routes import router as github_router from app.services.gitea_client import gitea @@ -52,7 +53,7 @@ async def lifespan(_app: FastAPI): from app.services.automations import automation_scheduler automation_task = asyncio.create_task(automation_scheduler()) - logger.info("FlowDeck v5.2.0 started on port %d", settings.app_port) + logger.info("FlowDeck v5.3.0 started on port %d", settings.app_port) try: yield finally: @@ -70,7 +71,7 @@ async def lifespan(_app: FastAPI): app = FastAPI( title="FlowDeck", - version="5.2.0", + version="5.3.0", docs_url="/docs" if settings.log_level == "DEBUG" else None, redoc_url=None, lifespan=lifespan, @@ -102,6 +103,7 @@ app.include_router(export.router) app.include_router(notifications_router) app.include_router(automations_router) app.include_router(collaboration_router) +app.include_router(realtime_router) app.include_router(agent.router) app.include_router(search.router) diff --git a/app/middleware/security.py b/app/middleware/security.py index 9c3018f..e6b1779 100644 --- a/app/middleware/security.py +++ b/app/middleware/security.py @@ -68,7 +68,7 @@ class ContentSecurityPolicyMiddleware(BaseHTTPMiddleware): "style-src 'self' 'unsafe-inline' https://fonts.googleapis.com; " "img-src 'self' data: blob: https:; " "font-src 'self' data: https://fonts.gstatic.com; " - "connect-src 'self' https: wss:; " + "connect-src 'self' https: wss: ws:; " "media-src 'self' blob:; " "frame-src 'self'; " "object-src 'none'; " diff --git a/app/routers/realtime.py b/app/routers/realtime.py new file mode 100644 index 0000000..f8f6467 --- /dev/null +++ b/app/routers/realtime.py @@ -0,0 +1,49 @@ +"""FlowDeck — v5.13.0 Realtime: WebSocket gateway /ws/pages/{page_id}. + +Auth via cookie session (flowdeck_session). Rooms in-memory par page. +""" +from __future__ import annotations + +import json +import logging + +from fastapi import APIRouter, WebSocket +from starlette.websockets import WebSocketDisconnect + +from app.auth.session import SessionManager +from app.services.realtime_server import manager + +logger = logging.getLogger(__name__) +router = APIRouter(tags=["realtime"]) + + +@router.websocket("/ws/pages/{page_id}") +async def ws_page(websocket: WebSocket, page_id: int): + await websocket.accept() + user = SessionManager.decode_session( + websocket.cookies.get("flowdeck_session", "") + ) + if not user or not user.get("id"): + try: + await websocket.close(code=4401) + except Exception: + pass + return + + conn = await manager.connect(websocket, page_id, user) + if not conn: + return + try: + while True: + raw = await websocket.receive_text() + try: + msg = json.loads(raw) + except (TypeError, ValueError): + continue + await manager.handle(conn, msg) + except WebSocketDisconnect: + pass + except Exception as e: # noqa: BLE001 + logger.debug("ws closed: %s", e) + finally: + await manager.disconnect(conn) \ No newline at end of file diff --git a/app/services/realtime_server.py b/app/services/realtime_server.py new file mode 100644 index 0000000..bf28f92 --- /dev/null +++ b/app/services/realtime_server.py @@ -0,0 +1,294 @@ +"""FlowDeck — v5.13.0 Realtime: WebSocket gateway, présences, curseurs live, +merge des opérations de blocs (last-write-wins par bloc) + version de page. + +Rooms in-memory (un seul worker uvicorn). Persistance en base (page.content) +avec debounce. Fallback polling côté client si le WS est indisponible. +""" +from __future__ import annotations + +import asyncio +import json +import logging + +from fastapi import WebSocket + +from app.db import get_conn + +logger = logging.getLogger(__name__) + +COLORS = ["#2383E2", "#46A758", "#E5484D", "#F76B15", "#8E4EC6", "#12A594", + "#FFC53D", "#D6409F", "#0091FF", "#3E63DD", "#30A46C", "#FF3333"] + + +def color_for(uid: int) -> str: + return COLORS[(uid or 0) % len(COLORS)] + + +def block_id() -> str: + import time + return f"b{int(time.time()*1000)}" + + +def apply_op(blocks: list[dict], op: dict) -> list[dict]: + """Apply one block op (insert/update/delete/move) — LWW par bloc.""" + t = op.get("type") + if t == "insert": + blk = op.get("block") or {} + if not blk.get("id"): + blk = dict(blk) + blk["id"] = block_id() + idx = op.get("index") + if not isinstance(idx, int): + idx = len(blocks) + idx = max(0, min(idx, len(blocks))) + return blocks[:idx] + [blk] + blocks[idx:] + if t == "update": + nb = op.get("block") or {} + if not nb.get("id"): + return blocks + return [nb if b.get("id") == nb["id"] else b for b in blocks] + if t == "delete": + bid = op.get("id") + return [b for b in blocks if b.get("id") != bid] + if t == "move": + bid = op.get("id") + idx = op.get("index", 0) or 0 + out = [b for b in blocks if b.get("id") != bid] + idx = max(0, min(idx, len(out))) + moved = next((b for b in blocks if b.get("id") == bid), None) + if moved is None: + return blocks + out.insert(idx, moved) + return out + return blocks + + +def merge_ops(blocks: list[dict], ops: list[dict]) -> list[dict]: + """Apply a batch of ops sequentially (arrival order).""" + out = blocks + for op in ops or []: + out = apply_op(out, op) + return out + + +class Room: + __slots__ = ("page_id", "blocks", "title", "version", "conns", + "persist_task", "dirty") + + def __init__(self, page_id: int): + self.page_id = page_id + self.blocks: list[dict] = [] + self.title = "" + self.version = 0 + self.conns: set["RTConn"] = set() + self.persist_task: asyncio.Task | None = None + self.dirty = False + + +class RTConn: + __slots__ = ("ws", "user", "page_id") + + def __init__(self, ws: WebSocket, user: dict, page_id: int): + self.ws = ws + self.user = user + self.page_id = page_id + + +class RealtimeManager: + def __init__(self): + self._rooms: dict[int, Room] = {} + + def room(self, page_id: int) -> Room: + return self._rooms.setdefault(page_id, Room(page_id)) + + @staticmethod + def _peer(user: dict) -> dict: + uid = user.get("id") or 0 + return { + "id": uid, + "login": user.get("login", ""), + "full_name": user.get("full_name", "") or user.get("login", ""), + "color": color_for(uid), + } + + async def load_room(self, room: Room) -> bool: + try: + with get_conn() as conn: + row = conn.execute( + "SELECT title, content, content_format FROM pages WHERE id=? AND deleted_at IS NULL", + (room.page_id,), + ).fetchone() + except Exception as e: + logger.warning("realtime load failed: %s", e) + return False + if not row: + return False + room.title = row["title"] or "" + if (row["content_format"] or "") == "blocks" and row["content"]: + try: + room.blocks = json.loads(row["content"]) + except Exception: + room.blocks = [] + return True + + async def connect(self, ws: WebSocket, page_id: int, user: dict) -> RTConn | None: + room = self.room(page_id) + if not room.conns and not await self.load_room(room): + await ws.close(code=4404) + return None + conn = RTConn(ws, user, page_id) + room.conns.add(conn) + me = self._peer(user) + peers = [self._peer(c.user) for c in room.conns if c is not conn] + await ws.send_json({"t": "welcome", "self": me, + "peers": peers, "color": me["color"]}) + await ws.send_json({"t": "sync", "blocks": room.blocks, + "title": room.title, "version": room.version}) + for c in room.conns: + if c is not conn: + try: + await c.ws.send_json({"t": "peer_join", "peer": me}) + except Exception: + pass + return conn + + async def disconnect(self, conn: RTConn): + room = self._rooms.get(conn.page_id) + if not room: + return + room.conns.discard(conn) + for c in room.conns: + try: + await c.ws.send_json({"t": "peer_leave", + "id": conn.user.get("id") or 0}) + except Exception: + pass + if not room.conns: + await self.flush(room) + self._rooms.pop(conn.page_id, None) + + async def flush(self, room: Room): + """Write current room state to DB (sync, used on idle + disconnect).""" + if room.persist_task and not room.persist_task.done(): + room.persist_task.cancel() + await self._persist(room) + + async def _persist(self, room: Room): + try: + with get_conn() as conn: + conn.execute( + "UPDATE pages SET content=?, title=?, updated_at=CURRENT_TIMESTAMP WHERE id=?", + (json.dumps(room.blocks, ensure_ascii=False), room.title, room.page_id), + ) + conn.commit() + room.dirty = False + except Exception as e: + logger.warning("realtime persist failed: %s", e) + + def _schedule_persist(self, room: Room): + if room.persist_task and not room.persist_task.done(): + return + room.dirty = True + + async def _run(): + try: + await asyncio.sleep(1.2) + await self._persist(room) + except asyncio.CancelledError: + pass + + room.persist_task = asyncio.create_task(_run()) + + async def _broadcast(self, room: Room, msg: dict, exclude: RTConn | None = None): + for c in room.conns: + if c is exclude: + continue + try: + await c.ws.send_json(msg) + except Exception: + pass + + async def handle(self, conn: RTConn, msg: dict): + room = self._rooms.get(conn.page_id) + if not room: + return + t = msg.get("t") + me = (conn.user.get("id") or 0) + + if t == "hello": + try: + await conn.ws.send_json({"t": "sync", "blocks": room.blocks, + "title": room.title, "version": room.version}) + except Exception: + pass + return + + if t == "sync_req": + try: + await conn.ws.send_json({"t": "sync", "blocks": room.blocks, + "title": room.title, "version": room.version}) + except Exception: + pass + return + + if t == "ping": + try: + await conn.ws.send_json({"t": "pong"}) + except Exception: + pass + return + + if t == "op": + op = msg.get("op") or {} + client_v = msg.get("v", 0) + room.blocks = apply_op(room.blocks, op) + room.version += 1 + self._schedule_persist(room) + stale = client_v < room.version - 1 + await self._broadcast(room, {"t": "op", "op": op, "from": me, + "v": room.version}, exclude=conn) + try: + await conn.ws.send_json({"t": "ack", "v": room.version, + "stale": stale}) + except Exception: + pass + if stale: + try: + await conn.ws.send_json({"t": "sync", + "blocks": room.blocks, + "title": room.title, + "version": room.version}) + except Exception: + pass + return + + if t == "title": + fmt = (msg.get("title") or "").strip() + if fmt and fmt != room.title: + room.title = fmt + room.version += 1 + self._schedule_persist(room) + await self._broadcast(room, {"t": "title", "title": room.title, + "from": me, "v": room.version}, exclude=conn) + try: + await conn.ws.send_json({"t": "ack", "v": room.version, + "stale": False}) + except Exception: + pass + return + + if t == "sel": + await self._broadcast(room, {"t": "sel", "from": me, + "peer": self._peer(conn.user), + "block": msg.get("block"), + "offset": msg.get("offset", 0)}, exclude=conn) + return + + async def room_state(self, page_id: int) -> dict: + room = self.room(page_id) + if not room.conns and not room.blocks: + await self.load_room(room) + return {"blocks": room.blocks, "title": room.title, "version": room.version} + + +manager = RealtimeManager() \ No newline at end of file diff --git a/app/templates/_page_editor_realtime.html b/app/templates/_page_editor_realtime.html new file mode 100644 index 0000000..d8736a5 --- /dev/null +++ b/app/templates/_page_editor_realtime.html @@ -0,0 +1,480 @@ + + \ No newline at end of file diff --git a/app/templates/_page_editor_scripts.html b/app/templates/_page_editor_scripts.html index 3697a0c..53ebcbc 100644 --- a/app/templates/_page_editor_scripts.html +++ b/app/templates/_page_editor_scripts.html @@ -715,7 +715,7 @@ init(){ _bid=Math.floor(Math.random()*10000); const dataEl=document.getElementById('page-data'); - if(dataEl){try{const data=JSON.parse(dataEl.textContent);this.pid=data.id;this.pageTitle=data.title||'';this.favorited=data.favorited||false;const fmt=data.content_format||'blocks';const raw=data.content||''; + if(dataEl){try{const data=JSON.parse(dataEl.textContent);this.pid=data.id;this.pageTitle=data.title||'';this.favorited=data.favorited||false;const fmt=data.content_format||'blocks';this.contentFormat=fmt;const raw=data.content||''; if(fmt==='file'){this.fileData=data;this.loadFileContent(data);} else if(fmt==='blocks'&&raw){try{this.blocks=JSON.parse(raw);this.blocks.forEach(b=>{if(!b.id)b.id=genId()});}catch(e){this.blocks=[];}} else if(raw&&raw.trim())this.blocks=this.md2b(raw);}catch(e){}} @@ -726,6 +726,7 @@ SM.init();_tblInitDelegate();this.render(); const tEl=document.getElementById('_titleEl'); if(tEl){const syncT=()=>{tEl.classList.toggle('empty',!(tEl.textContent||'').trim());};syncT();tEl.addEventListener('input',syncT);tEl.addEventListener('blur',syncT);} + if(this.contentFormat==='blocks'&&window.__fdRT&&window.__fdRT.start){setTimeout(()=>{try{window.__fdRT.start(this);}catch(e){}},150);} setTimeout(()=>{this.focusBlock();const ct=document.getElementById('_blocksCt');if(!ct)return; ct.addEventListener('keydown',e=>this.onKd(e)); ct.addEventListener('paste',e=>this.handleSmartPaste(e)); diff --git a/app/templates/page_editor.html b/app/templates/page_editor.html index a184603..9784f19 100644 --- a/app/templates/page_editor.html +++ b/app/templates/page_editor.html @@ -9,5 +9,6 @@ page_title %}{{ page.title }}{% endblock %} {% block topbar %} {% endblock %} {% block scripts %} {% include "_page_editor_scripts.html" %} +{% include "_page_editor_realtime.html" %} {% include "_database_table_scripts.html" %} {% endblock %} diff --git a/tests/test_realtime.py b/tests/test_realtime.py new file mode 100644 index 0000000..1204c2b --- /dev/null +++ b/tests/test_realtime.py @@ -0,0 +1,266 @@ +"""FlowDeck — v5.13.0 Realtime: WebSocket gateway, présence, curseurs live, +merge LWW des opérations de blocs + version de page + persistance debounce. + +Covers: la route WS /ws/pages/{page_id} (auth via cookie, page introuvable), +le protocole hello/sync/op/ack/sel/title/peer_*, le merge last-write-wins et +la resynchronisation des clients périmés. +""" +import json +import os +import tempfile + +import pytest +from fastapi.testclient import TestClient + + +from starlette.websockets import WebSocketDisconnect + + +@pytest.fixture +def client(): + db_file = tempfile.NamedTemporaryFile(suffix=".db", delete=False) + db_path = db_file.name + db_file.close() + + os.environ["DATABASE_URL"] = f"sqlite:///{db_path}" + os.environ["APP_SECRET_KEY"] = "test-secret-for-realtime" + os.environ["RATE_LIMIT_ENABLED"] = "false" + from app.config import settings + settings.database_url = f"sqlite:///{db_path}" + + from app.main import app + from app.db import init_db, get_conn + init_db() + with get_conn() as conn: + conn.execute("INSERT OR IGNORE INTO users (id, login, full_name, is_admin) VALUES (1, 'tester', 'Tester', 1)") + conn.execute("INSERT OR IGNORE INTO users (id, login, full_name, is_admin) VALUES (2, 'other', 'Other', 0)") + conn.commit() + + from app.services.realtime_server import manager + manager._rooms = {} + + tc = TestClient(app, raise_server_exceptions=False) + yield tc + + os.unlink(db_path) + + +def _make_page(title="Realtime Page", blocks=None): + from app.db import get_conn + with get_conn() as conn: + cur = conn.execute( + "INSERT INTO pages (workspace, title, content, content_format, parent_section) " + "VALUES ('Private', ?, ?, 'blocks', 'Private')", + (title, json.dumps(blocks or [], ensure_ascii=False)), + ) + conn.commit() + return cur.lastrowid + + +def _token(user_id, login): + from app.auth.session import SessionManager + return SessionManager.create_session({"id": user_id, "login": login, + "full_name": login.title(), "is_admin": 1}) + + +def _auth_client(client, user_id=1, login="tester"): + client.cookies.set("flowdeck_session", _token(user_id, login)) + + +# ── apply_op (service pur) ── + +def test_apply_op_insert(): + from app.services.realtime_server import apply_op + blocks = [{"id": "a", "content": "A"}] + out = apply_op(blocks, {"type": "insert", "index": 0, + "block": {"id": "b", "content": "B"}}) + assert [b["id"] for b in out] == ["b", "a"] + assert apply_op(blocks, {"type": "insert", "index": None, + "block": {"content": "C"}})[-1]["content"] == "C" + # index clampé + assert apply_op(blocks, {"type": "insert", "index": 99, + "block": {"id": "z"}})[-1]["id"] == "z" + + +def test_apply_op_update_delete_move(): + from app.services.realtime_server import apply_op + blocks = [{"id": "a", "content": "A"}, {"id": "b", "content": "B"}, + {"id": "c", "content": "C"}] + out = apply_op(blocks, {"type": "update", "block": {"id": "b", "content": "B2"}}) + assert next(x for x in out if x["id"] == "b")["content"] == "B2" + out = apply_op(blocks, {"type": "update", "block": {"id": "nope", "content": "X"}}) + assert out == blocks + out = apply_op(blocks, {"type": "delete", "id": "a"}) + assert [b["id"] for b in out] == ["b", "c"] + out = apply_op(blocks, {"type": "move", "id": "c", "index": 0}) + assert [b["id"] for b in out] == ["c", "a", "b"] + out = apply_op(blocks, {"type": "move", "id": "nope", "index": 0}) + assert out == blocks + assert apply_op(blocks, {"type": "unknown"}) == blocks + + +def test_merge_ops_sequential(): + from app.services.realtime_server import merge_ops + blocks = [] + ops = [ + {"type": "insert", "index": 0, "block": {"id": "a", "content": "A"}}, + {"type": "insert", "index": 1, "block": {"id": "b", "content": "B"}}, + {"type": "update", "block": {"id": "a", "content": "A2"}}, + {"type": "move", "id": "b", "index": 0}, + ] + assert [x["id"] for x in merge_ops(blocks, ops)] == ["b", "a"] + assert next(x for x in merge_ops(blocks, ops) if x["id"] == "a")["content"] == "A2" + + +# ── Auth & présence de page ── + +def test_ws_requires_auth(client): + _make_page() + with pytest.raises(WebSocketDisconnect) as exc: + with client.websocket_connect("/ws/pages/1") as ws: + ws.receive_json() + assert exc.value.code == 4401 + + +def test_ws_auth_rejected(client): + _make_page() + with pytest.raises(WebSocketDisconnect) as exc: + with client.websocket_connect("/ws/pages/1") as ws: + ws.receive_json() + assert exc.value.code == 4401 + + +def test_ws_missing_page_4404(client): + _auth_client(client) + with pytest.raises(WebSocketDisconnect) as exc: + with client.websocket_connect("/ws/pages/9999") as ws: + ws.receive_json() + assert exc.value.code == 4404 + + +def test_ws_hello_gets_sync(client): + pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "Hi"}]) + _auth_client(client) + with client.websocket_connect(f"/ws/pages/{pid}") as ws: + w = ws.receive_json() + assert w["t"] == "welcome" + assert w["self"]["id"] == 1 + s = ws.receive_json() + assert s["t"] == "sync" + assert s["version"] == 0 + assert s["blocks"][0]["id"] == "a" + assert s["title"] == "Realtime Page" + ws.send_json({"t": "hello"}) + s2 = ws.receive_json() + assert s2["t"] == "sync" and s2["blocks"][0]["content"] == "Hi" + + +# ── merge LWW + broadcast ── + +def test_ws_lww_merge(client): + pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "x"}, + {"id": "b", "type": "paragraph", "content": "y"}]) + _auth_client(client, 1, "tester") + with client.websocket_connect(f"/ws/pages/{pid}") as wa: + wa.receive_json() # welcome + wa.receive_json() # sync + _auth_client(client, 2, "other") + with client.websocket_connect(f"/ws/pages/{pid}") as wb: + wb.receive_json() # welcome (peer tester) + wb.receive_json() # sync + wa.receive_json() # peer_join other + wa.send_json({"t": "op", "v": 0, "op": {"type": "update", + "block": {"id": "a", "type": "paragraph", "content": "first"}}}) + ack = wa.receive_json() # ack de son propre op + assert ack["t"] == "ack" and ack["v"] == 1 + b_op = wb.receive_json() + assert b_op["t"] == "op" and b_op["from"] == 1 + assert b_op["op"]["block"]["content"] == "first" + wb.send_json({"t": "op", "v": 1, "op": {"type": "update", + "block": {"id": "a", "type": "paragraph", "content": "second"}}}) + wb.receive_json() # ack de son propre op + a_op = wa.receive_json() + assert a_op["t"] == "op" and a_op["from"] == 2 + assert a_op["op"]["block"]["content"] == "second" + # dernier arrivé gagne (LWW) et persistance au disconnect + from app.db import get_conn + with get_conn() as conn: + row = conn.execute("SELECT content FROM pages WHERE id=?", (pid,)).fetchone() + d = json.loads(row["content"]) + assert d[0]["content"] == "second" + + +# ── présence ── + +def test_ws_presence_join_leave(client): + pid = _make_page() + _auth_client(client, 1, "tester") + with client.websocket_connect(f"/ws/pages/{pid}") as wa: + wa.receive_json(); wa.receive_json() + _auth_client(client, 2, "other") + with client.websocket_connect(f"/ws/pages/{pid}") as wb: + w_f = wb.receive_json() # welcome + assert w_f["t"] == "welcome" + assert any(p["id"] == 1 for p in w_f["peers"]) + wb.receive_json() # sync + pj = wa.receive_json() # peer_join other + assert pj["t"] == "peer_join" and pj["peer"]["id"] == 2 + pl = wa.receive_json() # peer_leave other + assert pl["t"] == "peer_leave" and pl["id"] == 2 + + +def test_ws_cursor_broadcast(client): + pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "l"}]) + _auth_client(client, 1, "tester") + with client.websocket_connect(f"/ws/pages/{pid}") as wa: + wa.receive_json(); wa.receive_json() + _auth_client(client, 2, "other") + with client.websocket_connect(f"/ws/pages/{pid}") as wb: + wb.receive_json(); wb.receive_json() + wa.receive_json() # peer_join + wa.send_json({"t": "sel", "block": "a", "offset": 1}) + m = wb.receive_json() + assert m["t"] == "sel" and m["from"] == 1 + assert m["block"] == "a" and m["offset"] == 1 + wb.send_json({"t": "sel", "block": None, "offset": 0}) + m2 = wa.receive_json() + assert m2["t"] == "sel" and m2["block"] is None + + +def test_ws_title_broadcast(client): + pid = _make_page() + _auth_client(client, 1, "tester") + with client.websocket_connect(f"/ws/pages/{pid}") as wa: + wa.receive_json(); wa.receive_json() + _auth_client(client, 2, "other") + with client.websocket_connect(f"/ws/pages/{pid}") as wb: + wb.receive_json(); wb.receive_json() + wa.receive_json() + wa.send_json({"t": "title", "title": "Nouveau titre"}) + m = wb.receive_json() + assert m["t"] == "title" and m["title"] == "Nouveau titre" + + +# ── resynchronisation des clients périmés ── + +def test_ws_stale_client_gets_sync(client): + pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "x"}]) + _auth_client(client) + with client.websocket_connect(f"/ws/pages/{pid}") as ws: + ws.receive_json() # welcome + ws.receive_json() # sync + for i in range(3): + ws.send_json({"t": "op", "v": 0, "op": {"type": "update", + "block": {"id": "a", "type": "paragraph", "content": f"v{i}"}}}) + # séquence déterministe : ack(v1), ack(v2,stale)+sync(v2), ack(v3,stale)+sync(v3) + got_stale = False + syncs = [] + for _ in range(5): + m = ws.receive_json() + if m["t"] == "ack" and m.get("stale"): + got_stale = True + if m["t"] == "sync": + syncs.append(m) + assert got_stale + assert syncs and syncs[-1]["version"] >= 3 + assert syncs[-1]["blocks"][0]["content"] == "v2" \ No newline at end of file