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(¤t, 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(¤t, 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(¤t) > 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}