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