Files
inkwell/core/src/sync/pull.rs
T
bvandeusenandClaude Opus 5.5 0101b05487
CI & Build / Build now, or wait for Android? (push) Successful in 6s
Android / Build, or is the channel already serving this? (push) Successful in 6s
CI & Build / Python tests (push) Successful in 13s
CI & Build / Python lint (push) Successful in 2s
Desktop (Tauri) / Build, or is the channel already serving this? (push) Successful in 3s
CI & Build / Web typecheck and unit tests (push) Successful in 13s
CI & Build / integration (push) Successful in 2m6s
CI & Build / Build & push image (push) Skipped
Android / Core and FFI clippy and tests (push) Successful in 1m18s
Desktop (Tauri) / Web tests, clippy, Rust tests and rustfmt (push) Failing after 3m19s
Desktop (Tauri) / Tauri desktop (Linux) (push) Skipped
Desktop (Tauri) / Windows installer (cross-compiled) (push) Skipped
Desktop (Tauri) / Update manifest (push) Skipped
Android / Kotlin + Rust (APK) (push) Successful in 9m29s
Android / Build the server image (push) Successful in 2s
Stored timestamps have one writer: local::iso and local::now
The store's time format (RFC 3339, UTC, milliseconds, Z) is what makes lexical
order chronological. It was spelled out ten times as
to_rfc3339_opts(SecondsFormat::Millis, true), with two private now() copies
(store, pull). local::iso(t) and local::now() now hold it; store, pull, engine
and portable call them.

DRY pass #2, batch 2, F5 (#5372).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-08 14:18:33 -04:00

997 lines
37 KiB
Rust

