Skip to content

Commit 2e2fd0e

Browse files
committed
Deduplicate live mode remote refresh work
1 parent 3284255 commit 2e2fd0e

4 files changed

Lines changed: 293 additions & 54 deletions

File tree

apps/desktop/src-tauri/src/main.rs

Lines changed: 122 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -501,6 +501,9 @@ static LAUNCH_AT_LOGIN_STATE: OnceLock<Mutex<Option<bool>>> = OnceLock::new();
501501
static LIVE_MODE_REMOTE_PULL_CURSOR: OnceLock<Mutex<usize>> = OnceLock::new();
502502
static LIVE_MODE_REMOTE_PULL_SCAN_TIMES: OnceLock<Mutex<BTreeMap<MountId, Instant>>> =
503503
OnceLock::new();
504+
static LIVE_MODE_REMOTE_FAST_FORWARD_TIMES: OnceLock<
505+
Mutex<BTreeMap<(MountId, RemoteId), Instant>>,
506+
> = OnceLock::new();
504507
static LIVE_MODE_LOCAL_RECONCILE_TIMES: OnceLock<Mutex<BTreeMap<PathBuf, Instant>>> =
505508
OnceLock::new();
506509
static LIVE_MODE_TICK_IN_PROGRESS: AtomicBool = AtomicBool::new(false);
@@ -519,6 +522,7 @@ const TRAY_POPOVER_EDGE_MARGIN: f64 = 8.0;
519522
const TRAY_POPOVER_ANCHOR_OFFSET: f64 = 12.0;
520523
const LIVE_MODE_LOCAL_RECONCILE_INTERVAL: Duration = Duration::from_secs(5);
521524
const LIVE_MODE_ACTIVE_REMOTE_CHECK_INTERVAL: Duration = Duration::from_secs(5);
525+
const LIVE_MODE_REMOTE_FAST_FORWARD_LEASE: Duration = Duration::from_secs(30);
522526
const LIVE_MODE_ACTIVE_TARGET_WINDOW: Duration = Duration::from_secs(5 * 60);
523527
const LIVE_MODE_REMOTE_CHECK_ESTIMATED_REQUESTS_PER_PAGE: f64 = 1.0;
524528
const LIVE_MODE_REMOTE_CHECK_QUEUE_SHARE: f64 = 1.0 / 3.0;
@@ -2500,47 +2504,82 @@ fn live_mode_queue_remote_fast_forward_at_state_root(
25002504
state_root: &Path,
25012505
target: &LiveModeRemoteTarget,
25022506
) -> Result<(), String> {
2503-
live_mode_send_daemon_queue_request_at_state_root(
2507+
let key = (target.mount_id.clone(), target.remote_id.clone());
2508+
if !live_mode_claim_remote_fast_forward_key(&key, Instant::now()) {
2509+
return Ok(());
2510+
}
2511+
2512+
let response = send_request(
25042513
state_root,
2505-
DaemonRequest::RemoteFastForward {
2514+
&DaemonRequest::RemoteFastForward {
25062515
mount_id: target.mount_id.0.clone(),
25072516
remote_id: target.remote_id.0.clone(),
25082517
path: target.path.clone(),
25092518
},
2510-
)?;
2511-
desktop_log(
2512-
"info",
2513-
"live_mode_remote_fast_forward_queued",
2519+
)
2520+
.map_err(|error| {
2521+
live_mode_release_remote_fast_forward_key(&key);
25142522
format!(
2515-
"queued remote fast-forward for `{}` ({}/{})",
2516-
target.path.display(),
2517-
target.mount_id.as_str(),
2518-
target.remote_id.as_str()
2519-
),
2520-
);
2523+
"Live Mode could not reach the Locality daemon to queue remote sync work: {}",
2524+
error.message()
2525+
)
2526+
})?;
2527+
if !response.ok {
2528+
live_mode_release_remote_fast_forward_key(&key);
2529+
let message = response
2530+
.error
2531+
.map(|error| error.message)
2532+
.unwrap_or_else(|| "daemon rejected the Live Mode request".to_string());
2533+
return Err(format!(
2534+
"Live Mode could not queue remote sync work: {message}"
2535+
));
2536+
}
2537+
2538+
if response
2539+
.payload
2540+
.as_ref()
2541+
.and_then(|payload| payload.get("queued"))
2542+
.and_then(|queued| queued.as_bool())
2543+
.unwrap_or(true)
2544+
{
2545+
desktop_log(
2546+
"info",
2547+
"live_mode_remote_fast_forward_queued",
2548+
format!(
2549+
"queued remote fast-forward for `{}` ({}/{})",
2550+
target.path.display(),
2551+
target.mount_id.as_str(),
2552+
target.remote_id.as_str()
2553+
),
2554+
);
2555+
}
25212556
Ok(())
25222557
}
25232558

