diff --git a/target/ast10x0/tests/util_ipc/async_transaction/BUILD.bazel b/target/ast10x0/tests/util_ipc/async_transaction/BUILD.bazel index b2a0e272..9574b5ac 100644 --- a/target/ast10x0/tests/util_ipc/async_transaction/BUILD.bazel +++ b/target/ast10x0/tests/util_ipc/async_transaction/BUILD.bazel @@ -84,6 +84,7 @@ rust_app( target_compatible_with = TARGET_COMPATIBLE_WITH, deps = [ "//util/ipc", + "//util/service", "@pigweed//pw_kernel/userspace", "@pigweed//pw_log/rust:pw_log", "@pigweed//pw_status/rust:pw_status", diff --git a/target/ast10x0/tests/util_ipc/async_transaction/initiator_main.rs b/target/ast10x0/tests/util_ipc/async_transaction/initiator_main.rs index f9a7c29c..3bf1096e 100644 --- a/target/ast10x0/tests/util_ipc/async_transaction/initiator_main.rs +++ b/target/ast10x0/tests/util_ipc/async_transaction/initiator_main.rs @@ -18,6 +18,10 @@ //! | double start | `start` while pending | buffers returned | //! | try_recv early | `try_recv` before the response | `Ok(None)` | //! | send_len too big | `start` past the end of `send` | `OutOfRange` | +//! | transport roundtrip| `AsyncChannelTransport` start/poll | byte incremented | +//! | transport not ready| `poll` before the response | `Ok(None)` | +//! | transport too large| request past the send buffer | `TooLarge` | +//! | transport cancel | `cancel` then reuse | channel freed | #![no_main] #![no_std] @@ -27,7 +31,8 @@ use pw_status::{Error, Result}; use userspace::entry; use userspace::syscall::{self, Signals}; use userspace::time::{Clock, Instant, SystemClock}; -use util_ipc::{AsyncTransaction, IpcHandle, IpcInitiator}; +use util_ipc::{AsyncChannelTransport, AsyncTransaction, IpcHandle, IpcInitiator}; +use util_service::{AsyncTransport, TransportError}; static mut SEND_BUF: [u8; 1] = [0x10]; static mut RECV_BUF: [u8; 1] = [0u8; 1]; @@ -273,6 +278,139 @@ fn test_send_len_out_of_range() -> Result<()> { } } +/// Buffers the transport lends the kernel. One byte each, matching the +/// handler's request and response. +static mut TXN_SEND: [u8; 1] = [0u8; 1]; +static mut TXN_RECV: [u8; 1] = [0u8; 1]; + +/// # Safety +/// Only called from this single-threaded app, and only while no +/// `AsyncChannelTransport` still holds a prior borrow. +unsafe fn transport() -> AsyncChannelTransport { + // Safety: see function doc. + unsafe { + AsyncChannelTransport::new( + IpcHandle::new(handle::IPC), + &mut *core::ptr::addr_of_mut!(TXN_SEND), + &mut *core::ptr::addr_of_mut!(TXN_RECV), + ) + } +} + +/// The seam end to end: a request goes out as plain bytes, the response +/// comes back into the caller's own buffer, and the transport is idle +/// afterwards so the next request can reuse it. +fn test_transport_roundtrip() -> Result<()> { + // Safety: no other transport is live right now. + let mut transport = unsafe { transport() }; + let mut resp = [0u8; 1]; + + transport.start(&[0x10]).map_err(|_| Error::Internal)?; + syscall::object_wait(transport.as_raw(), Signals::READABLE, Instant::MAX)?; + + let Some(len) = transport.poll(&mut resp).map_err(|_| Error::Internal)? else { + pw_log::error!("transport roundtrip: not ready after READABLE"); + return Err(Error::Internal); + }; + if len != 1 || resp[0] != 0x11 { + pw_log::error!("transport roundtrip: unexpected response"); + return Err(Error::Internal); + } + + // Idle again: a second round-trip reuses the same buffers. + transport.start(&[0x10]).map_err(|_| Error::Internal)?; + syscall::object_wait(transport.as_raw(), Signals::READABLE, Instant::MAX)?; + transport.poll(&mut resp).map_err(|_| Error::Internal)?; + Ok(()) +} + +/// The case a loopback cannot reach: poll returns Ok(None) while the +/// handler is parked, then Ok(Some) once it answers. +fn test_transport_not_ready_then_ready() -> Result<()> { + let ipc = IpcHandle::new(handle::IPC); + if syscall::object_wait(ipc.as_raw(), Signals::USER, SystemClock::now()).is_ok() { + pw_log::error!("transport not ready: stale parked signal"); + return Err(Error::Internal); + } + + // Safety: no other transport is live right now. + let mut transport = unsafe { transport() }; + let mut resp = [0u8; 1]; + + transport.start(&[0x40]).map_err(|_| Error::Internal)?; + syscall::object_wait(transport.as_raw(), Signals::USER, Instant::MAX)?; + + match transport.poll(&mut resp) { + Ok(None) => {} + Ok(Some(_)) => { + pw_log::error!("transport not ready: completed before the handler responded"); + return Err(Error::Internal); + } + Err(_) => { + pw_log::error!("transport not ready: poll failed"); + return Err(Error::Internal); + } + } + + ipc.set_peer_user_signal(true)?; + syscall::object_wait(transport.as_raw(), Signals::READABLE, Instant::MAX)?; + let ready = transport.poll(&mut resp).map_err(|_| Error::Internal)?; + ipc.set_peer_user_signal(false)?; + + if ready != Some(1) || resp[0] != 0x41 { + pw_log::error!("transport not ready: unexpected response"); + return Err(Error::Internal); + } + Ok(()) +} + +/// A request longer than the send buffer is refused before anything goes +/// out, and the transport stays idle. +fn test_transport_request_too_large() -> Result<()> { + // Safety: no other transport is live right now. + let mut transport = unsafe { transport() }; + + match transport.start(&[0x10, 0x20]) { + Err(TransportError::TooLarge) => {} + Err(_) => { + pw_log::error!("transport too large: wrong error"); + return Err(Error::Internal); + } + Ok(()) => { + pw_log::error!("transport too large: start() succeeded"); + return Err(Error::Internal); + } + } + + // Still idle: a correctly sized request goes through. + transport.start(&[0x10]).map_err(|_| Error::Internal)?; + transport.cancel().map_err(|_| Error::Internal)?; + Ok(()) +} + +/// Cancel returns the buffers to the transport, so the next start works +/// and a poll with nothing in flight is WrongState. +fn test_transport_cancel() -> Result<()> { + // Safety: no other transport is live right now. + let mut transport = unsafe { transport() }; + let mut resp = [0u8; 1]; + + transport.start(&[0x10]).map_err(|_| Error::Internal)?; + transport.cancel().map_err(|_| Error::Internal)?; + + match transport.poll(&mut resp) { + Err(TransportError::WrongState) => {} + _ => { + pw_log::error!("transport cancel: poll after cancel was not WrongState"); + return Err(Error::Internal); + } + } + + transport.start(&[0x10]).map_err(|_| Error::Internal)?; + transport.cancel().map_err(|_| Error::Internal)?; + Ok(()) +} + /// How often the whole sequence runs. Every case leaves the channel idle /// and both USER signals lowered, so a repeat that fails means state leaked /// from the pass before it. @@ -286,7 +424,11 @@ fn run_all() -> Result<()> { .and_then(|_| test_drop_cancels()) .and_then(|_| test_double_start()) .and_then(|_| test_try_recv_before_response()) - .and_then(|_| test_send_len_out_of_range()); + .and_then(|_| test_send_len_out_of_range()) + .and_then(|_| test_transport_roundtrip()) + .and_then(|_| test_transport_not_ready_then_ready()) + .and_then(|_| test_transport_request_too_large()) + .and_then(|_| test_transport_cancel()); if ret.is_err() { pw_log::error!("failed in round {}", round as u32); diff --git a/util/ipc/BUILD.bazel b/util/ipc/BUILD.bazel index d684ccda..55a32387 100644 --- a/util/ipc/BUILD.bazel +++ b/util/ipc/BUILD.bazel @@ -7,6 +7,7 @@ rust_library( name = "ipc", srcs = [ "async_transaction.rs", + "channel_transport.rs", "lib.rs", "target.rs", ], @@ -19,6 +20,7 @@ rust_library( }), visibility = ["//visibility:public"], deps = [ + "//util/service", "@pigweed//pw_kernel/userspace", "@pigweed//pw_status/rust:pw_status", ], diff --git a/util/ipc/channel_transport.rs b/util/ipc/channel_transport.rs new file mode 100644 index 00000000..d1661f73 --- /dev/null +++ b/util/ipc/channel_transport.rs @@ -0,0 +1,121 @@ +// Licensed under the Apache-2.0 license +// SPDX-License-Identifier: Apache-2.0 + +//! `util_service::AsyncTransport` over a kernel channel. +//! +//! `AsyncTransaction` lends the kernel `'static` buffers for the duration of +//! a transaction, which a caller holding an ordinary `&[u8]` request cannot +//! satisfy. This type owns that pair of buffers, copies each request in and +//! each response out, and hands callers the plain-slice seam every service +//! shares. +//! +//! Buffer sizes are the wiring's choice: a request longer than the send +//! buffer, or a response longer than the caller's, is `TooLarge` rather than +//! a truncated frame. + +use util_service::{AsyncTransport, TransportError}; + +use super::async_transaction::{AsyncTransaction, Buffers}; +use super::IpcInitiator; + +/// One channel's worth of async round-trip, with the buffers it lends the +/// kernel. +pub struct AsyncChannelTransport { + txn: AsyncTransaction, + /// Held while idle, lent to `txn` while a round-trip is in flight. + idle: Option, +} + +impl AsyncChannelTransport { + /// Wrap an initiator with the buffers it lends the kernel. `send` must + /// fit the largest request this channel carries, `recv` the largest + /// response. + pub fn new(handle: H, send: &'static mut [u8], recv: &'static mut [u8]) -> Self { + Self { + txn: AsyncTransaction::new(handle), + idle: Some(Buffers { send, recv }), + } + } + + /// The raw channel handle, to register with a WaitGroup so the event + /// loop wakes when the response lands. + pub fn as_raw(&self) -> u32 { + self.txn.as_raw() + } + + /// Whether a round-trip is in flight. + pub fn is_pending(&self) -> bool { + self.txn.is_pending() + } +} + +impl AsyncTransport for AsyncChannelTransport { + fn start(&mut self, req: &[u8]) -> Result<(), TransportError> { + let Some(buffers) = self.idle.take() else { + return Err(TransportError::WrongState); + }; + let Buffers { send, recv } = buffers; + + if req.len() > send.len() { + self.idle = Some(Buffers { send, recv }); + return Err(TransportError::TooLarge); + } + send[..req.len()].copy_from_slice(req); + + match self.txn.start(send, req.len(), recv) { + Ok(()) => Ok(()), + Err(e) => { + // The buffers came back, so the next start can reuse them. + self.idle = Some(Buffers { + send: e.send, + recv: e.recv, + }); + Err(TransportError::Failed) + } + } + } + + fn poll(&mut self, resp: &mut [u8]) -> Result, TransportError> { + if self.idle.is_some() { + return Err(TransportError::WrongState); + } + + match self.txn.try_recv() { + Ok(None) => Ok(None), + Ok(Some(completion)) => { + let len = completion.len; + let too_large = len > resp.len(); + if !too_large { + resp[..len].copy_from_slice(&completion.recv[..len]); + } + self.idle = Some(Buffers { + send: completion.send, + recv: completion.recv, + }); + if too_large { + return Err(TransportError::TooLarge); + } + Ok(Some(len)) + } + Err(e) => { + // try_recv hands the buffers out on every failure, so the + // channel is idle again and the next call is start. + self.idle = e.buffers; + Err(TransportError::Failed) + } + } + } + + fn cancel(&mut self) -> Result<(), TransportError> { + if self.idle.is_some() { + return Err(TransportError::WrongState); + } + match self.txn.cancel() { + Ok(buffers) => { + self.idle = Some(buffers); + Ok(()) + } + Err(_) => Err(TransportError::Failed), + } + } +} diff --git a/util/ipc/lib.rs b/util/ipc/lib.rs index eebaab29..64430647 100644 --- a/util/ipc/lib.rs +++ b/util/ipc/lib.rs @@ -96,7 +96,9 @@ impl IpcHandle { } mod async_transaction; +mod channel_transport; mod target; pub use async_transaction::{AsyncTransaction, Buffers, Completion, RecvError, StartError}; +pub use channel_transport::AsyncChannelTransport; pub use target::{AsSyscallBuffer, Instant};