"""FlowDeck — v6.4.0 Realtime production: merge 3-voix (au-delà du LWW), broadcast non bloquant + coalescence des curseurs, corrections de fuites de rooms, anti-flood et observabilité. Covers : * ``realtime_merge`` — merge_text_3way / merge_block_3way (purs, sans WS) * protocole WS avec ``base`` → ack fusionné + drapeau conflict * convergence de deux clients sur le même bloc (les deux saisies survivent) * le LWW historique est préservé quand il n'y a pas de ``base`` * room 4404 non enregistrée (fuite corrigée), ``room_state`` page inexistante * coalescence des curseurs via la file sortante * ``GET /api/realtime/stats`` """ import json import os import tempfile import pytest from conftest import anon, login_test_client 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-v64" os.environ["RATE_LIMIT_ENABLED"] = "false" from app.config import settings settings.database_url = f"sqlite:///{db_path}" from app.db import get_conn, init_db from app.main import app 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 = {} manager.stat_ops = 0 manager.stat_merges = 0 manager.stat_conflicts = 0 manager.stat_slow_disconnects = 0 manager.stat_connections_total = 0 tc = TestClient(app, raise_server_exceptions=False) yield login_test_client(tc) os.unlink(db_path) def _make_page(title="V64 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)) # ── realtime_merge (purs) ──────────────────────────────────────────────── def test_merge_text_disjoint_edits_survive(): from app.services.realtime_merge import merge_text_3way base = "hello world" cur = "hello brave world" # serveur : insert "brave" inc = "hello world today" # client : append "today" out, conflict = merge_text_3way(base, cur, inc) assert conflict is False assert "brave" in out and "today" in out def test_merge_text_disjoint_edits_order_independent(): from app.services.realtime_merge import merge_text_3way base = "hello world" a, _ = merge_text_3way(base, "hello brave world", "hello world today") b, _ = merge_text_3way(base, "hello world today", "hello brave world") assert a == b == "hello brave world today" def test_merge_text_real_overlap_flags_conflict_lww(): from app.services.realtime_merge import merge_text_3way out, conflict = merge_text_3way("hello", "hxllo", "hyllo") assert conflict is True assert out == "hyllo" # LWW : incoming gagne def test_merge_text_shortcuts(): from app.services.realtime_merge import merge_text_3way assert merge_text_3way("a", "a", "b") == ("b", False) # client a changé assert merge_text_3way("a", "b", "a") == ("b", False) # serveur a changé assert merge_text_3way("a", "b", "b") == ("b", False) # même changement def test_merge_text_same_point_insertion_not_destructive(): from app.services.realtime_merge import merge_text_3way out, conflict = merge_text_3way("ab", "aXb", "aYb") # les deux insertions au même index sont conservées (concaténées) assert "X" in out and "Y" in out assert conflict is False def test_changed_region(): from app.services.realtime_merge import changed_region assert changed_region("hello", "hello") is None assert changed_region("hello world", "hello brave world") == (6, 6, "brave ") assert changed_region("abc", "abXc") == (2, 2, "X") def test_merge_block_field_level_no_conflict(): """Serveur et client changent des champs différents → pas de conflit.""" from app.services.realtime_merge import merge_block_3way base = {"id": "1", "text": "hi", "done": False} cur = {"id": "1", "text": "hi there", "done": False} # serveur édite text inc = {"id": "1", "text": "hi", "done": True} # client coche done merged, conflicts = merge_block_3way(base, cur, inc) assert merged["text"] == "hi there" assert merged["done"] is True assert conflicts == [] def test_merge_block_text_disjoint_merges(): from app.services.realtime_merge import merge_block_3way base = {"id": "1", "text": "hello world"} cur = {"id": "1", "text": "hello brave world"} inc = {"id": "1", "text": "hello world today"} merged, conflicts = merge_block_3way(base, cur, inc) assert conflicts == [] assert "brave" in merged["text"] and "today" in merged["text"] def test_merge_block_scalar_conflict_lww(): from app.services.realtime_merge import merge_block_3way merged, conflicts = merge_block_3way({"id": "1", "n": 1}, {"id": "1", "n": 2}, {"id": "1", "n": 3}) assert merged["n"] == 3 assert "n" in conflicts def test_merge_block_server_deleted_field_stays_deleted(): from app.services.realtime_merge import merge_block_3way base = {"id": "1", "tmp": "x"} cur = {"id": "1"} # serveur a supprimé tmp inc = {"id": "1", "tmp": "x"} # client n'a pas touché tmp merged, _ = merge_block_3way(base, cur, inc) assert "tmp" not in merged def test_merge_block_incoming_non_dict_flagged(): from app.services.realtime_merge import merge_block_3way _, conflicts = merge_block_3way({"id": "1"}, {"id": "1"}, "junk") assert conflicts == ["__block__"] def test_merge_block_preserves_id_when_base_empty(): from app.services.realtime_merge import merge_block_3way merged, _ = merge_block_3way({}, {}, {"id": "z", "text": "new"}) assert merged["id"] == "z" # ── apply_op (LWW historique préservé) ─────────────────────────────────── def test_apply_op_still_lww_without_base(): from app.services.realtime_server import apply_op blocks = [{"id": "a", "content": "A"}] out = apply_op(blocks, {"type": "update", "block": {"id": "a", "content": "B"}}) assert out[0]["content"] == "B" out = apply_op(blocks, {"type": "delete", "id": "a"}) assert out == [] out = apply_op(blocks, {"type": "move", "id": "a", "index": 0}) assert [b["id"] for b in out] == ["a"] # ── protocole WS ───────────────────────────────────────────────────────── def test_ws_requires_auth(client): anon(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_4404_does_not_leak_room(client): """Une tentative de connexion sur une page inexistante ne doit pas laisser de Room orpheline en mémoire (fuite corrigée en v6.4.0).""" from app.services.realtime_server import manager _auth_client(client) with pytest.raises(WebSocketDisconnect): with client.websocket_connect("/ws/pages/9999") as ws: ws.receive_json() assert 9999 not in manager._rooms, f"room leak: {list(manager._rooms)}" 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" s = ws.receive_json() assert s["t"] == "sync" assert s["blocks"][0]["id"] == "a" ws.send_json({"t": "hello"}) s2 = ws.receive_json() assert s2["t"] == "sync" and s2["blocks"][0]["content"] == "Hi" def test_ws_update_with_base_returns_merged(client): """Update avec base → ack fusionné, pas de conflit si seul le client a changé.""" 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 ws.send_json({"t": "op", "v": 0, "op": { "type": "update", "block": {"id": "a", "type": "paragraph", "content": "x!"}, "base": {"id": "a", "type": "paragraph", "content": "x"}, }}) ack = ws.receive_json() assert ack["t"] == "ack" assert ack.get("merged")["content"] == "x!" assert ack.get("conflict") is False def test_ws_two_clients_same_block_merge_disjoint(client): """Deux clients éditent le même bloc à des endroits différents → les deux saisies survivent (au-delà du LWW) et convergeamment.""" pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "hello world"}]) _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 wb.receive_json() # sync wa.receive_json() # peer_join # A édite la fin wa.send_json({"t": "op", "v": 1, "op": { "type": "update", "block": {"id": "a", "type": "paragraph", "content": "hello world today"}, "base": {"id": "a", "type": "paragraph", "content": "hello world"}, }}) ack_a = wa.receive_json() assert ack_a["t"] == "ack" op_b = wb.receive_json() assert op_b["t"] == "op" # B édite le début, sur la même base wb.send_json({"t": "op", "v": 1, "op": { "type": "update", "block": {"id": "a", "type": "paragraph", "content": "hello brave world"}, "base": {"id": "a", "type": "paragraph", "content": "hello world"}, }}) ack_b = wb.receive_json() assert ack_b["t"] == "ack" merged = ack_b["merged"]["content"] # les deux éditions disjointes doivent être conservées assert "brave" in merged, merged assert "today" in merged, merged assert ack_b["conflict"] is False # A reçoit aussi le bloc fusionné op_a2 = wa.receive_json() assert op_a2["t"] == "op" assert "brave" in op_a2["op"]["block"]["content"] assert "today" in op_a2["op"]["block"]["content"] # persistance du contenu fusionné 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 "brave" in d[0]["content"] and "today" in d[0]["content"] def test_ws_update_without_base_stays_lww(client): """Sans ``base`` (ancien client / protocole), le LWW historique s'applique.""" pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "x"}]) _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": "op", "v": 0, "op": {"type": "update", "block": {"id": "a", "type": "paragraph", "content": "first"}}}) wa.receive_json() # ack b_op = wb.receive_json() 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 (pas de merged sans base) a_op = wa.receive_json() assert a_op["op"]["block"]["content"] == "second" from app.db import get_conn with get_conn() as conn: row = conn.execute("SELECT content FROM pages WHERE id=?", (pid,)).fetchone() assert json.loads(row["content"])[0]["content"] == "second" def test_ws_conflicting_scalar_flagged(client): """Conflit sur un champ scalaire multi-états → drapeau conflict=True + LWW. Un booléen ne peut pas diverger depuis la même base (2 valeurs) ; on utilise un ``status`` à 3 états pour que les deux côtés changent *différemment*. """ pid = _make_page(blocks=[{"id": "a", "type": "todo", "content": "t", "status": "todo"}]) _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() # A passe todo → doing wa.send_json({"t": "op", "v": 1, "op": { "type": "update", "block": {"id": "a", "type": "todo", "content": "t", "status": "doing"}, "base": {"id": "a", "type": "todo", "content": "t", "status": "todo"}, }}) wa.receive_json() # ack wb.receive_json() # op broadcast # B passe todo → done (conflit scalaire réel, même base) wb.send_json({"t": "op", "v": 1, "op": { "type": "update", "block": {"id": "a", "type": "todo", "content": "t", "status": "done"}, "base": {"id": "a", "type": "todo", "content": "t", "status": "todo"}, }}) ack_b = wb.receive_json() assert ack_b["t"] == "ack" assert ack_b["conflict"] is True assert ack_b["merged"]["status"] == "done" # LWW : incoming gagne assert "status" in str(ack_b.get("merged")) def test_ws_stats_endpoint(client): _auth_client(client) pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "x"}]) with client.websocket_connect(f"/ws/pages/{pid}") as ws: ws.receive_json() ws.receive_json() ws.send_json({"t": "op", "v": 0, "op": { "type": "update", "block": {"id": "a", "type": "paragraph", "content": "y"}, "base": {"id": "a", "type": "paragraph", "content": "x"}, }}) ws.receive_json() # ack r = client.get("/api/realtime/stats") assert r.status_code == 200 d = r.json() assert d["rooms"] == 1 assert d["connections"] == 1 assert d["ops"] >= 1 assert d["merges"] >= 1 assert isinstance(d["pages"], list) and d["pages"][0]["page_id"] == pid def test_ws_stats_requires_auth(client): anon(client) r = client.get("/api/realtime/stats") assert r.status_code == 200 assert r.json() == {"error": "unauthorized"} def test_ws_op_budget_anti_flood(client): """Au-delà du budget d'ops, le serveur répond ack stale (sans appliquer).""" from app.services import realtime_server as rs pid = _make_page(blocks=[{"id": "a", "type": "paragraph", "content": "x"}]) _auth_client(client) old_max = rs.OP_WINDOW_MAX rs.OP_WINDOW_MAX = 3 try: with client.websocket_connect(f"/ws/pages/{pid}") as ws: ws.receive_json() ws.receive_json() for i in range(6): ws.send_json({"t": "op", "v": 0, "op": {"type": "update", "block": {"id": "a", "type": "paragraph", "content": f"v{i}"}}}) acks = [ws.receive_json() for _ in range(6)] stale = [a for a in acks if a.get("stale")] assert stale, "budget anti-flood n'a pas été déclenché" finally: rs.OP_WINDOW_MAX = old_max def test_room_state_missing_page_returns_empty(client): import asyncio from app.services.realtime_server import manager out = asyncio.run(manager.room_state(424242)) assert out == {"blocks": [], "title": "", "version": 0} assert 424242 not in manager._rooms def test_writer_coalesces_cursor_updates(): """La file sortante d'une connexion coalesce les curseurs (une seule position finale conservée) tout en préservant l'ordre des messages importants.""" import asyncio class FakeWS: def __init__(self): self.sent = [] async def send_json(self, msg): self.sent.append(msg) async def run(): from app.services.realtime_server import RealtimeManager, RTConn ws = FakeWS() conn = RTConn(ws, {"id": 1, "login": "u"}, 1) conn.writer = asyncio.create_task(RealtimeManager()._writer(conn)) # empiler plusieurs curseurs + un message important au milieu conn.out_q.put_nowait({"t": "sel", "block": "a", "offset": 1}) conn.out_q.put_nowait({"t": "sel", "block": "a", "offset": 2}) conn.out_q.put_nowait({"t": "title", "title": "T"}) conn.out_q.put_nowait({"t": "sel", "block": "a", "offset": 3}) await asyncio.sleep(0.05) conn.closed = True conn.out_q.put_nowait(None) await asyncio.sleep(0.02) if conn.writer and not conn.writer.done(): conn.writer.cancel() return ws.sent sent = asyncio.run(run()) sels = [m for m in sent if m.get("t") == "sel"] titles = [m for m in sent if m.get("t") == "title"] # les curseurs sont coalescés : on ne garde pas les positions 1 et 2 séparées assert len(sels) <= 2, f"curseurs non coalescés: {sels}" assert titles and titles[0]["title"] == "T" # l'offset final (3) doit être présent assert sels[-1]["offset"] == 3, sels