Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
146 changes: 144 additions & 2 deletions target/ast10x0/tests/util_ipc/async_transaction/initiator_main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand All @@ -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];
Expand Down Expand Up @@ -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<IpcHandle> {
// 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.
Expand All @@ -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);
Expand Down
2 changes: 2 additions & 0 deletions util/ipc/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ rust_library(
name = "ipc",
srcs = [
"async_transaction.rs",
"channel_transport.rs",
"lib.rs",
"target.rs",
],
Expand All @@ -19,6 +20,7 @@ rust_library(
}),
visibility = ["//visibility:public"],
deps = [
"//util/service",
"@pigweed//pw_kernel/userspace",
"@pigweed//pw_status/rust:pw_status",
],
Expand Down
121 changes: 121 additions & 0 deletions util/ipc/channel_transport.rs
Original file line number Diff line number Diff line change
@@ -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<H: IpcInitiator> {
txn: AsyncTransaction<H>,
/// Held while idle, lent to `txn` while a round-trip is in flight.
idle: Option<Buffers>,
}

impl<H: IpcInitiator> AsyncChannelTransport<H> {
/// 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<H: IpcInitiator> AsyncTransport for AsyncChannelTransport<H> {
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<Option<usize>, 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),
}
}
}
2 changes: 2 additions & 0 deletions util/ipc/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};