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;
23
24#[cfg(windows)]
25const CREATE_NO_WINDOW: u32 = 0x08000000;
26
27#[derive(Debug, Clone)]
28pub struct WorktreeWatchTargetOptions {
29    pub repo: Option<PathBuf>,
30    pub parent: Option<PathBuf>,
31}
32
33#[derive(Debug, Clone)]
34pub struct WorktreeWatchStartOptions {
35    pub target: WorktreeWatchTargetOptions,
36    pub detach: bool,
37    pub sync_interval_seconds: u64,
38    pub once: bool,
39}
40
41#[derive(Debug, Clone)]
42pub struct WorktreeWatchSyncOptions {
43    pub target: WorktreeWatchTargetOptions,
44    pub dry_run: bool,
45    pub resync: bool,
46}
47
48#[derive(Debug, Clone)]
49pub struct WorktreeWatchStatusOptions {
50    pub target: WorktreeWatchTargetOptions,
51    pub json: bool,
52    pub records: bool,
53    pub record_limit: usize,
54    pub stats: bool,
55    pub stats_gap_minutes: u64,
56}
57
58#[derive(Debug, Clone)]
59struct RepoIdentity {
60    owner: String,
61    name: String,
62    branch: String,
63    root: PathBuf,
64    head_sha: Option<String>,
65}
66
67#[derive(Debug, Clone)]
68struct SpoolLayout {
69    root: PathBuf,
70    events_file: PathBuf,
71    commits_file: PathBuf,
72    stats_file: PathBuf,
73    watcher_state_file: PathBuf,
74    sync_log_file: PathBuf,
75    event_dedupe_file: PathBuf,
76}
77
78#[derive(Debug, Serialize, Deserialize, Clone)]
79#[serde(rename_all = "camelCase")]
80struct WorktreeMutationEvent {
81    id: String,
82    repo_owner: String,
83    repo_name: String,
84    branch_name: String,
85    repo_root: String,
86    head_sha: Option<String>,
87    event_kind: String,
88    paths: Vec<String>,
89    primary_path: Option<String>,
90    old_path: Option<String>,
91    new_path: Option<String>,
92    added_lines: Option<u64>,
93    removed_lines: Option<u64>,
94    file_created: bool,
95    file_removed: bool,
96    folder_created: bool,
97    folder_removed: bool,
98    renamed_or_moved: bool,
99    raw_kind: String,
100    occurred_at: DateTime<Utc>,
101}
102
103#[derive(Debug, Serialize, Deserialize, Clone)]
104#[serde(rename_all = "camelCase")]
105struct WorktreeCommitEvent {
106    id: String,
107    repo_owner: String,
108    repo_name: String,
109    branch_name: String,
110    repo_root: String,
111    previous_head_sha: Option<String>,
112    head_sha: String,
113    subject: Option<String>,
114    author_name: Option<String>,
115    author_email: Option<String>,
116    committed_at: Option<String>,
117    occurred_at: DateTime<Utc>,
118}
119
120#[derive(Debug, Serialize)]
121#[serde(rename_all = "camelCase")]
122struct WorktreeMutationIngestPayload {
123    device: WorktreeMutationDevicePayload,
124    repository: WorktreeMutationRepositoryPayload,
125    events: Vec<WorktreeMutationEvent>,
126    commits: Vec<WorktreeCommitEvent>,
127}
128
129#[derive(Debug, Serialize)]
130#[serde(rename_all = "camelCase")]
131struct WorktreeMutationDevicePayload {
132    hardware_id: String,
133    device_name: Option<String>,
134    hostname: Option<String>,
135    platform: String,
136}
137
138#[derive(Debug, Serialize)]
139#[serde(rename_all = "camelCase")]
140struct WorktreeMutationRepositoryPayload {
141    owner: String,
142    name: String,
143    branch_name: String,
144    repo_root: String,
145}
146
147#[derive(Debug)]
148struct SyncCandidate {
149    path: PathBuf,
150    kind: SyncFileKind,
151    records: Vec<Value>,
152}
153
154#[derive(Debug, Clone, Copy)]
155enum SyncFileKind {
156    Events,
157    Commits,
158}
159
160#[derive(Debug, Serialize, Clone)]
161#[serde(rename_all = "camelCase")]
162struct WorktreeWatchStatsSummary {
163    generated_at: DateTime<Utc>,
164    repository_owner: String,
165    repository_name: String,
166    branch_name: String,
167    repo_root: String,
168    spool_root: String,
169    session_gap_minutes: u64,
170    total_events: u64,
171    total_commits: u64,
172    first_event_at: Option<DateTime<Utc>>,
173    last_event_at: Option<DateTime<Utc>>,
174    estimated_coding_seconds: u64,
175    added_lines: u64,
176    removed_lines: u64,
177    by_file: Vec<WorktreeWatchFileStats>,
178    by_filetype: Vec<WorktreeWatchBucketStats>,
179    by_event_kind: Vec<WorktreeWatchBucketStats>,
180}
181
182#[derive(Debug, Serialize, Clone, Default)]
183#[serde(rename_all = "camelCase")]
184struct WorktreeWatchFileStats {
185    path: String,
186    filetype: String,
187    event_count: u64,
188    estimated_coding_seconds: u64,
189    added_lines: u64,
190    removed_lines: u64,
191}
192
193#[derive(Debug, Serialize, Clone, Default)]
194#[serde(rename_all = "camelCase")]
195struct WorktreeWatchBucketStats {
196    name: String,
197    event_count: u64,
198    estimated_coding_seconds: u64,
199    added_lines: u64,
200    removed_lines: u64,
201}
202
203#[derive(Debug, Serialize, Deserialize, Clone)]
204#[serde(rename_all = "camelCase")]
205struct WorktreeWatcherState {
206    pid: u32,
207    repo_root: String,
208    repository_owner: String,
209    repository_name: String,
210    branch_name: String,
211    executable: String,
212    started_at: DateTime<Utc>,
213}
214
215#[derive(Debug, Serialize, Deserialize, Clone)]
216#[serde(rename_all = "camelCase")]
217struct WorktreeWatchSyncLogEntry {
218    id: String,
219    synced_at: DateTime<Utc>,
220    endpoint: String,
221    repository_owner: String,
222    repository_name: String,
223    branch_name: String,
224    repo_root: String,
225    resync: bool,
226    status_code: u16,
227    event_count: usize,
228    commit_count: usize,
229    spool_file_count: usize,
230    spool_files: Vec<String>,
231}
232
233#[derive(Debug, Serialize, Deserialize, Clone)]
234#[serde(rename_all = "camelCase")]
235struct WorktreeWatchEventDedupeRecord {
236    fingerprint: String,
237    first_seen_at: DateTime<Utc>,
238    event_kind: String,
239    primary_path: Option<String>,
240}
241
242#[derive(Debug, Default)]
243struct RecentEventDedupe {
244    fingerprints: BTreeSet<String>,
245}
246
247pub async fn run_worktree_watch_start(options: WorktreeWatchStartOptions) -> Result<(), String> {
248    if options.target.parent.is_some() {
249        let identities = resolve_target_identities(&options.target)?;
250        if identities.is_empty() {
251            println!("No git repositories found under parent folder.");
252            return Ok(());
253        }
254
255        if options.detach {
256            for identity in &identities {
257                let path = spawn_detached_worktree_watch_for_identity(identity)?;
258                println!("Started background worktree watcher for {}", path.display());
259            }
260            println!(
261                "Started {} background worktree watcher(s).",
262                identities.len()
263            );
264            return Ok(());
265        }
266
267        if options.once {
268            for identity in &identities {
269                let layout = prepare_spool_layout(identity)?;
270                record_commit_snapshot(identity, &layout, None)?;
271                println!("Recorded current git state for {}", identity.root.display());
272            }
273            return Ok(());
274        }
275
276        return Err(
277            "Parent-folder watching starts one watcher per repository. Pass `--detach` to run them in the background, or `--once` to snapshot each repo."
278                .to_string(),
279        );
280    }
281
282    if options.detach {
283        return spawn_detached_worktree_watch(options.target.repo.as_deref()).map(|path| {
284            println!("Started background worktree watcher for {}", path.display());
285        });
286    }
287
288    let identity = resolve_repo_identity(options.target.repo.as_deref())?;
289    let layout = prepare_spool_layout(&identity)?;
290    println!(
291        "Watching {} and spooling mutations to {}",
292        identity.root.display(),
293        layout.root.display()
294    );
295
296    if options.once {
297        record_commit_snapshot(&identity, &layout, None)?;
298        return Ok(());
299    }
300
301    watch_foreground(identity, layout, options.sync_interval_seconds).await
302}
303
304pub async fn run_worktree_watch_sync(options: WorktreeWatchSyncOptions) -> Result<(), String> {
305    let identities = resolve_target_identities(&options.target)?;
306    let repo_count = identities.len();
307    let mut plans = Vec::new();
308    let mut total_files = 0usize;
309    let mut total_records = 0usize;
310
311    for identity in identities {
312        let layout = prepare_spool_layout(&identity)?;
313        let candidates = collect_spool_candidates(&layout, options.resync)?;
314        total_files += candidates.len();
315        total_records += candidates
316            .iter()
317            .map(|candidate| candidate.records.len())
318            .sum::<usize>();
319        if !candidates.is_empty() {
320            plans.push((identity, candidates));
321        }
322    }
323
324    if options.dry_run {
325        println!(
326            "Would sync {} record(s) from {} spool file(s) across {} repo(s).",
327            total_records, total_files, repo_count
328        );
329        if options.resync {
330            println!("Resync mode includes files already marked `.synced.*.jsonl`.");
331        }
332        return Ok(());
333    }
334
335    if plans.is_empty() {
336        print_no_sync_candidates_message(&options.target, options.resync)?;
337        return Ok(());
338    }
339
340    let upload_repo_count = plans.len();
341    for (identity, candidates) in plans {
342        sync_candidates(&identity, candidates, options.resync).await?;
343    }
344    println!(
345        "Synced {} worktree mutation record(s) across {} repo(s).",
346        total_records, upload_repo_count
347    );
348    Ok(())
349}
350
351pub async fn run_worktree_watch_status(options: WorktreeWatchStatusOptions) -> Result<(), String> {
352    let identities = resolve_target_identities(&options.target)?;
353    let mut payloads = Vec::new();
354    let mut total_files = 0usize;
355    let mut total_records = 0usize;
356
357    for identity in identities {
358        let status =
359            worktree_watch_status_payload(&identity, options.records, options.record_limit)?;
360        let stats = if options.stats {
361            Some(generate_and_store_stats(
362                &identity,
363                options.stats_gap_minutes,
364            )?)
365        } else {
366            None
367        };
368        let status = attach_optional_stats(status, stats)?;
369        total_files += status
370            .get("unsyncedFiles")
371            .and_then(Value::as_u64)
372            .unwrap_or(0) as usize;
373        total_records += status
374            .get("unsyncedRecords")
375            .and_then(Value::as_u64)
376            .unwrap_or(0) as usize;
377        payloads.push(status);
378    }
379
380    if options.json {
381        let payload = if options.target.parent.is_some() {
382            json!({
383                "repositories": payloads,
384                "repositoryCount": payloads.len(),
385                "unsyncedFiles": total_files,
386                "unsyncedRecords": total_records,
387            })
388        } else {
389            payloads.into_iter().next().unwrap_or_else(|| json!({}))
390        };
391        println!(
392            "{}",
393            serde_json::to_string_pretty(&payload)
394                .map_err(|error| format!("Failed to render status JSON: {error}"))?
395        );
396    } else {
397        for (index, payload) in payloads.iter().enumerate() {
398            if index > 0 {
399                println!();
400            }
401            println!(
402                "repo: {}/{}",
403                payload
404                    .get("repositoryOwner")
405                    .and_then(Value::as_str)
406                    .unwrap_or("unknown"),
407                payload
408                    .get("repositoryName")
409                    .and_then(Value::as_str)
410                    .unwrap_or("repository")
411            );
412            println!(
413                "branch: {}",
414                payload
415                    .get("branchName")
416                    .and_then(Value::as_str)
417                    .unwrap_or("unknown")
418            );
419            println!(
420                "root: {}",
421                payload
422                    .get("repoRoot")
423                    .and_then(Value::as_str)
424                    .unwrap_or("")
425            );
426            println!(
427                "spool: {}",
428                payload
429                    .get("spoolRoot")
430                    .and_then(Value::as_str)
431                    .unwrap_or("")
432            );
433            println!(
434                "unsynced files: {}",
435                payload
436                    .get("unsyncedFiles")
437                    .and_then(Value::as_u64)
438                    .unwrap_or(0)
439            );
440            println!(
441                "unsynced records: {}",
442                payload
443                    .get("unsyncedRecords")
444                    .and_then(Value::as_u64)
445                    .unwrap_or(0)
446            );
447            println!(
448                "synced files: {}",
449                payload
450                    .get("syncedFiles")
451                    .and_then(Value::as_u64)
452                    .unwrap_or(0)
453            );
454            println!(
455                "synced records: {}",
456                payload
457                    .get("syncedRecords")
458                    .and_then(Value::as_u64)
459                    .unwrap_or(0)
460            );
461            print_last_sync(payload);
462            if options.records {
463                print_status_records(payload);
464            }
465            if options.stats {
466                print_status_stats(payload);
467            }
468        }
469        if options.target.parent.is_some() {
470            println!();
471            println!(
472                "total: {} repo(s), {} unsynced file(s), {} unsynced record(s)",
473                payloads.len(),
474                total_files,
475                total_records
476            );
477        }
478    }
479
480    Ok(())
481}
482
483fn attach_optional_stats(
484    mut payload: Value,
485    stats: Option<WorktreeWatchStatsSummary>,
486) -> Result<Value, String> {
487    if let Some(stats) = stats {
488        if let Some(object) = payload.as_object_mut() {
489            object.insert(
490                "stats".to_string(),
491                serde_json::to_value(stats)
492                    .map_err(|error| format!("Failed to render stats payload: {error}"))?,
493            );
494        }
495    }
496    Ok(payload)
497}
498
499pub fn spawn_detached_worktree_watch_for_current_repo() -> Result<Option<PathBuf>, String> {
500    if is_background_child() {
501        return Ok(None);
502    }
503
504    let identity = match resolve_repo_identity(None) {
505        Ok(identity) => identity,
506        Err(_) => return Ok(None),
507    };
508    spawn_detached_worktree_watch_for_identity(&identity).map(Some)
509}
510
511fn resolve_target_identities(
512    target: &WorktreeWatchTargetOptions,
513) -> Result<Vec<RepoIdentity>, String> {
514    if target.repo.is_some() && target.parent.is_some() {
515        return Err("Pass either `--repo` or `--parent`, not both.".to_string());
516    }
517
518    if let Some(parent) = target.parent.as_deref() {
519        return discover_repo_identities_under(parent);
520    }
521
522    resolve_repo_identity(target.repo.as_deref()).map(|identity| vec![identity])
523}
524
525fn worktree_watch_status_payload(
526    identity: &RepoIdentity,
527    include_records: bool,
528    record_limit: usize,
529) -> Result<Value, String> {
530    let layout = prepare_spool_layout(identity)?;
531    let candidates = collect_sync_candidates(&layout)?;
532    let all_candidates = collect_spool_candidates(&layout, true)?;
533    let record_count: usize = candidates
534        .iter()
535        .map(|candidate| candidate.records.len())
536        .sum();
537    let all_record_count: usize = all_candidates
538        .iter()
539        .map(|candidate| candidate.records.len())
540        .sum();
541    let synced_file_count = all_candidates.len().saturating_sub(candidates.len());
542    let synced_record_count = all_record_count.saturating_sub(record_count);
543
544    let mut payload = json!({
545        "repoRoot": identity.root.display().to_string(),
546        "repositoryOwner": identity.owner.clone(),
547        "repositoryName": identity.name.clone(),
548        "branchName": identity.branch.clone(),
549        "spoolRoot": layout.root.display().to_string(),
550        "unsyncedFiles": candidates.len(),
551        "unsyncedRecords": record_count,
552        "syncedFiles": synced_file_count,
553        "syncedRecords": synced_record_count,
554        "localSpoolFiles": all_candidates.len(),
555        "localSpoolRecords": all_record_count,
556    });
557
558    if let Some(last_sync) = read_last_sync_log_entry(&layout.sync_log_file)? {
559        if let Some(object) = payload.as_object_mut() {
560            object.insert(
561                "lastSync".to_string(),
562                serde_json::to_value(last_sync)
563                    .map_err(|error| format!("Failed to render sync log payload: {error}"))?,
564            );
565        }
566    }
567
568    if include_records {
569        if let Some(object) = payload.as_object_mut() {
570            object.insert(
571                "records".to_string(),
572                collect_status_records(&candidates, record_limit)?,
573            );
574        }
575    }
576
577    Ok(payload)
578}
579
580fn collect_status_records(candidates: &[SyncCandidate], limit: usize) -> Result<Value, String> {
581    let available: usize = candidates
582        .iter()
583        .map(|candidate| candidate.records.len())
584        .sum();
585    let mut shown = 0usize;
586    let mut events = Vec::new();
587    let mut commits = Vec::new();
588
589    'candidates: for candidate in candidates {
590        for record in &candidate.records {
591            if shown >= limit {
592                break 'candidates;
593            }
594
595            match candidate.kind {
596                SyncFileKind::Events => {
597                    let event = serde_json::from_value::<WorktreeMutationEvent>(record.clone())
598                        .map_err(|error| {
599                            format!(
600                                "Failed to decode event from {}: {error}",
601                                candidate.path.display()
602                            )
603                        })?;
604                    events.push(json!({
605                        "sourceFile": candidate.path.display().to_string(),
606                        "occurredAt": event.occurred_at,
607                        "eventKind": event.event_kind,
608                        "primaryPath": event.primary_path,
609                        "oldPath": event.old_path,
610                        "newPath": event.new_path,
611                        "paths": event.paths,
612                        "addedLines": event.added_lines,
613                        "removedLines": event.removed_lines,
614                        "headSha": event.head_sha,
615                        "rawKind": event.raw_kind,
616                    }));
617                }
618                SyncFileKind::Commits => {
619                    let commit = serde_json::from_value::<WorktreeCommitEvent>(record.clone())
620                        .map_err(|error| {
621                            format!(
622                                "Failed to decode commit from {}: {error}",
623                                candidate.path.display()
624                            )
625                        })?;
626                    commits.push(json!({
627                        "sourceFile": candidate.path.display().to_string(),
628                        "occurredAt": commit.occurred_at,
629                        "headSha": commit.head_sha,
630                        "previousHeadSha": commit.previous_head_sha,
631                        "subject": commit.subject,
632                        "authorName": commit.author_name,
633                        "authorEmail": commit.author_email,
634                        "committedAt": commit.committed_at,
635                    }));
636                }
637            }
638
639            shown += 1;
640        }
641    }
642
643    Ok(json!({
644        "available": available,
645        "shown": shown,
646        "truncated": shown < available,
647        "events": events,
648        "commits": commits,
649    }))
650}
651
652fn print_status_records(payload: &Value) {
653    let Some(records) = payload.get("records") else {
654        return;
655    };
656    let shown = records.get("shown").and_then(Value::as_u64).unwrap_or(0);
657    if shown == 0 {
658        println!("records: none");
659        return;
660    }
661
662    println!("records:");
663    if let Some(events) = records.get("events").and_then(Value::as_array) {
664        if !events.is_empty() {
665            println!("  mutations:");
666            for event in events {
667                print_status_event_record(event);
668            }
669        }
670    }
671    if let Some(commits) = records.get("commits").and_then(Value::as_array) {
672        if !commits.is_empty() {
673            println!("  commits:");
674            for commit in commits {
675                print_status_commit_record(commit);
676            }
677        }
678    }
679    if records
680        .get("truncated")
681        .and_then(Value::as_bool)
682        .unwrap_or(false)
683    {
684        let available = records
685            .get("available")
686            .and_then(Value::as_u64)
687            .unwrap_or(shown);
688        println!("  showing {shown} of {available} record(s)");
689    }
690}
691
692fn print_status_stats(payload: &Value) {
693    let Some(stats) = payload.get("stats") else {
694        return;
695    };
696    println!("stats:");
697    println!(
698        "  coding time: {}",
699        format_duration(
700            stats
701                .get("estimatedCodingSeconds")
702                .and_then(Value::as_u64)
703                .unwrap_or(0)
704        )
705    );
706    println!(
707        "  events: {}",
708        stats
709            .get("totalEvents")
710            .and_then(Value::as_u64)
711            .unwrap_or(0)
712    );
713    println!(
714        "  commits: {}",
715        stats
716            .get("totalCommits")
717            .and_then(Value::as_u64)
718            .unwrap_or(0)
719    );
720    println!(
721        "  lines: +{} -{}",
722        stats.get("addedLines").and_then(Value::as_u64).unwrap_or(0),
723        stats
724            .get("removedLines")
725            .and_then(Value::as_u64)
726            .unwrap_or(0)
727    );
728    if let Some(spool_root) = value_str(stats, "spoolRoot") {
729        println!(
730            "  stored: {}",
731            Path::new(spool_root).join("stats.json").display()
732        );
733    }
734    print_stats_bucket_section(stats, "byFiletype", "  by filetype:", 10);
735    print_stats_bucket_section(stats, "byEventKind", "  by event kind:", 10);
736    print_stats_file_section(stats, 10);
737}
738
739fn print_last_sync(payload: &Value) {
740    let Some(last_sync) = payload.get("lastSync") else {
741        return;
742    };
743    let synced_at = value_str(last_sync, "syncedAt").unwrap_or("unknown-time");
744    let status = last_sync
745        .get("statusCode")
746        .and_then(Value::as_u64)
747        .unwrap_or(0);
748    let events = last_sync
749        .get("eventCount")
750        .and_then(Value::as_u64)
751        .unwrap_or(0);
752    let commits = last_sync
753        .get("commitCount")
754        .and_then(Value::as_u64)
755        .unwrap_or(0);
756    let files = last_sync
757        .get("spoolFileCount")
758        .and_then(Value::as_u64)
759        .unwrap_or(0);
760    let resync = last_sync
761        .get("resync")
762        .and_then(Value::as_bool)
763        .unwrap_or(false);
764    println!(
765        "last sync: {synced_at}, HTTP {status}, {events} event(s), {commits} commit(s), {files} file(s), resync={resync}"
766    );
767}
768
769fn print_stats_bucket_section(stats: &Value, key: &str, title: &str, limit: usize) {
770    let Some(items) = stats.get(key).and_then(Value::as_array) else {
771        return;
772    };
773    if items.is_empty() {
774        return;
775    }
776    println!("{title}");
777    for item in items.iter().take(limit) {
778        let name = value_str(item, "name").unwrap_or("unknown");
779        let events = item.get("eventCount").and_then(Value::as_u64).unwrap_or(0);
780        let seconds = item
781            .get("estimatedCodingSeconds")
782            .and_then(Value::as_u64)
783            .unwrap_or(0);
784        let added = item.get("addedLines").and_then(Value::as_u64).unwrap_or(0);
785        let removed = item
786            .get("removedLines")
787            .and_then(Value::as_u64)
788            .unwrap_or(0);
789        println!(
790            "    {name}: {}, {} event(s), +{} -{}",
791            format_duration(seconds),
792            events,
793            added,
794            removed
795        );
796    }
797}
798
799fn print_stats_file_section(stats: &Value, limit: usize) {
800    let Some(items) = stats.get("byFile").and_then(Value::as_array) else {
801        return;
802    };
803    if items.is_empty() {
804        return;
805    }
806    println!("  by file:");
807    for item in items.iter().take(limit) {
808        let path = value_str(item, "path").unwrap_or("(no path)");
809        let filetype = value_str(item, "filetype").unwrap_or("(none)");
810        let events = item.get("eventCount").and_then(Value::as_u64).unwrap_or(0);
811        let seconds = item
812            .get("estimatedCodingSeconds")
813            .and_then(Value::as_u64)
814            .unwrap_or(0);
815        let added = item.get("addedLines").and_then(Value::as_u64).unwrap_or(0);
816        let removed = item
817            .get("removedLines")
818            .and_then(Value::as_u64)
819            .unwrap_or(0);
820        println!(
821            "    {path} [{filetype}]: {}, {} event(s), +{} -{}",
822            format_duration(seconds),
823            events,
824            added,
825            removed
826        );
827    }
828}
829
830fn print_status_event_record(event: &Value) {
831    let occurred_at = value_str(event, "occurredAt").unwrap_or("unknown-time");
832    let kind = value_str(event, "eventKind").unwrap_or("event");
833    let primary_path = event_display_path(event);
834    let added = event.get("addedLines").and_then(Value::as_u64).unwrap_or(0);
835    let removed = event
836        .get("removedLines")
837        .and_then(Value::as_u64)
838        .unwrap_or(0);
839    println!("    {occurred_at} {kind} {primary_path} (+{added} -{removed})");
840
841    if let Some(head_sha) = value_str(event, "headSha") {
842        println!("      head: {}", short_sha(head_sha));
843    }
844    if let Some(paths) = event.get("paths").and_then(Value::as_array) {
845        if paths.len() > 1 {
846            let rendered = paths
847                .iter()
848                .filter_map(Value::as_str)
849                .collect::<Vec<_>>()
850                .join(", ");
851            println!("      paths: {rendered}");
852        }
853    }
854}
855
856fn print_status_commit_record(commit: &Value) {
857    let occurred_at = value_str(commit, "occurredAt").unwrap_or("unknown-time");
858    let head_sha = value_str(commit, "headSha")
859        .map(short_sha)
860        .unwrap_or_else(|| "unknown".to_string());
861    let subject = value_str(commit, "subject").unwrap_or("(no subject)");
862    println!("    {occurred_at} {head_sha} {subject}");
863
864    if let Some(author) = value_str(commit, "authorName") {
865        println!("      author: {author}");
866    }
867}
868
869fn event_display_path(event: &Value) -> String {
870    let old_path = value_str(event, "oldPath");
871    let new_path = value_str(event, "newPath");
872    if let (Some(old_path), Some(new_path)) = (old_path, new_path) {
873        return format!("{old_path} -> {new_path}");
874    }
875    value_str(event, "primaryPath")
876        .map(str::to_string)
877        .or_else(|| {
878            event
879                .get("paths")
880                .and_then(Value::as_array)
881                .and_then(|paths| paths.first())
882                .and_then(Value::as_str)
883                .map(str::to_string)
884        })
885        .unwrap_or_else(|| "(no path)".to_string())
886}
887
888fn value_str<'a>(value: &'a Value, key: &str) -> Option<&'a str> {
889    value.get(key).and_then(Value::as_str)
890}
891
892fn short_sha(value: &str) -> String {
893    value.chars().take(12).collect()
894}
895
896fn generate_and_store_stats(
897    identity: &RepoIdentity,
898    session_gap_minutes: u64,
899) -> Result<WorktreeWatchStatsSummary, String> {
900    let layout = prepare_spool_layout(identity)?;
901    let candidates = collect_spool_candidates(&layout, true)?;
902    let summary = build_stats_summary(identity, &layout, &candidates, session_gap_minutes)?;
903    let rendered = serde_json::to_string_pretty(&summary)
904        .map_err(|error| format!("Failed to serialize stats: {error}"))?;
905    fs::write(&layout.stats_file, format!("{rendered}\n")).map_err(|error| {
906        format!(
907            "Failed to write stats file {}: {error}",
908            layout.stats_file.display()
909        )
910    })?;
911    Ok(summary)
912}
913
914fn build_stats_summary(
915    identity: &RepoIdentity,
916    layout: &SpoolLayout,
917    candidates: &[SyncCandidate],
918    session_gap_minutes: u64,
919) -> Result<WorktreeWatchStatsSummary, String> {
920    let mut events = Vec::new();
921    let mut total_commits = 0u64;
922    for candidate in candidates {
923        match candidate.kind {
924            SyncFileKind::Events => {
925                for record in &candidate.records {
926                    events.push(
927                        serde_json::from_value::<WorktreeMutationEvent>(record.clone()).map_err(
928                            |error| {
929                                format!(
930                                    "Failed to decode event from {}: {error}",
931                                    candidate.path.display()
932                                )
933                            },
934                        )?,
935                    );
936                }
937            }
938            SyncFileKind::Commits => {
939                total_commits += candidate.records.len() as u64;
940            }
941        }
942    }
943    events.sort_by_key(|event| event.occurred_at);
944
945    let gap_seconds = session_gap_minutes.saturating_mul(60).max(60);
946    let mut previous_at: Option<DateTime<Utc>> = None;
947    let mut total_seconds = 0u64;
948    let mut total_added = 0u64;
949    let mut total_removed = 0u64;
950    let mut by_file: BTreeMap<String, WorktreeWatchFileStats> = BTreeMap::new();
951    let mut by_filetype: BTreeMap<String, WorktreeWatchBucketStats> = BTreeMap::new();
952    let mut by_event_kind: BTreeMap<String, WorktreeWatchBucketStats> = BTreeMap::new();
953
954    for event in &events {
955        let seconds = coding_seconds_for_event(previous_at, event.occurred_at, gap_seconds);
956        previous_at = Some(event.occurred_at);
957        let added = event.added_lines.unwrap_or(0);
958        let removed = event.removed_lines.unwrap_or(0);
959        let path = stats_event_path(event);
960        let filetype = filetype_for_path(&path);
961
962        total_seconds += seconds;
963        total_added += added;
964        total_removed += removed;
965        update_file_stats(&mut by_file, &path, &filetype, seconds, added, removed);
966        update_bucket_stats(&mut by_filetype, &filetype, seconds, added, removed);
967        update_bucket_stats(
968            &mut by_event_kind,
969            &event.event_kind,
970            seconds,
971            added,
972            removed,
973        );
974    }
975
976    Ok(WorktreeWatchStatsSummary {
977        generated_at: Utc::now(),
978        repository_owner: identity.owner.clone(),
979        repository_name: identity.name.clone(),
980        branch_name: identity.branch.clone(),
981        repo_root: identity.root.display().to_string(),
982        spool_root: layout.root.display().to_string(),
983        session_gap_minutes,
984        total_events: events.len() as u64,
985        total_commits,
986        first_event_at: events.first().map(|event| event.occurred_at),
987        last_event_at: events.last().map(|event| event.occurred_at),
988        estimated_coding_seconds: total_seconds,
989        added_lines: total_added,
990        removed_lines: total_removed,
991        by_file: sorted_file_stats(by_file),
992        by_filetype: sorted_bucket_stats(by_filetype),
993        by_event_kind: sorted_bucket_stats(by_event_kind),
994    })
995}
996
997fn coding_seconds_for_event(
998    previous_at: Option<DateTime<Utc>>,
999    occurred_at: DateTime<Utc>,
1000    gap_seconds: u64,
1001) -> u64 {
1002    let Some(previous_at) = previous_at else {
1003        return 60;
1004    };
1005    let delta = occurred_at.signed_duration_since(previous_at).num_seconds();
1006    if delta <= 0 {
1007        0
1008    } else {
1009        (delta as u64).min(gap_seconds)
1010    }
1011}
1012
1013fn update_file_stats(
1014    by_file: &mut BTreeMap<String, WorktreeWatchFileStats>,
1015    path: &str,
1016    filetype: &str,
1017    seconds: u64,
1018    added: u64,
1019    removed: u64,
1020) {
1021    let stats = by_file
1022        .entry(path.to_string())
1023        .or_insert_with(|| WorktreeWatchFileStats {
1024            path: path.to_string(),
1025            filetype: filetype.to_string(),
1026            ..WorktreeWatchFileStats::default()
1027        });
1028    stats.event_count += 1;
1029    stats.estimated_coding_seconds += seconds;
1030    stats.added_lines += added;
1031    stats.removed_lines += removed;
1032}
1033
1034fn update_bucket_stats(
1035    buckets: &mut BTreeMap<String, WorktreeWatchBucketStats>,
1036    name: &str,
1037    seconds: u64,
1038    added: u64,
1039    removed: u64,
1040) {
1041    let stats = buckets
1042        .entry(name.to_string())
1043        .or_insert_with(|| WorktreeWatchBucketStats {
1044            name: name.to_string(),
1045            ..WorktreeWatchBucketStats::default()
1046        });
1047    stats.event_count += 1;
1048    stats.estimated_coding_seconds += seconds;
1049    stats.added_lines += added;
1050    stats.removed_lines += removed;
1051}
1052
1053fn sorted_file_stats(
1054    by_file: BTreeMap<String, WorktreeWatchFileStats>,
1055) -> Vec<WorktreeWatchFileStats> {
1056    let mut stats = by_file.into_values().collect::<Vec<_>>();
1057    stats.sort_by(|left, right| {
1058        right
1059            .estimated_coding_seconds
1060            .cmp(&left.estimated_coding_seconds)
1061            .then_with(|| right.event_count.cmp(&left.event_count))
1062            .then_with(|| left.path.cmp(&right.path))
1063    });
1064    stats
1065}
1066
1067fn sorted_bucket_stats(
1068    buckets: BTreeMap<String, WorktreeWatchBucketStats>,
1069) -> Vec<WorktreeWatchBucketStats> {
1070    let mut stats = buckets.into_values().collect::<Vec<_>>();
1071    stats.sort_by(|left, right| {
1072        right
1073            .estimated_coding_seconds
1074            .cmp(&left.estimated_coding_seconds)
1075            .then_with(|| right.event_count.cmp(&left.event_count))
1076            .then_with(|| left.name.cmp(&right.name))
1077    });
1078    stats
1079}
1080
1081fn stats_event_path(event: &WorktreeMutationEvent) -> String {
1082    event
1083        .new_path
1084        .as_deref()
1085        .or(event.primary_path.as_deref())
1086        .or_else(|| event.paths.first().map(String::as_str))
1087        .unwrap_or("(no path)")
1088        .to_string()
1089}
1090
1091fn filetype_for_path(path: &str) -> String {
1092    Path::new(path)
1093        .extension()
1094        .and_then(|value| value.to_str())
1095        .map(|value| value.to_ascii_lowercase())
1096        .filter(|value| !value.is_empty())
1097        .unwrap_or_else(|| "(none)".to_string())
1098}
1099
1100fn format_duration(seconds: u64) -> String {
1101    let hours = seconds / 3600;
1102    let minutes = (seconds % 3600) / 60;
1103    let remaining_seconds = seconds % 60;
1104    if hours > 0 {
1105        format!("{hours}h {minutes}m")
1106    } else if minutes > 0 {
1107        format!("{minutes}m {remaining_seconds}s")
1108    } else {
1109        format!("{remaining_seconds}s")
1110    }
1111}
1112
1113fn watch_foreground(
1114    identity: RepoIdentity,
1115    layout: SpoolLayout,
1116    sync_interval_seconds: u64,
1117) -> impl std::future::Future<Output = Result<(), String>> {
1118    async move {
1119        let (tx, rx) = mpsc::channel();
1120        let mut watcher = RecommendedWatcher::new(
1121            move |result| {
1122                let _ = tx.send(result);
1123            },
1124            Config::default(),
1125        )
1126        .map_err(|error| format!("Failed to create filesystem watcher: {error}"))?;
1127
1128        watcher
1129            .watch(&identity.root, RecursiveMode::Recursive)
1130            .map_err(|error| {
1131                format!(
1132                    "Failed to watch repository root {}: {error}",
1133                    identity.root.display()
1134                )
1135            })?;
1136
1137        let mut last_head = identity.head_sha.clone();
1138        let mut last_commit_check = Instant::now();
1139        let mut last_sync = Instant::now();
1140        let mut event_dedupe = RecentEventDedupe {
1141            fingerprints: load_event_dedupe_fingerprints(&layout.event_dedupe_file)?,
1142        };
1143
1144        loop {
1145            match rx.recv_timeout(Duration::from_millis(750)) {
1146                Ok(Ok(event)) => {
1147                    if let Some(mutation) = mutation_event_from_notify(&identity, event) {
1148                        append_unique_mutation_event(&layout, &mutation, &mut event_dedupe)?;
1149                    }
1150                }
1151                Ok(Err(error)) => {
1152                    eprintln!("worktree watcher error: {error}");
1153                }
1154                Err(mpsc::RecvTimeoutError::Timeout) => {}
1155                Err(mpsc::RecvTimeoutError::Disconnected) => {
1156                    return Err("Filesystem watcher disconnected.".to_string());
1157                }
1158            }
1159
1160            if last_commit_check.elapsed() >= Duration::from_secs(3) {
1161                last_head = record_commit_snapshot(&identity, &layout, last_head.as_deref())?;
1162                last_commit_check = Instant::now();
1163            }
1164
1165            if sync_interval_seconds > 0
1166                && last_sync.elapsed() >= Duration::from_secs(sync_interval_seconds)
1167            {
1168                if let Err(error) = sync_spool_for_identity(&identity, &layout).await {
1169                    eprintln!("worktree mutation sync failed: {error}");
1170                }
1171                last_sync = Instant::now();
1172            }
1173        }
1174    }
1175}
1176
1177fn mutation_event_from_notify(
1178    identity: &RepoIdentity,
1179    event: Event,
1180) -> Option<WorktreeMutationEvent> {
1181    if event.kind.is_access() || event.kind.is_other() {
1182        return None;
1183    }
1184
1185    let paths: Vec<String> = event
1186        .paths
1187        .iter()
1188        .filter(|path| !is_ignored_path(&identity.root, path))
1189        .filter_map(|path| relative_slash_path(&identity.root, path))
1190        .collect();
1191
1192    if paths.is_empty() {
1193        return None;
1194    }
1195
1196    let primary_path = paths.first().cloned();
1197    let (old_path, new_path) = rename_paths(&event.kind, &paths);
1198    let line_counts = primary_path
1199        .as_deref()
1200        .and_then(|path| git_numstat_for_path(&identity.root, path).ok());
1201    let event_kind = classify_event_kind(&event.kind);
1202    let is_dir = event
1203        .paths
1204        .first()
1205        .and_then(|path| fs::metadata(path).ok())
1206        .map(|metadata| metadata.is_dir())
1207        .unwrap_or(false);
1208
1209    Some(WorktreeMutationEvent {
1210        id: Uuid::new_v4().to_string(),
1211        repo_owner: identity.owner.clone(),
1212        repo_name: identity.name.clone(),
1213        branch_name: identity.branch.clone(),
1214        repo_root: identity.root.display().to_string(),
1215        head_sha: current_head(&identity.root).ok().flatten(),
1216        event_kind: event_kind.to_string(),
1217        paths,
1218        primary_path,
1219        old_path,
1220        new_path,
1221        added_lines: line_counts.map(|counts| counts.0),
1222        removed_lines: line_counts.map(|counts| counts.1),
1223        file_created: matches!(
1224            event.kind,
1225            EventKind::Create(CreateKind::File | CreateKind::Any)
1226        ) && !is_dir,
1227        file_removed: matches!(
1228            event.kind,
1229            EventKind::Remove(RemoveKind::File | RemoveKind::Any)
1230        ) && !is_dir,
1231        folder_created: matches!(
1232            event.kind,
1233            EventKind::Create(CreateKind::Folder | CreateKind::Any)
1234        ) && is_dir,
1235        folder_removed: matches!(
1236            event.kind,
1237            EventKind::Remove(RemoveKind::Folder | RemoveKind::Any)
1238        ) && is_dir,
1239        renamed_or_moved: is_rename_or_move(&event.kind),
1240        raw_kind: format!("{:?}", event.kind),
1241        occurred_at: Utc::now(),
1242    })
1243}
1244
1245fn append_unique_mutation_event(
1246    layout: &SpoolLayout,
1247    mutation: &WorktreeMutationEvent,
1248    dedupe: &mut RecentEventDedupe,
1249) -> Result<bool, String> {
1250    let fingerprint = mutation_event_fingerprint(mutation);
1251    if dedupe.fingerprints.contains(&fingerprint)
1252        || event_dedupe_file_contains(&layout.event_dedupe_file, &fingerprint)?
1253    {
1254        dedupe.fingerprints.insert(fingerprint);
1255        return Ok(false);
1256    }
1257
1258    append_json_line(&layout.events_file, mutation)?;
1259    append_json_line(
1260        &layout.event_dedupe_file,
1261        &WorktreeWatchEventDedupeRecord {
1262            fingerprint: fingerprint.clone(),
1263            first_seen_at: Utc::now(),
1264            event_kind: mutation.event_kind.clone(),
1265            primary_path: mutation.primary_path.clone(),
1266        },
1267    )?;
1268    dedupe.fingerprints.insert(fingerprint);
1269    Ok(true)
1270}
1271
1272fn mutation_event_fingerprint(mutation: &WorktreeMutationEvent) -> String {
1273    let bucket = mutation.occurred_at.timestamp() / EVENT_DEDUPE_WINDOW_SECONDS;
1274    let payload = json!({
1275        "repoOwner": mutation.repo_owner,
1276        "repoName": mutation.repo_name,
1277        "branchName": mutation.branch_name,
1278        "repoRoot": mutation.repo_root,
1279        "headSha": mutation.head_sha,
1280        "eventKind": mutation.event_kind,
1281        "paths": mutation.paths,
1282        "primaryPath": mutation.primary_path,
1283        "oldPath": mutation.old_path,
1284        "newPath": mutation.new_path,
1285        "addedLines": mutation.added_lines,
1286        "removedLines": mutation.removed_lines,
1287        "fileCreated": mutation.file_created,
1288        "fileRemoved": mutation.file_removed,
1289        "folderCreated": mutation.folder_created,
1290        "folderRemoved": mutation.folder_removed,
1291        "renamedOrMoved": mutation.renamed_or_moved,
1292        "timeBucket": bucket,
1293    });
1294    let encoded = serde_json::to_vec(&payload).unwrap_or_default();
1295    let mut hasher = Sha256::new();
1296    hasher.update(encoded);
1297    format!("{:x}", hasher.finalize())
1298}
1299
1300fn load_event_dedupe_fingerprints(path: &Path) -> Result<BTreeSet<String>, String> {
1301    let mut fingerprints = BTreeSet::new();
1302    if !path.exists() {
1303        return Ok(fingerprints);
1304    }
1305
1306    for value in read_jsonl_values(path)? {
1307        if let Some(fingerprint) = value.get("fingerprint").and_then(Value::as_str) {
1308            fingerprints.insert(fingerprint.to_string());
1309        }
1310    }
1311    Ok(fingerprints)
1312}
1313
1314fn event_dedupe_file_contains(path: &Path, fingerprint: &str) -> Result<bool, String> {
1315    if !path.exists() {
1316        return Ok(false);
1317    }
1318    for value in read_jsonl_values(path)? {
1319        if value
1320            .get("fingerprint")
1321            .and_then(Value::as_str)
1322            .map(|value| value == fingerprint)
1323            .unwrap_or(false)
1324        {
1325            return Ok(true);
1326        }
1327    }
1328    Ok(false)
1329}
1330
1331fn classify_event_kind(kind: &EventKind) -> &'static str {
1332    match kind {
1333        EventKind::Create(CreateKind::File) => "file_create",
1334        EventKind::Create(CreateKind::Folder) => "folder_create",
1335        EventKind::Create(_) => "create",
1336        EventKind::Remove(RemoveKind::File) => "file_remove",
1337        EventKind::Remove(RemoveKind::Folder) => "folder_remove",
1338        EventKind::Remove(_) => "remove",
1339        EventKind::Modify(ModifyKind::Name(
1340            RenameMode::From | RenameMode::To | RenameMode::Both,
1341        )) => "rename_or_move",
1342        EventKind::Modify(ModifyKind::Data(_)) => "file_modify",
1343        EventKind::Modify(ModifyKind::Metadata(_)) => "metadata_modify",
1344        EventKind::Modify(_) => "modify",
1345        _ => "other",
1346    }
1347}
1348
1349fn is_rename_or_move(kind: &EventKind) -> bool {
1350    matches!(
1351        kind,
1352        EventKind::Modify(ModifyKind::Name(
1353            RenameMode::From | RenameMode::To | RenameMode::Both
1354        ))
1355    )
1356}
1357
1358fn rename_paths(kind: &EventKind, paths: &[String]) -> (Option<String>, Option<String>) {
1359    if !is_rename_or_move(kind) {
1360        return (None, None);
1361    }
1362
1363    match paths {
1364        [old_path, new_path, ..] => (Some(old_path.clone()), Some(new_path.clone())),
1365        [path] => match kind {
1366            EventKind::Modify(ModifyKind::Name(RenameMode::From)) => (Some(path.clone()), None),
1367            EventKind::Modify(ModifyKind::Name(RenameMode::To)) => (None, Some(path.clone())),
1368            _ => (None, Some(path.clone())),
1369        },
1370        _ => (None, None),
1371    }
1372}
1373
1374fn record_commit_snapshot(
1375    identity: &RepoIdentity,
1376    layout: &SpoolLayout,
1377    previous_head: Option<&str>,
1378) -> Result<Option<String>, String> {
1379    let head = current_head(&identity.root)?;
1380    let Some(head_sha) = head else {
1381        return Ok(None);
1382    };
1383
1384    if previous_head == Some(head_sha.as_str()) {
1385        return Ok(Some(head_sha));
1386    }
1387
1388    let commit = WorktreeCommitEvent {
1389        id: Uuid::new_v4().to_string(),
1390        repo_owner: identity.owner.clone(),
1391        repo_name: identity.name.clone(),
1392        branch_name: identity.branch.clone(),
1393        repo_root: identity.root.display().to_string(),
1394        previous_head_sha: previous_head.map(str::to_string),
1395        head_sha: head_sha.clone(),
1396        subject: git_output(&identity.root, &["log", "-1", "--pretty=%s"]).ok(),
1397        author_name: git_output(&identity.root, &["log", "-1", "--pretty=%an"]).ok(),
1398        author_email: git_output(&identity.root, &["log", "-1", "--pretty=%ae"]).ok(),
1399        committed_at: git_output(&identity.root, &["log", "-1", "--pretty=%cI"]).ok(),
1400        occurred_at: Utc::now(),
1401    };
1402    append_json_line(&layout.commits_file, &commit)?;
1403    Ok(Some(head_sha))
1404}
1405
1406async fn sync_spool_for_identity(
1407    identity: &RepoIdentity,
1408    layout: &SpoolLayout,
1409) -> Result<(), String> {
1410    let candidates = collect_sync_candidates(layout)?;
1411    if candidates.is_empty() {
1412        return Ok(());
1413    }
1414
1415    sync_candidates(identity, candidates, false).await
1416}
1417
1418async fn sync_candidates(
1419    identity: &RepoIdentity,
1420    candidates: Vec<SyncCandidate>,
1421    resync: bool,
1422) -> Result<(), String> {
1423    let mut events = Vec::new();
1424    let mut commits = Vec::new();
1425    for candidate in &candidates {
1426        match candidate.kind {
1427            SyncFileKind::Events => {
1428                for record in &candidate.records {
1429                    events.push(
1430                        serde_json::from_value::<WorktreeMutationEvent>(record.clone()).map_err(
1431                            |error| {
1432                                format!(
1433                                    "Failed to decode event from {}: {error}",
1434                                    candidate.path.display()
1435                                )
1436                            },
1437                        )?,
1438                    );
1439                }
1440            }
1441            SyncFileKind::Commits => {
1442                for record in &candidate.records {
1443                    commits.push(
1444                        serde_json::from_value::<WorktreeCommitEvent>(record.clone()).map_err(
1445                            |error| {
1446                                format!(
1447                                    "Failed to decode commit from {}: {error}",
1448                                    candidate.path.display()
1449                                )
1450                            },
1451                        )?,
1452                    );
1453                }
1454            }
1455        }
1456    }
1457
1458    let token = resolve_cli_access_token()?;
1459    let device = resolve_device_identity()?;
1460    let client = cli_request_client()?;
1461    let api = ApiConfig::from_env();
1462    let endpoint = api.cli_worktree_mutations_endpoint();
1463    let event_count = events.len();
1464    let commit_count = commits.len();
1465    let payload = WorktreeMutationIngestPayload {
1466        device: WorktreeMutationDevicePayload {
1467            hardware_id: device.hardware_id,
1468            device_name: current_hostname(),
1469            hostname: current_hostname(),
1470            platform: std::env::consts::OS.to_string(),
1471        },
1472        repository: WorktreeMutationRepositoryPayload {
1473            owner: identity.owner.clone(),
1474            name: identity.name.clone(),
1475            branch_name: identity.branch.clone(),
1476            repo_root: identity.root.display().to_string(),
1477        },
1478        events,
1479        commits,
1480    };
1481
1482    let response = client
1483        .post(endpoint.clone())
1484        .bearer_auth(token)
1485        .json(&payload)
1486        .send()
1487        .await
1488        .map_err(|error| format!("Failed to upload worktree mutations: {error}"))?;
1489
1490    if response.status() == StatusCode::UNAUTHORIZED {
1491        return Err(
1492            "Your stored CLI session is no longer valid. Run `xbp login` again.".to_string(),
1493        );
1494    }
1495
1496    if !response.status().is_success() {
1497        let status = response.status();
1498        let body = response.text().await.unwrap_or_default();
1499        return Err(format!(
1500            "Worktree mutation upload failed with {status}: {body}"
1501        ));
1502    }
1503
1504    let status_code = response.status().as_u16();
1505    let layout = prepare_spool_layout(identity)?;
1506    append_sync_log(
1507        identity,
1508        &layout,
1509        &endpoint,
1510        status_code,
1511        event_count,
1512        commit_count,
1513        &candidates,
1514        resync,
1515    )?;
1516
1517    for candidate in candidates {
1518        if !is_synced_spool_path(&candidate.path) {
1519            mark_synced(&candidate.path)?;
1520        }
1521    }
1522
1523    Ok(())
1524}
1525
1526fn collect_sync_candidates(layout: &SpoolLayout) -> Result<Vec<SyncCandidate>, String> {
1527    collect_spool_candidates(layout, false)
1528}
1529
1530fn print_no_sync_candidates_message(
1531    target: &WorktreeWatchTargetOptions,
1532    resync: bool,
1533) -> Result<(), String> {
1534    if resync {
1535        println!("No local worktree mutation JSONL files found.");
1536        return Ok(());
1537    }
1538
1539    let identities = resolve_target_identities(target)?;
1540    let mut synced_files = 0usize;
1541    let mut synced_records = 0usize;
1542    for identity in identities {
1543        let layout = prepare_spool_layout(&identity)?;
1544        let all_candidates = collect_spool_candidates(&layout, true)?;
1545        let unsynced_candidates = collect_sync_candidates(&layout)?;
1546        synced_files += all_candidates
1547            .len()
1548            .saturating_sub(unsynced_candidates.len());
1549        let all_records = all_candidates
1550            .iter()
1551            .map(|candidate| candidate.records.len())
1552            .sum::<usize>();
1553        let unsynced_records = unsynced_candidates
1554            .iter()
1555            .map(|candidate| candidate.records.len())
1556            .sum::<usize>();
1557        synced_records += all_records.saturating_sub(unsynced_records);
1558    }
1559
1560    if synced_files > 0 {
1561        println!(
1562            "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."
1563        );
1564    } else {
1565        println!("No unsynced worktree mutation files found.");
1566    }
1567    Ok(())
1568}
1569
1570fn collect_spool_candidates(
1571    layout: &SpoolLayout,
1572    include_synced: bool,
1573) -> Result<Vec<SyncCandidate>, String> {
1574    let mut candidates = Vec::new();
1575    if !layout.root.exists() {
1576        return Ok(candidates);
1577    }
1578
1579    for entry in fs::read_dir(&layout.root).map_err(|error| {
1580        format!(
1581            "Failed to read spool directory {}: {error}",
1582            layout.root.display()
1583        )
1584    })? {
1585        let path = entry
1586            .map_err(|error| format!("Failed to read spool entry: {error}"))?
1587            .path();
1588        let Some(file_name) = path.file_name().and_then(|value| value.to_str()) else {
1589            continue;
1590        };
1591        if !file_name.ends_with(".jsonl") || (!include_synced && file_name.contains(".synced.")) {
1592            continue;
1593        }
1594
1595        let kind = if file_name.starts_with("events-") {
1596            SyncFileKind::Events
1597        } else if file_name.starts_with("commits-") {
1598            SyncFileKind::Commits
1599        } else {
1600            continue;
1601        };
1602        let records = read_jsonl_values(&path)?;
1603        if records.is_empty() {
1604            continue;
1605        }
1606        candidates.push(SyncCandidate {
1607            path,
1608            kind,
1609            records,
1610        });
1611    }
1612
1613    Ok(candidates)
1614}
1615
1616fn append_sync_log(
1617    identity: &RepoIdentity,
1618    layout: &SpoolLayout,
1619    endpoint: &str,
1620    status_code: u16,
1621    event_count: usize,
1622    commit_count: usize,
1623    candidates: &[SyncCandidate],
1624    resync: bool,
1625) -> Result<(), String> {
1626    let entry = WorktreeWatchSyncLogEntry {
1627        id: Uuid::new_v4().to_string(),
1628        synced_at: Utc::now(),
1629        endpoint: endpoint.to_string(),
1630        repository_owner: identity.owner.clone(),
1631        repository_name: identity.name.clone(),
1632        branch_name: identity.branch.clone(),
1633        repo_root: identity.root.display().to_string(),
1634        resync,
1635        status_code,
1636        event_count,
1637        commit_count,
1638        spool_file_count: candidates.len(),
1639        spool_files: candidates
1640            .iter()
1641            .map(|candidate| candidate.path.display().to_string())
1642            .collect(),
1643    };
1644    append_json_line(&layout.sync_log_file, &entry)
1645}
1646
1647fn read_last_sync_log_entry(path: &Path) -> Result<Option<WorktreeWatchSyncLogEntry>, String> {
1648    if !path.exists() {
1649        return Ok(None);
1650    }
1651    let values = read_jsonl_values(path)?;
1652    let Some(value) = values.into_iter().last() else {
1653        return Ok(None);
1654    };
1655    serde_json::from_value(value)
1656        .map(Some)
1657        .map_err(|error| format!("Failed to decode sync log {}: {error}", path.display()))
1658}
1659
1660fn is_synced_spool_path(path: &Path) -> bool {
1661    path.file_name()
1662        .and_then(|value| value.to_str())
1663        .map(|value| value.contains(".synced."))
1664        .unwrap_or(false)
1665}
1666
1667fn mark_synced(path: &Path) -> Result<(), String> {
1668    let Some(file_name) = path.file_name().and_then(|value| value.to_str()) else {
1669        return Ok(());
1670    };
1671    let synced_name = file_name.replace(
1672        ".jsonl",
1673        &format!(".synced.{}.jsonl", Utc::now().timestamp()),
1674    );
1675    let synced_path = path.with_file_name(synced_name);
1676    fs::rename(path, &synced_path).map_err(|error| {
1677        format!(
1678            "Failed to mark spool file {} as synced: {error}",
1679            path.display()
1680        )
1681    })
1682}
1683
1684fn read_jsonl_values(path: &Path) -> Result<Vec<Value>, String> {
1685    let file = File::open(path)
1686        .map_err(|error| format!("Failed to open spool file {}: {error}", path.display()))?;
1687    let reader = BufReader::new(file);
1688    let mut values = Vec::new();
1689    for line in reader.lines() {
1690        let line =
1691            line.map_err(|error| format!("Failed to read spool file {}: {error}", path.display()))?;
1692        if line.trim().is_empty() {
1693            continue;
1694        }
1695        values.push(serde_json::from_str(&line).map_err(|error| {
1696            format!(
1697                "Failed to parse JSONL record in {}: {error}",
1698                path.display()
1699            )
1700        })?);
1701    }
1702    Ok(values)
1703}
1704
1705fn spawn_detached_worktree_watch(repo: Option<&Path>) -> Result<PathBuf, String> {
1706    let identity = resolve_repo_identity(repo)?;
1707    spawn_detached_worktree_watch_for_identity(&identity)
1708}
1709
1710fn spawn_detached_worktree_watch_for_identity(identity: &RepoIdentity) -> Result<PathBuf, String> {
1711    let layout = prepare_spool_layout(identity)?;
1712    replace_existing_watcher(identity, &layout)?;
1713    let executable = std::env::current_exe()
1714        .map_err(|error| format!("Failed to resolve current XBP executable: {error}"))?;
1715    let mut command = Command::new(&executable);
1716    command
1717        .arg("worktree-watch")
1718        .arg("start")
1719        .arg("--repo")
1720        .arg(&identity.root)
1721        .env(BACKGROUND_CHILD_ENV, "1")
1722        .current_dir(&identity.root)
1723        .stdin(Stdio::null())
1724        .stdout(Stdio::null())
1725        .stderr(Stdio::null());
1726
1727    #[cfg(windows)]
1728    {
1729        use std::os::windows::process::CommandExt;
1730        command.creation_flags(CREATE_NO_WINDOW);
1731    }
1732
1733    let child = command
1734        .spawn()
1735        .map_err(|error| format!("Failed to spawn background worktree watcher: {error}"))?;
1736    write_watcher_state(identity, &layout, child.id(), &executable)?;
1737    Ok(identity.root.clone())
1738}
1739
1740fn replace_existing_watcher(identity: &RepoIdentity, layout: &SpoolLayout) -> Result<(), String> {
1741    let Some(state) = read_watcher_state(&layout.watcher_state_file)? else {
1742        return Ok(());
1743    };
1744
1745    if !same_repo_watcher_state(identity, &state) {
1746        return Ok(());
1747    }
1748
1749    if is_matching_watcher_process(&state) {
1750        stop_process(&state)?;
1751        println!(
1752            "Replaced existing background worktree watcher {} for {}",
1753            state.pid,
1754            identity.root.display()
1755        );
1756    }
1757
1758    let _ = fs::remove_file(&layout.watcher_state_file);
1759    Ok(())
1760}
1761
1762fn read_watcher_state(path: &Path) -> Result<Option<WorktreeWatcherState>, String> {
1763    if !path.exists() {
1764        return Ok(None);
1765    }
1766    let raw = fs::read_to_string(path)
1767        .map_err(|error| format!("Failed to read watcher state {}: {error}", path.display()))?;
1768    serde_json::from_str(&raw)
1769        .map(Some)
1770        .map_err(|error| format!("Failed to parse watcher state {}: {error}", path.display()))
1771}
1772
1773fn write_watcher_state(
1774    identity: &RepoIdentity,
1775    layout: &SpoolLayout,
1776    pid: u32,
1777    executable: &Path,
1778) -> Result<(), String> {
1779    let state = WorktreeWatcherState {
1780        pid,
1781        repo_root: identity.root.display().to_string(),
1782        repository_owner: identity.owner.clone(),
1783        repository_name: identity.name.clone(),
1784        branch_name: identity.branch.clone(),
1785        executable: executable.display().to_string(),
1786        started_at: Utc::now(),
1787    };
1788    let rendered = serde_json::to_string_pretty(&state)
1789        .map_err(|error| format!("Failed to serialize watcher state: {error}"))?;
1790    fs::write(&layout.watcher_state_file, format!("{rendered}\n")).map_err(|error| {
1791        format!(
1792            "Failed to write watcher state {}: {error}",
1793            layout.watcher_state_file.display()
1794        )
1795    })
1796}
1797
1798fn same_repo_watcher_state(identity: &RepoIdentity, state: &WorktreeWatcherState) -> bool {
1799    path_identity_key(&identity.root) == path_identity_key(Path::new(&state.repo_root))
1800        && identity.owner == state.repository_owner
1801        && identity.name == state.repository_name
1802        && identity.branch == state.branch_name
1803}
1804
1805fn is_matching_watcher_process(state: &WorktreeWatcherState) -> bool {
1806    if state.pid == std::process::id() {
1807        return false;
1808    }
1809    let system = System::new_all();
1810    let Some(process) = system.process(Pid::from_u32(state.pid)) else {
1811        return false;
1812    };
1813    let command_line = process
1814        .cmd()
1815        .iter()
1816        .map(|part| part.to_string_lossy())
1817        .collect::<Vec<_>>()
1818        .join(" ");
1819    command_line.contains("worktree-watch")
1820        && command_line.contains("start")
1821        && command_line.contains(&state.repo_root)
1822}
1823
1824fn stop_process(state: &WorktreeWatcherState) -> Result<(), String> {
1825    let system = System::new_all();
1826    let Some(process) = system.process(Pid::from_u32(state.pid)) else {
1827        return Ok(());
1828    };
1829    if !is_matching_watcher_process(state) {
1830        return Ok(());
1831    }
1832    if process.kill() {
1833        Ok(())
1834    } else {
1835        Err(format!(
1836            "Failed to stop existing watcher process {}",
1837            state.pid
1838        ))
1839    }
1840}
1841
1842fn discover_repo_identities_under(parent: &Path) -> Result<Vec<RepoIdentity>, String> {
1843    let parent = normalize_windows_verbatim_path(
1844        fs::canonicalize(parent).unwrap_or_else(|_| parent.to_path_buf()),
1845    );
1846    if !parent.is_dir() {
1847        return Err(format!(
1848            "Parent folder {} does not exist or is not a directory.",
1849            parent.display()
1850        ));
1851    }
1852
1853    let mut identities = Vec::new();
1854    let mut seen = BTreeSet::new();
1855    discover_repo_identities_in_dir(&parent, &parent, &mut seen, &mut identities)?;
1856    identities.sort_by(|left, right| left.root.cmp(&right.root));
1857    Ok(identities)
1858}
1859
1860fn discover_repo_identities_in_dir(
1861    parent: &Path,
1862    dir: &Path,
1863    seen: &mut BTreeSet<String>,
1864    identities: &mut Vec<RepoIdentity>,
1865) -> Result<(), String> {
1866    if dir != parent && is_skipped_discovery_dir(dir) {
1867        return Ok(());
1868    }
1869
1870    if dir.join(".git").exists() {
1871        if let Ok(identity) = resolve_repo_identity(Some(dir)) {
1872            let key = path_identity_key(&identity.root);
1873            if seen.insert(key) {
1874                identities.push(identity);
1875            }
1876            return Ok(());
1877        }
1878    }
1879
1880    let entries = fs::read_dir(dir)
1881        .map_err(|error| format!("Failed to read folder {}: {error}", dir.display()))?;
1882    for entry in entries {
1883        let entry = entry.map_err(|error| {
1884            format!(
1885                "Failed to read folder entry under {}: {error}",
1886                dir.display()
1887            )
1888        })?;
1889        let file_type = entry
1890            .file_type()
1891            .map_err(|error| format!("Failed to inspect {}: {error}", entry.path().display()))?;
1892        if file_type.is_dir() {
1893            discover_repo_identities_in_dir(parent, &entry.path(), seen, identities)?;
1894        }
1895    }
1896
1897    Ok(())
1898}
1899
1900fn resolve_repo_identity(repo: Option<&Path>) -> Result<RepoIdentity, String> {
1901    let start = match repo {
1902        Some(path) => path.to_path_buf(),
1903        None => std::env::current_dir()
1904            .map_err(|error| format!("Failed to resolve current directory: {error}"))?,
1905    };
1906    let root_raw = git_output(&start, &["rev-parse", "--show-toplevel"])?;
1907    let root = normalize_windows_verbatim_path(
1908        fs::canonicalize(root_raw.trim()).unwrap_or_else(|_| PathBuf::from(root_raw.trim())),
1909    );
1910    let branch = git_output(&root, &["rev-parse", "--abbrev-ref", "HEAD"])
1911        .unwrap_or_else(|_| "unknown".to_string());
1912    let remote = git_output(&root, &["remote", "get-url", DEFAULT_REMOTE]).unwrap_or_default();
1913    let (owner, name) = parse_remote_owner_repo(&remote).unwrap_or_else(|| {
1914        let name = root
1915            .file_name()
1916            .and_then(|value| value.to_str())
1917            .unwrap_or("repository")
1918            .to_string();
1919        ("unknown".to_string(), name)
1920    });
1921    let head_sha = current_head(&root)?;
1922
1923    Ok(RepoIdentity {
1924        owner,
1925        name,
1926        branch: sanitize_path_component(branch.trim()),
1927        root,
1928        head_sha,
1929    })
1930}
1931
1932fn prepare_spool_layout(identity: &RepoIdentity) -> Result<SpoolLayout, String> {
1933    let home = dirs::home_dir().ok_or_else(|| "Failed to resolve home directory.".to_string())?;
1934    let run_id = Uuid::new_v4().to_string();
1935    let root = home
1936        .join(".xbp")
1937        .join("mutations")
1938        .join(sanitize_path_component(&identity.owner))
1939        .join(sanitize_path_component(&identity.name))
1940        .join(sanitize_path_component(&identity.branch));
1941    fs::create_dir_all(&root).map_err(|error| {
1942        format!(
1943            "Failed to create worktree mutation spool {}: {error}",
1944            root.display()
1945        )
1946    })?;
1947    Ok(SpoolLayout {
1948        events_file: root.join(format!("events-{run_id}.jsonl")),
1949        commits_file: root.join(format!("commits-{run_id}.jsonl")),
1950        stats_file: root.join("stats.json"),
1951        watcher_state_file: root.join("watcher-state.json"),
1952        sync_log_file: root.join("sync-log.jsonl"),
1953        event_dedupe_file: root.join("event-dedupe.jsonl"),
1954        root,
1955    })
1956}
1957
1958fn append_json_line<T: Serialize>(path: &Path, value: &T) -> Result<(), String> {
1959    let mut file = OpenOptions::new()
1960        .create(true)
1961        .append(true)
1962        .open(path)
1963        .map_err(|error| format!("Failed to open spool file {}: {error}", path.display()))?;
1964    let line = serde_json::to_string(value)
1965        .map_err(|error| format!("Failed to serialize worktree event: {error}"))?;
1966    writeln!(file, "{line}")
1967        .map_err(|error| format!("Failed to append spool file {}: {error}", path.display()))
1968}
1969
1970fn git_numstat_for_path(repo_root: &Path, relative_path: &str) -> Result<(u64, u64), String> {
1971    let output = Command::new("git")
1972        .args(["diff", "--numstat", "--"])
1973        .arg(relative_path)
1974        .current_dir(repo_root)
1975        .output()
1976        .map_err(|error| format!("Failed to run git diff --numstat: {error}"))?;
1977    if !output.status.success() {
1978        return Ok((0, 0));
1979    }
1980    let stdout = String::from_utf8_lossy(&output.stdout);
1981    let mut added = 0;
1982    let mut removed = 0;
1983    for line in stdout.lines() {
1984        let mut parts = line.split_whitespace();
1985        added += parse_numstat_count(parts.next());
1986        removed += parse_numstat_count(parts.next());
1987    }
1988    Ok((added, removed))
1989}
1990
1991fn parse_numstat_count(value: Option<&str>) -> u64 {
1992    value.and_then(|raw| raw.parse::<u64>().ok()).unwrap_or(0)
1993}
1994
1995fn current_head(repo_root: &Path) -> Result<Option<String>, String> {
1996    match git_output(repo_root, &["rev-parse", "HEAD"]) {
1997        Ok(value) => Ok(Some(value)),
1998        Err(error)
1999            if error.contains("unknown revision") || error.contains("ambiguous argument") =>
2000        {
2001            Ok(None)
2002        }
2003        Err(error) => Err(error),
2004    }
2005}
2006
2007fn git_output(repo_root: &Path, args: &[&str]) -> Result<String, String> {
2008    let output = Command::new("git")
2009        .args(args)
2010        .current_dir(repo_root)
2011        .output()
2012        .map_err(|error| format!("Failed to run git {}: {error}", args.join(" ")))?;
2013    if !output.status.success() {
2014        return Err(String::from_utf8_lossy(&output.stderr).trim().to_string());
2015    }
2016    Ok(String::from_utf8_lossy(&output.stdout).trim().to_string())
2017}
2018
2019fn parse_remote_owner_repo(remote: &str) -> Option<(String, String)> {
2020    let trimmed = remote.trim().trim_end_matches(".git");
2021    if trimmed.is_empty() {
2022        return None;
2023    }
2024
2025    let path_part = if !trimmed.contains("://") {
2026        if let Some((_, path)) = trimmed.rsplit_once(':') {
2027            path
2028        } else {
2029            trimmed
2030        }
2031    } else {
2032        trimmed
2033            .trim_start_matches("https://")
2034            .trim_start_matches("http://")
2035            .trim_start_matches("ssh://")
2036            .split_once('/')
2037            .map(|(_, path)| path)
2038            .unwrap_or(trimmed)
2039    };
2040    let mut parts = path_part.rsplitn(2, '/');
2041    let name = parts.next()?.trim();
2042    let owner = parts.next()?.trim();
2043    if owner.is_empty() || name.is_empty() {
2044        return None;
2045    }
2046    Some((
2047        sanitize_path_component(owner),
2048        sanitize_path_component(name.trim_end_matches(".git")),
2049    ))
2050}
2051
2052fn sanitize_path_component(value: &str) -> String {
2053    let sanitized: String = value
2054        .chars()
2055        .map(|ch| {
2056            if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
2057                ch
2058            } else {
2059                '-'
2060            }
2061        })
2062        .collect();
2063    sanitized
2064        .trim_matches('-')
2065        .chars()
2066        .take(160)
2067        .collect::<String>()
2068}
2069
2070fn normalize_windows_verbatim_path(path: PathBuf) -> PathBuf {
2071    PathBuf::from(strip_windows_verbatim_prefix(&path.to_string_lossy()))
2072}
2073
2074fn strip_windows_verbatim_prefix(input: &str) -> &str {
2075    input.strip_prefix(r"\\?\").unwrap_or(input)
2076}
2077
2078fn path_identity_key(path: &Path) -> String {
2079    let value = path.to_string_lossy().to_string();
2080    if cfg!(windows) {
2081        value.to_ascii_lowercase()
2082    } else {
2083        value
2084    }
2085}
2086
2087fn is_skipped_discovery_dir(path: &Path) -> bool {
2088    let Some(name) = path.file_name().and_then(|value| value.to_str()) else {
2089        return false;
2090    };
2091    matches!(
2092        name,
2093        ".git"
2094            | ".hg"
2095            | ".svn"
2096            | ".next"
2097            | ".turbo"
2098            | ".vercel"
2099            | "node_modules"
2100            | "target"
2101            | "dist"
2102            | "build"
2103    )
2104}
2105
2106fn relative_slash_path(root: &Path, path: &Path) -> Option<String> {
2107    let relative = path.strip_prefix(root).ok().unwrap_or(path);
2108    Some(relative.to_string_lossy().replace('\\', "/"))
2109}
2110
2111fn is_ignored_path(root: &Path, path: &Path) -> bool {
2112    let Some(relative) = relative_slash_path(root, path) else {
2113        return false;
2114    };
2115    relative == ".git"
2116        || relative.starts_with(".git/")
2117        || relative == "target"
2118        || relative.starts_with("target/")
2119}
2120
2121fn current_hostname() -> Option<String> {
2122    std::env::var("COMPUTERNAME")
2123        .or_else(|_| std::env::var("HOSTNAME"))
2124        .ok()
2125        .map(|value| value.trim().to_string())
2126        .filter(|value| !value.is_empty())
2127}
2128
2129fn is_background_child() -> bool {
2130    std::env::var(BACKGROUND_CHILD_ENV)
2131        .ok()
2132        .map(|value| value == "1")
2133        .unwrap_or(false)
2134}
2135
2136#[allow(dead_code)]
2137fn content_sha256(path: &Path) -> Option<String> {
2138    let bytes = fs::read(path).ok()?;
2139    let mut hasher = Sha256::new();
2140    hasher.update(bytes);
2141    Some(format!("{:x}", hasher.finalize()))
2142}
2143
2144#[cfg(test)]
2145mod tests {
2146    use super::*;
2147
2148    #[test]
2149    fn parses_https_and_ssh_remote_urls() {
2150        assert_eq!(
2151            parse_remote_owner_repo("https://github.com/xylex-group/xbp.git"),
2152            Some(("xylex-group".to_string(), "xbp".to_string()))
2153        );
2154        assert_eq!(
2155            parse_remote_owner_repo("git@github.com:xylex-group/xbp.git"),
2156            Some(("xylex-group".to_string(), "xbp".to_string()))
2157        );
2158    }
2159
2160    #[test]
2161    fn sanitizes_branch_for_storage_path() {
2162        assert_eq!(
2163            sanitize_path_component("feature/worktree watcher"),
2164            "feature-worktree-watcher"
2165        );
2166    }
2167
2168    #[test]
2169    fn normalizes_windows_verbatim_repo_roots() {
2170        assert_eq!(
2171            normalize_windows_verbatim_path(PathBuf::from(
2172                r"\\?\C:\Users\floris\Documents\GitHub\mollie-api-rust"
2173            )),
2174            PathBuf::from(r"C:\Users\floris\Documents\GitHub\mollie-api-rust")
2175        );
2176    }
2177
2178    #[test]
2179    fn skips_heavy_discovery_directories() {
2180        assert!(is_skipped_discovery_dir(Path::new("node_modules")));
2181        assert!(is_skipped_discovery_dir(Path::new("target")));
2182        assert!(is_skipped_discovery_dir(Path::new(".next")));
2183        assert!(!is_skipped_discovery_dir(Path::new("mollie-api-rust")));
2184    }
2185
2186    #[test]
2187    fn builds_coding_stats_by_file_and_filetype() {
2188        let identity = RepoIdentity {
2189            owner: "xylex-group".to_string(),
2190            name: "xbp".to_string(),
2191            branch: "main".to_string(),
2192            root: PathBuf::from(r"C:\Users\floris\Documents\GitHub\xbp"),
2193            head_sha: None,
2194        };
2195        let layout = SpoolLayout {
2196            root: PathBuf::from(r"C:\Users\floris\.xbp\mutations\xylex-group\xbp\main"),
2197            events_file: PathBuf::from("events.jsonl"),
2198            commits_file: PathBuf::from("commits.jsonl"),
2199            stats_file: PathBuf::from("stats.json"),
2200            watcher_state_file: PathBuf::from("watcher-state.json"),
2201            sync_log_file: PathBuf::from("sync-log.jsonl"),
2202            event_dedupe_file: PathBuf::from("event-dedupe.jsonl"),
2203        };
2204        let candidates = vec![SyncCandidate {
2205            path: PathBuf::from("events.jsonl"),
2206            kind: SyncFileKind::Events,
2207            records: vec![
2208                stats_test_event("2026-07-08T10:00:00Z", "src/main.rs", 4, 1),
2209                stats_test_event("2026-07-08T10:05:00Z", "src/main.rs", 2, 0),
2210                stats_test_event("2026-07-08T11:00:00Z", "README.md", 1, 3),
2211            ],
2212        }];
2213
2214        let summary = build_stats_summary(&identity, &layout, &candidates, 15).unwrap();
2215
2216        assert_eq!(summary.total_events, 3);
2217        assert_eq!(summary.estimated_coding_seconds, 60 + 300 + 900);
2218        assert_eq!(summary.added_lines, 7);
2219        assert_eq!(summary.removed_lines, 4);
2220        assert_eq!(summary.by_file[0].path, "README.md");
2221        assert_eq!(summary.by_file[0].estimated_coding_seconds, 900);
2222        assert_eq!(summary.by_filetype[0].name, "md");
2223        assert_eq!(summary.by_filetype[0].estimated_coding_seconds, 900);
2224        assert_eq!(summary.by_filetype[1].name, "rs");
2225        assert_eq!(summary.by_filetype[1].estimated_coding_seconds, 360);
2226    }
2227
2228    #[test]
2229    fn mutation_fingerprint_ignores_random_event_id_but_keeps_time_bucket() {
2230        let first: WorktreeMutationEvent = serde_json::from_value(stats_test_event(
2231            "2026-07-08T10:00:00Z",
2232            "src/main.rs",
2233            4,
2234            1,
2235        ))
2236        .unwrap();
2237        let same_bucket: WorktreeMutationEvent = serde_json::from_value(stats_test_event(
2238            "2026-07-08T10:00:00Z",
2239            "src/main.rs",
2240            4,
2241            1,
2242        ))
2243        .unwrap();
2244        let next_bucket: WorktreeMutationEvent = serde_json::from_value(stats_test_event(
2245            "2026-07-08T10:00:06Z",
2246            "src/main.rs",
2247            4,
2248            1,
2249        ))
2250        .unwrap();
2251
2252        assert_eq!(
2253            mutation_event_fingerprint(&first),
2254            mutation_event_fingerprint(&same_bucket)
2255        );
2256        assert_ne!(
2257            mutation_event_fingerprint(&first),
2258            mutation_event_fingerprint(&next_bucket)
2259        );
2260    }
2261
2262    #[test]
2263    fn classifies_create_remove_and_rename_events() {
2264        assert_eq!(
2265            classify_event_kind(&EventKind::Create(CreateKind::File)),
2266            "file_create"
2267        );
2268        assert_eq!(
2269            classify_event_kind(&EventKind::Remove(RemoveKind::Folder)),
2270            "folder_remove"
2271        );
2272        assert_eq!(
2273            classify_event_kind(&EventKind::Modify(ModifyKind::Name(RenameMode::Both))),
2274            "rename_or_move"
2275        );
2276    }
2277
2278    fn stats_test_event(
2279        occurred_at: &str,
2280        path: &str,
2281        added_lines: u64,
2282        removed_lines: u64,
2283    ) -> Value {
2284        json!({
2285            "id": Uuid::new_v4().to_string(),
2286            "repoOwner": "xylex-group",
2287            "repoName": "xbp",
2288            "branchName": "main",
2289            "repoRoot": r"C:\Users\floris\Documents\GitHub\xbp",
2290            "headSha": null,
2291            "eventKind": "file_modify",
2292            "paths": [path],
2293            "primaryPath": path,
2294            "oldPath": null,
2295            "newPath": null,
2296            "addedLines": added_lines,
2297            "removedLines": removed_lines,
2298            "fileCreated": false,
2299            "fileRemoved": false,
2300            "folderCreated": false,
2301            "folderRemoved": false,
2302            "renamedOrMoved": false,
2303            "rawKind": "Modify(Data(Content))",
2304            "occurredAt": occurred_at,
2305        })
2306    }
2307}