use std::io::Write;
use std::net::{IpAddr, Ipv4Addr, SocketAddr, TcpStream};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex, OnceLock};
use std::time::Duration;
use crate::runtime::env_table::{EPICS_IOC_LOG_INET, EPICS_IOC_LOG_PORT};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IocLogError {
NoServerConfigured,
NoRestartThread,
}
impl std::fmt::Display for IocLogError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NoServerConfigured => {
f.write_str("no log server configured (EPICS_IOC_LOG_INET / EPICS_IOC_LOG_PORT)")
}
Self::NoRestartThread => f.write_str("could not start the log client thread"),
}
}
}
impl std::error::Error for IocLogError {}
const MSG_BUF_SIZE: usize = 0x4000;
const RESTART_DELAY: Duration = Duration::from_secs(5);
static IOC_LOG_DISABLE: AtomicBool = AtomicBool::new(false);
static PREFIX: Mutex<Option<String>> = Mutex::new(None);
static CLIENT: OnceLock<Arc<LogClient>> = OnceLock::new();
struct LogClientState {
sock: Option<TcpStream>,
msg_buf: Vec<u8>,
connect_count: u32,
shutdown: bool,
}
struct LogClient {
addr: SocketAddr,
name: String,
state: Mutex<LogClientState>,
wake: Condvar,
}
impl LogClient {
fn send_chunk(&self, state: &mut LogClientState, text: &[u8]) {
let mut rest = text;
while !rest.is_empty() {
let mut left = MSG_BUF_SIZE - state.msg_buf.len();
if left < rest.len() && !state.msg_buf.is_empty() && state.sock.is_some() {
self.flush_locked(state);
left = MSG_BUF_SIZE - state.msg_buf.len();
}
if left == 0 {
eprintln!("log client: messages to \"{}\" are lost", self.name);
break;
}
let take = left.min(rest.len());
state.msg_buf.extend_from_slice(&rest[..take]);
rest = &rest[take..];
}
}
fn send(&self, message: &str) {
let mut state = self.state.lock().expect("log client");
if let Some(prefix) = PREFIX.lock().expect("log client prefix").as_deref() {
self.send_chunk(&mut state, prefix.as_bytes());
}
self.send_chunk(&mut state, message.as_bytes());
}
fn flush_locked(&self, state: &mut LogClientState) {
let Some(sock) = state.sock.as_mut() else {
return;
};
match sock.write_all(&state.msg_buf) {
Ok(()) => {
let _ = sock.flush();
state.msg_buf.clear();
}
Err(e) => {
eprintln!(
"log client: lost contact with log server at '{}'\n because \"{e}\"",
self.name
);
state.sock = None;
}
}
}
fn connect(&self) {
let sock = TcpStream::connect_timeout(&self.addr, RESTART_DELAY);
let mut state = self.state.lock().expect("log client");
match sock {
Ok(sock) => {
let _ = sock.set_nodelay(true);
state.sock = Some(sock);
state.connect_count += 1;
eprintln!("log client: connected to log server at '{}'", self.name);
}
Err(_) => {
state.sock = None;
}
}
}
}
fn restart_thread(client: Arc<LogClient>) {
loop {
let (connected, shutdown) = {
let state = client.state.lock().expect("log client");
(state.sock.is_some(), state.shutdown)
};
if shutdown {
return;
}
if !connected {
client.connect();
}
{
let mut state = client.state.lock().expect("log client");
client.flush_locked(&mut state);
}
let state = client.state.lock().expect("log client");
let _ = client
.wake
.wait_timeout_while(state, RESTART_DELAY, |s| !s.shutdown);
}
}
fn get_config() -> Result<SocketAddr, IocLogError> {
let Some(port) = EPICS_IOC_LOG_PORT.long() else {
eprintln!(
"iocLog: EPICS environment variable \"{}\" undefined",
EPICS_IOC_LOG_PORT.name()
);
return Err(IocLogError::NoServerConfigured);
};
if !(0..=i64::from(u16::MAX)).contains(&port) {
eprintln!(
"iocLog: EPICS environment variable \"{}\" out of range",
EPICS_IOC_LOG_PORT.name()
);
return Err(IocLogError::NoServerConfigured);
}
let inet = EPICS_IOC_LOG_INET.get().unwrap_or_default();
let Ok(addr) = inet.trim().parse::<Ipv4Addr>() else {
eprintln!(
"iocLog: EPICS environment variable \"{}\" undefined",
EPICS_IOC_LOG_INET.name()
);
return Err(IocLogError::NoServerConfigured);
};
#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
Ok(SocketAddr::new(IpAddr::V4(addr), port as u16))
}
pub fn ioc_log_init() -> Result<(), IocLogError> {
if IOC_LOG_DISABLE.load(Ordering::Relaxed) {
return Ok(());
}
if CLIENT.get().is_some() {
return Ok(());
}
let addr = get_config()?;
let client = Arc::new(LogClient {
addr,
name: addr.to_string(),
state: Mutex::new(LogClientState {
sock: None,
msg_buf: Vec::with_capacity(MSG_BUF_SIZE),
connect_count: 0,
shutdown: false,
}),
wake: Condvar::new(),
});
if CLIENT.set(Arc::clone(&client)).is_err() {
return Ok(());
}
let worker = Arc::clone(&client);
if crate::runtime::task::spawn_dedicated_thread(
"logRestart".to_string(),
crate::runtime::task::ThreadPriority::Low,
crate::runtime::task::StackSizeClass::Small,
move || restart_thread(worker),
)
.is_err()
{
eprintln!("log client: unable to start reconnection thread");
return Err(IocLogError::NoServerConfigured);
}
let sender = Arc::clone(&client);
crate::runtime::log::errlog_add_listener(move |message| {
if !IOC_LOG_DISABLE.load(Ordering::Relaxed) {
sender.send(message);
}
});
Ok(())
}
pub fn ioc_log_prefix(prefix: &str) {
let mut current = PREFIX.lock().expect("log client prefix");
match current.as_deref() {
Some(existing) => {
if existing != prefix {
println!(
"{} iocLogPrefix: The prefix was already set to \"{existing}\" and can't be changed.",
crate::runtime::log::erl_warning()
);
}
}
None => *current = Some(prefix.to_string()),
}
}
#[must_use]
pub fn ioc_log_prefix_get() -> Option<String> {
PREFIX.lock().expect("log client prefix").clone()
}
pub fn set_ioc_log_disable(disable: bool) {
IOC_LOG_DISABLE.store(disable, Ordering::Relaxed);
}
#[must_use]
pub fn ioc_log_disabled() -> bool {
IOC_LOG_DISABLE.load(Ordering::Relaxed)
}
pub fn ioc_log_flush() {
if let Some(client) = CLIENT.get() {
let mut state = client.state.lock().expect("log client");
client.flush_locked(&mut state);
}
}
#[must_use]
pub fn ioc_log_show(level: u32) -> Vec<String> {
let Some(client) = CLIENT.get() else {
return Vec::new();
};
let state = client.state.lock().expect("log client");
let mut out = Vec::new();
if state.sock.is_some() {
out.push(format!(
"log client: connected to log server at '{}'",
client.name
));
} else {
out.push(format!(
"log client: disconnected from log server at '{}'",
client.name
));
}
if let Some(prefix) = PREFIX.lock().expect("log client prefix").as_deref() {
out.push(format!("log client: prefix is \"{prefix}\""));
}
if level > 0 {
out.push(format!(
"log client: sock {}, connect cycles = {}",
if state.sock.is_some() {
"OK"
} else {
"INVALID"
},
state.connect_count
));
}
if level > 1 {
out.push(format!(
"log client: {} bytes in buffer",
state.msg_buf.len()
));
if !state.msg_buf.is_empty() {
out.push("-------------------------".to_string());
out.push(String::from_utf8_lossy(&state.msg_buf).into_owned());
out.push("-------------------------".to_string());
}
}
out
}
#[cfg(test)]
mod tests {
use super::*;
use serial_test::serial;
use std::io::Read;
use std::net::TcpListener;
#[test]
#[serial(ioc_log)]
fn with_no_inet_configured_init_declines_instead_of_connecting() {
unsafe {
std::env::remove_var("EPICS_IOC_LOG_INET");
}
assert_eq!(get_config(), Err(IocLogError::NoServerConfigured));
}
#[test]
#[serial(ioc_log)]
fn the_port_defaults_to_7004_and_is_range_checked() {
unsafe {
std::env::set_var("EPICS_IOC_LOG_INET", "127.0.0.1");
std::env::remove_var("EPICS_IOC_LOG_PORT");
}
assert_eq!(get_config().expect("configured").port(), 7004);
unsafe {
std::env::set_var("EPICS_IOC_LOG_PORT", "70000");
}
assert_eq!(get_config(), Err(IocLogError::NoServerConfigured));
unsafe {
std::env::remove_var("EPICS_IOC_LOG_PORT");
}
}
#[test]
#[serial(ioc_log)]
fn the_prefix_is_write_once_and_a_repeat_of_the_same_value_is_silent() {
*PREFIX.lock().expect("prefix") = None;
ioc_log_prefix("fac=SR ");
assert_eq!(ioc_log_prefix_get().as_deref(), Some("fac=SR "));
ioc_log_prefix("fac=SR ");
assert_eq!(ioc_log_prefix_get().as_deref(), Some("fac=SR "));
ioc_log_prefix("fac=BTS ");
assert_eq!(
ioc_log_prefix_get().as_deref(),
Some("fac=SR "),
"the first prefix stands; C refuses to change one already in use"
);
*PREFIX.lock().expect("prefix") = None;
}
#[test]
#[serial(ioc_log)]
fn a_log_server_receives_the_ioc_messages_with_the_prefix_prepended() {
let server = TcpListener::bind("127.0.0.1:0").expect("log server");
let port = server.local_addr().expect("addr").port();
unsafe {
std::env::set_var("EPICS_IOC_LOG_INET", "127.0.0.1");
std::env::set_var("EPICS_IOC_LOG_PORT", port.to_string());
}
*PREFIX.lock().expect("prefix") = None;
ioc_log_prefix("ioc=TEST ");
set_ioc_log_disable(false);
ioc_log_init().expect("the client must start");
let (mut peer, _) = server.accept().expect("the client must connect");
peer.set_read_timeout(Some(Duration::from_secs(5)))
.expect("read timeout");
crate::runtime::log::errlog_printf("bind failed\n");
crate::runtime::log::errlog_flush();
let mut buf = [0u8; 256];
let mut got = String::new();
for _ in 0..50 {
ioc_log_flush();
match peer.read(&mut buf) {
Ok(0) => break,
Ok(n) => {
got.push_str(&String::from_utf8_lossy(&buf[..n]));
break;
}
Err(_) => std::thread::sleep(Duration::from_millis(20)),
}
}
assert_eq!(
got, "ioc=TEST bind failed\n",
"the server must see the prefix then the message"
);
unsafe {
std::env::remove_var("EPICS_IOC_LOG_INET");
std::env::remove_var("EPICS_IOC_LOG_PORT");
}
}
#[test]
#[serial(ioc_log)]
fn disabling_forwarding_silences_a_client_that_is_already_connected() {
let server = TcpListener::bind("127.0.0.1:0").expect("log server");
let port = server.local_addr().expect("addr").port();
unsafe {
std::env::set_var("EPICS_IOC_LOG_INET", "127.0.0.1");
std::env::set_var("EPICS_IOC_LOG_PORT", port.to_string());
}
*PREFIX.lock().expect("prefix") = None;
set_ioc_log_disable(false);
ioc_log_init().expect("the client must start");
let (mut peer, _) = server.accept().expect("the client must connect");
peer.set_read_timeout(Some(Duration::from_millis(300)))
.expect("read timeout");
set_ioc_log_disable(true);
crate::runtime::log::errlog_printf("suppressed\n");
crate::runtime::log::errlog_flush();
ioc_log_flush();
let mut buf = [0u8; 64];
let n = peer.read(&mut buf).unwrap_or(0);
assert_eq!(
n,
0,
"nothing may reach the server while iocLogDisable is set: {:?}",
String::from_utf8_lossy(&buf[..n])
);
set_ioc_log_disable(false);
unsafe {
std::env::remove_var("EPICS_IOC_LOG_INET");
std::env::remove_var("EPICS_IOC_LOG_PORT");
}
}
}