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