diff --git a/aw-sync/src/sync.rs b/aw-sync/src/sync.rs index 69d811f7..aef36647 100644 --- a/aw-sync/src/sync.rs +++ b/aw-sync/src/sync.rs @@ -547,6 +547,69 @@ fn get_or_create_sync_bucket( Err(e) => return Err(format!("Failed to get bucket '{sanitized_id}': {e:?}")), } } + + // Pre-#697 origin-based fallback (ActivityWatch/aw-server-rust#707): + // Both exact lookups missed. A desktop that imported the peer before + // ActivityWatch/aw-server-rust#697 landed holds the raw hostname from that + // day (e.g. `…-synced-from-POCO F8 Ultra`) with `$aw.sync.origin` set to + // that same raw value. After ActivityWatch/aw-android#273 migrates the + // phone's staging hostname to `poco_f8_ultra`, first-hand buckets carry no + // `$aw.sync.origin`, so both lookups above miss and the whole history gets + // re-imported. Scan the destination's -synced-from- buckets for one whose + // base ID matches and whose `$aw.sync.origin` sanitizes to the same target. + if !is_push { + let target_base = bucket_from + .id + .split("-synced-from-") + .next() + .unwrap_or(bucket_from.id.as_str()); + let target_sanitized = sanitize_hostname( + sync_origin + .as_deref() + .unwrap_or(bucket_from.hostname.as_str()), + ); + let all_dest = ds_to + .get_buckets() + .map_err(|e| format!("Failed to list dest buckets for origin scan: {e:?}"))?; + let mut candidates: Vec = all_dest + .into_values() + .filter(|b| { + let b_base = b.id.split("-synced-from-").next().unwrap_or(&b.id); + if b_base != target_base { + return false; + } + // $aw.sync.origin must be present and sanitize to the same value. + b.data + .get("$aw.sync.origin") + .and_then(|v| v.as_str()) + .map(|s| sanitize_hostname(s) == target_sanitized) + .unwrap_or(false) + }) + .collect(); + match candidates.len() { + 0 => {} // fall through to create a new bucket + 1 => { + let found = candidates.remove(0); + info!( + " ↩ Reusing pre-#697 bucket '{}' for '{}'", + found.id, bucket_from.id + ); + return Ok(found); + } + n => { + // Two distinct pre-#697 buckets share the same sanitized origin — + // ambiguous. Refuse rather than silently merging distinct histories + // (ActivityWatch/aw-server-rust#697 :368). + let ids: Vec<&str> = candidates.iter().map(|b| b.id.as_str()).collect(); + return Err(format!( + "Cannot resolve destination for '{}': {n} pre-#697 buckets share \ + sanitized origin '{}': {ids:?}; deduplicate manually", + bucket_from.id, target_sanitized + )); + } + } + } + let (final_id, final_hostname) = (sanitized_id, sanitized_hostname); let mut bucket_new = bucket_from.clone(); diff --git a/aw-sync/tests/sync.rs b/aw-sync/tests/sync.rs index d4e02669..14b46f3d 100644 --- a/aw-sync/tests/sync.rs +++ b/aw-sync/tests/sync.rs @@ -261,6 +261,246 @@ mod sync_tests { ); } + /// A desktop that imported the peer **before** ActivityWatch/aw-server-rust#697 + /// landed holds `…-synced-from-POCO F8 Ultra` (raw, with `$aw.sync.origin` set + /// to the raw value by #697's import stamp). After ActivityWatch/aw-android#273 + /// migrates the phone's hostname to `poco_f8_ultra`, first-hand buckets carry no + /// `$aw.sync.origin`, so the two direct lookups miss. The pre-#697 fallback scan + /// must find the legacy bucket and resume from it rather than creating a new one + /// that triggers a full re-import. + #[test] + fn test_pre697_origin_scan_resumes_legacy_bucket() { + let state = init_teststate(); + + // Post-migration phone bucket: sanitized hostname, no $aw.sync.origin. + let src_bucket: Bucket = serde_json::from_value(serde_json::json!({ + "id": "aw-watcher-android", + "type": "currentwindow", + "hostname": "poco_f8_ultra", + "client": "aw-android" + })) + .unwrap(); + state.ds_src.create_bucket(&src_bucket).unwrap(); + + // Pre-#697 destination bucket: raw ID + $aw.sync.origin stamped by #697. + let legacy_id = "aw-watcher-android-synced-from-POCO F8 Ultra"; + let legacy_bucket: Bucket = serde_json::from_value(serde_json::json!({ + "id": legacy_id, + "type": "currentwindow", + "hostname": "POCO F8 Ultra", + "client": "aw-android", + "data": {"$aw.sync.origin": "POCO F8 Ultra"} + })) + .unwrap(); + state.ds_dest.create_bucket(&legacy_bucket).unwrap(); + + // Seed the legacy destination bucket with one event at T0 — simulating + // previously-imported history. This establishes the resume cursor. + let t0 = Utc::now(); + let existing_event: Event = serde_json::from_value(serde_json::json!({ + "timestamp": t0.to_rfc3339(), + "duration": 1, + "data": {"app": "existing"} + })) + .unwrap(); + state + .ds_dest + .insert_events(legacy_id, &[existing_event]) + .unwrap(); + state.ds_dest.force_commit().unwrap(); + + // Two source events: one before T0 (already covered) and one after T0 (new). + // The sync must read the cursor from the reused bucket and import only the + // post-T0 event — not re-import everything from scratch. + let before_t0: Event = serde_json::from_value(serde_json::json!({ + "timestamp": (t0 - Duration::hours(1)).to_rfc3339(), + "duration": 1, + "data": {"app": "old"} + })) + .unwrap(); + let after_t0: Event = serde_json::from_value(serde_json::json!({ + "timestamp": (t0 + Duration::hours(1)).to_rfc3339(), + "duration": 1, + "data": {"app": "new"} + })) + .unwrap(); + state + .ds_src + .insert_events("aw-watcher-android", &[before_t0, after_t0]) + .unwrap(); + state.ds_src.force_commit().unwrap(); + + aw_sync::sync_datastores( + &state.ds_src, + &state.ds_dest, + false, // pull + None, + &SyncSpec::default(), + ) + .unwrap(); + + let dest_buckets = state.ds_dest.get_buckets().unwrap(); + + // The legacy bucket must be reused, not replaced. + assert!( + dest_buckets.contains_key(legacy_id), + "legacy bucket must be preserved" + ); + // No new sanitized fork must appear. + let forked_id = "aw-watcher-android-synced-from-poco_f8_ultra"; + assert!( + !dest_buckets.contains_key(forked_id), + "a sanitized fork must not be created; got: {:?}", + dest_buckets.keys().collect::>() + ); + // Exactly 2 events: the pre-existing one at T0 plus the new T0+1h event. + // A count of 3 would mean the cursor was NOT read (full re-import from scratch). + let event_count = state + .ds_dest + .get_event_count(legacy_id, None, None) + .unwrap(); + assert_eq!( + event_count, 2, + "legacy bucket must have exactly 2 events (existing + new); \ + 3 would mean the cursor was ignored and history was re-imported" + ); + } + + /// Two distinct pre-#697 buckets for the same base ID whose `$aw.sync.origin` + /// values sanitize to the same target must trigger an error rather than a + /// silent merge (ActivityWatch/aw-server-rust#697 :368). + #[test] + fn test_pre697_origin_scan_refuses_ambiguous_candidates() { + let state = init_teststate(); + + // Ambiguous source bucket: two pre-#697 destination buckets share the same + // sanitized origin — sync_one must skip this bucket with a warning, not abort. + let src_bucket: Bucket = serde_json::from_value(serde_json::json!({ + "id": "aw-watcher-android", + "type": "currentwindow", + "hostname": "poco_f8_ultra", + "client": "aw-android" + })) + .unwrap(); + state.ds_src.create_bucket(&src_bucket).unwrap(); + + // A second, healthy source bucket that must sync successfully even while the + // ambiguous bucket is being skipped — one unresolvable bucket must not abort + // the whole peer sync. + let healthy_id = "aw-watcher-window"; + let healthy_bucket: Bucket = serde_json::from_value(serde_json::json!({ + "id": healthy_id, + "type": "currentwindow", + "hostname": "poco_f8_ultra", + "client": "aw-qt" + })) + .unwrap(); + state.ds_src.create_bucket(&healthy_bucket).unwrap(); + + let ts = Utc::now(); + let ev: Event = serde_json::from_value(serde_json::json!({ + "timestamp": ts.to_rfc3339(), + "duration": 1, + "data": {"app": "test"} + })) + .unwrap(); + state.ds_src.insert_events(healthy_id, &[ev]).unwrap(); + + // The ambiguous bucket must ALSO carry an event. Otherwise a regression that + // silently reused one of the two legacy candidates would copy nothing, both + // legacy buckets would still read 0 events, and the refusal assertions below + // would pass vacuously. With a real event present, only an actual skip keeps + // them at 0. + let ev_ambiguous: Event = serde_json::from_value(serde_json::json!({ + "timestamp": ts.to_rfc3339(), + "duration": 1, + "data": {"app": "ambiguous"} + })) + .unwrap(); + state + .ds_src + .insert_events(&src_bucket.id, &[ev_ambiguous]) + .unwrap(); + state.ds_src.force_commit().unwrap(); + + // Premise guard: the refusal assertions below are only meaningful if the + // ambiguous source bucket actually has something to copy. + assert_eq!( + state + .ds_src + .get_event_count(&src_bucket.id, None, None) + .unwrap(), + 1, + "premise: the ambiguous source bucket must carry one event" + ); + + // Two legacy destination buckets whose origins both sanitize to "poco_f8_ultra". + for (legacy_id, raw_origin) in [ + ( + "aw-watcher-android-synced-from-POCO F8 Ultra", + "POCO F8 Ultra", + ), + ( + "aw-watcher-android-synced-from-Poco F8 Ultra", + "Poco F8 Ultra", + ), + ] { + let b: Bucket = serde_json::from_value(serde_json::json!({ + "id": legacy_id, + "type": "currentwindow", + "hostname": raw_origin, + "client": "aw-android", + "data": {"$aw.sync.origin": raw_origin} + })) + .unwrap(); + state.ds_dest.create_bucket(&b).unwrap(); + } + + // The ambiguous android bucket is skipped (warn+continue); the healthy window + // bucket syncs normally. sync_datastores must return Ok overall. + aw_sync::sync_datastores( + &state.ds_src, + &state.ds_dest, + false, // pull + None, + &SyncSpec::default(), + ) + .unwrap(); + + let dest_buckets = state.ds_dest.get_buckets().unwrap(); + + // The healthy bucket was synced — it has a destination and received events. + let healthy_dest_id = "aw-watcher-window-synced-from-poco_f8_ultra"; + assert!( + dest_buckets.contains_key(healthy_dest_id), + "healthy bucket must be synced even when another bucket is ambiguous; got: {:?}", + dest_buckets.keys().collect::>() + ); + let healthy_count = state + .ds_dest + .get_event_count(healthy_dest_id, None, None) + .unwrap(); + assert!(healthy_count > 0, "healthy bucket must have events synced"); + + // The two ambiguous legacy buckets were not written to — the conflict was + // skipped, not merged. The source bucket holds an event (premise guard + // above), so a silent pick-one-candidate regression would make one of these + // counts 1 and fail here. + for legacy_id in [ + "aw-watcher-android-synced-from-POCO F8 Ultra", + "aw-watcher-android-synced-from-Poco F8 Ultra", + ] { + let count = state + .ds_dest + .get_event_count(legacy_id, None, None) + .unwrap(); + assert_eq!( + count, 0, + "ambiguous legacy bucket '{legacy_id}' must not have received events" + ); + } + } + /// Case-only hostnames (`PIXEL8`) have no whitespace, so a whitespace-only /// guard would leave the destination as `…-synced-from-PIXEL8`. Android's /// later hostname migration produces `pixel8` and forks the history.