// Copyright (c) 2020 Huawei Technologies Co.,Ltd. All rights reserved.
//
// StratoVirt is licensed under Mulan PSL v2.
// You can use this software according to the terms and conditions of the Mulan
// PSL v2.
// You may obtain a copy of Mulan PSL v2 at:
//         http://license.coscl.org.cn/MulanPSL2
// THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY
// KIND, EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO
// NON-INFRINGEMENT, MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE.
// See the Mulan PSL v2 for more details.

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>>;

/// Provide the trait that helps handle the input data.
pub trait InputReceiver: Send {
    /// Handle the input data and trigger interrupt if necessary.
    fn receive(&mut self, buffer: &[u8]);
    /// Return the remain space size of receiver buffer.
    /// 0 if receiver is not ready or no space in FIFO
    fn remain_size(&mut self) -> usize;
    /// Tell receiver that RX is paused and receiver
    /// must unpause it when it becomes ready
    fn set_paused(&mut self);
}

/// Provide the trait that notifies device the socket is opened or closed.
pub trait ChardevNotifyDevice: Send {
    fn chardev_notify(&mut self, status: ChardevStatus);
}

pub enum ChardevStatus {
    Close,
    Open,
}

/// Character device structure.
pub struct Chardev {
    /// Id of chardev.
    id: String,
    /// Type of backend device.
    backend: ChardevType,
    /// Socket listener for chardev of socket type.
    listener: Option<SocketListener>,
    /// Chardev input.
    input: Option<Arc<Mutex<dyn CommunicatInInterface>>>,
    /// Chardev output.
    pub output: Option<Arc<Mutex<dyn CommunicatOutInterface>>>,
    /// Fd of socket stream.
    stream_fd: Option<i32>,
    /// Fd of pipe output.
    pipe_output_fd: Option<RawFd>,
    /// Input receiver.
    receiver: Option<Arc<Mutex<dyn InputReceiver>>>,
    /// Used to notify device the socket is opened or closed.
    dev: Option<Arc<Mutex<dyn ChardevNotifyDevice>>>,
    /// Whether event-handling of device is initialized
    /// and we wait for port to become available
    wait_port: bool,
    /// Scheduled DPC to unpause input stream.
    /// Unpause must be done inside event-loop
    unpause_timer: Option<u64>,
    /// output listener to notify when output stream fd can be written
    output_listener_fd: Option<Arc<EventFd>>,
    /// Event to kick output
    kick_out_evt: Arc<EventFd>,
    /// output buffer queue
    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);
                // SAFETY: master was created in the function of set_pty_raw_mode,
                // the value can be guaranteed to be legal.
                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);

                    // add file to temporary pool, so it could be cleaned when vm exit.
                    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 {
        // set_receiver() will unpause rx
        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) {
        // Receiver calls this if it returned 0 from remain_size()
        // and now it's ready to accept rx-data again
        if self.input.is_none() {
            error!("unpause called for non-initialized device \'{}\'", &self.id);
            return;
        }
        if self.unpause_timer.is_some() {
            return; // already set
        }

        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(())
}

// If write is blocked, return true. Otherwise return false.
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))
}

// Notification handling in case of stdio or pty usage.
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() {
        // Method `realize` expected to be called before we get here because to build event
        // notifier we need already valid file descriptors here.
        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(); // it will be rescheduled if needed

        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],
    ))
}

// Notification handling in case of pipe usage.
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 {
        // Keep the output fd registered so EPOLLOUT can be toggled with
        // AddEvents/DeleteEvents. HANG_UP is not used to track the FIFO peer.
        notifiers.push(EventNotifier::new(
            NotifierOperation::AddShared,
            output_fd,
            None,
            EventSet::HANG_UP,
            vec![outavail_handler],
        ));
    }

    notifiers
}

// Notification handling in case of listening (server) socket.
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() {
        // Method `realize` expected to be called before we get here because to build event
        // notifier we need already valid file descriptors here.
        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);
            }

            // Note: we use stream_arc variable here because we want to capture it and prolongate
            // its lifetime with this notifier callback lifetime. It allows us to ensure
            // that socket fd be valid until we unregister it from epoll_fd subscription.
            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])
    }
}

/// Provide backend trait object receiving the input from the guest.
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),
        }
    }
}

/// Provide backend trait object processing the output from the guest.
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>);
// SAFETY: Send and Sync is auto-implemented for Sender<T>,
// implementing them for SenderWrapper is safe too.
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();
    }
}