use std::env;
use std::fs::{self, File};
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use log::{info, warn};
use time::macros::format_description;
use time::OffsetDateTime;
static TAP_ID: AtomicUsize = AtomicUsize::new(0);
#[derive(Clone, Debug)]
pub(crate) struct RawFrameTap {
sink: Option<Arc<Sink>>,
}
impl RawFrameTap {
pub(crate) fn disabled() -> Self {
Self { sink: None }
}
pub(crate) fn from_env() -> Self {
match env::var("IBAPI_RAW_CAPTURE_DIR") {
Ok(dir) if !dir.is_empty() => Self::capturing_to(dir),
_ => Self::disabled(),
}
}
pub(crate) fn capturing_to(dir: impl AsRef<Path>) -> Self {
let dir = dir.as_ref().to_path_buf();
if let Err(err) = fs::create_dir_all(&dir) {
warn!("raw frame capture disabled: cannot create {}: {err}", dir.display());
return Self::disabled();
}
let stamp = OffsetDateTime::now_utc()
.format(&format_description!("[year]-[month]-[day]-[hour]-[minute]"))
.unwrap_or_else(|_| String::from("unknown"));
let sink = Sink {
prefix: format!("{stamp}-{}", TAP_ID.fetch_add(1, Ordering::SeqCst)),
dir,
state: Mutex::new(State {
segment: None,
next_number: 1,
}),
};
match sink.new_segment(0) {
Some(segment) => {
sink.lock().segment = Some(segment);
Self { sink: Some(Arc::new(sink)) }
}
None => Self::disabled(),
}
}
pub(crate) fn record_length_prefix(&self, prefix: &[u8; 4]) {
let Some(sink) = &self.sink else { return };
let declared = u32::from_be_bytes(*prefix);
sink.write(|segment| {
let line = format!("{},{},{},{declared}\n", segment.next_seq, timestamp(), segment.offset);
segment.frames.write_all(prefix)?;
segment.offset += prefix.len() as u64;
segment.next_seq += 1;
segment.index.write_all(line.as_bytes())
});
}
pub(crate) fn record_body(&self, body: &[u8]) {
let Some(sink) = &self.sink else { return };
sink.write(|segment| {
segment.frames.write_all(body)?;
segment.offset += body.len() as u64;
Ok(())
});
}
pub(crate) fn start_new_segment(&self) {
let Some(sink) = &self.sink else { return };
let mut state = sink.lock();
if state.segment.is_none() {
return;
}
let number = state.next_number;
state.next_number += 1;
state.segment = sink.new_segment(number);
}
}
fn timestamp() -> String {
let format = format_description!("[year]-[month]-[day]T[hour]:[minute]:[second].[subsecond digits:3]Z");
OffsetDateTime::now_utc().format(&format).unwrap_or_else(|_| String::from("unknown"))
}
#[derive(Debug)]
struct Sink {
dir: PathBuf,
prefix: String,
state: Mutex<State>,
}
#[derive(Debug)]
struct State {
segment: Option<Segment>,
next_number: usize,
}
#[derive(Debug)]
struct Segment {
frames: File,
index: File,
offset: u64,
next_seq: usize,
}
impl Sink {
fn lock(&self) -> std::sync::MutexGuard<'_, State> {
self.state.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn new_segment(&self, number: usize) -> Option<Segment> {
let base = self.dir.join(format!("{}-inbound-{number:03}", self.prefix));
let frames_path = base.with_extension("bin");
match (File::create(&frames_path), File::create(base.with_extension("idx"))) {
(Ok(frames), Ok(index)) => {
info!("raw frame capture: {}", frames_path.display());
Some(Segment {
frames,
index,
offset: 0,
next_seq: 0,
})
}
(Err(err), _) | (_, Err(err)) => {
warn!("raw frame capture disabled: cannot open {}: {err}", frames_path.display());
None
}
}
}
fn write(&self, record: impl FnOnce(&mut Segment) -> std::io::Result<()>) {
let mut state = self.lock();
let Some(segment) = state.segment.as_mut() else { return };
if let Err(err) = record(segment) {
warn!("raw frame capture disabled: write failed: {err}");
state.segment = None;
}
}
}
#[cfg(test)]
pub(crate) mod test_support {
use std::fs;
use std::path::{Path, PathBuf};
pub(crate) fn segments(dir: &Path, extension: &str) -> Vec<PathBuf> {
let mut paths: Vec<PathBuf> = fs::read_dir(dir)
.expect("capture directory must exist")
.map(|entry| entry.expect("readable entry").path())
.filter(|path| path.extension().is_some_and(|ext| ext == extension))
.collect();
paths.sort();
paths
}
pub(crate) fn frames(dir: &Path) -> Vec<u8> {
segments(dir, "bin")
.iter()
.flat_map(|path| fs::read(path).expect("readable capture"))
.collect()
}
pub(crate) fn index(dir: &Path) -> Vec<String> {
segments(dir, "idx")
.iter()
.flat_map(|path| {
fs::read_to_string(path)
.expect("readable index")
.lines()
.map(String::from)
.collect::<Vec<_>>()
})
.collect()
}
}
#[cfg(test)]
#[path = "raw_capture_tests.rs"]
mod tests;