Files
flowdeck/app/services/sync_engine.py
T
bruno b5207216f1 feat: v6.0.0 PWA offline support
- manifest + icones, service worker (precache, network-first, Background Sync)

- module client FlowOffline (IndexedDB, queue, delta, flush) + hook editeur

- endpoints /api/v2/sync/{delta,batch,status} + moteur de sync (conflits LWW/orpheline/copie offline)

- migrations offline_sync_queue + sync_version (triggers)

- UI offline (banner, badge sync, toasts, icone dirty) + doc /help

- tests pytest (sync, migrations, SW, offline) + E2E Playwright; bump 6.0.0
2026-09-18 13:05:40 -04:00

422 lines
19 KiB
Python

"""FlowDeck — offline synchronization engine (v6.0.0 PWA).
Reconciles offline mutations (queued on the client) with server state:
- ``get_delta`` — changes on the server since a given timestamp, for offline
clients to pull before pushing their own batch.
- ``apply_batch`` — replays a batch of offline mutations with optimistic
concurrency control. ``sync_version`` (auto-bumped by SQLite triggers added
in migration v17) is the version token.
Conflict model (from docs/V6_PWA_Progressive_Web_App.md):
- edit-edit → last-write-wins by default (applied + reported)
- edit-delete→ the page was deleted server-side → orphan copy created
- create-create → same title already exists server-side → renamed "… (copie offline)"
- delete → soft-delete (idempotent)
"""
from __future__ import annotations
import datetime as _dt
import json
import logging
import time
from app.db import get_conn
logger = logging.getLogger(__name__)
class ConflictError(Exception):
"""Raised when a mutation conflicts with server state."""
def __init__(self, detail: dict):
super().__init__(detail.get("type"))
self.details = detail
class SyncEngine:
"""Apply/read offline mutations. Stateless methods, thin sqlite access."""
# ── helpers ────────────────────────────────────────────────────────────
@staticmethod
def _now_epoch() -> float:
return time.time()
@staticmethod
def _can_access(conn, user_id: int, workspace_id: int | None) -> bool:
"""Owner or member of the workspace; legacy NULL workspace → allow."""
if not workspace_id:
return True
ws = conn.execute(
"SELECT owner_id FROM workspaces WHERE id=?", (workspace_id,)
).fetchone()
if ws and ws["owner_id"] == user_id:
return True
member = conn.execute(
"SELECT 1 FROM workspace_members WHERE workspace_id=? AND user_id=?",
(workspace_id, user_id),
).fetchone()
return bool(member)
@staticmethod
def _page_workspace_id(conn, page_id: int) -> int | None:
row = conn.execute(
"SELECT workspace_id FROM pages WHERE id=?", (page_id,)
).fetchone()
return row["workspace_id"] if row else None
# ── delta ─────────────────────────────────────────────────────────────
async def get_delta(self, user_id: int, since: float, workspace_id: int | None) -> dict:
"""Return server-side changes since `since` (epoch seconds)."""
if since and since > 1e12:
since /= 1000.0 # accept epoch-millis from legacy clients
changes: list[dict] = []
with get_conn() as conn:
if not self._can_access(conn, user_id, workspace_id):
return {"error": "forbidden"}
# ── pages (created / updated / soft-deleted) ──
rows = conn.execute(
"""SELECT * FROM pages
WHERE (? IS NULL OR workspace_id = ?)
AND (CAST(strftime('%s', COALESCE(updated_at, created_at)) AS REAL) > ?
OR (deleted_at IS NOT NULL
AND CAST(strftime('%s', deleted_at) AS REAL) > ?))""",
(workspace_id, workspace_id, since, since),
).fetchall()
for r in rows:
data = dict(r)
deleted = data.get("deleted_at") is not None
created_epoch = _iso_epoch(data.get("created_at"))
created_after = bool(created_epoch and created_epoch > since)
if deleted:
ctype = "page_deleted"
elif created_after:
ctype = "page_created"
else:
ctype = "page_updated"
changes.append({"change_type": ctype, "data": data})
# ── collections (created / updated) ──
coll_rows = conn.execute(
"SELECT * FROM collections WHERE CAST(strftime('%s', updated_at) AS REAL) > ?",
(since,),
).fetchall()
for r in coll_rows:
data = dict(r)
created_epoch = _iso_epoch(data.get("created_at"))
changes.append({
"change_type": "collection_created" if created_epoch and created_epoch > since
else "collection_updated",
"data": data,
})
# ── collection rows (informational for offline reading) ──
cp_rows = conn.execute(
"SELECT * FROM collection_pages WHERE CAST(strftime('%s', updated_at) AS REAL) > ?",
(since,),
).fetchall()
for r in cp_rows:
changes.append({"change_type": "collection_row_updated", "data": dict(r)})
return {
"changes": changes,
"server_time": self._now_epoch(),
"has_more": False,
}
# ── batch ──────────────────────────────────────────────────────────────
async def apply_batch(self, user_id: int, mutations: list[dict], device_id: str) -> dict:
"""Apply a batch of offline mutations; returns per-mutation results.
``mutations`` = [{"id"|"mutation_id", "type", "payload", "client_timestamp"}]
"""
results: list[dict] = []
conflicts: list[dict] = []
for mut in mutations:
mut_id = mut.get("id") or mut.get("mutation_id") or f"m{len(results)}"
mtype = mut.get("type", "")
payload = mut.get("payload") or {}
client_ts = float(mut.get("client_timestamp") or 0)
result = {
"mutation_id": mut_id,
"type": mtype,
"status": "synced",
"server_version": None,
}
try:
handler = getattr(self, f"_mut_{mtype}", None)
if handler is None:
raise ValueError(f"unknown mutation type: {mtype}")
outcome = handler(user_id, payload)
result.update(outcome)
status = "conflict" if outcome.get("conflict") else "synced"
result["status"] = status
if outcome.get("conflict"):
conflicts.append({
"mutation_id": mut_id,
"type": mtype,
**outcome["conflict"],
})
except ConflictError as exc:
result["status"] = "conflict"
result["conflict"] = exc.details
conflicts.append({"mutation_id": mut_id, "type": mtype, **exc.details})
except Exception as exc: # noqa: BLE001 — report, don't kill the batch
logger.warning("sync mutation %s failed: %s", mut_id, exc)
result["status"] = "failed"
result["error"] = str(exc)
finally:
self._record(user_id, device_id, mtype, mut, client_ts, result)
results.append(result)
return {"results": results, "conflicts": conflicts}
# ── page mutations ─────────────────────────────────────────────────────
def _mut_page_create(self, user_id: int, payload: dict) -> dict:
title = (payload.get("title") or "New page").strip() or "New page"
workspace_id = payload.get("workspace_id")
parent_id = payload.get("parent_id")
with get_conn() as conn:
if not self._can_access(conn, user_id, workspace_id):
raise ConflictError({"type": "forbidden", "workspace_id": workspace_id})
# create-create conflict: same title already present at same parent
new_title = title
dup = conn.execute(
"""SELECT id FROM pages
WHERE title=? AND (? IS NULL OR workspace_id = ?)
AND (? IS NULL OR parent_id IS ?)
AND deleted_at IS NULL
LIMIT 1""",
(title, workspace_id, workspace_id, parent_id, parent_id),
).fetchone()
if dup:
new_title = f"{title} (copie offline)"
cur = conn.execute(
"""INSERT INTO pages
(workspace, workspace_id, title, content, content_format,
parent_id, parent_section, sort_order)
VALUES ('', ?, ?, ?, ?, ?, 'Private', ?)""",
(workspace_id, new_title,
payload.get("content", ""),
payload.get("content_format", "blocks"),
parent_id,
payload.get("sort_order", 0)),
)
conn.commit()
page_id = cur.lastrowid
synced = payload.get("client_page_id")
return {"page_id": page_id, "server_version": 1, "client_page_id": synced,
"conflict": {"type": "create_create", "renamed": new_title != title}
if new_title != title else None}
def _mut_page_update(self, user_id: int, payload: dict) -> dict:
page_id = payload.get("page_id")
base_version = payload.get("base_version")
with get_conn() as conn:
row = conn.execute("SELECT * FROM pages WHERE id=?", (page_id,)).fetchone()
if not row:
raise ConflictError({"type": "page_not_found", "page_id": page_id})
ws_id = row["workspace_id"] or payload.get("workspace_id")
if not self._can_access(conn, user_id, ws_id):
raise ConflictError({"type": "forbidden", "page_id": page_id})
# edit-delete: page soft-deleted server-side → orphan copy
if row["deleted_at"] is not None:
cur = conn.execute(
"""INSERT INTO pages
(workspace, workspace_id, title, content, content_format, parent_section)
VALUES ('', ?, ?, ?, ?, 'Private')""",
(ws_id,
payload.get("title") or row["title"],
(payload.get("content")
if payload.get("content") is not None else row["content"]),
payload.get("content_format") or row["content_format"]),
)
conn.commit()
orphan_id = cur.lastrowid
raise ConflictError({
"type": "edit_delete",
"page_id": page_id,
"new_page_id": orphan_id,
"detail": "La page a été supprimée côté serveur — copie récréée en page orpheline",
})
# edit-edit: version mismatch → last-write-wins + conflict report
server_version = row["sync_version"]
conflict = None
if base_version is not None and server_version != base_version:
conflict = {
"type": "edit_edit",
"page_id": page_id,
"client_version": base_version,
"server_version": server_version,
}
sets, params = [], []
if payload.get("title") is not None:
sets.append("title=?")
params.append(payload["title"])
if payload.get("content") is not None:
sets.append("content=?")
params.append(payload["content"])
if payload.get("content_format") is not None:
sets.append("content_format=?")
params.append(payload["content_format"])
sets_str = ", ".join(sets) if sets else "updated_at=updated_at"
params.append(page_id)
conn.execute(
f"UPDATE pages SET {sets_str}, updated_at=CURRENT_TIMESTAMP WHERE id=?",
params,
)
conn.commit()
new_version = conn.execute(
"SELECT sync_version FROM pages WHERE id=?", (page_id,)
).fetchone()["sync_version"]
return {"page_id": page_id, "server_version": new_version,
"conflict": conflict}
def _mut_page_delete(self, user_id: int, payload: dict) -> dict:
page_id = payload.get("page_id")
with get_conn() as conn:
row = conn.execute("SELECT id, workspace_id FROM pages WHERE id=?", (page_id,)).fetchone()
if not row:
return {"page_id": page_id, "server_version": None} # idempotent / already hard-deleted
ws_id = row["workspace_id"]
if not self._can_access(conn, user_id, ws_id):
raise ConflictError({"type": "forbidden", "page_id": page_id})
conn.execute(
"UPDATE pages SET deleted_at=CURRENT_TIMESTAMP, "
"updated_at=CURRENT_TIMESTAMP WHERE id=? AND deleted_at IS NULL",
(page_id,),
)
conn.commit()
new_version = conn.execute(
"SELECT sync_version FROM pages WHERE id=?", (page_id,)
).fetchone()["sync_version"]
return {"page_id": page_id, "server_version": new_version}
def _mut_page_move(self, user_id: int, payload: dict) -> dict:
page_id = payload.get("page_id")
with get_conn() as conn:
row = conn.execute("SELECT id, workspace_id FROM pages WHERE id=?", (page_id,)).fetchone()
if not row:
return {"page_id": page_id, "server_version": None}
ws_id = row["workspace_id"]
if not self._can_access(conn, user_id, ws_id):
raise ConflictError({"type": "forbidden", "page_id": page_id})
sets, params = [], []
if payload.get("parent_id") is not None:
sets.append("parent_id=?")
params.append(payload["parent_id"])
if payload.get("sort_order") is not None:
sets.append("sort_order=?")
params.append(payload["sort_order"])
if sets:
params.append(page_id)
conn.execute(
f"UPDATE pages SET {', '.join(sets)}, updated_at=CURRENT_TIMESTAMP WHERE id=?",
params,
)
conn.commit()
new_version = conn.execute(
"SELECT sync_version FROM pages WHERE id=?", (page_id,)
).fetchone()["sync_version"]
return {"page_id": page_id, "server_version": new_version}
# ── collection mutations ───────────────────────────────────────────────
def _mut_collection_create(self, user_id: int, payload: dict) -> dict:
with get_conn() as conn:
cur = conn.execute(
"""INSERT INTO collections (name, description, icon, schema_json)
VALUES (?, ?, ?, ?)""",
(payload.get("name", ""), payload.get("description", ""),
payload.get("icon", "📋"), json.dumps(payload.get("schema", []))),
)
conn.commit()
return {"collection_id": cur.lastrowid, "server_version": 1}
def _mut_collection_update(self, user_id: int, payload: dict) -> dict:
cid = payload.get("collection_id")
with get_conn() as conn:
row = conn.execute("SELECT * FROM collections WHERE id=?", (cid,)).fetchone()
if not row:
raise ConflictError({"type": "collection_not_found", "collection_id": cid})
sets, params = [], []
if payload.get("name") is not None:
sets.append("name=?")
params.append(payload["name"])
if payload.get("description") is not None:
sets.append("description=?")
params.append(payload["description"])
if payload.get("icon") is not None:
sets.append("icon=?")
params.append(payload["icon"])
if payload.get("schema") is not None:
sets.append("schema_json=?")
params.append(json.dumps(payload["schema"]))
if sets:
params.append(cid)
conn.execute(
f"UPDATE collections SET {', '.join(sets)}, updated_at=CURRENT_TIMESTAMP WHERE id=?",
params,
)
conn.commit()
new_version = conn.execute(
"SELECT sync_version FROM collections WHERE id=?", (cid,)
).fetchone()["sync_version"]
return {"collection_id": cid, "server_version": new_version}
def _mut_collection_delete(self, user_id: int, payload: dict) -> dict:
cid = payload.get("collection_id")
with get_conn() as conn:
conn.execute("DELETE FROM collections WHERE id=?", (cid,))
conn.commit()
return {"collection_id": cid, "server_version": None}
# ── audit ──────────────────────────────────────────────────────────────
def _record(self, user_id: int, device_id: str, mtype: str, mut: dict,
client_ts: float, result: dict) -> None:
"""Persist every mutation in the server-side audit queue."""
try:
with get_conn() as conn:
conn.execute(
"""INSERT INTO offline_sync_queue
(user_id, device_id, type, payload, client_timestamp,
server_version, status, error)
VALUES (?, ?, ?, ?, ?, ?, ?, ?)""",
(user_id, device_id, mtype, json.dumps(mut),
client_ts, result.get("server_version"),
result.get("status", "synced"),
result.get("error") or (json.dumps(result["conflict"], ensure_ascii=False)
if result.get("conflict") else None)),
)
conn.commit()
except Exception: # noqa: BLE001 — audit must never break the batch
logger.warning("failed to record sync audit row", exc_info=True)
def _iso_epoch(value) -> float | None:
"""Best-effort ISO→epoch. SQLite CURRENT_TIMESTAMP → 'YYYY-MM-DD HH:MM:SS'."""
try:
if value is None:
return None
# append a timezone so strptime behaves on naive timestamps
text = str(value).replace("T", " ").split(".")[0]
parsed = _dt.datetime.strptime(text, "%Y-%m-%d %H:%M:%S")
return parsed.replace(tzinfo=_dt.UTC).timestamp()
except ValueError:
return None