attachments sync: attach offline, upload when linked, removals stick (#5168)
CI & Build / Python lint (push) Successful in 3s
CI & Build / Build now, or wait for Android? (push) Successful in 4s
Android / Build, or is the channel already serving this? (push) Successful in 4s
Desktop (Tauri) / Build, or is the channel already serving this? (push) Successful in 2s
CI & Build / Web typecheck and unit tests (push) Successful in 9s
CI & Build / Python tests (push) Successful in 11s
CI & Build / integration (push) Successful in 35s
CI & Build / Build & push image (push) Skipped
Desktop (Tauri) / Web tests, clippy, Rust tests and rustfmt (push) Failing after 2m23s
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 7m17s
CI & Build / Python lint (push) Successful in 3s
CI & Build / Build now, or wait for Android? (push) Successful in 4s
Android / Build, or is the channel already serving this? (push) Successful in 4s
Desktop (Tauri) / Build, or is the channel already serving this? (push) Successful in 2s
CI & Build / Web typecheck and unit tests (push) Successful in 9s
CI & Build / Python tests (push) Successful in 11s
CI & Build / integration (push) Successful in 35s
CI & Build / Build & push image (push) Skipped
Desktop (Tauri) / Web tests, clippy, Rust tests and rustfmt (push) Failing after 2m23s
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 7m17s
Desktop could not create an attachment at all, and a removed attachment or dismissed preview came back on the next pull. Now: - core: add_attachment keeps the bytes in the blob store and queues the row (schema v10: attachments.uploaded / upload_error). Push uploads it once its note has landed. A refusal that retrying won't fix (too large, id clash, hash mismatch) is recorded on the file and not re-sent every cycle; the editor shows it. - core: removing a synced attachment or dismissing a preview leaves a tombstone in pending_deletes; push sends it as an `attachment`/`preview` delete, and a pull while it waits doesn't put the row back. A pull also keeps files still waiting to upload instead of replacing them wholesale. - server: PUT /api/sync/attachments/<id> (raw body, sha256-checked, idempotent, size-capped) and child deletes in push, which apply regardless of LWW and answer noop for rows the caller can't see. One store_attachment helper for the upload route, the importer and sync. Protocol 5, feature attachment_sync; the client sends neither to a server without it. - server: migration 0031 makes a link preview's insert/delete bump its note, so background-fetched previews and web dismissals reach linked devices. - desktop: Attach and paste-image work offline (raw-bytes IPC command). - SVG is served as a download by the desktop blob scheme too (as #1981 did for the web), and drawn as a file chip on both. - autosync: drop the catch_unwind; release builds abort on panic, so it only ever worked in debug builds. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -55,6 +55,10 @@ pub struct Attachment {
|
||||
pub mime: String,
|
||||
pub size: Option<i64>,
|
||||
pub sha256: Option<String>,
|
||||
/// Why the server refused a file attached on this device, in words to show on it.
|
||||
/// Absent for everything else, including a file still waiting to upload.
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pub upload_error: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
|
||||
@@ -274,6 +274,23 @@ UPDATE saved_filters
|
||||
AND params LIKE '%"color"%';
|
||||
"#;
|
||||
|
||||
// v10 (#5168): an attachment can be created on THIS device.
|
||||
//
|
||||
// Until now every attachment row arrived on the delta feed, so every one was already on
|
||||
// the server. A file attached here, offline, is not, and `uploaded = 0` is the queue
|
||||
// push drains once the note has landed. The rows already here all came from the
|
||||
// server, which is why the default is 1.
|
||||
//
|
||||
// `upload_error` holds a refusal that retrying won't fix — over the size limit, an id
|
||||
// already in use, bytes that don't match their hash. A row carrying one leaves the
|
||||
// queue: a file refused for its size would otherwise be re-sent in full on every
|
||||
// cycle, every five minutes, for as long as the app was open. The message is what the
|
||||
// editor shows on the file instead.
|
||||
const SCHEMA_V10: &str = r#"
|
||||
ALTER TABLE attachments ADD COLUMN uploaded INTEGER NOT NULL DEFAULT 1;
|
||||
ALTER TABLE attachments ADD COLUMN upload_error TEXT;
|
||||
"#;
|
||||
|
||||
pub fn migrate(conn: &Connection) -> rusqlite::Result<()> {
|
||||
conn.execute_batch("PRAGMA foreign_keys = ON;")?;
|
||||
let version: i64 = conn.query_row("PRAGMA user_version", [], |r| r.get(0))?;
|
||||
@@ -313,6 +330,10 @@ pub fn migrate(conn: &Connection) -> rusqlite::Result<()> {
|
||||
conn.execute_batch(SCHEMA_V9)?;
|
||||
conn.execute_batch("PRAGMA user_version = 9;")?;
|
||||
}
|
||||
if version < 10 {
|
||||
conn.execute_batch(SCHEMA_V10)?;
|
||||
conn.execute_batch("PRAGMA user_version = 10;")?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -427,7 +448,34 @@ mod tests {
|
||||
let version: i64 = conn
|
||||
.query_row("PRAGMA user_version", [], |r| r.get(0))
|
||||
.expect("version");
|
||||
assert_eq!(version, 9);
|
||||
assert_eq!(version, 10);
|
||||
}
|
||||
|
||||
/// Every attachment that predates v10 came down the feed, so it is already on the
|
||||
/// server. Defaulting it to "waiting to upload" would re-send every file once.
|
||||
#[test]
|
||||
fn v10_counts_existing_attachments_as_already_on_the_server() {
|
||||
let conn = v7_db();
|
||||
migrate_v8(&conn).expect("v8");
|
||||
conn.execute_batch(SCHEMA_V9).expect("v9");
|
||||
conn.execute_batch("PRAGMA user_version = 9;").expect("v9");
|
||||
add_note(&conn, "n", "a note");
|
||||
conn.execute(
|
||||
"INSERT INTO attachments (id, note_id, url) VALUES ('a', 'n', '/x')",
|
||||
[],
|
||||
)
|
||||
.expect("seed");
|
||||
|
||||
migrate(&conn).expect("migrate");
|
||||
let (uploaded, error): (bool, Option<String>) = conn
|
||||
.query_row(
|
||||
"SELECT uploaded, upload_error FROM attachments WHERE id = 'a'",
|
||||
[],
|
||||
|r| Ok((r.get(0)?, r.get(1)?)),
|
||||
)
|
||||
.expect("row");
|
||||
assert!(uploaded);
|
||||
assert_eq!(error, None);
|
||||
}
|
||||
|
||||
/// The column is gone, not merely unread. Asserted by asking SQLite rather than by
|
||||
|
||||
+215
-2
@@ -92,7 +92,7 @@ fn items_of(body: &str) -> Vec<ChecklistItem> {
|
||||
|
||||
fn load_attachments(conn: &Connection, note_id: &str) -> rusqlite::Result<Vec<Attachment>> {
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT id, url, filename, mime, size, sha256 FROM attachments WHERE note_id = ?1 ORDER BY position ASC",
|
||||
"SELECT id, url, filename, mime, size, sha256, upload_error FROM attachments WHERE note_id = ?1 ORDER BY position ASC",
|
||||
)?;
|
||||
let rows = stmt.query_map([note_id], |r| {
|
||||
let server_url: String = r.get(1)?;
|
||||
@@ -117,6 +117,7 @@ fn load_attachments(conn: &Connection, note_id: &str) -> rusqlite::Result<Vec<At
|
||||
mime,
|
||||
size: r.get(4)?,
|
||||
sha256,
|
||||
upload_error: r.get(6)?,
|
||||
})
|
||||
})?;
|
||||
rows.collect()
|
||||
@@ -660,20 +661,110 @@ pub fn delete_item(conn: &Connection, id: &str, item_id: &str) -> rusqlite::Resu
|
||||
set_body(conn, id, derive::remove_item(&body, index))
|
||||
}
|
||||
|
||||
/// The stored form of a declared content type: the bare media type, lowercase.
|
||||
/// Matches the server's `normalize_mime`, so both sides file a type the same way.
|
||||
fn normalize_mime(raw: &str) -> String {
|
||||
let bare = raw
|
||||
.split(';')
|
||||
.next()
|
||||
.unwrap_or("")
|
||||
.trim()
|
||||
.to_ascii_lowercase();
|
||||
if bare.is_empty() {
|
||||
"application/octet-stream".to_string()
|
||||
} else {
|
||||
bare
|
||||
}
|
||||
}
|
||||
|
||||
/// A file's name reduced to its last path component, as the server's
|
||||
/// `_safe_filename` does — the name is for display and download, never a path.
|
||||
fn safe_filename(raw: &str) -> String {
|
||||
let base = raw.trim().replace('\\', "/");
|
||||
let base = base.rsplit('/').next().unwrap_or("").trim();
|
||||
let capped: String = base.chars().take(255).collect();
|
||||
if capped.is_empty() {
|
||||
"file".to_string()
|
||||
} else {
|
||||
capped
|
||||
}
|
||||
}
|
||||
|
||||
/// Attach a file on this device, linked or not.
|
||||
///
|
||||
/// The bytes go into the blob store under their hash, and the row waits with
|
||||
/// `uploaded = 0` until push sends it — after the note itself has landed, since the
|
||||
/// server files an attachment under its note. The note is touched, so it is dirty
|
||||
/// too: that is what a background sync keys on to send it promptly.
|
||||
pub fn add_attachment(
|
||||
conn: &Connection,
|
||||
blobs: &crate::sync::blobs::BlobStore,
|
||||
note_id: &str,
|
||||
filename: &str,
|
||||
mime: &str,
|
||||
bytes: &[u8],
|
||||
) -> Result<Note, String> {
|
||||
let exists: Option<i64> = conn
|
||||
.query_row("SELECT 1 FROM notes WHERE id = ?1", [note_id], |r| r.get(0))
|
||||
.optional()
|
||||
.map_err(|e| e.to_string())?;
|
||||
if exists.is_none() {
|
||||
return Err("That note no longer exists.".to_string());
|
||||
}
|
||||
let sha256 = blobs.put(bytes)?;
|
||||
let id = new_id();
|
||||
conn.execute(
|
||||
"INSERT INTO attachments (id, note_id, url, filename, mime, size, sha256, position, uploaded)
|
||||
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7,
|
||||
(SELECT COALESCE(MAX(position) + 1, 0) FROM attachments WHERE note_id = ?2),
|
||||
0)",
|
||||
params![
|
||||
id,
|
||||
note_id,
|
||||
// The server's route for it, which is what the row holds once it has
|
||||
// synced too. Nothing renders it: `load_attachments` serves the local bytes.
|
||||
format!("/api/notes/{note_id}/attachments/{id}"),
|
||||
safe_filename(filename),
|
||||
normalize_mime(mime),
|
||||
bytes.len() as i64,
|
||||
sha256,
|
||||
],
|
||||
)
|
||||
.map_err(|e| e.to_string())?;
|
||||
touch(conn, note_id).map_err(|e| e.to_string())?;
|
||||
load_note(conn, note_id).map_err(|e| e.to_string())
|
||||
}
|
||||
|
||||
pub fn delete_attachment(conn: &Connection, id: &str, att_id: &str) -> rusqlite::Result<Note> {
|
||||
let uploaded: Option<bool> = conn
|
||||
.query_row(
|
||||
"SELECT uploaded FROM attachments WHERE id = ?1 AND note_id = ?2",
|
||||
params![att_id, id],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.optional()?;
|
||||
conn.execute(
|
||||
"DELETE FROM attachments WHERE id = ?1 AND note_id = ?2",
|
||||
params![att_id, id],
|
||||
)?;
|
||||
// Only a file the server holds needs a tombstone; without one the next pull would
|
||||
// bring it straight back. One still waiting to upload leaves the queue with its row.
|
||||
if uploaded == Some(true) {
|
||||
record_pending_delete(conn, "attachment", att_id)?;
|
||||
}
|
||||
touch(conn, id)?;
|
||||
load_note(conn, id)
|
||||
}
|
||||
|
||||
pub fn delete_preview(conn: &Connection, id: &str, preview_id: &str) -> rusqlite::Result<Note> {
|
||||
conn.execute(
|
||||
let removed = conn.execute(
|
||||
"DELETE FROM link_previews WHERE id = ?1 AND note_id = ?2",
|
||||
params![preview_id, id],
|
||||
)?;
|
||||
// Previews are made by the server, so every one it has is one it would send back.
|
||||
if removed > 0 {
|
||||
record_pending_delete(conn, "preview", preview_id)?;
|
||||
}
|
||||
touch(conn, id)?;
|
||||
load_note(conn, id)
|
||||
}
|
||||
@@ -1460,4 +1551,126 @@ mod tests {
|
||||
trash(&conn, &binned.id).expect("trash");
|
||||
assert_eq!(list_labels(&conn).expect("labels")[0].count, Some(1));
|
||||
}
|
||||
|
||||
fn blobs(tag: &str) -> crate::sync::blobs::BlobStore {
|
||||
let dir = std::env::temp_dir().join(format!("ts-store-blobs-{}-{tag}", std::process::id()));
|
||||
let _ = std::fs::remove_dir_all(&dir);
|
||||
crate::sync::blobs::BlobStore::new(dir).expect("blobs")
|
||||
}
|
||||
|
||||
fn tombstones(conn: &Connection) -> Vec<(String, String)> {
|
||||
let mut stmt = conn
|
||||
.prepare("SELECT entity, id FROM pending_deletes ORDER BY entity, id")
|
||||
.expect("prepare");
|
||||
let rows = stmt
|
||||
.query_map([], |r| Ok((r.get(0)?, r.get(1)?)))
|
||||
.expect("query");
|
||||
rows.collect::<rusqlite::Result<_>>().expect("rows")
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_attached_file_is_kept_here_and_queued_to_upload() {
|
||||
let conn = db();
|
||||
let blobs = blobs("attach");
|
||||
let n = note(&conn, "with a receipt");
|
||||
conn.execute("UPDATE notes SET dirty = 0", [])
|
||||
.expect("clean");
|
||||
|
||||
let got = add_attachment(
|
||||
&conn,
|
||||
&blobs,
|
||||
&n.id,
|
||||
"C:\\scans\\receipt.PDF",
|
||||
"Application/PDF; name=x",
|
||||
b"%PDF",
|
||||
)
|
||||
.expect("attach");
|
||||
|
||||
let att = &got.attachments[0];
|
||||
assert_eq!(
|
||||
att.filename.as_deref(),
|
||||
Some("receipt.PDF"),
|
||||
"a name, not a path"
|
||||
);
|
||||
assert_eq!(att.mime, "application/pdf");
|
||||
assert_eq!(att.size, Some(4));
|
||||
let hash = att.sha256.clone().expect("hashed");
|
||||
assert!(blobs.has(&hash), "the bytes are on this device");
|
||||
assert!(att.url.contains(&hash), "and served from here: {}", att.url);
|
||||
assert_eq!(att.upload_error, None);
|
||||
let (uploaded, dirty): (bool, bool) = conn
|
||||
.query_row(
|
||||
"SELECT a.uploaded, n.dirty FROM attachments a JOIN notes n ON n.id = a.note_id",
|
||||
[],
|
||||
|r| Ok((r.get(0)?, r.get(1)?)),
|
||||
)
|
||||
.expect("row");
|
||||
assert!(!uploaded, "it waits for push");
|
||||
assert!(
|
||||
dirty,
|
||||
"the note is touched, which is what starts a background sync"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn attaching_to_a_note_that_is_gone_writes_nothing() {
|
||||
let conn = db();
|
||||
let blobs = blobs("gone");
|
||||
assert!(add_attachment(&conn, &blobs, "missing", "a.txt", "text/plain", b"hi").is_err());
|
||||
let rows: i64 = conn
|
||||
.query_row("SELECT COUNT(*) FROM attachments", [], |r| r.get(0))
|
||||
.expect("count");
|
||||
assert_eq!(rows, 0);
|
||||
let files = std::fs::read_dir(blobs.root()).expect("dir").count();
|
||||
assert_eq!(files, 0, "and no bytes were filed for it");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn only_a_file_the_server_holds_leaves_a_tombstone_when_removed() {
|
||||
let conn = db();
|
||||
let blobs = blobs("remove");
|
||||
let n = note(&conn, "two files");
|
||||
conn.execute(
|
||||
"INSERT INTO attachments (id, note_id, url) VALUES ('synced', ?1, '/x')",
|
||||
[&n.id],
|
||||
)
|
||||
.expect("synced row");
|
||||
let with_local =
|
||||
add_attachment(&conn, &blobs, &n.id, "new.txt", "text/plain", b"hi").expect("attach");
|
||||
let local = with_local
|
||||
.attachments
|
||||
.iter()
|
||||
.find(|a| a.id != "synced")
|
||||
.expect("local")
|
||||
.id
|
||||
.clone();
|
||||
|
||||
delete_attachment(&conn, &n.id, "synced").expect("remove synced");
|
||||
delete_attachment(&conn, &n.id, &local).expect("remove local");
|
||||
|
||||
// Without the tombstone the next pull would put the synced file straight back.
|
||||
// The local one never reached the server, so there is nothing to tell it.
|
||||
assert_eq!(
|
||||
tombstones(&conn),
|
||||
vec![("attachment".to_string(), "synced".to_string())]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn dismissing_a_preview_leaves_a_tombstone() {
|
||||
let conn = db();
|
||||
let n = note(&conn, "https://example.com");
|
||||
conn.execute(
|
||||
"INSERT INTO link_previews (id, note_id, url) VALUES ('p1', ?1, 'https://example.com')",
|
||||
[&n.id],
|
||||
)
|
||||
.expect("preview");
|
||||
|
||||
delete_preview(&conn, &n.id, "p1").expect("dismiss");
|
||||
delete_preview(&conn, &n.id, "not-there").expect("a miss is fine");
|
||||
assert_eq!(
|
||||
tombstones(&conn),
|
||||
vec![("preview".to_string(), "p1".to_string())]
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
+27
-5
@@ -78,6 +78,15 @@ impl BlobStore {
|
||||
Ok(path)
|
||||
}
|
||||
|
||||
/// File bytes this device produced (a file attached here) and return their hash.
|
||||
/// The hash is computed from the bytes, so unlike [`store`](Self::store) there is
|
||||
/// nothing to verify against.
|
||||
pub fn put(&self, bytes: &[u8]) -> Result<String, String> {
|
||||
let hash = digest(bytes);
|
||||
self.store(&hash, bytes)?;
|
||||
Ok(hash)
|
||||
}
|
||||
|
||||
pub fn read(&self, sha256: &str) -> Option<Vec<u8>> {
|
||||
fs::read(self.path(sha256)?).ok()
|
||||
}
|
||||
@@ -140,7 +149,9 @@ fn urlencode(value: &str) -> String {
|
||||
out
|
||||
}
|
||||
|
||||
fn urldecode(value: &str) -> String {
|
||||
/// Undo percent-encoding. Also used for the filename the desktop's attach command
|
||||
/// receives in a header, which can only carry ASCII.
|
||||
pub fn urldecode(value: &str) -> String {
|
||||
let bytes = value.as_bytes();
|
||||
let mut out: Vec<u8> = Vec::with_capacity(bytes.len());
|
||||
let mut i = 0;
|
||||
@@ -163,12 +174,14 @@ fn urldecode(value: &str) -> String {
|
||||
///
|
||||
/// The mime rides in the URL and this scheme is an origin of its own, so echoing an
|
||||
/// arbitrary type would let an attachment claiming `text/html` run as a document
|
||||
/// there. Echoing is safe only because of the FAMILY check: nothing starting with
|
||||
/// `image/` can name a scriptable type. Everything else is served as an opaque
|
||||
/// download — the right treatment for an arbitrary file regardless.
|
||||
/// there. Echoing is safe only for the families that can't carry script, and
|
||||
/// `image/` is not quite one of them: `image/svg+xml` is a document that runs its own
|
||||
/// `<script>`, so it is served as an opaque download like everything else unfamiliar.
|
||||
/// The web made the same exclusion for the same reason (#1981).
|
||||
fn content_type_for(mime: &str) -> String {
|
||||
const RENDERABLE: &[&str] = &["image/", "audio/", "video/"];
|
||||
let familiar = RENDERABLE.iter().any(|p| mime.starts_with(p)) || mime == "application/pdf";
|
||||
let familiar = (RENDERABLE.iter().any(|p| mime.starts_with(p)) && mime != "image/svg+xml")
|
||||
|| mime == "application/pdf";
|
||||
// A header value can't carry control characters, and a mime type has no business
|
||||
// being long — both would only arrive from a malformed or hostile feed.
|
||||
let printable = mime.len() <= 100 && mime.bytes().all(|b| b.is_ascii_graphic());
|
||||
@@ -295,6 +308,8 @@ mod tests {
|
||||
let opaque = "application/octet-stream";
|
||||
assert_eq!(content_type_for("text/html"), opaque);
|
||||
assert_eq!(content_type_for("application/javascript"), opaque);
|
||||
// An image family member that is really a document with script in it.
|
||||
assert_eq!(content_type_for("image/svg+xml"), opaque);
|
||||
assert_eq!(content_type_for(""), opaque);
|
||||
// A control character can't reach a header value even under a safe family.
|
||||
assert_eq!(content_type_for("image/png\r\nX-Evil: 1"), opaque);
|
||||
@@ -309,6 +324,13 @@ mod tests {
|
||||
assert!(body.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn put_files_bytes_under_their_own_hash() {
|
||||
let store = store("put");
|
||||
assert_eq!(store.put(b"hello").expect("put"), HELLO);
|
||||
assert_eq!(store.read(HELLO).as_deref(), Some(&b"hello"[..]));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_blob_reads_as_none() {
|
||||
let store = store("missing");
|
||||
|
||||
@@ -27,6 +27,11 @@ const REQUEST_TIMEOUT: Duration = Duration::from_secs(10);
|
||||
/// failing one at ten seconds would make a large store impossible to ever pull.
|
||||
const SYNC_TIMEOUT: Duration = Duration::from_secs(120);
|
||||
|
||||
/// One file going up can be far larger than a page of notes, and a timeout shorter
|
||||
/// than the transfer is not a failure that retrying ever fixes — the same upload
|
||||
/// would time out again every cycle.
|
||||
const UPLOAD_TIMEOUT: Duration = Duration::from_secs(600);
|
||||
|
||||
/// Shared by every call that presents a token, so a revoked one reads the same way
|
||||
/// wherever it surfaces.
|
||||
const TOKEN_REJECTED: &str = "This server rejected the device token — it may have been \
|
||||
@@ -314,6 +319,71 @@ pub async fn fetch_attachment(
|
||||
.map_err(|e| format!("Couldn't download an attachment from {base_url}: {e}"))
|
||||
}
|
||||
|
||||
/// Why an upload didn't land, split by whether trying again could help.
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub enum UploadError {
|
||||
/// Worth another attempt next cycle: the network, a server error, or a note the
|
||||
/// server hasn't got yet.
|
||||
Retry(String),
|
||||
/// The server refused this file and will refuse it the same way every time —
|
||||
/// too large, an id in use, bytes that don't match their hash.
|
||||
Refused(String),
|
||||
}
|
||||
|
||||
/// Upload one file attached on this device, under the id it was given here.
|
||||
///
|
||||
/// `PUT /api/sync/attachments/<id>` takes the raw bytes; idempotent, so a retry of
|
||||
/// an upload whose reply was lost answers 200 and stores nothing twice.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn upload_attachment(
|
||||
base_url: &str,
|
||||
token: &str,
|
||||
note_id: &str,
|
||||
attachment_id: &str,
|
||||
filename: &str,
|
||||
mime: &str,
|
||||
sha256: &str,
|
||||
bytes: Vec<u8>,
|
||||
) -> Result<(), UploadError> {
|
||||
let url = format!("{base_url}/api/sync/attachments/{attachment_id}");
|
||||
let client = http_with(UPLOAD_TIMEOUT).map_err(UploadError::Retry)?;
|
||||
let request = prepare(client.put(url), Some(token))
|
||||
.query(&[
|
||||
("note_id", note_id),
|
||||
("filename", filename),
|
||||
("sha256", sha256),
|
||||
])
|
||||
.header(reqwest::header::CONTENT_TYPE, mime)
|
||||
.body(bytes);
|
||||
let response = request
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| UploadError::Retry(describe_transport_error(base_url, &e)))?;
|
||||
|
||||
let status = response.status();
|
||||
if status.is_success() {
|
||||
return Ok(());
|
||||
}
|
||||
if status == StatusCode::UNAUTHORIZED {
|
||||
return Err(UploadError::Retry(TOKEN_REJECTED.to_string()));
|
||||
}
|
||||
// The server's own words, when it gave some: "file is too large (max 25 MB)" is
|
||||
// what a person can act on, and a bare status code is not.
|
||||
let reason = response
|
||||
.json::<serde_json::Value>()
|
||||
.await
|
||||
.ok()
|
||||
.and_then(|v| v.get("error").and_then(|e| e.as_str()).map(String::from))
|
||||
.unwrap_or_else(|| format!("HTTP {}", status.as_u16()));
|
||||
if status.is_client_error() && status != StatusCode::NOT_FOUND {
|
||||
Err(UploadError::Refused(reason))
|
||||
} else {
|
||||
// 404: the note isn't on the server yet (its push was rejected, say). 5xx:
|
||||
// the server's problem, and likely a passing one.
|
||||
Err(UploadError::Retry(reason))
|
||||
}
|
||||
}
|
||||
|
||||
/// Send a batch of changes and hand back the raw reply.
|
||||
///
|
||||
/// Returns text rather than parsed results so this module stays pure transport —
|
||||
|
||||
@@ -23,7 +23,9 @@ use std::sync::OnceLock;
|
||||
///
|
||||
/// v4 (M315): `color` left the note. NOT a floor raise on either side — see the note
|
||||
/// on [`MIN_SERVER_PROTOCOL_VERSION`].
|
||||
pub const CLIENT_PROTOCOL_VERSION: u32 = 4;
|
||||
/// v5 (#5168): files attached here upload, and removed attachments and previews are
|
||||
/// pushed. Additive: both go only to a server advertising `attachment_sync`.
|
||||
pub const CLIENT_PROTOCOL_VERSION: u32 = 5;
|
||||
|
||||
/// The oldest server protocol this client can drive — the symmetric half of the
|
||||
/// server's `min_client_protocol_version`.
|
||||
@@ -47,7 +49,8 @@ pub const REQUIRED_FEATURES: &[&str] = &["notes", "labels"];
|
||||
/// Capabilities whose absence costs a feature but not the link. Listing these
|
||||
/// explicitly (rather than diffing against whatever the server happens to send) is
|
||||
/// what lets the UI name exactly what the user will be missing.
|
||||
pub const OPTIONAL_FEATURES: &[&str] = &["attachments", "tombstones", "revisions"];
|
||||
pub const OPTIONAL_FEATURES: &[&str] =
|
||||
&["attachments", "tombstones", "revisions", "attachment_sync"];
|
||||
|
||||
/// The handshake fields of `GET /api/config`.
|
||||
///
|
||||
@@ -75,7 +78,7 @@ pub struct ServerInfo {
|
||||
}
|
||||
|
||||
impl ServerInfo {
|
||||
fn has_feature(&self, name: &str) -> bool {
|
||||
pub fn has_feature(&self, name: &str) -> bool {
|
||||
self.sync_features.iter().any(|f| f.as_str() == name)
|
||||
}
|
||||
|
||||
|
||||
+16
-10
@@ -38,7 +38,17 @@ pub async fn run_cycle(
|
||||
base_url: &str,
|
||||
token: &str,
|
||||
) -> Result<SyncOutcome, String> {
|
||||
let push = push::run(db, base_url, token).await?;
|
||||
// What this server can do, asked before anything is sent: files attached here,
|
||||
// and removed attachments and previews, go only to a server advertising
|
||||
// `attachment_sync`. An older one keeps them queued on this device until it is
|
||||
// updated, rather than refusing each one on every cycle. Best-effort — an
|
||||
// unanswered probe reads as "not advertised", and the cycle carries on.
|
||||
let server = super::client::probe(base_url).await.ok().map(|p| p.server);
|
||||
let attachment_sync = server
|
||||
.as_ref()
|
||||
.is_some_and(|s| s.has_feature("attachment_sync"));
|
||||
|
||||
let push = push::run(db, blobs, base_url, token, attachment_sync).await?;
|
||||
let pull = pull::run(db, blobs, base_url, token).await?;
|
||||
|
||||
if pull.clobbered_dirty > 0 {
|
||||
@@ -51,15 +61,11 @@ pub async fn run_cycle(
|
||||
);
|
||||
}
|
||||
|
||||
// While we're already talking to this server, re-read what it says about itself.
|
||||
// Today that's the trash-retention window the Trash view counts down against, and
|
||||
// it can change under us whenever an admin edits the setting. Best-effort on
|
||||
// purpose: a config blip must not fail a cycle whose actual work already
|
||||
// succeeded, and the stored value simply stays as it was.
|
||||
let retention = super::client::probe(base_url)
|
||||
.await
|
||||
.ok()
|
||||
.and_then(|p| p.server.trash_retention_days);
|
||||
// The same answer carries the trash-retention window the Trash view counts down
|
||||
// against, which can change whenever an admin edits the setting. Best-effort on
|
||||
// purpose: a config blip must not fail a cycle whose actual work succeeded, and
|
||||
// the stored value simply stays as it was.
|
||||
let retention = server.and_then(|s| s.trash_retention_days);
|
||||
|
||||
let status = {
|
||||
let conn = db.0.lock().map_err(|e| e.to_string())?;
|
||||
|
||||
+116
-2
@@ -281,14 +281,41 @@ fn upsert_note(conn: &Connection, note: &wire::Note) -> rusqlite::Result<()> {
|
||||
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",
|
||||
"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.
|
||||
// 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)",
|
||||
@@ -313,6 +340,9 @@ fn replace_previews(conn: &Connection, note: &wire::Note) -> rusqlite::Result<()
|
||||
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)
|
||||
@@ -778,4 +808,88 @@ mod tests {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
+343
-28
@@ -1,8 +1,9 @@
|
||||
//! Push: send local changes to the server and apply what it says (M10.7c).
|
||||
//!
|
||||
//! Two sources feed a push: rows flagged `dirty` (created or edited locally) and rows
|
||||
//! Three sources feed a push: rows flagged `dirty` (created or edited locally), rows
|
||||
//! in `pending_deletes` (permanently deleted locally — see `local::schema` v2 for why
|
||||
//! a delete needs its own record).
|
||||
//! a delete needs its own record), and files attached on this device that are still
|
||||
//! waiting to go up (`attachments.uploaded = 0`, schema v10).
|
||||
//!
|
||||
//! Sync is **whole-note**: an upsert carries the client's full current state, not a
|
||||
//! patch (docs/sync.md). The server resolves conflicts last-write-wins by the client's
|
||||
@@ -11,7 +12,8 @@
|
||||
use rusqlite::{params, Connection, OptionalExtension};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use super::client;
|
||||
use super::blobs::BlobStore;
|
||||
use super::client::{self, UploadError};
|
||||
use super::state;
|
||||
use crate::local::Db;
|
||||
|
||||
@@ -35,6 +37,11 @@ pub struct PushSummary {
|
||||
/// realistic case). Silently retrying forever would be the wrong shape.
|
||||
pub rejected: usize,
|
||||
pub errors: Vec<String>,
|
||||
/// Files attached on this device that reached the server this cycle.
|
||||
pub uploaded: usize,
|
||||
/// Files that didn't. Their reasons are in `errors`; one the server refused for
|
||||
/// good also carries its reason on the attachment and isn't tried again.
|
||||
pub upload_failed: usize,
|
||||
}
|
||||
|
||||
impl PushSummary {
|
||||
@@ -47,6 +54,8 @@ impl PushSummary {
|
||||
self.noop += other.noop;
|
||||
self.rejected += other.rejected;
|
||||
self.errors.extend(other.errors);
|
||||
self.uploaded += other.uploaded;
|
||||
self.upload_failed += other.upload_failed;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -136,9 +145,17 @@ pub struct PushResult {
|
||||
|
||||
/// Everything waiting to go up, oldest edit first so a truncated batch still makes
|
||||
/// forward progress in a sensible order.
|
||||
pub fn collect(conn: &Connection, limit: usize) -> rusqlite::Result<Vec<Change>> {
|
||||
///
|
||||
/// `attachment_sync` is whether the server takes attachment and preview removals.
|
||||
/// Without it they stay queued here, untouched, until a server that does — an older
|
||||
/// one would answer "unknown entity" and the removal would read as a failure.
|
||||
pub fn collect(
|
||||
conn: &Connection,
|
||||
limit: usize,
|
||||
attachment_sync: bool,
|
||||
) -> rusqlite::Result<Vec<Change>> {
|
||||
let mut out = Vec::new();
|
||||
collect_deletes(conn, &mut out, limit)?;
|
||||
collect_deletes(conn, &mut out, limit, attachment_sync)?;
|
||||
if out.len() < limit {
|
||||
collect_labels(conn, &mut out, limit)?;
|
||||
}
|
||||
@@ -148,11 +165,21 @@ pub fn collect(conn: &Connection, limit: usize) -> rusqlite::Result<Vec<Change>>
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
fn collect_deletes(conn: &Connection, out: &mut Vec<Change>, limit: usize) -> rusqlite::Result<()> {
|
||||
fn collect_deletes(
|
||||
conn: &Connection,
|
||||
out: &mut Vec<Change>,
|
||||
limit: usize,
|
||||
attachment_sync: bool,
|
||||
) -> rusqlite::Result<()> {
|
||||
// Filtered in SQL rather than skipped below, so held-back removals can't use up
|
||||
// the LIMIT and starve the deletes that can go.
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT entity, id, deleted_at FROM pending_deletes ORDER BY deleted_at LIMIT ?1",
|
||||
"SELECT entity, id, deleted_at FROM pending_deletes
|
||||
WHERE entity IN ('note', 'label')
|
||||
OR (?2 AND entity IN ('attachment', 'preview'))
|
||||
ORDER BY deleted_at LIMIT ?1",
|
||||
)?;
|
||||
let rows = stmt.query_map(params![limit as i64], |r| {
|
||||
let rows = stmt.query_map(params![limit as i64, attachment_sync], |r| {
|
||||
Ok((
|
||||
r.get::<_, String>(0)?,
|
||||
r.get::<_, String>(1)?,
|
||||
@@ -161,11 +188,13 @@ fn collect_deletes(conn: &Connection, out: &mut Vec<Change>, limit: usize) -> ru
|
||||
})?;
|
||||
for row in rows {
|
||||
let (entity, id, deleted_at) = row?;
|
||||
// Only 'note' and 'label' exist on the wire; anything else is a bug in a
|
||||
// writer, and shipping it would earn a blanket rejection for the batch.
|
||||
// Only these exist on the wire; anything else is a bug in a writer, and
|
||||
// shipping it would earn a rejection that no retry could clear.
|
||||
let entity: &'static str = match entity.as_str() {
|
||||
"note" => "note",
|
||||
"label" => "label",
|
||||
"attachment" => "attachment",
|
||||
"preview" => "preview",
|
||||
_ => continue,
|
||||
};
|
||||
out.push(Change::delete(entity, id, deleted_at));
|
||||
@@ -414,7 +443,9 @@ pub fn has_pending(conn: &Connection) -> rusqlite::Result<bool> {
|
||||
.query_row(
|
||||
"SELECT 1 FROM notes WHERE dirty = 1
|
||||
UNION ALL SELECT 1 FROM labels WHERE dirty = 1
|
||||
UNION ALL SELECT 1 FROM pending_deletes LIMIT 1",
|
||||
UNION ALL SELECT 1 FROM pending_deletes
|
||||
UNION ALL SELECT 1 FROM attachments WHERE uploaded = 0 AND upload_error IS NULL
|
||||
LIMIT 1",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)
|
||||
@@ -434,22 +465,141 @@ pub fn pending_fingerprint(conn: &Connection) -> rusqlite::Result<Option<String>
|
||||
let fingerprint: String = conn.query_row(
|
||||
"SELECT (SELECT COUNT(*) || ':' || IFNULL(MAX(updated_at), '') FROM notes WHERE dirty = 1)
|
||||
|| '|' || (SELECT COUNT(*) || ':' || IFNULL(MAX(updated_at), '') FROM labels WHERE dirty = 1)
|
||||
|| '|' || (SELECT COUNT(*) || ':' || IFNULL(MAX(deleted_at), '') FROM pending_deletes)",
|
||||
|| '|' || (SELECT COUNT(*) || ':' || IFNULL(MAX(deleted_at), '') FROM pending_deletes)
|
||||
|| '|' || (SELECT COUNT(*) FROM attachments WHERE uploaded = 0 AND upload_error IS NULL)",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)?;
|
||||
Ok((fingerprint != "0:|0:|0:").then_some(fingerprint))
|
||||
Ok((fingerprint != "0:|0:|0:|0").then_some(fingerprint))
|
||||
}
|
||||
|
||||
/// A file attached on this device, ready to go up.
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub struct PendingUpload {
|
||||
pub note_id: String,
|
||||
pub id: String,
|
||||
pub filename: String,
|
||||
pub mime: String,
|
||||
pub sha256: String,
|
||||
}
|
||||
|
||||
/// Files waiting to upload whose note the server already holds. A note still dirty
|
||||
/// after the push (its change was rejected) keeps its files back too: the server files
|
||||
/// an attachment under its note, and would answer 404 for one it doesn't have.
|
||||
pub fn pending_uploads(conn: &Connection) -> rusqlite::Result<Vec<PendingUpload>> {
|
||||
let mut stmt = conn.prepare(
|
||||
"SELECT a.note_id, a.id, IFNULL(a.filename, 'file'), a.mime, a.sha256
|
||||
FROM attachments a JOIN notes n ON n.id = a.note_id
|
||||
WHERE a.uploaded = 0 AND a.upload_error IS NULL AND a.sha256 IS NOT NULL
|
||||
AND n.dirty = 0
|
||||
ORDER BY a.note_id, a.position",
|
||||
)?;
|
||||
let rows = stmt.query_map([], |r| {
|
||||
Ok(PendingUpload {
|
||||
note_id: r.get(0)?,
|
||||
id: r.get(1)?,
|
||||
filename: r.get(2)?,
|
||||
mime: r.get(3)?,
|
||||
sha256: r.get(4)?,
|
||||
})
|
||||
})?;
|
||||
rows.collect()
|
||||
}
|
||||
|
||||
/// Record what became of one upload.
|
||||
fn settle_upload(
|
||||
conn: &Connection,
|
||||
id: &str,
|
||||
result: &Result<(), UploadError>,
|
||||
) -> rusqlite::Result<()> {
|
||||
match result {
|
||||
Ok(()) => conn.execute(
|
||||
"UPDATE attachments SET uploaded = 1, upload_error = NULL WHERE id = ?1",
|
||||
params![id],
|
||||
)?,
|
||||
Err(UploadError::Refused(reason)) => conn.execute(
|
||||
"UPDATE attachments SET upload_error = ?2 WHERE id = ?1",
|
||||
params![id, reason],
|
||||
)?,
|
||||
Err(UploadError::Retry(_)) => 0,
|
||||
};
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Send the bytes of every file attached here whose note has landed.
|
||||
///
|
||||
/// One file failing never fails the cycle, the same as a download: the notes have
|
||||
/// already gone, and one unreachable or oversized file must not hold up every sync
|
||||
/// after it. A refusal is recorded on the attachment instead, and a passing failure
|
||||
/// is simply tried again next cycle.
|
||||
async fn upload_pending(
|
||||
db: &Db,
|
||||
blobs: &BlobStore,
|
||||
base_url: &str,
|
||||
token: &str,
|
||||
) -> Result<PushSummary, String> {
|
||||
let wanted = {
|
||||
let conn = db.0.lock().map_err(|e| e.to_string())?;
|
||||
pending_uploads(&conn).map_err(|e| e.to_string())?
|
||||
};
|
||||
let mut summary = PushSummary::default();
|
||||
for upload in wanted {
|
||||
let result = match blobs.read(&upload.sha256) {
|
||||
Some(bytes) => {
|
||||
client::upload_attachment(
|
||||
base_url,
|
||||
token,
|
||||
&upload.note_id,
|
||||
&upload.id,
|
||||
&upload.filename,
|
||||
&upload.mime,
|
||||
&upload.sha256,
|
||||
bytes,
|
||||
)
|
||||
.await
|
||||
}
|
||||
// Nothing to send and nothing that could bring it back: the file was
|
||||
// only ever on this device.
|
||||
None => Err(UploadError::Refused(
|
||||
"the file's bytes are missing from this device".to_string(),
|
||||
)),
|
||||
};
|
||||
{
|
||||
let conn = db.0.lock().map_err(|e| e.to_string())?;
|
||||
settle_upload(&conn, &upload.id, &result).map_err(|e| e.to_string())?;
|
||||
}
|
||||
match result {
|
||||
Ok(()) => summary.uploaded += 1,
|
||||
Err(UploadError::Refused(reason)) | Err(UploadError::Retry(reason)) => {
|
||||
log::warn!("attachment {}: {reason}", upload.id);
|
||||
summary.upload_failed += 1;
|
||||
summary
|
||||
.errors
|
||||
.push(format!("{} didn't upload: {reason}", upload.filename));
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(summary)
|
||||
}
|
||||
|
||||
/// Send everything pending, in batches, applying each batch's results before the
|
||||
/// next is collected.
|
||||
pub async fn run(db: &Db, base_url: &str, token: &str) -> Result<PushSummary, String> {
|
||||
/// next is collected — then the files, once their notes are there to hold them.
|
||||
///
|
||||
/// `attachment_sync`: whether the server advertises it (see [`collect`]). Without
|
||||
/// it, uploads wait too.
|
||||
pub async fn run(
|
||||
db: &Db,
|
||||
blobs: &BlobStore,
|
||||
base_url: &str,
|
||||
token: &str,
|
||||
attachment_sync: bool,
|
||||
) -> Result<PushSummary, String> {
|
||||
let mut total = PushSummary::default();
|
||||
|
||||
loop {
|
||||
let batch = {
|
||||
let conn = db.0.lock().map_err(|e| e.to_string())?;
|
||||
collect(&conn, BATCH).map_err(|e| e.to_string())?
|
||||
collect(&conn, BATCH, attachment_sync).map_err(|e| e.to_string())?
|
||||
};
|
||||
if batch.is_empty() {
|
||||
break;
|
||||
@@ -477,6 +627,10 @@ pub async fn run(db: &Db, base_url: &str, token: &str) -> Result<PushSummary, St
|
||||
}
|
||||
}
|
||||
|
||||
if attachment_sync {
|
||||
total.absorb(upload_pending(db, blobs, base_url, token).await?);
|
||||
}
|
||||
|
||||
if total.rejected > 0 {
|
||||
log::warn!(
|
||||
"push: {} change(s) rejected by the server: {}",
|
||||
@@ -485,13 +639,16 @@ pub async fn run(db: &Db, base_url: &str, token: &str) -> Result<PushSummary, St
|
||||
);
|
||||
}
|
||||
log::info!(
|
||||
"push complete: {} sent ({} created, {} applied, {} kept, {} noop, {} rejected)",
|
||||
"push complete: {} sent ({} created, {} applied, {} kept, {} noop, {} rejected), \
|
||||
{} file(s) uploaded, {} not",
|
||||
total.sent,
|
||||
total.created,
|
||||
total.applied,
|
||||
total.kept,
|
||||
total.noop,
|
||||
total.rejected
|
||||
total.rejected,
|
||||
total.uploaded,
|
||||
total.upload_failed
|
||||
);
|
||||
Ok(total)
|
||||
}
|
||||
@@ -548,7 +705,7 @@ mod tests {
|
||||
let conn = db();
|
||||
seed_note(&conn, "clean", 0);
|
||||
seed_note(&conn, "dirty", 1);
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
assert_eq!(batch.len(), 1);
|
||||
assert_eq!(batch[0].id, "dirty");
|
||||
assert_eq!(batch[0].op, "upsert");
|
||||
@@ -573,7 +730,7 @@ mod tests {
|
||||
)
|
||||
.expect("seed membership");
|
||||
}
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
let note = batch.iter().find(|c| c.entity == "note").expect("note");
|
||||
assert_eq!(note.label_ids.as_deref(), Some(&["manual".to_string()][..]));
|
||||
}
|
||||
@@ -583,7 +740,7 @@ mod tests {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 0);
|
||||
store::delete_forever(&conn, "n1").expect("delete");
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
assert_eq!(batch.len(), 1);
|
||||
assert_eq!(batch[0].op, "delete");
|
||||
assert_eq!(batch[0].entity, "note");
|
||||
@@ -594,7 +751,7 @@ mod tests {
|
||||
fn applied_clears_dirty_and_records_the_revision() {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 1);
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
apply_results(&conn, &batch, &[ok("applied", Some(42))]).expect("apply");
|
||||
assert_eq!(dirty_count(&conn), 0);
|
||||
let rev: i64 = conn
|
||||
@@ -611,7 +768,7 @@ mod tests {
|
||||
// every time; the following pull adopts the server's version instead.
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 1);
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
let summary = apply_results(&conn, &batch, &[ok("kept", Some(99))]).expect("apply");
|
||||
assert_eq!(summary.kept, 1);
|
||||
assert_eq!(dirty_count(&conn), 0);
|
||||
@@ -625,7 +782,7 @@ mod tests {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 1);
|
||||
state::set_cursor(&conn, 100).expect("cursor");
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
apply_results(&conn, &batch, &[ok("kept", Some(40))]).expect("apply");
|
||||
assert_eq!(state::read(&conn).expect("state").last_cursor, 39);
|
||||
}
|
||||
@@ -635,7 +792,7 @@ mod tests {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 1);
|
||||
state::set_cursor(&conn, 10).expect("cursor");
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
apply_results(&conn, &batch, &[ok("kept", Some(40))]).expect("apply");
|
||||
assert_eq!(
|
||||
state::read(&conn).expect("state").last_cursor,
|
||||
@@ -648,7 +805,7 @@ mod tests {
|
||||
fn rejected_stays_dirty_and_is_reported() {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 1);
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
let mut bad = ok("rejected", None);
|
||||
bad.error = Some("name in use".into());
|
||||
let summary = apply_results(&conn, &batch, &[bad]).expect("apply");
|
||||
@@ -662,7 +819,7 @@ mod tests {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 0);
|
||||
store::delete_forever(&conn, "n1").expect("delete");
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
apply_results(&conn, &batch, &[ok("applied", Some(7))]).expect("apply");
|
||||
assert!(!has_pending(&conn).expect("pending"));
|
||||
}
|
||||
@@ -673,7 +830,7 @@ mod tests {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n1", 1);
|
||||
store::delete_forever(&conn, "n1").expect("delete");
|
||||
let batch = collect(&conn, 100).expect("collect");
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
apply_results(&conn, &batch, &[ok("noop", None)]).expect("apply");
|
||||
assert!(!has_pending(&conn).expect("pending"));
|
||||
}
|
||||
@@ -761,4 +918,162 @@ mod tests {
|
||||
conn.execute("UPDATE notes SET dirty = 0", []).unwrap();
|
||||
assert_eq!(pending_fingerprint(&conn).unwrap(), None);
|
||||
}
|
||||
|
||||
fn queued_upload(conn: &Connection, id: &str, note_id: &str, error: Option<&str>) {
|
||||
conn.execute(
|
||||
"INSERT INTO attachments (id, note_id, url, filename, mime, sha256, uploaded, upload_error)
|
||||
VALUES (?1, ?2, '/x', 'f.txt', 'text/plain', 'h', 0, ?3)",
|
||||
params![id, note_id, error],
|
||||
)
|
||||
.expect("queued upload");
|
||||
}
|
||||
|
||||
fn upload_ids(conn: &Connection) -> Vec<String> {
|
||||
pending_uploads(conn)
|
||||
.expect("uploads")
|
||||
.into_iter()
|
||||
.map(|u| u.id)
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn removals_of_files_and_previews_wait_for_a_server_that_takes_them() {
|
||||
let conn = db();
|
||||
store::record_pending_delete(&conn, "attachment", "a1").expect("tombstone");
|
||||
store::record_pending_delete(&conn, "preview", "p1").expect("tombstone");
|
||||
store::record_pending_delete(&conn, "note", "n9").expect("tombstone");
|
||||
|
||||
// An older server would answer "unknown entity" to each of them.
|
||||
let older: Vec<_> = collect(&conn, 100, false)
|
||||
.expect("collect")
|
||||
.iter()
|
||||
.map(|c| c.entity)
|
||||
.collect();
|
||||
assert_eq!(older, vec!["note"]);
|
||||
|
||||
let mut newer: Vec<_> = collect(&conn, 100, true)
|
||||
.expect("collect")
|
||||
.iter()
|
||||
.map(|c| (c.entity, c.op))
|
||||
.collect();
|
||||
newer.sort();
|
||||
assert_eq!(
|
||||
newer,
|
||||
vec![
|
||||
("attachment", "delete"),
|
||||
("note", "delete"),
|
||||
("preview", "delete")
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_removed_attachment_stays_removed_through_a_whole_cycle() {
|
||||
use crate::sync::{pull, wire};
|
||||
let conn = db();
|
||||
let server_note = |revision: i64, attachments: Vec<wire::Attachment>| wire::Note {
|
||||
id: "n1".into(),
|
||||
body: "a note".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,
|
||||
previews: vec![],
|
||||
};
|
||||
let page = |note: wire::Note, cursor: i64| wire::ChangesPage {
|
||||
notes: vec![note],
|
||||
labels: vec![],
|
||||
cursor,
|
||||
has_more: false,
|
||||
};
|
||||
let a1 = wire::Attachment {
|
||||
id: "a1".into(),
|
||||
url: "/x".into(),
|
||||
filename: None,
|
||||
mime: "image/png".into(),
|
||||
size: None,
|
||||
sha256: None,
|
||||
};
|
||||
pull::apply_page(&conn, &page(server_note(1, vec![a1]), 1)).expect("first pull");
|
||||
|
||||
store::delete_attachment(&conn, "n1", "a1").expect("remove");
|
||||
|
||||
// Push: the removal and the touched note go up, and the server applies both.
|
||||
let batch = collect(&conn, 100, true).expect("collect");
|
||||
let results: Vec<PushResult> = batch
|
||||
.iter()
|
||||
.map(|c| ok("applied", (c.entity == "note").then_some(5)))
|
||||
.collect();
|
||||
apply_results(&conn, &batch, &results).expect("results");
|
||||
assert!(
|
||||
!has_pending(&conn).expect("pending"),
|
||||
"the tombstone is settled"
|
||||
);
|
||||
|
||||
// Pull: the server's note now comes without it.
|
||||
pull::apply_page(&conn, &page(server_note(6, vec![]), 6)).expect("second pull");
|
||||
let left: i64 = conn
|
||||
.query_row("SELECT COUNT(*) FROM attachments", [], |r| r.get(0))
|
||||
.expect("count");
|
||||
assert_eq!(left, 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_file_uploads_only_once_its_note_is_on_the_server() {
|
||||
let conn = db();
|
||||
seed_note(&conn, "landed", 0);
|
||||
seed_note(&conn, "unsent", 1);
|
||||
queued_upload(&conn, "a", "landed", None);
|
||||
queued_upload(&conn, "b", "unsent", None);
|
||||
queued_upload(&conn, "c", "landed", Some("file is too large (max 25 MB)"));
|
||||
assert_eq!(upload_ids(&conn), vec!["a"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_refused_upload_leaves_the_queue_and_a_passing_failure_stays_in_it() {
|
||||
let conn = db();
|
||||
seed_note(&conn, "n", 0);
|
||||
queued_upload(&conn, "big", "n", None);
|
||||
queued_upload(&conn, "flaky", "n", None);
|
||||
assert!(
|
||||
pending_fingerprint(&conn).expect("fp").is_some(),
|
||||
"uploads are pending work"
|
||||
);
|
||||
|
||||
let refused = Err(UploadError::Refused("file is too large (max 25 MB)".into()));
|
||||
settle_upload(&conn, "big", &refused).expect("settle");
|
||||
settle_upload(&conn, "flaky", &Err(UploadError::Retry("offline".into()))).expect("settle");
|
||||
assert_eq!(
|
||||
upload_ids(&conn),
|
||||
vec!["flaky"],
|
||||
"the refusal is not resent every cycle"
|
||||
);
|
||||
assert!(has_pending(&conn).expect("pending"));
|
||||
|
||||
settle_upload(&conn, "flaky", &Ok(())).expect("settle");
|
||||
assert!(upload_ids(&conn).is_empty());
|
||||
assert!(!has_pending(&conn).expect("pending"));
|
||||
assert_eq!(pending_fingerprint(&conn).expect("fp"), None);
|
||||
let reason: Option<String> = conn
|
||||
.query_row(
|
||||
"SELECT upload_error FROM attachments WHERE id = 'big'",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.expect("row");
|
||||
assert_eq!(
|
||||
reason.as_deref(),
|
||||
Some("file is too large (max 25 MB)"),
|
||||
"shown on the file"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user