- app/services/realtime_merge.py : merge à 3 voix diff3-lite, regions disjointes conservees, conflit par champ + drapeau - protocole base (client embarque la base de sa saisie) ; sans base -> LWW historique (retro-compat) - ack renvoie le bloc fusionne + conflict ; adoption cote client + toast ; broadcast du resultat fusionne - broadcast non bloquant : file sortante + tache writer par connexion, coalescence des curseurs - clients trop lents deconnectes (4413), budget ops anti-flood (400/10s) - fix fuite room 4404 + room_state() sur page inexistante - GET /api/realtime/stats (observabilite) - 26 tests test_realtime_v64.py ; suite 725 verte ; ruff + eslint OK ; version 6.4.0
467 lines
19 KiB
Python
467 lines
19 KiB
Python
"""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 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 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):
|
|
_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):
|
|
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
|