diff --git a/CHANGELOG.md b/CHANGELOG.md index 261a35e7c..8e077cac1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -132,6 +132,17 @@ released versions carry their date on the heading. #### Importing +- 2026-10-07: **A message an import edits is no longer hidden behind a copy + of its old text.** When an import set not to hide duplicates gave a stored + message the text of a later edit, or added an attachment to it, the message + stayed hidden behind another backup's copy of what it said before, so a + search for its new text found nothing. A copy hidden behind it stayed + hidden too, though their texts no longer matched. Such an import now checks + the duplicates of the messages it changed, and of the copies around them, + whatever its own setting. The messages it adds stay as they came. A stored + message that only imports with dedupe off brought stays as it came too, + even when a later import changes its text. + - 2026-10-05: **A WhatsApp import from Android that fills the disk holding the Scratch Directory stops with the free-space sentence.** The encrypted WhatsApp backup is decrypted into the Scratch Directory, and its decrypted diff --git a/crates/server/server/src/cli.rs b/crates/server/server/src/cli.rs index e991907fc..9509cfe4d 100644 --- a/crates/server/server/src/cli.rs +++ b/crates/server/server/src/cli.rs @@ -138,7 +138,8 @@ pub struct ImportArgs { #[arg(long, default_value = "replace")] pub mode: ImportMode, - /// Skip the cross-source soft-dedupe pass after import + /// Skip the cross-source soft-dedupe pass after import. The import then + /// writes no content keys either; a later full dedupe writes them #[arg(long)] pub skip_dedupe: bool, diff --git a/crates/server/server/src/db/staging.rs b/crates/server/server/src/db/staging.rs index c2229174a..880f356b6 100644 --- a/crates/server/server/src/db/staging.rs +++ b/crates/server/server/src/db/staging.rs @@ -21,6 +21,8 @@ use sqlx::query::Query; use sqlx::sqlite::SqliteArguments; use sqlx::{Row, Sqlite, SqliteConnection}; +use crate::dedupe::HAS_CONTENT_KEY_SQL; + use super::sql::{SQLITE_IN_CHUNK, max_rows_for_bind_limit, values_tuples}; // ── Staging: what one import writes before promotion ───────────────────── @@ -792,11 +794,17 @@ pub async fn add_copy_tapbacks( // temp tables mapping staging ids to production ids, which the map writers // below fill in the order `imports_api::promote` calls them. -/// Create, or empty, a temp table mapping staging ids to production ids. +/// Create, or empty, a temp table mapping staging ids to production ids, +/// with the column definitions `extra_columns` after those two. /// Two statements on purpose: one prepared statement holds one command. -async fn reset_id_map(conn: &mut SqliteConnection, table: &str) -> Result<()> { +async fn reset_id_map( + conn: &mut SqliteConnection, + table: &str, + extra_columns: &[&str], +) -> Result<()> { + let extra: String = extra_columns.iter().map(|c| format!(", {c}")).collect(); let create = format!( - "CREATE TEMP TABLE IF NOT EXISTS {table} (staging_id BIGINT PRIMARY KEY, prod_id BIGINT NOT NULL)" + "CREATE TEMP TABLE IF NOT EXISTS {table} (staging_id BIGINT PRIMARY KEY, prod_id BIGINT NOT NULL{extra})" ); sqlx::query(&create).execute(&mut *conn).await?; let clear = format!("DELETE FROM {table}"); @@ -893,7 +901,7 @@ pub async fn upsert_conversations(conn: &mut SqliteConnection, account_id: i64) /// /// Returns an error when a statement fails. pub async fn write_conversation_map(conn: &mut SqliteConnection, account_id: i64) -> Result<()> { - reset_id_map(conn, "_promote_conv_map").await?; + reset_id_map(conn, "_promote_conv_map", &[]).await?; sqlx::query( r" INSERT INTO _promote_conv_map (staging_id, prod_id) @@ -1201,7 +1209,7 @@ pub async fn write_message_map( account_id: i64, pairs: &HashMap, ) -> Result<()> { - reset_id_map(conn, "_promote_msg_map").await?; + reset_id_map(conn, "_promote_msg_map", &[]).await?; let pairs: Vec<(i64, i64)> = pairs.iter().map(|(&s, &p)| (s, p)).collect(); for chunk in pairs.chunks(SQLITE_IN_CHUNK) { let sql = format!( @@ -1307,7 +1315,7 @@ pub fn later_backup<'a>(staged: Option<&'a str>, held: Option<&str>) -> BackupOr /// /// Returns an error when the update fails. pub async fn promote_deletion_marks(conn: &mut SqliteConnection) -> Result { - reset_id_map(conn, "_promote_mark_map").await?; + reset_id_map(conn, "_promote_mark_map", &[]).await?; let sql = format!( r" INSERT INTO _promote_mark_map (staging_id, prod_id) @@ -1448,6 +1456,12 @@ fn later_edit_sql(n: &str, newest: &str, held_n: &str, held_newest: &str) -> Str /// promotion, those at or below `messages_before`, whose staged row gives /// it a later text. Returns how many messages it names. /// +/// Each row says too whether the text itself changes (`body_changed`), +/// rather than only the earlier versions, and whether the message had a +/// content key before (`keyed`), because [`promote_later_edits`] clears the +/// key of a message whose text changes and nothing says afterwards that it +/// had one. [`stored_messages_with_new_content`] reads both. +/// /// When both backups have a date ([`later_backup_sql`]), a staged row from /// a later backup gives its text and earlier versions whatever their times /// say, when either differs from what the message holds (#1804); one from @@ -1464,7 +1478,12 @@ fn later_edit_sql(n: &str, newest: &str, held_n: &str, held_newest: &str) -> Str /// /// Returns an error when a statement fails. pub async fn write_edit_map(conn: &mut SqliteConnection, messages_before: i64) -> Result { - reset_id_map(conn, "_promote_edit_map").await?; + reset_id_map( + conn, + "_promote_edit_map", + &["body_changed BOOLEAN NOT NULL", "keyed BOOLEAN NOT NULL"], + ) + .await?; // A version list as one value, in the order its rows were written, so // two lists compare whole. let versions = |table: &str, id: &str| { @@ -1475,13 +1494,15 @@ pub async fn write_edit_map(conn: &mut SqliteConnection, messages_before: i64) - }; let sql = format!( r" - INSERT INTO _promote_edit_map (staging_id, prod_id) - SELECT staging_id, prod_id + INSERT INTO _promote_edit_map (staging_id, prod_id, body_changed, keyed) + SELECT staging_id, prod_id, body_changed, keyed FROM ( SELECT mm.staging_id, mm.prod_id, {later_backup} AS later_backup, + sm.body IS NOT m.body AS body_changed, + {has_content_key} AS keyed, sm.body IS NOT m.body OR {staged_versions} IS NOT {held_versions} AS differs, sv.message_id IS NOT NULL AS has_versions, @@ -1510,6 +1531,7 @@ pub async fn write_edit_map(conn: &mut SqliteConnection, messages_before: i64) - staged_versions = versions("staging_message_versions", "mm.staging_id"), held_versions = versions("message_versions", "mm.prod_id"), later_edit = later_edit_sql("n", "newest", "held_n", "held_newest"), + has_content_key = HAS_CONTENT_KEY_SQL, ); Ok(sqlx::query(&sql) .bind(messages_before) @@ -1530,8 +1552,9 @@ pub struct PromotedEdits { /// Give each message `_promote_edit_map` names the text of its staged row, /// and delete the earlier versions it held, which /// [`promote_earlier_versions`] then replaces with the staged ones. The -/// content key is cleared, because it hashes the text: the content-key fill -/// computes it again. +/// content key of a message whose text changes is cleared, because it +/// hashes the text: the content-key fill computes it again. A message given +/// only other earlier versions keeps its key, which does not hash them. /// /// The versions' search entries are not touched here: the caller removes /// them first (`schema::unindex_versions_of_edited_messages`), while the @@ -1544,8 +1567,8 @@ pub async fn promote_later_edits(conn: &mut SqliteConnection) -> Result Result, /// New rows inserted. pub inserted: u64, } @@ -1705,15 +1730,67 @@ fn staged_by_production_message(table: &str, columns: &[&str]) -> String { /// Returns an error when a statement fails. pub async fn promote_attachments(conn: &mut SqliteConnection) -> Result { let staged = staged_by_production_message("staging_attachments", ATTACHMENT_COLUMNS); - let filled = sqlx::query(&fill_attachments_sql("attachments", &staged)) - .execute(&mut *conn) - .await? - .rows_affected(); + let mut filled_messages: Vec = sqlx::query_scalar(&format!( + "{} RETURNING message_id", + fill_attachments_sql("attachments", &staged) + )) + .fetch_all(&mut *conn) + .await?; + let filled = filled_messages.len() as u64; + filled_messages.sort_unstable(); + filled_messages.dedup(); let inserted = sqlx::query(&insert_new_attachments_sql("attachments", &staged)) .execute(&mut *conn) .await? .rows_affected(); - Ok(PromotedAttachments { filled, inserted }) + Ok(PromotedAttachments { + filled, + filled_messages, + inserted, + }) +} + +/// The messages production held before this promotion, those at or below +/// `messages_before`, whose content key this promotion changed and that had +/// one before, sorted: those whose text a later edit changed +/// (`_promote_edit_map`), those that gained an attachment, one above +/// `attachments_before`, and those of `filled`, the stored attachments +/// given their file, which [`PromotedAttachments::filled_messages`] names +/// because nothing in the row says so afterwards. +/// +/// A message without a content key is left out: an import with dedupe off +/// brought it, and no dedupe has compared it. An edit that changes only +/// the earlier versions is left out too, because the key does not hash +/// them. +/// +/// # Errors +/// +/// Returns an error when the query fails. +pub async fn stored_messages_with_new_content( + conn: &mut SqliteConnection, + messages_before: i64, + attachments_before: i64, + filled: &[i64], +) -> Result> { + Ok(sqlx::query_scalar(&format!( + r" + SELECT prod_id FROM _promote_edit_map WHERE body_changed AND keyed + UNION + SELECT m.id + FROM messages m + WHERE {HAS_CONTENT_KEY_SQL} + AND ( + m.id IN (SELECT message_id FROM attachments WHERE id > $2 AND message_id <= $1) + OR m.id IN (SELECT value FROM json_each($3)) + ) + ORDER BY 1 + ", + )) + .bind(messages_before) + .bind(attachments_before) + .bind(serde_json::to_string(filled)?) + .fetch_all(&mut *conn) + .await?) } /// Insert the staged tapbacks under their production messages, skipping any diff --git a/crates/server/server/src/dedupe.rs b/crates/server/server/src/dedupe.rs index 87d499234..c0adefafe 100644 --- a/crates/server/server/src/dedupe.rs +++ b/crates/server/server/src/dedupe.rs @@ -1,6 +1,6 @@ //! Cross-source content fingerprint and soft-hide dedupe. -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeSet, HashMap, HashSet}; use std::fmt::Write as _; use std::io::{self, Write}; use std::time::Instant; @@ -209,6 +209,7 @@ pub async fn dedupe_cross_source( .map(|(i, s)| (s.as_str(), i)) .collect(); let started = Instant::now(); + let exact_hidden: HashSet; { println!(" Refreshing the content keys that match the same message across sources…"); @@ -236,9 +237,10 @@ pub async fn dedupe_cross_source( { println!(" Hiding exact duplicates, the messages that share a content key…"); let _ = io::stdout().flush(); - let (groups, flagged) = flag_exact_content_key_dupes(&mut tx, account_id, &prio).await?; + let (groups, flags) = flag_exact_content_key_dupes(&mut tx, account_id, &prio).await?; stats.exact_groups = groups; - stats.exact_flagged = flagged; + stats.exact_flagged = flags.len() as u64; + exact_hidden = flags.into_iter().map(|(loser, _)| loser).collect(); println!( " Found {} and hid {}, {:.1} s in all", words( @@ -255,7 +257,8 @@ pub async fn dedupe_cross_source( println!(" Flagging near duplicates, sent within {near_window_secs} s of each other…"); let _ = io::stdout().flush(); stats.near_flagged = - flag_near_time_dupes(&mut tx, account_id, &prio, near_window_secs).await?; + flag_near_time_dupes(&mut tx, account_id, &prio, near_window_secs, &exact_hidden) + .await?; println!( " Flagged {}, {:.1} s in all", words( @@ -271,12 +274,148 @@ pub async fn dedupe_cross_source( Ok(stats) } +/// The near-time window, in seconds, of a dedupe nobody chose one for: an +/// import's, and the pass an import runs for the messages it changed. +pub const NEAR_WINDOW_SECS: i64 = 2; + +/// What [`dedupe_changed_messages`] did. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub struct ChangedDedupe { + /// Messages that were hidden and are shown now. + pub shown: u64, + /// Messages that were shown and are hidden now, or hidden behind + /// another message than before. + pub hidden: u64, +} + +/// Put right the duplicate flags of the messages `changed`, whose content +/// an import changed in the transaction `conn` is in, and of every message +/// whose flag depends on theirs (#1805). +/// +/// The import's dedupe setting governs the rows it brings, so this runs +/// whatever the setting is. `changed` names only messages that had a +/// content key before the import changed them: a message an import with +/// dedupe off added has none, no dedupe has compared it, and it stays as it +/// came. This computes the changed messages' content keys again, then runs +/// both passes of [`dedupe_cross_source`] over the messages a dedupe has +/// seen, those with a content key, so a message without one is neither +/// hidden nor a winner here. Only the flags of the messages tied to a +/// changed one are written: those it was hidden behind or hid, before or +/// now, and the messages tied to those in turn. Every other flag stays as +/// it is. +/// +/// # Errors +/// +/// Returns an error when a statement fails or the hashing task panics. +pub async fn dedupe_changed_messages( + conn: &mut SqliteConnection, + account_id: i64, + changed: &[i64], + near_window_secs: i64, +) -> Result { + if changed.is_empty() { + return Ok(ChangedDedupe::default()); + } + let changed_json = serde_json::to_string(changed)?; + sqlx::query( + "UPDATE messages SET content_key = NULL WHERE id IN (SELECT value FROM json_each($1))", + ) + .bind(&changed_json) + .execute(&mut *conn) + .await?; + recompute_content_keys(conn, KeyScope::Changed(&changed_json), account_id).await?; + + let priority = source_priority_from_db(conn, account_id).await?; + let prio: HashMap<&str, usize> = priority + .iter() + .enumerate() + .map(|(i, s)| (s.as_str(), i)) + .collect(); + let (_, mut flags) = exact_flags(conn, account_id, &prio).await?; + let exact_hidden: HashSet = flags.iter().map(|&(loser, _)| loser).collect(); + let by_conversation = load_near_rows(conn, account_id, NearRows::Keyed, &exact_hidden).await?; + flags.extend(cluster_near_dupes(by_conversation, &prio, near_window_secs)); + let now: HashMap = flags.into_iter().collect(); + + let stored: HashMap = sqlx::query_as( + r" + SELECT m.id, m.duplicate_of + FROM messages m + JOIN conversations c ON c.id = m.conversation_id + WHERE c.account_id = $1 AND m.duplicate_of IS NOT NULL + ", + ) + .bind(account_id) + .fetch_all(&mut *conn) + .await? + .into_iter() + .collect(); + + let tied = tied_messages(changed, &stored, &now); + let mut result = ChangedDedupe::default(); + let mut flags: Vec<(i64, i64)> = Vec::new(); + for &id in &tied { + match (stored.get(&id), now.get(&id)) { + (Some(_), None) => result.shown += 1, + (before, Some(&winner)) => { + if before != Some(&winner) { + result.hidden += 1; + } + flags.push((id, winner)); + } + (None, None) => {} + } + } + let tied: Vec = tied.into_iter().collect(); + sqlx::query( + "UPDATE messages SET duplicate_of = NULL WHERE id IN (SELECT value FROM json_each($1))", + ) + .bind(serde_json::to_string(&tied)?) + .execute(&mut *conn) + .await?; + if !flags.is_empty() { + apply_duplicate_flags(conn, "_changed_flags", &flags).await?; + } + Ok(result) +} + +/// The messages tied to one in `changed`: `changed` itself, and every +/// message one hides or is hidden behind, in the `(loser → winner)` maps +/// `stored` (the flags as they are) or `now` (as the passes would set +/// them), followed to the end. +/// +/// The set is closed both ways, so writing the flags `now` gives its +/// messages leaves no stored flag outside it pointing at a message inside +/// it, and no new flag inside it pointing outside: every hidden message +/// still points at a shown one. +fn tied_messages( + changed: &[i64], + stored: &HashMap, + now: &HashMap, +) -> BTreeSet { + let mut neighbours: HashMap> = HashMap::new(); + for (&loser, &winner) in stored.iter().chain(now.iter()) { + neighbours.entry(loser).or_default().push(winner); + neighbours.entry(winner).or_default().push(loser); + } + let mut tied: BTreeSet = changed.iter().copied().collect(); + let mut queue: Vec = changed.to_vec(); + while let Some(id) = queue.pop() { + for &next in neighbours.get(&id).into_iter().flatten() { + if tied.insert(next) { + queue.push(next); + } + } + } + tied +} + /// Compute `content_key` for production rows that still lack one (after attachments exist). pub async fn fill_missing_content_keys( conn: &mut SqliteConnection, account_id: i64, ) -> Result { - recompute_content_keys(conn, true, account_id).await + recompute_content_keys(conn, KeyScope::Missing, account_id).await } /// Recompute every content key of the account and write the ones that differ @@ -287,7 +426,7 @@ pub async fn fill_missing_content_keys( /// group, and both are part of the key. A key left as it was stops matching /// the same message from another source. async fn refresh_content_keys(conn: &mut SqliteConnection, account_id: i64) -> Result { - recompute_content_keys(conn, false, account_id).await + recompute_content_keys(conn, KeyScope::All, account_id).await } /// Bulk-insert fingerprints into the `_content_keys` temp table in chunks that fit the bind limit. @@ -322,8 +461,19 @@ async fn insert_content_key_rows( Ok(()) } -/// Compute and store content keys for the account's messages: every message -/// (writing only the keys that changed), or only those without a key. +/// Which of the account's messages [`recompute_content_keys`] computes a key +/// for. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum KeyScope<'a> { + /// Every message, writing only the keys that changed. + All, + /// The messages without a key. + Missing, + /// The messages whose ids the JSON array names, writing every key. + Changed(&'a str), +} + +/// Compute and store content keys for the account's messages in `scope`. /// Returns how many were written. /// /// # Errors @@ -331,10 +481,10 @@ async fn insert_content_key_rows( /// Returns an error when a query fails or the hashing task panics. async fn recompute_content_keys( conn: &mut SqliteConnection, - missing_only: bool, + scope: KeyScope<'_>, account_id: i64, ) -> Result { - let Some(inputs) = ContentKeyInputs::load(conn, account_id, missing_only).await? else { + let Some(inputs) = ContentKeyInputs::load(conn, account_id, scope).await? else { return Ok(0); }; println!( @@ -345,7 +495,7 @@ async fn recompute_content_keys( let keys = tokio::task::spawn_blocking(move || inputs.hash()) .await .context("content-key hash task panicked")?; - let keys = if missing_only { + let keys = if scope != KeyScope::All { keys } else { let stored: HashMap = sqlx::query_as( @@ -390,6 +540,10 @@ fn sender_for_key_sql() -> String { ) } +/// Whether the message `m` has a content key, as a SQL condition: whether a +/// dedupe, or an import that filled the keys, has seen it. +pub(crate) const HAS_CONTENT_KEY_SQL: &str = "m.content_key IS NOT NULL AND m.content_key != ''"; + /// Everything the content-key hash reads, loaded in three queries so the /// hashing runs off the database thread with no lookups of its own. struct ContentKeyInputs { @@ -406,12 +560,16 @@ impl ContentKeyInputs { async fn load( conn: &mut SqliteConnection, account_id: i64, - missing_only: bool, + scope: KeyScope<'_>, ) -> Result> { - let filter = if missing_only { - "WHERE (m.content_key IS NULL OR m.content_key = '') AND c.account_id = $1" - } else { - "WHERE c.account_id = $1" + let filter = match scope { + KeyScope::All => "WHERE c.account_id = $1", + KeyScope::Missing => { + "WHERE (m.content_key IS NULL OR m.content_key = '') AND c.account_id = $1" + } + KeyScope::Changed(_) => { + "WHERE m.id IN (SELECT value FROM json_each($2)) AND c.account_id = $1" + } }; let sql = format!( r" @@ -427,10 +585,11 @@ impl ContentKeyInputs { ", sender = sender_for_key_sql(), ); - let rows: Vec = sqlx::query_as(&sql) - .bind(account_id) - .fetch_all(&mut *conn) - .await?; + let mut query = sqlx::query_as(&sql).bind(account_id); + if let KeyScope::Changed(ids) = scope { + query = query.bind(ids); + } + let rows: Vec = query.fetch_all(&mut *conn).await?; if rows.is_empty() { return Ok(None); } @@ -549,16 +708,32 @@ struct KeyedCand { /// Hide the messages that share a fingerprint with a preferred-source twin, /// keeping as many as one source holds (see [`exact_group_flags`]), and /// each whole-second message whose own source holds it with milliseconds -/// too (see [`content_key_group_flags`]). Returns (groups, hidden). +/// too (see [`content_key_group_flags`]). Returns the groups in which a +/// message was hidden, and the `(loser, winner)` pairs. async fn flag_exact_content_key_dupes( tx: &mut WriteTx<'_>, account_id: i64, prio: &HashMap<&str, usize>, -) -> Result<(u64, u64)> { +) -> Result<(u64, Vec<(i64, i64)>)> { let conn: &mut SqliteConnection = tx; + let (groups, flags) = exact_flags(conn, account_id, prio).await?; + if !flags.is_empty() { + apply_duplicate_flags(conn, "_pass_a_flags", &flags).await?; + } + Ok((groups, flags)) +} + +/// The exact pass over every message of the account with a content key, +/// as though none were hidden: the groups in which a message would be +/// hidden, and the `(loser, winner)` pairs. Nothing is written. +async fn exact_flags( + conn: &mut SqliteConnection, + account_id: i64, + prio: &HashMap<&str, usize>, +) -> Result<(u64, Vec<(i64, i64)>)> { // One scan of messages + one aggregated attachment pass, then group in Rust. // Avoids N round-trips (one SELECT + several UPDATEs per duplicate key). - let rows: Vec<(i64, String, String, i64, String)> = sqlx::query_as( + let rows: Vec<(i64, String, String, i64, String)> = sqlx::query_as(&format!( r" SELECT m.id, m.source, m.content_key, COALESCE(ac.n, 0), m.time_precision FROM messages m @@ -573,9 +748,9 @@ async fn flag_exact_content_key_dupes( GROUP BY a.message_id ) ac ON ac.message_id = m.id WHERE c.account_id = $1 - AND m.content_key IS NOT NULL AND m.content_key != '' + AND {HAS_CONTENT_KEY_SQL} ", - ) + )) .bind(account_id) .fetch_all(&mut *conn) .await?; @@ -605,15 +780,7 @@ async fn flag_exact_content_key_dupes( } flags.extend(group_flags); } - let flagged = flags.len() as u64; - - if flags.is_empty() { - return Ok((groups, 0)); - } - - apply_duplicate_flags(conn, "_pass_a_flags", &flags).await?; - - Ok((groups, flagged)) + Ok((groups, flags)) } /// The `(loser, winner)` pairs of the messages that share one content key. @@ -765,15 +932,17 @@ impl NearRow { } /// Flag messages that match one from another source within `window_secs` on chat, -/// direction, sender, and body or attachments. Returns how many were flagged. +/// direction, sender, and body or attachments, leaving out those the exact +/// pass hid (`exact_hidden`). Returns how many were flagged. async fn flag_near_time_dupes( tx: &mut WriteTx<'_>, account_id: i64, prio: &HashMap<&str, usize>, window_secs: i64, + exact_hidden: &HashSet, ) -> Result { let conn: &mut SqliteConnection = tx; - let by_conversation = load_near_rows(conn, account_id).await?; + let by_conversation = load_near_rows(conn, account_id, NearRows::All, exact_hidden).await?; let flags = cluster_near_dupes(by_conversation, prio, window_secs); if flags.is_empty() { return Ok(0); @@ -782,12 +951,24 @@ async fn flag_near_time_dupes( Ok(flags.len() as u64) } -/// Every unflagged message of the account with its attachment fingerprint, -/// grouped by conversation. Two queries and the grouping happen here so the -/// clustering does no per-message lookups. +/// Which messages the near-time pass reads. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum NearRows { + /// Every message of the account. + All, + /// The messages with a content key: those a dedupe has already seen, + /// and the changed ones [`dedupe_changed_messages`] gave one. + Keyed, +} + +/// The messages of the account in `rows` that are not in `hidden`, with +/// their attachment fingerprints, grouped by conversation. Two queries and +/// the grouping happen here so the clustering does no per-message lookups. async fn load_near_rows( conn: &mut SqliteConnection, account_id: i64, + rows: NearRows, + hidden: &HashSet, ) -> Result>> { type NearDedupeRow = ( i64, @@ -806,10 +987,13 @@ async fn load_near_rows( FROM messages m JOIN conversations c ON c.id = m.conversation_id LEFT JOIN handles hs ON hs.id = m.sender_handle_id - WHERE c.account_id = $1 - AND m.duplicate_of IS NULL + WHERE c.account_id = $1 {keyed} ", sender = sender_for_key_sql(), + keyed = match rows { + NearRows::All => String::new(), + NearRows::Keyed => format!("AND {HAS_CONTENT_KEY_SQL}"), + }, ); let msg_rows: Vec = sqlx::query_as(&msg_sql) .bind(account_id) @@ -837,6 +1021,9 @@ async fn load_near_rows( let mut by_conversation: HashMap> = HashMap::new(); for (id, conversation_id, source, is_from_me, ts, body, sender_norm, content_key) in msg_rows { + if hidden.contains(&id) { + continue; + } let Some(secs) = parse_rfc3339_utc_secs(ts.trim()) else { continue; }; diff --git a/crates/server/server/src/dedupe/tests.rs b/crates/server/server/src/dedupe/tests.rs index 84ccd9965..beb474cf7 100644 --- a/crates/server/server/src/dedupe/tests.rs +++ b/crates/server/server/src/dedupe/tests.rs @@ -1723,13 +1723,14 @@ fn survivor(rows: &HashMap, id: i64, ctx: &str) -> i64 { async fn assert_dedupe_invariants(conn: &mut SqliteConnection, ctx: &str) { let rows = load_gen_rows(conn).await; - let expected: HashMap = ContentKeyInputs::load(conn, TEST_ACCOUNT_ID, false) - .await - .unwrap() - .expect("the database has messages") - .hash() - .into_iter() - .collect(); + let expected: HashMap = + ContentKeyInputs::load(conn, TEST_ACCOUNT_ID, KeyScope::All) + .await + .unwrap() + .expect("the database has messages") + .hash() + .into_iter() + .collect(); for (id, row) in &rows { assert_eq!( row.content_key.as_deref(), diff --git a/crates/server/server/src/import_cli.rs b/crates/server/server/src/import_cli.rs index e0f23801f..b2b5866a7 100644 --- a/crates/server/server/src/import_cli.rs +++ b/crates/server/server/src/import_cli.rs @@ -195,7 +195,10 @@ async fn import_under_session( mode: opts.mode, source: opts.source_override.as_deref().unwrap_or(""), account_id, - fill_content_keys: true, + // The full dedupe after the import computes every content key, so + // only an import it does not follow puts the changed messages' + // flags right on its own. + fill_content_keys: !opts.skip_dedupe, import_id: Some(import_run.id), source_from_jsonl: plan.from_jsonl, media: opts.media, diff --git a/crates/server/server/src/imports_api/mod.rs b/crates/server/server/src/imports_api/mod.rs index 5ad274a3d..89d459c1b 100644 --- a/crates/server/server/src/imports_api/mod.rs +++ b/crates/server/server/src/imports_api/mod.rs @@ -67,7 +67,11 @@ pub struct ImportOptions<'a> { pub source: &'a str, /// Account the import writes into. pub account_id: i64, - /// Fill missing `content_key` values during promote (needed before cross-source dedupe). + /// Fill missing `content_key` values during promote (needed before + /// cross-source dedupe). With the fill, promote leaves the duplicate + /// flags of the stored messages whose content it changed to the full + /// dedupe, [`crate::dedupe::dedupe_cross_source`], which the caller runs + /// afterwards; without it, promote puts those flags right itself. pub fill_content_keys: bool, /// Optional Import Run id (messages stamped on promote). pub import_id: Option, @@ -92,7 +96,8 @@ pub struct FixedImportArgs<'a> { pub source: &'a str, /// Account the import writes into. pub account_id: i64, - /// Fill missing `content_key` values during promote. + /// Fill missing `content_key` values during promote, as + /// [`ImportOptions::fill_content_keys`] says. pub fill_content_keys: bool, /// Optional Import Run id (messages stamped on promote). pub import_id: Option, @@ -1798,7 +1803,7 @@ async fn run_import_path( .await; let counts = import_result?; let dedupe_stats = if do_dedupe { - Some(dedupe::dedupe_cross_source(&mut conn, account, None, 2).await?) + Some(dedupe::dedupe_cross_source(&mut conn, account, None, dedupe::NEAR_WINDOW_SECS).await?) } else { None }; diff --git a/crates/server/server/src/imports_api/promote.rs b/crates/server/server/src/imports_api/promote.rs index 7b7c52d0d..49446f4c8 100644 --- a/crates/server/server/src/imports_api/promote.rs +++ b/crates/server/server/src/imports_api/promote.rs @@ -92,9 +92,39 @@ const PROMOTE_MESSAGE_BATCH: i64 = 50_000; /// Below this many staged messages the secondary indexes are cheaper to keep than to rebuild. const PROMOTE_INDEX_DROP_MIN_STAGING: i64 = 5_000; +/// What the promotion did to the content of the messages production held +/// before it, as [`PromotedContent::changed_messages`] reads it. +struct PromotedContent { + /// The highest message id before the promotion: every message at or + /// below it was stored already. + messages_before: i64, + /// The highest attachment id before the promotion: every attachment + /// above it is new. + attachments_before: i64, + /// The stored messages whose stored attachments took their files. + filled_messages: Vec, +} + +impl PromotedContent { + /// The stored messages whose content key the promotion changed and that + /// had one before, sorted ([`staging::stored_messages_with_new_content`]). + async fn changed_messages(&self, conn: &mut SqliteConnection) -> Result> { + staging::stored_messages_with_new_content( + conn, + self.messages_before, + self.attachments_before, + &self.filled_messages, + ) + .await + } +} + impl Promote<'_> { /// Every phase, in order: the source wipe in replace mode, then each /// table, the search index, and the content keys when asked for. + /// Without the content keys, the duplicate flags of the stored messages + /// whose content changed are put right; with them, the caller runs the + /// full dedupe after the commit, which puts every flag right. async fn run(&mut self, wipe_sources: &[String], fill_content_keys: bool) -> Result<()> { if self.mode == ImportMode::Replace { self.wipe_sources(wipe_sources).await?; @@ -102,13 +132,20 @@ impl Promote<'_> { self.promote_conversations().await?; self.promote_participants().await?; let messages_before = self.promote_messages().await?; - let attachments_before = self.promote_attachments().await?; + let (attachments_before, filled_messages) = self.promote_attachments().await?; self.promote_tapbacks().await?; let versions_before = self.promote_earlier_versions(messages_before).await?; self.index_fts(messages_before, attachments_before, versions_before) .await?; if fill_content_keys { self.fill_content_keys().await?; + } else { + self.dedupe_changed_messages(&PromotedContent { + messages_before, + attachments_before, + filled_messages, + }) + .await?; } Ok(()) } @@ -417,8 +454,9 @@ impl Promote<'_> { /// Insert the staged attachments under their production messages. /// Returns the highest attachment id that existed before the insert: /// every new row lands above it, which is how [`Self::index_fts`] finds - /// the existing messages that gained an attachment. - async fn promote_attachments(&mut self) -> Result { + /// the existing messages that gained an attachment. Returns too the + /// stored messages whose attachment rows took their missing files. + async fn promote_attachments(&mut self) -> Result<(i64, Vec)> { let phase = Self::begin("Writing the import's attachments…"); let attachments_before = staging::max_attachment_id(self.tx).await?; let promoted = staging::promote_attachments(self.tx).await?; @@ -439,7 +477,7 @@ impl Promote<'_> { ), ), ); - Ok(attachments_before) + Ok((attachments_before, promoted.filled_messages)) } /// Insert the staged tapbacks under their production messages. @@ -527,6 +565,48 @@ impl Promote<'_> { Ok(()) } + /// Put right the duplicate flags of the stored messages whose content + /// this promotion changed, and of the messages tied to them + /// ([`crate::dedupe::dedupe_changed_messages`]), whatever the import's + /// dedupe setting: the setting governs the rows the import brings, and + /// a flag the import itself made wrong is put right (#1805). A + /// promotion that changed no stored message's content runs nothing. + /// + /// Only an import that fills no content keys runs this: one that fills + /// them is followed by the full dedupe once the caller commits + /// ([`crate::dedupe::dedupe_cross_source`]), which would write every + /// flag this writes again. + async fn dedupe_changed_messages(&mut self, promoted: &PromotedContent) -> Result<()> { + let changed = promoted.changed_messages(self.tx).await?; + if changed.is_empty() { + return Ok(()); + } + let phase = Self::begin(format_args!( + "Checking the duplicates of the {} whose content changed…", + words( + as_count(changed.len()), + "1 stored message", + "{n} stored messages" + ) + )); + let result = crate::dedupe::dedupe_changed_messages( + self.tx, + self.account_id, + &changed, + crate::dedupe::NEAR_WINDOW_SECS, + ) + .await?; + self.done( + phase, + format!( + "{} shown again, and {} hidden", + words(result.shown, "1 message is", "{n} messages are"), + words(result.hidden, "1 is", "{n} are"), + ), + ); + Ok(()) + } + /// Log the counts and return them. The caller commits. fn finish(self) -> PromoteStats { let Promote { stats, started, .. } = self; diff --git a/crates/server/server/src/imports_api/tests.rs b/crates/server/server/src/imports_api/tests.rs index 55a998362..01620b141 100644 --- a/crates/server/server/src/imports_api/tests.rs +++ b/crates/server/server/src/imports_api/tests.rs @@ -5810,4 +5810,5 @@ async fn a_page_of_import_runs_is_read_without_a_statement_per_row() { } mod backup_dates; +mod changed_content_dedupe; mod time_precision; diff --git a/crates/server/server/src/imports_api/tests/changed_content_dedupe.rs b/crates/server/server/src/imports_api/tests/changed_content_dedupe.rs new file mode 100644 index 000000000..bae96808f --- /dev/null +++ b/crates/server/server/src/imports_api/tests/changed_content_dedupe.rs @@ -0,0 +1,402 @@ +//! An import that changes a stored message's content puts its duplicate +//! flag right, whatever the import's dedupe setting (#1805). + +use super::*; + +/// 2015-03-12T18:04:22Z, the second every message here falls in. +const SECOND: i64 = 1_426_183_462_000; + +/// The chat every message here is in, and its one other participant, who +/// sends every message. +const CHAT: &str = "+15555550123"; + +/// Create an Import Run for `source` with `dedupe` on or off, post `lines` +/// under one conversation header as its one batch, and complete it. +async fn import( + state: &crate::server::AppState, + token: &str, + source: &str, + dedupe: bool, + lines: &[MessageLine], +) { + let header = conversation_header(source, CHAT).participant(CHAT, None); + let mut body = format!("{header}\n"); + for line in lines { + body.push_str(&format!("{}\n", line.clone().sender(CHAT).at(SECOND))); + } + let (_, created): (String, serde_json::Value) = post_created_json( + state, + "/v1/imports", + token, + serde_json::json!({ "source": source, "mode": "append", "dedupe": dedupe }), + ) + .await; + let id = created["id"].as_i64().unwrap(); + let (status, text) = crate::test_support::post_raw( + state, + &format!("/v1/imports/{id}/batches"), + token, + "application/jsonl", + body, + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK, "{text}"); + let _: serde_json::Value = post_json( + state, + &format!("/v1/imports/{id}/complete"), + token, + serde_json::json!({ "status": "completed" }), + ) + .await; +} + +/// The message `guid` edited once, from "see you at six" to `text`. +fn edited(guid: &str, text: &str) -> MessageLine { + message_line(guid, text).edit(EarlierVersion { + part_index: 0, + text: "see you at six".into(), + edited_at_unix_ms: Some(SECOND + 60_000), + }) +} + +/// The guid of the message `guid` is hidden behind, or `None` when it is +/// shown. +async fn hidden_behind(state: &crate::server::AppState, guid: &str) -> Option { + let mut conn = state.db.acquire().await.unwrap(); + sqlx::query_scalar( + r" + SELECT w.guid + FROM messages m + LEFT JOIN messages w ON w.id = m.duplicate_of + WHERE m.guid = $1 + ", + ) + .bind(guid) + .fetch_one(&mut *conn) + .await + .unwrap() +} + +/// The texts a search for `word` finds. +async fn found(state: &crate::server::AppState, token: &str, word: &str) -> Vec { + let page: serde_json::Value = get_json(state, &format!("/v1/messages?q={word}"), token).await; + page["items"] + .as_array() + .unwrap() + .iter() + .map(|m| m["text"].as_str().unwrap().to_string()) + .collect() +} + +/// The issue's scenario: a message hidden behind another source's copy of +/// its text takes a later edit from an append with dedupe off. It no longer +/// matches that copy, so it is shown, and a search for its new text finds +/// it. A new message the same import brings stays shown beside the copy +/// it duplicates, because the import's dedupe is off. +#[tokio::test] +async fn an_edit_with_dedupe_off_shows_a_message_hidden_behind_a_copy_of_its_old_text() { + let (state, _fixture, token) = importer().await; + import( + &state, + &token, + "sms", + true, + &[ + message_line("n-six", "see you at six"), + message_line("n-lunch", "lunch?"), + ], + ) + .await; + import( + &state, + &token, + "imessage", + true, + &[message_line("m-six", "see you at six")], + ) + .await; + assert_eq!( + hidden_behind(&state, "m-six").await.as_deref(), + Some("n-six") + ); + + import( + &state, + &token, + "imessage", + false, + &[ + edited("m-six", "see you at seven"), + message_line("m-lunch", "lunch?"), + ], + ) + .await; + + assert_eq!(hidden_behind(&state, "m-six").await, None); + assert_eq!(found(&state, &token, "seven").await, ["see you at seven"]); + assert_eq!( + hidden_behind(&state, "m-lunch").await, + None, + "dedupe off leaves the import's new messages as they came" + ); +} + +/// A message other copies are hidden behind takes a later edit from an +/// append with dedupe off. The copy of its old text is shown again, and a +/// copy of its new text from a third source is hidden behind it. +#[tokio::test] +async fn an_edit_with_dedupe_off_evaluates_the_copies_around_the_message_again() { + let (state, _fixture, token) = importer().await; + import( + &state, + &token, + "imessage", + true, + &[message_line("m-six", "see you at six")], + ) + .await; + import( + &state, + &token, + "sms", + true, + &[message_line("r-six", "see you at six")], + ) + .await; + import( + &state, + &token, + "whatsapp", + true, + &[message_line("q-seven", "see you at seven")], + ) + .await; + assert_eq!( + hidden_behind(&state, "r-six").await.as_deref(), + Some("m-six") + ); + assert_eq!(hidden_behind(&state, "q-seven").await, None); + + import( + &state, + &token, + "imessage", + false, + &[edited("m-six", "see you at seven")], + ) + .await; + + assert_eq!(hidden_behind(&state, "m-six").await, None); + assert_eq!( + hidden_behind(&state, "r-six").await, + None, + "the old text no longer matches" + ); + assert_eq!( + hidden_behind(&state, "q-seven").await.as_deref(), + Some("m-six"), + "the new text matches" + ); + assert_eq!( + found(&state, &token, "seven").await, + ["see you at seven"], + "the two copies of the new text are found once" + ); +} + +/// A message another copy is hidden behind gains its attachment's file from +/// an append with dedupe off. The copy, with the same text and no +/// attachment, still matches it within the near-time window and stays +/// hidden, and the message's content key hashes the file it now has. +/// +/// The duplicate-flag assertions hold without the fix too, because no flag +/// changes here: only the content key computed again, the last assertion, +/// fails without it. +#[tokio::test] +async fn an_attachment_with_dedupe_off_keeps_a_copy_that_still_matches_hidden() { + let (state, _fixture, token) = importer().await; + let missing = IrAttachment { + size_bytes: Some(12), + missing_reason: Some("not_found".into()), + ..attachment( + "attachments/photo.bin", + "photo.bin", + "application/octet-stream", + ) + }; + import( + &state, + &token, + "imessage", + true, + &[message_line("m-photo", "see attached").attachment(missing)], + ) + .await; + import( + &state, + &token, + "sms", + true, + &[message_line("r-photo", "see attached")], + ) + .await; + assert_eq!( + hidden_behind(&state, "r-photo").await.as_deref(), + Some("m-photo") + ); + let key = |state: crate::server::AppState| async move { + let mut conn = state.db.acquire().await.unwrap(); + sqlx::query_scalar::<_, Option>( + "SELECT content_key FROM messages WHERE guid = 'm-photo'", + ) + .fetch_one(&mut *conn) + .await + .unwrap() + }; + let before = key(state.clone()).await; + + let bytes = b"photo-bytes".to_vec(); + let sha = assets_api::sha256_hex(&bytes); + let (status, text) = crate::test_support::put_raw( + &state, + &format!("/v1/assets/{sha}"), + &token, + "application/octet-stream", + bytes.clone(), + ) + .await; + assert!(status.is_success(), "{status}: {text}"); + let found_file = IrAttachment { + digest_sha256: Some(sha), + size_bytes: Some(bytes.len() as u64), + ..attachment( + "attachments/photo.bin", + "photo.bin", + "application/octet-stream", + ) + }; + import( + &state, + &token, + "imessage", + false, + &[message_line("m-photo", "see attached").attachment(found_file)], + ) + .await; + + let mut conn = state.db.acquire().await.unwrap(); + let files: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM attachments a JOIN messages m ON m.id = a.message_id \ + WHERE m.guid = 'm-photo' AND a.sha256 IS NOT NULL", + ) + .fetch_one(&mut *conn) + .await + .unwrap(); + drop(conn); + assert_eq!(files, 1, "the append gives the stored attachment its file"); + assert_eq!(hidden_behind(&state, "m-photo").await, None); + assert_eq!( + hidden_behind(&state, "r-photo").await.as_deref(), + Some("m-photo") + ); + assert_ne!( + key(state.clone()).await, + before, + "the content key hashes the attachment's file" + ); +} + +/// An append with dedupe off that changes no stored message's content runs +/// no dedupe: a flag the full pass would set stays unset. +/// +/// A guard, not a test of the fix: it passes without the fix, and fails +/// when the fix runs a dedupe over the whole account. +#[tokio::test] +async fn an_append_with_dedupe_off_that_changes_nothing_stored_runs_no_dedupe() { + let (state, _fixture, token) = importer().await; + import( + &state, + &token, + "sms", + true, + &[message_line("n-six", "see you at six")], + ) + .await; + import( + &state, + &token, + "imessage", + true, + &[message_line("m-six", "see you at six")], + ) + .await; + { + let mut conn = state.db.acquire().await.unwrap(); + let mut tx = crate::db::begin_write(&mut conn).await.unwrap(); + sqlx::query("UPDATE messages SET duplicate_of = NULL WHERE guid = 'm-six'") + .execute(&mut *tx) + .await + .unwrap(); + tx.commit().await.unwrap(); + } + + import( + &state, + &token, + "imessage", + false, + &[message_line("m-six", "see you at six")], + ) + .await; + + assert_eq!(hidden_behind(&state, "m-six").await, None); +} + +/// A message an import with dedupe off brought, with no content key, takes +/// a later edit from another append with dedupe off. Its new text matches +/// a copy from a source imported with dedupe on, but no dedupe was asked +/// for its own source, so it stays as it came: shown, without a content +/// key, and the copy stays shown beside it. +#[tokio::test] +async fn an_edit_with_dedupe_off_leaves_a_message_no_dedupe_has_seen_as_it_came() { + let (state, _fixture, token) = importer().await; + import( + &state, + &token, + "sms", + true, + &[ + message_line("n-six", "see you at six"), + message_line("n-seven", "see you at seven"), + ], + ) + .await; + import( + &state, + &token, + "imessage", + false, + &[message_line("m-six", "see you at six")], + ) + .await; + assert_eq!(hidden_behind(&state, "m-six").await, None); + + import( + &state, + &token, + "imessage", + false, + &[edited("m-six", "see you at seven")], + ) + .await; + + assert_eq!(hidden_behind(&state, "m-six").await, None); + assert_eq!(hidden_behind(&state, "n-seven").await, None); + let mut conn = state.db.acquire().await.unwrap(); + let key: Option = + sqlx::query_scalar("SELECT content_key FROM messages WHERE guid = 'm-six'") + .fetch_one(&mut *conn) + .await + .unwrap(); + assert_eq!(key, None, "no dedupe has seen the message"); +} diff --git a/docs/architecture/contacts-identities-and-messages.md b/docs/architecture/contacts-identities-and-messages.md index 51d48c934..6713654ad 100644 --- a/docs/architecture/contacts-identities-and-messages.md +++ b/docs/architecture/contacts-identities-and-messages.md @@ -439,9 +439,24 @@ date says which state is the newer; the message's own times record when it was written, not when a part was unsent, so they cannot tell ([#1741](https://github.com/messagecrate/message-crate/issues/1741), [#1804](https://github.com/messagecrate/message-crate/issues/1804), -[#1924](https://github.com/messagecrate/message-crate/issues/1924)). An append -with dedupe off still leaves a changed message's duplicate flag until the next -dedupe ([#1805](https://github.com/messagecrate/message-crate/issues/1805)). +[#1924](https://github.com/messagecrate/message-crate/issues/1924)). + +**An import that changes a stored message's content puts its duplicate flag +right, whatever the import's dedupe setting.** A later edit changes a stored +message's text, and an attachment added or given its file changes what its +content key hashes, so its duplicate flag can stop matching. The import runs +the dedupe for those messages inside the import's own write transaction, and +writes the flags of the messages tied to them: those a changed message hid or +was hidden behind, before or now, followed to the end +(`dedupe_changed_messages` in `dedupe.rs`). An import with dedupe on leaves +this to the full dedupe that follows it. The pass compares only messages with +a content key, the ones a dedupe has seen, and takes up a changed message only +when it had one, so a message an import with dedupe off brought stays as it +came, even when a later import changes it. Why: the dedupe +setting governs the rows an import brings, and a flag the import itself made +wrong would hide a message behind a copy whose text no longer matches, where +no search finds it +([#1805](https://github.com/messagecrate/message-crate/issues/1805)). **Within one source, a whole-second message is the duplicate of its millisecond twin.** Each message says whether its source recorded its time to diff --git a/docs/src/content/docs/docs/developer/reference/server-cli.md b/docs/src/content/docs/docs/developer/reference/server-cli.md index 2dd4634af..b220f7b7a 100644 --- a/docs/src/content/docs/docs/developer/reference/server-cli.md +++ b/docs/src/content/docs/docs/developer/reference/server-cli.md @@ -71,7 +71,7 @@ Import a message-ir JSONL directory, one Import Run per source (source from expo * `--mode ` — Import mode: replace (wipe sources found in input) or append Default value: `replace` -* `--skip-dedupe` — Skip the cross-source soft-dedupe pass after import +* `--skip-dedupe` — Skip the cross-source soft-dedupe pass after import. The import then writes no content keys either; a later full dedupe writes them * `--window-secs ` — Near-time window in seconds for dedupe Pass B (default 2) Default value: `2`