use std::io;
use std::os::windows::io::{AsRawHandle, BorrowedHandle};
use std::sync::Arc;
use std::time::Duration;
use tokio::net::windows::named_pipe::NamedPipeServer;
use tokio::sync::OwnedSemaphorePermit;
use windows_sys::Win32::Storage::FileSystem::FlushFileBuffers;
use super::dispatch::Dispatcher;
struct ConnectedPipe(NamedPipeServer);
impl Drop for ConnectedPipe {
fn drop(&mut self) {
let _ = self.0.disconnect();
}
}
pub(super) async fn serve_named_pipe(
server: NamedPipeServer,
dispatcher: Arc<Dispatcher>,
permit: OwnedSemaphorePermit,
flush_timeout: Duration,
) -> io::Result<()> {
let mut server = ConnectedPipe(server);
let served = super::server::serve(&mut server.0, dispatcher).await;
let flushed = drain(&server.0, permit, flush_timeout).await;
drop(server);
served.and(flushed)
}
async fn drain(
server: &NamedPipeServer,
permit: OwnedSemaphorePermit,
flush_timeout: Duration,
) -> io::Result<()> {
let handle =
unsafe { BorrowedHandle::borrow_raw(server.as_raw_handle()) }.try_clone_to_owned()?;
let flushed = tokio::task::spawn_blocking(move || {
let _permit = permit;
if unsafe { FlushFileBuffers(handle.as_raw_handle()) } == 0 {
Err(io::Error::last_os_error())
} else {
Ok(())
}
});
tokio::time::timeout(flush_timeout, flushed)
.await
.map_err(|_| {
io::Error::new(
io::ErrorKind::TimedOut,
"control pipe reply drain timed out",
)
})?
.map_err(io::Error::other)?
}