use std::collections::VecDeque;
use std::fs::{read_link, File, OpenOptions};
use std::io::{ErrorKind, Stdin, Stdout};
use std::os::unix::fs::OpenOptionsExt;
use std::os::unix::io::{AsRawFd, FromRawFd, RawFd};
use std::path::{Path, PathBuf};
use std::rc::Rc;
use std::sync::{
mpsc::{channel, Sender},
Arc, Mutex,
};
use std::thread;
use std::time::{Duration, UNIX_EPOCH};
use anyhow::{bail, Context, Result};
use log::{error, info, warn};
use nix::fcntl::{fcntl, FcntlArg, OFlag};
use nix::pty::openpty;
use nix::sys::termios::{cfmakeraw, tcgetattr, tcsetattr, SetArg, Termios};
use vmm_sys_util::epoll::EventSet;
use vmm_sys_util::eventfd::EventFd;
use machine_manager::event_loop::EventLoop;
use machine_manager::machine::{PathInfo, PTY_PATH};
use machine_manager::{
config::{ChardevConfig, ChardevType, SocketType},
temp_cleaner::TempCleaner,
};
use util::file::clear_file;
use util::loop_context::{
create_new_eventfd, gen_delete_notifiers, read_fd, EventNotifier, EventNotifierHelper,
NotifierCallback, NotifierOperation,
};
use util::set_termi_raw_mode;
use util::socket::{SocketListener, SocketStream};
use util::time::{get_format_time, gettime};
use util::unix::limit_permission;
const BUF_QUEUE_SIZE: usize = 128;
const LOG_ROTATE_SIZE_MAX: u64 = 100 * 1024 * 1024;
const LOG_ROTATE_COUNT_MAX: u32 = 7;
type PipeFile = Arc<Mutex<File>>;
pub trait InputReceiver: Send {
fn receive(&mut self, buffer: &[u8]);
fn remain_size(&mut self) -> usize;
fn set_paused(&mut self);
}
pub trait ChardevNotifyDevice: Send {
fn chardev_notify(&mut self, status: ChardevStatus);
}
pub enum ChardevStatus {
Close,
Open,
}
pub struct Chardev {
id: String,
backend: ChardevType,
listener: Option<SocketListener>,
input: Option<Arc<Mutex<dyn CommunicatInInterface>>>,
pub output: Option<Arc<Mutex<dyn CommunicatOutInterface>>>,
stream_fd: Option<i32>,
pipe_output_fd: Option<RawFd>,
receiver: Option<Arc<Mutex<dyn InputReceiver>>>,
dev: Option<Arc<Mutex<dyn ChardevNotifyDevice>>>,
wait_port: bool,
unpause_timer: Option<u64>,
output_listener_fd: Option<Arc<EventFd>>,
kick_out_evt: Arc<EventFd>,
outbuf: VecDeque<Vec<u8>>,
}
impl Chardev {
pub fn new(chardev_cfg: ChardevConfig) -> Result<Self> {
Ok(Chardev {
id: chardev_cfg.id(),
backend: chardev_cfg.classtype,
listener: None,
input: None,
output: None,
stream_fd: None,
pipe_output_fd: None,
receiver: None,
dev: None,
wait_port: false,
unpause_timer: None,
output_listener_fd: None,
kick_out_evt: Arc::new(create_new_eventfd()?),
outbuf: VecDeque::with_capacity(BUF_QUEUE_SIZE),
})
}
pub fn realize(&mut self) -> Result<()> {
match &self.backend {
ChardevType::Stdio { .. } => {
set_termi_raw_mode().with_context(|| "Failed to set terminal to raw mode")?;
self.input = Some(Arc::new(Mutex::new(std::io::stdin())));
self.output = Some(Arc::new(Mutex::new(std::io::stdout())));
}
ChardevType::Pty { .. } => {
let (master, path) =
set_pty_raw_mode().with_context(|| "Failed to set pty to raw mode")?;
info!("Pty path is: {:?}", path);
let path_info = PathInfo {
path: format!("pty:{:?}", &path),
label: self.id.clone(),
};
PTY_PATH.lock().unwrap().push(path_info);
let master_arc = Arc::new(Mutex::new(unsafe { File::from_raw_fd(master) }));
self.input = Some(master_arc.clone());
self.output = Some(master_arc);
}
ChardevType::Socket {
server,
nowait,
wait,
..
} => {
if !*server || !(*nowait || wait.as_deref() == Some("off")) {
bail!(
"Argument \'server\' and \'nowait\' or \'wait=off\' are both required for chardev \'{}\'",
&self.id
);
}
let socket_type = self.backend.socket_type()?;
if let SocketType::Tcp { host, port } = socket_type {
let listener = SocketListener::bind_by_tcp(&host, port).with_context(|| {
format!(
"Failed to bind socket for chardev \'{}\', address: {}:{}",
&self.id, host, port
)
})?;
self.listener = Some(listener);
} else if let SocketType::Unix { path } = socket_type {
clear_file(path.clone())?;
let listener = SocketListener::bind_by_uds(&path).with_context(|| {
format!(
"Failed to bind socket for chardev \'{}\', path: {}",
&self.id, path
)
})?;
self.listener = Some(listener);
TempCleaner::add_path(path.clone());
limit_permission(&path).with_context(|| {
format!(
"Failed to change file permission for chardev \'{}\', path: {}",
&self.id, path
)
})?;
}
}
ChardevType::File { path, .. } => {
let writer = RotatingFileWriter::new(path).with_context(|| {
format!("Failed to open file chardev '{}': {}", self.id, path)
})?;
self.output = Some(Arc::new(Mutex::new(writer)));
}
ChardevType::Pipe { path, .. } => {
let (input, output) = open_pipe_backend(path)?;
self.pipe_output_fd = Some(output.lock().unwrap().as_raw_fd());
self.input = Some(input);
self.output = Some(output);
}
ChardevType::Null { .. } => (),
ChardevType::RedirectToLog { .. } => {
let (sender, receiver) = channel::<u8>();
self.output = Some(Arc::new(Mutex::new(SenderWrapper(sender))));
let res = thread::Builder::new()
.name("Redirect to log".to_string())
.spawn(move || {
let mut buffer = String::new();
loop {
match receiver.recv() {
Ok(ch) => {
if ch == b'\n' {
info!("{}", buffer);
buffer.clear();
} else {
buffer.push(ch.into());
}
}
Err(e) => {
warn!("Failed to receive message: {}", e);
break;
}
}
}
});
if let Err(e) = res {
error!("Failed to start Redirect to log thread: {:?}", e);
}
}
};
Ok(())
}
pub fn set_receiver<T: 'static + InputReceiver>(&mut self, dev: &Arc<Mutex<T>>) {
self.receiver = Some(dev.clone());
if self.wait_port {
warn!("Serial port for chardev \'{}\' appeared.", &self.id);
self.wait_port = false;
self.unpause_rx();
}
}
fn wait_for_port(&mut self, input_fd: RawFd) -> EventNotifier {
warn!(
"Serial port for chardev \'{}\' is not ready yet, waiting for port.",
&self.id
);
self.wait_port = true;
EventNotifier::new(
NotifierOperation::Modify,
input_fd,
None,
EventSet::HANG_UP,
vec![],
)
}
pub fn set_device(&mut self, dev: Arc<Mutex<dyn ChardevNotifyDevice>>) {
self.dev = Some(dev.clone());
}
pub fn unpause_rx(&mut self) {
if self.input.is_none() {
error!("unpause called for non-initialized device \'{}\'", &self.id);
return;
}
if self.unpause_timer.is_some() {
return;
}
let input_fd = self.input.clone().unwrap().lock().unwrap().as_raw_fd();
let unpause_fn = Box::new(move || {
let res = EventLoop::update_event(
vec![EventNotifier::new(
NotifierOperation::AddEvents,
input_fd,
None,
EventSet::IN | EventSet::HANG_UP,
vec![],
)],
None,
);
if let Err(e) = res {
error!("Failed to unpause on fd {input_fd}: {e:?}");
}
});
let main_loop = EventLoop::get_ctx(None).unwrap();
let timer_id = main_loop.timer_add(unpause_fn, Duration::ZERO);
self.unpause_timer = Some(timer_id);
}
fn cancel_unpause_timer(&mut self) {
if let Some(timer_id) = self.unpause_timer {
let main_loop = EventLoop::get_ctx(None).unwrap();
main_loop.timer_del(timer_id);
self.unpause_timer = None;
}
}
fn clear_outbuf(&mut self) {
self.outbuf.clear();
}
pub fn outbuf_is_full(&self) -> bool {
self.outbuf.len() == self.outbuf.capacity()
}
pub fn outbuf_free_size(&self) -> usize {
self.outbuf.capacity() - self.outbuf.len()
}
pub fn fill_outbuf(&mut self, buf: Vec<u8>, listener_fd: Option<Arc<EventFd>>) -> Result<()> {
match self.backend {
ChardevType::File { .. } | ChardevType::Pty { .. } | ChardevType::Stdio { .. } => {
if self.output.is_none() {
bail!("chardev has no output");
}
return write_buffer_sync(self.output.as_ref().unwrap().clone(), buf);
}
ChardevType::Socket { .. } | ChardevType::Pipe { .. } => (),
ChardevType::Null { .. } => return Ok(()),
ChardevType::RedirectToLog { .. } => {
if self.output.is_none() {
bail!("Channel has no sender");
}
return write_buffer_sync(self.output.as_ref().unwrap().clone(), buf);
}
}
if self.output.is_none() {
return Ok(());
}
if self.outbuf_is_full() {
bail!("Failed to append buffer because output buffer queue is full");
}
self.outbuf.push_back(buf);
self.output_listener_fd = listener_fd;
let _ = self.kick_out_evt.as_ref().write(1);
Ok(())
}
pub fn set_outbuf_listener(&mut self, listener_fd: Option<Arc<EventFd>>) {
self.output_listener_fd = listener_fd;
}
fn consume_outbuf(&mut self) -> Result<()> {
if self.output.is_none() {
bail!("no output interface");
}
let output = self.output.as_ref().unwrap();
while !self.outbuf.is_empty() {
if write_buffer_async(output.clone(), self.outbuf.front_mut().unwrap())? {
break;
}
self.outbuf.pop_front();
}
Ok(())
}
}
fn write_buffer_sync(writer: Arc<Mutex<dyn CommunicatOutInterface>>, buf: Vec<u8>) -> Result<()> {
let len = buf.len();
let mut written = 0_usize;
let mut locked_writer = writer.lock().unwrap();
while written < len {
match locked_writer.write(&buf[written..len]) {
Ok(n) => written += n,
Err(e) => bail!("chardev failed to write file with error {:?}", e),
}
}
locked_writer
.flush()
.with_context(|| "chardev failed to flush")?;
Ok(())
}
fn write_buffer_async(
writer: Arc<Mutex<dyn CommunicatOutInterface>>,
buf: &mut Vec<u8>,
) -> Result<bool> {
let len = buf.len();
let mut locked_writer = writer.lock().unwrap();
let mut written = 0_usize;
while written < len {
match locked_writer.write(&buf[written..len]) {
Ok(0) => break,
Ok(n) => written += n,
Err(e) => {
let err_type = e.kind();
if err_type != ErrorKind::WouldBlock && err_type != ErrorKind::Interrupted {
bail!("chardev failed to write data with error {:?}", e);
}
break;
}
}
}
locked_writer
.flush()
.with_context(|| "chardev failed to flush")?;
if written == len {
return Ok(false);
}
buf.drain(0..written);
Ok(true)
}
fn set_pty_raw_mode() -> Result<(i32, PathBuf)> {
let (master, slave) = match openpty(None, None) {
Ok(res) => (res.master, res.slave),
Err(e) => bail!("Failed to open pty, error is {:?}", e),
};
let proc_path = PathBuf::from(format!("/proc/self/fd/{}", slave));
let path = read_link(proc_path).with_context(|| "Failed to read slave pty link")?;
let mut new_termios: Termios = match tcgetattr(slave) {
Ok(tm) => tm,
Err(e) => bail!("Failed to get mode of pty, error is {:?}", e),
};
cfmakeraw(&mut new_termios);
if let Err(e) = tcsetattr(slave, SetArg::TCSAFLUSH, &new_termios) {
bail!("Failed to set pty to raw mode, error is {:?}", e);
}
let fcnt_arg = FcntlArg::F_SETFL(OFlag::from_bits(libc::O_NONBLOCK).unwrap());
if let Err(e) = fcntl(master, fcnt_arg) {
bail!(
"Failed to set pty master to nonblocking mode, error is {:?}",
e
);
}
Ok((master, path))
}
fn open_pipe_file(path: &str) -> Result<File> {
let file = OpenOptions::new()
.read(true)
.write(true)
.custom_flags(libc::O_CLOEXEC | libc::O_NONBLOCK)
.open(path)
.with_context(|| format!("Failed to open pipe file {path}"))?;
Ok(file)
}
fn open_pipe_backend(path: &str) -> Result<(PipeFile, PipeFile)> {
let input_path = format!("{path}.in");
let output_path = format!("{path}.out");
if let (Ok(input), Ok(output)) = (open_pipe_file(&input_path), open_pipe_file(&output_path)) {
return Ok((Arc::new(Mutex::new(input)), Arc::new(Mutex::new(output))));
}
let file = Arc::new(Mutex::new(open_pipe_file(path)?));
Ok((file.clone(), file))
}
fn get_terminal_notifier(chardev: Arc<Mutex<Chardev>>) -> Option<EventNotifier> {
let locked_chardev = chardev.lock().unwrap();
let input = locked_chardev.input.clone();
if input.is_none() {
error!(
"Failed to initialize input events for chardev \'{}\', chardev not initialized",
&locked_chardev.id
);
return None;
}
let cloned_chardev = chardev.clone();
let input_fd = input.unwrap().lock().unwrap().as_raw_fd();
let event_handler: Rc<NotifierCallback> = Rc::new(move |_, _| {
let mut locked_chardev = cloned_chardev.lock().unwrap();
if locked_chardev.receiver.is_none() {
let wait_port = locked_chardev.wait_for_port(input_fd);
return Some(vec![wait_port]);
}
locked_chardev.cancel_unpause_timer();
let receiver = locked_chardev.receiver.clone().unwrap();
let input = locked_chardev.input.clone().unwrap();
drop(locked_chardev);
let mut locked_receiver = receiver.lock().unwrap();
let buff_size = locked_receiver.remain_size();
if buff_size == 0 {
locked_receiver.set_paused();
return Some(vec![EventNotifier::new(
NotifierOperation::Modify,
input_fd,
None,
EventSet::HANG_UP,
vec![],
)]);
}
let mut buffer = vec![0_u8; buff_size];
if let Ok(bytes_count) = input.lock().unwrap().chr_read_raw(&mut buffer) {
locked_receiver.receive(&buffer[..bytes_count]);
} else {
let os_error = std::io::Error::last_os_error();
let locked_chardev = cloned_chardev.lock().unwrap();
error!(
"Failed to read input data from chardev \'{}\', {}",
&locked_chardev.id, &os_error
);
}
None
});
Some(EventNotifier::new(
NotifierOperation::AddShared,
input_fd,
None,
EventSet::IN,
vec![event_handler],
))
}
fn get_pipe_notifiers(chardev: Arc<Mutex<Chardev>>) -> Vec<EventNotifier> {
let locked_chardev = chardev.lock().unwrap();
let input = locked_chardev.input.clone();
if input.is_none() || locked_chardev.output.is_none() {
error!(
"Failed to initialize pipe events for chardev \'{}\', chardev not initialized",
&locked_chardev.id
);
return Vec::new();
}
let input_fd = input.unwrap().lock().unwrap().as_raw_fd();
let output_fd = match locked_chardev.pipe_output_fd {
Some(fd) => fd,
None => {
error!("Failed to initialize pipe output event, output fd is missing");
return Vec::new();
}
};
let kick_out_fd = locked_chardev.kick_out_evt.as_ref().as_raw_fd();
drop(locked_chardev);
let handling_chardev = chardev.clone();
let send_buffers = Rc::new(move || {
let mut locked_chardev = handling_chardev.lock().unwrap();
if let Err(e) = locked_chardev.consume_outbuf() {
error!("Failed to consume outbuf with error {:?}", e);
locked_chardev.clear_outbuf();
return Some(vec![EventNotifier::new(
NotifierOperation::DeleteEvents,
output_fd,
None,
EventSet::OUT,
Vec::new(),
)]);
}
if locked_chardev.output_listener_fd.is_some() {
let fd = locked_chardev.output_listener_fd.as_ref().unwrap();
if let Err(e) = fd.write(1) {
error!("Failed to write eventfd with error {:?}", e);
return None;
}
locked_chardev.output_listener_fd = None;
}
if locked_chardev.outbuf.is_empty() {
Some(vec![EventNotifier::new(
NotifierOperation::DeleteEvents,
output_fd,
None,
EventSet::OUT,
Vec::new(),
)])
} else {
Some(vec![EventNotifier::new(
NotifierOperation::AddEvents,
output_fd,
None,
EventSet::OUT,
Vec::new(),
)])
}
});
let handling_chardev = chardev.clone();
let input_handler: Rc<NotifierCallback> = Rc::new(move |event, _| {
if event & EventSet::IN != EventSet::IN {
return None;
}
let mut locked_chardev = handling_chardev.lock().unwrap();
locked_chardev.cancel_unpause_timer();
if locked_chardev.receiver.is_none() {
let wait_port = locked_chardev.wait_for_port(input_fd);
return Some(vec![wait_port]);
}
let receiver = locked_chardev.receiver.clone().unwrap();
let input = locked_chardev.input.clone().unwrap();
drop(locked_chardev);
let mut locked_receiver = receiver.lock().unwrap();
let buff_size = locked_receiver.remain_size();
if buff_size == 0 {
locked_receiver.set_paused();
return Some(vec![EventNotifier::new(
NotifierOperation::DeleteEvents,
input_fd,
None,
EventSet::IN,
vec![],
)]);
}
let mut buffer = vec![0_u8; buff_size];
let mut locked_input = input.lock().unwrap();
match locked_input.chr_read_raw(&mut buffer) {
Ok(bytes_count) => {
if bytes_count > 0 {
locked_receiver.receive(&buffer[..bytes_count]);
} else {
return Some(vec![EventNotifier::new(
NotifierOperation::DeleteEvents,
input_fd,
None,
EventSet::IN,
vec![],
)]);
}
}
Err(_) => {
let os_error = std::io::Error::last_os_error();
if os_error.kind() != std::io::ErrorKind::WouldBlock {
let locked_chardev = handling_chardev.lock().unwrap();
error!(
"Failed to read input data from chardev \'{}\', {}",
&locked_chardev.id, &os_error
);
}
}
}
None
});
let send_buffers_cb = send_buffers.clone();
let outavail_handler = Rc::new(move |event, _| {
if event & EventSet::OUT != EventSet::OUT {
return None;
}
send_buffers_cb()
});
let send_handler = Rc::new(move |_event, fd| {
read_fd(fd);
send_buffers()
});
let mut handlers = vec![input_handler];
if input_fd == output_fd {
handlers.push(outavail_handler.clone());
}
let mut notifiers = vec![
EventNotifier::new(
NotifierOperation::AddShared,
input_fd,
None,
EventSet::IN | EventSet::HANG_UP,
handlers,
),
EventNotifier::new(
NotifierOperation::AddShared,
kick_out_fd,
None,
EventSet::IN,
vec![send_handler],
),
];
if input_fd != output_fd {
notifiers.push(EventNotifier::new(
NotifierOperation::AddShared,
output_fd,
None,
EventSet::HANG_UP,
vec![outavail_handler],
));
}
notifiers
}
fn get_socket_notifier(chardev: Arc<Mutex<Chardev>>) -> Option<EventNotifier> {
let locked_chardev = chardev.lock().unwrap();
let listener = &locked_chardev.listener;
if listener.is_none() {
error!(
"Failed to setup io-event notifications for chardev \'{}\', device not initialized",
&locked_chardev.id
);
return None;
}
let cloned_chardev = chardev.clone();
let event_handler: Rc<NotifierCallback> = Rc::new(move |_, _| {
let mut locked_chardev = cloned_chardev.lock().unwrap();
let stream = locked_chardev.listener.as_ref().unwrap().accept().unwrap();
let connection_info = stream.link_description();
info!(
"Chardev \'{}\' event, connection opened: {}",
&locked_chardev.id, connection_info
);
let stream_fd = stream.as_raw_fd();
let stream_arc = Arc::new(Mutex::new(stream));
let listener_fd = locked_chardev.listener.as_ref().unwrap().as_raw_fd();
let kick_out_fd = locked_chardev.kick_out_evt.as_ref().as_raw_fd();
let notify_dev = locked_chardev.dev.clone();
locked_chardev.stream_fd = Some(stream_fd);
locked_chardev.input = Some(stream_arc.clone());
locked_chardev.output = Some(stream_arc.clone());
drop(locked_chardev);
if let Some(dev) = notify_dev {
dev.lock().unwrap().chardev_notify(ChardevStatus::Open);
}
let handling_chardev = cloned_chardev.clone();
let close_connection = Rc::new(move || {
let mut locked_chardev = handling_chardev.lock().unwrap();
let notify_dev = locked_chardev.dev.clone();
locked_chardev.input = None;
locked_chardev.output = None;
locked_chardev.stream_fd = None;
locked_chardev.cancel_unpause_timer();
locked_chardev.outbuf.clear();
if locked_chardev.output_listener_fd.is_some() {
let _ = locked_chardev.output_listener_fd.as_ref().unwrap().write(1);
locked_chardev.output_listener_fd = None;
}
info!(
"Chardev \'{}\' event, connection closed: {}",
&locked_chardev.id, connection_info
);
drop(locked_chardev);
if let Some(dev) = notify_dev {
dev.lock().unwrap().chardev_notify(ChardevStatus::Close);
}
let stream_fd = stream_arc.lock().unwrap().as_raw_fd();
Some(gen_delete_notifiers(&[stream_fd, kick_out_fd]))
});
let handling_chardev = cloned_chardev.clone();
let input_handler: Rc<NotifierCallback> = Rc::new(move |event, _| {
let mut locked_chardev = handling_chardev.lock().unwrap();
let peer_disconnected = event & EventSet::HANG_UP == EventSet::HANG_UP;
if peer_disconnected && locked_chardev.receiver.is_none() {
drop(locked_chardev);
return close_connection();
}
let input_ready = event & EventSet::IN == EventSet::IN;
if input_ready {
locked_chardev.cancel_unpause_timer();
if locked_chardev.receiver.is_none() {
let wait_port = locked_chardev.wait_for_port(stream_fd);
return Some(vec![wait_port]);
}
let receiver = locked_chardev.receiver.clone().unwrap();
let input = locked_chardev.input.clone().unwrap();
drop(locked_chardev);
let mut locked_receiver = receiver.lock().unwrap();
let buff_size = locked_receiver.remain_size();
if buff_size == 0 {
locked_receiver.set_paused();
return Some(vec![EventNotifier::new(
NotifierOperation::DeleteEvents,
stream_fd,
None,
EventSet::IN,
vec![],
)]);
}
let mut buffer = vec![0_u8; buff_size];
let mut locked_input = input.lock().unwrap();
if let Ok(bytes_count) = locked_input.chr_read_raw(&mut buffer) {
if bytes_count > 0 {
locked_receiver.receive(&buffer[..bytes_count]);
} else {
drop(locked_receiver);
drop(locked_input);
return close_connection();
}
} else {
let os_error = std::io::Error::last_os_error();
if os_error.kind() != std::io::ErrorKind::WouldBlock {
let locked_chardev = handling_chardev.lock().unwrap();
error!(
"Failed to read input data from chardev \'{}\', {}",
&locked_chardev.id, &os_error
);
}
}
}
None
});
let handling_chardev = cloned_chardev.clone();
let send_buffers = Rc::new(move || {
let mut locked_chardev = handling_chardev.lock().unwrap();
if let Err(e) = locked_chardev.consume_outbuf() {
error!("Failed to consume outbuf with error {:?}", e);
locked_chardev.clear_outbuf();
return Some(vec![EventNotifier::new(
NotifierOperation::DeleteEvents,
stream_fd,
None,
EventSet::OUT,
Vec::new(),
)]);
}
if locked_chardev.output_listener_fd.is_some() {
let fd = locked_chardev.output_listener_fd.as_ref().unwrap();
if let Err(e) = fd.write(1) {
error!("Failed to write eventfd with error {:?}", e);
return None;
}
locked_chardev.output_listener_fd = None;
}
if locked_chardev.outbuf.is_empty() {
Some(vec![EventNotifier::new(
NotifierOperation::DeleteEvents,
stream_fd,
None,
EventSet::OUT,
Vec::new(),
)])
} else {
Some(vec![EventNotifier::new(
NotifierOperation::AddEvents,
stream_fd,
None,
EventSet::OUT,
Vec::new(),
)])
}
});
let send_buffers_cb = send_buffers.clone();
let outavail_handler = Rc::new(move |event, _| {
if event & EventSet::OUT != EventSet::OUT {
return None;
}
send_buffers_cb()
});
let send_handler = Rc::new(move |_event, fd| {
read_fd(fd);
send_buffers()
});
Some(vec![
EventNotifier::new(
NotifierOperation::AddShared,
stream_fd,
Some(listener_fd),
EventSet::IN | EventSet::HANG_UP,
vec![input_handler, outavail_handler],
),
EventNotifier::new(
NotifierOperation::AddShared,
kick_out_fd,
None,
EventSet::IN,
vec![send_handler],
),
])
});
let listener_fd = listener.as_ref().unwrap().as_raw_fd();
Some(EventNotifier::new(
NotifierOperation::AddShared,
listener_fd,
None,
EventSet::IN,
vec![event_handler],
))
}
impl EventNotifierHelper for Chardev {
fn internal_notifiers(chardev: Arc<Mutex<Self>>) -> Vec<EventNotifier> {
let notifier = {
let backend = chardev.lock().unwrap().backend.clone();
match backend {
ChardevType::Stdio { .. } => get_terminal_notifier(chardev),
ChardevType::Pty { .. } => get_terminal_notifier(chardev),
ChardevType::Socket { .. } => get_socket_notifier(chardev),
ChardevType::Pipe { .. } => return get_pipe_notifiers(chardev),
ChardevType::Null { .. } => None,
ChardevType::File { .. } => None,
ChardevType::RedirectToLog { .. } => None,
}
};
notifier.map_or(Vec::new(), |value| vec![value])
}
}
pub trait CommunicatInInterface: std::marker::Send + std::os::unix::io::AsRawFd {
fn chr_read_raw(&mut self, buf: &mut [u8]) -> Result<usize> {
match nix::unistd::read(self.as_raw_fd(), buf) {
Err(e) => bail!("Failed to read buffer: {:?}", e),
Ok(bytes) => Ok(bytes),
}
}
}
pub trait CommunicatOutInterface: std::io::Write + std::marker::Send {}
impl CommunicatInInterface for SocketStream {}
impl CommunicatInInterface for File {}
impl CommunicatInInterface for Stdin {}
impl CommunicatOutInterface for SocketStream {}
impl CommunicatOutInterface for File {}
impl CommunicatOutInterface for Stdout {}
impl CommunicatOutInterface for SenderWrapper {}
impl CommunicatOutInterface for RotatingFileWriter {}
struct RotatingFileWriter {
file: File,
path: String,
current_size: u64,
create_day_key: i32,
}
impl RotatingFileWriter {
fn new(path: &str) -> Result<Self> {
let file = open_rotating_file(path)?;
let metadata = file.metadata()?;
let modified = metadata.modified()?.duration_since(UNIX_EPOCH)?.as_secs();
Ok(Self {
file,
path: path.to_string(),
current_size: metadata.len(),
create_day_key: rotating_day_key(i64::try_from(modified)?),
})
}
fn rotate_if_needed(&mut self, written: usize) -> Result<()> {
self.current_size = self.current_size.saturating_add(written as u64);
let today = rotating_day_key(gettime()?.0);
if self.current_size < LOG_ROTATE_SIZE_MAX && self.create_day_key == today {
return Ok(());
}
let oldest = format!("{}{}", self.path, LOG_ROTATE_COUNT_MAX - 1);
if Path::new(&oldest).exists() {
std::fs::remove_file(&oldest)
.with_context(|| format!("Failed to remove log file {}", oldest))?;
}
for suffix in (0..LOG_ROTATE_COUNT_MAX - 1).rev() {
let source = if suffix == 0 {
self.path.clone()
} else {
format!("{}{}", self.path, suffix)
};
if !Path::new(&source).exists() {
continue;
}
let target = format!("{}{}", self.path, suffix + 1);
std::fs::rename(&source, &target).with_context(|| {
format!("Failed to rename log file from {} to {}", source, target)
})?;
}
self.file = open_rotating_file(&self.path)?;
self.current_size = 0;
self.create_day_key = today;
Ok(())
}
}
impl std::io::Write for RotatingFileWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.file.write_all(buf)?;
self.rotate_if_needed(buf.len())
.map_err(std::io::Error::other)?;
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
self.file.flush()
}
}
fn open_rotating_file(path: &str) -> Result<File> {
OpenOptions::new()
.append(true)
.create(true)
.mode(0o640)
.open(path)
.with_context(|| format!("Failed to open log file {}", path))
}
fn rotating_day_key(sec: i64) -> i32 {
let time = get_format_time(sec);
time[0] * 10_000 + time[1] * 100 + time[2]
}
struct SenderWrapper(Sender<u8>);
unsafe impl std::marker::Send for SenderWrapper {}
impl std::io::Write for SenderWrapper {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
for i in buf {
self.0.send(*i).map_err(std::io::Error::other)?;
}
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[cfg(test)]
mod tests {
use std::ffi::CString;
use std::fs;
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
use super::*;
fn test_dir() -> PathBuf {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos();
let path = std::env::temp_dir().join(format!(
"stratovirt-chardev-pipe-{}-{}",
std::process::id(),
now
));
fs::create_dir(&path).unwrap();
path
}
fn mkfifo(path: &Path) {
let path = CString::new(path.to_str().unwrap()).unwrap();
let ret = unsafe { libc::mkfifo(path.as_ptr(), 0o600) };
assert_eq!(ret, 0);
}
#[test]
fn test_open_pipe_backend_prefers_split_paths() {
let dir = test_dir();
let base = dir.join("charpipe");
let input = PathBuf::from(format!("{}.in", base.to_str().unwrap()));
let output = PathBuf::from(format!("{}.out", base.to_str().unwrap()));
mkfifo(&input);
mkfifo(&output);
let (input, output) = open_pipe_backend(base.to_str().unwrap()).unwrap();
let input_fd = input.lock().unwrap().as_raw_fd();
let output_fd = output.lock().unwrap().as_raw_fd();
assert_ne!(input_fd, output_fd);
fs::remove_dir_all(dir).unwrap();
}
#[test]
fn test_open_pipe_backend_falls_back_to_single_path() {
let dir = test_dir();
let base = dir.join("charpipe");
mkfifo(&base);
let (input, output) = open_pipe_backend(base.to_str().unwrap()).unwrap();
let input_fd = input.lock().unwrap().as_raw_fd();
let output_fd = output.lock().unwrap().as_raw_fd();
assert_eq!(input_fd, output_fd);
fs::remove_dir_all(dir).unwrap();
}
#[test]
fn test_open_pipe_backend_missing_path_fails() {
let dir = test_dir();
let base = dir.join("charpipe");
assert!(open_pipe_backend(base.to_str().unwrap()).is_err());
fs::remove_dir_all(dir).unwrap();
}
}