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