1use crate::commands::cli_session::{cli_request_client, resolve_cli_access_token};
2use crate::config::{
3 resolve_device_identity, resolve_worktree_watch_config, ApiConfig, WorktreeWatchConfig,
4};
5use chrono::{DateTime, Utc};
6use notify::event::{CreateKind, ModifyKind, RemoveKind, RenameMode};
7use notify::{Config, Event, EventKind, RecommendedWatcher, RecursiveMode, Watcher};
8use reqwest::StatusCode;
9use serde::{Deserialize, Serialize};
10use serde_json::{json, Value};
11use sha2::{Digest, Sha256};
12use std::collections::{BTreeMap, BTreeSet};
13use std::fs::{self, File, OpenOptions};
14use std::io::{BufRead, BufReader, Write};
15use std::path::{Path, PathBuf};
16use std::process::{Command, Stdio};
17use std::sync::{mpsc, OnceLock};
18use std::time::{Duration, Instant};
19use sysinfo::{Pid, System};
20use uuid::Uuid;
21
22const BACKGROUND_CHILD_ENV: &str = "XBP_WORKTREE_WATCH_BACKGROUND_CHILD";
23const DEFAULT_REMOTE: &str = "origin";
24const EVENT_DEDUPE_WINDOW_SECONDS: i64 = 5;
25const WORKTREE_MUTATION_SYNC_BATCH_RECORD_LIMIT: usize = 250;
26const WORKTREE_MUTATION_SYNC_BATCH_JSON_BYTES: usize = 768 * 1024;
27const WORKTREE_MUTATION_SYNC_MAX_ATTEMPTS: usize = 3;
28const WORKTREE_MUTATION_SYNC_RETRY_DELAY_SECONDS: u64 = 2;
29
30#[derive(Debug, Clone, Default)]
33struct WatchIgnoreRules {
34 forbidden_path_prefixes: Vec<String>,
36 forbidden_folders: BTreeSet<String>,
38 banned_words: Vec<String>,
40}
41
42impl WatchIgnoreRules {
43 fn load() -> Self {
44 Self::from_config(&resolve_worktree_watch_config())
45 }
46
47 fn from_config(config: &WorktreeWatchConfig) -> Self {
48 let mut forbidden_folders = built_in_forbidden_folders();
49 for folder in &config.forbidden_folders {
50 let normalized = normalize_folder_token(folder);
51 if !normalized.is_empty() {
52 forbidden_folders.insert(normalized);
53 }
54 }
55
56 let forbidden_path_prefixes = config
57 .forbidden_paths
58 .iter()
59 .map(|path| normalize_path_prefix(path))
60 .filter(|path| !path.is_empty())
61 .collect::<Vec<_>>();
62
63 let banned_words = config
64 .banned_words
65 .iter()
66 .map(|word| word.trim().to_ascii_lowercase())
67 .filter(|word| !word.is_empty())
68 .collect::<Vec<_>>();
69
70 Self {
71 forbidden_path_prefixes,
72 forbidden_folders,
73 banned_words,
74 }
75 }
76
77 fn ignores_relative(&self, relative: &str) -> bool {
78 let normalized = normalize_relative_watch_path(relative);
79 if normalized.is_empty() {
80 return false;
81 }
82
83 let lower = normalized.to_ascii_lowercase();
84 let components = lower
85 .split('/')
86 .filter(|component| !component.is_empty())
87 .collect::<Vec<_>>();
88
89 if components
90 .iter()
91 .any(|component| self.forbidden_folders.contains(*component))
92 {
93 return true;
94 }
95
96 for prefix in &self.forbidden_path_prefixes {
97 if lower == *prefix || lower.starts_with(&format!("{prefix}/")) {
98 return true;
99 }
100 if !prefix.contains('/') && components.iter().any(|component| *component == prefix) {
102 return true;
103 }
104 }
105
106 for word in &self.banned_words {
107 if lower.contains(word) {
108 return true;
109 }
110 }
111
112 false
113 }
114
115 fn ignores_path(&self, root: &Path, path: &Path) -> bool {
116 let Ok(relative) = path.strip_prefix(root) else {
117 return false;
118 };
119 self.ignores_relative(&relative.to_string_lossy())
120 }
121
122 fn skips_discovery_dir(&self, path: &Path) -> bool {
123 let Some(name) = path.file_name().and_then(|value| value.to_str()) else {
124 return false;
125 };
126 self.forbidden_folders
127 .contains(&name.trim().to_ascii_lowercase())
128 }
129}
130
131fn built_in_forbidden_folders() -> BTreeSet<String> {
132 [
133 ".git",
134 ".hg",
135 ".svn",
136 ".next",
137 ".turbo",
138 ".vercel",
139 "node_modules",
140 "target",
141 "dist",
142 "build",
143 ]
144 .into_iter()
145 .map(str::to_string)
146 .collect()
147}
148
149fn normalize_folder_token(raw: &str) -> String {
150 raw.trim()
151 .trim_matches(|ch| ch == '/' || ch == '\\')
152 .to_ascii_lowercase()
153}
154
155fn normalize_path_prefix(raw: &str) -> String {
156 let trimmed = raw.trim().replace('\\', "/");
157 let without_dot = trimmed
158 .trim_start_matches("./")
159 .trim_matches(|ch| ch == '/' || ch == '\\');
160 without_dot.to_ascii_lowercase()
161}
162
163fn normalize_relative_watch_path(raw: &str) -> String {
164 raw.trim()
165 .replace('\\', "/")
166 .trim_start_matches("./")
167 .trim_matches('/')
168 .to_string()
169}
170
171fn watch_ignore_rules() -> &'static WatchIgnoreRules {
172 static RULES: OnceLock<WatchIgnoreRules> = OnceLock::new();
173 RULES.get_or_init(WatchIgnoreRules::load)
174}
175
176#[cfg(windows)]
177const CREATE_NO_WINDOW: u32 = 0x08000000;
178
179#[derive(Debug, Clone)]
180pub struct WorktreeWatchTargetOptions {
181 pub repo: Option<PathBuf>,
182 pub parent: Option<PathBuf>,
183}
184
185#[derive(Debug, Clone)]
186pub struct WorktreeWatchStartOptions {
187 pub target: WorktreeWatchTargetOptions,
188 pub detach: bool,
189 pub sync_interval_seconds: u64,
190 pub once: bool,
191}
192
193#[derive(Debug, Clone)]
194pub struct WorktreeWatchSyncOptions {
195 pub target: WorktreeWatchTargetOptions,
196 pub dry_run: bool,
197 pub resync: bool,
198}
199
200#[derive(Debug, Clone)]
201pub struct WorktreeWatchStopOptions {
202 pub target: WorktreeWatchTargetOptions,
203 pub force: bool,
204}
205
206#[derive(Debug, Clone)]
207pub struct WorktreeWatchStatusOptions {
208 pub target: WorktreeWatchTargetOptions,
209 pub json: bool,
210 pub records: bool,
211 pub record_limit: usize,
212 pub stats: bool,
213 pub stats_gap_minutes: u64,
214}
215
216#[derive(Debug, Clone)]
217struct RepoIdentity {
218 owner: String,
219 name: String,
220 branch: String,
221 root: PathBuf,
222 head_sha: Option<String>,
223}
224
225#[derive(Debug, Clone)]
226struct SpoolLayout {
227 root: PathBuf,
228 events_file: PathBuf,
229 commits_file: PathBuf,
230 stats_file: PathBuf,
231 watcher_state_file: PathBuf,
232 sync_log_file: PathBuf,
233 event_dedupe_file: PathBuf,
234 event_dedupe_dir: PathBuf,
235 commit_dedupe_dir: PathBuf,
236}
237
238#[derive(Debug, Clone)]
239struct ParentSpoolLayout {
240 watcher_state_file: PathBuf,
241}
242
243#[derive(Debug)]
244struct WatchedRepo {
245 identity: RepoIdentity,
246 layout: SpoolLayout,
247 last_head: Option<String>,
248 event_dedupe: RecentEventDedupe,
249}
250
251#[derive(Debug, Serialize, Deserialize, Clone)]
252#[serde(rename_all = "camelCase")]
253struct WorktreeMutationEvent {
254 id: String,
255 #[serde(default, skip_serializing_if = "Option::is_none")]
256 fingerprint: Option<String>,
257 repo_owner: String,
258 repo_name: String,
259 branch_name: String,
260 repo_root: String,
261 head_sha: Option<String>,
262 event_kind: String,
263 paths: Vec<String>,
264 primary_path: Option<String>,
265 old_path: Option<String>,
266 new_path: Option<String>,
267 added_lines: Option<u64>,
268 removed_lines: Option<u64>,
269 total_lines: Option<u64>,
270 file_created: bool,
271 file_removed: bool,
272 folder_created: bool,
273 folder_removed: bool,
274 renamed_or_moved: bool,
275 raw_kind: String,
276 occurred_at: DateTime<Utc>,
277}
278
279#[derive(Debug, Serialize, Deserialize, Clone)]
280#[serde(rename_all = "camelCase")]
281struct WorktreeCommitEvent {
282 id: String,
283 #[serde(default, skip_serializing_if = "Option::is_none")]
284 fingerprint: Option<String>,
285 repo_owner: String,
286 repo_name: String,
287 branch_name: String,
288 repo_root: String,
289 previous_head_sha: Option<String>,
290 head_sha: String,
291 subject: Option<String>,
292 author_name: Option<String>,
293 author_email: Option<String>,
294 committed_at: Option<String>,
295 occurred_at: DateTime<Utc>,
296}
297
298#[derive(Debug, Serialize)]
299#[serde(rename_all = "camelCase")]
300struct WorktreeMutationIngestPayload {
301 device: WorktreeMutationDevicePayload,
302 repository: WorktreeMutationRepositoryPayload,
303 events: Vec<WorktreeMutationEvent>,
304 commits: Vec<WorktreeCommitEvent>,
305}
306
307#[derive(Debug, Serialize)]
308#[serde(rename_all = "camelCase")]
309struct WorktreeMutationDevicePayload {
310 hardware_id: String,
311 device_name: Option<String>,
312 hostname: Option<String>,
313 platform: String,
314}
315
316#[derive(Debug, Serialize)]
317#[serde(rename_all = "camelCase")]
318struct WorktreeMutationRepositoryPayload {
319 owner: String,
320 name: String,
321 branch_name: String,
322 repo_root: String,
323}
324
325#[derive(Debug)]
326struct SyncCandidate {
327 path: PathBuf,
328 records: SyncRecords,
329}
330
331#[derive(Debug)]
332struct WorktreeMutationUploadBatch {
333 events: Vec<WorktreeMutationEvent>,
334 commits: Vec<WorktreeCommitEvent>,
335 estimated_json_bytes: usize,
336}
337
338#[derive(Debug, Clone, Copy, PartialEq, Eq)]
339enum WorktreeMutationSyncErrorKind {
340 Retryable,
341 NonRetryable,
342}
343
344#[derive(Debug)]
345struct WorktreeMutationSyncError {
346 kind: WorktreeMutationSyncErrorKind,
347 message: String,
348}
349
350struct SyncLogAppend<'a> {
351 identity: &'a RepoIdentity,
352 layout: &'a SpoolLayout,
353 endpoint: &'a str,
354 status_code: u16,
355 event_count: usize,
356 commit_count: usize,
357 candidates: &'a [SyncCandidate],
358 resync: bool,
359}
360
361#[derive(Debug, Clone, Copy)]
362enum SyncFileKind {
363 Events,
364 Commits,
365}
366
367#[derive(Debug)]
368enum SyncRecords {
369 Events(Vec<WorktreeMutationEvent>),
370 Commits(Vec<WorktreeCommitEvent>),
371}
372
373impl SyncRecords {
374 fn len(&self) -> usize {
375 match self {
376 Self::Events(records) => records.len(),
377 Self::Commits(records) => records.len(),
378 }
379 }
380
381 fn is_empty(&self) -> bool {
382 self.len() == 0
383 }
384}
385
386#[derive(Debug, Serialize, Clone)]
387#[serde(rename_all = "camelCase")]
388struct WorktreeWatchStatsSummary {
389 generated_at: DateTime<Utc>,
390 repository_owner: String,
391 repository_name: String,
392 branch_name: String,
393 repo_root: String,
394 spool_root: String,
395 session_gap_minutes: u64,
396 total_events: u64,
397 total_commits: u64,
398 first_event_at: Option<DateTime<Utc>>,
399 last_event_at: Option<DateTime<Utc>>,
400 estimated_coding_seconds: u64,
401 added_lines: u64,
402 removed_lines: u64,
403 by_file: Vec<WorktreeWatchFileStats>,
404 by_filetype: Vec<WorktreeWatchBucketStats>,
405 by_event_kind: Vec<WorktreeWatchBucketStats>,
406}
407
408#[derive(Debug, Serialize, Clone, Default)]
409#[serde(rename_all = "camelCase")]
410struct WorktreeWatchFileStats {
411 path: String,
412 filetype: String,
413 event_count: u64,
414 estimated_coding_seconds: u64,
415 added_lines: u64,
416 removed_lines: u64,
417}
418
419#[derive(Debug, Serialize, Clone, Default)]
420#[serde(rename_all = "camelCase")]
421struct WorktreeWatchBucketStats {
422 name: String,
423 event_count: u64,
424 estimated_coding_seconds: u64,
425 added_lines: u64,
426 removed_lines: u64,
427}
428
429#[derive(Debug, Serialize, Deserialize, Clone)]
430#[serde(rename_all = "camelCase")]
431struct WorktreeWatcherState {
432 pid: u32,
433 repo_root: String,
434 repository_owner: String,
435 repository_name: String,
436 branch_name: String,
437 executable: String,
438 started_at: DateTime<Utc>,
439}
440
441#[derive(Debug, Serialize, Deserialize, Clone)]
442#[serde(rename_all = "camelCase")]
443struct ParentWorktreeWatcherState {
444 pid: u32,
445 parent_root: String,
446 executable: String,
447 started_at: DateTime<Utc>,
448}
449
450#[derive(Debug, Serialize, Deserialize, Clone)]
451#[serde(rename_all = "camelCase")]
452struct WorktreeWatchSyncLogEntry {
453 id: String,
454 synced_at: DateTime<Utc>,
455 endpoint: String,
456 repository_owner: String,
457 repository_name: String,
458 branch_name: String,
459 repo_root: String,
460 resync: bool,
461 status_code: u16,
462 event_count: usize,
463 commit_count: usize,
464 spool_file_count: usize,
465 spool_files: Vec<String>,
466}
467
468#[derive(Debug, Clone, Copy, PartialEq, Eq)]
469enum StopWatcherOutcome {
470 Stopped(u32),
471 RemovedStaleState,
472 NoState,
473}
474
475#[derive(Debug, Serialize, Deserialize, Clone)]
476#[serde(rename_all = "camelCase")]
477struct WorktreeWatchEventDedupeRecord {
478 fingerprint: String,
479 first_seen_at: DateTime<Utc>,
480 event_kind: String,
481 primary_path: Option<String>,
482}
483
484#[derive(Debug, Default)]
485struct RecentEventDedupe {
486 fingerprints: BTreeSet<String>,
487}
488
489pub async fn run_worktree_watch_start(options: WorktreeWatchStartOptions) -> Result<(), String> {
490 if let Some(parent) = options.target.parent.as_deref() {
491 let identities = resolve_target_identities(&options.target)?;
492 if identities.is_empty() {
493 println!("No git repositories found under parent folder.");
494 return Ok(());
495 }
496
497 if options.detach {
498 let path = spawn_detached_parent_worktree_watch(parent)?;
499 println!(
500 "Started one background worktree watcher for {} covering {} repo(s).",
501 path.display(),
502 identities.len()
503 );
504 return Ok(());
505 }
506
507 if options.once {
508 for identity in &identities {
509 let layout = prepare_spool_layout(identity)?;
510 record_commit_snapshot(identity, &layout, None)?;
511 println!("Recorded current git state for {}", identity.root.display());
512 }
513 return Ok(());
514 }
515
516 let parent = canonical_parent_path(parent)?;
517 println!(
518 "Watching {} recursively and routing mutations for {} repo(s).",
519 parent.display(),
520 identities.len()
521 );
522 return watch_parent_foreground(parent, identities, options.sync_interval_seconds).await;
523 }
524
525 if options.detach {
526 return spawn_detached_worktree_watch(options.target.repo.as_deref()).map(|path| {
527 println!("Started background worktree watcher for {}", path.display());
528 });
529 }
530
531 let identity = resolve_repo_identity(options.target.repo.as_deref())?;
532 let layout = prepare_spool_layout(&identity)?;
533 println!(
534 "Watching {} and spooling mutations to {}",
535 identity.root.display(),
536 layout.root.display()
537 );
538
539 if options.once {
540 record_commit_snapshot(&identity, &layout, None)?;
541 return Ok(());
542 }
543
544 watch_foreground(identity, layout, options.sync_interval_seconds).await
545}
546
547pub async fn run_worktree_watch_sync(options: WorktreeWatchSyncOptions) -> Result<(), String> {
548 let identities = resolve_target_identities(&options.target)?;
549 let repo_count = identities.len();
550 let mut plans = Vec::new();
551 let mut total_files = 0usize;
552 let mut total_records = 0usize;
553
554 for identity in identities {
555 let layout = prepare_spool_layout(&identity)?;
556 let candidates = collect_spool_candidates(&layout, options.resync)?;
557 total_files += candidates.len();
558 total_records += candidates
559 .iter()
560 .map(|candidate| candidate.records.len())
561 .sum::<usize>();
562 if !candidates.is_empty() {
563 plans.push((identity, candidates));
564 }
565 }
566
567 if options.dry_run {
568 println!(
569 "Would sync {} record(s) from {} spool file(s) across {} repo(s).",
570 total_records, total_files, repo_count
571 );
572 if options.resync {
573 println!("Resync mode includes files already marked `.synced.*.jsonl`.");
574 }
575 return Ok(());
576 }
577
578 if plans.is_empty() {
579 print_no_sync_candidates_message(&options.target, options.resync)?;
580 return Ok(());
581 }
582
583 let upload_repo_count = plans.len();
584 for (identity, candidates) in plans {
585 sync_candidates(&identity, candidates, options.resync).await?;
586 }
587 println!(
588 "Synced {} worktree mutation record(s) across {} repo(s).",
589 total_records, upload_repo_count
590 );
591 Ok(())
592}
593
594pub fn run_worktree_watch_stop(options: WorktreeWatchStopOptions) -> Result<(), String> {
595 if let Some(parent) = options.target.parent.as_deref() {
596 let parent = canonical_parent_path(parent)?;
597 match stop_existing_parent_watcher(&parent, options.force)? {
598 StopWatcherOutcome::Stopped(pid) => {
599 println!(
600 "Stopped background parent worktree watcher {} for {}",
601 pid,
602 parent.display()
603 );
604 }
605 StopWatcherOutcome::RemovedStaleState => {
606 println!(
607 "Removed stale parent worktree watcher state for {}",
608 parent.display()
609 );
610 }
611 StopWatcherOutcome::NoState => {
612 println!("No parent worktree watcher state for {}", parent.display());
613 }
614 }
615 }
616
617 let identities = resolve_target_identities(&options.target)?;
618 if identities.is_empty() {
619 println!("No git repositories found for worktree watcher stop.");
620 return Ok(());
621 }
622
623 let mut stopped = 0usize;
624 let mut stale = 0usize;
625 let mut missing = 0usize;
626 for identity in identities {
627 let layout = prepare_spool_layout(&identity)?;
628 match stop_existing_watcher(&identity, &layout, options.force)? {
629 StopWatcherOutcome::Stopped(pid) => {
630 stopped += 1;
631 println!(
632 "Stopped background worktree watcher {} for {}",
633 pid,
634 identity.root.display()
635 );
636 }
637 StopWatcherOutcome::RemovedStaleState => {
638 stale += 1;
639 println!(
640 "Removed stale worktree watcher state for {}",
641 identity.root.display()
642 );
643 }
644 StopWatcherOutcome::NoState => {
645 missing += 1;
646 println!(
647 "No background worktree watcher state for {}",
648 identity.root.display()
649 );
650 }
651 }
652 }
653
654 println!(
655 "Worktree watcher stop complete: {stopped} stopped, {stale} stale state file(s) removed, {missing} without state."
656 );
657 Ok(())
658}
659
660pub async fn run_worktree_watch_status(options: WorktreeWatchStatusOptions) -> Result<(), String> {
661 let identities = resolve_target_identities(&options.target)?;
662 let mut payloads = Vec::new();
663 let mut total_files = 0usize;
664 let mut total_records = 0usize;
665
666 for identity in identities {
667 let status =
668 worktree_watch_status_payload(&identity, options.records, options.record_limit)?;
669 let stats = if options.stats {
670 Some(generate_and_store_stats(
671 &identity,
672 options.stats_gap_minutes,
673 )?)
674 } else {
675 None
676 };
677 let status = attach_optional_stats(status, stats)?;
678 total_files += status
679 .get("unsyncedFiles")
680 .and_then(Value::as_u64)
681 .unwrap_or(0) as usize;
682 total_records += status
683 .get("unsyncedRecords")
684 .and_then(Value::as_u64)
685 .unwrap_or(0) as usize;
686 payloads.push(status);
687 }
688
689 if options.json {
690 let payload = if options.target.parent.is_some() {
691 json!({
692 "repositories": payloads,
693 "repositoryCount": payloads.len(),
694 "unsyncedFiles": total_files,
695 "unsyncedRecords": total_records,
696 })
697 } else {
698 payloads.into_iter().next().unwrap_or_else(|| json!({}))
699 };
700 println!(
701 "{}",
702 serde_json::to_string_pretty(&payload)
703 .map_err(|error| format!("Failed to render status JSON: {error}"))?
704 );
705 } else {
706 for (index, payload) in payloads.iter().enumerate() {
707 if index > 0 {
708 println!();
709 }
710 println!(
711 "repo: {}/{}",
712 payload
713 .get("repositoryOwner")
714 .and_then(Value::as_str)
715 .unwrap_or("unknown"),
716 payload
717 .get("repositoryName")
718 .and_then(Value::as_str)
719 .unwrap_or("repository")
720 );
721 println!(
722 "branch: {}",
723 payload
724 .get("branchName")
725 .and_then(Value::as_str)
726 .unwrap_or("unknown")
727 );
728 println!(
729 "root: {}",
730 payload
731 .get("repoRoot")
732 .and_then(Value::as_str)
733 .unwrap_or("")
734 );
735 println!(
736 "spool: {}",
737 payload
738 .get("spoolRoot")
739 .and_then(Value::as_str)
740 .unwrap_or("")
741 );
742 println!(
743 "unsynced files: {}",
744 payload
745 .get("unsyncedFiles")
746 .and_then(Value::as_u64)
747 .unwrap_or(0)
748 );
749 println!(
750 "unsynced records: {}",
751 payload
752 .get("unsyncedRecords")
753 .and_then(Value::as_u64)
754 .unwrap_or(0)
755 );
756 println!(
757 "synced files: {}",
758 payload
759 .get("syncedFiles")
760 .and_then(Value::as_u64)
761 .unwrap_or(0)
762 );
763 println!(
764 "synced records: {}",
765 payload
766 .get("syncedRecords")
767 .and_then(Value::as_u64)
768 .unwrap_or(0)
769 );
770 print_last_sync(payload);
771 if options.records {
772 print_status_records(payload);
773 }
774 if options.stats {
775 print_status_stats(payload);
776 }
777 }
778 if options.target.parent.is_some() {
779 println!();
780 println!(
781 "total: {} repo(s), {} unsynced file(s), {} unsynced record(s)",
782 payloads.len(),
783 total_files,
784 total_records
785 );
786 }
787 }
788
789 Ok(())
790}
791
792fn attach_optional_stats(
793 mut payload: Value,
794 stats: Option<WorktreeWatchStatsSummary>,
795) -> Result<Value, String> {
796 if let Some(stats) = stats {
797 if let Some(object) = payload.as_object_mut() {
798 object.insert(
799 "stats".to_string(),
800 serde_json::to_value(stats)
801 .map_err(|error| format!("Failed to render stats payload: {error}"))?,
802 );
803 }
804 }
805 Ok(payload)
806}
807
808pub fn spawn_detached_worktree_watch_for_current_repo() -> Result<Option<PathBuf>, String> {
809 if is_background_child() {
810 return Ok(None);
811 }
812
813 let identity = match resolve_repo_identity(None) {
814 Ok(identity) => identity,
815 Err(_) => return Ok(None),
816 };
817 spawn_detached_worktree_watch_for_identity(&identity).map(Some)
818}
819
820fn resolve_target_identities(
821 target: &WorktreeWatchTargetOptions,
822) -> Result<Vec<RepoIdentity>, String> {
823 if target.repo.is_some() && target.parent.is_some() {
824 return Err("Pass either `--repo` or `--parent`, not both.".to_string());
825 }
826
827 if let Some(parent) = target.parent.as_deref() {
828 return discover_repo_identities_under(parent);
829 }
830
831 resolve_repo_identity(target.repo.as_deref()).map(|identity| vec![identity])
832}
833
834fn worktree_watch_status_payload(
835 identity: &RepoIdentity,
836 include_records: bool,
837 record_limit: usize,
838) -> Result<Value, String> {
839 let layout = prepare_spool_layout(identity)?;
840 let candidates = collect_sync_candidates(&layout)?;
841 let all_candidates = collect_spool_candidates(&layout, true)?;
842 let record_count: usize = candidates
843 .iter()
844 .map(|candidate| candidate.records.len())
845 .sum();
846 let all_record_count: usize = all_candidates
847 .iter()
848 .map(|candidate| candidate.records.len())
849 .sum();
850 let synced_file_count = all_candidates.len().saturating_sub(candidates.len());
851 let synced_record_count = all_record_count.saturating_sub(record_count);
852
853 let mut payload = json!({
854 "repoRoot": identity.root.display().to_string(),
855 "repositoryOwner": identity.owner.clone(),
856 "repositoryName": identity.name.clone(),
857 "branchName": identity.branch.clone(),
858 "spoolRoot": layout.root.display().to_string(),
859 "unsyncedFiles": candidates.len(),
860 "unsyncedRecords": record_count,
861 "syncedFiles": synced_file_count,
862 "syncedRecords": synced_record_count,
863 "localSpoolFiles": all_candidates.len(),
864 "localSpoolRecords": all_record_count,
865 });
866
867 if let Some(last_sync) = read_last_sync_log_entry(&layout.sync_log_file)? {
868 if let Some(object) = payload.as_object_mut() {
869 object.insert(
870 "lastSync".to_string(),
871 serde_json::to_value(last_sync)
872 .map_err(|error| format!("Failed to render sync log payload: {error}"))?,
873 );
874 }
875 }
876
877 if include_records {
878 if let Some(object) = payload.as_object_mut() {
879 object.insert(
880 "records".to_string(),
881 collect_status_records(&candidates, record_limit)?,
882 );
883 }
884 }
885
886 Ok(payload)
887}
888
889fn collect_status_records(candidates: &[SyncCandidate], limit: usize) -> Result<Value, String> {
890 let available: usize = candidates
891 .iter()
892 .map(|candidate| candidate.records.len())
893 .sum();
894 let mut shown = 0usize;
895 let mut events = Vec::new();
896 let mut commits = Vec::new();
897
898 'candidates: for candidate in candidates {
899 match &candidate.records {
900 SyncRecords::Events(records) => {
901 for event in records {
902 if shown >= limit {
903 break 'candidates;
904 }
905 events.push(json!({
906 "sourceFile": candidate.path.display().to_string(),
907 "occurredAt": event.occurred_at,
908 "eventKind": event.event_kind.clone(),
909 "primaryPath": event.primary_path.clone(),
910 "oldPath": event.old_path.clone(),
911 "newPath": event.new_path.clone(),
912 "paths": event.paths.clone(),
913 "addedLines": event.added_lines,
914 "removedLines": event.removed_lines,
915 "totalLines": event.total_lines,
916 "headSha": event.head_sha.clone(),
917 "rawKind": event.raw_kind.clone(),
918 }));
919 shown += 1;
920 }
921 }
922 SyncRecords::Commits(records) => {
923 for commit in records {
924 if shown >= limit {
925 break 'candidates;
926 }
927 commits.push(json!({
928 "sourceFile": candidate.path.display().to_string(),
929 "occurredAt": commit.occurred_at,
930 "headSha": commit.head_sha.clone(),
931 "previousHeadSha": commit.previous_head_sha.clone(),
932 "subject": commit.subject.clone(),
933 "authorName": commit.author_name.clone(),
934 "authorEmail": commit.author_email.clone(),
935 "committedAt": commit.committed_at.clone(),
936 }));
937 shown += 1;
938 }
939 }
940 }
941 }
942
943 Ok(json!({
944 "available": available,
945 "shown": shown,
946 "truncated": shown < available,
947 "events": events,
948 "commits": commits,
949 }))
950}
951
952fn print_status_records(payload: &Value) {
953 let Some(records) = payload.get("records") else {
954 return;
955 };
956 let shown = records.get("shown").and_then(Value::as_u64).unwrap_or(0);
957 if shown == 0 {
958 println!("records: none");
959 return;
960 }
961
962 println!("records:");
963 if let Some(events) = records.get("events").and_then(Value::as_array) {
964 if !events.is_empty() {
965 println!(" mutations:");
966 for event in events {
967 print_status_event_record(event);
968 }
969 }
970 }
971 if let Some(commits) = records.get("commits").and_then(Value::as_array) {
972 if !commits.is_empty() {
973 println!(" commits:");
974 for commit in commits {
975 print_status_commit_record(commit);
976 }
977 }
978 }
979 if records
980 .get("truncated")
981 .and_then(Value::as_bool)
982 .unwrap_or(false)
983 {
984 let available = records
985 .get("available")
986 .and_then(Value::as_u64)
987 .unwrap_or(shown);
988 println!(" showing {shown} of {available} record(s)");
989 }
990}
991
992fn print_status_stats(payload: &Value) {
993 let Some(stats) = payload.get("stats") else {
994 return;
995 };
996 println!("stats:");
997 println!(
998 " coding time: {}",
999 format_duration(
1000 stats
1001 .get("estimatedCodingSeconds")
1002 .and_then(Value::as_u64)
1003 .unwrap_or(0)
1004 )
1005 );
1006 println!(
1007 " events: {}",
1008 stats
1009 .get("totalEvents")
1010 .and_then(Value::as_u64)
1011 .unwrap_or(0)
1012 );
1013 println!(
1014 " commits: {}",
1015 stats
1016 .get("totalCommits")
1017 .and_then(Value::as_u64)
1018 .unwrap_or(0)
1019 );
1020 println!(
1021 " lines: +{} -{}",
1022 stats.get("addedLines").and_then(Value::as_u64).unwrap_or(0),
1023 stats
1024 .get("removedLines")
1025 .and_then(Value::as_u64)
1026 .unwrap_or(0)
1027 );
1028 if let Some(spool_root) = value_str(stats, "spoolRoot") {
1029 println!(
1030 " stored: {}",
1031 Path::new(spool_root).join("stats.json").display()
1032 );
1033 }
1034 print_stats_bucket_section(stats, "byFiletype", " by filetype:", 10);
1035 print_stats_bucket_section(stats, "byEventKind", " by event kind:", 10);
1036 print_stats_file_section(stats, 10);
1037}
1038
1039fn print_last_sync(payload: &Value) {
1040 let Some(last_sync) = payload.get("lastSync") else {
1041 return;
1042 };
1043 let synced_at = value_str(last_sync, "syncedAt").unwrap_or("unknown-time");
1044 let status = last_sync
1045 .get("statusCode")
1046 .and_then(Value::as_u64)
1047 .unwrap_or(0);
1048 let events = last_sync
1049 .get("eventCount")
1050 .and_then(Value::as_u64)
1051 .unwrap_or(0);
1052 let commits = last_sync
1053 .get("commitCount")
1054 .and_then(Value::as_u64)
1055 .unwrap_or(0);
1056 let files = last_sync
1057 .get("spoolFileCount")
1058 .and_then(Value::as_u64)
1059 .unwrap_or(0);
1060 let resync = last_sync
1061 .get("resync")
1062 .and_then(Value::as_bool)
1063 .unwrap_or(false);
1064 println!(
1065 "last sync: {synced_at}, HTTP {status}, {events} event(s), {commits} commit(s), {files} file(s), resync={resync}"
1066 );
1067}
1068
1069fn print_stats_bucket_section(stats: &Value, key: &str, title: &str, limit: usize) {
1070 let Some(items) = stats.get(key).and_then(Value::as_array) else {
1071 return;
1072 };
1073 if items.is_empty() {
1074 return;
1075 }
1076 println!("{title}");
1077 for item in items.iter().take(limit) {
1078 let name = value_str(item, "name").unwrap_or("unknown");
1079 let events = item.get("eventCount").and_then(Value::as_u64).unwrap_or(0);
1080 let seconds = item
1081 .get("estimatedCodingSeconds")
1082 .and_then(Value::as_u64)
1083 .unwrap_or(0);
1084 let added = item.get("addedLines").and_then(Value::as_u64).unwrap_or(0);
1085 let removed = item
1086 .get("removedLines")
1087 .and_then(Value::as_u64)
1088 .unwrap_or(0);
1089 println!(
1090 " {name}: {}, {} event(s), +{} -{}",
1091 format_duration(seconds),
1092 events,
1093 added,
1094 removed
1095 );
1096 }
1097}
1098
1099fn print_stats_file_section(stats: &Value, limit: usize) {
1100 let Some(items) = stats.get("byFile").and_then(Value::as_array) else {
1101 return;
1102 };
1103 if items.is_empty() {
1104 return;
1105 }
1106 println!(" by file:");
1107 for item in items.iter().take(limit) {
1108 let path = value_str(item, "path").unwrap_or("(no path)");
1109 let filetype = value_str(item, "filetype").unwrap_or("(none)");
1110 let events = item.get("eventCount").and_then(Value::as_u64).unwrap_or(0);
1111 let seconds = item
1112 .get("estimatedCodingSeconds")
1113 .and_then(Value::as_u64)
1114 .unwrap_or(0);
1115 let added = item.get("addedLines").and_then(Value::as_u64).unwrap_or(0);
1116 let removed = item
1117 .get("removedLines")
1118 .and_then(Value::as_u64)
1119 .unwrap_or(0);
1120 println!(
1121 " {path} [{filetype}]: {}, {} event(s), +{} -{}",
1122 format_duration(seconds),
1123 events,
1124 added,
1125 removed
1126 );
1127 }
1128}
1129
1130fn print_status_event_record(event: &Value) {
1131 let occurred_at = value_str(event, "occurredAt").unwrap_or("unknown-time");
1132 let kind = value_str(event, "eventKind").unwrap_or("event");
1133 let primary_path = event_display_path(event);
1134 let added = event.get("addedLines").and_then(Value::as_u64).unwrap_or(0);
1135 let removed = event
1136 .get("removedLines")
1137 .and_then(Value::as_u64)
1138 .unwrap_or(0);
1139 let total_lines = event.get("totalLines").and_then(Value::as_u64);
1140 if let Some(total_lines) = total_lines {
1141 println!(
1142 " {occurred_at} {kind} {primary_path} (+{added} -{removed}, {total_lines} lines)"
1143 );
1144 } else {
1145 println!(" {occurred_at} {kind} {primary_path} (+{added} -{removed})");
1146 }
1147
1148 if let Some(head_sha) = value_str(event, "headSha") {
1149 println!(" head: {}", short_sha(head_sha));
1150 }
1151 if let Some(paths) = event.get("paths").and_then(Value::as_array) {
1152 if paths.len() > 1 {
1153 let rendered = paths
1154 .iter()
1155 .filter_map(Value::as_str)
1156 .collect::<Vec<_>>()
1157 .join(", ");
1158 println!(" paths: {rendered}");
1159 }
1160 }
1161}
1162
1163fn print_status_commit_record(commit: &Value) {
1164 let occurred_at = value_str(commit, "occurredAt").unwrap_or("unknown-time");
1165 let head_sha = value_str(commit, "headSha")
1166 .map(short_sha)
1167 .unwrap_or_else(|| "unknown".to_string());
1168 let subject = value_str(commit, "subject").unwrap_or("(no subject)");
1169 println!(" {occurred_at} {head_sha} {subject}");
1170
1171 if let Some(author) = value_str(commit, "authorName") {
1172 println!(" author: {author}");
1173 }
1174}
1175
1176fn event_display_path(event: &Value) -> String {
1177 let old_path = value_str(event, "oldPath");
1178 let new_path = value_str(event, "newPath");
1179 if let (Some(old_path), Some(new_path)) = (old_path, new_path) {
1180 return format!("{old_path} -> {new_path}");
1181 }
1182 value_str(event, "primaryPath")
1183 .map(str::to_string)
1184 .or_else(|| {
1185 event
1186 .get("paths")
1187 .and_then(Value::as_array)
1188 .and_then(|paths| paths.first())
1189 .and_then(Value::as_str)
1190 .map(str::to_string)
1191 })
1192 .unwrap_or_else(|| "(no path)".to_string())
1193}
1194
1195fn value_str<'a>(value: &'a Value, key: &str) -> Option<&'a str> {
1196 value.get(key).and_then(Value::as_str)
1197}
1198
1199fn short_sha(value: &str) -> String {
1200 value.chars().take(12).collect()
1201}
1202
1203fn generate_and_store_stats(
1204 identity: &RepoIdentity,
1205 session_gap_minutes: u64,
1206) -> Result<WorktreeWatchStatsSummary, String> {
1207 let layout = prepare_spool_layout(identity)?;
1208 let candidates = collect_spool_candidates(&layout, true)?;
1209 let summary = build_stats_summary(identity, &layout, &candidates, session_gap_minutes)?;
1210 let rendered = serde_json::to_string_pretty(&summary)
1211 .map_err(|error| format!("Failed to serialize stats: {error}"))?;
1212 fs::write(&layout.stats_file, format!("{rendered}\n")).map_err(|error| {
1213 format!(
1214 "Failed to write stats file {}: {error}",
1215 layout.stats_file.display()
1216 )
1217 })?;
1218 Ok(summary)
1219}
1220
1221fn build_stats_summary(
1222 identity: &RepoIdentity,
1223 layout: &SpoolLayout,
1224 candidates: &[SyncCandidate],
1225 session_gap_minutes: u64,
1226) -> Result<WorktreeWatchStatsSummary, String> {
1227 let mut events = Vec::new();
1228 let mut total_commits = 0u64;
1229 for candidate in candidates {
1230 match &candidate.records {
1231 SyncRecords::Events(records) => events.extend(records.iter().cloned()),
1232 SyncRecords::Commits(records) => total_commits += records.len() as u64,
1233 }
1234 }
1235 events.sort_by_key(|event| event.occurred_at);
1236
1237 let gap_seconds = session_gap_minutes.saturating_mul(60).max(60);
1238 let mut previous_at: Option<DateTime<Utc>> = None;
1239 let mut total_seconds = 0u64;
1240 let mut total_added = 0u64;
1241 let mut total_removed = 0u64;
1242 let mut by_file: BTreeMap<String, WorktreeWatchFileStats> = BTreeMap::new();
1243 let mut by_filetype: BTreeMap<String, WorktreeWatchBucketStats> = BTreeMap::new();
1244 let mut by_event_kind: BTreeMap<String, WorktreeWatchBucketStats> = BTreeMap::new();
1245
1246 for event in &events {
1247 let seconds = coding_seconds_for_event(previous_at, event.occurred_at, gap_seconds);
1248 previous_at = Some(event.occurred_at);
1249 let added = event.added_lines.unwrap_or(0);
1250 let removed = event.removed_lines.unwrap_or(0);
1251 let path = stats_event_path(event);
1252 let filetype = filetype_for_path(&path);
1253
1254 total_seconds += seconds;
1255 total_added += added;
1256 total_removed += removed;
1257 update_file_stats(&mut by_file, &path, &filetype, seconds, added, removed);
1258 update_bucket_stats(&mut by_filetype, &filetype, seconds, added, removed);
1259 update_bucket_stats(
1260 &mut by_event_kind,
1261 &event.event_kind,
1262 seconds,
1263 added,
1264 removed,
1265 );
1266 }
1267
1268 Ok(WorktreeWatchStatsSummary {
1269 generated_at: Utc::now(),
1270 repository_owner: identity.owner.clone(),
1271 repository_name: identity.name.clone(),
1272 branch_name: identity.branch.clone(),
1273 repo_root: identity.root.display().to_string(),
1274 spool_root: layout.root.display().to_string(),
1275 session_gap_minutes,
1276 total_events: events.len() as u64,
1277 total_commits,
1278 first_event_at: events.first().map(|event| event.occurred_at),
1279 last_event_at: events.last().map(|event| event.occurred_at),
1280 estimated_coding_seconds: total_seconds,
1281 added_lines: total_added,
1282 removed_lines: total_removed,
1283 by_file: sorted_file_stats(by_file),
1284 by_filetype: sorted_bucket_stats(by_filetype),
1285 by_event_kind: sorted_bucket_stats(by_event_kind),
1286 })
1287}
1288
1289fn coding_seconds_for_event(
1290 previous_at: Option<DateTime<Utc>>,
1291 occurred_at: DateTime<Utc>,
1292 gap_seconds: u64,
1293) -> u64 {
1294 let Some(previous_at) = previous_at else {
1295 return 60;
1296 };
1297 let delta = occurred_at.signed_duration_since(previous_at).num_seconds();
1298 if delta <= 0 {
1299 0
1300 } else {
1301 (delta as u64).min(gap_seconds)
1302 }
1303}
1304
1305fn update_file_stats(
1306 by_file: &mut BTreeMap<String, WorktreeWatchFileStats>,
1307 path: &str,
1308 filetype: &str,
1309 seconds: u64,
1310 added: u64,
1311 removed: u64,
1312) {
1313 let stats = by_file
1314 .entry(path.to_string())
1315 .or_insert_with(|| WorktreeWatchFileStats {
1316 path: path.to_string(),
1317 filetype: filetype.to_string(),
1318 ..WorktreeWatchFileStats::default()
1319 });
1320 stats.event_count += 1;
1321 stats.estimated_coding_seconds += seconds;
1322 stats.added_lines += added;
1323 stats.removed_lines += removed;
1324}
1325
1326fn update_bucket_stats(
1327 buckets: &mut BTreeMap<String, WorktreeWatchBucketStats>,
1328 name: &str,
1329 seconds: u64,
1330 added: u64,
1331 removed: u64,
1332) {
1333 let stats = buckets
1334 .entry(name.to_string())
1335 .or_insert_with(|| WorktreeWatchBucketStats {
1336 name: name.to_string(),
1337 ..WorktreeWatchBucketStats::default()
1338 });
1339 stats.event_count += 1;
1340 stats.estimated_coding_seconds += seconds;
1341 stats.added_lines += added;
1342 stats.removed_lines += removed;
1343}
1344
1345fn sorted_file_stats(
1346 by_file: BTreeMap<String, WorktreeWatchFileStats>,
1347) -> Vec<WorktreeWatchFileStats> {
1348 let mut stats = by_file.into_values().collect::<Vec<_>>();
1349 stats.sort_by(|left, right| {
1350 right
1351 .estimated_coding_seconds
1352 .cmp(&left.estimated_coding_seconds)
1353 .then_with(|| right.event_count.cmp(&left.event_count))
1354 .then_with(|| left.path.cmp(&right.path))
1355 });
1356 stats
1357}
1358
1359fn sorted_bucket_stats(
1360 buckets: BTreeMap<String, WorktreeWatchBucketStats>,
1361) -> Vec<WorktreeWatchBucketStats> {
1362 let mut stats = buckets.into_values().collect::<Vec<_>>();
1363 stats.sort_by(|left, right| {
1364 right
1365 .estimated_coding_seconds
1366 .cmp(&left.estimated_coding_seconds)
1367 .then_with(|| right.event_count.cmp(&left.event_count))
1368 .then_with(|| left.name.cmp(&right.name))
1369 });
1370 stats
1371}
1372
1373fn stats_event_path(event: &WorktreeMutationEvent) -> String {
1374 event
1375 .new_path
1376 .as_deref()
1377 .or(event.primary_path.as_deref())
1378 .or_else(|| event.paths.first().map(String::as_str))
1379 .unwrap_or("(no path)")
1380 .to_string()
1381}
1382
1383fn filetype_for_path(path: &str) -> String {
1384 Path::new(path)
1385 .extension()
1386 .and_then(|value| value.to_str())
1387 .map(|value| value.to_ascii_lowercase())
1388 .filter(|value| !value.is_empty())
1389 .unwrap_or_else(|| "(none)".to_string())
1390}
1391
1392fn format_duration(seconds: u64) -> String {
1393 let hours = seconds / 3600;
1394 let minutes = (seconds % 3600) / 60;
1395 let remaining_seconds = seconds % 60;
1396 if hours > 0 {
1397 format!("{hours}h {minutes}m")
1398 } else if minutes > 0 {
1399 format!("{minutes}m {remaining_seconds}s")
1400 } else {
1401 format!("{remaining_seconds}s")
1402 }
1403}
1404
1405async fn watch_foreground(
1406 identity: RepoIdentity,
1407 layout: SpoolLayout,
1408 sync_interval_seconds: u64,
1409) -> Result<(), String> {
1410 let (tx, rx) = mpsc::channel();
1411 let mut watcher = RecommendedWatcher::new(
1412 move |result| {
1413 let _ = tx.send(result);
1414 },
1415 Config::default(),
1416 )
1417 .map_err(|error| format!("Failed to create filesystem watcher: {error}"))?;
1418
1419 watcher
1420 .watch(&identity.root, RecursiveMode::Recursive)
1421 .map_err(|error| {
1422 format!(
1423 "Failed to watch repository root {}: {error}",
1424 identity.root.display()
1425 )
1426 })?;
1427
1428 let mut last_head = identity.head_sha.clone();
1429 let mut last_commit_check = Instant::now();
1430 let mut last_sync = Instant::now();
1431 let mut event_dedupe = RecentEventDedupe {
1432 fingerprints: load_event_dedupe_fingerprints(&layout.event_dedupe_file)?,
1433 };
1434
1435 loop {
1436 match rx.recv_timeout(Duration::from_millis(750)) {
1437 Ok(Ok(event)) => {
1438 if let Some(mutation) = mutation_event_from_notify(&identity, event) {
1439 append_unique_mutation_event(&layout, &mutation, &mut event_dedupe)?;
1440 }
1441 }
1442 Ok(Err(error)) => {
1443 eprintln!("worktree watcher error: {error}");
1444 }
1445 Err(mpsc::RecvTimeoutError::Timeout) => {}
1446 Err(mpsc::RecvTimeoutError::Disconnected) => {
1447 return Err("Filesystem watcher disconnected.".to_string());
1448 }
1449 }
1450
1451 if last_commit_check.elapsed() >= Duration::from_secs(3) {
1452 last_head = record_commit_snapshot(&identity, &layout, last_head.as_deref())?;
1453 last_commit_check = Instant::now();
1454 }
1455
1456 if sync_interval_seconds > 0
1457 && last_sync.elapsed() >= Duration::from_secs(sync_interval_seconds)
1458 {
1459 if let Err(error) = sync_spool_for_identity(&identity, &layout).await {
1460 eprintln!("worktree mutation sync failed: {error}");
1461 }
1462 last_sync = Instant::now();
1463 }
1464 }
1465}
1466
1467async fn watch_parent_foreground(
1468 parent_root: PathBuf,
1469 identities: Vec<RepoIdentity>,
1470 sync_interval_seconds: u64,
1471) -> Result<(), String> {
1472 let (tx, rx) = mpsc::channel();
1473 let mut watcher = RecommendedWatcher::new(
1474 move |result| {
1475 let _ = tx.send(result);
1476 },
1477 Config::default(),
1478 )
1479 .map_err(|error| format!("Failed to create parent filesystem watcher: {error}"))?;
1480
1481 watcher
1482 .watch(&parent_root, RecursiveMode::Recursive)
1483 .map_err(|error| {
1484 format!(
1485 "Failed to watch parent folder {}: {error}",
1486 parent_root.display()
1487 )
1488 })?;
1489
1490 let mut repos = prepare_watched_repos(identities)?;
1491 let mut last_commit_check = Instant::now();
1492 let mut last_sync = Instant::now();
1493
1494 loop {
1495 match rx.recv_timeout(Duration::from_millis(750)) {
1496 Ok(Ok(event)) => {
1497 for (repo_index, routed_event) in route_parent_event(&repos, &event) {
1498 let repo = &mut repos[repo_index];
1499 if let Some(mutation) = mutation_event_from_notify(&repo.identity, routed_event)
1500 {
1501 append_unique_mutation_event(
1502 &repo.layout,
1503 &mutation,
1504 &mut repo.event_dedupe,
1505 )?;
1506 }
1507 }
1508 }
1509 Ok(Err(error)) => {
1510 eprintln!("parent worktree watcher error: {error}");
1511 }
1512 Err(mpsc::RecvTimeoutError::Timeout) => {}
1513 Err(mpsc::RecvTimeoutError::Disconnected) => {
1514 return Err("Parent filesystem watcher disconnected.".to_string());
1515 }
1516 }
1517
1518 if last_commit_check.elapsed() >= Duration::from_secs(3) {
1519 for repo in &mut repos {
1520 repo.last_head = record_commit_snapshot(
1521 &repo.identity,
1522 &repo.layout,
1523 repo.last_head.as_deref(),
1524 )?;
1525 }
1526 last_commit_check = Instant::now();
1527 }
1528
1529 if sync_interval_seconds > 0
1530 && last_sync.elapsed() >= Duration::from_secs(sync_interval_seconds)
1531 {
1532 for repo in &repos {
1533 if let Err(error) = sync_spool_for_identity(&repo.identity, &repo.layout).await {
1534 eprintln!(
1535 "worktree mutation sync failed for {}: {error}",
1536 repo.identity.root.display()
1537 );
1538 }
1539 }
1540 last_sync = Instant::now();
1541 }
1542 }
1543}
1544
1545fn prepare_watched_repos(identities: Vec<RepoIdentity>) -> Result<Vec<WatchedRepo>, String> {
1546 let mut repos = Vec::with_capacity(identities.len());
1547 for identity in identities {
1548 let layout = prepare_spool_layout(&identity)?;
1549 let event_dedupe = RecentEventDedupe {
1550 fingerprints: load_event_dedupe_fingerprints(&layout.event_dedupe_file)?,
1551 };
1552 repos.push(WatchedRepo {
1553 last_head: identity.head_sha.clone(),
1554 identity,
1555 layout,
1556 event_dedupe,
1557 });
1558 }
1559 repos.sort_by_key(|repo| std::cmp::Reverse(repo.identity.root.components().count()));
1560 Ok(repos)
1561}
1562
1563fn route_parent_event(repos: &[WatchedRepo], event: &Event) -> Vec<(usize, Event)> {
1564 let mut paths_by_repo: BTreeMap<usize, Vec<PathBuf>> = BTreeMap::new();
1565 for path in &event.paths {
1566 for (index, repo) in repos.iter().enumerate() {
1567 if path_belongs_to_repo(path, &repo.identity.root) {
1568 let paths = paths_by_repo.entry(index).or_default();
1569 if !paths.iter().any(|existing| existing == path) {
1570 paths.push(path.clone());
1571 }
1572 break;
1573 }
1574 }
1575 }
1576
1577 let mut routed = Vec::new();
1578 for (index, paths) in paths_by_repo {
1579 let mut repo_event = event.clone();
1580 repo_event.paths = paths;
1581 routed.push((index, repo_event));
1582 }
1583 routed
1584}
1585
1586fn path_belongs_to_repo(path: &Path, repo_root: &Path) -> bool {
1587 let path_key = path_identity_key(path).replace('\\', "/");
1588 let root_key = path_identity_key(repo_root).replace('\\', "/");
1589 path_key == root_key || path_key.starts_with(&format!("{root_key}/"))
1590}
1591
1592fn mutation_event_from_notify(
1593 identity: &RepoIdentity,
1594 event: Event,
1595) -> Option<WorktreeMutationEvent> {
1596 if event.kind.is_access() || event.kind.is_other() {
1597 return None;
1598 }
1599
1600 let paths: Vec<String> = event
1601 .paths
1602 .iter()
1603 .filter(|path| !is_ignored_path(&identity.root, path))
1604 .filter_map(|path| relative_slash_path(&identity.root, path))
1605 .fold(Vec::new(), |mut paths, path| {
1606 if !paths.iter().any(|existing| existing == &path) {
1607 paths.push(path);
1608 }
1609 paths
1610 });
1611
1612 if paths.is_empty() {
1613 return None;
1614 }
1615
1616 let primary_path = paths.first().cloned();
1617 let (old_path, new_path) = rename_paths(&event.kind, &paths);
1618 let line_counts = primary_path
1619 .as_deref()
1620 .and_then(|path| git_numstat_for_path(&identity.root, path).ok());
1621 let total_lines = primary_path
1622 .as_deref()
1623 .and_then(|path| file_line_count_for_path(&identity.root, path).ok());
1624 let event_kind = classify_event_kind(&event.kind);
1625 let is_dir = event
1626 .paths
1627 .first()
1628 .and_then(|path| fs::metadata(path).ok())
1629 .map(|metadata| metadata.is_dir())
1630 .unwrap_or(false);
1631
1632 Some(WorktreeMutationEvent {
1633 id: Uuid::new_v4().to_string(),
1634 fingerprint: None,
1635 repo_owner: identity.owner.clone(),
1636 repo_name: identity.name.clone(),
1637 branch_name: identity.branch.clone(),
1638 repo_root: identity.root.display().to_string(),
1639 head_sha: current_head(&identity.root).ok().flatten(),
1640 event_kind: event_kind.to_string(),
1641 paths,
1642 primary_path,
1643 old_path,
1644 new_path,
1645 added_lines: line_counts.map(|counts| counts.0),
1646 removed_lines: line_counts.map(|counts| counts.1),
1647 total_lines,
1648 file_created: matches!(
1649 event.kind,
1650 EventKind::Create(CreateKind::File | CreateKind::Any)
1651 ) && !is_dir,
1652 file_removed: matches!(
1653 event.kind,
1654 EventKind::Remove(RemoveKind::File | RemoveKind::Any)
1655 ) && !is_dir,
1656 folder_created: matches!(
1657 event.kind,
1658 EventKind::Create(CreateKind::Folder | CreateKind::Any)
1659 ) && is_dir,
1660 folder_removed: matches!(
1661 event.kind,
1662 EventKind::Remove(RemoveKind::Folder | RemoveKind::Any)
1663 ) && is_dir,
1664 renamed_or_moved: is_rename_or_move(&event.kind),
1665 raw_kind: format!("{:?}", event.kind),
1666 occurred_at: Utc::now(),
1667 })
1668}
1669
1670fn append_unique_mutation_event(
1671 layout: &SpoolLayout,
1672 mutation: &WorktreeMutationEvent,
1673 dedupe: &mut RecentEventDedupe,
1674) -> Result<bool, String> {
1675 let fingerprint = mutation_event_fingerprint(mutation);
1676 if dedupe.fingerprints.contains(&fingerprint)
1677 || event_dedupe_file_contains(&layout.event_dedupe_file, &fingerprint)?
1678 {
1679 dedupe.fingerprints.insert(fingerprint);
1680 return Ok(false);
1681 }
1682
1683 if !claim_event_fingerprint(layout, &fingerprint)? {
1684 dedupe.fingerprints.insert(fingerprint);
1685 return Ok(false);
1686 }
1687
1688 let mut mutation = mutation.clone();
1689 mutation.fingerprint = Some(fingerprint.clone());
1690 if let Err(error) = append_json_line(&layout.events_file, &mutation) {
1691 release_event_fingerprint_claim(layout, &fingerprint);
1692 return Err(error);
1693 }
1694 append_json_line(
1695 &layout.event_dedupe_file,
1696 &WorktreeWatchEventDedupeRecord {
1697 fingerprint: fingerprint.clone(),
1698 first_seen_at: Utc::now(),
1699 event_kind: mutation.event_kind.clone(),
1700 primary_path: mutation.primary_path.clone(),
1701 },
1702 )?;
1703 dedupe.fingerprints.insert(fingerprint);
1704 Ok(true)
1705}
1706
1707fn claim_event_fingerprint(layout: &SpoolLayout, fingerprint: &str) -> Result<bool, String> {
1708 fs::create_dir_all(&layout.event_dedupe_dir).map_err(|error| {
1709 format!(
1710 "Failed to create event dedupe directory {}: {error}",
1711 layout.event_dedupe_dir.display()
1712 )
1713 })?;
1714 let marker = event_dedupe_marker_path(layout, fingerprint);
1715 match OpenOptions::new()
1716 .write(true)
1717 .create_new(true)
1718 .open(&marker)
1719 {
1720 Ok(mut file) => {
1721 writeln!(file, "{}", Utc::now().to_rfc3339()).map_err(|error| {
1722 format!(
1723 "Failed to write event dedupe marker {}: {error}",
1724 marker.display()
1725 )
1726 })?;
1727 Ok(true)
1728 }
1729 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
1730 Err(error) => Err(format!(
1731 "Failed to create event dedupe marker {}: {error}",
1732 marker.display()
1733 )),
1734 }
1735}
1736
1737fn release_event_fingerprint_claim(layout: &SpoolLayout, fingerprint: &str) {
1738 let _ = fs::remove_file(event_dedupe_marker_path(layout, fingerprint));
1739}
1740
1741fn event_dedupe_marker_path(layout: &SpoolLayout, fingerprint: &str) -> PathBuf {
1742 layout.event_dedupe_dir.join(format!("{fingerprint}.seen"))
1743}
1744
1745fn mutation_event_fingerprint(mutation: &WorktreeMutationEvent) -> String {
1746 let bucket = mutation.occurred_at.timestamp() / EVENT_DEDUPE_WINDOW_SECONDS;
1747 let payload = json!({
1748 "repoOwner": mutation.repo_owner,
1749 "repoName": mutation.repo_name,
1750 "branchName": mutation.branch_name,
1751 "repoRoot": mutation.repo_root,
1752 "headSha": mutation.head_sha,
1753 "eventKind": mutation.event_kind,
1754 "paths": mutation.paths,
1755 "primaryPath": mutation.primary_path,
1756 "oldPath": mutation.old_path,
1757 "newPath": mutation.new_path,
1758 "addedLines": mutation.added_lines,
1759 "removedLines": mutation.removed_lines,
1760 "fileCreated": mutation.file_created,
1761 "fileRemoved": mutation.file_removed,
1762 "folderCreated": mutation.folder_created,
1763 "folderRemoved": mutation.folder_removed,
1764 "renamedOrMoved": mutation.renamed_or_moved,
1765 "timeBucket": bucket,
1766 });
1767 let encoded = serde_json::to_vec(&payload).unwrap_or_default();
1768 let mut hasher = Sha256::new();
1769 hasher.update(encoded);
1770 format!("{:x}", hasher.finalize())
1771}
1772
1773fn load_event_dedupe_fingerprints(path: &Path) -> Result<BTreeSet<String>, String> {
1774 let mut fingerprints = BTreeSet::new();
1775 if !path.exists() {
1776 return Ok(fingerprints);
1777 }
1778
1779 for value in read_jsonl_values(path)? {
1780 if let Some(fingerprint) = value.get("fingerprint").and_then(Value::as_str) {
1781 fingerprints.insert(fingerprint.to_string());
1782 }
1783 }
1784 Ok(fingerprints)
1785}
1786
1787fn event_dedupe_file_contains(path: &Path, fingerprint: &str) -> Result<bool, String> {
1788 if !path.exists() {
1789 return Ok(false);
1790 }
1791 for value in read_jsonl_values(path)? {
1792 if value
1793 .get("fingerprint")
1794 .and_then(Value::as_str)
1795 .map(|value| value == fingerprint)
1796 .unwrap_or(false)
1797 {
1798 return Ok(true);
1799 }
1800 }
1801 Ok(false)
1802}
1803
1804fn classify_event_kind(kind: &EventKind) -> &'static str {
1805 match kind {
1806 EventKind::Create(CreateKind::File) => "file_create",
1807 EventKind::Create(CreateKind::Folder) => "folder_create",
1808 EventKind::Create(_) => "create",
1809 EventKind::Remove(RemoveKind::File) => "file_remove",
1810 EventKind::Remove(RemoveKind::Folder) => "folder_remove",
1811 EventKind::Remove(_) => "remove",
1812 EventKind::Modify(ModifyKind::Name(
1813 RenameMode::From | RenameMode::To | RenameMode::Both,
1814 )) => "rename_or_move",
1815 EventKind::Modify(ModifyKind::Data(_)) => "file_modify",
1816 EventKind::Modify(ModifyKind::Metadata(_)) => "metadata_modify",
1817 EventKind::Modify(_) => "modify",
1818 _ => "other",
1819 }
1820}
1821
1822fn is_rename_or_move(kind: &EventKind) -> bool {
1823 matches!(
1824 kind,
1825 EventKind::Modify(ModifyKind::Name(
1826 RenameMode::From | RenameMode::To | RenameMode::Both
1827 ))
1828 )
1829}
1830
1831fn rename_paths(kind: &EventKind, paths: &[String]) -> (Option<String>, Option<String>) {
1832 if !is_rename_or_move(kind) {
1833 return (None, None);
1834 }
1835
1836 match paths {
1837 [old_path, new_path, ..] => (Some(old_path.clone()), Some(new_path.clone())),
1838 [path] => match kind {
1839 EventKind::Modify(ModifyKind::Name(RenameMode::From)) => (Some(path.clone()), None),
1840 EventKind::Modify(ModifyKind::Name(RenameMode::To)) => (None, Some(path.clone())),
1841 _ => (None, Some(path.clone())),
1842 },
1843 _ => (None, None),
1844 }
1845}
1846
1847fn record_commit_snapshot(
1848 identity: &RepoIdentity,
1849 layout: &SpoolLayout,
1850 previous_head: Option<&str>,
1851) -> Result<Option<String>, String> {
1852 let head = current_head(&identity.root)?;
1853 let Some(head_sha) = head else {
1854 return Ok(None);
1855 };
1856
1857 if previous_head == Some(head_sha.as_str()) {
1858 return Ok(Some(head_sha));
1859 }
1860
1861 let commit = WorktreeCommitEvent {
1862 id: Uuid::new_v4().to_string(),
1863 fingerprint: None,
1864 repo_owner: identity.owner.clone(),
1865 repo_name: identity.name.clone(),
1866 branch_name: identity.branch.clone(),
1867 repo_root: identity.root.display().to_string(),
1868 previous_head_sha: previous_head.map(str::to_string),
1869 head_sha: head_sha.clone(),
1870 subject: git_output(&identity.root, &["log", "-1", "--pretty=%s"]).ok(),
1871 author_name: git_output(&identity.root, &["log", "-1", "--pretty=%an"]).ok(),
1872 author_email: git_output(&identity.root, &["log", "-1", "--pretty=%ae"]).ok(),
1873 committed_at: git_output(&identity.root, &["log", "-1", "--pretty=%cI"]).ok(),
1874 occurred_at: Utc::now(),
1875 };
1876 append_unique_commit_event(layout, &commit)?;
1877 Ok(Some(head_sha))
1878}
1879
1880fn append_unique_commit_event(
1881 layout: &SpoolLayout,
1882 commit: &WorktreeCommitEvent,
1883) -> Result<bool, String> {
1884 let fingerprint = commit_event_fingerprint(commit);
1885 if !claim_commit_fingerprint(layout, &fingerprint)? {
1886 return Ok(false);
1887 }
1888
1889 let mut commit = commit.clone();
1890 commit.fingerprint = Some(fingerprint.clone());
1891 if let Err(error) = append_json_line(&layout.commits_file, &commit) {
1892 release_commit_fingerprint_claim(layout, &fingerprint);
1893 return Err(error);
1894 }
1895 Ok(true)
1896}
1897
1898fn commit_event_fingerprint(commit: &WorktreeCommitEvent) -> String {
1899 let payload = json!({
1900 "repoOwner": commit.repo_owner,
1901 "repoName": commit.repo_name,
1902 "branchName": commit.branch_name,
1903 "repoRoot": commit.repo_root,
1904 "previousHeadSha": commit.previous_head_sha,
1905 "headSha": commit.head_sha,
1906 });
1907 let encoded = serde_json::to_vec(&payload).unwrap_or_default();
1908 let mut hasher = Sha256::new();
1909 hasher.update(encoded);
1910 format!("{:x}", hasher.finalize())
1911}
1912
1913fn claim_commit_fingerprint(layout: &SpoolLayout, fingerprint: &str) -> Result<bool, String> {
1914 fs::create_dir_all(&layout.commit_dedupe_dir).map_err(|error| {
1915 format!(
1916 "Failed to create commit dedupe directory {}: {error}",
1917 layout.commit_dedupe_dir.display()
1918 )
1919 })?;
1920 let marker = commit_dedupe_marker_path(layout, fingerprint);
1921 match OpenOptions::new()
1922 .write(true)
1923 .create_new(true)
1924 .open(&marker)
1925 {
1926 Ok(mut file) => {
1927 writeln!(file, "{}", Utc::now().to_rfc3339()).map_err(|error| {
1928 format!(
1929 "Failed to write commit dedupe marker {}: {error}",
1930 marker.display()
1931 )
1932 })?;
1933 Ok(true)
1934 }
1935 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
1936 Err(error) => Err(format!(
1937 "Failed to create commit dedupe marker {}: {error}",
1938 marker.display()
1939 )),
1940 }
1941}
1942
1943fn release_commit_fingerprint_claim(layout: &SpoolLayout, fingerprint: &str) {
1944 let _ = fs::remove_file(commit_dedupe_marker_path(layout, fingerprint));
1945}
1946
1947fn commit_dedupe_marker_path(layout: &SpoolLayout, fingerprint: &str) -> PathBuf {
1948 layout.commit_dedupe_dir.join(format!("{fingerprint}.seen"))
1949}
1950
1951async fn sync_spool_for_identity(
1952 identity: &RepoIdentity,
1953 layout: &SpoolLayout,
1954) -> Result<(), String> {
1955 let candidates = collect_sync_candidates(layout)?;
1956 if candidates.is_empty() {
1957 return Ok(());
1958 }
1959
1960 sync_candidates(identity, candidates, false).await
1961}
1962
1963async fn sync_candidates(
1964 identity: &RepoIdentity,
1965 candidates: Vec<SyncCandidate>,
1966 resync: bool,
1967) -> Result<(), String> {
1968 let mut events = Vec::new();
1969 let mut commits = Vec::new();
1970 for candidate in &candidates {
1971 match &candidate.records {
1972 SyncRecords::Events(records) => events.extend(records.iter().cloned()),
1973 SyncRecords::Commits(records) => commits.extend(records.iter().cloned()),
1974 }
1975 }
1976
1977 let token = resolve_cli_access_token()?;
1978 let device = resolve_device_identity()?;
1979 let client = cli_request_client()?;
1980 let api = ApiConfig::from_env();
1981 let endpoint = api.cli_worktree_mutations_endpoint();
1982 let event_count = events.len();
1983 let commit_count = commits.len();
1984 let hardware_id = device.hardware_id;
1985 let hostname = current_hostname();
1986 let platform = std::env::consts::OS.to_string();
1987 let repo_root = identity.root.display().to_string();
1988 let batches = split_worktree_mutation_upload_batches(events.clone(), commits.clone());
1989 let batch_count = batches.len();
1990
1991 if batch_count > 1 {
1992 println!(
1993 "Uploading {} worktree mutation record(s) in {} batch(es) for {}/{} on {}.",
1994 event_count + commit_count,
1995 batch_count,
1996 identity.owner,
1997 identity.name,
1998 identity.branch
1999 );
2000 }
2001
2002 let mut status_code = 202;
2003 for attempt in 1..=WORKTREE_MUTATION_SYNC_MAX_ATTEMPTS {
2004 let attempt_batches =
2005 split_worktree_mutation_upload_batches(events.clone(), commits.clone());
2006 match upload_worktree_mutation_batches(
2007 &client,
2008 &endpoint,
2009 &token,
2010 &hardware_id,
2011 hostname.as_deref(),
2012 &platform,
2013 identity,
2014 &repo_root,
2015 attempt_batches,
2016 )
2017 .await
2018 {
2019 Ok(code) => {
2020 status_code = code;
2021 break;
2022 }
2023 Err(error)
2024 if error.kind == WorktreeMutationSyncErrorKind::Retryable
2025 && attempt < WORKTREE_MUTATION_SYNC_MAX_ATTEMPTS =>
2026 {
2027 eprintln!(
2028 "Worktree mutation sync attempt {attempt}/{} failed for {}/{} on {}: {} Restarting from the first batch in {} second(s).",
2029 WORKTREE_MUTATION_SYNC_MAX_ATTEMPTS,
2030 identity.owner,
2031 identity.name,
2032 identity.branch,
2033 error.message,
2034 WORKTREE_MUTATION_SYNC_RETRY_DELAY_SECONDS
2035 );
2036 tokio::time::sleep(Duration::from_secs(
2037 WORKTREE_MUTATION_SYNC_RETRY_DELAY_SECONDS,
2038 ))
2039 .await;
2040 }
2041 Err(error) => return Err(error.message),
2042 }
2043 }
2044
2045 let layout = prepare_spool_layout(identity)?;
2046 append_sync_log(SyncLogAppend {
2047 identity,
2048 layout: &layout,
2049 endpoint: &endpoint,
2050 status_code,
2051 event_count,
2052 commit_count,
2053 candidates: &candidates,
2054 resync,
2055 })?;
2056
2057 for candidate in candidates {
2058 if !is_synced_spool_path(&candidate.path) {
2059 mark_synced(&candidate.path)?;
2060 }
2061 }
2062
2063 Ok(())
2064}
2065
2066async fn upload_worktree_mutation_batches(
2067 client: &reqwest::Client,
2068 endpoint: &str,
2069 token: &str,
2070 hardware_id: &str,
2071 hostname: Option<&str>,
2072 platform: &str,
2073 identity: &RepoIdentity,
2074 repo_root: &str,
2075 batches: Vec<WorktreeMutationUploadBatch>,
2076) -> Result<u16, WorktreeMutationSyncError> {
2077 let batch_count = batches.len();
2078 let mut status_code = 202;
2079
2080 for (batch_index, batch) in batches.into_iter().enumerate() {
2081 let batch_event_count = batch.events.len();
2082 let batch_commit_count = batch.commits.len();
2083 let payload = WorktreeMutationIngestPayload {
2084 device: WorktreeMutationDevicePayload {
2085 hardware_id: hardware_id.to_string(),
2086 device_name: hostname.map(str::to_string),
2087 hostname: hostname.map(str::to_string),
2088 platform: platform.to_string(),
2089 },
2090 repository: WorktreeMutationRepositoryPayload {
2091 owner: identity.owner.clone(),
2092 name: identity.name.clone(),
2093 branch_name: identity.branch.clone(),
2094 repo_root: repo_root.to_string(),
2095 },
2096 events: batch.events,
2097 commits: batch.commits,
2098 };
2099
2100 let response = client
2101 .post(endpoint)
2102 .bearer_auth(token)
2103 .json(&payload)
2104 .send()
2105 .await
2106 .map_err(|error| WorktreeMutationSyncError {
2107 kind: WorktreeMutationSyncErrorKind::Retryable,
2108 message: format!(
2109 "Failed to upload worktree mutation batch {}/{}: {error}",
2110 batch_index + 1,
2111 batch_count
2112 ),
2113 })?;
2114
2115 if response.status() == StatusCode::UNAUTHORIZED {
2116 return Err(WorktreeMutationSyncError {
2117 kind: WorktreeMutationSyncErrorKind::NonRetryable,
2118 message: "Your stored CLI session is no longer valid. Run `xbp login` again."
2119 .to_string(),
2120 });
2121 }
2122
2123 if !response.status().is_success() {
2124 let status = response.status();
2125 let body = response.text().await.unwrap_or_default();
2126 return Err(WorktreeMutationSyncError {
2127 kind: if should_retry_worktree_mutation_sync_status(status) {
2128 WorktreeMutationSyncErrorKind::Retryable
2129 } else {
2130 WorktreeMutationSyncErrorKind::NonRetryable
2131 },
2132 message: format!(
2133 "Worktree mutation upload batch {}/{} failed with {status}: {body}",
2134 batch_index + 1,
2135 batch_count
2136 ),
2137 });
2138 }
2139
2140 status_code = response.status().as_u16();
2141 if batch_count > 1 {
2142 println!(
2143 "Uploaded worktree mutation batch {}/{} ({} event(s), {} commit(s)).",
2144 batch_index + 1,
2145 batch_count,
2146 batch_event_count,
2147 batch_commit_count
2148 );
2149 }
2150 }
2151
2152 Ok(status_code)
2153}
2154
2155fn should_retry_worktree_mutation_sync_status(status: StatusCode) -> bool {
2156 status.is_server_error()
2157 || status == StatusCode::REQUEST_TIMEOUT
2158 || status == StatusCode::TOO_MANY_REQUESTS
2159}
2160
2161fn split_worktree_mutation_upload_batches(
2162 events: Vec<WorktreeMutationEvent>,
2163 commits: Vec<WorktreeCommitEvent>,
2164) -> Vec<WorktreeMutationUploadBatch> {
2165 let mut batches = Vec::new();
2166 let mut current = empty_worktree_mutation_upload_batch();
2167
2168 for event in events {
2169 let estimated_json_bytes = serialized_json_len(&event);
2170 if should_start_next_worktree_mutation_upload_batch(¤t, estimated_json_bytes) {
2171 batches.push(current);
2172 current = empty_worktree_mutation_upload_batch();
2173 }
2174 current.estimated_json_bytes += estimated_json_bytes;
2175 current.events.push(event);
2176 }
2177
2178 for commit in commits {
2179 let estimated_json_bytes = serialized_json_len(&commit);
2180 if should_start_next_worktree_mutation_upload_batch(¤t, estimated_json_bytes) {
2181 batches.push(current);
2182 current = empty_worktree_mutation_upload_batch();
2183 }
2184 current.estimated_json_bytes += estimated_json_bytes;
2185 current.commits.push(commit);
2186 }
2187
2188 if current_record_count(¤t) > 0 {
2189 batches.push(current);
2190 }
2191
2192 batches
2193}
2194
2195fn empty_worktree_mutation_upload_batch() -> WorktreeMutationUploadBatch {
2196 WorktreeMutationUploadBatch {
2197 events: Vec::new(),
2198 commits: Vec::new(),
2199 estimated_json_bytes: 0,
2200 }
2201}
2202
2203fn should_start_next_worktree_mutation_upload_batch(
2204 batch: &WorktreeMutationUploadBatch,
2205 next_record_json_bytes: usize,
2206) -> bool {
2207 if current_record_count(batch) == 0 {
2208 return false;
2209 }
2210
2211 current_record_count(batch) >= WORKTREE_MUTATION_SYNC_BATCH_RECORD_LIMIT
2212 || batch.estimated_json_bytes + next_record_json_bytes
2213 > WORKTREE_MUTATION_SYNC_BATCH_JSON_BYTES
2214}
2215
2216fn current_record_count(batch: &WorktreeMutationUploadBatch) -> usize {
2217 batch.events.len() + batch.commits.len()
2218}
2219
2220fn serialized_json_len<T: Serialize>(value: &T) -> usize {
2221 serde_json::to_vec(value)
2222 .map(|encoded| encoded.len())
2223 .unwrap_or(0)
2224}
2225
2226fn collect_sync_candidates(layout: &SpoolLayout) -> Result<Vec<SyncCandidate>, String> {
2227 collect_spool_candidates(layout, false)
2228}
2229
2230fn print_no_sync_candidates_message(
2231 target: &WorktreeWatchTargetOptions,
2232 resync: bool,
2233) -> Result<(), String> {
2234 if resync {
2235 println!("No local worktree mutation JSONL files found.");
2236 return Ok(());
2237 }
2238
2239 let identities = resolve_target_identities(target)?;
2240 let mut synced_files = 0usize;
2241 let mut synced_records = 0usize;
2242 for identity in identities {
2243 let layout = prepare_spool_layout(&identity)?;
2244 let all_candidates = collect_spool_candidates(&layout, true)?;
2245 let unsynced_candidates = collect_sync_candidates(&layout)?;
2246 synced_files += all_candidates
2247 .len()
2248 .saturating_sub(unsynced_candidates.len());
2249 let all_records = all_candidates
2250 .iter()
2251 .map(|candidate| candidate.records.len())
2252 .sum::<usize>();
2253 let unsynced_records = unsynced_candidates
2254 .iter()
2255 .map(|candidate| candidate.records.len())
2256 .sum::<usize>();
2257 synced_records += all_records.saturating_sub(unsynced_records);
2258 }
2259
2260 if synced_files > 0 {
2261 println!(
2262 "No unsynced worktree mutation files found. Found {synced_files} already-synced local spool file(s) with {synced_records} record(s). Use `xbp worktree-watch sync --resync` to upload them again."
2263 );
2264 } else {
2265 println!("No unsynced worktree mutation files found.");
2266 }
2267 Ok(())
2268}
2269
2270fn collect_spool_candidates(
2271 layout: &SpoolLayout,
2272 include_synced: bool,
2273) -> Result<Vec<SyncCandidate>, String> {
2274 let mut candidates = Vec::new();
2275 if !layout.root.exists() {
2276 return Ok(candidates);
2277 }
2278
2279 for entry in fs::read_dir(&layout.root).map_err(|error| {
2280 format!(
2281 "Failed to read spool directory {}: {error}",
2282 layout.root.display()
2283 )
2284 })? {
2285 let path = entry
2286 .map_err(|error| format!("Failed to read spool entry: {error}"))?
2287 .path();
2288 let Some(file_name) = path.file_name().and_then(|value| value.to_str()) else {
2289 continue;
2290 };
2291 if !file_name.ends_with(".jsonl") || (!include_synced && file_name.contains(".synced.")) {
2292 continue;
2293 }
2294
2295 let kind = if file_name.starts_with("events-") {
2296 SyncFileKind::Events
2297 } else if file_name.starts_with("commits-") {
2298 SyncFileKind::Commits
2299 } else {
2300 continue;
2301 };
2302 let records = sanitize_spool_records(&path, kind, read_spool_records(&path, kind)?)?;
2303 if records.is_empty() {
2304 continue;
2305 }
2306 candidates.push(SyncCandidate { path, records });
2307 }
2308
2309 Ok(candidates)
2310}
2311
2312fn append_sync_log(request: SyncLogAppend<'_>) -> Result<(), String> {
2313 let entry = WorktreeWatchSyncLogEntry {
2314 id: Uuid::new_v4().to_string(),
2315 synced_at: Utc::now(),
2316 endpoint: request.endpoint.to_string(),
2317 repository_owner: request.identity.owner.clone(),
2318 repository_name: request.identity.name.clone(),
2319 branch_name: request.identity.branch.clone(),
2320 repo_root: request.identity.root.display().to_string(),
2321 resync: request.resync,
2322 status_code: request.status_code,
2323 event_count: request.event_count,
2324 commit_count: request.commit_count,
2325 spool_file_count: request.candidates.len(),
2326 spool_files: request
2327 .candidates
2328 .iter()
2329 .map(|candidate| candidate.path.display().to_string())
2330 .collect(),
2331 };
2332 append_json_line(&request.layout.sync_log_file, &entry)
2333}
2334
2335fn sanitize_spool_records(
2336 path: &Path,
2337 kind: SyncFileKind,
2338 records: SyncRecords,
2339) -> Result<SyncRecords, String> {
2340 match (kind, records) {
2341 (SyncFileKind::Events, SyncRecords::Events(records)) => {
2342 let mut sanitized = Vec::with_capacity(records.len());
2343 let mut changed = false;
2344
2345 for record in records {
2346 let original = serde_json::to_value(&record).ok();
2347 let Some(record) = sanitize_spooled_event(record) else {
2348 changed = true;
2349 continue;
2350 };
2351 if !changed {
2352 changed = serde_json::to_value(&record).ok() != original;
2353 }
2354 sanitized.push(record);
2355 }
2356
2357 if changed {
2358 rewrite_jsonl_records(path, &sanitized)?;
2359 }
2360
2361 Ok(SyncRecords::Events(sanitized))
2362 }
2363 (_, records) => Ok(records),
2364 }
2365}
2366
2367fn sanitize_spooled_event(mut event: WorktreeMutationEvent) -> Option<WorktreeMutationEvent> {
2368 event.paths = dedupe_filtered_spool_paths(event.paths);
2369 event.primary_path = sanitize_optional_spool_path(event.primary_path);
2370 event.old_path = sanitize_optional_spool_path(event.old_path);
2371 event.new_path = sanitize_optional_spool_path(event.new_path);
2372
2373 if event.paths.is_empty() {
2374 if let Some(path) = event.primary_path.clone() {
2375 event.paths.push(path);
2376 }
2377 if let Some(path) = event.old_path.clone() {
2378 if !event.paths.iter().any(|existing| existing == &path) {
2379 event.paths.push(path);
2380 }
2381 }
2382 if let Some(path) = event.new_path.clone() {
2383 if !event.paths.iter().any(|existing| existing == &path) {
2384 event.paths.push(path);
2385 }
2386 }
2387 }
2388
2389 if event.primary_path.is_none() {
2390 event.primary_path = event.paths.first().cloned();
2391 }
2392
2393 if event.paths.is_empty() {
2394 return None;
2395 }
2396
2397 Some(event)
2398}
2399
2400fn sanitize_optional_spool_path(path: Option<String>) -> Option<String> {
2401 path.and_then(sanitize_spool_path)
2402}
2403
2404fn sanitize_spool_path(path: String) -> Option<String> {
2405 let trimmed = path.trim();
2406 if trimmed.is_empty() || watch_ignore_rules().ignores_relative(trimmed) {
2407 return None;
2408 }
2409 Some(trimmed.replace('\\', "/"))
2410}
2411
2412fn dedupe_filtered_spool_paths(paths: Vec<String>) -> Vec<String> {
2413 let mut sanitized = Vec::with_capacity(paths.len());
2414 for path in paths {
2415 let Some(path) = sanitize_spool_path(path) else {
2416 continue;
2417 };
2418 if !sanitized.iter().any(|existing| existing == &path) {
2419 sanitized.push(path);
2420 }
2421 }
2422 sanitized
2423}
2424
2425fn rewrite_jsonl_records<T: Serialize>(path: &Path, records: &[T]) -> Result<(), String> {
2426 if records.is_empty() {
2427 if path.exists() {
2428 fs::remove_file(path).map_err(|error| {
2429 format!(
2430 "Failed to remove emptied spool file {}: {error}",
2431 path.display()
2432 )
2433 })?;
2434 }
2435 return Ok(());
2436 }
2437
2438 let temp_path = path.with_extension(format!(
2439 "{}.{}.tmp",
2440 path.extension()
2441 .and_then(|value| value.to_str())
2442 .unwrap_or("jsonl"),
2443 Uuid::new_v4()
2444 ));
2445 {
2446 let mut file = File::create(&temp_path).map_err(|error| {
2447 format!(
2448 "Failed to create temporary spool file {}: {error}",
2449 temp_path.display()
2450 )
2451 })?;
2452 for record in records {
2453 let line = serde_json::to_string(record)
2454 .map_err(|error| format!("Failed to serialize worktree event: {error}"))?;
2455 writeln!(file, "{line}").map_err(|error| {
2456 format!(
2457 "Failed to write temporary spool file {}: {error}",
2458 temp_path.display()
2459 )
2460 })?;
2461 }
2462 }
2463
2464 if path.exists() {
2465 fs::remove_file(path).map_err(|error| {
2466 format!(
2467 "Failed to replace rewritten spool file {}: {error}",
2468 path.display()
2469 )
2470 })?;
2471 }
2472
2473 fs::rename(&temp_path, path)
2474 .map_err(|error| format!("Failed to rewrite spool file {}: {error}", path.display()))
2475}
2476
2477fn read_last_sync_log_entry(path: &Path) -> Result<Option<WorktreeWatchSyncLogEntry>, String> {
2478 if !path.exists() {
2479 return Ok(None);
2480 }
2481 let values = read_jsonl_values(path)?;
2482 let Some(value) = values.into_iter().last() else {
2483 return Ok(None);
2484 };
2485 serde_json::from_value(value)
2486 .map(Some)
2487 .map_err(|error| format!("Failed to decode sync log {}: {error}", path.display()))
2488}
2489
2490fn is_synced_spool_path(path: &Path) -> bool {
2491 path.file_name()
2492 .and_then(|value| value.to_str())
2493 .map(|value| value.contains(".synced."))
2494 .unwrap_or(false)
2495}
2496
2497fn mark_synced(path: &Path) -> Result<(), String> {
2498 let Some(file_name) = path.file_name().and_then(|value| value.to_str()) else {
2499 return Ok(());
2500 };
2501 let synced_name = file_name.replace(
2502 ".jsonl",
2503 &format!(".synced.{}.jsonl", Utc::now().timestamp()),
2504 );
2505 let synced_path = path.with_file_name(synced_name);
2506 fs::rename(path, &synced_path).map_err(|error| {
2507 format!(
2508 "Failed to mark spool file {} as synced: {error}",
2509 path.display()
2510 )
2511 })
2512}
2513
2514fn read_spool_records(path: &Path, kind: SyncFileKind) -> Result<SyncRecords, String> {
2515 match kind {
2516 SyncFileKind::Events => read_jsonl_records::<WorktreeMutationEvent>(path)
2517 .map(SyncRecords::Events)
2518 .map_err(|error| format!("Failed to decode event from {}: {error}", path.display())),
2519 SyncFileKind::Commits => read_jsonl_records::<WorktreeCommitEvent>(path)
2520 .map(SyncRecords::Commits)
2521 .map_err(|error| format!("Failed to decode commit from {}: {error}", path.display())),
2522 }
2523}
2524
2525fn read_jsonl_records<T>(path: &Path) -> Result<Vec<T>, String>
2526where
2527 T: for<'de> Deserialize<'de>,
2528{
2529 let file = File::open(path)
2530 .map_err(|error| format!("Failed to open spool file {}: {error}", path.display()))?;
2531 let reader = BufReader::new(file);
2532 let mut records = Vec::new();
2533 for line in reader.lines() {
2534 let line =
2535 line.map_err(|error| format!("Failed to read spool file {}: {error}", path.display()))?;
2536 if line.trim().is_empty() {
2537 continue;
2538 }
2539 records.push(serde_json::from_str(&line).map_err(|error| {
2540 format!(
2541 "Failed to parse JSONL record in {}: {error}",
2542 path.display()
2543 )
2544 })?);
2545 }
2546 Ok(records)
2547}
2548
2549fn read_jsonl_values(path: &Path) -> Result<Vec<Value>, String> {
2550 read_jsonl_records(path)
2551}
2552
2553fn spawn_detached_worktree_watch(repo: Option<&Path>) -> Result<PathBuf, String> {
2554 let identity = resolve_repo_identity(repo)?;
2555 spawn_detached_worktree_watch_for_identity(&identity)
2556}
2557
2558fn spawn_detached_parent_worktree_watch(parent: &Path) -> Result<PathBuf, String> {
2559 let parent = canonical_parent_path(parent)?;
2560 let layout = prepare_parent_spool_layout(&parent)?;
2561 replace_existing_parent_watcher(&parent, &layout)?;
2562 let executable = std::env::current_exe()
2563 .map_err(|error| format!("Failed to resolve current XBP executable: {error}"))?;
2564 let mut command = Command::new(&executable);
2565 command
2566 .arg("worktree-watch")
2567 .arg("start")
2568 .arg("--parent")
2569 .arg(&parent)
2570 .env(BACKGROUND_CHILD_ENV, "1")
2571 .current_dir(&parent)
2572 .stdin(Stdio::null())
2573 .stdout(Stdio::null())
2574 .stderr(Stdio::null());
2575
2576 #[cfg(windows)]
2577 {
2578 use std::os::windows::process::CommandExt;
2579 command.creation_flags(CREATE_NO_WINDOW);
2580 }
2581
2582 let child = command
2583 .spawn()
2584 .map_err(|error| format!("Failed to spawn background parent worktree watcher: {error}"))?;
2585 write_parent_watcher_state(&parent, &layout, child.id(), &executable)?;
2586 Ok(parent)
2587}
2588
2589fn spawn_detached_worktree_watch_for_identity(identity: &RepoIdentity) -> Result<PathBuf, String> {
2590 let layout = prepare_spool_layout(identity)?;
2591 replace_existing_watcher(identity, &layout)?;
2592 let executable = std::env::current_exe()
2593 .map_err(|error| format!("Failed to resolve current XBP executable: {error}"))?;
2594 let mut command = Command::new(&executable);
2595 command
2596 .arg("worktree-watch")
2597 .arg("start")
2598 .arg("--repo")
2599 .arg(&identity.root)
2600 .env(BACKGROUND_CHILD_ENV, "1")
2601 .current_dir(&identity.root)
2602 .stdin(Stdio::null())
2603 .stdout(Stdio::null())
2604 .stderr(Stdio::null());
2605
2606 #[cfg(windows)]
2607 {
2608 use std::os::windows::process::CommandExt;
2609 command.creation_flags(CREATE_NO_WINDOW);
2610 }
2611
2612 let child: std::process::Child = command
2613 .spawn()
2614 .map_err(|error| format!("Failed to spawn background worktree watcher: {error}"))?;
2615 write_watcher_state(identity, &layout, child.id(), &executable)?;
2616 Ok(identity.root.clone())
2617}
2618
2619fn replace_existing_parent_watcher(
2620 parent: &Path,
2621 layout: &ParentSpoolLayout,
2622) -> Result<(), String> {
2623 let Some(state) = read_parent_watcher_state(&layout.watcher_state_file)? else {
2624 return Ok(());
2625 };
2626
2627 if path_identity_key(parent) != path_identity_key(Path::new(&state.parent_root)) {
2628 return Ok(());
2629 }
2630
2631 if is_matching_parent_watcher_process(&state) {
2632 stop_parent_process(&state, false)?;
2633 println!(
2634 "Replaced existing background parent worktree watcher {} for {}",
2635 state.pid,
2636 parent.display()
2637 );
2638 let _ = fs::remove_file(&layout.watcher_state_file);
2639 return Ok(());
2640 }
2641
2642 if parent_watcher_state_pid_exists_with_same_executable(&state) {
2643 return Err(format!(
2644 "Existing parent watcher state for {} points to running XBP process {}, but its command line could not be verified. Run `xbp worktree-watch stop --parent \"{}\" --force` before starting another detached parent watcher.",
2645 parent.display(),
2646 state.pid,
2647 parent.display()
2648 ));
2649 }
2650
2651 let _ = fs::remove_file(&layout.watcher_state_file);
2652 Ok(())
2653}
2654
2655fn replace_existing_watcher(identity: &RepoIdentity, layout: &SpoolLayout) -> Result<(), String> {
2656 let Some(state) = read_watcher_state(&layout.watcher_state_file)? else {
2657 return Ok(());
2658 };
2659
2660 if !same_repo_watcher_state(identity, &state) {
2661 return Ok(());
2662 }
2663
2664 if is_matching_watcher_process(&state) {
2665 stop_process(&state, false)?;
2666 println!(
2667 "Replaced existing background worktree watcher {} for {}",
2668 state.pid,
2669 identity.root.display()
2670 );
2671 let _ = fs::remove_file(&layout.watcher_state_file);
2672 return Ok(());
2673 }
2674
2675 if watcher_state_pid_exists_with_same_executable(&state) {
2676 return Err(format!(
2677 "Existing watcher state for {} points to running XBP process {}, but its command line could not be verified. Run `xbp worktree-watch stop --repo \"{}\" --force` before starting another detached watcher.",
2678 identity.root.display(),
2679 state.pid,
2680 identity.root.display()
2681 ));
2682 }
2683
2684 let _ = fs::remove_file(&layout.watcher_state_file);
2685 Ok(())
2686}
2687
2688fn stop_existing_watcher(
2689 identity: &RepoIdentity,
2690 layout: &SpoolLayout,
2691 force: bool,
2692) -> Result<StopWatcherOutcome, String> {
2693 let Some(state) = read_watcher_state(&layout.watcher_state_file)? else {
2694 return Ok(StopWatcherOutcome::NoState);
2695 };
2696
2697 if !same_repo_watcher_state(identity, &state) {
2698 return Ok(StopWatcherOutcome::NoState);
2699 }
2700
2701 if is_matching_watcher_process(&state) || force {
2702 let pid = state.pid;
2703 stop_process(&state, force)?;
2704 let _ = fs::remove_file(&layout.watcher_state_file);
2705 return Ok(StopWatcherOutcome::Stopped(pid));
2706 }
2707
2708 if watcher_state_pid_exists_with_same_executable(&state) {
2709 return Err(format!(
2710 "Watcher state for {} points to running XBP process {}, but its command line could not be verified. Re-run with `--force` to stop that PID.",
2711 identity.root.display(),
2712 state.pid
2713 ));
2714 }
2715
2716 let _ = fs::remove_file(&layout.watcher_state_file);
2717 Ok(StopWatcherOutcome::RemovedStaleState)
2718}
2719
2720fn stop_existing_parent_watcher(parent: &Path, force: bool) -> Result<StopWatcherOutcome, String> {
2721 let layout = prepare_parent_spool_layout(parent)?;
2722 let Some(state) = read_parent_watcher_state(&layout.watcher_state_file)? else {
2723 return Ok(StopWatcherOutcome::NoState);
2724 };
2725
2726 if path_identity_key(parent) != path_identity_key(Path::new(&state.parent_root)) {
2727 return Ok(StopWatcherOutcome::NoState);
2728 }
2729
2730 if is_matching_parent_watcher_process(&state) || force {
2731 let pid = state.pid;
2732 stop_parent_process(&state, force)?;
2733 let _ = fs::remove_file(&layout.watcher_state_file);
2734 return Ok(StopWatcherOutcome::Stopped(pid));
2735 }
2736
2737 if parent_watcher_state_pid_exists_with_same_executable(&state) {
2738 return Err(format!(
2739 "Parent watcher state for {} points to running XBP process {}, but its command line could not be verified. Re-run with `--force` to stop that PID.",
2740 parent.display(),
2741 state.pid
2742 ));
2743 }
2744
2745 let _ = fs::remove_file(&layout.watcher_state_file);
2746 Ok(StopWatcherOutcome::RemovedStaleState)
2747}
2748
2749fn read_watcher_state(path: &Path) -> Result<Option<WorktreeWatcherState>, String> {
2750 if !path.exists() {
2751 return Ok(None);
2752 }
2753 let raw = fs::read_to_string(path)
2754 .map_err(|error| format!("Failed to read watcher state {}: {error}", path.display()))?;
2755 serde_json::from_str(&raw)
2756 .map(Some)
2757 .map_err(|error| format!("Failed to parse watcher state {}: {error}", path.display()))
2758}
2759
2760fn read_parent_watcher_state(path: &Path) -> Result<Option<ParentWorktreeWatcherState>, String> {
2761 if !path.exists() {
2762 return Ok(None);
2763 }
2764 let raw = fs::read_to_string(path).map_err(|error| {
2765 format!(
2766 "Failed to read parent watcher state {}: {error}",
2767 path.display()
2768 )
2769 })?;
2770 serde_json::from_str(&raw).map(Some).map_err(|error| {
2771 format!(
2772 "Failed to parse parent watcher state {}: {error}",
2773 path.display()
2774 )
2775 })
2776}
2777
2778fn write_watcher_state(
2779 identity: &RepoIdentity,
2780 layout: &SpoolLayout,
2781 pid: u32,
2782 executable: &Path,
2783) -> Result<(), String> {
2784 let state = WorktreeWatcherState {
2785 pid,
2786 repo_root: identity.root.display().to_string(),
2787 repository_owner: identity.owner.clone(),
2788 repository_name: identity.name.clone(),
2789 branch_name: identity.branch.clone(),
2790 executable: executable.display().to_string(),
2791 started_at: Utc::now(),
2792 };
2793 let rendered = serde_json::to_string_pretty(&state)
2794 .map_err(|error| format!("Failed to serialize watcher state: {error}"))?;
2795 fs::write(&layout.watcher_state_file, format!("{rendered}\n")).map_err(|error| {
2796 format!(
2797 "Failed to write watcher state {}: {error}",
2798 layout.watcher_state_file.display()
2799 )
2800 })
2801}
2802
2803fn write_parent_watcher_state(
2804 parent: &Path,
2805 layout: &ParentSpoolLayout,
2806 pid: u32,
2807 executable: &Path,
2808) -> Result<(), String> {
2809 let state = ParentWorktreeWatcherState {
2810 pid,
2811 parent_root: parent.display().to_string(),
2812 executable: executable.display().to_string(),
2813 started_at: Utc::now(),
2814 };
2815 let rendered = serde_json::to_string_pretty(&state)
2816 .map_err(|error| format!("Failed to serialize parent watcher state: {error}"))?;
2817 fs::write(&layout.watcher_state_file, format!("{rendered}\n")).map_err(|error| {
2818 format!(
2819 "Failed to write parent watcher state {}: {error}",
2820 layout.watcher_state_file.display()
2821 )
2822 })
2823}
2824
2825fn same_repo_watcher_state(identity: &RepoIdentity, state: &WorktreeWatcherState) -> bool {
2826 path_identity_key(&identity.root) == path_identity_key(Path::new(&state.repo_root))
2827 && identity.owner == state.repository_owner
2828 && identity.name == state.repository_name
2829 && identity.branch == state.branch_name
2830}
2831
2832fn is_matching_watcher_process(state: &WorktreeWatcherState) -> bool {
2833 if state.pid == std::process::id() {
2834 return false;
2835 }
2836 let system = System::new_all();
2837 let Some(process) = system.process(Pid::from_u32(state.pid)) else {
2838 return false;
2839 };
2840 let command_line = process
2841 .cmd()
2842 .iter()
2843 .map(|part| part.to_string_lossy())
2844 .collect::<Vec<_>>()
2845 .join(" ");
2846 command_line.contains("worktree-watch")
2847 && command_line.contains("start")
2848 && command_line.contains(&state.repo_root)
2849 || watcher_state_process_identity_matches(process, state)
2850}
2851
2852fn is_matching_parent_watcher_process(state: &ParentWorktreeWatcherState) -> bool {
2853 if state.pid == std::process::id() {
2854 return false;
2855 }
2856 let system = System::new_all();
2857 let Some(process) = system.process(Pid::from_u32(state.pid)) else {
2858 return false;
2859 };
2860 let command_line = process
2861 .cmd()
2862 .iter()
2863 .map(|part| part.to_string_lossy())
2864 .collect::<Vec<_>>()
2865 .join(" ");
2866 (command_line.contains("worktree-watch")
2867 && command_line.contains("start")
2868 && command_line.contains("--parent")
2869 && command_line.contains(&state.parent_root))
2870 || parent_watcher_state_process_identity_matches(process, state)
2871}
2872
2873fn watcher_state_pid_exists_with_same_executable(state: &WorktreeWatcherState) -> bool {
2874 if state.pid == std::process::id() {
2875 return false;
2876 }
2877 let system = System::new_all();
2878 let Some(process) = system.process(Pid::from_u32(state.pid)) else {
2879 return false;
2880 };
2881 process_executable_matches_state(process, state)
2882}
2883
2884fn parent_watcher_state_pid_exists_with_same_executable(
2885 state: &ParentWorktreeWatcherState,
2886) -> bool {
2887 if state.pid == std::process::id() {
2888 return false;
2889 }
2890 let system = System::new_all();
2891 let Some(process) = system.process(Pid::from_u32(state.pid)) else {
2892 return false;
2893 };
2894 process_executable_matches_path(process, Path::new(&state.executable))
2895}
2896
2897fn watcher_state_process_identity_matches(
2898 process: &sysinfo::Process,
2899 state: &WorktreeWatcherState,
2900) -> bool {
2901 process_executable_matches_state(process, state)
2902 && process_start_time_matches_state(process, state)
2903}
2904
2905fn parent_watcher_state_process_identity_matches(
2906 process: &sysinfo::Process,
2907 state: &ParentWorktreeWatcherState,
2908) -> bool {
2909 process_executable_matches_path(process, Path::new(&state.executable))
2910 && parent_process_start_time_matches_state(process, state)
2911}
2912
2913fn process_executable_matches_state(
2914 process: &sysinfo::Process,
2915 state: &WorktreeWatcherState,
2916) -> bool {
2917 process_executable_matches_path(process, Path::new(&state.executable))
2918}
2919
2920fn process_executable_matches_path(process: &sysinfo::Process, executable: &Path) -> bool {
2921 let expected = path_identity_key(executable);
2922 process
2923 .exe()
2924 .map(|path| path_identity_key(path) == expected)
2925 .unwrap_or(false)
2926}
2927
2928fn process_start_time_matches_state(
2929 process: &sysinfo::Process,
2930 state: &WorktreeWatcherState,
2931) -> bool {
2932 let process_started = process.start_time() as i64;
2933 let state_started = state.started_at.timestamp();
2934 (process_started - state_started).abs() <= 10
2935}
2936
2937fn parent_process_start_time_matches_state(
2938 process: &sysinfo::Process,
2939 state: &ParentWorktreeWatcherState,
2940) -> bool {
2941 let process_started = process.start_time() as i64;
2942 let state_started = state.started_at.timestamp();
2943 (process_started - state_started).abs() <= 10
2944}
2945
2946fn stop_process(state: &WorktreeWatcherState, force: bool) -> Result<(), String> {
2947 let system = System::new_all();
2948 let Some(process) = system.process(Pid::from_u32(state.pid)) else {
2949 return Ok(());
2950 };
2951 if !force && !is_matching_watcher_process(state) {
2952 return Ok(());
2953 }
2954 if process.kill() {
2955 Ok(())
2956 } else {
2957 Err(format!(
2958 "Failed to stop existing watcher process {}",
2959 state.pid
2960 ))
2961 }
2962}
2963
2964fn stop_parent_process(state: &ParentWorktreeWatcherState, force: bool) -> Result<(), String> {
2965 let system = System::new_all();
2966 let Some(process) = system.process(Pid::from_u32(state.pid)) else {
2967 return Ok(());
2968 };
2969 if !force && !is_matching_parent_watcher_process(state) {
2970 return Ok(());
2971 }
2972 if process.kill() {
2973 Ok(())
2974 } else {
2975 Err(format!(
2976 "Failed to stop existing parent watcher process {}",
2977 state.pid
2978 ))
2979 }
2980}
2981
2982fn discover_repo_identities_under(parent: &Path) -> Result<Vec<RepoIdentity>, String> {
2983 let parent = canonical_parent_path(parent)?;
2984
2985 let mut identities = Vec::new();
2986 let mut seen = BTreeSet::new();
2987 discover_repo_identities_in_dir(&parent, &parent, &mut seen, &mut identities)?;
2988 identities.sort_by(|left, right| left.root.cmp(&right.root));
2989 Ok(identities)
2990}
2991
2992fn discover_repo_identities_in_dir(
2993 parent: &Path,
2994 dir: &Path,
2995 seen: &mut BTreeSet<String>,
2996 identities: &mut Vec<RepoIdentity>,
2997) -> Result<(), String> {
2998 if dir != parent && is_skipped_discovery_dir(dir) {
2999 return Ok(());
3000 }
3001
3002 if dir.join(".git").exists() {
3003 if let Ok(identity) = resolve_repo_identity(Some(dir)) {
3004 let key = path_identity_key(&identity.root);
3005 if seen.insert(key) {
3006 identities.push(identity);
3007 }
3008 return Ok(());
3009 }
3010 }
3011
3012 let entries = fs::read_dir(dir)
3013 .map_err(|error| format!("Failed to read folder {}: {error}", dir.display()))?;
3014 for entry in entries {
3015 let entry = entry.map_err(|error| {
3016 format!(
3017 "Failed to read folder entry under {}: {error}",
3018 dir.display()
3019 )
3020 })?;
3021 let file_type = entry
3022 .file_type()
3023 .map_err(|error| format!("Failed to inspect {}: {error}", entry.path().display()))?;
3024 if file_type.is_dir() {
3025 discover_repo_identities_in_dir(parent, &entry.path(), seen, identities)?;
3026 }
3027 }
3028
3029 Ok(())
3030}
3031
3032fn resolve_repo_identity(repo: Option<&Path>) -> Result<RepoIdentity, String> {
3033 let start = match repo {
3034 Some(path) => path.to_path_buf(),
3035 None => std::env::current_dir()
3036 .map_err(|error| format!("Failed to resolve current directory: {error}"))?,
3037 };
3038 let root_raw = git_output(&start, &["rev-parse", "--show-toplevel"])?;
3039 let root = normalize_windows_verbatim_path(
3040 fs::canonicalize(root_raw.trim()).unwrap_or_else(|_| PathBuf::from(root_raw.trim())),
3041 );
3042 let branch = git_output(&root, &["rev-parse", "--abbrev-ref", "HEAD"])
3043 .unwrap_or_else(|_| "unknown".to_string());
3044 let remote = git_output(&root, &["remote", "get-url", DEFAULT_REMOTE]).unwrap_or_default();
3045 let (owner, name) = parse_remote_owner_repo(&remote).unwrap_or_else(|| {
3046 let name = root
3047 .file_name()
3048 .and_then(|value| value.to_str())
3049 .unwrap_or("repository")
3050 .to_string();
3051 ("unknown".to_string(), name)
3052 });
3053 let head_sha = current_head(&root)?;
3054
3055 Ok(RepoIdentity {
3056 owner,
3057 name,
3058 branch: sanitize_path_component(branch.trim()),
3059 root,
3060 head_sha,
3061 })
3062}
3063
3064fn prepare_spool_layout(identity: &RepoIdentity) -> Result<SpoolLayout, String> {
3065 let home = dirs::home_dir().ok_or_else(|| "Failed to resolve home directory.".to_string())?;
3066 let run_id = Uuid::new_v4().to_string();
3067 let root = home
3068 .join(".xbp")
3069 .join("mutations")
3070 .join(sanitize_path_component(&identity.owner))
3071 .join(sanitize_path_component(&identity.name))
3072 .join(sanitize_path_component(&identity.branch));
3073 fs::create_dir_all(&root).map_err(|error| {
3074 format!(
3075 "Failed to create worktree mutation spool {}: {error}",
3076 root.display()
3077 )
3078 })?;
3079 let event_dedupe_dir = root.join("event-dedupe");
3080 fs::create_dir_all(&event_dedupe_dir).map_err(|error| {
3081 format!(
3082 "Failed to create worktree mutation dedupe dir {}: {error}",
3083 event_dedupe_dir.display()
3084 )
3085 })?;
3086 let commit_dedupe_dir = root.join("commit-dedupe");
3087 fs::create_dir_all(&commit_dedupe_dir).map_err(|error| {
3088 format!(
3089 "Failed to create worktree commit dedupe dir {}: {error}",
3090 commit_dedupe_dir.display()
3091 )
3092 })?;
3093 Ok(SpoolLayout {
3094 events_file: root.join(format!("events-{run_id}.jsonl")),
3095 commits_file: root.join(format!("commits-{run_id}.jsonl")),
3096 stats_file: root.join("stats.json"),
3097 watcher_state_file: root.join("watcher-state.json"),
3098 sync_log_file: root.join("sync-log.jsonl"),
3099 event_dedupe_file: root.join("event-dedupe.jsonl"),
3100 event_dedupe_dir,
3101 commit_dedupe_dir,
3102 root,
3103 })
3104}
3105
3106fn prepare_parent_spool_layout(parent: &Path) -> Result<ParentSpoolLayout, String> {
3107 let home = dirs::home_dir().ok_or_else(|| "Failed to resolve home directory.".to_string())?;
3108 let parent = canonical_parent_path(parent)?;
3109 let parent_key = path_identity_key(&parent);
3110 let mut hasher = Sha256::new();
3111 hasher.update(parent_key.as_bytes());
3112 let digest = format!("{:x}", hasher.finalize());
3113 let label = parent
3114 .file_name()
3115 .and_then(|value| value.to_str())
3116 .map(sanitize_path_component)
3117 .filter(|value| !value.is_empty())
3118 .unwrap_or_else(|| "parent".to_string());
3119 let root = home
3120 .join(".xbp")
3121 .join("mutations")
3122 .join("parents")
3123 .join(format!("{}-{}", label, &digest[..16]));
3124 fs::create_dir_all(&root).map_err(|error| {
3125 format!(
3126 "Failed to create parent worktree watcher state dir {}: {error}",
3127 root.display()
3128 )
3129 })?;
3130 Ok(ParentSpoolLayout {
3131 watcher_state_file: root.join("watcher-state.json"),
3132 })
3133}
3134
3135fn canonical_parent_path(parent: &Path) -> Result<PathBuf, String> {
3136 let normalized = normalize_windows_verbatim_path(
3137 fs::canonicalize(parent).unwrap_or_else(|_| parent.to_path_buf()),
3138 );
3139 if !normalized.is_dir() {
3140 return Err(format!(
3141 "Parent folder {} does not exist or is not a directory.",
3142 normalized.display()
3143 ));
3144 }
3145 Ok(normalized)
3146}
3147
3148fn append_json_line<T: Serialize>(path: &Path, value: &T) -> Result<(), String> {
3149 let mut file = OpenOptions::new()
3150 .create(true)
3151 .append(true)
3152 .open(path)
3153 .map_err(|error| format!("Failed to open spool file {}: {error}", path.display()))?;
3154 let line = serde_json::to_string(value)
3155 .map_err(|error| format!("Failed to serialize worktree event: {error}"))?;
3156 writeln!(file, "{line}")
3157 .map_err(|error| format!("Failed to append spool file {}: {error}", path.display()))
3158}
3159
3160fn git_numstat_for_path(repo_root: &Path, relative_path: &str) -> Result<(u64, u64), String> {
3161 let output = Command::new("git")
3162 .args(["diff", "--numstat", "--"])
3163 .arg(relative_path)
3164 .current_dir(repo_root)
3165 .output()
3166 .map_err(|error| format!("Failed to run git diff --numstat: {error}"))?;
3167 if !output.status.success() {
3168 return Ok((0, 0));
3169 }
3170 let stdout = String::from_utf8_lossy(&output.stdout);
3171 let mut added = 0;
3172 let mut removed = 0;
3173 for line in stdout.lines() {
3174 let mut parts = line.split_whitespace();
3175 added += parse_numstat_count(parts.next());
3176 removed += parse_numstat_count(parts.next());
3177 }
3178 Ok((added, removed))
3179}
3180
3181fn file_line_count_for_path(repo_root: &Path, relative_path: &str) -> Result<u64, String> {
3182 let path = repo_root.join(relative_path);
3183 let content = fs::read_to_string(&path)
3184 .map_err(|error| format!("Failed to read file {}: {error}", path.display()))?;
3185 Ok(content.lines().count() as u64)
3186}
3187
3188fn parse_numstat_count(value: Option<&str>) -> u64 {
3189 value.and_then(|raw| raw.parse::<u64>().ok()).unwrap_or(0)
3190}
3191
3192fn current_head(repo_root: &Path) -> Result<Option<String>, String> {
3193 match git_output(repo_root, &["rev-parse", "HEAD"]) {
3194 Ok(value) => Ok(Some(value)),
3195 Err(error)
3196 if error.contains("unknown revision") || error.contains("ambiguous argument") =>
3197 {
3198 Ok(None)
3199 }
3200 Err(error) => Err(error),
3201 }
3202}
3203
3204fn git_output(repo_root: &Path, args: &[&str]) -> Result<String, String> {
3205 let output = Command::new("git")
3206 .args(args)
3207 .current_dir(repo_root)
3208 .output()
3209 .map_err(|error| format!("Failed to run git {}: {error}", args.join(" ")))?;
3210 if !output.status.success() {
3211 return Err(String::from_utf8_lossy(&output.stderr).trim().to_string());
3212 }
3213 Ok(String::from_utf8_lossy(&output.stdout).trim().to_string())
3214}
3215
3216fn parse_remote_owner_repo(remote: &str) -> Option<(String, String)> {
3217 let trimmed = remote.trim().trim_end_matches(".git");
3218 if trimmed.is_empty() {
3219 return None;
3220 }
3221
3222 let path_part = if !trimmed.contains("://") {
3223 if let Some((_, path)) = trimmed.rsplit_once(':') {
3224 path
3225 } else {
3226 trimmed
3227 }
3228 } else {
3229 trimmed
3230 .trim_start_matches("https://")
3231 .trim_start_matches("http://")
3232 .trim_start_matches("ssh://")
3233 .split_once('/')
3234 .map(|(_, path)| path)
3235 .unwrap_or(trimmed)
3236 };
3237 let mut parts = path_part.rsplitn(2, '/');
3238 let name = parts.next()?.trim();
3239 let owner = parts.next()?.trim();
3240 if owner.is_empty() || name.is_empty() {
3241 return None;
3242 }
3243 Some((
3244 sanitize_path_component(owner),
3245 sanitize_path_component(name.trim_end_matches(".git")),
3246 ))
3247}
3248
3249fn sanitize_path_component(value: &str) -> String {
3250 let sanitized: String = value
3251 .chars()
3252 .map(|ch| {
3253 if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
3254 ch
3255 } else {
3256 '-'
3257 }
3258 })
3259 .collect();
3260 sanitized
3261 .trim_matches('-')
3262 .chars()
3263 .take(160)
3264 .collect::<String>()
3265}
3266
3267fn normalize_windows_verbatim_path(path: PathBuf) -> PathBuf {
3268 PathBuf::from(strip_windows_verbatim_prefix(&path.to_string_lossy()))
3269}
3270
3271fn strip_windows_verbatim_prefix(input: &str) -> &str {
3272 input.strip_prefix(r"\\?\").unwrap_or(input)
3273}
3274
3275fn path_identity_key(path: &Path) -> String {
3276 let value = path.to_string_lossy().to_string();
3277 if cfg!(windows) {
3278 value.to_ascii_lowercase()
3279 } else {
3280 value
3281 }
3282}
3283
3284fn path_string_has_ignored_component(path: &str) -> bool {
3285 watch_ignore_rules().ignores_relative(path)
3286}
3287
3288fn is_skipped_discovery_dir(path: &Path) -> bool {
3289 watch_ignore_rules().skips_discovery_dir(path)
3290}
3291
3292fn relative_slash_path(root: &Path, path: &Path) -> Option<String> {
3293 let relative = path.strip_prefix(root).ok().unwrap_or(path);
3294 Some(relative.to_string_lossy().replace('\\', "/"))
3295}
3296
3297fn is_ignored_path(root: &Path, path: &Path) -> bool {
3298 watch_ignore_rules().ignores_path(root, path)
3299}
3300
3301fn current_hostname() -> Option<String> {
3302 std::env::var("COMPUTERNAME")
3303 .or_else(|_| std::env::var("HOSTNAME"))
3304 .ok()
3305 .map(|value| value.trim().to_string())
3306 .filter(|value| !value.is_empty())
3307}
3308
3309fn is_background_child() -> bool {
3310 std::env::var(BACKGROUND_CHILD_ENV)
3311 .ok()
3312 .map(|value| value == "1")
3313 .unwrap_or(false)
3314}
3315
3316#[allow(dead_code)]
3317fn content_sha256(path: &Path) -> Option<String> {
3318 let bytes = fs::read(path).ok()?;
3319 let mut hasher = Sha256::new();
3320 hasher.update(bytes);
3321 Some(format!("{:x}", hasher.finalize()))
3322}
3323
3324#[cfg(test)]
3325mod tests {
3326 use super::*;
3327
3328 #[test]
3329 fn parses_https_and_ssh_remote_urls() {
3330 assert_eq!(
3331 parse_remote_owner_repo("https://github.com/xylex-group/xbp.git"),
3332 Some(("xylex-group".to_string(), "xbp".to_string()))
3333 );
3334 assert_eq!(
3335 parse_remote_owner_repo("git@github.com:xylex-group/xbp.git"),
3336 Some(("xylex-group".to_string(), "xbp".to_string()))
3337 );
3338 }
3339
3340 #[test]
3341 fn sanitizes_branch_for_storage_path() {
3342 assert_eq!(
3343 sanitize_path_component("feature/worktree watcher"),
3344 "feature-worktree-watcher"
3345 );
3346 }
3347
3348 #[test]
3349 fn normalizes_windows_verbatim_repo_roots() {
3350 assert_eq!(
3351 normalize_windows_verbatim_path(PathBuf::from(
3352 r"\\?\C:\Users\floris\Documents\GitHub\mollie-api-rust"
3353 )),
3354 PathBuf::from(r"C:\Users\floris\Documents\GitHub\mollie-api-rust")
3355 );
3356 }
3357
3358 #[test]
3359 fn parent_path_matching_uses_repo_boundaries() {
3360 let repo = Path::new(r"C:\Users\floris\Documents\GitHub\xbp");
3361 assert!(path_belongs_to_repo(
3362 Path::new(r"C:\Users\floris\Documents\GitHub\xbp\src\main.rs"),
3363 repo
3364 ));
3365 assert!(!path_belongs_to_repo(
3366 Path::new(r"C:\Users\floris\Documents\GitHub\xbp-other\src\main.rs"),
3367 repo
3368 ));
3369 }
3370
3371 #[test]
3372 fn skips_heavy_discovery_directories() {
3373 assert!(is_skipped_discovery_dir(Path::new("node_modules")));
3374 assert!(is_skipped_discovery_dir(Path::new("target")));
3375 assert!(is_skipped_discovery_dir(Path::new(".next")));
3376 assert!(!is_skipped_discovery_dir(Path::new("mollie-api-rust")));
3377 }
3378
3379 #[test]
3380 fn ignores_nested_generated_mutation_paths() {
3381 let root = Path::new(r"C:\Users\floris\Documents\GitHub\athena");
3382 assert!(is_ignored_path(
3383 root,
3384 Path::new(
3385 r"C:\Users\floris\Documents\GitHub\athena\apps\docs\node_modules\@aws-sdk\client-lambda\dist-types\schemas"
3386 )
3387 ));
3388 assert!(is_ignored_path(
3389 root,
3390 Path::new(r"C:\Users\floris\Documents\GitHub\athena\target\debug\build")
3391 ));
3392 assert!(!is_ignored_path(
3393 root,
3394 Path::new(r"C:\Users\floris\Documents\GitHub\athena\apps\web\src\main.ts")
3395 ));
3396 }
3397
3398 #[test]
3399 fn global_config_forbidden_paths_folders_and_banned_words() {
3400 let rules = WatchIgnoreRules::from_config(&WorktreeWatchConfig {
3401 forbidden_paths: vec![
3402 "secrets/".to_string(),
3403 "apps/web/.env.local".to_string(),
3404 ],
3405 forbidden_folders: vec![".cache".to_string(), "coverage".to_string()],
3406 banned_words: vec!["private-key".to_string(), "CREDENTIAL".to_string()],
3407 });
3408
3409 assert!(rules.ignores_relative("secrets/prod/api.key"));
3411 assert!(rules.ignores_relative("apps/web/.env.local"));
3412 assert!(!rules.ignores_relative("apps/web/.env.local.bak"));
3414
3415 assert!(rules.ignores_relative("apps/web/.cache/tmp"));
3417 assert!(rules.ignores_relative("packages/foo/coverage/index.html"));
3418
3419 assert!(rules.ignores_relative("docs/private-key-notes.md"));
3421 assert!(rules.ignores_relative("config/my-credential-store.json"));
3422
3423 assert!(!rules.ignores_relative("apps/web/src/main.ts"));
3425 assert!(!rules.skips_discovery_dir(Path::new("mollie-api-rust")));
3426 assert!(rules.skips_discovery_dir(Path::new(".cache")));
3427 assert!(rules.skips_discovery_dir(Path::new("node_modules"))); }
3429
3430 #[test]
3431 fn builds_coding_stats_by_file_and_filetype() {
3432 let identity = RepoIdentity {
3433 owner: "xylex-group".to_string(),
3434 name: "xbp".to_string(),
3435 branch: "main".to_string(),
3436 root: PathBuf::from(r"C:\Users\floris\Documents\GitHub\xbp"),
3437 head_sha: None,
3438 };
3439 let layout = SpoolLayout {
3440 root: PathBuf::from(r"C:\Users\floris\.xbp\mutations\xylex-group\xbp\main"),
3441 events_file: PathBuf::from("events.jsonl"),
3442 commits_file: PathBuf::from("commits.jsonl"),
3443 stats_file: PathBuf::from("stats.json"),
3444 watcher_state_file: PathBuf::from("watcher-state.json"),
3445 sync_log_file: PathBuf::from("sync-log.jsonl"),
3446 event_dedupe_file: PathBuf::from("event-dedupe.jsonl"),
3447 event_dedupe_dir: PathBuf::from("event-dedupe"),
3448 commit_dedupe_dir: PathBuf::from("commit-dedupe"),
3449 };
3450 let candidates = vec![SyncCandidate {
3451 path: PathBuf::from("events.jsonl"),
3452 records: SyncRecords::Events(vec![
3453 stats_test_event_record("2026-07-08T10:00:00Z", "src/main.rs", 4, 1),
3454 stats_test_event_record("2026-07-08T10:05:00Z", "src/main.rs", 2, 0),
3455 stats_test_event_record("2026-07-08T11:00:00Z", "README.md", 1, 3),
3456 ]),
3457 }];
3458
3459 let summary = build_stats_summary(&identity, &layout, &candidates, 15).unwrap();
3460
3461 assert_eq!(summary.total_events, 3);
3462 assert_eq!(summary.estimated_coding_seconds, 60 + 300 + 900);
3463 assert_eq!(summary.added_lines, 7);
3464 assert_eq!(summary.removed_lines, 4);
3465 assert_eq!(summary.by_file[0].path, "README.md");
3466 assert_eq!(summary.by_file[0].estimated_coding_seconds, 900);
3467 assert_eq!(summary.by_filetype[0].name, "md");
3468 assert_eq!(summary.by_filetype[0].estimated_coding_seconds, 900);
3469 assert_eq!(summary.by_filetype[1].name, "rs");
3470 assert_eq!(summary.by_filetype[1].estimated_coding_seconds, 360);
3471 }
3472
3473 #[test]
3474 fn mutation_fingerprint_ignores_random_event_id_but_keeps_time_bucket() {
3475 let first: WorktreeMutationEvent = serde_json::from_value(stats_test_event(
3476 "2026-07-08T10:00:00Z",
3477 "src/main.rs",
3478 4,
3479 1,
3480 ))
3481 .unwrap();
3482 let same_bucket: WorktreeMutationEvent = serde_json::from_value(stats_test_event(
3483 "2026-07-08T10:00:00Z",
3484 "src/main.rs",
3485 4,
3486 1,
3487 ))
3488 .unwrap();
3489 let next_bucket: WorktreeMutationEvent = serde_json::from_value(stats_test_event(
3490 "2026-07-08T10:00:06Z",
3491 "src/main.rs",
3492 4,
3493 1,
3494 ))
3495 .unwrap();
3496
3497 assert_eq!(
3498 mutation_event_fingerprint(&first),
3499 mutation_event_fingerprint(&same_bucket)
3500 );
3501 assert_ne!(
3502 mutation_event_fingerprint(&first),
3503 mutation_event_fingerprint(&next_bucket)
3504 );
3505 }
3506
3507 #[test]
3508 fn append_unique_mutation_event_is_idempotent_across_dedupe_instances() {
3509 let temp_root = std::env::temp_dir().join(format!("xbp-watch-test-{}", Uuid::new_v4()));
3510 fs::create_dir_all(&temp_root).unwrap();
3511 let layout = SpoolLayout {
3512 root: temp_root.clone(),
3513 events_file: temp_root.join("events.jsonl"),
3514 commits_file: temp_root.join("commits.jsonl"),
3515 stats_file: temp_root.join("stats.json"),
3516 watcher_state_file: temp_root.join("watcher-state.json"),
3517 sync_log_file: temp_root.join("sync-log.jsonl"),
3518 event_dedupe_file: temp_root.join("event-dedupe.jsonl"),
3519 event_dedupe_dir: temp_root.join("event-dedupe"),
3520 commit_dedupe_dir: temp_root.join("commit-dedupe"),
3521 };
3522 let mutation: WorktreeMutationEvent = serde_json::from_value(stats_test_event(
3523 "2026-07-08T10:00:00Z",
3524 "src/main.rs",
3525 4,
3526 1,
3527 ))
3528 .unwrap();
3529
3530 let first =
3531 append_unique_mutation_event(&layout, &mutation, &mut RecentEventDedupe::default())
3532 .unwrap();
3533 let second =
3534 append_unique_mutation_event(&layout, &mutation, &mut RecentEventDedupe::default())
3535 .unwrap();
3536
3537 assert!(first);
3538 assert!(!second);
3539 let events = read_jsonl_values(&layout.events_file).unwrap();
3540 assert_eq!(events.len(), 1);
3541 assert_eq!(
3542 events[0].get("fingerprint").and_then(Value::as_str),
3543 Some(mutation_event_fingerprint(&mutation).as_str())
3544 );
3545 let _ = fs::remove_dir_all(temp_root);
3546 }
3547
3548 #[test]
3549 fn reads_spool_records_as_typed_events() {
3550 let temp_root = std::env::temp_dir().join(format!("xbp-watch-test-{}", Uuid::new_v4()));
3551 fs::create_dir_all(&temp_root).unwrap();
3552 let events_file = temp_root.join("events.jsonl");
3553 append_json_line(
3554 &events_file,
3555 &stats_test_event_record("2026-07-08T10:00:00Z", "src/main.rs", 4, 1),
3556 )
3557 .unwrap();
3558
3559 let records = read_spool_records(&events_file, SyncFileKind::Events).unwrap();
3560
3561 match records {
3562 SyncRecords::Events(events) => {
3563 assert_eq!(events.len(), 1);
3564 assert_eq!(events[0].primary_path.as_deref(), Some("src/main.rs"));
3565 }
3566 SyncRecords::Commits(_) => panic!("expected typed event records"),
3567 }
3568 let _ = fs::remove_dir_all(temp_root);
3569 }
3570
3571 #[test]
3572 fn collect_spool_candidates_rewrites_legacy_generated_events() {
3573 let temp_root = std::env::temp_dir().join(format!("xbp-watch-test-{}", Uuid::new_v4()));
3574 fs::create_dir_all(&temp_root).unwrap();
3575 let events_file = temp_root.join("events-legacy.jsonl");
3576 append_json_line(
3577 &events_file,
3578 &stats_test_event_record(
3579 "2026-07-08T10:00:00Z",
3580 "apps/docs/node_modules/@aws-sdk/client-lambda/dist-types/schemas",
3581 0,
3582 0,
3583 ),
3584 )
3585 .unwrap();
3586 append_json_line(
3587 &events_file,
3588 &stats_test_event_record("2026-07-08T10:01:00Z", "apps/web/src/main.ts", 4, 1),
3589 )
3590 .unwrap();
3591 let layout = SpoolLayout {
3592 root: temp_root.clone(),
3593 events_file: temp_root.join("events.jsonl"),
3594 commits_file: temp_root.join("commits.jsonl"),
3595 stats_file: temp_root.join("stats.json"),
3596 watcher_state_file: temp_root.join("watcher-state.json"),
3597 sync_log_file: temp_root.join("sync-log.jsonl"),
3598 event_dedupe_file: temp_root.join("event-dedupe.jsonl"),
3599 event_dedupe_dir: temp_root.join("event-dedupe"),
3600 commit_dedupe_dir: temp_root.join("commit-dedupe"),
3601 };
3602
3603 let candidates = collect_spool_candidates(&layout, false).unwrap();
3604
3605 assert_eq!(candidates.len(), 1);
3606 match &candidates[0].records {
3607 SyncRecords::Events(events) => {
3608 assert_eq!(events.len(), 1);
3609 assert_eq!(
3610 events[0].primary_path.as_deref(),
3611 Some("apps/web/src/main.ts")
3612 );
3613 }
3614 SyncRecords::Commits(_) => panic!("expected event spool candidate"),
3615 }
3616 let rewritten = read_jsonl_records::<WorktreeMutationEvent>(&events_file).unwrap();
3617 assert_eq!(rewritten.len(), 1);
3618 assert_eq!(
3619 rewritten[0].primary_path.as_deref(),
3620 Some("apps/web/src/main.ts")
3621 );
3622 let _ = fs::remove_dir_all(temp_root);
3623 }
3624
3625 #[test]
3626 fn collect_spool_candidates_deletes_fully_ignored_event_files() {
3627 let temp_root = std::env::temp_dir().join(format!("xbp-watch-test-{}", Uuid::new_v4()));
3628 fs::create_dir_all(&temp_root).unwrap();
3629 let events_file = temp_root.join("events-legacy.jsonl");
3630 append_json_line(
3631 &events_file,
3632 &stats_test_event_record(
3633 "2026-07-08T10:00:00Z",
3634 "apps/docs/node_modules/@aws-sdk/client-lambda/dist-types/schemas",
3635 0,
3636 0,
3637 ),
3638 )
3639 .unwrap();
3640 let layout = SpoolLayout {
3641 root: temp_root.clone(),
3642 events_file: temp_root.join("events.jsonl"),
3643 commits_file: temp_root.join("commits.jsonl"),
3644 stats_file: temp_root.join("stats.json"),
3645 watcher_state_file: temp_root.join("watcher-state.json"),
3646 sync_log_file: temp_root.join("sync-log.jsonl"),
3647 event_dedupe_file: temp_root.join("event-dedupe.jsonl"),
3648 event_dedupe_dir: temp_root.join("event-dedupe"),
3649 commit_dedupe_dir: temp_root.join("commit-dedupe"),
3650 };
3651
3652 let candidates = collect_spool_candidates(&layout, false).unwrap();
3653
3654 assert!(candidates.is_empty());
3655 assert!(!events_file.exists());
3656 let _ = fs::remove_dir_all(temp_root);
3657 }
3658
3659 #[test]
3660 fn append_unique_commit_event_is_idempotent_across_watcher_instances() {
3661 let temp_root = std::env::temp_dir().join(format!("xbp-watch-test-{}", Uuid::new_v4()));
3662 fs::create_dir_all(&temp_root).unwrap();
3663 let layout = SpoolLayout {
3664 root: temp_root.clone(),
3665 events_file: temp_root.join("events.jsonl"),
3666 commits_file: temp_root.join("commits.jsonl"),
3667 stats_file: temp_root.join("stats.json"),
3668 watcher_state_file: temp_root.join("watcher-state.json"),
3669 sync_log_file: temp_root.join("sync-log.jsonl"),
3670 event_dedupe_file: temp_root.join("event-dedupe.jsonl"),
3671 event_dedupe_dir: temp_root.join("event-dedupe"),
3672 commit_dedupe_dir: temp_root.join("commit-dedupe"),
3673 };
3674 let commit = WorktreeCommitEvent {
3675 id: Uuid::new_v4().to_string(),
3676 fingerprint: None,
3677 repo_owner: "xylex-group".to_string(),
3678 repo_name: "xbp".to_string(),
3679 branch_name: "main".to_string(),
3680 repo_root: r"C:\Users\floris\Documents\GitHub\xbp".to_string(),
3681 previous_head_sha: Some("previous".to_string()),
3682 head_sha: "current".to_string(),
3683 subject: Some("test commit".to_string()),
3684 author_name: Some("Floris".to_string()),
3685 author_email: Some("floris@xylex.group".to_string()),
3686 committed_at: Some("2026-07-08T10:00:00Z".to_string()),
3687 occurred_at: Utc::now(),
3688 };
3689
3690 assert!(append_unique_commit_event(&layout, &commit).unwrap());
3691 assert!(!append_unique_commit_event(&layout, &commit).unwrap());
3692 let commits = read_jsonl_values(&layout.commits_file).unwrap();
3693 assert_eq!(commits.len(), 1);
3694 assert_eq!(
3695 commits[0].get("fingerprint").and_then(Value::as_str),
3696 Some(commit_event_fingerprint(&commit).as_str())
3697 );
3698 let _ = fs::remove_dir_all(temp_root);
3699 }
3700
3701 #[test]
3702 fn splits_worktree_mutation_uploads_into_bounded_batches() {
3703 let events = (0..(WORKTREE_MUTATION_SYNC_BATCH_RECORD_LIMIT + 1))
3704 .map(|index| {
3705 serde_json::from_value::<WorktreeMutationEvent>(stats_test_event(
3706 "2026-07-08T10:00:00Z",
3707 &format!("src/file-{index}.rs"),
3708 1,
3709 0,
3710 ))
3711 .unwrap()
3712 })
3713 .collect::<Vec<_>>();
3714 let commits = vec![
3715 worktree_commit_test_event("previous-a", "current-a"),
3716 worktree_commit_test_event("previous-b", "current-b"),
3717 ];
3718
3719 let batches = split_worktree_mutation_upload_batches(events, commits);
3720
3721 assert_eq!(batches.len(), 2);
3722 assert_eq!(
3723 batches[0].events.len() + batches[0].commits.len(),
3724 WORKTREE_MUTATION_SYNC_BATCH_RECORD_LIMIT
3725 );
3726 assert_eq!(batches[1].events.len(), 1);
3727 assert_eq!(batches[1].commits.len(), 2);
3728 }
3729
3730 #[test]
3731 fn splits_worktree_mutation_uploads_by_estimated_json_bytes() {
3732 let long_path = format!("src/{}.rs", "x".repeat(8_000));
3733 let events = (0..200)
3734 .map(|index| {
3735 serde_json::from_value::<WorktreeMutationEvent>(stats_test_event(
3736 "2026-07-08T10:00:00Z",
3737 &format!("{long_path}-{index}"),
3738 1,
3739 0,
3740 ))
3741 .unwrap()
3742 })
3743 .collect::<Vec<_>>();
3744
3745 let batches = split_worktree_mutation_upload_batches(events, Vec::new());
3746
3747 assert!(batches.len() > 1);
3748 for batch in batches {
3749 assert!(
3750 batch.estimated_json_bytes <= WORKTREE_MUTATION_SYNC_BATCH_JSON_BYTES,
3751 "batch estimated JSON bytes exceeded limit: {}",
3752 batch.estimated_json_bytes
3753 );
3754 }
3755 }
3756
3757 #[test]
3758 fn classifies_create_remove_and_rename_events() {
3759 assert_eq!(
3760 classify_event_kind(&EventKind::Create(CreateKind::File)),
3761 "file_create"
3762 );
3763 assert_eq!(
3764 classify_event_kind(&EventKind::Remove(RemoveKind::Folder)),
3765 "folder_remove"
3766 );
3767 assert_eq!(
3768 classify_event_kind(&EventKind::Modify(ModifyKind::Name(RenameMode::Both))),
3769 "rename_or_move"
3770 );
3771 }
3772
3773 #[test]
3774 fn retries_only_retryable_worktree_sync_statuses() {
3775 assert!(should_retry_worktree_mutation_sync_status(
3776 StatusCode::INTERNAL_SERVER_ERROR
3777 ));
3778 assert!(should_retry_worktree_mutation_sync_status(
3779 StatusCode::REQUEST_TIMEOUT
3780 ));
3781 assert!(should_retry_worktree_mutation_sync_status(
3782 StatusCode::TOO_MANY_REQUESTS
3783 ));
3784 assert!(!should_retry_worktree_mutation_sync_status(
3785 StatusCode::BAD_REQUEST
3786 ));
3787 assert!(!should_retry_worktree_mutation_sync_status(
3788 StatusCode::UNAUTHORIZED
3789 ));
3790 assert!(!should_retry_worktree_mutation_sync_status(
3791 StatusCode::PAYLOAD_TOO_LARGE
3792 ));
3793 }
3794
3795 fn stats_test_event(
3796 occurred_at: &str,
3797 path: &str,
3798 added_lines: u64,
3799 removed_lines: u64,
3800 ) -> Value {
3801 json!({
3802 "id": Uuid::new_v4().to_string(),
3803 "repoOwner": "xylex-group",
3804 "repoName": "xbp",
3805 "branchName": "main",
3806 "repoRoot": r"C:\Users\floris\Documents\GitHub\xbp",
3807 "headSha": null,
3808 "eventKind": "file_modify",
3809 "paths": [path],
3810 "primaryPath": path,
3811 "oldPath": null,
3812 "newPath": null,
3813 "addedLines": added_lines,
3814 "removedLines": removed_lines,
3815 "totalLines": added_lines + 42,
3816 "fileCreated": false,
3817 "fileRemoved": false,
3818 "folderCreated": false,
3819 "folderRemoved": false,
3820 "renamedOrMoved": false,
3821 "rawKind": "Modify(Data(Content))",
3822 "occurredAt": occurred_at,
3823 })
3824 }
3825
3826 fn stats_test_event_record(
3827 occurred_at: &str,
3828 path: &str,
3829 added_lines: u64,
3830 removed_lines: u64,
3831 ) -> WorktreeMutationEvent {
3832 serde_json::from_value(stats_test_event(
3833 occurred_at,
3834 path,
3835 added_lines,
3836 removed_lines,
3837 ))
3838 .unwrap()
3839 }
3840
3841 fn worktree_commit_test_event(previous_head_sha: &str, head_sha: &str) -> WorktreeCommitEvent {
3842 WorktreeCommitEvent {
3843 id: Uuid::new_v4().to_string(),
3844 fingerprint: None,
3845 repo_owner: "xylex-group".to_string(),
3846 repo_name: "xbp".to_string(),
3847 branch_name: "main".to_string(),
3848 repo_root: r"C:\Users\floris\Documents\GitHub\xbp".to_string(),
3849 previous_head_sha: Some(previous_head_sha.to_string()),
3850 head_sha: head_sha.to_string(),
3851 subject: Some("test commit".to_string()),
3852 author_name: Some("Floris".to_string()),
3853 author_email: Some("floris@xylex.group".to_string()),
3854 committed_at: Some("2026-07-08T10:00:00Z".to_string()),
3855 occurred_at: Utc::now(),
3856 }
3857 }
3858}