1use 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
19const LOG_FILE_BYTES: u64 = 32 * 1024 * 1024;
21const LOG_GENERATIONS: usize = 1;
23const ROTATION_CHECK_EVERY: u64 = 1;
25const LOG_CHANNEL_CAPACITY: usize = 4096;
26const DEAD_PROCESS_LOG_MAX_AGE: Duration = Duration::from_secs(24 * 60 * 60);
28const LOG_DIRECTORY_BUDGET_BYTES: u64 = 200 * 1024 * 1024;
30const 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
37pub 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 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
104fn 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 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 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
377fn 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 "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 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 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
487pub 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 window: 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
590pub 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
613pub fn note_watcher_events(count: usize) {
615 PERF.watcher_ingested
616 .fetch_add(count as u64, Ordering::Relaxed);
617}
618
619pub fn note_drain_paths(count: usize) {
621 PERF.watcher_paths
622 .fetch_add(count as u64, Ordering::Relaxed);
623}
624
625pub fn note_watcher_overflow() {
627 PERF.watcher_dropped.fetch_add(1, Ordering::Relaxed);
628}
629
630pub fn note_drain_slice() {
632 PERF.drain_slices.fetch_add(1, Ordering::Relaxed);
633}
634
635pub 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
645pub 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
654pub fn note_callgraph_invalidations(files: usize) {
656 PERF.callgraph_invalidations
657 .fetch_add(files as u64, Ordering::Relaxed);
658}
659
660pub 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
733pub 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:{}}} {} 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 format_tool_call_summary(new_tool_calls, tool_calls),
850 file_lines_dropped,
851 );
852}
853
854fn duration_millis_f64(duration: Duration) -> f64 {
855 duration.as_secs_f64() * 1_000.0
856}
857
858fn duration_millis_u64(duration: Duration) -> u64 {
859 duration.as_millis().min(u64::MAX as u128) as u64
860}
861
862fn summarize_tool_calls(samples: &VecDeque<ToolCallPerfSample>) -> ToolCallPerfSummary {
863 if samples.is_empty() {
864 return ToolCallPerfSummary::default();
865 }
866 let mut totals = samples
867 .iter()
868 .map(|sample| sample.total_ms)
869 .collect::<Vec<_>>();
870 let mut queues = samples
871 .iter()
872 .map(|sample| sample.queue_ms)
873 .collect::<Vec<_>>();
874 totals.sort_unstable();
875 queues.sort_unstable();
876 let median_index = (samples.len() - 1) / 2;
877 ToolCallPerfSummary {
878 window: samples.len(),
879 p50_total_ms: totals[median_index],
880 max_total_ms: totals[totals.len() - 1],
881 p50_queue_ms: queues[median_index],
882 max_queue_ms: queues[queues.len() - 1],
883 }
884}
885
886fn format_tool_call_summary(new_tool_calls: u64, summary: ToolCallPerfSummary) -> String {
887 format!(
888 "toolcall={{new:{new_tool_calls},window:{},p50_total_ms:{},max_total_ms:{},p50_queue_ms:{},max_queue_ms:{}}}",
889 summary.window,
890 summary.p50_total_ms,
891 summary.max_total_ms,
892 summary.p50_queue_ms,
893 summary.max_queue_ms,
894 )
895}
896
897fn format_optional_ms(value: Option<u64>) -> String {
898 value
899 .map(|value| value.to_string())
900 .unwrap_or_else(|| "none".to_string())
901}
902
903fn perf_sample_due() -> bool {
904 static ORIGIN: LazyLock<Instant> = LazyLock::new(Instant::now);
905 let now_ns = ORIGIN.elapsed().as_nanos().min(u64::MAX as u128) as u64;
906 let mut deadline = PERF.next_sample_ns.load(Ordering::Relaxed);
907 loop {
908 if now_ns < deadline {
909 return false;
910 }
911 let next = now_ns.saturating_add(PERF_SAMPLE_INTERVAL.as_nanos() as u64);
912 match PERF.next_sample_ns.compare_exchange_weak(
913 deadline,
914 next,
915 Ordering::Relaxed,
916 Ordering::Relaxed,
917 ) {
918 Ok(_) => return true,
919 Err(observed) => deadline = observed,
920 }
921 }
922}
923
924fn perf_tick_interval() -> Duration {
925 static INTERVAL: OnceLock<Duration> = OnceLock::new();
926 *INTERVAL.get_or_init(|| {
927 std::env::var("AFT_PERF_TICK_INTERVAL_MS")
928 .ok()
929 .and_then(|value| value.parse::<u64>().ok())
930 .filter(|value| *value > 0)
931 .map(Duration::from_millis)
932 .unwrap_or(DEFAULT_PERF_TICK_INTERVAL)
933 })
934}
935
936#[cfg(test)]
937mod tests {
938 use super::*;
939 use filetime::{set_file_mtime, FileTime};
940 use tempfile::TempDir;
941
942 fn line(value: &str) -> Vec<Vec<u8>> {
943 vec![format!("{value}\n").into_bytes()]
944 }
945
946 #[test]
947 fn epoch_timestamp_renders_known_dates() {
948 assert_eq!(format_epoch_secs(0), "1970-01-01T00:00:00Z");
951 assert_eq!(format_epoch_secs(1_704_067_200), "2024-01-01T00:00:00Z");
952 assert_eq!(format_epoch_secs(1_709_251_199), "2024-02-29T23:59:59Z");
953 assert_eq!(format_epoch_secs(4_102_444_800), "2100-01-01T00:00:00Z");
954 assert_eq!(format_epoch_secs(4_107_542_399), "2100-02-28T23:59:59Z");
955 }
956
957 #[test]
958 fn rotation_rolls_once_and_replaces_the_single_backup_generation() {
959 let temp = TempDir::new().unwrap();
960 let path = temp.path().join("aft-123.log");
961 fs::write(rotated_path(&path, 1), "stale backup\n").unwrap();
962 let mut sink = RotatingFile::open(path.clone(), 10, 1, 1).unwrap();
963 sink.write_batch(&line("aaaa")).unwrap();
964 sink.write_batch(&line("bbbb")).unwrap();
965 sink.write_batch(&line("cccc")).unwrap();
966 sink.write_batch(&line("dddd")).unwrap();
967 sink.write_batch(&line("eeee")).unwrap();
968
969 assert_eq!(fs::read_to_string(&path).unwrap(), "eeee\n");
970 assert_eq!(
971 fs::read_to_string(rotated_path(&path, 1)).unwrap(),
972 "cccc\ndddd\n"
973 );
974 assert!(!rotated_path(&path, 2).exists());
975 }
976
977 #[test]
978 fn dead_pid_sweep_respects_age_liveness_and_explicit_plugin_exclusion() {
979 let temp = TempDir::new().unwrap();
980 let dead = temp.path().join("aft-4294967294.log");
981 let dead_rotated = temp.path().join("aft-4294967294.log.1");
982 let fresh_dead = temp.path().join("aft-4294967293.log");
983 let own = temp.path().join(format!("aft-{}.log", std::process::id()));
984 let live_rotated = rotated_path(&own, 1);
985 let plugin = temp.path().join("aft-plugin.log");
986 let now = SystemTime::UNIX_EPOCH + Duration::from_secs(10 * 24 * 60 * 60);
987 for path in [&dead, &dead_rotated, &own, &live_rotated, &plugin] {
988 fs::write(path, "log").unwrap();
989 set_file_mtime(path, FileTime::from_unix_time(1, 0)).unwrap();
990 }
991 fs::write(&fresh_dead, "fresh").unwrap();
992 set_file_mtime(
993 &fresh_dead,
994 FileTime::from_unix_time(
995 (now - DEAD_PROCESS_LOG_MAX_AGE + Duration::from_secs(1))
996 .duration_since(SystemTime::UNIX_EPOCH)
997 .unwrap()
998 .as_secs() as i64,
999 0,
1000 ),
1001 )
1002 .unwrap();
1003
1004 let summary = sweep_logs(temp.path(), now, DEAD_PROCESS_LOG_MAX_AGE, u64::MAX).unwrap();
1005
1006 assert_eq!(summary.removed_files, 2);
1007 assert!(!dead.exists());
1008 assert!(!dead_rotated.exists());
1009 assert!(fresh_dead.exists());
1010 assert!(own.exists());
1011 assert!(live_rotated.exists());
1012 assert!(plugin.exists());
1013 }
1014
1015 #[test]
1016 fn budget_backstop_deletes_oldest_dead_files_but_not_live_files() {
1017 let temp = TempDir::new().unwrap();
1018 let oldest = temp.path().join("aft-4294967294.log");
1019 let newest = temp.path().join("aft-4294967293.log");
1020 let live = temp.path().join(format!("aft-{}.log", std::process::id()));
1021 let now = SystemTime::UNIX_EPOCH + Duration::from_secs(10 * 24 * 60 * 60);
1022 fs::write(&oldest, "oldest").unwrap();
1023 fs::write(&newest, "newest").unwrap();
1024 fs::write(&live, "live-live").unwrap();
1025 set_file_mtime(&oldest, FileTime::from_unix_time(1, 0)).unwrap();
1026 set_file_mtime(&newest, FileTime::from_unix_time(2, 0)).unwrap();
1027 set_file_mtime(&live, FileTime::from_unix_time(1, 0)).unwrap();
1028
1029 let summary = sweep_logs(
1030 temp.path(),
1031 now,
1032 Duration::from_secs(365 * 24 * 60 * 60),
1033 15,
1034 )
1035 .unwrap();
1036
1037 assert_eq!(summary.removed_files, 1);
1038 assert!(!oldest.exists());
1039 assert!(newest.exists());
1040 assert!(live.exists());
1041 }
1042
1043 #[test]
1044 fn tool_call_summary_uses_bounded_window_median_and_maxima() {
1045 let samples = VecDeque::from([
1046 ToolCallPerfSample {
1047 total_ms: 9,
1048 queue_ms: 5,
1049 },
1050 ToolCallPerfSample {
1051 total_ms: 3,
1052 queue_ms: 1,
1053 },
1054 ToolCallPerfSample {
1055 total_ms: 7,
1056 queue_ms: 2,
1057 },
1058 ToolCallPerfSample {
1059 total_ms: 5,
1060 queue_ms: 4,
1061 },
1062 ]);
1063
1064 let summary = summarize_tool_calls(&samples);
1065
1066 assert_eq!(summary.window, 4);
1067 assert_eq!(summary.p50_total_ms, 5);
1068 assert_eq!(summary.max_total_ms, 9);
1069 assert_eq!(summary.p50_queue_ms, 2);
1070 assert_eq!(summary.max_queue_ms, 5);
1071 }
1072
1073 #[test]
1074 fn tool_call_tick_labels_interval_count_and_rolling_window() {
1075 let samples = VecDeque::from([ToolCallPerfSample {
1076 total_ms: 3_000,
1077 queue_ms: 2_900,
1078 }]);
1079 let summary = summarize_tool_calls(&samples);
1080
1081 assert_eq!(
1082 format_tool_call_summary(1, summary),
1083 "toolcall={new:1,window:1,p50_total_ms:3000,max_total_ms:3000,p50_queue_ms:2900,max_queue_ms:2900}"
1084 );
1085 assert_eq!(
1086 format_tool_call_summary(0, summary),
1087 "toolcall={new:0,window:1,p50_total_ms:3000,max_total_ms:3000,p50_queue_ms:2900,max_queue_ms:2900}"
1088 );
1089 }
1090}