use std::path::{Path, PathBuf};
use tokio::sync::broadcast::error::{RecvError, TryRecvError};
use tokio::sync::{mpsc, watch};
use tokio::task::JoinHandle;
use clankers_core::{RobotContext, RobotError, RobotResult, Timestamp};
use clankers_data::McapWriter;
use clankers_ros2::sim::SimBus;
const ENCODING: &str = "json";
const DEFAULT_SCHEMA: &str = "clankeRS/Json";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RecordingPlan {
pub path: PathBuf,
pub topics: Vec<(String, String)>,
}
impl RecordingPlan {
pub fn resolve(
ctx: &RobotContext,
env_record: Option<&str>,
env_output: Option<&str>,
) -> Option<Self> {
let enabled = match env_record.map(|s| s.trim().to_ascii_lowercase()).as_deref() {
Some("1") | Some("true") => true,
Some("0") | Some("false") => false,
_ => ctx.config.logging.record_mcap,
};
if !enabled {
return None;
}
let path = match env_output.map(str::trim) {
Some(p) if !p.is_empty() => PathBuf::from(p),
_ => ctx
.resolve_path(&ctx.config.logging.output_dir)
.join(format!(
"{}_{}.mcap",
ctx.node_name(),
Timestamp::now().as_nanos()
)),
};
let mut topics: Vec<(String, String)> = ctx
.config
.topics
.input
.values()
.chain(ctx.config.topics.output.values())
.map(|t| {
(
t.name.clone(),
t.r#type.clone().unwrap_or_else(|| DEFAULT_SCHEMA.into()),
)
})
.collect();
topics.sort();
topics.dedup_by(|a, b| a.0 == b.0);
Some(Self { path, topics })
}
}
#[derive(Debug, Clone)]
pub struct RecordingSummary {
pub path: PathBuf,
pub messages: u64,
}
pub struct McapRecorder {
stop_tx: watch::Sender<bool>,
forwarders: Vec<JoinHandle<()>>,
writer_task: JoinHandle<RobotResult<u64>>,
path: PathBuf,
}
impl McapRecorder {
pub async fn start(ctx: &RobotContext) -> RobotResult<Option<Self>> {
let plan = RecordingPlan::resolve(
ctx,
std::env::var("CLANKERS_RECORD_MCAP").ok().as_deref(),
std::env::var("CLANKERS_RECORD_OUTPUT").ok().as_deref(),
);
match plan {
Some(plan) => Ok(Some(Self::start_with_plan(plan).await?)),
None => Ok(None),
}
}
pub async fn start_with_plan(plan: RecordingPlan) -> RobotResult<Self> {
if let Some(parent) = plan.path.parent() {
std::fs::create_dir_all(parent)
.map_err(|e| RobotError::Data(format!("create {}: {e}", parent.display())))?;
}
let mut writer = McapWriter::create(&plan.path)?;
if plan.topics.is_empty() {
tracing::warn!(
"record_mcap is enabled but no [topics.*] are configured; \
the recording will be empty"
);
}
let (msg_tx, mut msg_rx) = mpsc::channel::<(usize, u64, Vec<u8>)>(1024);
let (stop_tx, stop_rx) = watch::channel(false);
let bus = SimBus::global();
let mut forwarders = Vec::with_capacity(plan.topics.len());
for (idx, (topic, _)) in plan.topics.iter().enumerate() {
let mut rx = bus.subscribe(topic).await;
let tx = msg_tx.clone();
let mut stop = stop_rx.clone();
let topic = topic.clone();
forwarders.push(tokio::spawn(async move {
loop {
tokio::select! {
biased;
res = rx.recv() => match res {
Ok((stamp, data)) => {
if tx.send((idx, stamp, data)).await.is_err() {
break;
}
}
Err(RecvError::Lagged(n)) => {
tracing::warn!(dropped = n, topic = %topic, "recorder lagged");
}
Err(RecvError::Closed) => break,
},
_ = stop.changed() => {
loop {
match rx.try_recv() {
Ok((stamp, data)) => {
if tx.send((idx, stamp, data)).await.is_err() {
break;
}
}
Err(TryRecvError::Lagged(n)) => {
tracing::warn!(dropped = n, topic = %topic, "recorder lagged");
}
Err(_) => break,
}
}
break;
}
}
}
}));
}
drop(msg_tx);
let topics = plan.topics.clone();
let writer_task = tokio::spawn(async move {
let mut count = 0u64;
while let Some((idx, stamp, data)) = msg_rx.recv().await {
let (topic, schema) = &topics[idx];
writer.write_message(
topic,
schema,
ENCODING,
&data,
Timestamp::from_nanos(stamp),
)?;
count += 1;
}
writer.finish()?;
Ok(count)
});
Ok(Self {
stop_tx,
forwarders,
writer_task,
path: plan.path,
})
}
pub fn path(&self) -> &Path {
&self.path
}
pub async fn finish(self) -> RobotResult<RecordingSummary> {
let _ = self.stop_tx.send(true);
for forwarder in self.forwarders {
let _ = forwarder.await;
}
let messages = self
.writer_task
.await
.map_err(|e| RobotError::Data(format!("recorder task failed: {e}")))??;
tracing::info!(path = %self.path.display(), messages, "MCAP recording finished");
Ok(RecordingSummary {
path: self.path,
messages,
})
}
}
pub async fn run_with_recording<F>(ctx: RobotContext, node: F) -> RobotResult<()>
where
F: std::future::Future<Output = RobotResult<()>>,
{
let recorder = McapRecorder::start(&ctx).await?;
if let Some(rec) = &recorder {
tracing::info!(path = %rec.path().display(), "recording node I/O to MCAP");
}
let result = tokio::select! {
r = node => r,
_ = tokio::signal::ctrl_c() => {
tracing::info!("Ctrl-C received; shutting down");
Ok(())
}
};
match recorder {
Some(rec) => result.and(rec.finish().await.map(|_| ())),
None => result,
}
}