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 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
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:{}}} 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 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}