Skip to main content

xbp_cli/commands/
worktree_watch.rs

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/// Client-side ignore rules: built-in VCS/generated dirs + global config filters.
31/// Never uploaded; applied before events are spooled and again before sync.
32#[derive(Debug, Clone, Default)]
33struct WatchIgnoreRules {
34    /// Normalized lowercase prefixes without leading `./` (e.g. `secrets`, `apps/web/.env`).
35    forbidden_path_prefixes: Vec<String>,
36    /// Lowercase path segment names blocked anywhere (includes built-ins + config).
37    forbidden_folders: BTreeSet<String>,
38    /// Lowercase substrings matched against the full relative path.
39    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            // Also treat a single-segment prefix as a folder ban anywhere.
101            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(&current, 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(&current, 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(&current) > 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        // Nested under a watched tree still blocked by prefix.
3753        assert!(rules.ignores_relative("secrets/prod/api.key"));
3754        assert!(rules.ignores_relative("apps/web/.env.local"));
3755        // Exact file ban is not a string-prefix of other filenames.
3756        assert!(!rules.ignores_relative("apps/web/.env.local.bak"));
3757
3758        // Folder name anywhere
3759        assert!(rules.ignores_relative("apps/web/.cache/tmp"));
3760        assert!(rules.ignores_relative("packages/foo/coverage/index.html"));
3761
3762        // Banned word (case-insensitive substring)
3763        assert!(rules.ignores_relative("docs/private-key-notes.md"));
3764        assert!(rules.ignores_relative("config/my-credential-store.json"));
3765
3766        // Still allow normal source
3767        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"))); // built-in
3771    }
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}