diff --git a/rust/babblerd/src/babel/command.rs b/rust/babblerd/src/babel/command.rs new file mode 100644 index 000000000..330fbae5d --- /dev/null +++ b/rust/babblerd/src/babel/command.rs @@ -0,0 +1,53 @@ +//! Typed representation of commands sent to `babeld`'s local socket. +//! +//! This is the outbound counterpart to [`crate::babel::line`]: +//! +//! - [`crate::babel::line`] models what `babeld` emits +//! - this module models the runtime control lines that `babblerd` sends +//! +//! The scope here is intentionally narrow: this module only models the local-socket +//! commands that `babblerd` currently issues at runtime. +//! +//! NOTE: spawn-time `-C` configuration strings are still assembled in the process layer for now; +//! they are the next obvious thing to extract once the session loop is no longer +//! stringly-typed. + +use std::fmt; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum BabelCommand { + Dump, + Monitor, + Unmonitor, + Quit, + Interface(Box), +} + +impl fmt::Display for BabelCommand { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Dump => f.write_str("dump"), + Self::Monitor => f.write_str("monitor"), + Self::Unmonitor => f.write_str("unmonitor"), + Self::Quit => f.write_str("quit"), + Self::Interface(ifname) => write!(f, "interface {ifname}"), + } + } +} + +#[cfg(test)] +mod tests { + use super::BabelCommand; + + #[test] + fn renders_commands() { + assert_eq!(BabelCommand::Dump.to_string(), "dump"); + assert_eq!(BabelCommand::Monitor.to_string(), "monitor"); + assert_eq!(BabelCommand::Unmonitor.to_string(), "unmonitor"); + assert_eq!(BabelCommand::Quit.to_string(), "quit"); + assert_eq!( + BabelCommand::Interface("en2".into()).to_string(), + "interface en2" + ); + } +} diff --git a/rust/babblerd/src/babel/mod.rs b/rust/babblerd/src/babel/mod.rs index b32558209..8826a70ce 100644 --- a/rust/babblerd/src/babel/mod.rs +++ b/rust/babblerd/src/babel/mod.rs @@ -6,13 +6,15 @@ use tokio::time::Duration; use futures_lite::FutureExt; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader, Lines}; -use tokio::net::UnixStream; use tokio::net::unix::{OwnedReadHalf, OwnedWriteHalf}; +use tokio::net::UnixStream; use tokio::process::Command; use tokio::sync::{broadcast, mpsc}; +use crate::babel::line::{BabelLine, Status}; use crate::{BabbleError, Result}; +pub mod command; pub mod line; #[tracing::instrument(skip_all)] @@ -113,33 +115,41 @@ impl BabeldProcess { read: &mut Lines>, write: &mut OwnedWriteHalf, send: &broadcast::Sender, - cmd: &str, + cmd: &command::BabelCommand, ) -> io::Result> { - write.write_all(cmd.as_bytes()).await?; + let encoded = cmd.to_string(); + write.write_all(encoded.as_bytes()).await?; + write.write_all(b"\n").await?; loop { + // Ok(None) is only ever returned when stream is closed. + // When parsing errors occur from babeld, those Ok(None)'s are ignored in the loop + let Some(line) = read.next_line().await? else { tracing::warn!("babeld closed unexpectedly"); return Ok(None); }; tracing::info!("[babel] {:?}", line); - // TODO: replace later with propper parsing - match line::parse::parse_line(&line) { - Ok(line) => tracing::info!("[parsed] {:?}", line), - Err(err) => tracing::error!(error=%err, "failed to parse babeld line"), - } - - let ret = match line.as_str() { - "ok" => Ok(Some(true)), - "bad" => { - tracing::warn!("malformed message sent to babeld"); - Ok(Some(false)) + let ret = match line::parse::parse_line(&line) { + Ok(parsed) => { + tracing::info!("[parsed] {:?}", parsed); + match parsed { + BabelLine::Status(Status::Ok) => Ok(Some(true)), + BabelLine::Status(Status::Bad) => { + tracing::warn!("malformed message sent to babeld"); + Ok(Some(false)) + } + BabelLine::Status(Status::No(rest)) => { + tracing::warn!("message rejected: {rest:?}"); + Ok(Some(false)) + } + _ => Ok(None), + } } - _ if line.starts_with("no") => { - tracing::warn!("message rejected"); - Ok(Some(false)) + Err(err) => { + tracing::error!(error=%err, "failed to parse babeld line"); + Ok(None) } - _ => Ok(None), }; let Ok(_) = send.send(line) else { return Ok(None); @@ -168,7 +178,7 @@ impl BabeldProcess { } tracing::info!("babeld ok"); /* TODO(evan): push rather than pull - if Self::query(&mut babel_lines, &mut writer, "monitor\n") + if Self::query(&mut babel_lines, &mut writer, &command::BabelCommand::Monitor) .await? .is_none() { @@ -181,7 +191,7 @@ impl BabeldProcess { loop { tokio::select! { _ = interval.tick() => { - if Self::query(&mut babel_lines, &mut writer, &send, "dump\n") + if Self::query(&mut babel_lines, &mut writer, &send, &command::BabelCommand::Dump) .await? .is_none() { @@ -195,7 +205,8 @@ impl BabeldProcess { }; match babble { Babble::AddIface(iface) => { - Self::query(&mut babel_lines, &mut writer, &send, format!("interface {iface}\n").as_ref()).await?; + let cmd = command::BabelCommand::Interface(iface.into_boxed_str()); + Self::query(&mut babel_lines, &mut writer, &send, &cmd).await?; } } },