use std::{
borrow::Cow,
future::Future,
io,
sync::{Arc, Mutex},
};
use remoc::prelude::ServerShared;
use serde::{Deserialize, Serialize};
use smallvec::smallvec;
use tokio::net::UnixStream;
use tokio_util::task::AbortOnDropHandle;
use tracing::{Instrument, debug};
use crate::{
error::Code,
ipc::{
error::{IpcAcceptError, IpcOpenError, IpcPlumbingError},
quic::stream::{
IpcBiHandle, IpcUniHandle,
reader::{self as ipc_reader, IpcReadHypervisorIo},
writer::{self as ipc_writer, IpcWriteHypervisorIo},
},
transport::{FdTransfer, ReceivedFds, ReservedFdDelivery},
},
quic::{
self, BoxQuicStreamReader, BoxQuicStreamWriter, ConnectionError, GetStreamIdExt,
ManageStream, ReadStream, ResetStreamExt, StopStreamExt, StreamError, WriteStream,
},
rpc::{
lifecycle::{ConnectionErrorLatch, HasLatch, LifecycleExt},
quic::{
CachedLocalAuthority, CachedRemoteAuthority, LocalAuthorityClient,
LocalAuthorityServerShared, RemoteAuthorityClient, RemoteAuthorityServerShared,
},
},
util::deferred::Resolved,
varint::VarInt,
};
#[remoc::rtc::remote]
pub trait IpcConnection: Send + Sync {
async fn open_bi(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcBiHandle, StreamError>, IpcOpenError>;
async fn accept_bi(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcBiHandle, StreamError>, IpcAcceptError>;
async fn open_uni(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcUniHandle, StreamError>, IpcOpenError>;
async fn accept_uni(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcUniHandle, StreamError>, IpcAcceptError>;
async fn local_authority(&self) -> Result<Option<LocalAuthorityClient>, ConnectionError>;
async fn remote_authority(&self) -> Result<Option<RemoteAuthorityClient>, ConnectionError>;
async fn close(&self, code: Code, reason: Cow<'static, str>) -> Result<(), ConnectionError>;
async fn closed(&self) -> Result<ConnectionError, ConnectionError>;
}
#[derive(Serialize, Deserialize)]
pub struct ConnectionBootstrap {
pub connection: IpcConnectionClient,
}
pub async fn bridge_reader(quic_reader: impl ReadStream + Unpin, pipe: UnixStream) {
crate::rpc::stream::hypervisor::read::run_read_bridge(
quic_reader,
IpcReadHypervisorIo::new(pipe),
)
.await;
}
pub async fn bridge_writer(pipe: UnixStream, quic_writer: impl WriteStream + Unpin) {
crate::rpc::stream::hypervisor::write::run_write_bridge(
quic_writer,
IpcWriteHypervisorIo::new(pipe),
)
.await;
}
pub struct ConnectionAdapter<M> {
inner: Arc<M>,
fd_transfer: FdTransfer,
tasks: Mutex<Vec<AbortOnDropHandle<()>>>,
}
impl<M> ConnectionAdapter<M> {
pub fn new(inner: Arc<M>, fd_transfer: FdTransfer) -> Self {
Self {
inner,
fd_transfer,
tasks: Mutex::new(Vec::new()),
}
}
fn spawn_task(
&self,
task_kind: &'static str,
stream_id: Option<VarInt>,
task: impl Future<Output = ()> + Send + 'static,
) {
let handle = AbortOnDropHandle::new(tokio::spawn(task.in_current_span()));
let mut tasks = self
.tasks
.lock()
.expect("connection adapter task registry should not be poisoned");
let before = tasks.len();
tasks.retain(|task| !task.is_finished());
let reaped = before - tasks.len();
tasks.push(handle);
tracing::trace!(
boundary = "ipc-root-task-registry",
task_kind,
stream_id = ?stream_id.map(|id| id.into_inner()),
active_tasks = tasks.len(),
reaped_tasks = reaped,
"IPC root bridge task registered"
);
}
}
impl<M> IpcConnection for ConnectionAdapter<M>
where
M: ManageStream
+ quic::Lifecycle
+ quic::WithLocalAuthority
+ quic::WithRemoteAuthority
+ Send
+ Sync
+ 'static,
M::StreamReader: Unpin + 'static,
M::StreamWriter: Unpin + 'static,
M::LocalAuthority: Send + Sync,
M::RemoteAuthority: Send + Sync,
{
async fn open_bi(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcBiHandle, StreamError>, IpcOpenError> {
self.open_bi_impl(fd_id).await
}
async fn accept_bi(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcBiHandle, StreamError>, IpcAcceptError> {
self.accept_bi_impl(fd_id).await
}
async fn open_uni(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcUniHandle, StreamError>, IpcOpenError> {
self.open_uni_impl(fd_id).await
}
async fn accept_uni(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcUniHandle, StreamError>, IpcAcceptError> {
self.accept_uni_impl(fd_id).await
}
async fn local_authority(&self) -> Result<Option<LocalAuthorityClient>, ConnectionError> {
match quic::WithLocalAuthority::local_authority(self.inner.as_ref()).await? {
Some(agent) => {
let (server, client) = LocalAuthorityServerShared::new(Arc::new(agent), 1);
self.spawn_task("local-authority", None, async move {
let _ = server.serve(true).await;
});
Ok(Some(client))
}
None => Ok(None),
}
}
async fn remote_authority(&self) -> Result<Option<RemoteAuthorityClient>, ConnectionError> {
match quic::WithRemoteAuthority::remote_authority(self.inner.as_ref()).await? {
Some(agent) => {
let (server, client) = RemoteAuthorityServerShared::new(Arc::new(agent), 1);
self.spawn_task("remote-authority", None, async move {
let _ = server.serve(true).await;
});
Ok(Some(client))
}
None => Ok(None),
}
}
async fn close(&self, code: Code, reason: Cow<'static, str>) -> Result<(), ConnectionError> {
quic::Lifecycle::close(self.inner.as_ref(), code, reason);
Ok(())
}
async fn closed(&self) -> Result<ConnectionError, ConnectionError> {
Ok(quic::Lifecycle::closed(self.inner.as_ref()).await)
}
}
impl<M> ConnectionAdapter<M>
where
M: ManageStream
+ quic::Lifecycle
+ quic::WithLocalAuthority
+ quic::WithRemoteAuthority
+ Send
+ Sync
+ 'static,
M::StreamReader: Unpin + 'static,
M::StreamWriter: Unpin + 'static,
M::LocalAuthority: Send + Sync,
M::RemoteAuthority: Send + Sync,
{
async fn open_bi_impl(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcBiHandle, StreamError>, IpcOpenError> {
let delivery = self
.fd_transfer
.delivery(fd_id)
.reserve()
.await
.map_err(|error| ipc_io_plumbing(error, "reserve fd delivery"))
.map_err(IpcOpenError::from)?;
let (mut reader, writer) = ManageStream::open_bi(self.inner.as_ref()).await?;
let stream_id = match reader.stream_id().await {
Ok(id) => id,
Err(stream_err) => return Ok(Resolved::err(stream_err)),
};
self.bridge_bi(delivery, reader, writer, stream_id)
.await
.map(Resolved::ok)
.map_err(IpcOpenError::from)
}
async fn accept_bi_impl(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcBiHandle, StreamError>, IpcAcceptError> {
let delivery = self
.fd_transfer
.delivery(fd_id)
.reserve()
.await
.map_err(|error| ipc_io_plumbing(error, "reserve fd delivery"))
.map_err(IpcAcceptError::from)?;
let (mut reader, writer) = ManageStream::accept_bi(self.inner.as_ref()).await?;
let stream_id = match reader.stream_id().await {
Ok(id) => id,
Err(stream_err) => return Ok(Resolved::err(stream_err)),
};
self.bridge_bi(delivery, reader, writer, stream_id)
.await
.map(Resolved::ok)
.map_err(IpcAcceptError::from)
}
async fn bridge_bi(
&self,
delivery: ReservedFdDelivery,
mut reader: M::StreamReader,
mut writer: M::StreamWriter,
stream_id: VarInt,
) -> Result<IpcBiHandle, IpcPlumbingError> {
let fd_id = delivery.id();
let (srv_a, cli_a) = UnixStream::pair().map_err(|e| ipc_io_plumbing(e, "socketpair"))?;
let (srv_b, cli_b) = UnixStream::pair().map_err(|e| ipc_io_plumbing(e, "socketpair"))?;
let cli_a_std = cli_a
.into_std()
.map_err(|e| ipc_io_plumbing(e, "into_std"))?;
let cli_b_std = cli_b
.into_std()
.map_err(|e| ipc_io_plumbing(e, "into_std"))?;
if let Err(error) = delivery
.deliver(smallvec![cli_a_std.into(), cli_b_std.into()])
.await
{
let code = Code::H3_REQUEST_CANCELLED.into_inner();
let _ = reader.stop(code).await;
let _ = writer.reset(code).await;
return Err(ipc_io_plumbing(error, "deliver fds"));
}
tracing::trace!(
boundary = "ipc-root-fd-delivery",
fd_id = fd_id.into_inner(),
stream_id = stream_id.into_inner(),
fd_count = 2,
"IPC bidirectional stream FDs queued"
);
self.spawn_task("stream-reader", Some(stream_id), async move {
bridge_reader(reader, srv_a).await;
tracing::trace!(
boundary = "ipc-root-task-registry",
stream_id = stream_id.into_inner(),
"IPC root stream reader bridge task finished"
);
});
self.spawn_task("stream-writer", Some(stream_id), async move {
bridge_writer(srv_b, writer).await;
tracing::trace!(
boundary = "ipc-root-task-registry",
stream_id = stream_id.into_inner(),
"IPC root stream writer bridge task finished"
);
});
Ok(IpcBiHandle { stream_id })
}
async fn open_uni_impl(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcUniHandle, StreamError>, IpcOpenError> {
let delivery = self
.fd_transfer
.delivery(fd_id)
.reserve()
.await
.map_err(|error| ipc_io_plumbing(error, "reserve fd delivery"))
.map_err(IpcOpenError::from)?;
let mut writer = ManageStream::open_uni(self.inner.as_ref()).await?;
let stream_id = match writer.stream_id().await {
Ok(id) => id,
Err(stream_err) => return Ok(Resolved::err(stream_err)),
};
let (srv, cli) = UnixStream::pair().map_err(|e| ipc_io_plumbing(e, "socketpair"))?;
let cli_std = cli.into_std().map_err(|e| ipc_io_plumbing(e, "into_std"))?;
if let Err(error) = delivery.deliver(smallvec![cli_std.into()]).await {
let _ = writer.reset(Code::H3_REQUEST_CANCELLED.into_inner()).await;
return Err(IpcOpenError::from(ipc_io_plumbing(error, "deliver fds")));
}
self.spawn_task("stream-writer", Some(stream_id), async move {
bridge_writer(srv, writer).await;
tracing::trace!(
boundary = "ipc-root-task-registry",
stream_id = stream_id.into_inner(),
"IPC root stream writer bridge task finished"
);
});
Ok(Resolved::ok(IpcUniHandle { stream_id }))
}
async fn accept_uni_impl(
&self,
fd_id: VarInt,
) -> Result<Resolved<IpcUniHandle, StreamError>, IpcAcceptError> {
let delivery = self
.fd_transfer
.delivery(fd_id)
.reserve()
.await
.map_err(|error| ipc_io_plumbing(error, "reserve fd delivery"))
.map_err(IpcAcceptError::from)?;
let mut reader = ManageStream::accept_uni(self.inner.as_ref()).await?;
let stream_id = match reader.stream_id().await {
Ok(id) => id,
Err(stream_err) => return Ok(Resolved::err(stream_err)),
};
let (srv, cli) = UnixStream::pair().map_err(|e| ipc_io_plumbing(e, "socketpair"))?;
let cli_std = cli.into_std().map_err(|e| ipc_io_plumbing(e, "into_std"))?;
if let Err(error) = delivery.deliver(smallvec![cli_std.into()]).await {
let _ = reader.stop(Code::H3_REQUEST_CANCELLED.into_inner()).await;
return Err(IpcAcceptError::from(ipc_io_plumbing(error, "deliver fds")));
}
self.spawn_task("stream-reader", Some(stream_id), async move {
bridge_reader(reader, srv).await;
tracing::trace!(
boundary = "ipc-root-task-registry",
stream_id = stream_id.into_inner(),
"IPC root stream reader bridge task finished"
);
});
Ok(Resolved::ok(IpcUniHandle { stream_id }))
}
}
pub struct IpcConnectionHandle {
rpc: IpcConnectionClient,
fd_transfer: FdTransfer,
_remoc_task: AbortOnDropHandle<()>,
lifecycle: Arc<IpcLifecycle>,
}
struct IpcLifecycle {
connection: IpcConnectionClient,
latch: ConnectionErrorLatch,
close_tasks: Mutex<Vec<AbortOnDropHandle<()>>>,
}
impl HasLatch for IpcLifecycle {
fn latch(&self) -> &ConnectionErrorLatch {
&self.latch
}
}
impl quic::Lifecycle for IpcLifecycle {
fn close(&self, code: Code, reason: Cow<'static, str>) {
let rpc = self.connection.clone();
let handle = AbortOnDropHandle::new(tokio::spawn(
(async move {
let _ = IpcConnection::close(&rpc, code, reason).await;
})
.in_current_span(),
));
let mut tasks = self
.close_tasks
.lock()
.expect("ipc lifecycle close task registry should not be poisoned");
tasks.retain(|task| !task.is_finished());
tasks.push(handle);
}
fn check(&self) -> Result<(), ConnectionError> {
self.check_with_probe(|| {
remoc::rtc::Client::is_closed(&self.connection)
.then(IpcConnectionHandle::ipc_channel_error)
})
}
async fn closed(&self) -> ConnectionError {
self.resolve_closed(async {
IpcConnection::closed(&self.connection)
.await
.unwrap_or_else(|_| IpcConnectionHandle::ipc_channel_error())
})
.await
}
}
impl IpcConnectionHandle {
pub fn new(
rpc: IpcConnectionClient,
fd_transfer: FdTransfer,
remoc_task: AbortOnDropHandle<()>,
) -> Self {
let lifecycle = Arc::new(IpcLifecycle {
connection: rpc.clone(),
latch: ConnectionErrorLatch::new(),
close_tasks: Mutex::new(Vec::new()),
});
Self {
rpc,
fd_transfer,
_remoc_task: remoc_task,
lifecycle,
}
}
fn ipc_channel_error() -> ConnectionError {
quic::ConnectionError::Transport {
source: quic::TransportError {
kind: IPC_CHANNEL_ERROR_KIND,
frame_type: IPC_FRAME_TYPE,
reason: "ipc connection channel closed".into(),
},
}
}
}
impl quic::ManageStream for IpcConnectionHandle {
type StreamReader = Resolved<BoxQuicStreamReader, StreamError>;
type StreamWriter = Resolved<BoxQuicStreamWriter, StreamError>;
async fn open_bi(&self) -> Result<(Self::StreamReader, Self::StreamWriter), ConnectionError> {
let (resolved, received) = self.open_bi_with_fds().await?;
match resolved {
Resolved::Value { value: handle } => {
let received = received.expect("value response must include received fds");
let (r, w) = self.fds_to_bi(handle, received).await?;
Ok((Resolved::ok(r), Resolved::ok(w)))
}
Resolved::Error { error } => Ok((Resolved::err(error.clone()), Resolved::err(error))),
}
}
async fn accept_bi(&self) -> Result<(Self::StreamReader, Self::StreamWriter), ConnectionError> {
let (resolved, received) = self.accept_bi_with_fds().await?;
match resolved {
Resolved::Value { value: handle } => {
let received = received.expect("value response must include received fds");
let (r, w) = self.fds_to_bi(handle, received).await?;
Ok((Resolved::ok(r), Resolved::ok(w)))
}
Resolved::Error { error } => Ok((Resolved::err(error.clone()), Resolved::err(error))),
}
}
async fn open_uni(&self) -> Result<Self::StreamWriter, ConnectionError> {
let (resolved, received) = self.open_uni_with_fds().await?;
match resolved {
Resolved::Value { value: handle } => {
let received = received.expect("value response must include received fds");
let w = self.fds_to_uni_writer(handle, received).await?;
Ok(Resolved::ok(w))
}
Resolved::Error { error } => Ok(Resolved::err(error)),
}
}
async fn accept_uni(&self) -> Result<Self::StreamReader, ConnectionError> {
let (resolved, received) = self.accept_uni_with_fds().await?;
match resolved {
Resolved::Value { value: handle } => {
let received = received.expect("value response must include received fds");
let r = self.fds_to_uni_reader(handle, received).await?;
Ok(Resolved::ok(r))
}
Resolved::Error { error } => Ok(Resolved::err(error)),
}
}
}
impl IpcConnectionHandle {
async fn open_bi_with_fds(
&self,
) -> Result<(Resolved<IpcBiHandle, StreamError>, Option<ReceivedFds>), ConnectionError> {
let receiver = self.fd_transfer.receive();
let fd_id = receiver.id();
self.lifecycle
.guard(self.resolve_open_with_fds(receiver, IpcConnection::open_bi(&self.rpc, fd_id)))
.await
}
async fn accept_bi_with_fds(
&self,
) -> Result<(Resolved<IpcBiHandle, StreamError>, Option<ReceivedFds>), ConnectionError> {
let receiver = self.fd_transfer.receive();
let fd_id = receiver.id();
self.lifecycle
.guard(
self.resolve_accept_with_fds(receiver, IpcConnection::accept_bi(&self.rpc, fd_id)),
)
.await
}
async fn open_uni_with_fds(
&self,
) -> Result<(Resolved<IpcUniHandle, StreamError>, Option<ReceivedFds>), ConnectionError> {
let receiver = self.fd_transfer.receive();
let fd_id = receiver.id();
self.lifecycle
.guard(self.resolve_open_with_fds(receiver, IpcConnection::open_uni(&self.rpc, fd_id)))
.await
}
async fn accept_uni_with_fds(
&self,
) -> Result<(Resolved<IpcUniHandle, StreamError>, Option<ReceivedFds>), ConnectionError> {
let receiver = self.fd_transfer.receive();
let fd_id = receiver.id();
self.lifecycle
.guard(
self.resolve_accept_with_fds(receiver, IpcConnection::accept_uni(&self.rpc, fd_id)),
)
.await
}
async fn resolve_open_with_fds<H>(
&self,
receiver: crate::ipc::transport::FdReceiver,
rpc: impl Future<Output = Result<Resolved<H, StreamError>, IpcOpenError>>,
) -> Result<(Resolved<H, StreamError>, Option<ReceivedFds>), ConnectionError> {
let receive = receiver.into_future();
tokio::pin!(receive);
tokio::pin!(rpc);
tokio::select! {
biased;
receive_result = &mut receive => {
let received = receive_result.map_err(|e| ipc_transport_error(e, "receive fds"))?;
let resolved = rpc.await.map_err(map_open_err)?;
Ok((resolved, Some(received)))
}
rpc_result = &mut rpc => {
let resolved = rpc_result.map_err(map_open_err)?;
match resolved {
Resolved::Value { value } => {
let received = receive.await.map_err(|e| ipc_transport_error(e, "receive fds"))?;
Ok((Resolved::ok(value), Some(received)))
}
Resolved::Error { error } => Ok((Resolved::err(error), None)),
}
}
}
}
async fn resolve_accept_with_fds<H>(
&self,
receiver: crate::ipc::transport::FdReceiver,
rpc: impl Future<Output = Result<Resolved<H, StreamError>, IpcAcceptError>>,
) -> Result<(Resolved<H, StreamError>, Option<ReceivedFds>), ConnectionError> {
let receive = receiver.into_future();
tokio::pin!(receive);
tokio::pin!(rpc);
tokio::select! {
biased;
receive_result = &mut receive => {
let received = receive_result.map_err(|e| ipc_transport_error(e, "receive fds"))?;
let resolved = rpc.await.map_err(map_accept_err)?;
Ok((resolved, Some(received)))
}
rpc_result = &mut rpc => {
let resolved = rpc_result.map_err(map_accept_err)?;
match resolved {
Resolved::Value { value } => {
let received = receive.await.map_err(|e| ipc_transport_error(e, "receive fds"))?;
Ok((Resolved::ok(value), Some(received)))
}
Resolved::Error { error } => Ok((Resolved::err(error), None)),
}
}
}
}
async fn fds_to_bi(
&self,
handle: IpcBiHandle,
received: ReceivedFds,
) -> Result<(BoxQuicStreamReader, BoxQuicStreamWriter), ConnectionError> {
let IpcBiHandle { stream_id } = handle;
self.lifecycle.guard_sync(|| {
let (fd_a, fd_b) = received
.into_pair()
.map_err(|e| ipc_transport_error(e, "fd count"))?;
let lifecycle = self.lifecycle.clone();
let sock_a = UnixStream::from_std(std::os::unix::net::UnixStream::from(fd_a))
.map_err(|e| ipc_io_error(e, "UnixStream::from_std"))?;
let reader = Box::pin(ipc_reader::reader(stream_id, sock_a, lifecycle.clone()))
as BoxQuicStreamReader;
let sock_b = UnixStream::from_std(std::os::unix::net::UnixStream::from(fd_b))
.map_err(|e| ipc_io_error(e, "UnixStream::from_std"))?;
let writer =
Box::pin(ipc_writer::writer(stream_id, sock_b, lifecycle)) as BoxQuicStreamWriter;
Ok((reader, writer))
})
}
async fn fds_to_uni_writer(
&self,
handle: IpcUniHandle,
received: ReceivedFds,
) -> Result<BoxQuicStreamWriter, ConnectionError> {
let IpcUniHandle { stream_id } = handle;
self.lifecycle.guard_sync(|| {
let fd = received
.into_one()
.map_err(|e| ipc_transport_error(e, "fd count"))?;
let lifecycle = self.lifecycle.clone();
let sock = UnixStream::from_std(std::os::unix::net::UnixStream::from(fd))
.map_err(|e| ipc_io_error(e, "UnixStream::from_std"))?;
Ok(Box::pin(ipc_writer::writer(stream_id, sock, lifecycle)) as BoxQuicStreamWriter)
})
}
async fn fds_to_uni_reader(
&self,
handle: IpcUniHandle,
received: ReceivedFds,
) -> Result<BoxQuicStreamReader, ConnectionError> {
let IpcUniHandle { stream_id } = handle;
self.lifecycle.guard_sync(|| {
let fd = received
.into_one()
.map_err(|e| ipc_transport_error(e, "fd count"))?;
let lifecycle = self.lifecycle.clone();
let sock = UnixStream::from_std(std::os::unix::net::UnixStream::from(fd))
.map_err(|e| ipc_io_error(e, "UnixStream::from_std"))?;
Ok(Box::pin(ipc_reader::reader(stream_id, sock, lifecycle)) as BoxQuicStreamReader)
})
}
}
impl quic::WithLocalAuthority for IpcConnectionHandle {
type LocalAuthority = CachedLocalAuthority;
async fn local_authority(&self) -> Result<Option<CachedLocalAuthority>, ConnectionError> {
match self
.lifecycle
.guard(IpcConnection::local_authority(&self.rpc))
.await?
{
Some(client) => Ok(Some(
self.lifecycle
.guard(CachedLocalAuthority::from_client(client))
.await?,
)),
None => Ok(None),
}
}
}
impl quic::WithRemoteAuthority for IpcConnectionHandle {
type RemoteAuthority = CachedRemoteAuthority;
async fn remote_authority(&self) -> Result<Option<CachedRemoteAuthority>, ConnectionError> {
match self
.lifecycle
.guard(IpcConnection::remote_authority(&self.rpc))
.await?
{
Some(client) => Ok(Some(
self.lifecycle
.guard(CachedRemoteAuthority::from_client(client))
.await?,
)),
None => Ok(None),
}
}
}
impl quic::Lifecycle for IpcConnectionHandle {
fn close(&self, code: Code, reason: Cow<'static, str>) {
quic::Lifecycle::close(self.lifecycle.as_ref(), code, reason);
}
fn check(&self) -> Result<(), ConnectionError> {
quic::Lifecycle::check(self.lifecycle.as_ref())
}
async fn closed(&self) -> ConnectionError {
quic::Lifecycle::closed(self.lifecycle.as_ref()).await
}
}
pub(crate) const IPC_ERROR_KIND: VarInt = VarInt::from_u32(0x0a);
const IPC_CHANNEL_ERROR_KIND: VarInt = VarInt::from_u32(0x01);
pub(crate) const IPC_FRAME_TYPE: VarInt = VarInt::from_u32(0x00);
fn ipc_io_error(err: io::Error, context: &str) -> ConnectionError {
debug!(error = %snafu::Report::from_error(err), context, "ipc i/o error");
ConnectionError::Transport {
source: quic::TransportError {
kind: IPC_ERROR_KIND,
frame_type: IPC_FRAME_TYPE,
reason: format!("ipc: {context}").into(),
},
}
}
fn ipc_transport_error(err: impl std::error::Error, context: &str) -> ConnectionError {
debug!(error = %snafu::Report::from_error(&err), context, "ipc transport error");
ConnectionError::Transport {
source: quic::TransportError {
kind: IPC_ERROR_KIND,
frame_type: IPC_FRAME_TYPE,
reason: format!("ipc: {context}").into(),
},
}
}
fn ipc_io_plumbing(err: impl std::error::Error, context: &str) -> IpcPlumbingError {
debug!(error = %snafu::Report::from_error(&err), context, "ipc plumbing i/o error");
IpcPlumbingError::Io {
message: format!("{context}: {err}"),
}
}
fn plumbing_to_conn(err: &IpcPlumbingError) -> ConnectionError {
ConnectionError::Transport {
source: quic::TransportError {
kind: IPC_ERROR_KIND,
frame_type: IPC_FRAME_TYPE,
reason: err.to_string().into(),
},
}
}
fn map_open_err(error: IpcOpenError) -> ConnectionError {
match error {
IpcOpenError::Connection { source } => source,
IpcOpenError::Plumbing { source } => plumbing_to_conn(&source),
}
}
fn map_accept_err(error: IpcAcceptError) -> ConnectionError {
match error {
IpcAcceptError::Connection { source } => source,
IpcAcceptError::Plumbing { source } => plumbing_to_conn(&source),
}
}
impl IpcConnectionClient {
pub fn into_handle(
self,
fd_transfer: FdTransfer,
remoc_task: AbortOnDropHandle<()>,
) -> IpcConnectionHandle {
IpcConnectionHandle::new(self, fd_transfer, remoc_task)
}
}