//! Pull: bring a server's changes into the local store (M10.7b).
//!
//! The feed is a single monotonic sequence shared by notes and labels, so one
//! integer cursor is a total-order watermark over both (docs/sync.md). We loop pages
//! until the server says there are no more, persisting the cursor **in the same
//! transaction** as the page it describes — a cursor committed ahead of its data
//! would silently skip those rows forever, which reads as a clean sync.
use rusqlite::{params, Connection, OptionalExtension};
use serde::Serialize;
use super::blobs::BlobStore;
use super::client;
use super::state;
use super::wire;
use crate::local::{now, Db};
/// Backstop against a server that never stops saying `has_more`. At the server's
/// 1000-row page cap this is 10M rows — far past any real store, so hitting it means
/// something is wrong, not that someone has a lot of notes.
const MAX_PAGES: usize = 10_000;
/// What a pull did — for the UI, and for the log when something looks off.
#[derive(Debug, Clone, Default, Serialize, PartialEq, Eq)]
pub struct PullSummary {
pub pages: usize,
pub notes_applied: usize,
pub notes_deleted: usize,
pub labels_applied: usize,
pub labels_deleted: usize,
pub cursor: i64,
/// Rows that still held unpushed local edits when the server's version landed on
/// top. Should be 0 in the normal cycle, because push runs first; anything higher
/// means local work was overwritten, which is worth saying out loud.
pub clobbered_dirty: usize,
pub blobs_downloaded: usize,
/// Attachments whose bytes couldn't be fetched or failed verification. Counted
/// rather than fatal — see `download_missing_blobs`.
pub blobs_failed: usize,
}
impl PullSummary {
fn absorb(&mut self, other: PullSummary) {
self.pages += other.pages;
self.notes_applied += other.notes_applied;
self.notes_deleted += other.notes_deleted;
self.labels_applied += other.labels_applied;
self.labels_deleted += other.labels_deleted;
self.clobbered_dirty += other.clobbered_dirty;
self.blobs_downloaded += other.blobs_downloaded;
self.blobs_failed += other.blobs_failed;
self.cursor = other.cursor;
}
}
/// `(note_id, attachment_id, sha256)` for every attachment that advertises a hash.
/// The caller filters against the blob store — which blobs we hold isn't a SQL
/// question.
pub fn hashed_attachments(conn: &Connection) -> rusqlite::Result<Vec<(String, String, String)>> {
let mut stmt = conn.prepare(
"SELECT note_id, id, sha256 FROM attachments
WHERE sha256 IS NOT NULL AND sha256 <> ''",
)?;
let rows = stmt.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?;
rows.collect()
}
/// Fetch the bytes for any attachment we have metadata for but no blob.
///
/// A failed attachment NEVER fails the sync. Notes are the primary data and they've
/// already landed; an image that didn't arrive is retried on the next cycle simply
/// because its blob still counts as missing. Aborting here would mean one unreachable
/// file could block every future sync.
async fn download_missing_blobs(
db: &Db,
blobs: &BlobStore,
base_url: &str,
token: &str,
) -> Result<(usize, usize), String> {
let wanted = {
let conn = db.conn()?;
hashed_attachments(&conn).map_err(|e| e.to_string())?
};
let mut downloaded = 0;
let mut failed = 0;
for (note_id, attachment_id, sha256) in wanted {
// Content-addressed, so this skips blobs we already hold — including the same
// image attached to a different note.
if blobs.has(&sha256) {
continue;
}
match client::fetch_attachment(base_url, token, &note_id, &attachment_id).await {
Ok(bytes) => match blobs.store(&sha256, &bytes) {
Ok(_) => downloaded += 1,
Err(e) => {
log::warn!("attachment {attachment_id}: {e}");
failed += 1;
}
},
Err(e) => {
log::warn!("attachment {attachment_id}: {e}");
failed += 1;
}
}
}
Ok((downloaded, failed))
}
/// Apply one page and advance the cursor, atomically.
///
/// Labels are applied before notes so a membership never references a label row that
/// doesn't exist yet.
pub fn apply_page(conn: &Connection, page: &wire::ChangesPage) -> rusqlite::Result<PullSummary> {
let tx = conn.unchecked_transaction()?;
let mut summary = PullSummary {
pages: 1,
cursor: page.cursor,
..Default::default()
};
for label in &page.labels {
if label.is_tombstone() {
tx.execute("DELETE FROM labels WHERE id = ?1", params![label.id])?;
summary.labels_deleted += 1;
} else {
upsert_label(&tx, label)?;
summary.labels_applied += 1;
}
}
for note in &page.notes {
if note.is_tombstone() {
// A purge tombstone carries no content — its only job is to say "delete
// your copy". Children go with it via ON DELETE CASCADE.
tx.execute("DELETE FROM notes WHERE id = ?1", params![note.id])?;
summary.notes_deleted += 1;
continue;
}
if is_dirty(&tx, &note.id)? {
summary.clobbered_dirty += 1;
}
upsert_note(&tx, note)?;
summary.notes_applied += 1;
}
// After the notes, never before: a page can carry a note AND its revocation only
// when the revocation is the newer of the two (sharing again deletes an older
// one on the server), so the revocation is the one that has to win. Only a note
// someone else owns can be revoked; this account's own are never touched.
for id in &page.revoked {
let removed = tx.execute(
"DELETE FROM notes WHERE id = ?1 AND permission <> 'owner'",
params![id],
)?;
summary.notes_deleted += removed;
}
state::set_cursor(&tx, page.cursor)?;
tx.commit()?;
Ok(summary)
}
fn is_dirty(conn: &Connection, note_id: &str) -> rusqlite::Result<bool> {
let dirty: Option<i64> = conn
.query_row(
"SELECT dirty FROM notes WHERE id = ?1",
params![note_id],
|r| r.get(0),
)
.optional()?;
Ok(dirty == Some(1))
}
fn upsert_label(conn: &Connection, label: &wire::Label) -> rusqlite::Result<()> {
// One label per name is enforced on both sides (locally a UNIQUE index on
// lower(name); on the server, per owner). A label created offline can therefore
// collide with one the server already had under a different id — "work" typed on
// this machine and "work" that already existed.
//
// The server's row wins, but its MEMBERSHIPS have to survive the swap. Just
// deleting the local duplicate would cascade its note_labels away, stripping the
// label off notes that this pull never even mentions — silent loss that no later
// page would repair. So: free the name, insert the server's row, re-point the
// memberships onto it, then drop the husk.
let duplicates: Vec<String> = {
let mut stmt =
conn.prepare("SELECT id FROM labels WHERE lower(name) = lower(?1) AND id <> ?2")?;
let rows = stmt.query_map(params![label.name, label.id], |r| r.get::<_, String>(0))?;
rows.collect::<rusqlite::Result<Vec<String>>>()?
};
// Renaming first is what makes the insert possible at all — the unique index
// would otherwise reject the server's row before anything could be merged.
for old in &duplicates {
conn.execute(
"UPDATE labels SET name = name || ' (superseded ' || id || ')' WHERE id = ?1",
params![old],
)?;
}
let created = label.created_at.clone().unwrap_or_else(now);
conn.execute(
"INSERT INTO labels (id, name, color, created_at, updated_at, sync_revision, dirty)
VALUES (?1, ?2, ?3, ?4, ?4, ?5, 0)
ON CONFLICT(id) DO UPDATE SET
name = excluded.name,
color = excluded.color,
sync_revision = excluded.sync_revision,
dirty = 0",
params![
label.id,
label.name,
label.color,
created,
label.sync_revision
],
)?;
for old in &duplicates {
// OR IGNORE guards a (note_id, label_id) collision. Today the unique index on
// lower(name) makes that unreachable — two same-name labels can't coexist
// locally — so this is belt-and-braces against that index changing, not a
// case we've seen. Anything it skips cascades away with the husk below, which
// is correct: those are duplicates of a membership that now exists.
conn.execute(
"UPDATE OR IGNORE note_labels SET label_id = ?1 WHERE label_id = ?2",
params![label.id, old],
)?;
conn.execute("DELETE FROM labels WHERE id = ?1", params![old])?;
}
Ok(())
}
fn upsert_note(conn: &Connection, note: &wire::Note) -> rusqlite::Result<()> {
let created = note.created_at.clone().unwrap_or_else(now);
let updated = note.updated_at.clone().unwrap_or_else(|| created.clone());
// The server's `deleted_at` is the authority on trash AGE. Taking it from the feed
// rather than stamping "now" locally is what keeps a note trashed three weeks ago
// from looking brand-new to a device that only just heard about it — otherwise
// every fresh install would silently reset the whole retention clock. Falls back
// to the note's updated_at only if an older server omits the field.
let trashed_at = if note.trashed {
note.deleted_at.clone().or_else(|| Some(updated.clone()))
} else {
None
};
// `created_at` is deliberately absent from the UPDATE clause: a note's birth time
// never changes, and the server's copy is the same value anyway.
// A server without `shares` sends no permission, and every note it sends is ours.
let permission = match note.permission.as_deref() {
Some("edit") => "edit",
Some("view") => "view",
_ => "owner",
};
let (shared_by_id, shared_by_name) = match &note.shared_by {
Some(by) => (Some(by.id.as_str()), Some(by.display_name.as_str())),
None => (None, None),
};
conn.execute(
"INSERT INTO notes (id, body, position, pinned, archived,
trashed, remind_at, recurrence, created_at, updated_at,
sync_revision, trashed_at, dirty,
permission, shared, shared_by_id, shared_by_name)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, 0, ?13, ?14, ?15, ?16)
ON CONFLICT(id) DO UPDATE SET
body = excluded.body,
position = excluded.position,
pinned = excluded.pinned,
archived = excluded.archived,
trashed = excluded.trashed,
remind_at = excluded.remind_at,
recurrence = excluded.recurrence,
updated_at = excluded.updated_at,
sync_revision = excluded.sync_revision,
trashed_at = excluded.trashed_at,
dirty = 0,
permission = excluded.permission,
shared = excluded.shared,
shared_by_id = excluded.shared_by_id,
shared_by_name = excluded.shared_by_name",
params![
note.id,
note.body,
note.position,
note.pinned,
note.archived,
note.trashed,
note.remind_at,
note.recurrence,
created,
updated,
note.sync_revision,
trashed_at,
permission,
note.shared,
shared_by_id,
shared_by_name,
],
)?;
// Children are replaced wholesale: a delta carries the note's FULL current state,
// so "what the server sent" IS the complete set. Diffing would be more code and
// could leave behind a row the server no longer has.
replace_attachments(conn, note)?;
replace_previews(conn, note)?;
replace_labels(conn, note)?;
Ok(())
}
/// Whether this device removed the row and the server hasn't acknowledged it yet.
/// Such a row is still on the server, so it is still in the feed — and putting it
/// back would undo a removal that is only waiting for the next push.
fn removed_here(conn: &Connection, entity: &str, id: &str) -> rusqlite::Result<bool> {
let found: Option<i64> = conn
.query_row(
"SELECT 1 FROM pending_deletes WHERE entity = ?1 AND id = ?2",
params![entity, id],
|r| r.get(0),
)
.optional()?;
Ok(found.is_some())
}
fn replace_attachments(conn: &Connection, note: &wire::Note) -> rusqlite::Result<()> {
// Only rows the server already had are replaced. A file attached here and still
// waiting to upload exists nowhere else yet, and deleting it would lose it; once
// it has gone up, the server's copy arrives under the same id and takes its place.
conn.execute(
"DELETE FROM attachments WHERE note_id = ?1 AND uploaded = 1",
params![note.id],
)?;
for (index, att) in note.attachments.iter().enumerate() {
if removed_here(conn, "attachment", &att.id)? {
continue;
}
// The server listing an id this device is still waiting to upload means the
// upload landed and only its reply was lost. The server's row replaces it.
conn.execute(
"DELETE FROM attachments WHERE id = ?1 AND uploaded = 0",
params![att.id],
)?;
// The feed carries no explicit position for attachments — they arrive in
// creation order, so the index preserves it. A plain INSERT, so an id listed
// twice fails the page rather than being quietly merged.
conn.execute(
"INSERT INTO attachments (id, note_id, url, filename, mime, size, sha256, position)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
att.id,
note.id,
att.url,
att.filename,
att.mime,
att.size,
att.sha256,
index as i64
],
)?;
}
Ok(())
}
fn replace_previews(conn: &Connection, note: &wire::Note) -> rusqlite::Result<()> {
conn.execute(
"DELETE FROM link_previews WHERE note_id = ?1",
params![note.id],
)?;
for (index, preview) in note.previews.iter().enumerate() {
if removed_here(conn, "preview", &preview.id)? {
continue;
}
conn.execute(
"INSERT INTO link_previews (id, note_id, url, title, description, image_url,
site_name, position)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
preview.id,
note.id,
preview.url,
preview.title,
preview.description,
preview.image_url,
preview.site_name,
index as i64
],
)?;
}
Ok(())
}
fn replace_labels(conn: &Connection, note: &wire::Note) -> rusqlite::Result<()> {
conn.execute(
"DELETE FROM note_labels WHERE note_id = ?1",
params![note.id],
)?;
for label in &note.labels {
ensure_label_stub(conn, label)?;
// `via_tag` is applied verbatim rather than re-derived from the body. The
// server already reconciled tags when it saved the note, and re-deriving here
// would call the local find-or-create path, which marks new labels dirty and
// would push them straight back — sync churn out of nothing.
conn.execute(
"INSERT OR IGNORE INTO note_labels (note_id, label_id, via_tag)
VALUES (?1, ?2, ?3)",
params![note.id, label.id, label.via_tag],
)?;
}
Ok(())
}
/// Materialize a label referenced by a note, if we don't have it yet.
///
/// Notes and labels page from one shared sequence, so a note can reference a label
/// whose own delta landed in an earlier page — or, right at a page boundary, hasn't
/// landed. The note carries enough of the label to create it, so a membership never
/// fails on a missing row. `OR IGNORE` because the label's real delta (later in this
/// page or a future one) is the authority on its name and color.
fn ensure_label_stub(conn: &Connection, label: &wire::NoteLabel) -> rusqlite::Result<()> {
let ts = now();
conn.execute(
"INSERT OR IGNORE INTO labels (id, name, color, created_at, updated_at, dirty)
VALUES (?1, ?2, ?3, ?4, ?4, 0)",
params![label.id, label.name, label.color, ts],
)?;
Ok(())
}
/// Loop the feed to exhaustion, starting from the persisted cursor.
///
/// NOTE ON ORDERING: the full cycle is push-then-pull (docs/sync.md). Running this
/// against a store with unpushed edits lets the server's version land on top of them
/// — counted as `clobbered_dirty` and logged, rather than hidden.
pub async fn run(
db: &Db,
blobs: &BlobStore,
base_url: &str,
token: &str,
shares: bool,
) -> Result<PullSummary, String> {
let mut total = PullSummary::default();
loop {
let since = {
let conn = db.conn()?;
state::read(&conn).map_err(|e| e.to_string())?.last_cursor
};
let page = client::fetch_changes(base_url, token, since, shares).await?;
// Trust the data over the flag: a server that claims more pages without
// advancing the cursor would spin this loop forever.
if page.has_more && page.cursor <= since {
return Err(format!(
"The server reported more changes but its cursor didn't advance past \
{since}. Stopping rather than looping forever."
));
}
let has_more = page.has_more;
let applied = {
let conn = db.conn()?;
apply_page(&conn, &page).map_err(|e| e.to_string())?
};
total.absorb(applied);
if !has_more {
break;
}
if total.pages >= MAX_PAGES {
return Err(format!(
"Stopped after {MAX_PAGES} pages without reaching the end of the \
server's changes. Something is wrong with the feed."
));
}
}
// Notes first, bytes after: the metadata is what makes the attachments knowable,
// and knowing one is missing is what lets the next cycle retry it.
let (downloaded, failed) = download_missing_blobs(db, blobs, base_url, token).await?;
total.blobs_downloaded = downloaded;
total.blobs_failed = failed;
if total.clobbered_dirty > 0 {
log::warn!(
"pull overwrote {} note(s) that still had unpushed local edits",
total.clobbered_dirty
);
}
if total.blobs_failed > 0 {
log::warn!(
"pull: {} attachment(s) couldn't be downloaded; will retry next sync",
total.blobs_failed
);
}
log::info!(
"pull complete: {} page(s), {} note(s) applied, {} deleted, {} label(s) applied, cursor {}",
total.pages,
total.notes_applied,
total.notes_deleted,
total.labels_applied,
total.cursor
);
Ok(total)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::local::schema;
fn db() -> Connection {
let conn = Connection::open_in_memory().expect("in-memory db");
schema::migrate(&conn).expect("migrate");
conn
}
fn note(id: &str, revision: i64) -> wire::Note {
wire::Note {
id: id.to_string(),
body: "Body".into(),
position: 0,
pinned: false,
archived: false,
trashed: false,
deleted_at: None,
remind_at: None,
recurrence: None,
created_at: Some("2026-07-26T00:00:00.000Z".into()),
updated_at: Some("2026-07-26T00:00:00.000Z".into()),
sync_revision: revision,
purged_at: None,
labels: vec![],
attachments: vec![],
previews: vec![],
permission: None,
shared: false,
shared_by: None,
}
}
fn attachment(id: &str) -> wire::Attachment {
wire::Attachment {
id: id.to_string(),
url: "/blob/x".into(),
filename: None,
mime: "image/png".into(),
size: None,
sha256: None,
}
}
fn page(notes: Vec<wire::Note>, labels: Vec<wire::Label>, cursor: i64) -> wire::ChangesPage {
wire::ChangesPage {
notes,
labels,
cursor,
has_more: false,
revoked: vec![],
}
}
fn count(conn: &Connection, sql: &str) -> i64 {
conn.query_row(sql, [], |r| r.get(0)).expect("count")
}
fn trash_stamp(conn: &Connection, id: &str) -> Option<String> {
let sql = "SELECT trashed_at FROM notes WHERE id = ?1";
conn.query_row(sql, [id], |r| r.get(0)).expect("stamp")
}
#[test]
fn applies_a_note_and_advances_the_cursor() {
let conn = db();
let summary = apply_page(&conn, &page(vec![note("n1", 7)], vec![], 7)).expect("apply");
assert_eq!(summary.notes_applied, 1);
assert_eq!(count(&conn, "SELECT COUNT(*) FROM notes"), 1);
assert_eq!(state::read(&conn).expect("state").last_cursor, 7);
}
#[test]
fn pulled_rows_are_not_dirty() {
// They came FROM the server, so pushing them back would be pure churn.
let conn = db();
apply_page(&conn, &page(vec![note("n1", 1)], vec![], 1)).expect("apply");
assert_eq!(count(&conn, "SELECT dirty FROM notes WHERE id = 'n1'"), 0);
}
#[test]
fn tombstone_deletes_the_local_note() {
let conn = db();
apply_page(&conn, &page(vec![note("n1", 1)], vec![], 1)).expect("apply");
let mut dead = note("n1", 2);
dead.purged_at = Some("2026-07-26T01:00:00.000Z".into());
let summary = apply_page(&conn, &page(vec![dead], vec![], 2)).expect("apply");
assert_eq!(summary.notes_deleted, 1);
assert_eq!(count(&conn, "SELECT COUNT(*) FROM notes"), 0);
}
fn shared_note(id: &str, revision: i64, permission: &str) -> wire::Note {
let mut n = note(id, revision);
n.permission = Some(permission.into());
n.shared = true;
n.shared_by = Some(wire::SharedBy {
id: "u-owner".into(),
display_name: "Robin".into(),
});
n
}
#[test]
fn a_shared_note_lands_saying_who_shared_it_and_how() {
let conn = db();
apply_page(&conn, &page(vec![shared_note("n1", 1, "edit")], vec![], 1)).expect("apply");
let held: (String, i64, String) = conn
.query_row(
"SELECT permission, shared, shared_by_name FROM notes WHERE id = 'n1'",
[],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.expect("row");
assert_eq!(held, ("edit".to_string(), 1, "Robin".to_string()));
// A server without shares sends no permission: the note is this account's own.
apply_page(&conn, &page(vec![note("n2", 2)], vec![], 2)).expect("apply");
assert_eq!(
count(
&conn,
"SELECT COUNT(*) FROM notes WHERE id = 'n2' AND permission = 'owner'"
),
1
);
}
#[test]
fn a_revoked_note_leaves_and_an_owned_one_never_does() {
let conn = db();
apply_page(
&conn,
&page(
vec![shared_note("theirs", 1, "view"), note("mine", 2)],
vec![],
2,
),
)
.expect("apply");
let mut revoked = page(vec![], vec![], 3);
revoked.revoked = vec!["theirs".into(), "mine".into(), "unknown".into()];
let summary = apply_page(&conn, &revoked).expect("apply");
assert_eq!(summary.notes_deleted, 1);
let left: Vec<String> = {
let mut stmt = conn.prepare("SELECT id FROM notes").unwrap();
let rows = stmt.query_map([], |r| r.get(0)).unwrap();
rows.collect::<rusqlite::Result<_>>().unwrap()
};
assert_eq!(left, ["mine"]);
}
#[test]
fn a_revocation_beside_its_note_in_one_page_wins() {
// The server deletes an older revocation when it shares again, so a page holds
// both only when the revocation is the newer: the note must not survive it.
let conn = db();
let mut both = page(vec![shared_note("n1", 4, "view")], vec![], 5);
both.revoked = vec!["n1".into()];
apply_page(&conn, &both).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM notes"), 0);
}
#[test]
fn trashed_is_not_a_tombstone() {
// `trashed` is ordinary state that keeps syncing; only `purged_at` deletes.
let conn = db();
let mut trashed = note("n1", 1);
trashed.trashed = true;
apply_page(&conn, &page(vec![trashed], vec![], 1)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM notes"), 1);
assert_eq!(count(&conn, "SELECT trashed FROM notes WHERE id = 'n1'"), 1);
}
#[test]
fn trash_age_comes_from_the_server_not_from_now() {
// The retention countdown runs off this timestamp. Stamping it locally would
// hand every note a fresh 30 days on any device that syncs it for the first
// time — a note trashed last month would never expire anywhere.
let conn = db();
let mut trashed = note("n1", 1);
trashed.trashed = true;
trashed.deleted_at = Some("2026-06-01T09:30:00+00:00".into());
apply_page(&conn, &page(vec![trashed], vec![], 1)).expect("apply");
let stamped = trash_stamp(&conn, "n1");
assert_eq!(stamped.as_deref(), Some("2026-06-01T09:30:00+00:00"));
}
#[test]
fn restoring_a_note_server_side_clears_its_trash_stamp() {
let conn = db();
let mut trashed = note("n1", 1);
trashed.trashed = true;
trashed.deleted_at = Some("2026-06-01T09:30:00+00:00".into());
apply_page(&conn, &page(vec![trashed], vec![], 1)).expect("apply");
apply_page(&conn, &page(vec![note("n1", 2)], vec![], 2)).expect("apply");
let stamped = trash_stamp(&conn, "n1");
assert_eq!(stamped, None, "an untrashed note keeps no trash stamp");
}
#[test]
fn an_older_server_without_deleted_at_still_ages_the_trash() {
// Falls back to updated_at rather than leaving the stamp null, which would
// make the note un-expirable and its countdown blank.
let conn = db();
let mut trashed = note("n1", 1);
trashed.trashed = true;
trashed.deleted_at = None;
apply_page(&conn, &page(vec![trashed], vec![], 1)).expect("apply");
let stamped = trash_stamp(&conn, "n1");
assert_eq!(stamped.as_deref(), Some("2026-07-26T00:00:00.000Z"));
}
#[test]
fn children_are_replaced_not_merged() {
// Was written over checklist items; they are lines of the body now (M304), so
// attachments carry the point instead. It is the same property either way: a
// delta is the note's FULL current state, so a child the server dropped has to
// disappear locally rather than linger.
let conn = db();
let mut first = note("n1", 1);
first.attachments = vec![attachment("a1"), attachment("a2")];
apply_page(&conn, &page(vec![first], vec![], 1)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM attachments"), 2);
let mut second = note("n1", 2);
second.attachments = vec![attachment("a1")];
apply_page(&conn, &page(vec![second], vec![], 2)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM attachments"), 1);
}
#[test]
fn note_label_membership_materializes_a_missing_label() {
// The label's own delta may have landed in an earlier page, or not yet.
let conn = db();
let mut n = note("n1", 1);
n.labels = vec![wire::NoteLabel {
id: "l1".into(),
name: "work".into(),
color: "blue".into(),
via_tag: true,
}];
apply_page(&conn, &page(vec![n], vec![], 1)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM labels"), 1);
assert_eq!(
count(
&conn,
"SELECT via_tag FROM note_labels WHERE note_id = 'n1'"
),
1,
"via_tag is applied verbatim, not re-derived"
);
}
#[test]
fn server_label_replaces_a_local_duplicate_by_name() {
let conn = db();
conn.execute(
"INSERT INTO labels (id, name, color, created_at, updated_at, dirty)
VALUES ('local-id', 'Work', 'default', '2026-01-01', '2026-01-01', 1)",
[],
)
.expect("seed local label");
let server = wire::Label {
id: "server-id".into(),
name: "work".into(),
color: "blue".into(),
sync_revision: 5,
purged_at: None,
created_at: Some("2026-07-26T00:00:00.000Z".into()),
};
apply_page(&conn, &page(vec![], vec![server], 5)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM labels"), 1);
let id: String = conn
.query_row("SELECT id FROM labels", [], |r| r.get(0))
.expect("label");
assert_eq!(id, "server-id", "the server's row wins on pull");
}
#[test]
fn merging_a_duplicate_label_keeps_its_note_memberships() {
// The notes carrying the local label may not be in this page at all, so a
// plain delete would strip the label off them with nothing to repair it.
let conn = db();
apply_page(&conn, &page(vec![note("n1", 1)], vec![], 1)).expect("seed note");
conn.execute(
"INSERT INTO labels (id, name, color, created_at, updated_at, dirty)
VALUES ('local-id', 'Work', 'default', '2026-01-01', '2026-01-01', 1)",
[],
)
.expect("seed local label");
conn.execute(
"INSERT INTO note_labels (note_id, label_id, via_tag)
VALUES ('n1', 'local-id', 0)",
[],
)
.expect("seed membership");
let server = wire::Label {
id: "server-id".into(),
name: "work".into(),
color: "blue".into(),
sync_revision: 5,
purged_at: None,
created_at: None,
};
apply_page(&conn, &page(vec![], vec![server], 5)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM labels"), 1);
let label_id: String = conn
.query_row(
"SELECT label_id FROM note_labels WHERE note_id = 'n1'",
[],
|r| r.get(0),
)
.expect("membership survived");
assert_eq!(label_id, "server-id", "membership re-pointed, not dropped");
}
#[test]
fn label_tombstone_deletes_and_cascades_memberships() {
let conn = db();
let mut n = note("n1", 1);
n.labels = vec![wire::NoteLabel {
id: "l1".into(),
name: "work".into(),
color: "blue".into(),
via_tag: false,
}];
apply_page(&conn, &page(vec![n], vec![], 1)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM note_labels"), 1);
let dead = wire::Label {
id: "l1".into(),
name: "work".into(),
color: "blue".into(),
sync_revision: 2,
purged_at: Some("2026-07-26T01:00:00.000Z".into()),
created_at: None,
};
apply_page(&conn, &page(vec![], vec![dead], 2)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM labels"), 0);
assert_eq!(
count(&conn, "SELECT COUNT(*) FROM note_labels"),
0,
"membership should cascade with the label"
);
}
#[test]
fn overwriting_a_dirty_note_is_counted() {
let conn = db();
conn.execute(
"INSERT INTO notes (id, body, created_at, updated_at, dirty)
VALUES ('n1', 'local edit', '2026-01-01', '2026-01-01', 1)",
[],
)
.expect("seed dirty note");
let summary = apply_page(&conn, &page(vec![note("n1", 9)], vec![], 9)).expect("apply");
assert_eq!(summary.clobbered_dirty, 1);
}
#[test]
fn applying_a_fresh_note_reports_no_clobber() {
let conn = db();
let summary = apply_page(&conn, &page(vec![note("n1", 1)], vec![], 1)).expect("apply");
assert_eq!(summary.clobbered_dirty, 0);
}
#[test]
fn empty_page_still_advances_the_cursor() {
// The server can page past rows that were trimmed to the shared watermark.
let conn = db();
apply_page(&conn, &page(vec![], vec![], 42)).expect("apply");
assert_eq!(state::read(&conn).expect("state").last_cursor, 42);
}
#[test]
fn note_upsert_preserves_the_original_created_at() {
let conn = db();
apply_page(&conn, &page(vec![note("n1", 1)], vec![], 1)).expect("apply");
let mut later = note("n1", 2);
later.created_at = Some("2099-01-01T00:00:00.000Z".into());
apply_page(&conn, &page(vec![later], vec![], 2)).expect("apply");
let created: String = conn
.query_row("SELECT created_at FROM notes WHERE id = 'n1'", [], |r| {
r.get(0)
})
.expect("created_at");
assert_eq!(created, "2026-07-26T00:00:00.000Z");
}
#[test]
fn a_page_that_fails_leaves_the_cursor_untouched() {
// Atomicity is the whole resumability story: a cursor committed ahead of its
// data would skip those rows forever. Force a failure with a duplicate
// attachment id inside one page.
let conn = db();
let mut n = note("n1", 3);
n.attachments = vec![attachment("dup"), attachment("dup")];
assert!(apply_page(&conn, &page(vec![n], vec![], 3)).is_err());
assert_eq!(state::read(&conn).expect("state").last_cursor, 0);
assert_eq!(count(&conn, "SELECT COUNT(*) FROM notes"), 0);
}
fn preview(id: &str) -> wire::Preview {
wire::Preview {
id: id.to_string(),
url: "https://example.com/".into(),
title: None,
description: None,
image_url: None,
site_name: None,
}
}
fn queued_upload(conn: &Connection, id: &str, note_id: &str) {
conn.execute(
"INSERT INTO attachments (id, note_id, url, uploaded) VALUES (?1, ?2, '/x', 0)",
params![id, note_id],
)
.expect("queued upload");
}
fn attachment_ids(conn: &Connection) -> Vec<String> {
let mut stmt = conn
.prepare("SELECT id FROM attachments ORDER BY id")
.expect("prepare");
let rows = stmt.query_map([], |r| r.get(0)).expect("query");
rows.collect::<rusqlite::Result<_>>().expect("rows")
}
#[test]
fn a_pull_keeps_a_file_still_waiting_to_upload() {
// It exists only on this device until push sends it; a pull that replaced the
// note's attachments wholesale would delete the only copy.
let conn = db();
apply_page(&conn, &page(vec![note("n1", 1)], vec![], 1)).expect("apply");
queued_upload(&conn, "mine", "n1");
let mut changed = note("n1", 2);
changed.attachments = vec![attachment("theirs")];
apply_page(&conn, &page(vec![changed], vec![], 2)).expect("apply");
assert_eq!(attachment_ids(&conn), vec!["mine", "theirs"]);
}
#[test]
fn the_servers_copy_takes_over_from_an_upload_whose_reply_was_lost() {
let conn = db();
apply_page(&conn, &page(vec![note("n1", 1)], vec![], 1)).expect("apply");
queued_upload(&conn, "a1", "n1");
let mut listed = note("n1", 2);
listed.attachments = vec![attachment("a1")];
apply_page(&conn, &page(vec![listed], vec![], 2)).expect("apply");
assert_eq!(attachment_ids(&conn), vec!["a1"]);
let uploaded: bool = conn
.query_row(
"SELECT uploaded FROM attachments WHERE id = 'a1'",
[],
|r| r.get(0),
)
.expect("row");
assert!(uploaded, "not sent a second time");
}
#[test]
fn a_removal_waiting_to_be_pushed_is_not_undone_by_a_pull() {
let conn = db();
let mut first = note("n1", 1);
first.attachments = vec![attachment("a1")];
first.previews = vec![preview("p1")];
apply_page(&conn, &page(vec![first], vec![], 1)).expect("apply");
crate::local::store::delete_attachment(&conn, "n1", "a1").expect("remove");
crate::local::store::delete_preview(&conn, "n1", "p1").expect("dismiss");
// The server hasn't heard yet, so its copy of the note still lists both.
let mut stale = note("n1", 2);
stale.attachments = vec![attachment("a1")];
stale.previews = vec![preview("p1")];
apply_page(&conn, &page(vec![stale], vec![], 2)).expect("apply");
assert_eq!(count(&conn, "SELECT COUNT(*) FROM attachments"), 0);
assert_eq!(count(&conn, "SELECT COUNT(*) FROM link_previews"), 0);
}
}