use crate::communication::ipc_types::Announce as AnnounceBlob;
use std::io::Write;
use thiserror::Error;
use tracing::info;
#[derive(Error, Debug)]
pub enum AnnounceError {
#[error("Error serializing announcement: {0}")]
Serialization(#[from] serde_json::Error),
#[error("Error writing to stdout: {0}")]
Io(#[from] std::io::Error),
}
pub fn send_announce(announce: &AnnounceBlob) -> Result<(), AnnounceError> {
let message_to_orchestrator =
crate::communication::ipc_types::ModuleToOrchestrator::Announce(announce.clone());
let json = serde_json::to_string(&message_to_orchestrator)?;
match std::io::Write::write_all(&mut std::io::stdout(), json.as_bytes()) {
Ok(_) => {
if let Err(e) = std::io::Write::write_all(&mut std::io::stdout(), b"\n") {
if e.kind() == std::io::ErrorKind::BrokenPipe {
tracing::warn!("Broken pipe when writing newline to stdout - orchestrator may have closed connection");
return Ok(()); }
return Err(AnnounceError::Io(e));
}
if let Err(e) = std::io::stdout().flush() {
if e.kind() == std::io::ErrorKind::BrokenPipe {
tracing::warn!("Broken pipe when flushing stdout - orchestrator may have closed connection");
return Ok(()); }
return Err(AnnounceError::Io(e));
}
info!("Sent announcement to orchestrator");
Ok(())
}
Err(e) => {
if e.kind() == std::io::ErrorKind::BrokenPipe {
tracing::warn!("Broken pipe when sending announcement - orchestrator may have closed connection");
tracing::info!("Announcement that failed to send: {}", json);
Ok(()) } else {
Err(AnnounceError::Io(e))
}
}
}
}