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