"""Delta-sync API for the local-first native clients (M8 sync hub). Pull: `GET /api/sync/changes?since=` returns every note + label the caller owns whose sync_revision advanced past the cursor, newest-revision last, paginated. Notes and labels both draw from ONE shared sequence (sync_revision_seq), so the cursor is a single monotonic watermark across both entity types. The web app does NOT use this — it stays on the live REST API. This surface exists purely so a native client can mirror the server into its local store and resume from where it left off (since=0 = full initial sync). """ from __future__ import annotations import uuid from datetime import datetime, timezone from quart import Blueprint, g, jsonify, request from sqlalchemy import delete as sa_delete from sqlalchemy import func, select from .auth import login_required from .common import iso, parse_dt from .db import session_scope from .labeling import reconcile_manual_labels, resolve_owned_label_ids from .models.label import Label, NoteLabel from .models.note import Note from .models.note_revision import NoteRevision from .revisions import should_snapshot from .notes import ( _reconcile_tags, _serialize_notes, derive_display_title, normalize_color, normalize_recurrence, ) from .retention import purge_note from .serialize import serialize_label_sync from .unfurl_queue import schedule as schedule_unfurls bp = Blueprint("sync", __name__, url_prefix="/api/sync") DEFAULT_LIMIT = 500 MAX_LIMIT = 1000 MAX_PUSH = 1000 # per-batch change cap # --- protocol versioning (M10.6) -------------------------------------------- # # The client<->server compatibility contract. These integers version the WIRE # PROTOCOL, deliberately separate from the app's release version, so a client and # server on different releases can still work out whether they can talk. Without # that separation every protocol change would force app<->server lockstep. # # SYNC_PROTOCOL_VERSION what this server speaks. # MIN_CLIENT_PROTOCOL_VERSION the oldest client protocol it still accepts. # # Bump SYNC_PROTOCOL_VERSION for ANY wire change. Raise # MIN_CLIENT_PROTOCOL_VERSION only for a genuinely BREAKING one: it is the switch # that hard-blocks older clients, so additive changes must leave it alone. # v2 (M13): `kind` and `title` both left the wire. Dropping a field a v1 client sends # and expects back is breaking, so the FLOOR moves too — a v1 client would keep pushing # both and would read back notes carrying neither. # # One bump for the pair: they landed in the same protocol generation, and nothing ever # ran against a half-applied v2. SYNC_PROTOCOL_VERSION = 3 MIN_CLIENT_PROTOCOL_VERSION = 3 # Named capabilities beyond the base protocol. An ADDITIVE change earns a name # here rather than a min-version bump, so a newer client meeting an older server # can degrade to "some features unavailable" instead of refusing to sync. Clients # test for the name, never infer a capability from a version number — that's what # keeps feature gating independent of release lockstep. SYNC_FEATURES: tuple[str, ...] = ( "notes", # note delta sync (pull + push) "labels", # the label catalog as its own entity "attachments", # blob upload/download, deduped by sha256 "tombstones", # purge propagates as a content-less row "revisions", # an overwritten version snapshots into note history ) def protocol_advertisement() -> dict: """What the server publishes about the sync protocol, merged into `/api/config`. DB-free and unauthenticated on purpose: a client has to be able to ask "can I talk to you at all?" before it holds a device token — or even has an account. """ return { "sync_protocol_version": SYNC_PROTOCOL_VERSION, "min_client_protocol_version": MIN_CLIENT_PROTOCOL_VERSION, "sync_features": list(SYNC_FEATURES), } def _parse_since(raw: str | None) -> int: """The pull cursor: a non-negative revision watermark. Bad/absent → 0 (full sync).""" try: return max(int(raw), 0) if raw is not None else 0 except (ValueError, TypeError): return 0 def _clamp_limit(raw: str | None) -> int: try: return max(1, min(int(raw), MAX_LIMIT)) if raw is not None else DEFAULT_LIMIT except (ValueError, TypeError): return DEFAULT_LIMIT def _page_cursor(note_revs: list[int], label_revs: list[int], since: int, limit: int) -> tuple[int, bool]: """Compute the next cursor + has_more when paging TWO revision streams that share one sequence. Each stream is fetched `rev > since ORDER BY rev LIMIT limit`. If either stream came back FULL (== limit) we're truncating, so the safe cursor is the SMALLER of the two page boundaries — advancing only to where BOTH streams are fully drained, so nothing between the cursor and the next pull is skipped. If neither is full, everything ≤ max(returned) is drained. Both lists are ascending. """ boundaries = [] if len(note_revs) == limit: boundaries.append(note_revs[-1]) if len(label_revs) == limit: boundaries.append(label_revs[-1]) if boundaries: return min(boundaries), True all_revs = note_revs + label_revs return (max(all_revs) if all_revs else since), False @bp.get("/changes") @login_required async def changes(): since = _parse_since(request.args.get("since")) limit = _clamp_limit(request.args.get("limit")) async with session_scope() as db: # ALL of the owner's notes/labels (any state — active/archived/trash/purged), # since a client mirrors everything; ordered by the shared revision. note_rows = ( await db.scalars( select(Note) .where(Note.owner_id == g.user_id, Note.sync_revision > since) .order_by(Note.sync_revision) .limit(limit) ) ).all() label_rows = ( await db.scalars( select(Label) .where(Label.owner_id == g.user_id, Label.sync_revision > since) .order_by(Label.sync_revision) .limit(limit) ) ).all() cursor, has_more = _page_cursor( [n.sync_revision for n in note_rows], [lb.sync_revision for lb in label_rows], since, limit, ) # Trim each stream to the shared watermark so the two feeds stay aligned. note_rows = [n for n in note_rows if n.sync_revision <= cursor] label_rows = [lb for lb in label_rows if lb.sync_revision <= cursor] # Note bodies come from the shared note serializer; sync adds the two # delta-only fields on top (folding these into the serializer itself waits on # the notes.py serialization split). notes_out = await _serialize_notes(db, note_rows) for data, n in zip(notes_out, note_rows): data["sync_revision"] = n.sync_revision data["purged_at"] = iso(n.purged_at) return jsonify( { "notes": notes_out, "labels": [serialize_label_sync(lb) for lb in label_rows], "cursor": cursor, "has_more": has_more, } ) # --- Push: apply client mutations (LWW by client edit-time, non-destructive) --- def client_wins(client_edited_at: datetime | None, server_edited_at: datetime | None) -> bool: """Last-write-wins: the client's version is applied iff its edit-time is at least the server's. A missing client time never overwrites a real server edit; a missing server time (new/unknown row) always yields to a present client edit.""" if client_edited_at is None: return server_edited_at is None if server_edited_at is None: return True return client_edited_at >= server_edited_at def _assign_note_fields(note: Note, ch: dict) -> None: """Overwrite a note's scalar fields from a client's FULL-state change (sync is whole-note, not a partial patch — the client sends its authoritative version).""" note.body = ch["body"] if isinstance(ch.get("body"), str) else "" note.color = normalize_color(ch.get("color")) note.pinned = bool(ch.get("pinned")) note.archived = bool(ch.get("archived")) if ch.get("trashed"): if note.deleted_at is None: note.deleted_at = datetime.now(timezone.utc) else: note.deleted_at = None note.remind_at = parse_dt(ch.get("remind_at")) note.recurrence = normalize_recurrence(ch.get("recurrence")) if isinstance(ch.get("position"), int): note.position = ch["position"] # `_first_item_text` and `_apply_note_items` lived here until M304. Both existed for # one reason — a checklist was a table beside the body — and both are gone with it. A # pushed change carries its items as `- [ ] ` lines inside `body`, so applying them is # applying the body, and naming the note is reading its first line. A client that still # sends an `items` array is a v2 client, and the version floor below turns it away # before any of this runs. async def _apply_note_manual_labels(db, note: Note, ch: dict) -> None: """Set the note's MANUAL (picker) label memberships from client label_ids, leaving tag-sourced (via_tag) rows to _reconcile_tags. Only labels the caller owns count.""" raw = ch.get("label_ids") if not isinstance(raw, list): return wanted: set = set() for r in raw: try: wanted.add(uuid.UUID(str(r))) except (ValueError, TypeError): continue # sync is lenient: skip a malformed id rather than reject the push owned = await resolve_owned_label_ids(db, wanted, g.user_id) await reconcile_manual_labels(db, note, owned) async def _apply_note(db, ch: dict) -> dict: raw_id = ch.get("id") try: nid = uuid.UUID(str(raw_id)) except (ValueError, TypeError): return {"id": raw_id, "entity": "note", "status": "rejected", "error": "invalid id"} op = ch.get("op", "upsert") edited_at = parse_dt(ch.get("edited_at")) note = await db.scalar(select(Note).where(Note.id == nid)) if note is not None and note.owner_id != g.user_id: # A client only ever pushes ids of notes IT created, so this branch is only # reached by a probe (or a ~0-probability UUID collision). Reject with a # GENERIC message so the response doesn't confirm the id belongs to another # user (don't leak existence via a distinctive "not yours"). return {"id": str(nid), "entity": "note", "status": "rejected", "error": "cannot apply"} if op == "delete": if note is None: return {"id": str(nid), "entity": "note", "status": "noop"} if not client_wins(edited_at, note.updated_at): return {"id": str(nid), "entity": "note", "status": "kept", "sync_revision": note.sync_revision} await purge_note(db, note, edited_at) await db.flush() await db.refresh(note, ["sync_revision"]) return {"id": str(nid), "entity": "note", "status": "applied", "sync_revision": note.sync_revision} creating = note is None if creating: note = Note(id=nid, owner_id=g.user_id, body="", display_title="") created = parse_dt(ch.get("created_at")) if created is not None: note.created_at = created db.add(note) elif not client_wins(edited_at, note.updated_at): return {"id": str(nid), "entity": "note", "status": "kept", "sync_revision": note.sync_revision} elif note.purged_at is not None: note.purged_at = None # client re-created/edited → clear the tombstone old_body = note.body _assign_note_fields(note, ch) note.display_title = derive_display_title(note.body) if edited_at is not None: note.updated_at = edited_at # Non-destructive LWW: snapshot the overwritten server body into history — # subject to the same session window as a direct edit (revisions.should_snapshot). # This path is why the window is a time rule rather than a flag on the wire: a # client autosaving every second pushes a body change every second, and without # the check the SERVER would snapshot each one no matter how restrained the # client's own store was being. if not creating and await should_snapshot(db, note.id, old_body, note.body): db.add(NoteRevision(note_id=note.id, body=old_body)) await db.flush() # assign note.id before items/labels/links await _reconcile_tags(db, note) await _apply_note_manual_labels(db, note, ch) await db.flush() await db.refresh(note, ["sync_revision"]) # A note pushed from a linked client gets the same link previews as one typed into # the web app — the client picks them up on its next pull. Scheduled rather than # awaited: a push batch must not wait on somebody else's website. if creating or note.body != old_body: schedule_unfurls(note.id, note.body) return { "id": str(nid), "entity": "note", "status": "created" if creating else "applied", "sync_revision": note.sync_revision, } async def _apply_label(db, ch: dict) -> dict: raw_id = ch.get("id") try: lid = uuid.UUID(str(raw_id)) except (ValueError, TypeError): return {"id": raw_id, "entity": "label", "status": "rejected", "error": "invalid id"} op = ch.get("op", "upsert") edited_at = parse_dt(ch.get("edited_at")) label = await db.scalar(select(Label).where(Label.id == lid)) if label is not None and label.owner_id != g.user_id: # Generic rejection (see _apply_note): don't confirm a foreign-owned id exists. return {"id": str(lid), "entity": "label", "status": "rejected", "error": "cannot apply"} if op == "delete": if label is None: return {"id": str(lid), "entity": "label", "status": "noop"} if not client_wins(edited_at, label.updated_at): return {"id": str(lid), "entity": "label", "status": "kept", "sync_revision": label.sync_revision} await db.execute(sa_delete(NoteLabel).where(NoteLabel.label_id == label.id)) label.purged_at = datetime.now(timezone.utc) if edited_at is not None: label.updated_at = edited_at await db.flush() await db.refresh(label, ["sync_revision"]) return {"id": str(lid), "entity": "label", "status": "applied", "sync_revision": label.sync_revision} name = (ch.get("name") or "").strip() creating = label is None if not creating and not client_wins(edited_at, label.updated_at): return {"id": str(lid), "entity": "label", "status": "kept", "sync_revision": label.sync_revision} # Names are unique per owner — a same-name clash on a DIFFERENT id can't be an insert. if name: clash = await db.scalar( select(Label.id).where( Label.owner_id == g.user_id, func.lower(Label.name) == name.lower(), Label.id != lid ) ) if clash is not None: return {"id": str(lid), "entity": "label", "status": "rejected", "error": "name in use"} if creating: if not name: return {"id": str(lid), "entity": "label", "status": "rejected", "error": "name required"} label = Label(id=lid, owner_id=g.user_id, name=name, color=normalize_color(ch.get("color"))) db.add(label) else: if label.purged_at is not None: label.purged_at = None if name: label.name = name label.color = normalize_color(ch.get("color")) if edited_at is not None: label.updated_at = edited_at await db.flush() await db.refresh(label, ["sync_revision"]) return { "id": str(lid), "entity": "label", "status": "created" if creating else "applied", "sync_revision": label.sync_revision, } @bp.post("/push") @login_required async def push(): """Apply a batch of client changes. Additive + owner-scoped; LWW by client edit-time with a version-history snapshot on any overwrite (nothing is lost).""" body = await request.get_json(silent=True) or {} changes = body.get("changes") if not isinstance(changes, list): return jsonify({"error": "changes must be a list"}), 400 if len(changes) > MAX_PUSH: return jsonify({"error": f"too many changes in one push (max {MAX_PUSH})"}), 400 results = [] async with session_scope() as db: for ch in changes: if not isinstance(ch, dict): results.append({"status": "rejected", "error": "not an object"}) continue entity = ch.get("entity") if entity == "note": results.append(await _apply_note(db, ch)) elif entity == "label": results.append(await _apply_label(db, ch)) else: results.append({"id": ch.get("id"), "status": "rejected", "error": "unknown entity"}) await db.commit() return jsonify({"results": results})