2524-
fn live_mode_send_daemon_queue_request_at_state_root(
2525-
state_root: &Path,
2526-
request: DaemonRequest,
2527-
) -> Result<(), String> {
2528-
match send_request(state_root, &request) {
2529-
Ok(response) if response.ok => Ok(()),
2530-
Ok(response) => {
2531-
let message = response
2532-
.error
2533-
.map(|error| error.message)
2534-
.unwrap_or_else(|| "daemon rejected the Live Mode request".to_string());
2535-
Err(format!(
2536-
"Live Mode could not queue remote sync work: {message}"
2537-
))
2538-
}
2539-
Err(error) => Err(format!(
2540-
"Live Mode could not reach the Locality daemon to queue remote sync work: {}",
2541-
error.message()
2542-
)),
2559+
fn live_mode_claim_remote_fast_forward_key(key: &(MountId, RemoteId), now: Instant) -> bool {
2560+
let Ok(mut times) = live_mode_remote_fast_forward_times().lock() else {
2561+
return true;
2562+
};
2563+
times.retain(|_, last| {
2564+
now.checked_duration_since(*last)
2565+
.is_some_and(|elapsed| elapsed < LIVE_MODE_REMOTE_FAST_FORWARD_LEASE)
2566+
});
2567+
if times.contains_key(key) {
2568+
return false;
25432569
}
2570+
times.insert(key.clone(), now);
2571+
true
2572+
}
2573+
2574+
fn live_mode_release_remote_fast_forward_key(key: &(MountId, RemoteId)) {
2575+
let Ok(mut times) = live_mode_remote_fast_forward_times().lock() else {
2576+
return;
2577+
};
2578+
times.remove(key);
2579+
}
2580+
2581+
fn live_mode_remote_fast_forward_times() -> &'static Mutex<BTreeMap<(MountId, RemoteId), Instant>> {
2582+
LIVE_MODE_REMOTE_FAST_FORWARD_TIMES.get_or_init(|| Mutex::new(BTreeMap::new()))
25442583
}
25452584

25462585
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
@@ -10246,17 +10285,18 @@ mod tests {
1024610285

1024710286
use super::{
1024810287
ActionReport, DESKTOP_ACTIVITY_LIMIT, DESKTOP_INSTALL_MARKER_VERSION,
10249-
LIVE_MODE_RUNNER_ACTIVE_INTERVAL, LIVE_MODE_RUNNER_PERIODIC_RECHECK, MonitorScreenBounds,
10250-
PendingChange, ScreenBounds, TerminalCliLinkState, TrayVisualState,
10251-
acknowledge_install_state_at, activity_timestamp, clear_mount_cached_projection,
10252-
clear_state_root_contents, clear_visible_projection_paths, conflict_preview,
10253-
connection_metadata_changed, current_daemon_build_id, current_desktop_build_id,
10254-
diff_report_message, exact_located_entity_record, exact_notion_entry_matches,
10255-
failed_push_summary, has_unresolved_conflict_markers, hydration_after_editor_write,
10256-
inspect_install_state, install_terminal_cli_link_at,
10288+
LIVE_MODE_REMOTE_FAST_FORWARD_LEASE, LIVE_MODE_RUNNER_ACTIVE_INTERVAL,
10289+
LIVE_MODE_RUNNER_PERIODIC_RECHECK, MonitorScreenBounds, PendingChange, ScreenBounds,
10290+
TerminalCliLinkState, TrayVisualState, acknowledge_install_state_at, activity_timestamp,
10291+
clear_mount_cached_projection, clear_state_root_contents, clear_visible_projection_paths,
10292+
conflict_preview, connection_metadata_changed, current_daemon_build_id,
10293+
current_desktop_build_id, diff_report_message, exact_located_entity_record,
10294+
exact_notion_entry_matches, failed_push_summary, has_unresolved_conflict_markers,
10295+
hydration_after_editor_write, inspect_install_state, install_terminal_cli_link_at,
1025710296
install_terminal_cli_link_in_path_dirs, is_notion_access_lost_message,
10258-
is_unsupported_schema_version_message, live_mode_local_reconcile_targets_for_mount_at,
10259-
live_mode_merge_remote_drift_markdown, live_mode_remote_check_page_budget_for_rate,
10297+
is_unsupported_schema_version_message, live_mode_claim_remote_fast_forward_key,
10298+
live_mode_local_reconcile_targets_for_mount_at, live_mode_merge_remote_drift_markdown,
10299+
live_mode_release_remote_fast_forward_key, live_mode_remote_check_page_budget_for_rate,
1026010300
live_mode_remote_pull_candidates, live_mode_remote_pull_scan_is_due_for_key,
1026110301
live_mode_runner_should_tick, live_mode_should_reconcile_local_target_for_key,
1026210302
live_mode_target, live_mode_tick_from_snapshot, live_mode_wake_generation,
@@ -14285,6 +14325,47 @@ mod tests {
1428514325
));
1428614326
}
1428714327

14328+
#[test]
14329+
fn live_mode_remote_fast_forward_lease_suppresses_recent_duplicate() {
14330+
let key = (
14331+
MountId::new("lease-mount-duplicate"),
14332+
RemoteId::new("lease-remote-duplicate"),
14333+
);
14334+
live_mode_release_remote_fast_forward_key(&key);
14335+
let now = Instant::now();
14336+
14337+
assert!(live_mode_claim_remote_fast_forward_key(&key, now));
14338+
assert!(!live_mode_claim_remote_fast_forward_key(
14339+
&key,
14340+
now + Duration::from_millis(1)
14341+
));
14342+
assert!(live_mode_claim_remote_fast_forward_key(
14343+
&key,
14344+
now + LIVE_MODE_REMOTE_FAST_FORWARD_LEASE + Duration::from_millis(1)
14345+
));
14346+
14347+
live_mode_release_remote_fast_forward_key(&key);
14348+
}
14349+
14350+
#[test]
14351+
fn live_mode_remote_fast_forward_lease_can_release_for_retry() {
14352+
let key = (
14353+
MountId::new("lease-mount-release"),
14354+
RemoteId::new("lease-remote-release"),
14355+
);
14356+
live_mode_release_remote_fast_forward_key(&key);
14357+
let now = Instant::now();
14358+
14359+
assert!(live_mode_claim_remote_fast_forward_key(&key, now));
14360+
live_mode_release_remote_fast_forward_key(&key);
14361+
assert!(live_mode_claim_remote_fast_forward_key(
14362+
&key,
14363+
now + Duration::from_millis(1)
14364+
));
14365+
14366+
live_mode_release_remote_fast_forward_key(&key);
14367+
}
14368+
1428814369
#[test]
1428914370
fn desktop_activity_records_newest_first_and_caps_entries() {
1429014371
let temp = TestTempDir::new("desktop-activity");

crates/localityd/src/hydration.rs

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -611,6 +611,11 @@ impl HydrationQueue {
611611
self.pending.len()
612612
}
613613

614+
pub fn contains_target(&self, mount_id: &MountId, remote_id: &RemoteId) -> bool {
615+
self.pending
616+
.contains_key(&HydrationKey::new(mount_id.clone(), remote_id.clone()))
617+
}
618+
614619
pub fn is_empty(&self) -> bool {
615620
self.pending.is_empty()
616621
}
@@ -750,12 +755,16 @@ struct HydrationKey {
750755
}
751756

752757
impl HydrationKey {
753-
fn from_request(request: &HydrationRequest) -> Self {
758+
fn new(mount_id: MountId, remote_id: RemoteId) -> Self {
754759
Self {
755-
mount_id: request.mount_id.clone(),
756-
remote_id: request.remote_id.clone(),
760+
mount_id,
761+
remote_id,
757762
}
758763
}
764+
765+
fn from_request(request: &HydrationRequest) -> Self {
766+
Self::new(request.mount_id.clone(), request.remote_id.clone())
767+
}
759768
}
760769

761770
fn merge_request(existing: &mut HydrationRequest, incoming: HydrationRequest) {

0 commit comments

Comments
 (0)