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