Skip to main content

aft/
logging.rs

1//! Durable process logging and low-cost periodic performance summaries.
2//!
3//! Rust module processes use one file per PID. That avoids cross-process rename
4//! races while preserving a single greppable directory for all AFT activity.
5
6use crate::executor::Executor;
7use crate::run_tool_call::ToolCallPhaseDurations;
8use std::collections::{BTreeMap, VecDeque};
9use std::fs::{self, File, OpenOptions};
10use std::io::{self, BufWriter, Write};
11use std::path::{Path, PathBuf};
12use std::sync::atomic::{AtomicU64, Ordering};
13use std::sync::mpsc::{self, SyncSender, TrySendError};
14use std::sync::{LazyLock, Mutex, OnceLock};
15use std::thread;
16use std::time::{Duration, Instant, SystemTime};
17
18const LOG_FILE_BYTES: u64 = 20 * 1024 * 1024;
19const LOG_GENERATIONS: usize = 5;
20const ROTATION_CHECK_EVERY: u64 = 64;
21const LOG_CHANNEL_CAPACITY: usize = 4096;
22const DEAD_PROCESS_LOG_MAX_AGE: Duration = Duration::from_secs(7 * 24 * 60 * 60);
23const DEFAULT_PERF_TICK_INTERVAL: Duration = Duration::from_secs(60);
24const PERF_SAMPLE_INTERVAL: Duration = Duration::from_millis(250);
25const SLOW_TOOL_CALL_THRESHOLD: Duration = Duration::from_millis(50);
26const TOOL_CALL_SAMPLE_CAPACITY: usize = 256;
27
28/// Initialize the `RUST_LOG`-filtered stderr logger and its additive file sink.
29pub fn init() {
30    let storage_root = crate::bash_background::storage_dir(None);
31    let logs_dir = storage_root.join("logs");
32    let file_name = format!("aft-{}.log", std::process::id());
33    let file_path = logs_dir.join(file_name);
34
35    let file_tx = match prepare_file_sink(&logs_dir, &file_path) {
36        Ok(sink) => {
37            let (tx, rx) = mpsc::sync_channel(LOG_CHANNEL_CAPACITY);
38            thread::Builder::new()
39                .name("aft-log-writer".to_string())
40                .spawn(move || run_file_writer(sink, rx))
41                .map(|_| {
42                    if let Ok(mut control) = FILE_CONTROL.lock() {
43                        control.tx = Some(tx.clone());
44                        control.storage_root = Some(storage_root.clone());
45                    }
46                    Some(tx)
47                })
48                .unwrap_or_else(|error| {
49                    write_stderr_once(&format!(
50                        "[aft] durable log disabled: cannot start writer thread: {error}\n"
51                    ));
52                    None
53                })
54        }
55        Err(error) => {
56            write_stderr_once(&format!(
57                "[aft] durable log disabled for {}: {error}\n",
58                file_path.display()
59            ));
60            None
61        }
62    };
63
64    env_logger::Builder::from_env(env_logger::Env::default().default_filter_or("info"))
65        .target(env_logger::Target::Pipe(Box::new(TeeWriter { file_tx })))
66        .format(|buf, record| {
67            let prefix = if record.target().starts_with("aft::lsp")
68                || record.target().starts_with("aft_lsp")
69            {
70                "[aft-lsp]"
71            } else {
72                "[aft]"
73            };
74            // Wall-clock stamp so post-hoc log forensics can correlate
75            // lines with external events (health probes, module bounces).
76            // Seconds precision is enough; chrono is avoided on purpose —
77            // this hand-rolls UTC from the epoch to keep deps flat.
78            writeln!(
79                buf,
80                "{} {} {}",
81                format_utc_timestamp(),
82                prefix,
83                record.args()
84            )
85        })
86        .init();
87}
88
89/// Render `now` as `YYYY-MM-DDTHH:MM:SSZ` without a date-time dependency.
90///
91/// Civil-date math uses the days-from-epoch algorithm (Howard Hinnant's
92/// `civil_from_days`); u64 seconds keep it valid far past 2100.
93fn format_utc_timestamp() -> String {
94    let secs = SystemTime::now()
95        .duration_since(SystemTime::UNIX_EPOCH)
96        .map(|d| d.as_secs())
97        .unwrap_or(0);
98    format_epoch_secs(secs)
99}
100
101fn format_epoch_secs(secs: u64) -> String {
102    let (days, rem) = (secs / 86_400, secs % 86_400);
103    let (hh, mm, ss) = (rem / 3600, (rem % 3600) / 60, rem % 60);
104    // Howard Hinnant's civil_from_days: adding 719,468 shifts Unix epoch day 0
105    // into the algorithm's era, which begins on 0000-03-01 (putting the leap
106    // day last in each year simplifies the month/day arithmetic below).
107    let z = days as i64 + 719_468;
108    let era = z.div_euclid(146_097);
109    let doe = z.rem_euclid(146_097);
110    let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
111    let y = yoe + era * 400;
112    let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
113    let mp = (5 * doy + 2) / 153;
114    let d = doy - (153 * mp + 2) / 5 + 1;
115    let m = if mp < 10 { mp + 3 } else { mp - 9 };
116    let y = if m <= 2 { y + 1 } else { y };
117    format!("{y:04}-{m:02}-{d:02}T{hh:02}:{mm:02}:{ss:02}Z")
118}
119
120fn prepare_file_sink(logs_dir: &Path, file_path: &Path) -> io::Result<RotatingFile> {
121    fs::create_dir_all(logs_dir)?;
122    sweep_dead_process_logs(logs_dir, SystemTime::now(), DEAD_PROCESS_LOG_MAX_AGE)?;
123    RotatingFile::open(
124        file_path.to_path_buf(),
125        LOG_FILE_BYTES,
126        LOG_GENERATIONS,
127        ROTATION_CHECK_EVERY,
128    )
129}
130
131enum LogMessage {
132    Write(Vec<u8>),
133    Reconfigure(PathBuf),
134}
135
136#[derive(Default)]
137struct FileControl {
138    tx: Option<SyncSender<LogMessage>>,
139    storage_root: Option<PathBuf>,
140}
141
142static FILE_CONTROL: LazyLock<Mutex<FileControl>> =
143    LazyLock::new(|| Mutex::new(FileControl::default()));
144
145struct TeeWriter {
146    file_tx: Option<SyncSender<LogMessage>>,
147}
148
149impl Write for TeeWriter {
150    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
151        io::stderr().write_all(buf)?;
152        if let Some(tx) = self.file_tx.as_ref() {
153            match tx.try_send(LogMessage::Write(buf.to_vec())) {
154                Ok(()) => {}
155                Err(TrySendError::Full(_)) => {
156                    PERF.file_lines_dropped.fetch_add(1, Ordering::Relaxed);
157                }
158                Err(TrySendError::Disconnected(_)) => self.file_tx = None,
159            }
160        }
161        Ok(buf.len())
162    }
163
164    fn flush(&mut self) -> io::Result<()> {
165        io::stderr().flush()
166    }
167}
168
169fn run_file_writer(mut sink: RotatingFile, rx: mpsc::Receiver<LogMessage>) {
170    while let Ok(message) = rx.recv() {
171        let mut lines = Vec::new();
172        let mut reconfigure = None;
173        match message {
174            LogMessage::Write(line) => {
175                lines.push(line);
176                while lines.len() < 256 {
177                    match rx.try_recv() {
178                        Ok(LogMessage::Write(line)) => lines.push(line),
179                        Ok(LogMessage::Reconfigure(storage_root)) => {
180                            reconfigure = Some(storage_root);
181                            break;
182                        }
183                        Err(_) => break,
184                    }
185                }
186            }
187            LogMessage::Reconfigure(storage_root) => reconfigure = Some(storage_root),
188        }
189        if !lines.is_empty() {
190            if let Err(error) = sink.write_batch(&lines) {
191                write_stderr_once(&format!(
192                    "[aft] durable log disabled after write failure for {}: {error}\n",
193                    sink.path.display()
194                ));
195                break;
196            }
197        }
198        if let Some(storage_root) = reconfigure {
199            let logs_dir = storage_root.join("logs");
200            let path = logs_dir.join(format!("aft-{}.log", std::process::id()));
201            match prepare_file_sink(&logs_dir, &path) {
202                Ok(new_sink) => sink = new_sink,
203                Err(error) => write_stderr_once(&format!(
204                    "[aft] durable log could not switch to {}: {error}\n",
205                    path.display()
206                )),
207            }
208        }
209    }
210}
211
212fn write_stderr_once(message: &str) {
213    let _ = io::stderr().write_all(message.as_bytes());
214}
215
216struct RotatingFile {
217    path: PathBuf,
218    writer: Option<BufWriter<File>>,
219    size: u64,
220    threshold: u64,
221    generations: usize,
222    check_every: u64,
223    writes_since_check: u64,
224}
225
226impl RotatingFile {
227    fn open(
228        path: PathBuf,
229        threshold: u64,
230        generations: usize,
231        check_every: u64,
232    ) -> io::Result<Self> {
233        let file = OpenOptions::new().create(true).append(true).open(&path)?;
234        let size = file.metadata()?.len();
235        Ok(Self {
236            path,
237            writer: Some(BufWriter::new(file)),
238            size,
239            threshold,
240            generations,
241            check_every: check_every.max(1),
242            writes_since_check: 0,
243        })
244    }
245
246    fn write_batch(&mut self, lines: &[Vec<u8>]) -> io::Result<()> {
247        let batch_bytes = lines.iter().map(Vec::len).sum::<usize>() as u64;
248        self.writes_since_check = self.writes_since_check.saturating_add(lines.len() as u64);
249        if self.writes_since_check >= self.check_every
250            && self.size > 0
251            && self.size.saturating_add(batch_bytes) > self.threshold
252        {
253            self.rotate()?;
254        }
255        let writer = self
256            .writer
257            .as_mut()
258            .ok_or_else(|| io::Error::other("log writer unavailable"))?;
259        for line in lines {
260            writer.write_all(line)?;
261        }
262        // The worker batches channel messages before this flush. File I/O never
263        // runs on request, watcher, executor, or transport threads.
264        writer.flush()?;
265        self.size = self.size.saturating_add(batch_bytes);
266        if self.writes_since_check >= self.check_every {
267            self.writes_since_check = 0;
268        }
269        Ok(())
270    }
271
272    fn rotate(&mut self) -> io::Result<()> {
273        if let Some(mut writer) = self.writer.take() {
274            writer.flush()?;
275        }
276        if self.generations > 0 {
277            let oldest = rotated_path(&self.path, self.generations);
278            remove_file_if_present(&oldest)?;
279            for generation in (1..self.generations).rev() {
280                let from = rotated_path(&self.path, generation);
281                let to = rotated_path(&self.path, generation + 1);
282                rename_if_present(&from, &to)?;
283            }
284            rename_if_present(&self.path, &rotated_path(&self.path, 1))?;
285        } else {
286            remove_file_if_present(&self.path)?;
287        }
288        let file = OpenOptions::new()
289            .create(true)
290            .write(true)
291            .truncate(true)
292            .open(&self.path)?;
293        self.writer = Some(BufWriter::new(file));
294        self.size = 0;
295        self.writes_since_check = 0;
296        Ok(())
297    }
298}
299
300fn rotated_path(base: &Path, generation: usize) -> PathBuf {
301    let mut path = base.as_os_str().to_os_string();
302    path.push(format!(".{generation}"));
303    PathBuf::from(path)
304}
305
306fn remove_file_if_present(path: &Path) -> io::Result<()> {
307    match fs::remove_file(path) {
308        Ok(()) => Ok(()),
309        Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
310        Err(error) => Err(error),
311    }
312}
313
314fn rename_if_present(from: &Path, to: &Path) -> io::Result<()> {
315    match fs::rename(from, to) {
316        Ok(()) => Ok(()),
317        Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
318        Err(error) => Err(error),
319    }
320}
321
322fn sweep_dead_process_logs(dir: &Path, now: SystemTime, max_age: Duration) -> io::Result<()> {
323    for entry in fs::read_dir(dir)? {
324        let entry = match entry {
325            Ok(entry) => entry,
326            Err(_) => continue,
327        };
328        let name = entry.file_name();
329        let Some(pid) = process_log_pid(name.to_string_lossy().as_ref()) else {
330            continue;
331        };
332        let metadata = match entry.metadata() {
333            Ok(metadata) => metadata,
334            Err(_) => continue,
335        };
336        let age = metadata
337            .modified()
338            .ok()
339            .and_then(|modified| now.duration_since(modified).ok());
340        if age.is_some_and(|age| age >= max_age) && !process_is_alive(pid) {
341            let _ = fs::remove_file(entry.path());
342        }
343    }
344    Ok(())
345}
346
347fn process_log_pid(name: &str) -> Option<u32> {
348    let rest = name.strip_prefix("aft-")?;
349    let (pid, suffix) = rest.split_once(".log")?;
350    if !suffix.is_empty()
351        && !(suffix.starts_with('.') && suffix[1..].chars().all(|ch| ch.is_ascii_digit()))
352    {
353        return None;
354    }
355    pid.parse().ok()
356}
357
358#[cfg(unix)]
359fn process_is_alive(pid: u32) -> bool {
360    if pid == std::process::id() {
361        return true;
362    }
363    let result = unsafe { libc::kill(pid as libc::pid_t, 0) };
364    result == 0 || io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
365}
366
367#[cfg(windows)]
368fn process_is_alive(pid: u32) -> bool {
369    use std::ffi::c_void;
370
371    const PROCESS_QUERY_LIMITED_INFORMATION: u32 = 0x1000;
372    const STILL_ACTIVE: u32 = 259;
373    const ERROR_INVALID_PARAMETER: u32 = 87;
374    #[link(name = "kernel32")]
375    unsafe extern "system" {
376        #[link_name = "OpenProcess"]
377        fn open_process(access: u32, inherit_handle: i32, process_id: u32) -> *mut c_void;
378        #[link_name = "GetExitCodeProcess"]
379        fn get_exit_code_process(process: *mut c_void, exit_code: *mut u32) -> i32;
380        #[link_name = "CloseHandle"]
381        fn close_handle(object: *mut c_void) -> i32;
382        #[link_name = "GetLastError"]
383        fn get_last_error() -> u32;
384    }
385
386    if pid == std::process::id() {
387        return true;
388    }
389    let handle = unsafe { open_process(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid) };
390    if handle.is_null() {
391        // ERROR_INVALID_PARAMETER means no process owns this PID. For access
392        // errors, retain the file rather than risk deleting a live process log.
393        return unsafe { get_last_error() } != ERROR_INVALID_PARAMETER;
394    }
395    let mut exit_code = 0;
396    let queried = unsafe { get_exit_code_process(handle, &mut exit_code) } != 0;
397    unsafe { close_handle(handle) };
398    queried && exit_code == STILL_ACTIVE
399}
400
401#[cfg(not(any(unix, windows)))]
402fn process_is_alive(pid: u32) -> bool {
403    pid == std::process::id()
404}
405
406#[derive(Default)]
407struct PerfMetrics {
408    watcher_ingested: AtomicU64,
409    watcher_paths: AtomicU64,
410    watcher_dropped: AtomicU64,
411    drain_slices: AtomicU64,
412    semantic_collects: AtomicU64,
413    semantic_files: AtomicU64,
414    semantic_chunks: AtomicU64,
415    semantic_ms: AtomicU64,
416    callgraph_invalidations: AtomicU64,
417    file_lines_dropped: AtomicU64,
418    tool_call_count: AtomicU64,
419    tool_calls: Mutex<VecDeque<ToolCallPerfSample>>,
420    tier2: Mutex<BTreeMap<String, (u64, u64)>>,
421    next_sample_ns: AtomicU64,
422    reporter: Mutex<PerfReporter>,
423}
424
425struct PerfReporter {
426    last_report: Instant,
427    last_completed_interactive: u64,
428    last_completed_maintenance: u64,
429    last_tool_call_count: u64,
430}
431
432impl Default for PerfReporter {
433    fn default() -> Self {
434        Self {
435            last_report: Instant::now(),
436            last_completed_interactive: 0,
437            last_completed_maintenance: 0,
438            last_tool_call_count: 0,
439        }
440    }
441}
442
443#[derive(Clone, Copy)]
444struct ToolCallPerfSample {
445    total_ms: u64,
446    queue_ms: u64,
447}
448
449#[derive(Clone, Copy, Default)]
450struct ToolCallPerfSummary {
451    count: usize,
452    p50_total_ms: u64,
453    max_total_ms: u64,
454    p50_queue_ms: u64,
455    max_queue_ms: u64,
456}
457
458#[derive(Clone, Copy, Default)]
459struct ExecutorSample {
460    interactive_running: usize,
461    maintenance_running: usize,
462    interactive_queued: usize,
463    maintenance_queued: usize,
464    interactive_oldest_ms: Option<u64>,
465    maintenance_oldest_ms: Option<u64>,
466}
467
468static PERF: LazyLock<PerfMetrics> = LazyLock::new(PerfMetrics::default);
469
470/// Move subsequent file log writes to a newly configured storage root.
471///
472/// Reconfiguration is queued behind existing writes and is a no-op when the
473/// root has not changed. Initialization and explicit configure changes call
474/// this directly, avoiding storage-root polling on transport drain turns.
475pub fn sync_storage_root(storage_root: PathBuf) {
476    let Ok(mut control) = FILE_CONTROL.lock() else {
477        return;
478    };
479    if control.storage_root.as_ref() == Some(&storage_root) {
480        return;
481    }
482    let Some(tx) = control.tx.as_ref() else {
483        return;
484    };
485    if tx
486        .try_send(LogMessage::Reconfigure(storage_root.clone()))
487        .is_ok()
488    {
489        control.storage_root = Some(storage_root);
490    }
491}
492
493/// Called by `drain_watcher_events_bounded` for dispatch events actually received.
494pub fn note_watcher_events(count: usize) {
495    PERF.watcher_ingested
496        .fetch_add(count as u64, Ordering::Relaxed);
497}
498
499/// Called when a watcher drain slice takes paths from dispatch continuation state.
500pub fn note_drain_paths(count: usize) {
501    PERF.watcher_paths
502        .fetch_add(count as u64, Ordering::Relaxed);
503}
504
505/// Called when `drain_watcher_events_bounded` receives a rescan-required overflow signal.
506pub fn note_watcher_overflow() {
507    PERF.watcher_dropped.fetch_add(1, Ordering::Relaxed);
508}
509
510/// Called by the standalone request loop before a request-triggered runtime drain.
511pub fn note_drain_slice() {
512    PERF.drain_slices.fetch_add(1, Ordering::Relaxed);
513}
514
515/// Called after `SemanticIndex::collect_chunks` has collected one real file batch.
516pub fn note_semantic_collect(chunks: usize, files: usize, elapsed_ms: u64) {
517    PERF.semantic_collects.fetch_add(1, Ordering::Relaxed);
518    PERF.semantic_chunks
519        .fetch_add(chunks as u64, Ordering::Relaxed);
520    PERF.semantic_files
521        .fetch_add(files as u64, Ordering::Relaxed);
522    PERF.semantic_ms.fetch_add(elapsed_ms, Ordering::Relaxed);
523}
524
525/// Called by `Tier2PhaseTimings::log` after a Tier-2 scan performs measurable work.
526pub fn note_tier2_scan(category: String, elapsed_ms: u64) {
527    if let Ok(mut tier2) = PERF.tier2.lock() {
528        let entry = tier2.entry(category).or_default();
529        entry.0 = entry.0.saturating_add(1);
530        entry.1 = entry.1.saturating_add(elapsed_ms);
531    }
532}
533
534/// Called after watcher-driven callgraph `refresh_files` succeeds for concrete paths.
535pub fn note_callgraph_invalidations(files: usize) {
536    PERF.callgraph_invalidations
537        .fetch_add(files as u64, Ordering::Relaxed);
538}
539
540/// Record a completed subc tool call for slow-call diagnostics and the standing
541/// perf-tick window. The writer calls this only after `write_all` has handed the
542/// complete response frame to the transport.
543pub fn note_tool_call_trace(
544    name: &str,
545    root: &Path,
546    channel: u16,
547    corr: u64,
548    phases: ToolCallPhaseDurations,
549) {
550    let sample = ToolCallPerfSample {
551        total_ms: duration_millis_u64(phases.total),
552        queue_ms: duration_millis_u64(phases.queue),
553    };
554    if let Ok(mut samples) = PERF.tool_calls.lock() {
555        if samples.len() == TOOL_CALL_SAMPLE_CAPACITY {
556            samples.pop_front();
557        }
558        samples.push_back(sample);
559        PERF.tool_call_count.fetch_add(1, Ordering::Relaxed);
560    }
561
562    crate::slog_debug!(
563        "tool_call phase name={} channel={} corr={} total_ms={:.3} queue_ms={:.3} translate_ms={:.3} exec_ms={:.3} format_ms={:.3} finalize_ms={:.3} egress_ms={:.3} egress_enqueue_ms={:.3} egress_queue_ms={:.3} egress_prepare_ms={:.3} egress_write_ms={:.3} frame_bytes={} writer_queue_depth={} writer_active={} writer_queue_full={} reserve_timeouts={} root={}",
564        name,
565        channel,
566        corr,
567        duration_millis_f64(phases.total),
568        duration_millis_f64(phases.queue),
569        duration_millis_f64(phases.translate),
570        duration_millis_f64(phases.execute),
571        duration_millis_f64(phases.format),
572        duration_millis_f64(phases.finalize),
573        duration_millis_f64(phases.egress),
574        duration_millis_f64(phases.egress_enqueue),
575        duration_millis_f64(phases.egress_queue),
576        duration_millis_f64(phases.egress_prepare),
577        duration_millis_f64(phases.egress_write),
578        phases.frame_bytes,
579        phases.writer_queue_depth,
580        phases.writer_active_at_enqueue,
581        phases.writer_queue_was_full,
582        phases.writer_reserve_timeouts,
583        root.display(),
584    );
585
586    if phases.total > SLOW_TOOL_CALL_THRESHOLD {
587        crate::slog_warn!(
588            "slow tool_call name={} channel={} corr={} total={}ms queue={} translate={} exec={} format={} finalize={} egress={} egress_enqueue={} egress_queue={} egress_prepare={} egress_write={} frame_bytes={} writer_queue_depth={} writer_active={} writer_queue_full={} reserve_timeouts={} root={}",
589            name,
590            channel,
591            corr,
592            duration_millis_u64(phases.total),
593            duration_millis_u64(phases.queue),
594            duration_millis_u64(phases.translate),
595            duration_millis_u64(phases.execute),
596            duration_millis_u64(phases.format),
597            duration_millis_u64(phases.finalize),
598            duration_millis_u64(phases.egress),
599            duration_millis_u64(phases.egress_enqueue),
600            duration_millis_u64(phases.egress_queue),
601            duration_millis_u64(phases.egress_prepare),
602            duration_millis_u64(phases.egress_write),
603            phases.frame_bytes,
604            phases.writer_queue_depth,
605            phases.writer_active_at_enqueue,
606            phases.writer_queue_was_full,
607            phases.writer_reserve_timeouts,
608            root.display(),
609        );
610    }
611}
612
613/// Sample executor liveness and emit one busy-only aggregate at the configured cadence.
614///
615/// The transport may call this every loop turn; an atomic deadline keeps all
616/// executor sampling and reporter locking off that path between drain ticks.
617pub fn perf_tick(executor: Option<&Executor>) {
618    if !perf_sample_due() {
619        return;
620    }
621
622    let sample = executor.and_then(|executor| {
623        executor
624            .try_dispatch_liveness_snapshot()
625            .map(|snapshot| ExecutorSample {
626                interactive_running: snapshot.running.interactive,
627                maintenance_running: snapshot.running.maintenance,
628                interactive_queued: snapshot.interactive.queued,
629                maintenance_queued: snapshot.maintenance.queued,
630                interactive_oldest_ms: snapshot.interactive.oldest_age_ms,
631                maintenance_oldest_ms: snapshot.maintenance.oldest_age_ms,
632            })
633    });
634
635    let completion_counts = executor.map_or((0, 0), Executor::completion_counts);
636    let tool_call_count = PERF.tool_call_count.load(Ordering::Relaxed);
637    let (completed_interactive, completed_maintenance, new_tool_calls) = {
638        let Ok(mut reporter) = PERF.reporter.lock() else {
639            return;
640        };
641        if reporter.last_report.elapsed() < perf_tick_interval() {
642            return;
643        }
644        reporter.last_report = Instant::now();
645        let completed = (
646            completion_counts
647                .0
648                .saturating_sub(reporter.last_completed_interactive),
649            completion_counts
650                .1
651                .saturating_sub(reporter.last_completed_maintenance),
652            tool_call_count.saturating_sub(reporter.last_tool_call_count),
653        );
654        reporter.last_completed_interactive = completion_counts.0;
655        reporter.last_completed_maintenance = completion_counts.1;
656        reporter.last_tool_call_count = tool_call_count;
657        completed
658    };
659
660    let watcher_ingested = PERF.watcher_ingested.swap(0, Ordering::Relaxed);
661    let watcher_paths = PERF.watcher_paths.swap(0, Ordering::Relaxed);
662    let watcher_dropped = PERF.watcher_dropped.swap(0, Ordering::Relaxed);
663    let drain_slices = PERF.drain_slices.swap(0, Ordering::Relaxed);
664    let semantic_collects = PERF.semantic_collects.swap(0, Ordering::Relaxed);
665    let semantic_files = PERF.semantic_files.swap(0, Ordering::Relaxed);
666    let semantic_chunks = PERF.semantic_chunks.swap(0, Ordering::Relaxed);
667    let semantic_ms = PERF.semantic_ms.swap(0, Ordering::Relaxed);
668    let callgraph_invalidations = PERF.callgraph_invalidations.swap(0, Ordering::Relaxed);
669    let file_lines_dropped = PERF.file_lines_dropped.swap(0, Ordering::Relaxed);
670    let tier2 = PERF
671        .tier2
672        .lock()
673        .map(|mut tier2| std::mem::take(&mut *tier2))
674        .unwrap_or_default();
675    let tool_calls = PERF
676        .tool_calls
677        .lock()
678        .map(|samples| summarize_tool_calls(&samples))
679        .unwrap_or_default();
680
681    let executor_busy = sample.is_some_and(|sample| {
682        sample.interactive_running > 0
683            || sample.maintenance_running > 0
684            || sample.interactive_queued > 0
685            || sample.maintenance_queued > 0
686    });
687    let active = watcher_ingested > 0
688        || watcher_paths > 0
689        || watcher_dropped > 0
690        || drain_slices > 0
691        || semantic_collects > 0
692        || callgraph_invalidations > 0
693        || completed_interactive > 0
694        || completed_maintenance > 0
695        || new_tool_calls > 0
696        || file_lines_dropped > 0
697        || !tier2.is_empty()
698        || executor_busy;
699    if !active {
700        return;
701    }
702
703    let tier2_summary = if tier2.is_empty() {
704        "none".to_string()
705    } else {
706        tier2
707            .into_iter()
708            .map(|(category, (count, ms))| format!("{category}:{count}/{ms}ms"))
709            .collect::<Vec<_>>()
710            .join(",")
711    };
712    let sample = sample.unwrap_or_default();
713    crate::slog_info!(
714        "perf tick: watcher={{ingested:{},paths:{},dropped:{}}} drains={} tier2=[{}] semantic={{collects:{},files:{},chunks:{},ms:{}}} callgraph_invalidations={} executor_completed={{interactive:{},maintenance:{}}} oldest_queued_ms={{interactive:{},maintenance:{}}} toolcall={{count:{},p50_total_ms:{},max_total_ms:{},p50_queue_ms:{},max_queue_ms:{}}} file_log_dropped={}",
715        watcher_ingested,
716        watcher_paths,
717        watcher_dropped,
718        drain_slices,
719        tier2_summary,
720        semantic_collects,
721        semantic_files,
722        semantic_chunks,
723        semantic_ms,
724        callgraph_invalidations,
725        completed_interactive,
726        completed_maintenance,
727        format_optional_ms(sample.interactive_oldest_ms),
728        format_optional_ms(sample.maintenance_oldest_ms),
729        tool_calls.count,
730        tool_calls.p50_total_ms,
731        tool_calls.max_total_ms,
732        tool_calls.p50_queue_ms,
733        tool_calls.max_queue_ms,
734        file_lines_dropped,
735    );
736}
737
738fn duration_millis_f64(duration: Duration) -> f64 {
739    duration.as_secs_f64() * 1_000.0
740}
741
742fn duration_millis_u64(duration: Duration) -> u64 {
743    duration.as_millis().min(u64::MAX as u128) as u64
744}
745
746fn summarize_tool_calls(samples: &VecDeque<ToolCallPerfSample>) -> ToolCallPerfSummary {
747    if samples.is_empty() {
748        return ToolCallPerfSummary::default();
749    }
750    let mut totals = samples
751        .iter()
752        .map(|sample| sample.total_ms)
753        .collect::<Vec<_>>();
754    let mut queues = samples
755        .iter()
756        .map(|sample| sample.queue_ms)
757        .collect::<Vec<_>>();
758    totals.sort_unstable();
759    queues.sort_unstable();
760    let median_index = (samples.len() - 1) / 2;
761    ToolCallPerfSummary {
762        count: samples.len(),
763        p50_total_ms: totals[median_index],
764        max_total_ms: totals[totals.len() - 1],
765        p50_queue_ms: queues[median_index],
766        max_queue_ms: queues[queues.len() - 1],
767    }
768}
769
770fn format_optional_ms(value: Option<u64>) -> String {
771    value
772        .map(|value| value.to_string())
773        .unwrap_or_else(|| "none".to_string())
774}
775
776fn perf_sample_due() -> bool {
777    static ORIGIN: LazyLock<Instant> = LazyLock::new(Instant::now);
778    let now_ns = ORIGIN.elapsed().as_nanos().min(u64::MAX as u128) as u64;
779    let mut deadline = PERF.next_sample_ns.load(Ordering::Relaxed);
780    loop {
781        if now_ns < deadline {
782            return false;
783        }
784        let next = now_ns.saturating_add(PERF_SAMPLE_INTERVAL.as_nanos() as u64);
785        match PERF.next_sample_ns.compare_exchange_weak(
786            deadline,
787            next,
788            Ordering::Relaxed,
789            Ordering::Relaxed,
790        ) {
791            Ok(_) => return true,
792            Err(observed) => deadline = observed,
793        }
794    }
795}
796
797fn perf_tick_interval() -> Duration {
798    static INTERVAL: OnceLock<Duration> = OnceLock::new();
799    *INTERVAL.get_or_init(|| {
800        std::env::var("AFT_PERF_TICK_INTERVAL_MS")
801            .ok()
802            .and_then(|value| value.parse::<u64>().ok())
803            .filter(|value| *value > 0)
804            .map(Duration::from_millis)
805            .unwrap_or(DEFAULT_PERF_TICK_INTERVAL)
806    })
807}
808
809#[cfg(test)]
810mod tests {
811    use super::*;
812    use filetime::{set_file_mtime, FileTime};
813    use tempfile::TempDir;
814
815    fn line(value: &str) -> Vec<Vec<u8>> {
816        vec![format!("{value}\n").into_bytes()]
817    }
818
819    #[test]
820    fn epoch_timestamp_renders_known_dates() {
821        // Epoch start, a modern date, a post-2038 date (u64 range), and the
822        // 2100 non-leap century boundary that naive leap logic gets wrong.
823        assert_eq!(format_epoch_secs(0), "1970-01-01T00:00:00Z");
824        assert_eq!(format_epoch_secs(1_704_067_200), "2024-01-01T00:00:00Z");
825        assert_eq!(format_epoch_secs(1_709_251_199), "2024-02-29T23:59:59Z");
826        assert_eq!(format_epoch_secs(4_102_444_800), "2100-01-01T00:00:00Z");
827        assert_eq!(format_epoch_secs(4_107_542_399), "2100-02-28T23:59:59Z");
828    }
829
830    #[test]
831    fn rotation_rolls_at_threshold_and_preserves_newest_generations() {
832        let temp = TempDir::new().unwrap();
833        let path = temp.path().join("aft-123.log");
834        let mut sink = RotatingFile::open(path.clone(), 10, 2, 1).unwrap();
835        sink.write_batch(&line("aaaa")).unwrap();
836        sink.write_batch(&line("bbbb")).unwrap();
837        sink.write_batch(&line("cccc")).unwrap();
838        sink.write_batch(&line("dddd")).unwrap();
839        sink.write_batch(&line("eeee")).unwrap();
840
841        assert_eq!(fs::read_to_string(&path).unwrap(), "eeee\n");
842        assert_eq!(
843            fs::read_to_string(rotated_path(&path, 1)).unwrap(),
844            "cccc\ndddd\n"
845        );
846        assert_eq!(
847            fs::read_to_string(rotated_path(&path, 2)).unwrap(),
848            "aaaa\nbbbb\n"
849        );
850        assert!(!rotated_path(&path, 3).exists());
851    }
852
853    #[test]
854    fn dead_pid_sweep_removes_only_old_process_logs() {
855        let temp = TempDir::new().unwrap();
856        let dead = temp.path().join("aft-4294967294.log");
857        let dead_rotated = temp.path().join("aft-4294967294.log.1");
858        let live = temp.path().join(format!("aft-{}.log", std::process::id()));
859        let unrelated = temp.path().join("aft-plugin.log");
860        for path in [&dead, &dead_rotated, &live, &unrelated] {
861            fs::write(path, "log").unwrap();
862            set_file_mtime(path, FileTime::from_unix_time(1, 0)).unwrap();
863        }
864
865        sweep_dead_process_logs(
866            temp.path(),
867            SystemTime::UNIX_EPOCH + Duration::from_secs(10 * 24 * 60 * 60),
868            DEAD_PROCESS_LOG_MAX_AGE,
869        )
870        .unwrap();
871
872        assert!(!dead.exists());
873        assert!(!dead_rotated.exists());
874        assert!(live.exists());
875        assert!(unrelated.exists());
876    }
877
878    #[test]
879    fn tool_call_summary_uses_bounded_window_median_and_maxima() {
880        let samples = VecDeque::from([
881            ToolCallPerfSample {
882                total_ms: 9,
883                queue_ms: 5,
884            },
885            ToolCallPerfSample {
886                total_ms: 3,
887                queue_ms: 1,
888            },
889            ToolCallPerfSample {
890                total_ms: 7,
891                queue_ms: 2,
892            },
893            ToolCallPerfSample {
894                total_ms: 5,
895                queue_ms: 4,
896            },
897        ]);
898
899        let summary = summarize_tool_calls(&samples);
900
901        assert_eq!(summary.count, 4);
902        assert_eq!(summary.p50_total_ms, 5);
903        assert_eq!(summary.max_total_ms, 9);
904        assert_eq!(summary.p50_queue_ms, 2);
905        assert_eq!(summary.max_queue_ms, 5);
906    }
907}