diff --git a/Cargo.lock b/Cargo.lock index d580649dc5..3d51a3aa6a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1188,8 +1188,8 @@ version = "0.1.0" name = "toyos-blockring" version = "0.1.0" dependencies = [ - "loom", "toyos-blockhold", + "toyos-transport", ] [[package]] @@ -1431,6 +1431,14 @@ version = "0.1.0" name = "toyos-tmpdir" version = "0.1.0" +[[package]] +name = "toyos-transport" +version = "0.1.0" +dependencies = [ + "loom", + "toyos-untrusted", +] + [[package]] name = "toyos-untrusted" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index e71be2c363..e7de9c08ab 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -56,6 +56,7 @@ members = [ "toyos-symbols", "toyos-tco", "toyos-tmpdir", + "toyos-transport", "toyos-untrusted", "toyos-update", "toyos-userbound", diff --git a/issues/design-debt/the-transport-heads-orderings-have-no-oracle.md b/issues/design-debt/the-transport-heads-orderings-have-no-oracle.md new file mode 100644 index 0000000000..29f06c94e3 --- /dev/null +++ b/issues/design-debt/the-transport-heads-orderings-have-no-oracle.md @@ -0,0 +1,18 @@ +--- +status: open +kind: tooling +opened: 2026-09-27 +--- + +# The transport head's orderings have no oracle + +A consumer stores its head `Release` after it has loaded the entries below it, +and a producer loads the head `Acquire` before it writes over them +(`toyos-transport/src/queue.rs`, `Consumer::release` and `Producer::space`). +No test reds if either is `Relaxed`: what the pair forbids is load buffering — +a consumer's load of an entry reading the producer's later overwrite — which +loom does not model, and every other test runs both ends on one thread. The +tail's edge has its control (`publish-relaxed`); the head's has none. + +**Exit condition.** A control that relaxes the head's two orderings, and a +model or a run on a weakly ordered CPU that goes red under it. diff --git a/issues/filesystem/a-read-or-write-answered-lost-is-never-answered.md b/issues/filesystem/a-read-or-write-answered-lost-is-never-answered.md new file mode 100644 index 0000000000..15b6be578c --- /dev/null +++ b/issues/filesystem/a-read-or-write-answered-lost-is-never-answered.md @@ -0,0 +1,24 @@ +--- +status: assigned +kind: defect +opened: 2026-09-27 +--- + +# A read or write answered Lost is never answered + +`Client::complete` (`toyos-blockring/src/client.rs`) takes a completion's tag +off the wire before it looks at the status. A user read or write answered +`Status::Lost` is then refused as `Violation::Entry`, and the session ends; but +the tag is no longer on the wire, so `Client::session_ended` does not answer it +`Refused`, and nothing ever answers its ticket. blockd's client +(`userland/blockd/src/session.rs`) keeps that ticket in `pending` for good, +across every reconnect. That breaks the module's own "every request asked for +is answered exactly once". Only a server that breaks the protocol writes that +completion; blockd does not. + +Held by the orchestrator. + +**Exit condition.** A completion the client refuses leaves its tag answered by +the session's end — refused before the tag is taken off the wire, or answered +where it is refused — and a model or unit test in which a server answers a +write `Lost` goes red without the fix and green with it. diff --git a/src/build.rs b/src/build.rs index 34dba99c00..27a9fb9520 100644 --- a/src/build.rs +++ b/src/build.rs @@ -2758,6 +2758,7 @@ mod tests { ("toyos-sched-sim", "toyos-sched/sim/Cargo.toml"), ("toyos-proclife", "toyos-proclife/Cargo.toml"), ("toyos-blockring", "toyos-blockring/Cargo.toml"), + ("toyos-transport", "toyos-transport/Cargo.toml"), ] { let path = root.join(manifest); let text = fs::read_to_string(&path) diff --git a/src/ci.rs b/src/ci.rs index 5518b76c6b..3c56f35500 100644 --- a/src/ci.rs +++ b/src/ci.rs @@ -219,6 +219,7 @@ const SCHED_LOOM: &[&str] = &["-p", "toyos-sched-loom"]; const SCHED_SIM: &[&str] = &["-p", "toyos-sched-sim"]; const PROCLIFE: &[&str] = &["-p", "toyos-proclife"]; const BLOCKRING: &[&str] = &["-p", "toyos-blockring"]; +const TRANSPORT: &[&str] = &["-p", "toyos-transport"]; const fn red( krate: &'static [&'static str], @@ -367,8 +368,10 @@ pub(crate) const CONTROLS: &[Control] = &[ red(BLOCKRING, "mutate-no-reissue-after-loss", None, &[ "what_a_flush_calls_durable_is_on_the_medium ... FAILED", ]), - red(BLOCKRING, "mutate-ring-publish-relaxed", Some("loom_ring"), &[ - "a_published_request_is_read_whole ... FAILED", + red(TRANSPORT, "publish-relaxed", Some("loom"), &["a_published_entry_is_read_whole ... FAILED"]), + red(TRANSPORT, "no-clamp", Some("loom"), &["a_hostile_producer_yields_entries_or_a_violation ... FAILED"]), + red(TRANSPORT, "end-keeps-inflight", None, &[ + "an_end_answers_every_tag_once_and_a_late_completion_nothing ... FAILED", ]), ]; diff --git a/tests/common/blockd.rs b/tests/common/blockd.rs index 10390e75b8..40520cc34e 100644 --- a/tests/common/blockd.rs +++ b/tests/common/blockd.rs @@ -432,8 +432,12 @@ pub fn blockd_serves_partitions( Ok(()) } -/// blockd's two failures, each survived by its client. +/// blockd's two failures, each survived by its client, and a client that +/// breaks the protocol, survived by blockd. /// +/// - `hostile-head`: with a write on the device, a client moves its completion +/// ring's head a ring behind blockd's tail; the answer finds no room, blockd +/// ends that session, and serves the next. /// - `reset`: blockd withholds its second write's answer; the silence ends in /// a controller reset; the withheld write is answered not done; and the /// write acknowledged before it, which the reset may have lost, is on the @@ -454,6 +458,7 @@ pub fn blockd_survives_its_death( rust_bins: &[(String, Vec)], ) -> Result<(), String> { let (mut qemu, layout, disk, trace, before) = boot(c_bins, rust_bins, "blockd-death", &[])?; + let hostile = role(&mut qemu, "hostile-head", Duration::from_secs(120))?; let reset = role(&mut qemu, "reset", Duration::from_secs(240))?; for want in [ "blockd: WITHHELD the device's answer to a write", @@ -477,7 +482,8 @@ pub fn blockd_survives_its_death( } let tail = partclaim::shut_down(qemu); partclaim::no_panic("on the way down", &tail)?; - let mut log = Serial::named("blockd_survives_its_death", format!("{}{}", reset.serial, crash.serial)); + let mut log = + Serial::named("blockd_survives_its_death", format!("{}{}{}", hostile.serial, reset.serial, crash.serial)); log.push(&tail); log.must_be_clean()?; diff --git a/tests/toyos-rust-tests/Cargo.lock b/tests/toyos-rust-tests/Cargo.lock index bc2b2c2b2b..5879429cd2 100644 --- a/tests/toyos-rust-tests/Cargo.lock +++ b/tests/toyos-rust-tests/Cargo.lock @@ -2102,6 +2102,7 @@ name = "toyos-blockring" version = "0.1.0" dependencies = [ "toyos-blockhold", + "toyos-transport", ] [[package]] @@ -2174,6 +2175,17 @@ dependencies = [ name = "toyos-tco" version = "0.1.0" +[[package]] +name = "toyos-transport" +version = "0.1.0" +dependencies = [ + "toyos-untrusted", +] + +[[package]] +name = "toyos-untrusted" +version = "0.1.0" + [[package]] name = "toyos-wallclock" version = "0.1.0" diff --git a/tests/toyos-rust-tests/src/bin/blockd_io.rs b/tests/toyos-rust-tests/src/bin/blockd_io.rs index c80803d460..0fea755af1 100644 --- a/tests/toyos-rust-tests/src/bin/blockd_io.rs +++ b/tests/toyos-rust-tests/src/bin/blockd_io.rs @@ -15,6 +15,9 @@ //! was answered; //! - `bench` — the same bytes through the kernel's driver (a partition claim on //! the first controller) and through blockd, timed; +//! - `hostile-head` — a client that, with a write on the device, moves its +//! completion ring's head a ring behind blockd's tail: the session is +//! ended, and blockd serves the next one; //! - `reset` — blockd started withholding its second answer: the silence ends //! in a controller reset, the withheld write is answered not done, and the //! write acknowledged before it is on the medium after the next flush; @@ -32,9 +35,12 @@ use std::io::{BufRead, BufReader, Write}; use std::os::toyos::process::{ChildExt, CommandExt}; use std::process::{Child, Command, Stdio}; +use std::sync::atomic::Ordering; +use std::sync::mpsc::{self, Receiver, RecvTimeoutError}; use std::time::{Duration, Instant}; use blockd::nvme::{Controller, Owner}; +use blockd::region::Region; use blockd::{Error, Outcome, Session, Unsent}; use toyos::endow::Endowments; use toyos::poller::{Poller, READABLE}; @@ -45,8 +51,9 @@ use toyos::syscap::SysCap; use toyos::AsHandle; use toyos_abi::part::PartGuid; use toyos_abi::syscall::{DeviceType, PciId, SyscallError, DEV_PREFIX, SERVE_PREFIX, SYSCAP_LABEL}; +use toyos_blockring::layout::{ARENA, CQ_HEAD, CQ_TAIL, DEPTH, SQ_BASE, SQ_TAIL}; use toyos_blockring::wire::{self, Refusal}; -use toyos_blockring::{BLOCK_BYTES, MAX_REQUEST_BLOCKS, PORT}; +use toyos_blockring::{Op, Request, BLOCK_BYTES, MAX_REQUEST_BLOCKS, PORT}; const SELF: &str = "/system/bin/test_rs_blockd_io"; @@ -91,6 +98,10 @@ const NARROW: u64 = 128 * 1024 * 1024; /// a liveness bound, far past what one block takes. const AIMED: Duration = Duration::from_secs(10); +/// How long `hostile-head` waits for blockd to withhold its write's answer, and +/// then for the reset that ends its session: a liveness bound. +const SILENCE_ENDS: Duration = Duration::from_secs(30); + fn guid(text: &str) -> [u8; 16] { PartGuid::parse(text).unwrap_or_else(|| panic!("{text} is no GUID")).0 } @@ -121,6 +132,8 @@ struct Blockd { acceptor: Acceptor, connector: Connector, child: Option, + /// Every line the running blockd says. + said: Option>, } impl Blockd { @@ -130,7 +143,7 @@ impl Blockd { fn with(syscap: SysCap, args: &[&str]) -> Self { let (acceptor, connector) = port::create().unwrap_or_else(|e| fail(format!("no port: {e:?}"))); - let mut blockd = Self { syscap, acceptor, connector, child: None }; + let mut blockd = Self { syscap, acceptor, connector, child: None, said: None }; blockd.spawn(args, false); blockd } @@ -164,6 +177,7 @@ impl Blockd { command.endow(&format!("{SERVE_PREFIX}{PORT}"), acceptor.0); let mut child = command.spawn().unwrap_or_else(|e| fail(format!("blockd did not start: {e}"))); let out = child.stdout.take().expect("piped"); + let (says, said) = mpsc::channel(); let mut kill = kill_on_withheld.then(|| { toyos_abi::syscall::dup(toyos_abi::RawHandle(child.as_raw_handle())) .unwrap_or_else(|e| fail(format!("blockd's handle would not duplicate: {e:?}"))) @@ -177,9 +191,27 @@ impl Blockd { println!("blockd_io: blockd killed with the withheld write done on the device"); } } + let _ = says.send(line); } }); self.child = Some(child); + self.said = Some(said); + } + + /// Wait, at most `bound`, for the running blockd to say a line holding + /// `needle`. + fn says(&self, needle: &str, bound: Duration) { + let said = self.said.as_ref().expect("spawned"); + let asked = Instant::now(); + loop { + let left = bound.saturating_sub(asked.elapsed()); + match said.recv_timeout(left) { + Ok(line) if line.contains(needle) => return, + Ok(_) => {} + Err(RecvTimeoutError::Timeout) => fail(format!("blockd did not say {needle:?} in {bound:?}")), + Err(RecvTimeoutError::Disconnected) => fail(format!("blockd ended before it said {needle:?}")), + } + } } /// End the running blockd, if one is, and wait for it to be gone. @@ -497,6 +529,54 @@ fn reset() { println!("blockd_io: PASS reset"); } +/// A client whose write is on the device moves its completion ring's head a +/// ring's depth behind the tail blockd published, so the answer finds no room: +/// blockd ends that session, and serves the next. +fn hostile_head() { + let blockd = Blockd::start(&["--silence-write", "1"]); + let region = Region::create().unwrap_or_else(|e| fail(format!("a region: {e:?}"))); + let conn = blockd.names().open(PORT).unwrap_or_else(|e| fail(format!("the port: {e:?}"))); + let shared = region.share().unwrap_or_else(|e| fail(format!("a second handle: {e:?}"))); + conn.send_bytes_with_handles(&[shared], wire::MSG_OPEN, &guid(TARGET)) + .unwrap_or_else(|e| fail(format!("the open: {e:?}"))); + let header = conn.recv_header().unwrap_or_else(|e| fail(format!("the answer: {e:?}"))); + let mut payload = [0u8; 64]; + conn.recv_bytes(&header, &mut payload).unwrap_or_else(|e| fail(format!("the answer: {e:?}"))); + if header.msg_type != wire::MSG_OPENED { + fail(format!("the slot's open was answered {}", header.msg_type)); + } + // A write of the slot's block 0 from arena block 0 under tag 1, as the + // words a client puts on the request ring, published and rung. + let words = region.words(); + let run = ARENA.run(0, 1).unwrap_or_else(|| fail("arena block 0 is no run".into())); + let write = Request { op: Op::Write { run, lba: 0 }, tag: 1 }; + for (at, word) in write.encode().into_iter().enumerate() { + words[SQ_BASE + at].store(word, Ordering::Relaxed); + } + words[SQ_TAIL].store(1, Ordering::Release); + conn.write_nonblock(&[1]).unwrap_or_else(|e| fail(format!("the doorbell: {e:?}"))); + blockd.says("WITHHELD", SILENCE_ENDS); + let tail = words[CQ_TAIL].load(Ordering::Acquire); + words[CQ_HEAD].store(tail.wrapping_sub(DEPTH), Ordering::Release); + conn.write_nonblock(&[1]).unwrap_or_else(|e| fail(format!("the doorbell: {e:?}"))); + println!("blockd_io: with a write on the device, the client moved its completion head {DEPTH} behind the tail"); + blockd.says("closed after", SILENCE_ENDS); + println!("blockd_io: blockd ended the session and runs on"); + let mut next = open(blockd.names(), TARGET); + let block = pattern(0x6B, 0); + match next.write(0, &block) { + Ok(Outcome::Done) => {} + other => fail(format!("the next session's write was answered {other:?}")), + } + flushed(&mut next); + match next.read(0, 1) { + Ok((Outcome::Done, Some(data))) if data == block => {} + other => fail(format!("the next session read back {:?}", other.map(|(o, _)| o))), + } + println!("blockd_io: the next session wrote, flushed and read back the slot's block 0"); + println!("blockd_io: PASS hostile-head"); +} + /// The FAT32 volume's device: a session, with blockd's supervisor beside it. struct Volume { blockd: Blockd, @@ -964,6 +1044,7 @@ fn main() { Some("claims") => claims(), Some("holder") => holder_role(args.get(2).map_or("", String::as_str)), Some("bench") => bench(), + Some("hostile-head") => hostile_head(), Some("reset") => reset(), Some("crash") => crash(), Some(role @ ("dma-inside" | "dma-outside" | "dma-revoked" | "dma-after")) => dma(role), diff --git a/toyos-blockring/Cargo.toml b/toyos-blockring/Cargo.toml index ae73456fa9..f9cbb2b18c 100644 --- a/toyos-blockring/Cargo.toml +++ b/toyos-blockring/Cargo.toml @@ -1,15 +1,15 @@ # A member of the host workspace (root `Cargo.toml`). What lives here is the # block protocol between a block service and its client — the shared session -# page's layout, the two rings on it, the words a request and a completion are, -# the control frames a session is opened with, and every decision -# either end makes about what a completion means — with nothing that touches a -# device, a handle or a mapping. `userland/blockd` serves it and its client -# glue speaks it; both map the page and hand this crate the words. +# page's layout, where its two `toyos-transport` rings are, the words a request +# and a completion are, the control frames a session is opened with, and every +# decision either end makes about what a completion means — with nothing that +# touches a device, a handle or a mapping. `userland/blockd` serves it and its +# client glue speaks it; both map the page and hand this crate the words. # # Its defects are orders — a completion racing a crash, a flush racing a reset, # a reconnect racing a reissue — so `src/model.rs` enumerates every ordering of -# a scripted client against a server, a device and a crash, and -# `tests/loom_ring.rs` checks the rings' publication under loom. +# a scripted client against a server, a device and a crash; the rings' +# publication is `toyos-transport`'s, and checked under loom there. [package] name = "toyos-blockring" @@ -32,19 +32,14 @@ mutate-abort-keeps-inflight = [] # re-issues nothing: writes it had acknowledged are gone and a later flush # still says durable. mutate-no-reissue-after-loss = [] -# The rings publish a tail with `Relaxed`, so the consumer can see the index -# before the entry's words. -mutate-ring-publish-relaxed = [] [dependencies] # Whose flush answers for which writes the disk lost, and who holds which span: # the server half's bookkeeping, decided where the kernel's block layer decides # it today. toyos-blockhold = { path = "../toyos-blockhold" } - -[dev-dependencies] -# Already resolved in this workspace for `kernel-loom` and `toyos-sched/loom`. -loom = "0.7" +# The rings on the session page. +toyos-transport = { path = "../toyos-transport" } [lints.rust] warnings = "deny" diff --git a/toyos-blockring/src/client.rs b/toyos-blockring/src/client.rs index 933665c785..191cf9c300 100644 --- a/toyos-blockring/src/client.rs +++ b/toyos-blockring/src/client.rs @@ -5,7 +5,7 @@ //! **Every request asked for is answered exactly once**, by [`Client::complete`] //! or by [`Client::session_ended`], and never twice: a completion for a tag //! that is not on the wire is the server breaking the session -//! ([`Violation`]), not a second answer. +//! ([`Violation::Tag`]), not a second answer. //! //! **An acknowledged write is kept until a flush covers it.** Its arena blocks //! stay pinned ([`Client::take_released`] is when they come back) because a @@ -36,7 +36,10 @@ use alloc::collections::{BTreeMap, VecDeque}; use alloc::vec::Vec; +use toyos_transport::{Inflight, Run, Untrusted, Violation}; + use crate::entry::{Completion, Op, Request, Status}; +use crate::layout::DEPTH; /// The caller's name for one thing it asked for. pub type Ticket = u64; @@ -65,17 +68,11 @@ pub enum Outcome { Refused, } -/// The server said something no server of this protocol says. The session is -/// over, and [`Client::session_ended`] is what the caller does next. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub struct Violation; - /// A write acknowledged and not yet covered by a flush. #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] struct Acked { lba: u64, - blocks: u32, - arena: u32, + run: Run, /// When it was first acknowledged, which is the order it is issued again /// in. first: u64, @@ -86,7 +83,7 @@ struct Acked { impl Acked { fn overlaps_range(&self, lba: u64, blocks: u32) -> bool { - self.lba < lba + u64::from(blocks) && lba < self.lba + u64::from(self.blocks) + self.lba < lba + u64::from(blocks) && lba < self.lba + u64::from(self.run.count()) } } @@ -105,9 +102,6 @@ enum Kind { #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] struct Queued { op: Op, - lba: u64, - blocks: u32, - arena: u32, kind: Kind, } @@ -120,47 +114,42 @@ struct Sent { losses: u64, } +/// `D` is the most requests it has on the wire at once: a session's [`DEPTH`]. #[derive(Clone, Debug, PartialEq, Eq, Hash, Default)] -pub struct Client { +pub struct Client { up: bool, /// Not yet on the wire, in order. Writes being issued again are always at /// the front, then any flush attempt waiting on them. outbox: VecDeque, - wire: BTreeMap, + wire: Inflight, acked: VecDeque, - next_tag: u32, next_seq: u64, /// Losses taken so far: a flush sent before one says nothing about what it /// lost. losses: u64, attempts: BTreeMap, answers: VecDeque<(Ticket, Outcome)>, - released: VecDeque<(u32, u32)>, + released: VecDeque, /// Writes put on the wire again after a loss, over the client's life. reissued: u64, /// An acknowledged write was given up: no flush is durable again. gave_up: bool, } -impl Client { +impl Client { /// A client with no session yet: requests wait until /// [`Self::session_started`]. pub fn new() -> Self { Self::default() } - /// Ask for a read or write of `blocks` at `lba`, through arena blocks - /// from `arena`, or for a flush (whose range is ignored). - pub fn submit(&mut self, ticket: Ticket, op: Op, lba: u64, blocks: u32, arena: u32) { - let (lba, blocks, arena) = match op { - Op::Flush => (0, 0, 0), - Op::Read | Op::Write => (lba, blocks, arena), - }; + /// Ask for a read, a write or a flush. + pub fn submit(&mut self, ticket: Ticket, op: Op) { let kind = match op { Op::Flush => Kind::Flush(ticket), - Op::Read | Op::Write => Kind::User(ticket), + Op::Read { .. } | Op::Write { .. } => Kind::User(ticket), }; - self.outbox.push_back(Queued { op, lba, blocks, arena, kind }); + self.outbox.push_back(Queued { op, kind }); } fn reissuing(&self) -> bool { @@ -169,8 +158,8 @@ impl Client { } /// The next request to put on the wire, tagged; `None` while there is no - /// session, nothing to send, or what is next must wait for writes being - /// issued again. + /// session, nothing to send, `D` on the wire, or what is next must + /// wait for writes being issued again. pub fn next_request(&mut self) -> Option { if !self.up { return None; @@ -183,8 +172,7 @@ impl Client { // later one of the caller's. Kind::Reissue(acked) => { let blocked = self.wire.values().any(|sent| { - let q = sent.queued; - q.op == Op::Write && acked.overlaps_range(q.lba, q.blocks) + matches!(sent.queued.op, Op::Write { run, lba } if acked.overlaps_range(lba, run.count())) }); if blocked { return None; @@ -196,28 +184,24 @@ impl Client { } } } + let tag = self.wire.insert(Sent { queued: front, covers: self.next_seq, losses: self.losses }).ok()?; self.outbox.pop_front(); if matches!(front.kind, Kind::Reissue(_)) { self.reissued += 1; } - let tag = self.next_tag; - self.next_tag = self.next_tag.wrapping_add(1); - let (covers, losses) = (self.next_seq, self.losses); - self.wire.insert(tag, Sent { queued: front, covers, losses }); - Some(Request { op: front.op, tag, lba: front.lba, blocks: front.blocks, arena: front.arena }) + Some(Request { op: front.op, tag }) } /// What the server answered. pub fn complete(&mut self, completion: Completion) -> Result<(), Violation> { - let sent = self.wire.remove(&completion.tag).ok_or(Violation)?; + let sent = self.wire.answer(Untrusted::new(completion.tag))?; let q = sent.queued; match (q.kind, completion.status) { (Kind::User(ticket), Status::Ok) => { match q.op { - Op::Write => { + Op::Write { run, lba } => { let seq = self.bump(); - let acked = - Acked { lba: q.lba, blocks: q.blocks, arena: q.arena, first: seq, seq, attempts: 0 }; + let acked = Acked { lba, run, first: seq, seq, attempts: 0 }; // Sent before a loss it did not see: the device may // have taken it and lost it, and an earlier write it // overlaps is being issued again — so it is issued @@ -228,19 +212,21 @@ impl Client { self.acked.push_back(acked); } } - Op::Read => self.released.push_back((q.arena, q.blocks)), - Op::Flush => return Err(Violation), + Op::Read { run, .. } => self.released.push_back(run), + Op::Flush => return Err(Violation::Entry), } self.answers.push_back((ticket, Outcome::Done)); } (Kind::User(ticket), status) => { - self.released.push_back((q.arena, q.blocks)); + if let Op::Read { run, .. } | Op::Write { run, .. } = q.op { + self.released.push_back(run); + } let outcome = match status { Status::Invalid => Outcome::Invalid, Status::Device => Outcome::Device, // A read or a write answered `Lost` is a server that does // not know which op it was. - Status::Lost | Status::Ok => return Err(Violation), + Status::Lost | Status::Ok => return Err(Violation::Entry), }; self.answers.push_back((ticket, outcome)); } @@ -259,7 +245,7 @@ impl Client { if acked.attempts > MAX_ATTEMPTS { // The device will not take the only copy there is: every // flush waiting is told so, and the write is gone. - self.released.push_back((acked.arena, acked.blocks)); + self.released.push_back(acked.run); self.gave_up = true; self.fail_waiting_flushes(); } else { @@ -274,7 +260,7 @@ impl Client { (Kind::Flush(ticket), Status::Ok) => { while self.acked.front().is_some_and(|a| a.seq < sent.covers) { let a = self.acked.pop_front().expect("just seen"); - self.released.push_back((a.arena, a.blocks)); + self.released.push_back(a.run); } self.attempts.remove(&ticket); let outcome = if self.gave_up { Outcome::Device } else { Outcome::Durable }; @@ -298,15 +284,18 @@ impl Client { /// The session is over: the server hung up, broke the protocol, or died. pub fn session_ended(&mut self) { self.up = false; - let wire = core::mem::take(&mut self.wire); + let mut wire = Vec::new(); + self.wire.end(|_, sent| wire.push(sent)); let mut flushes = Vec::new(); - for sent in wire.into_values() { + for sent in wire { let q = sent.queued; match q.kind { Kind::User(ticket) => { #[cfg(not(feature = "mutate-session-end-forgets"))] { - self.released.push_back((q.arena, q.blocks)); + if let Op::Read { run, .. } | Op::Write { run, .. } = q.op { + self.released.push_back(run); + } self.answers.push_back((ticket, Outcome::Refused)); } #[cfg(feature = "mutate-session-end-forgets")] @@ -338,19 +327,19 @@ impl Client { self.answers.drain(..) } - /// Arena runs `(first, blocks)` nothing will read or write again. - pub fn take_released(&mut self) -> impl Iterator + '_ { + /// Arena runs nothing will read or write again. + pub fn take_released(&mut self) -> impl Iterator + '_ { self.released.drain(..) } /// Nothing is waiting, on the wire, or held for a flush. pub fn quiet(&self) -> bool { - self.outbox.is_empty() && self.wire.is_empty() && self.acked.is_empty() + self.outbox.is_empty() && self.wire.values().next().is_none() && self.acked.is_empty() } - /// How many requests are on the wire. - pub fn on_the_wire(&self) -> usize { - self.wire.len() + /// The tags on the wire. + pub fn on_the_wire(&self) -> impl Iterator + '_ { + self.wire.tags() } /// How many writes have gone on the wire again after a loss. @@ -380,13 +369,7 @@ impl Client { Kind::User(_) | Kind::Flush(_) => true, }) .unwrap_or(self.outbox.len()); - self.outbox.insert(at, Queued { - op: Op::Write, - lba: acked.lba, - blocks: acked.blocks, - arena: acked.arena, - kind: Kind::Reissue(acked), - }); + self.outbox.insert(at, Queued { op: Op::Write { run: acked.run, lba: acked.lba }, kind: Kind::Reissue(acked) }); } /// Ask the flush `ticket` again, once every write being issued again is @@ -397,7 +380,7 @@ impl Client { .iter() .position(|q| !matches!(q.kind, Kind::Reissue(_) | Kind::Flush(_))) .unwrap_or(self.outbox.len()); - self.outbox.insert(at, Queued { op: Op::Flush, lba: 0, blocks: 0, arena: 0, kind: Kind::Flush(ticket) }); + self.outbox.insert(at, Queued { op: Op::Flush, kind: Kind::Flush(ticket) }); } /// The device may no longer hold any write acknowledged and not covered: @@ -433,6 +416,11 @@ impl Client { #[cfg(test)] mod tests { use super::*; + use crate::layout::ARENA; + + fn run(first: u32, count: u32) -> Run { + ARENA.run(first, count).unwrap() + } fn answer(client: &mut Client, request: Request, status: Status) { client.complete(Completion { tag: request.tag, status }).unwrap(); @@ -440,8 +428,8 @@ mod tests { #[test] fn nothing_goes_out_before_a_session() { - let mut client = Client::new(); - client.submit(1, Op::Read, 0, 1, 0); + let mut client: Client = Client::new(); + client.submit(1, Op::Read { run: run(0, 1), lba: 0 }); assert_eq!(client.next_request(), None); client.session_started(); assert!(client.next_request().is_some()); @@ -449,12 +437,12 @@ mod tests { #[test] fn a_second_answer_for_one_tag_is_a_violation() { - let mut client = Client::new(); + let mut client: Client = Client::new(); client.session_started(); - client.submit(1, Op::Read, 0, 1, 0); + client.submit(1, Op::Read { run: run(0, 1), lba: 0 }); let r = client.next_request().unwrap(); answer(&mut client, r, Status::Ok); - assert_eq!(client.complete(Completion { tag: r.tag, status: Status::Ok }), Err(Violation)); + assert_eq!(client.complete(Completion { tag: r.tag, status: Status::Ok }), Err(Violation::Tag)); } /// A write acknowledged, then a flush answered `Lost`: the write goes out @@ -462,16 +450,16 @@ mod tests { /// `Durable` only from the second. #[test] fn a_lost_flush_reissues_then_asks_again() { - let mut client = Client::new(); + let mut client: Client = Client::new(); client.session_started(); - client.submit(1, Op::Write, 5, 2, 3); + client.submit(1, Op::Write { run: run(3, 2), lba: 5 }); let w = client.next_request().unwrap(); answer(&mut client, w, Status::Ok); - client.submit(2, Op::Flush, 0, 0, 0); + client.submit(2, Op::Flush); let f = client.next_request().unwrap(); answer(&mut client, f, Status::Lost); let again = client.next_request().unwrap(); - assert_eq!((again.op, again.lba, again.blocks, again.arena), (Op::Write, 5, 2, 3)); + assert_eq!(again.op, Op::Write { run: run(3, 2), lba: 5 }); assert_eq!(client.next_request(), None, "the flush waits for the write"); answer(&mut client, again, Status::Ok); let f2 = client.next_request().unwrap(); @@ -479,7 +467,7 @@ mod tests { answer(&mut client, f2, Status::Ok); let answers: Vec<_> = client.take_answers().collect(); assert_eq!(answers, [(1, Outcome::Done), (2, Outcome::Durable)]); - assert_eq!(client.take_released().collect::>(), [(3, 2)]); + assert_eq!(client.take_released().collect::>(), [run(3, 2)]); assert!(client.quiet()); } @@ -487,22 +475,22 @@ mod tests { /// other, in the order they were first acknowledged. #[test] fn overlapping_writes_go_out_again_in_order() { - let mut client = Client::new(); + let mut client: Client = Client::new(); client.session_started(); - client.submit(1, Op::Write, 0, 2, 0); + client.submit(1, Op::Write { run: run(0, 2), lba: 0 }); let a = client.next_request().unwrap(); answer(&mut client, a, Status::Ok); - client.submit(2, Op::Write, 1, 1, 2); + client.submit(2, Op::Write { run: run(2, 1), lba: 1 }); let b = client.next_request().unwrap(); answer(&mut client, b, Status::Ok); client.session_ended(); client.session_started(); let first = client.next_request().unwrap(); - assert_eq!((first.lba, first.arena), (0, 0)); + assert_eq!(first.op, Op::Write { run: run(0, 2), lba: 0 }); assert_eq!(client.next_request(), None, "the overlapping one waits"); answer(&mut client, first, Status::Ok); let second = client.next_request().unwrap(); - assert_eq!((second.lba, second.arena), (1, 2)); + assert_eq!(second.op, Op::Write { run: run(2, 1), lba: 1 }); } /// A write the device refused on every reissue is gone, and its caller was @@ -510,24 +498,24 @@ mod tests { /// never durable. #[test] fn a_write_given_up_poisons_every_later_flush() { - let mut client = Client::new(); + let mut client: Client = Client::new(); client.session_started(); - client.submit(1, Op::Write, 0, 1, 0); + client.submit(1, Op::Write { run: run(0, 1), lba: 0 }); let w = client.next_request().unwrap(); answer(&mut client, w, Status::Ok); client.session_ended(); for _ in 0..=MAX_ATTEMPTS { client.session_started(); let again = client.next_request().unwrap(); - assert_eq!((again.op, again.lba), (Op::Write, 0)); + assert_eq!(again.op, Op::Write { run: run(0, 1), lba: 0 }); answer(&mut client, again, Status::Device); } - client.submit(2, Op::Flush, 0, 0, 0); + client.submit(2, Op::Flush); let f = client.next_request().unwrap(); assert_eq!(f.op, Op::Flush); answer(&mut client, f, Status::Ok); assert_eq!(client.take_answers().collect::>(), [(1, Outcome::Done), (2, Outcome::Device)]); - client.submit(3, Op::Flush, 0, 0, 0); + client.submit(3, Op::Flush); let f = client.next_request().unwrap(); answer(&mut client, f, Status::Ok); assert_eq!(client.take_answers().collect::>(), [(3, Outcome::Device)]); @@ -535,16 +523,16 @@ mod tests { #[test] fn what_was_on_the_wire_at_the_end_is_refused_and_what_was_not_waits() { - let mut client = Client::new(); + let mut client: Client = Client::new(); client.session_started(); - client.submit(1, Op::Write, 0, 1, 0); - client.submit(2, Op::Read, 4, 1, 1); + client.submit(1, Op::Write { run: run(0, 1), lba: 0 }); + client.submit(2, Op::Read { run: run(1, 1), lba: 4 }); let _ = client.next_request().unwrap(); client.session_ended(); assert_eq!(client.take_answers().collect::>(), [(1, Outcome::Refused)]); - assert_eq!(client.take_released().collect::>(), [(0, 1)]); + assert_eq!(client.take_released().collect::>(), [run(0, 1)]); client.session_started(); let r = client.next_request().unwrap(); - assert_eq!((r.op, r.lba), (Op::Read, 4)); + assert_eq!(r.op, Op::Read { run: run(1, 1), lba: 4 }); } } diff --git a/toyos-blockring/src/entry.rs b/toyos-blockring/src/entry.rs index 1e8b8b8e17..f39788b4c0 100644 --- a/toyos-blockring/src/entry.rs +++ b/toyos-blockring/src/entry.rs @@ -4,51 +4,34 @@ //! request is bounded here against the arena and the partition, and a word //! this protocol does not define is a refusal, never a default. -use crate::layout::{ARENA_BLOCKS, CQE_WORDS, MAX_REQUEST_BLOCKS, SQE_WORDS}; +use toyos_transport::{Run, Untrusted}; -/// What a request asks of the partition. +use crate::layout::{ARENA, CQE_WORDS, MAX_REQUEST_BLOCKS, SQE_WORDS}; + +/// What a request asks of the partition, and the arena blocks the data is in +/// or goes to. `lba` is the partition's own block number, from 0: nothing in +/// this protocol names a device block, so a neighbour's blocks have no +/// spelling. #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] pub enum Op { - /// `blocks` from the partition's block `lba` into the arena. - Read, - /// `blocks` from the arena to the partition's block `lba`. - Write, + /// The run's blocks from the partition's block `lba` into the arena. + Read { run: Run, lba: u64 }, + /// The run's blocks from the arena to the partition's block `lba`. + Write { run: Run, lba: u64 }, /// Every write acknowledged before this was submitted, onto the medium. Flush, } -impl Op { - const fn word(self) -> u32 { - match self { - Self::Read => 1, - Self::Write => 2, - Self::Flush => 3, - } - } +const READ: u32 = 1; +const WRITE: u32 = 2; +const FLUSH: u32 = 3; - const fn from_word(word: u32) -> Option { - match word { - 1 => Some(Self::Read), - 2 => Some(Self::Write), - 3 => Some(Self::Flush), - _ => None, - } - } -} - -/// One request. `lba` is the partition's own block number, from 0: nothing in -/// this protocol names a device block, so a neighbour's blocks have no -/// spelling. A flush carries no range, and its `lba`, `blocks` and `arena` are -/// zero. +/// One request. #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] pub struct Request { pub op: Op, /// The client's name for it, echoed by its completion. pub tag: u32, - pub lba: u64, - pub blocks: u32, - /// The first arena block the data is in or goes to. - pub arena: u32, } /// Why a request was answered without being done. @@ -62,53 +45,53 @@ pub enum Refused { impl Request { pub fn encode(&self) -> [u32; SQE_WORDS] { - [ - self.op.word(), - self.tag, - self.lba as u32, - (self.lba >> 32) as u32, - self.blocks, - self.arena, - 0, - 0, - ] + let (op, lba, first, count) = match self.op { + Op::Read { run, lba } => (READ, lba, run.first(), run.count()), + Op::Write { run, lba } => (WRITE, lba, run.first(), run.count()), + Op::Flush => (FLUSH, 0, 0, 0), + }; + [op, self.tag, lba as u32, (lba >> 32) as u32, count, first, 0, 0] } /// The request these words are, bounded against a partition of /// `partition_blocks`: every block it names is inside both the arena and /// the partition, a transfer moves at least one block and at most /// [`MAX_REQUEST_BLOCKS`], and a flush names nothing. - pub fn decode(words: [u32; SQE_WORDS], partition_blocks: u64) -> Result { - let tag = words[1]; + pub fn decode(words: [Untrusted; SQE_WORDS], partition_blocks: u64) -> Result { + let [op, tag, lba_low, lba_high, blocks, arena, reserved @ ..] = words; + let tag = opaque(tag); let refused = Refused::Malformed { tag }; - let op = Op::from_word(words[0]).ok_or(refused)?; - let lba = u64::from(words[2]) | (u64::from(words[3]) << 32); - let (blocks, arena) = (words[4], words[5]); - if words[6] != 0 || words[7] != 0 { + if !reserved.iter().all(|word| word.is(0)) { return Err(refused); } - match op { - Op::Flush => { - if lba != 0 || blocks != 0 || arena != 0 { - return Err(refused); - } - } - Op::Read | Op::Write => { - if blocks == 0 || blocks > MAX_REQUEST_BLOCKS { - return Err(refused); - } - if arena.checked_add(blocks).is_none_or(|end| end > ARENA_BLOCKS) { - return Err(refused); - } - if lba.checked_add(u64::from(blocks)).is_none_or(|end| end > partition_blocks) { - return Err(refused); - } - } + let lba = u64::from(opaque(lba_low)) | (u64::from(opaque(lba_high)) << 32); + let transfer: fn(Run, u64) -> Op = if op.is(READ) { + |run, lba| Op::Read { run, lba } + } else if op.is(WRITE) { + |run, lba| Op::Write { run, lba } + } else if op.is(FLUSH) && lba == 0 && blocks.is(0) && arena.is(0) { + return Ok(Self { op: Op::Flush, tag }); + } else { + return Err(refused); + }; + let run = Run::decode(arena, blocks, &ARENA).map_err(|_| refused)?; + let blocks = run.count(); + if blocks > MAX_REQUEST_BLOCKS || lba.checked_add(u64::from(blocks)).is_none_or(|end| end > partition_blocks) { + return Err(refused); } - Ok(Self { op, tag, lba, blocks, arena }) + Ok(Self { op: transfer(run, lba), tag }) } } +/// A word every value of which means something: a tag its answer echoes, or +/// half of an `lba` the partition bounds whole. +fn opaque(word: Untrusted) -> u32 { + word.at_most(u32::MAX.into()) + .ok() + .and_then(|word| u32::try_from(word).ok()) + .expect("every u32 is at most u32::MAX") +} + /// What a completion says of its request. #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] pub enum Status { @@ -135,16 +118,6 @@ impl Status { Self::Lost => 3, } } - - const fn from_word(word: u32) -> Option { - match word { - 0 => Some(Self::Ok), - 1 => Some(Self::Invalid), - 2 => Some(Self::Device), - 3 => Some(Self::Lost), - _ => None, - } - } } /// One completion. @@ -161,11 +134,14 @@ impl Completion { /// `None` for words no server of this protocol writes: the client treats a /// server that wrote them as one that broke the session. - pub fn decode(words: [u32; CQE_WORDS]) -> Option { - if words[2] != 0 || words[3] != 0 { + pub fn decode(words: [Untrusted; CQE_WORDS]) -> Option { + let [tag, status, reserved @ ..] = words; + if !reserved.iter().all(|word| word.is(0)) { return None; } - Some(Self { tag: words[0], status: Status::from_word(words[1])? }) + let status = + [Status::Ok, Status::Invalid, Status::Device, Status::Lost].into_iter().find(|s| status.is(s.word()))?; + Some(Self { tag: opaque(tag), status }) } } @@ -175,39 +151,54 @@ mod tests { const PARTITION: u64 = 1000; + fn run(first: u32, count: u32) -> Run { + ARENA.run(first, count).unwrap() + } + + fn peer(words: [u32; N]) -> [Untrusted; N] { + words.map(Untrusted::new) + } + #[test] fn a_request_survives_its_words() { + let last = ARENA.slots() - MAX_REQUEST_BLOCKS; for request in [ - Request { op: Op::Read, tag: 7, lba: 999, blocks: 1, arena: 0 }, - Request { op: Op::Write, tag: u32::MAX, lba: 0, blocks: MAX_REQUEST_BLOCKS, arena: ARENA_BLOCKS - MAX_REQUEST_BLOCKS }, - Request { op: Op::Flush, tag: 0, lba: 0, blocks: 0, arena: 0 }, + Request { op: Op::Read { run: run(0, 1), lba: 999 }, tag: 7 }, + Request { op: Op::Write { run: run(last, MAX_REQUEST_BLOCKS), lba: 0 }, tag: u32::MAX }, + Request { op: Op::Flush, tag: 0 }, ] { - assert_eq!(Request::decode(request.encode(), PARTITION), Ok(request)); + assert_eq!(Request::decode(peer(request.encode()), PARTITION), Ok(request)); } - let wide = Request { op: Op::Read, tag: 1, lba: 1 << 40, blocks: 2, arena: 3 }; - assert_eq!(Request::decode(wide.encode(), u64::MAX), Ok(wide)); + let wide = Request { op: Op::Read { run: run(3, 2), lba: 1 << 40 }, tag: 1 }; + assert_eq!(Request::decode(peer(wide.encode()), u64::MAX), Ok(wide)); } /// Every field a hostile client can set out of range is refused, and the /// refusal still carries its tag. #[test] fn a_request_outside_its_bounds_is_refused_by_tag() { - let good = Request { op: Op::Write, tag: 42, lba: 10, blocks: 2, arena: 5 }; + let write = |lba| Request { op: Op::Write { run: run(5, 2), lba }, tag: 42 }; + let good = write(10); + let with = |at: usize, word: u32| { + let mut words = good.encode(); + words[at] = word; + words + }; let cases: [(&str, [u32; SQE_WORDS]); 10] = [ - ("op 0", { let mut w = good.encode(); w[0] = 0; w }), - ("op 4", { let mut w = good.encode(); w[0] = 4; w }), - ("no blocks", Request { blocks: 0, ..good }.encode()), - ("too many blocks", Request { blocks: MAX_REQUEST_BLOCKS + 1, ..good }.encode()), - ("past the arena", Request { arena: ARENA_BLOCKS - 1, ..good }.encode()), - ("arena wraps", Request { arena: u32::MAX, ..good }.encode()), - ("past the partition", Request { lba: PARTITION - 1, ..good }.encode()), - ("lba wraps", Request { lba: u64::MAX, ..good }.encode()), - ("reserved word", { let mut w = good.encode(); w[7] = 1; w }), - ("a flush with a range", Request { op: Op::Flush, ..good }.encode()), + ("op 0", with(0, 0)), + ("op 4", with(0, 4)), + ("no blocks", with(4, 0)), + ("too many blocks", with(4, MAX_REQUEST_BLOCKS + 1)), + ("past the arena", with(5, ARENA.slots() - 1)), + ("arena wraps", with(5, u32::MAX)), + ("past the partition", write(PARTITION - 1).encode()), + ("lba wraps", write(u64::MAX).encode()), + ("reserved word", with(7, 1)), + ("a flush with a range", with(0, 3)), ]; for (what, words) in cases { assert_eq!( - Request::decode(words, PARTITION), + Request::decode(peer(words), PARTITION), Err(Refused::Malformed { tag: 42 }), "{what} was not refused" ); @@ -218,9 +209,9 @@ mod tests { fn a_completion_survives_its_words_and_refuses_what_it_does_not_define() { for status in [Status::Ok, Status::Invalid, Status::Device, Status::Lost] { let c = Completion { tag: 9, status }; - assert_eq!(Completion::decode(c.encode()), Some(c)); + assert_eq!(Completion::decode(peer(c.encode())), Some(c)); } - assert_eq!(Completion::decode([1, 4, 0, 0]), None); - assert_eq!(Completion::decode([1, 0, 1, 0]), None); + assert_eq!(Completion::decode(peer([1, 4, 0, 0])), None); + assert_eq!(Completion::decode(peer([1, 0, 1, 0])), None); } } diff --git a/toyos-blockring/src/layout.rs b/toyos-blockring/src/layout.rs index 1d36547d02..05f60c90d8 100644 --- a/toyos-blockring/src/layout.rs +++ b/toyos-blockring/src/layout.rs @@ -3,14 +3,22 @@ //! The four ring indices sit on cache lines of their own, so the client's //! stores to its two and the server's to its two never share a line. +use toyos_transport::{Consumer, Geometry, Place, Producer, Word}; + /// A session's whole region: the one size shared memory comes in. -pub const SESSION_BYTES: usize = 2 * 1024 * 1024; +pub const SESSION_BYTES: usize = Geometry::BYTES as usize; /// The unit every request is in, and the unit the arena is cut into. pub const BLOCK_BYTES: usize = 4096; +/// The arena: whole blocks after the rings' page, which a request names by +/// run. +pub const ARENA: Geometry = match Geometry::new(BLOCK_BYTES as u32) { + Some(arena) => arena, + None => panic!("a block is no longer than the arena"), +}; + /// How many requests, and so how many completions, one session has in flight. -/// A power of two, so an index is its ring position masked. pub const DEPTH: u32 = 64; /// The most blocks one request moves. A driver whose device takes less in one @@ -35,17 +43,28 @@ pub const CQ_BASE: usize = SQ_BASE + DEPTH as usize * SQE_WORDS; /// Every word the rings use; the page they are on is the first block. pub const RING_WORDS: usize = CQ_BASE + DEPTH as usize * CQE_WORDS; -/// Where the arena starts, in bytes: the block after the rings' page. -pub const ARENA_OFFSET: usize = BLOCK_BYTES; +/// The request ring and the completion ring. +pub const REQUESTS: Place = Place::new::(); +pub const COMPLETIONS: Place = Place::new::(); -/// The arena's blocks; a request's `arena` is an index below this. -pub const ARENA_BLOCKS: u32 = ((SESSION_BYTES - ARENA_OFFSET) / BLOCK_BYTES) as u32; +/// A client's two ends: requests out, completions in. +pub type ClientRings = (Producer, Consumer); -const _: () = assert!(DEPTH.is_power_of_two()); -const _: () = assert!(RING_WORDS * 4 <= ARENA_OFFSET); -const _: () = assert!(MAX_REQUEST_BLOCKS <= ARENA_BLOCKS); +/// A server's two ends: requests in, completions out. +pub type ServerRings = (Consumer, Producer); -/// The byte offset of arena block `block` in the region. -pub const fn arena_byte(block: u32) -> usize { - ARENA_OFFSET + block as usize * BLOCK_BYTES +/// The client's ends of a session page, every word it owns set to 0. Done +/// before the page is sent to a server, and again before it is sent to the +/// next one. +pub fn client(page: &[W; RING_WORDS]) -> ClientRings { + (Producer::new(page, REQUESTS), Consumer::new(page, COMPLETIONS)) } + +/// The server's ends of a session page it was sent, every word it owns set to +/// 0. Whatever the client left in its own is bounded when first looked at. +pub fn server(page: &[W; RING_WORDS]) -> ServerRings { + (Consumer::new(page, REQUESTS), Producer::new(page, COMPLETIONS)) +} + +const _: () = assert!(RING_WORDS * 4 <= Geometry::HEADER_BYTES as usize); +const _: () = assert!(MAX_REQUEST_BLOCKS <= ARENA.slots()); diff --git a/toyos-blockring/src/lib.rs b/toyos-blockring/src/lib.rs index 57ee1dbd95..dca1b010a1 100644 --- a/toyos-blockring/src/lib.rs +++ b/toyos-blockring/src/lib.rs @@ -3,12 +3,12 @@ //! **A session is one shared region**, [`SESSION_BYTES`] long, that the client //! makes and sends. Its first page holds two single-producer rings — requests //! from the client ([`layout::SQ_BASE`]), completions from the server -//! ([`layout::CQ_BASE`]) — and the rest is the *arena*: whole -//! [`BLOCK_BYTES`] blocks a request names by index and the device moves data -//! into and out of directly. Nothing on the page is a pointer and nothing on it +//! ([`layout::CQ_BASE`]) — and the rest is the *arena* ([`layout::ARENA`]): +//! whole [`BLOCK_BYTES`] blocks a request names by run and the device moves +//! data into and out of directly. Nothing on the page is a pointer and nothing on it //! is trusted by the end that did not write it: a consumer bounds every index -//! and every field before it acts ([`entry::Request::decode`], -//! [`ring::Consumer`]). +//! and every field before it acts ([`entry::Request::decode`], the +//! transport's cursor bounds). //! //! **A doorbell is a byte on the session's connection**, written after the //! entries it announces are published. The connection is also what tells each @@ -28,9 +28,9 @@ //! **Requests in flight at once are unordered**, as on the device: a client //! that needs one to follow another waits for the first's answer. //! -//! Pure: `alloc` and `toyos-blockhold`, no `unsafe`. The ends that map the -//! page — `userland/blockd` and its client — hand this crate the page as -//! words ([`ring::Word`]) and act on what it answers. +//! Pure: `alloc`, `toyos-blockhold` and `toyos-transport`, no `unsafe`. The +//! ends that map the page — `userland/blockd` and its client — hand it the +//! page as words ([`toyos_transport::Word`]) and act on what it answers. #![cfg_attr(not(test), no_std)] #![forbid(unsafe_code)] @@ -40,7 +40,6 @@ extern crate alloc; pub mod client; pub mod entry; pub mod layout; -pub mod ring; pub mod server; pub mod wire; @@ -48,7 +47,8 @@ pub mod wire; mod model; pub use entry::{Completion, Op, Request, Status}; -pub use layout::{ARENA_BLOCKS, BLOCK_BYTES, DEPTH, MAX_REQUEST_BLOCKS, SESSION_BYTES}; +pub use toyos_transport::Run; +pub use layout::{BLOCK_BYTES, DEPTH, MAX_REQUEST_BLOCKS, SESSION_BYTES}; /// The name a block service is served under. A holder of its connector may /// open any partition the service has. diff --git a/toyos-blockring/src/model.rs b/toyos-blockring/src/model.rs index 72786ae985..0edc7115fb 100644 --- a/toyos-blockring/src/model.rs +++ b/toyos-blockring/src/model.rs @@ -1,12 +1,13 @@ //! Every ordering of a client, a server, a device and their failures. //! //! A scripted caller asks for writes and flushes; the rings between the client -//! and the server are queues; the server is [`ServerSession`] over a device of -//! two blocks with a volatile cache. [`explore`] runs, depth first and -//! exhaustively, every interleaving of: the caller asking for its next step, -//! the server taking a request, the device completing any one it holds, **the -//! device failing any one it holds** (not done, answered so), the client -//! reading a completion, **the device being reset** under whatever is in +//! and the server are queues, and rings of the transport's own are driven +//! beside them and held to them both ways ([`Queues`]); the server is [`ServerSession`] +//! over a device of two blocks with a volatile cache. [`explore`] runs, depth +//! first and exhaustively, every interleaving of: the caller asking for its +//! next step, the server taking a request, the device completing any one it +//! holds, **the device failing any one it holds** (not done, answered so), the +//! client reading a completion, **the device being reset** under whatever is in //! flight (each command dropped, or run before the stop with its completion //! read or not; the cache kept or dropped), **the server dying** (each command //! dropped or applied; the cache kept or dropped; the rings left as they @@ -29,12 +30,16 @@ use alloc::collections::{BTreeMap, VecDeque}; use alloc::format; use alloc::string::String; use alloc::vec::Vec; +use core::cell::Cell; +use core::sync::atomic::Ordering; use std::collections::HashSet; use toyos_blockhold::Holds; +use toyos_transport::{Consumer, Place, Producer, Untrusted, Word}; use crate::client::{Client, Outcome, Ticket, MAX_ATTEMPTS}; use crate::entry::{Completion, Op, Request}; +use crate::layout::{ARENA, CQE_WORDS, SQE_WORDS}; use crate::server::{ServerSession, Taken}; const BLOCKS: usize = 2; @@ -63,6 +68,8 @@ pub struct Explored { /// End states where a flush was answered the device's refusal: the client /// gave a write or a flush up. pub given_up: usize, + /// Whether the request ring, and the completion ring, held [`DEPTH`]. + pub filled: [bool; 2], } /// One thing the scripted caller asks for. @@ -79,23 +86,125 @@ pub enum Law { Durable, } -#[derive(Clone, Debug)] +/// One word of the session page. The model runs on one thread, so every order +/// is the program's. +#[derive(Clone)] +struct Shared(Cell); + +impl Word for Shared { + fn load(&self, _: Ordering) -> u32 { + self.0.get() + } + fn store(&self, value: u32, _: Ordering) { + self.0.set(value) + } +} + +/// How deep the model's rings are: shallow enough that a script fills and +/// wraps each. +const DEPTH: u32 = 4; +const SQ_BASE: usize = 4; +const CQ_BASE: usize = SQ_BASE + DEPTH as usize * SQE_WORDS; +const WORDS: usize = CQ_BASE + DEPTH as usize * CQE_WORDS; +const SQ: Place = Place::new::<0, 1, SQ_BASE>(); +const CQ: Place = Place::new::<2, 3, CQ_BASE>(); + +type ClientEnds = (Producer, Consumer); +type ServerEnds = (Consumer, Producer); + +/// The rings between the client and the server, twice: as queues, the +/// reference, and as the transport's rings, which give either end exactly what +/// the queue does — every entry it holds, and nothing when it holds none. A +/// state is keyed by its queues alone ([`key`]), so the search tells apart +/// what the protocol can, not where on the page it is. +#[derive(Clone)] +struct Queues { + sq: VecDeque, + cq: VecDeque, + page: [Shared; WORDS], + client: ClientEnds, + server: ServerEnds, +} + +impl Queues { + fn new() -> Self { + let page = core::array::from_fn(|_| Shared(Cell::new(0))); + let (client, server) = Self::ends(&page); + Self { sq: VecDeque::new(), cq: VecDeque::new(), page, client, server } + } + + /// Both ends over `page`, every cursor stored 0. + fn ends(page: &[Shared; WORDS]) -> (ClientEnds, ServerEnds) { + ((Producer::new(page, SQ), Consumer::new(page, CQ)), (Consumer::new(page, SQ), Producer::new(page, CQ))) + } + + /// The session is over: what either queue held is gone, and the next + /// session's two ends start over the same page. + fn reset(&mut self) { + self.sq.clear(); + self.cq.clear(); + (self.client, self.server) = Self::ends(&self.page); + } + + fn send(&mut self, request: Request) { + assert_eq!(self.client.0.push(&self.page, request.encode()), Ok(true), "a request past the ring's depth"); + self.client.0.publish(&self.page); + self.sq.push_back(request); + } + + /// The oldest request, as the server's ring gives it. + fn take(&mut self) -> Option<[Untrusted; SQE_WORDS]> { + let want = self.sq.pop_front()?; + let words = self.server.0.pop(&self.page).expect("an honest client's tail"); + let words = words.expect("the ring holds what the queue does"); + self.server.0.release(&self.page); + assert_eq!(Request::decode(words, u64::MAX), Ok(want), "the ring gave the server what the queue holds"); + Some(words) + } + + fn post(&mut self, c: Completion) { + assert_eq!(self.server.1.push(&self.page, c.encode()), Ok(true), "a completion past the ring's depth"); + self.server.1.publish(&self.page); + self.cq.push_back(c); + } + + /// The oldest completion, as the client's ring gives it. + fn read(&mut self) -> Option { + let want = self.cq.pop_front()?; + let words: [Untrusted; CQE_WORDS] = + self.client.1.pop(&self.page).expect("an honest server's tail").expect("the ring holds what the queue does"); + self.client.1.release(&self.page); + assert_eq!(Completion::decode(words), Some(want), "the ring gave the client what the queue holds"); + Some(want) + } + + /// A ring whose queue is empty gives nothing. + fn hold_empty(&mut self) { + if self.sq.is_empty() { + assert_eq!(self.server.0.pop(&self.page), Ok(None), "the request ring gave what the queue does not hold"); + } + if self.cq.is_empty() { + assert_eq!(self.client.1.pop(&self.page), Ok(None), "the completion ring gave what the queue does not hold"); + } + } +} + +#[derive(Clone)] struct Server { session: ServerSession, holds: Holds, } -#[derive(Clone, Debug)] +#[derive(Clone)] struct World { - client: Client, + client: Client<{ DEPTH as usize }>, next: usize, /// Every answer each ticket has had. answers: BTreeMap>, /// For each flush ticket, the write tickets answered `Done` before it was /// asked for. acked_before: BTreeMap>, - sq: VecDeque, - cq: VecDeque, + queues: Queues, server: Option, /// The connection: false once the server has died, until a reconnect. alive: bool, @@ -112,12 +221,12 @@ struct World { struct Run<'a> { script: &'a [Step], - /// Visited states, keyed by their whole rendering: `Holds` carries no - /// `Hash`, and a rendering is exact. + /// Visited states, by [`key`]. seen: HashSet, broken: Option<(Law, String)>, ends: usize, given_up: usize, + filled: [bool; 2], /// Whether a flush may end answered `Device`. may_give_up: bool, /// The steps from the start to here, for a failure to name. @@ -126,13 +235,28 @@ struct Run<'a> { /// Explore `script` against at most `failures`. pub fn explore(script: &[Step], failures: Failures) -> Explored { + let mut run = Run { + script, + seen: HashSet::new(), + broken: None, + ends: 0, + given_up: 0, + filled: [false; 2], + may_give_up: failures.total() > MAX_ATTEMPTS, + path: Vec::new(), + }; + dfs(&mut run, start(failures)); + Explored { broken: run.broken, ends: run.ends, given_up: run.given_up, filled: run.filled } +} + +/// The client connected to a fresh server, with nothing asked yet. +fn start(failures: Failures) -> World { let mut world = World { client: Client::new(), next: 0, answers: BTreeMap::new(), acked_before: BTreeMap::new(), - sq: VecDeque::new(), - cq: VecDeque::new(), + queues: Queues::new(), server: None, alive: false, losses: 0, @@ -143,11 +267,7 @@ pub fn explore(script: &[Step], failures: Failures) -> Explored { left: failures, }; connect(&mut world); - let may_give_up = failures.total() > MAX_ATTEMPTS; - let mut run = - Run { script, seen: HashSet::new(), broken: None, ends: 0, given_up: 0, may_give_up, path: Vec::new() }; - dfs(&mut run, world); - Explored { broken: run.broken, ends: run.ends, given_up: run.given_up } + world } fn connect(world: &mut World) { @@ -160,11 +280,10 @@ fn connect(world: &mut World) { pump(world); } -/// What the glue does after every event: every request the client will send -/// goes onto the ring. +/// Every request the client will send goes onto the ring. fn pump(world: &mut World) { while let Some(request) = world.client.next_request() { - world.sq.push_back(request); + world.queues.send(request); } } @@ -180,9 +299,9 @@ fn write_of(script: &[Step], ticket: Ticket) -> Option<(u64, u8)> { /// the cache onto the medium. fn apply(script: &[Step], cache: &mut [Option; BLOCKS], media: &mut [u8; BLOCKS], request: Request) { match request.op { - Op::Write => { - let (_, value) = write_of(script, u64::from(request.arena)).expect("a write's arena is its ticket"); - cache[request.lba as usize] = Some(value); + Op::Write { run, lba } => { + let (_, value) = write_of(script, u64::from(run.first())).expect("a write's arena block is its ticket"); + cache[lba as usize] = Some(value); } Op::Flush => { for (b, slot) in cache.iter_mut().enumerate() { @@ -191,7 +310,7 @@ fn apply(script: &[Step], cache: &mut [Option; BLOCKS], media: &mut [u8; BLO } } } - Op::Read => {} + Op::Read { .. } => {} } } @@ -235,56 +354,44 @@ fn fail(run: &mut Run, law: Law, why: String) { } } -fn go(run: &mut Run, step: String, world: World) { - run.path.push(step); - dfs(run, world); - run.path.pop(); -} - -/// Take the client's answers and hold each against the laws. -fn collect(run: &mut Run, world: &mut World) { +/// What the glue does after every event: take the client's answers, hold each +/// against the laws, and pump. +fn settle(script: &[Step], world: &mut World) -> Result<(), (Law, String)> { let answers: Vec<_> = world.client.take_answers().collect(); let _ = world.client.take_released().count(); for (ticket, outcome) in answers { let had = world.answers.entry(ticket).or_default(); had.push(outcome); if had.len() > 1 { - let had = had.clone(); - fail(run, Law::Answers, format!("ticket {ticket} answered twice: {had:?}")); - return; + return Err((Law::Answers, format!("ticket {ticket} answered twice: {had:?}"))); } if outcome == Outcome::Durable { - durable(run, world, ticket); + durable(script, world, ticket).map_err(|why| (Law::Durable, why))?; } } pump(world); + Ok(()) } /// What the flush `ticket`, just answered durable, promised is on the medium. -fn durable(run: &mut Run, world: &World, flush: Ticket) { +fn durable(script: &[Step], world: &World, flush: Ticket) -> Result<(), String> { let before = &world.acked_before[&flush]; for block in 0..BLOCKS as u64 { - let on_block = |t: &Ticket| write_of(run.script, *t).is_some_and(|(b, _)| b == block); + let on_block = |t: &Ticket| write_of(script, *t).is_some_and(|(b, _)| b == block); let Some(last) = before.iter().copied().filter(on_block).max() else { continue }; - let (_, want) = write_of(run.script, last).expect("a write"); + let (_, want) = write_of(script, last).expect("a write"); let allowed: Vec = core::iter::once(want) - .chain((last + 1..run.script.len() as u64).filter(on_block).filter_map(|t| { - write_of(run.script, t).map(|(_, v)| v) - })) + .chain((last + 1..script.len() as u64).filter(on_block).filter_map(|t| write_of(script, t).map(|(_, v)| v))) .collect(); let on = world.media[block as usize]; if !allowed.contains(&on) { - fail( - run, - Law::Durable, - format!( - "flush {flush} was answered durable with block {block} holding {on}, not the \ - {want} acknowledged before it (or a later one of {allowed:?})" - ), - ); - return; + return Err(format!( + "flush {flush} was answered durable with block {block} holding {on}, not the \ + {want} acknowledged before it (or a later one of {allowed:?})" + )); } } + Ok(()) } /// Nothing more can happen: every ticket was answered, and every flush said @@ -311,12 +418,80 @@ fn end(run: &mut Run, world: &World) { } } -fn dfs(run: &mut Run, world: World) { - if run.broken.is_some() || !run.seen.insert(format!("{world:?}")) { +/// Names a state's tags by first appearance. +#[derive(Default)] +struct Names(Vec); + +impl Names { + fn of(&mut self, tag: u32) -> u32 { + let at = self.0.iter().position(|&t| t == tag).unwrap_or_else(|| { + self.0.push(tag); + self.0.len() - 1 + }); + at as u32 + } +} + +/// A state's name in the search: everything it holds, with every tag renamed +/// by first appearance, the client's first in slot order, and the client's +/// table rendered by its filled slots. +fn key(world: &World) -> String { + let mut names = Names::default(); + let wire: Vec = world.client.on_the_wire().map(|t| names.of(t)).collect(); + let mut request = |r: &Request| Request { tag: names.of(r.tag), ..*r }; + let sq: Vec = world.queues.sq.iter().map(&mut request).collect(); + let device: Vec = world.device.iter().map(|(_, r)| request(r)).collect(); + let cq: Vec = world.queues.cq.iter().map(|c| Completion { tag: names.of(c.tag), ..*c }).collect(); + let server = world.server.as_ref().map(|s| { + let inflight: Vec<(u32, Op)> = s.session.inflight().map(|(t, op)| (names.of(t), op)).collect(); + format!("{inflight:?} {:?}", s.holds) + }); + let posted: Vec = world.posted.iter().map(|&t| names.of(t)).collect(); + format!( + "{:?} {wire:?} {sq:?} {device:?} {cq:?} {server:?} {posted:?} {} {:?} {:?} {} {} {:?} {:?} {:?}", + world.client, + world.next, + world.answers, + world.acked_before, + world.alive, + world.losses, + world.cache, + world.media, + world.left + ) +} + +fn dfs(run: &mut Run, mut world: World) { + world.queues.hold_empty(); + if run.broken.is_some() || !run.seen.insert(key(&world)) { return; } - let mut moved = false; - let script = run.script; + run.filled[0] |= world.queues.sq.len() == DEPTH as usize; + run.filled[1] |= world.queues.cq.len() == DEPTH as usize; + let next = next(run.script, &world); + if next.is_empty() { + run.ends += 1; + end(run, &world); + } + for (step, after) in next { + run.path.push(step); + match after { + Ok(world) => dfs(run, world), + Err((law, why)) => fail(run, law, why), + } + run.path.pop(); + } +} + +/// The world after an event, or the law it broke. +type After = Result; + +/// Every event that can happen next, named, and what it leads to. +fn next(script: &[Step], world: &World) -> Vec<(String, After)> { + let mut next = Vec::new(); + let mut after = |step: String, after: After| { + next.push((step, after.and_then(|mut w| settle(script, &mut w).map(|()| w)))); + }; // The caller asks for its next step, never a write to a block with a write // still unanswered: two writes in flight to one block are unordered. @@ -330,7 +505,10 @@ fn dfs(run: &mut Run, world: World) { if !busy { let mut w = world.clone(); match script[world.next] { - Step::Write { block, .. } => w.client.submit(ticket, Op::Write, block, 1, ticket as u32), + Step::Write { block, .. } => { + let run = ARENA.run(ticket as u32, 1).expect("a script's ticket is an arena block"); + w.client.submit(ticket, Op::Write { run, lba: block }); + } Step::Flush => { let acked = w .answers @@ -339,27 +517,30 @@ fn dfs(run: &mut Run, world: World) { .map(|(t, _)| *t) .collect(); w.acked_before.insert(ticket, acked); - w.client.submit(ticket, Op::Flush, 0, 0, 0); + w.client.submit(ticket, Op::Flush); } } w.next += 1; - pump(&mut w); - moved = true; - go(run, format!("ask {}", world.next), w); + after(format!("ask {}", world.next), Ok(w)); } } // The server takes the oldest request. - if world.alive && !world.sq.is_empty() { + if world.alive && !world.queues.sq.is_empty() { let mut w = world.clone(); - let request = w.sq.pop_front().expect("just seen"); + let words = w.queues.take().expect("just seen"); let server = w.server.as_mut().expect("alive"); - match server.session.take(request.encode()) { - Taken::Issue(request) => w.device.push((request.tag, request)), - Taken::Answer(c) => w.cq.push_back(c), - } - moved = true; - go(run, format!("take {:?}#{}", request.op, request.tag), w); + let step = match server.session.take(words) { + Taken::Issue(request) => { + w.device.push((request.tag, request)); + format!("take {:?}#{}", request.op, request.tag) + } + Taken::Answer(c) => { + w.queues.post(c); + format!("refuse #{}", c.tag) + } + }; + after(step, Ok(w)); } // The device completes any one command it holds. @@ -370,11 +551,10 @@ fn dfs(run: &mut Run, world: World) { let losses = w.losses; if let Some(server) = w.server.as_mut() { if let Some(c) = server.session.complete(tag, true, &mut server.holds, losses) { - w.cq.push_back(c); + w.queues.post(c); } } - moved = true; - go(run, format!("done {:?}#{tag}", request.op), w); + after(format!("done {:?}#{tag}", request.op), Ok(w)); } // The device fails any one command it holds: not done, and answered so. @@ -386,25 +566,21 @@ fn dfs(run: &mut Run, world: World) { let losses = w.losses; if let Some(server) = w.server.as_mut() { if let Some(c) = server.session.complete(tag, false, &mut server.holds, losses) { - w.cq.push_back(c); + w.queues.post(c); } } - moved = true; - go(run, format!("fail {:?}#{tag}", request.op), w); + after(format!("fail {:?}#{tag}", request.op), Ok(w)); } } // The client reads the oldest completion. - if world.client.up() && !world.cq.is_empty() { + if world.client.up() && !world.queues.cq.is_empty() { let mut w = world.clone(); - let c = w.cq.pop_front().expect("just seen"); - if w.client.complete(c).is_err() { - fail(run, Law::Answers, format!("the client met a second completion for tag {}", c.tag)); - return; - } - collect(run, &mut w); - moved = true; - go(run, format!("read #{} {:?}", c.tag, c.status), w); + let c = w.queues.read().expect("just seen"); + let read = w.client.complete(c).map(|()| w).map_err(|_| { + (Law::Answers, format!("the client met a second completion for tag {}", c.tag)) + }); + after(format!("read #{} {:?}", c.tag, c.status), read); } // The server reads a completion the device posted before it was reset: @@ -415,15 +591,14 @@ fn dfs(run: &mut Run, world: World) { let losses = w.losses; let server = w.server.as_mut().expect("alive"); if let Some(c) = server.session.complete(tag, true, &mut server.holds, losses) { - w.cq.push_back(c); + w.queues.post(c); } - moved = true; - go(run, format!("posted #{tag}"), w); + after(format!("posted #{tag}"), Ok(w)); } // The device is reset under everything in flight. if world.left.resets > 0 && world.alive { - for (posted, cache, media) in fates(script, &world, true) { + for (posted, cache, media) in fates(script, world, true) { let mut w = world.clone(); w.left.resets -= 1; w.device.clear(); @@ -432,15 +607,16 @@ fn dfs(run: &mut Run, world: World) { w.media = media; w.losses += 1; let server = w.server.as_mut().expect("alive"); - w.cq.extend(server.session.abort_all()); - moved = true; - go(run, format!("reset({posted:?} {cache:?} {media:?})"), w); + for c in server.session.abort_all() { + w.queues.post(c); + } + after(format!("reset({posted:?} {cache:?} {media:?})"), Ok(w)); } } // The server dies. if world.left.crashes > 0 && world.alive { - for (_, cache, media) in fates(script, &world, false) { + for (_, cache, media) in fates(script, world, false) { let mut w = world.clone(); w.left.crashes -= 1; w.device.clear(); @@ -449,8 +625,7 @@ fn dfs(run: &mut Run, world: World) { w.media = media; w.server = None; w.alive = false; - moved = true; - go(run, format!("crash({cache:?} {media:?})"), w); + after(format!("crash({cache:?} {media:?})"), Ok(w)); } } @@ -458,25 +633,18 @@ fn dfs(run: &mut Run, world: World) { if !world.alive && world.client.up() { let mut w = world.clone(); w.client.session_ended(); - w.sq.clear(); - w.cq.clear(); - collect(run, &mut w); - moved = true; - go(run, "notice".into(), w); + w.queues.reset(); + after("notice".into(), Ok(w)); } // It reconnects to the server started in its place. if !world.alive && !world.client.up() { let mut w = world.clone(); connect(&mut w); - moved = true; - go(run, "reconnect".into(), w); + after("reconnect".into(), Ok(w)); } - if !moved { - run.ends += 1; - end(run, &world); - } + next } #[cfg(test)] @@ -490,6 +658,7 @@ mod tests { Step::Write { block: 0, value: 3 }, Step::Flush, ]; + const _: () = assert!(SCRIPT.len() > DEPTH as usize, "a ring never wraps"); const fn at_most(resets: u8, crashes: u8, errors: u8) -> Failures { Failures { resets, crashes, errors } @@ -567,12 +736,15 @@ mod tests { } /// The model is not vacuous: with no failure at all, a run ends with both - /// flushes durable and the medium holding the script's last values. + /// flushes durable and the medium holding the script's last values; each + /// ring holds its depth, and wraps, because one session carries more than + /// that. #[test] fn the_model_reaches_the_end_it_should() { let explored = verdict(&SCRIPT, at_most(0, 0, 0)); assert_eq!(explored.broken, None); assert!(explored.ends >= 1); + assert_eq!(explored.filled, [true, true], "a ring never held its depth"); } /// Nor is the give-up path out of its reach: past [`MAX_ATTEMPTS`] failures @@ -584,4 +756,33 @@ mod tests { assert_eq!(explored.broken, None); assert!(explored.given_up > 0, "no run gave a write up"); } + + /// The world after each event in turn, each the first whose name starts with + /// the word given. + fn walk(mut world: World, events: &[&str]) -> World { + for event in events { + let (_, after) = next(&SCRIPT, &world) + .into_iter() + .find(|(step, _)| step.starts_with(event)) + .unwrap_or_else(|| panic!("no {event} after {}", key(&world))); + world = after.expect("no law is broken on the way"); + } + world + } + + /// Two states a key that orders tags by value merges, though a reset parts + /// them: F2 on the device and slot 1 free, reached with W1 asked after W0 + /// was answered, and with the two in flight together. Named alike, they act + /// alike: with W3 taken and the device reset, they are still named alike. + #[test] + fn states_named_alike_act_alike() { + let fresh = start(at_most(1, 0, 0)); + let answered_first = + walk(fresh.clone(), &["ask", "take", "done", "read", "ask", "take", "done", "read", "ask", "take"]); + let together = walk(fresh, &["ask", "ask", "take", "done", "read", "take", "done", "read", "ask", "take"]); + assert_ne!(answered_first.device, together.device, "the flush on the device has one tag in both"); + assert_eq!(key(&answered_first), key(&together), "the search tells the two apart"); + let [answered_first, together] = [answered_first, together].map(|w| walk(w, &["ask", "take", "reset"])); + assert_eq!(key(&answered_first), key(&together), "a reset answered the two in different orders"); + } } diff --git a/toyos-blockring/src/ring.rs b/toyos-blockring/src/ring.rs deleted file mode 100644 index 9b148d4452..0000000000 --- a/toyos-blockring/src/ring.rs +++ /dev/null @@ -1,241 +0,0 @@ -//! The two single-producer rings on a session's page. -//! -//! **A producer writes an entry's words, then publishes the tail with -//! `Release`; a consumer loads the tail with `Acquire`, then reads the words.** -//! The head goes back the other way: a consumer publishes it with `Release` -//! only after the entries below it are read, and a producer loads it with -//! `Acquire` before it writes over them. `tests/loom_ring.rs` holds the -//! first edge; the second is the same pair turned round. -//! -//! Each end keeps its own index in a local and only ever *stores* the shared -//! one, so nothing the peer writes into its own index can move ours; what the -//! peer's index claims is bounded against ours before it is believed -//! ([`Violation`]). - -use core::sync::atomic::{AtomicU32, Ordering}; - -use crate::layout::{CQE_WORDS, CQ_BASE, CQ_HEAD, CQ_TAIL, DEPTH, RING_WORDS, SQE_WORDS, SQ_BASE, SQ_HEAD, SQ_TAIL}; - -/// One shared 32-bit word: an atomic over the mapped page, or a model's. -pub trait Word { - fn load(&self, order: Ordering) -> u32; - fn store(&self, value: u32, order: Ordering); -} - -impl Word for AtomicU32 { - fn load(&self, order: Ordering) -> u32 { - AtomicU32::load(self, order) - } - fn store(&self, value: u32, order: Ordering) { - AtomicU32::store(self, value, order) - } -} - -/// The peer's index says something no peer of this protocol can: more entries -/// published than the ring holds, or more consumed than were produced. The -/// session is over. -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub struct Violation; - -/// The publishing order a tail is stored with. -#[cfg(not(feature = "mutate-ring-publish-relaxed"))] -const PUBLISH: Ordering = Ordering::Release; -#[cfg(feature = "mutate-ring-publish-relaxed")] -const PUBLISH: Ordering = Ordering::Relaxed; - -/// Where one ring is on the page, in words. -#[derive(Clone, Copy, Debug)] -struct Place { - head: usize, - tail: usize, - entries: usize, -} - -const REQUESTS: Place = Place { head: SQ_HEAD, tail: SQ_TAIL, entries: SQ_BASE }; -const COMPLETIONS: Place = Place { head: CQ_HEAD, tail: CQ_TAIL, entries: CQ_BASE }; - -fn checked(page: &[W]) -> &[W] { - assert!(page.len() >= RING_WORDS, "a session page holds every ring word"); - page -} - -/// The end of a ring that writes entries. It holds its indices and not the -/// page, so its owner keeps the mapping beside it; every call is given the -/// page. -#[derive(Debug)] -pub struct Producer { - place: Place, - local: u32, - published: u32, -} - -impl Producer { - fn new(page: &[W], place: Place) -> Self { - checked(page)[place.tail].store(0, Ordering::Release); - Self { place, local: 0, published: 0 } - } - - /// How many entries may be pushed before the consumer frees more. - pub fn space(&self, page: &[W]) -> Result { - let head = checked(page)[self.place.head].load(Ordering::Acquire); - let used = self.local.wrapping_sub(head); - if used > DEPTH { - return Err(Violation); - } - Ok(DEPTH - used) - } - - /// Write one entry. It is the consumer's only once [`Self::publish`] runs. - /// - /// # Panics - /// When the ring has no space: the caller asks [`Self::space`] first. - pub fn push(&mut self, page: &[W], words: [u32; N]) { - assert!(self.space(page).is_ok_and(|space| space > 0), "a push into a full ring"); - let at = self.place.entries + (self.local % DEPTH) as usize * N; - for (i, word) in words.into_iter().enumerate() { - page[at + i].store(word, Ordering::Relaxed); - } - self.local = self.local.wrapping_add(1); - } - - /// Publish every entry pushed so far; answers whether there was any. - pub fn publish(&mut self, page: &[W]) -> bool { - if self.published == self.local { - return false; - } - checked(page)[self.place.tail].store(self.local, PUBLISH); - self.published = self.local; - true - } -} - -/// The end of a ring that reads entries; like [`Producer`], it holds indices -/// and is given the page. -#[derive(Debug)] -pub struct Consumer { - place: Place, - local: u32, - released: u32, -} - -impl Consumer { - fn new(page: &[W], place: Place) -> Self { - checked(page)[place.head].store(0, Ordering::Release); - Self { place, local: 0, released: 0 } - } - - /// The next published entry, `None` for none, or the producer's tail - /// claiming more than the ring holds. - pub fn pop(&mut self, page: &[W]) -> Result, Violation> { - let tail = checked(page)[self.place.tail].load(Ordering::Acquire); - let ready = tail.wrapping_sub(self.local); - if ready > DEPTH { - return Err(Violation); - } - if ready == 0 { - return Ok(None); - } - let at = self.place.entries + (self.local % DEPTH) as usize * N; - let words = core::array::from_fn(|i| page[at + i].load(Ordering::Relaxed)); - self.local = self.local.wrapping_add(1); - Ok(Some(words)) - } - - /// Give every entry popped so far back to the producer. - pub fn release(&mut self, page: &[W]) { - if self.released != self.local { - checked(page)[self.place.head].store(self.local, Ordering::Release); - self.released = self.local; - } - } -} - -/// A client's two ends: requests out, completions in. -pub type ClientRings = (Producer, Consumer); - -/// A server's two ends: requests in, completions out. -pub type ServerRings = (Consumer, Producer); - -/// The client's ends of a session page, both indices it owns set to 0. Done -/// before the page is sent to a server, and again before it is sent to the -/// next one. -pub fn client(page: &[W]) -> ClientRings { - (Producer::new(page, REQUESTS), Consumer::new(page, COMPLETIONS)) -} - -/// The server's ends of a session page it was sent, both indices it owns set -/// to 0. Whatever the client left in its own two is bounded by the first -/// [`Consumer::pop`] and [`Producer::space`]. -pub fn server(page: &[W]) -> ServerRings { - (Consumer::new(page, REQUESTS), Producer::new(page, COMPLETIONS)) -} - -#[cfg(test)] -mod tests { - use super::*; - use crate::entry::{Completion, Op, Request, Status}; - use alloc::vec::Vec; - - fn page() -> Vec { - (0..RING_WORDS).map(|_| AtomicU32::new(0)).collect() - } - - #[test] - fn a_request_crosses_and_its_completion_comes_back() { - let page = page(); - let (mut sq, mut cq) = client(&page); - let (mut rq, mut cp) = server(&page); - let request = Request { op: Op::Write, tag: 3, lba: 8, blocks: 2, arena: 1 }; - sq.push(&page, request.encode()); - assert_eq!(rq.pop(&page), Ok(None), "an entry is nobody's before it is published"); - assert!(sq.publish(&page)); - assert_eq!(rq.pop(&page).map(|w| w.map(|w| Request::decode(w, 100))), Ok(Some(Ok(request)))); - rq.release(&page); - cp.push(&page, Completion { tag: 3, status: Status::Ok }.encode()); - cp.publish(&page); - assert_eq!(cq.pop(&page).map(|w| w.and_then(Completion::decode)), Ok(Some(Completion { tag: 3, status: Status::Ok }))); - } - - /// The ring wraps many times over, and space is exactly what the consumer - /// has given back. - #[test] - fn the_ring_wraps_and_counts_its_space() { - let page = page(); - let (mut sq, _) = client(&page); - let (mut rq, _) = server(&page); - let (mut pushed, mut popped) = (0u32, 0u32); - for round in 0..5 * DEPTH { - // Uneven batches both ways, so the two indices cross every slot - // at every distance. - for _ in 0..1 + round % 9 { - assert_eq!(sq.space(&page), Ok(DEPTH - (pushed - popped))); - if pushed - popped == DEPTH { - break; - } - sq.push(&page, Request { op: Op::Read, tag: pushed, lba: 0, blocks: 1, arena: 0 }.encode()); - pushed += 1; - } - sq.publish(&page); - for _ in 0..1 + round % 7 { - let Ok(Some(words)) = rq.pop(&page) else { break }; - assert_eq!(words[1], popped, "entries come out in the order they went in"); - popped += 1; - } - rq.release(&page); - } - assert!(pushed > 2 * DEPTH, "the ring wrapped"); - } - - /// A hostile producer's tail far ahead of the consumer, and a hostile - /// consumer's head ahead of what was produced, both end the session. - #[test] - fn a_peer_index_out_of_reach_is_a_violation() { - let page = page(); - let (sq, _) = client(&page); - let (mut rq, _) = server(&page); - page[SQ_TAIL].store(DEPTH + 1, Ordering::Release); - assert_eq!(rq.pop(&page), Err(Violation)); - page[SQ_HEAD].store(5, Ordering::Release); - assert_eq!(sq.space(&page), Err(Violation)); - } -} diff --git a/toyos-blockring/src/server.rs b/toyos-blockring/src/server.rs index 6157ffd6ad..24ca1703bf 100644 --- a/toyos-blockring/src/server.rs +++ b/toyos-blockring/src/server.rs @@ -15,9 +15,10 @@ //! own since it last heard. The loss count is the device's, bumped by whoever //! resets it. -use alloc::collections::BTreeMap; +use alloc::vec::Vec; use toyos_blockhold::{Holds, Writer}; +use toyos_transport::Untrusted; use crate::entry::{Completion, Op, Refused, Request, Status}; use crate::layout::SQE_WORDS; @@ -30,8 +31,9 @@ pub struct ServerSession { blocks: u64, /// The span of the device this session holds, which is its writer. first: u64, - /// By tag: the op each request in flight asked for. - inflight: BTreeMap, + /// Each request in flight, in the order it was taken: a tag is only ever + /// compared for equality. + inflight: Vec<(u32, Op)>, } /// A request the device is to be asked for. @@ -47,7 +49,7 @@ impl ServerSession { /// A session over the partition at device block `first`, `blocks` long, /// which `holds` already holds for it. pub fn new(first: u64, blocks: u64) -> Self { - Self { blocks, first, inflight: BTreeMap::new() } + Self { blocks, first, inflight: Vec::new() } } /// The writer this session's writes and flushes are accounted to. @@ -65,9 +67,9 @@ impl ServerSession { self.blocks } - /// Requests taken and not yet answered. - pub fn inflight(&self) -> usize { - self.inflight.len() + /// Requests taken and not yet answered, in the order they were taken. + pub fn inflight(&self) -> impl ExactSizeIterator + '_ { + self.inflight.iter().copied() } /// Decide what one entry the client published is. @@ -75,17 +77,17 @@ impl ServerSession { /// A malformed entry, and a tag already in flight, are answered at once /// and never reach the device: the second would make one tag two /// requests, and the client could not tell which answer was whose. - pub fn take(&mut self, words: [u32; SQE_WORDS]) -> Taken { + pub fn take(&mut self, words: [Untrusted; SQE_WORDS]) -> Taken { let request = match Request::decode(words, self.blocks) { Ok(request) => request, Err(Refused::Malformed { tag }) => { return Taken::Answer(Completion { tag, status: Status::Invalid }) } }; - if self.inflight.contains_key(&request.tag) { + if self.inflight.iter().any(|&(tag, _)| tag == request.tag) { return Taken::Answer(Completion { tag: request.tag, status: Status::Invalid }); } - self.inflight.insert(request.tag, request.op); + self.inflight.push((request.tag, request.op)); Taken::Issue(request) } @@ -99,11 +101,12 @@ impl ServerSession { holds: &mut Holds, losses: u64, ) -> Option { - let op = self.inflight.remove(&tag)?; + let at = self.inflight.iter().position(|&(t, _)| t == tag)?; + let (_, op) = self.inflight.remove(at); let status = match (op, done) { (_, false) => Status::Device, - (Op::Read, true) => Status::Ok, - (Op::Write, true) => { + (Op::Read { .. }, true) => Status::Ok, + (Op::Write { .. }, true) => { holds.wrote(self.writer(), losses); Status::Ok } @@ -116,14 +119,10 @@ impl ServerSession { } /// The device was reset under every request in flight: each is answered - /// [`Status::Device`], and a late answer from the device for any of them - /// finds nothing to complete. - pub fn abort_all(&mut self) -> alloc::vec::Vec { - let answered = self - .inflight - .keys() - .map(|&tag| Completion { tag, status: Status::Device }) - .collect(); + /// [`Status::Device`], in the order it was taken, and a late answer from + /// the device for any of them finds nothing to complete. + pub fn abort_all(&mut self) -> Vec { + let answered = self.inflight.iter().map(|&(tag, _)| Completion { tag, status: Status::Device }).collect(); #[cfg(not(feature = "mutate-abort-keeps-inflight"))] self.inflight.clear(); answered @@ -133,12 +132,13 @@ impl ServerSession { #[cfg(test)] mod tests { use super::*; + use crate::layout::ARENA; - fn write(tag: u32) -> [u32; SQE_WORDS] { - Request { op: Op::Write, tag, lba: 0, blocks: 1, arena: 0 }.encode() + fn write(tag: u32) -> [Untrusted; SQE_WORDS] { + Request { op: Op::Write { run: ARENA.run(0, 1).unwrap(), lba: 0 }, tag }.encode().map(Untrusted::new) } - fn flush(tag: u32) -> [u32; SQE_WORDS] { - Request { op: Op::Flush, tag, lba: 0, blocks: 0, arena: 0 }.encode() + fn flush(tag: u32) -> [Untrusted; SQE_WORDS] { + Request { op: Op::Flush, tag }.encode().map(Untrusted::new) } #[test] @@ -151,12 +151,24 @@ mod tests { assert_eq!(session.complete(1, true, &mut holds, 1), None); } + /// A reset answers in the order the requests were taken, whatever their + /// tags' values: nothing orders a tag but its arrival. + #[test] + fn a_reset_answers_in_the_order_taken() { + let mut session = ServerSession::new(0, 10); + for tag in [9, 4] { + assert!(matches!(session.take(write(tag)), Taken::Issue(_))); + } + let answered: Vec = session.abort_all().iter().map(|c| c.tag).collect(); + assert_eq!(answered, [9, 4]); + } + #[test] fn a_tag_in_flight_twice_is_refused_unissued() { let mut session = ServerSession::new(0, 10); assert!(matches!(session.take(write(4)), Taken::Issue(_))); assert_eq!(session.take(write(4)), Taken::Answer(Completion { tag: 4, status: Status::Invalid })); - assert_eq!(session.inflight(), 1); + assert_eq!(session.inflight().len(), 1); } /// A write acknowledged before a reset bumped the loss count is a loss the diff --git a/toyos-blockring/tests/loom_ring.rs b/toyos-blockring/tests/loom_ring.rs deleted file mode 100644 index 26bcef0b3e..0000000000 --- a/toyos-blockring/tests/loom_ring.rs +++ /dev/null @@ -1,69 +0,0 @@ -//! The rings' publication edge under loom: a consumer that sees a tail sees -//! every word of the entries below it. -//! -//! The client and the server are two processes on two CPUs over one shared -//! page, so this is the one property of the rings no host test that runs both -//! ends on one thread can reach. `mutate-ring-publish-relaxed` takes the edge -//! away and this must red: -//! -//! cargo test -p toyos-blockring --features mutate-ring-publish-relaxed --test loom_ring - -use core::sync::atomic::Ordering; - -use loom::sync::atomic::AtomicU32; -use loom::sync::Arc; -use toyos_blockring::entry::{Op, Request}; -use toyos_blockring::layout::RING_WORDS; -use toyos_blockring::ring::{self, Word}; - -/// A loom atomic as a page word: the trait is this crate's and the type is -/// loom's, so the two meet through a wrapper. -struct Shared(AtomicU32); - -impl Word for Shared { - fn load(&self, order: Ordering) -> u32 { - self.0.load(order) - } - fn store(&self, value: u32, order: Ordering) { - self.0.store(value, order) - } -} - -fn page() -> Arc> { - Arc::new((0..RING_WORDS).map(|_| Shared(AtomicU32::new(0))).collect()) -} - -const PARTITION: u64 = 1 << 20; - -fn request(n: u32) -> Request { - Request { op: Op::Write, tag: 100 + n, lba: 7 + u64::from(n), blocks: 1 + n, arena: 3 * n } -} - -/// Two requests published one at a time, read by the other end as they -/// arrive: each is whole, in order, and exactly what was written. -#[test] -fn a_published_request_is_read_whole() { - loom::model(|| { - let page = page(); - let server_page = Arc::clone(&page); - let server = loom::thread::spawn(move || { - let (mut requests, _) = ring::server(&server_page); - let mut read = Vec::new(); - while read.len() < 2 { - match requests.pop(&server_page).expect("the client keeps the protocol") { - Some(words) => read.push(Request::decode(words, PARTITION)), - None => loom::thread::yield_now(), - } - } - requests.release(&server_page); - read - }); - let (mut requests, _) = ring::client(&page); - for n in 0..2 { - requests.push(&page, request(n).encode()); - requests.publish(&page); - } - let read = server.join().expect("the server thread"); - assert_eq!(read, [Ok(request(0)), Ok(request(1))], "a request was read before its words"); - }); -} diff --git a/toyos-transport/Cargo.toml b/toyos-transport/Cargo.toml new file mode 100644 index 0000000000..cd1fbf10a0 --- /dev/null +++ b/toyos-transport/Cargo.toml @@ -0,0 +1,38 @@ +# A member of the host workspace (root `Cargo.toml`). The one transport every +# client/server session runs over: the rings, the arena's runs and the tag +# table, as decisions over words an adapter hands in. Nothing here maps, copies +# or blocks. +# +# Its defects are orders and hostile peers, so `tests/loom.rs` holds the +# publication edge and a hostile peer. + +[package] +name = "toyos-transport" +version = "0.1.0" +edition = "2021" +license = "MIT OR Apache-2.0" +publish = false + +[features] +# Declared, never enabled by any build that ships: each reverts one decision so +# the model that refuses it is shown able to red at all. `src/ci.rs`'s +# `CONTROLS` runs each and demands the named test's own FAILED line. +# +# A producer publishes its tail `Relaxed`, so a consumer can see the index +# before the entry's words. +publish-relaxed = [] +# A peer's cursor is believed however far it claims to be, so a hostile +# producer hands the consumer more entries than the ring holds. +no-clamp = [] +# `Inflight::end` answers nothing and keeps every tag in flight, so a session's +# requests outlive it unanswered and a late completion answers one. +end-keeps-inflight = [] + +[dependencies] +toyos-untrusted = { path = "../toyos-untrusted" } + +[dev-dependencies] +loom = "0.7" + +[lints.rust] +warnings = "deny" diff --git a/toyos-transport/src/arena.rs b/toyos-transport/src/arena.rs new file mode 100644 index 0000000000..d85fbe9fe9 --- /dev/null +++ b/toyos-transport/src/arena.rs @@ -0,0 +1,121 @@ +//! A region's arena, cut into slots, and the runs of it an entry names. +//! +//! **A run exists only once it is bounded by the [`Geometry`]**: a peer's is +//! decoded ([`Run::decode`]) and this side's is cut ([`Geometry::run`]). The +//! adapter reaches a run's bytes only through its [`Span`]. + +use core::fmt; + +use crate::{Untrusted, Violation}; + +/// Bytes `offset..offset + len` of the region. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub struct Span { + pub offset: usize, + pub len: usize, +} + +/// A region's arena: whole slots after its header page. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub struct Geometry { + slot_bytes: u32, + slots: u32, +} + +impl Geometry { + /// Every region is this long: the granule shared memory comes in. + pub const BYTES: u32 = 0x20_0000; + /// The header page the queues and cursors are on; the arena follows it. + pub const HEADER_BYTES: u32 = 0x1000; + const ARENA_BYTES: u32 = Self::BYTES - Self::HEADER_BYTES; + + /// The arena cut into slots of `slot_bytes`; `None` if not one fits. + pub const fn new(slot_bytes: u32) -> Option { + match Self::ARENA_BYTES.checked_div(slot_bytes) { + Some(slots) if slots > 0 => Some(Self { slot_bytes, slots }), + _ => None, + } + } + + pub const fn slots(&self) -> u32 { + self.slots + } + + /// Slots `first..first + count`; `None` for none, or past the arena. + pub fn run(&self, first: u32, count: u32) -> Option { + if count == 0 || first.checked_add(count)? > self.slots { + return None; + } + let offset = first.checked_mul(self.slot_bytes)?.checked_add(Self::HEADER_BYTES)?; + let len = count.checked_mul(self.slot_bytes)?; + Some(Run { first, count, span: Span { offset: offset.try_into().ok()?, len: len.try_into().ok()? } }) + } +} + +/// Slots of the arena, bounded by its geometry. +#[derive(Clone, Copy, PartialEq, Eq, Hash)] +pub struct Run { + first: u32, + count: u32, + span: Span, +} + +/// Its slots alone: the span follows from them. +impl fmt::Debug for Run { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("Run").field("first", &self.first).field("count", &self.count).finish() + } +} + +impl Run { + /// The run a peer's two words name. + pub fn decode(first: Untrusted, count: Untrusted, geometry: &Geometry) -> Result { + let bounded = |word: Untrusted| word.at_most(u64::from(geometry.slots)).ok()?.try_into().ok(); + bounded(first) + .zip(bounded(count)) + .and_then(|(first, count)| geometry.run(first, count)) + .ok_or(Violation::Run) + } + + pub fn first(&self) -> u32 { + self.first + } + + pub fn count(&self) -> u32 { + self.count + } + + /// Its bytes in the region. + pub fn span(&self) -> Span { + self.span + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const PAGES: Geometry = Geometry::new(4096).unwrap(); + + #[test] + fn an_arena_is_whole_slots_after_the_header() { + assert_eq!(PAGES.slots(), 511); + assert_eq!(Geometry::new(0), None); + assert_eq!(Geometry::new(Geometry::BYTES), None, "a slot longer than the arena"); + } + + /// Every run a peer can name is inside the arena, whole slots of it; the + /// last slot is a run and one past it is not, however the sum wraps. + #[test] + fn a_run_a_peer_names_is_inside_the_arena() { + let run = Run::decode(Untrusted::new(510), Untrusted::new(1), &PAGES).unwrap(); + assert_eq!(run.span(), Span { offset: 0x1000 + 510 * 4096, len: 4096 }); + for (first, count) in [(510, 2), (0, 0), (511, 1), (1, u32::MAX), (u32::MAX, 1), (0, 512)] { + assert_eq!( + Run::decode(Untrusted::new(first), Untrusted::new(count), &PAGES), + Err(Violation::Run), + "{first}+{count}" + ); + } + } +} diff --git a/toyos-transport/src/inflight.rs b/toyos-transport/src/inflight.rs new file mode 100644 index 0000000000..9a45176bdb --- /dev/null +++ b/toyos-transport/src/inflight.rs @@ -0,0 +1,158 @@ +//! The requests a client has on the wire, by tag. +//! +//! **Each tag is answered exactly once**: by the completion that names it +//! ([`Inflight::answer`]) or, when the session ends, by [`Inflight::end`] — +//! never by both, and never by a completion naming a tag from before its +//! slot was last filled. A tag is the slot's index under the slot's own +//! sequence, so a tag answered, ended or replayed names nothing. The index +//! takes the bits the table needs and the sequence the rest, and wraps: a tag +//! replayed a wrap later answers the request then in its slot, which a server +//! could answer by naming that request's own tag, so it grants the peer +//! nothing. + +use core::fmt; + +use crate::{Untrusted, Violation}; + +#[derive(Clone, PartialEq, Eq, Hash)] +struct Slot { + seq: u32, + value: Option, +} + +/// At most `D` requests in flight, each carrying what its answer needs. +#[derive(Clone, PartialEq, Eq, Hash)] +pub struct Inflight { + slots: [Slot; D], +} + +impl Inflight { + const INDEX_BITS: u32 = usize::BITS.wrapping_sub(D.saturating_sub(1).leading_zeros()); + const INDEX_MASK: u32 = 1u32.wrapping_shl(Self::INDEX_BITS).wrapping_sub(1); + const SEQ_MASK: u32 = u32::MAX.wrapping_shr(Self::INDEX_BITS); + + pub fn new() -> Self { + const { assert!(D > 0 && Self::INDEX_BITS <= 16, "a tag keeps at least 16 bits of sequence") }; + Self { slots: core::array::from_fn(|_| Slot { seq: 0, value: None }) } + } + + fn tag(seq: u32, index: u32) -> u32 { + seq.wrapping_shl(Self::INDEX_BITS) | index + } + + /// Keep `value` in flight under a tag no earlier request had in this slot; + /// `Err(value)` with `D` in flight. + pub fn insert(&mut self, value: T) -> Result { + let Some((index, slot)) = (0u32..).zip(self.slots.iter_mut()).find(|(_, slot)| slot.value.is_none()) else { + return Err(value); + }; + slot.seq = slot.seq.wrapping_add(1) & Self::SEQ_MASK; + slot.value = Some(value); + Ok(Self::tag(slot.seq, index)) + } + + /// What the peer's `tag` answers; a tag nothing is in flight under is a + /// violation. + pub fn answer(&mut self, tag: Untrusted) -> Result { + let index = tag.map(|t| t & Self::INDEX_MASK).index(D).map_err(|_| Violation::Tag)?; + let slot = self.slots.get_mut(index).ok_or(Violation::Tag)?; + if !tag.map(|t| t.wrapping_shr(Self::INDEX_BITS)).is(slot.seq) { + return Err(Violation::Tag); + } + slot.value.take().ok_or(Violation::Tag) + } + + /// What is in flight, in no order. + pub fn values(&self) -> impl Iterator { + self.slots.iter().filter_map(|slot| slot.value.as_ref()) + } + + /// Every tag in flight, by its slot. + pub fn tags(&self) -> impl Iterator + '_ { + (0u32..).zip(&self.slots).filter(|(_, slot)| slot.value.is_some()).map(|(index, slot)| Self::tag(slot.seq, index)) + } + + /// The session ended: `each` is given every tag in flight, once, with what + /// it carried, and none of them is answered again. + pub fn end(&mut self, mut each: impl FnMut(u32, T)) { + for (index, slot) in (0u32..).zip(self.slots.iter_mut()) { + #[cfg(not(feature = "end-keeps-inflight"))] + let value = slot.value.take(); + #[cfg(feature = "end-keeps-inflight")] + let value = None; + if let Some(value) = value { + each(Self::tag(slot.seq, index), value); + } + } + } +} + +impl Default for Inflight { + fn default() -> Self { + Self::new() + } +} + +/// Only what is in flight, by slot. +impl fmt::Debug for Inflight { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_map().entries((0u32..).zip(&self.slots).filter_map(|(i, slot)| Some((i, slot.value.as_ref()?)))).finish() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_tag_is_answered_once_and_its_slot_takes_a_new_one() { + let mut inflight = Inflight::<&str, 2>::new(); + let a = inflight.insert("a").unwrap(); + let b = inflight.insert("b").unwrap(); + assert_eq!(inflight.insert("c"), Err("c")); + assert_eq!(inflight.answer(Untrusted::new(a)), Ok("a")); + assert_eq!(inflight.answer(Untrusted::new(a)), Err(Violation::Tag), "answered twice"); + let c = inflight.insert("c").unwrap(); + assert_ne!(c, a, "the slot's next tag is not its last"); + let mut held: Vec<&str> = inflight.values().copied().collect(); + held.sort_unstable(); + assert_eq!(held, ["b", "c"], "what is in flight, and only that"); + assert_eq!(inflight.tags().collect::>(), [c, b], "the tags in flight, by slot"); + assert_eq!(inflight.answer(Untrusted::new(a)), Err(Violation::Tag), "a replay answers nothing"); + assert_eq!(inflight.answer(Untrusted::new(2)), Err(Violation::Tag), "an index past the table"); + assert_eq!(inflight.answer(Untrusted::new(b)), Ok("b")); + assert_eq!(inflight.answer(Untrusted::new(c)), Ok("c")); + assert!(inflight.insert("d").is_ok() && inflight.insert("e").is_ok(), "every slot is free again"); + } + + /// A slot filled 2^16 times over still refuses its first tag: the index + /// takes six bits of a 64-slot table's tag, not sixteen. + #[test] + fn a_slot_refilled_past_sixteen_bits_refuses_its_first_tag() { + let mut inflight = Inflight::<(), 64>::new(); + let first = inflight.insert(()).unwrap(); + inflight.answer(Untrusted::new(first)).unwrap(); + for _ in 0..u16::MAX { + let tag = inflight.insert(()).unwrap(); + inflight.answer(Untrusted::new(tag)).unwrap(); + } + inflight.insert(()).unwrap(); + assert_eq!(inflight.answer(Untrusted::new(first)), Err(Violation::Tag)); + } + + #[test] + fn an_end_answers_every_tag_once_and_a_late_completion_nothing() { + let mut inflight = Inflight::::new(); + let tags: Vec = (0..3).map(|n| inflight.insert(n).unwrap()).collect(); + assert_eq!(inflight.answer(Untrusted::new(tags[1])), Ok(1)); + let mut ended = Vec::new(); + inflight.end(|tag, value| ended.push((tag, value))); + assert_eq!(ended, [(tags[0], 0), (tags[2], 2)]); + for tag in tags { + assert_eq!(inflight.answer(Untrusted::new(tag)), Err(Violation::Tag)); + } + let mut again = 0; + inflight.end(|_, _| again += 1); + assert_eq!(again, 0, "an end answers nothing a first one answered"); + } +} diff --git a/toyos-transport/src/lib.rs b/toyos-transport/src/lib.rs new file mode 100644 index 0000000000..1fbfcbfecb --- /dev/null +++ b/toyos-transport/src/lib.rs @@ -0,0 +1,71 @@ +//! The transport every client/server session runs over. +//! +//! **A session is one connection and, at most, one region.** The connection +//! carries the hello, handles, doorbells and the hang-up; the region carries +//! [`Producer`]/[`Consumer`] queues of fixed-size entries and an arena of +//! [`Run`]s. A protocol says what the entries mean; this crate says only what +//! is safe. +//! +//! **Nothing here holds the region.** An adapter hands every call the `N` +//! [`Word`]s an end was made over, and copies bytes itself, through a run's +//! [`Span`], once: no reference to the peer's bytes is formed. +//! +//! **What the peer writes is untrusted until decoded.** An entry comes out as +//! [`Untrusted`] words; a peer's cursor is bounded against the ring before it +//! is believed, and the variant of [`Violation`] names what it claimed; a run +//! is bounded against the [`Geometry`], and a tag against the [`Inflight`] +//! table. An end keeps its own cursor locally and only stores the shared one, +//! so nothing the peer writes moves it, and a cursor the peer moves backwards +//! within bounds costs only the peer. +//! +//! **A restart is seen only as a hang-up**, and [`Inflight::end`] answers every +//! tag that was in flight, once; a tag from before it answers nothing. + +#![cfg_attr(not(test), no_std)] +#![forbid(unsafe_code)] +#![cfg_attr(not(test), forbid(clippy::arithmetic_side_effects, clippy::unwrap_used, clippy::expect_used, clippy::panic))] +#![cfg_attr(not(test), deny(clippy::indexing_slicing, clippy::as_conversions))] + +mod arena; +mod inflight; +mod queue; + +use core::sync::atomic::{AtomicU32, Ordering}; + +pub use arena::{Geometry, Run, Span}; +pub use inflight::Inflight; +pub use queue::{Consumer, Place, Producer}; +pub use toyos_untrusted::Untrusted; + +/// One shared 32-bit word of a region: an atomic over the mapping, or a +/// model's. +pub trait Word { + fn load(&self, order: Ordering) -> u32; + fn store(&self, value: u32, order: Ordering); +} + +impl Word for AtomicU32 { + fn load(&self, order: Ordering) -> u32 { + AtomicU32::load(self, order) + } + fn store(&self, value: u32, order: Ordering) { + AtomicU32::store(self, value, order) + } +} + +/// The peer wrote what no peer of the protocol writes. The session is over; +/// the variant is its name. +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub enum Violation { + /// A producer's tail more than the ring holds past what was released. + TailPastDepth, + /// A consumer's head past what was published, or more than the ring holds + /// behind it. + HeadPastTail, + /// A reply no server of the protocol writes. + Entry, + /// A run outside the arena. + Run, + /// A tag nothing is in flight under. + Tag, +} diff --git a/toyos-transport/src/queue.rs b/toyos-transport/src/queue.rs new file mode 100644 index 0000000000..681b101b24 --- /dev/null +++ b/toyos-transport/src/queue.rs @@ -0,0 +1,309 @@ +//! A single-producer queue of `E`-word entries, `D` deep, among `N` words. +//! +//! **A producer writes an entry's words, then publishes the tail with +//! `Release`; a consumer loads the tail with `Acquire`, then reads the words.** +//! The head goes back the other way. Each end looks at the peer's cursor when +//! what it last saw is spent, not once per entry, and stores its own once per +//! batch. + +use core::sync::atomic::Ordering; + +use crate::{Untrusted, Violation, Word}; + +#[cfg(not(feature = "publish-relaxed"))] +const PUBLISH: Ordering = Ordering::Release; +#[cfg(feature = "publish-relaxed")] +const PUBLISH: Ordering = Ordering::Relaxed; + +/// Word `at` of an end's words, and the one index into them: every `at` an end +/// forms is below `N`, because its [`Place`] is and a ring position is masked +/// below the depth. +#[allow(clippy::indexing_slicing)] +fn word(page: &[W; N], at: usize) -> &W { + &page[at] +} + +/// A distance the peer's cursor claims, believed only up to `cap`. +fn clamp(claimed: Untrusted, cap: u32, broken: Violation) -> Result { + let cap = if cfg!(feature = "no-clamp") { u32::MAX } else { cap }; + claimed.at_most(u64::from(cap)).ok().and_then(|n| u32::try_from(n).ok()).ok_or(broken) +} + +/// Where a queue of `E`-word entries, `D` deep, is among `N` words: the word +/// its consumer stores its head in, the word its producer stores its tail in, +/// and its first entry. +/// +/// ``` +/// use toyos_transport::Place; +/// const QUEUE: Place<2, 4, 10> = Place::new::<0, 1, 2>(); +/// ``` +/// +/// A place with a word at `N` or past it does not compile: +/// +/// ```compile_fail,E0080 +/// use toyos_transport::Place; +/// const QUEUE: Place<2, 4, 10> = Place::new::<0, 1, 3>(); +/// ``` +/// +/// ```compile_fail,E0080 +/// use toyos_transport::Place; +/// const QUEUE: Place<2, 4, 10> = Place::new::<10, 1, 2>(); +/// ``` +/// +/// ```compile_fail,E0080 +/// use toyos_transport::Place; +/// const QUEUE: Place<2, 4, 10> = Place::new::<0, 10, 2>(); +/// ``` +#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] +pub struct Place { + head: usize, + tail: usize, + entries: usize, +} + +impl Place { + pub const fn new() -> Self { + const { + assert!(D.is_power_of_two() && E > 0, "a queue is a power of two deep, of entries of a word or more"); + assert!(usize::BITS >= u32::BITS, "a ring position is a u32"); + #[allow(clippy::as_conversions)] + let entries_end = ENTRIES + D as usize * E; + assert!(HEAD < N && TAIL < N && entries_end <= N, "a place names a word past its words"); + } + Self { head: HEAD, tail: TAIL, entries: ENTRIES } + } + + /// The words of the entry at ring position `at`. + fn entry<'a, W>(&self, page: &'a [W; N], at: u32) -> impl Iterator { + // `usize` holds every `u32`, so the slot is exact. + #[allow(clippy::as_conversions)] + let slot = (at & D.wrapping_sub(1)) as usize; + let first = self.entries.wrapping_add(slot.wrapping_mul(E)); + (0..E).map(move |k| word(page, first.wrapping_add(k))) + } + + /// How far past `released` the producer's tail is: at most `D`. + fn published(&self, page: &[W; N], released: u32) -> Result { + let tail = word(page, self.tail).load(Ordering::Acquire); + clamp(Untrusted::new(tail).map(|t| t.wrapping_sub(released)), D, Violation::TailPastDepth) + } + + /// How far behind `published` the consumer's head is: at most `D`. + fn unreleased(&self, page: &[W; N], published: u32) -> Result { + let head = word(page, self.head).load(Ordering::Acquire); + clamp(Untrusted::new(head).map(|h| published.wrapping_sub(h)), D, Violation::HeadPastTail) + } +} + +/// The end of a queue that writes entries. It holds its cursors and not the +/// region; every call is given the region's words. +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct Producer { + place: Place, + local: u32, + published: u32, + /// Entries that may be pushed before the head is looked at again. + room: u32, +} + +impl Producer { + /// This end of the queue at `place`, its tail stored 0. + pub fn new(page: &[W; N], place: Place) -> Self { + word(page, place.tail).store(0, Ordering::Release); + Self { place, local: 0, published: 0, room: D } + } + + /// How many entries may be pushed before the consumer frees more. + pub fn space(&mut self, page: &[W; N]) -> Result { + let unreleased = self.place.unreleased(page, self.published)?; + let pending = self.local.wrapping_sub(self.published); + self.room = D.saturating_sub(pending.saturating_add(unreleased)); + Ok(self.room) + } + + /// Write one entry, or answer `false` for a full queue and write nothing. + /// It is the consumer's once [`Self::publish`] runs. + pub fn push(&mut self, page: &[W; N], words: [u32; E]) -> Result { + if self.room == 0 && self.space(page)? == 0 { + return Ok(false); + } + for (shared, value) in self.place.entry(page, self.local).zip(words) { + shared.store(value, Ordering::Relaxed); + } + self.local = self.local.wrapping_add(1); + self.room = self.room.wrapping_sub(1); + Ok(true) + } + + /// Publish every entry pushed so far; `false` if there was none. + pub fn publish(&mut self, page: &[W; N]) -> bool { + if self.published == self.local { + return false; + } + word(page, self.place.tail).store(self.local, PUBLISH); + self.published = self.local; + true + } +} + +/// The end of a queue that reads entries; like [`Producer`], it holds cursors +/// and is given the region's words. +#[derive(Clone, Debug, PartialEq, Eq, Hash)] +pub struct Consumer { + place: Place, + local: u32, + released: u32, + /// Entries published and not yet popped, as last seen. + ready: u32, +} + +impl Consumer { + /// This end of the queue at `place`, its head stored 0. + pub fn new(page: &[W; N], place: Place) -> Self { + word(page, place.head).store(0, Ordering::Release); + Self { place, local: 0, released: 0, ready: 0 } + } + + /// The next published entry, or `None` for none. + pub fn pop(&mut self, page: &[W; N]) -> Result; E]>, Violation> { + if self.ready == 0 { + let published = self.place.published(page, self.released)?; + self.ready = published.saturating_sub(self.local.wrapping_sub(self.released)); + if self.ready == 0 { + return Ok(None); + } + } + let mut words = [Untrusted::new(0); E]; + for (out, shared) in words.iter_mut().zip(self.place.entry(page, self.local)) { + *out = Untrusted::new(shared.load(Ordering::Relaxed)); + } + self.local = self.local.wrapping_add(1); + self.ready = self.ready.wrapping_sub(1); + Ok(Some(words)) + } + + /// Give every entry popped so far back to the producer. + pub fn release(&mut self, page: &[W; N]) { + if self.released != self.local { + word(page, self.place.head).store(self.local, Ordering::Release); + self.released = self.local; + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use core::sync::atomic::AtomicU32; + + const D: u32 = 8; + const WORDS: usize = 32 + 2 * D as usize; + const PLACE: Place<2, D, WORDS> = Place::new::<0, 16, 32>(); + + fn page() -> [AtomicU32; WORDS] { + core::array::from_fn(|_| AtomicU32::new(0)) + } + + fn ends(page: &[AtomicU32; WORDS]) -> (Producer<2, D, WORDS>, Consumer<2, D, WORDS>) { + (Producer::new(page, PLACE), Consumer::new(page, PLACE)) + } + + fn plain(words: Option<[Untrusted; 2]>) -> Option<[u32; 2]> { + words.map(|w| w.map(|w| w.at_most(u32::MAX.into()).unwrap() as u32)) + } + + #[test] + fn an_entry_is_nobodys_before_it_is_published() { + let page = page(); + let (mut tx, mut rx) = ends(&page); + assert_eq!(tx.push(&page, [3, 4]), Ok(true)); + assert_eq!(rx.pop(&page).map(plain), Ok(None)); + assert!(tx.publish(&page)); + assert!(!tx.publish(&page), "nothing new, nothing published"); + assert_eq!(rx.pop(&page).map(plain), Ok(Some([3, 4]))); + assert_eq!(rx.pop(&page).map(plain), Ok(None)); + } + + /// The ring wraps many times over, and space is exactly what the consumer + /// has given back. + #[test] + fn the_ring_wraps_and_counts_its_space() { + let page = page(); + let (mut tx, mut rx) = ends(&page); + let (mut pushed, mut popped) = (0u32, 0u32); + for round in 0..5 * D { + for _ in 0..1 + round % 9 { + assert_eq!(tx.space(&page), Ok(D - (pushed - popped))); + if !tx.push(&page, [pushed, !pushed]).unwrap() { + assert_eq!(pushed - popped, D); + break; + } + pushed += 1; + } + tx.publish(&page); + for _ in 0..1 + round % 7 { + let Some(words) = plain(rx.pop(&page).unwrap()) else { break }; + assert_eq!(words, [popped, !popped], "entries come out whole, in the order they went in"); + popped += 1; + } + rx.release(&page); + } + assert!(pushed > 2 * D, "the ring wrapped"); + } + + /// A tail more than the ring past what was released, and a head past what + /// was published, each end the session by name; a cursor moved backwards + /// within bounds costs only its owner. + #[test] + fn a_peer_cursor_out_of_reach_is_a_violation() { + let page = page(); + let (mut tx, mut rx) = ends(&page); + page[16].store(D + 1, Ordering::Release); + assert_eq!(rx.pop(&page), Err(Violation::TailPastDepth)); + page[16].store(u32::MAX, Ordering::Release); + assert_eq!(rx.pop(&page), Err(Violation::TailPastDepth), "a tail behind what was released"); + page[0].store(1, Ordering::Release); + assert_eq!(tx.space(&page), Err(Violation::HeadPastTail)); + page[0].store(u32::MAX, Ordering::Release); + assert_eq!(tx.space(&page), Ok(D - 1), "a head moved back costs its consumer the room"); + } + + /// A producer that steps its tail one past each entry popped, none of them + /// released: the ring's depth is taken, and the next tail is past it. + #[test] + fn a_tail_stepped_past_each_pop_is_refused_at_the_depth() { + let page = page(); + let (_, mut rx) = ends(&page); + let mut taken = 0; + let refused = loop { + page[16].store(taken + 1, Ordering::Release); + match rx.pop(&page) { + Ok(Some(_)) => taken += 1, + Ok(None) => panic!("a published entry was not taken"), + Err(violation) => break violation, + } + assert!(taken <= D, "took {taken} entries from a ring of {D} without releasing one"); + }; + assert_eq!((taken, refused), (D, Violation::TailPastDepth)); + } + + /// A consumer that steps its head onto each entry pushed, none of them + /// published: the ring's depth is pushed, and the next head is past what + /// was published. + #[test] + fn a_head_stepped_past_each_push_is_refused_at_the_depth() { + let page = page(); + let (mut tx, _) = ends(&page); + let mut pushed = 0; + let refused = loop { + match tx.push(&page, [pushed, pushed]) { + Ok(true) => pushed += 1, + Ok(false) => panic!("a ring with nothing published was full"), + Err(violation) => break violation, + } + page[0].store(pushed, Ordering::Release); + assert!(pushed <= D, "pushed {pushed} into a ring of {D} with none published"); + }; + assert_eq!((pushed, refused), (D, Violation::HeadPastTail)); + } +} diff --git a/toyos-transport/tests/loom.rs b/toyos-transport/tests/loom.rs new file mode 100644 index 0000000000..8051dc92a5 --- /dev/null +++ b/toyos-transport/tests/loom.rs @@ -0,0 +1,141 @@ +//! The transport's two ends on two CPUs, under loom: what no host test that +//! runs both ends on one thread can reach. +//! +//! - **Publication**: a consumer that sees a tail sees every word of the +//! entries below it. `publish-relaxed` takes the edge away. +//! - **A hostile peer** leaves the honest end with entries, nothing, or a +//! named violation, and never more entries than the ring holds. `no-clamp` +//! believes the peer. +//! +//! cargo test -p toyos-transport --features --test loom + +use core::sync::atomic::Ordering; + +use loom::sync::atomic::AtomicU32; +use loom::sync::Arc; +use toyos_transport::{Consumer, Place, Producer, Untrusted, Violation, Word}; + +/// A loom atomic as a region word: the trait is this crate's and the type is +/// loom's, so the two meet through a wrapper. +struct Shared(AtomicU32); + +impl Word for Shared { + fn load(&self, order: Ordering) -> u32 { + self.0.load(order) + } + fn store(&self, value: u32, order: Ordering) { + self.0.store(value, order) + } +} + +const E: usize = 2; +const D: u32 = 2; +const WORDS: usize = 2 + E * D as usize; +const PLACE: Place = Place::new::<0, 1, 2>(); + +fn page() -> Arc<[Shared; WORDS]> { + Arc::new(core::array::from_fn(|_| Shared(AtomicU32::new(0)))) +} + +fn ends(page: &[Shared; WORDS]) -> (Producer, Consumer) { + (Producer::new(page, PLACE), Consumer::new(page, PLACE)) +} + +fn entry(n: u32) -> [u32; E] { + [100 + n, 7 * n + 3] +} + +fn plain(words: [Untrusted; E]) -> [u32; E] { + words.map(|w| w.at_most(u32::MAX.into()).unwrap() as u32) +} + +/// Two entries published one at a time, read by the other end as they +/// arrive: each is whole, in order, and exactly what was written. +#[test] +fn a_published_entry_is_read_whole() { + loom::model(|| { + let page = page(); + let (mut tx, mut rx) = ends(&page); + let consumer_page = Arc::clone(&page); + let consumer = loom::thread::spawn(move || { + let mut read = Vec::new(); + while read.len() < 2 { + match rx.pop(&consumer_page).expect("the producer keeps the protocol") { + Some(words) => read.push(plain(words)), + None => loom::thread::yield_now(), + } + } + rx.release(&consumer_page); + read + }); + for n in 0..2 { + assert!(tx.push(&page, entry(n)).unwrap()); + tx.publish(&page); + } + let read = consumer.join().expect("the consumer thread"); + assert_eq!(read, [entry(0), entry(1)], "an entry was read before its words"); + }); +} + +/// A producer that stores a tail past the ring, one behind what was released, +/// and garbage into the entries and into the consumer's own head: every pop is +/// an entry, nothing, or [`Violation::TailPastDepth`], and no more than the +/// ring's depth of entries is taken without a release. +#[test] +fn a_hostile_producer_yields_entries_or_a_violation() { + loom::model(|| { + let page = page(); + let (_, mut rx) = ends(&page); + let hostile_page = Arc::clone(&page); + let hostile = loom::thread::spawn(move || { + for (at, value) in [(1, D + 1), (3, 9), (0, 5), (1, u32::MAX), (2, 7)] { + hostile_page[at].store(value, Ordering::Release); + } + }); + let mut taken = 0; + for _ in 0..D + 2 { + match rx.pop(&page) { + Ok(Some(_)) => taken += 1, + Ok(None) => {} + Err(violation) => { + assert_eq!(violation, Violation::TailPastDepth); + break; + } + } + } + hostile.join().unwrap(); + assert!(taken <= D, "took {taken} entries from a ring of {D} without releasing one"); + }); +} + +/// A consumer that stores a head past what was published, one wrapped far +/// behind, and garbage into the producer's own tail and the entries: every +/// push is room, a full ring, or [`Violation::HeadPastTail`], and no more +/// than the ring's depth is ever pushed ahead of a head that was never moved. +#[test] +fn a_hostile_consumer_yields_room_or_a_violation() { + loom::model(|| { + let page = page(); + let (mut tx, _) = ends(&page); + let hostile_page = Arc::clone(&page); + let hostile = loom::thread::spawn(move || { + for (at, value) in [(0, 5), (1, 9), (0, u32::MAX), (2, 1), (4, 6)] { + hostile_page[at].store(value, Ordering::Release); + } + }); + let mut pushed = 0; + for n in 0..D + 2 { + match tx.push(&page, entry(n)) { + Ok(true) => pushed += 1, + Ok(false) => {} + Err(violation) => { + assert_eq!(violation, Violation::HeadPastTail); + break; + } + } + tx.publish(&page); + } + hostile.join().unwrap(); + assert!(pushed <= D, "pushed {pushed} into a ring of {D} nobody released"); + }); +} diff --git a/userland/Cargo.lock b/userland/Cargo.lock index a1147e1c2b..30245ea0ba 100644 --- a/userland/Cargo.lock +++ b/userland/Cargo.lock @@ -4058,6 +4058,7 @@ name = "toyos-blockring" version = "0.1.0" dependencies = [ "toyos-blockhold", + "toyos-transport", ] [[package]] @@ -4150,6 +4151,17 @@ version = "0.1.0" name = "toyos-tmpdir" version = "0.1.0" +[[package]] +name = "toyos-transport" +version = "0.1.0" +dependencies = [ + "toyos-untrusted", +] + +[[package]] +name = "toyos-untrusted" +version = "0.1.0" + [[package]] name = "toyos-update" version = "0.1.0" diff --git a/userland/blockd/src/main.rs b/userland/blockd/src/main.rs index 1a5d0802d9..ad65b88890 100644 --- a/userland/blockd/src/main.rs +++ b/userland/blockd/src/main.rs @@ -47,8 +47,7 @@ use toyos_abi::part::{PartGuid, GUID_TEXT_LEN}; use toyos_abi::syscall::{DEV_PREFIX, SyscallError}; use toyos_blockhold::Holds; use toyos_blockring::entry::{Completion, Op}; -use toyos_blockring::layout::{arena_byte, DEPTH}; -use toyos_blockring::ring::{self, ServerRings}; +use toyos_blockring::layout::{self, ServerRings, DEPTH}; use toyos_blockring::server::{ServerSession, Taken}; use toyos_blockring::wire::{self, Opened, Refusal}; use toyos_blockring::{BLOCK_BYTES, PORT, SESSION_BYTES}; @@ -246,7 +245,7 @@ impl Service { fn admit(&mut self, opening: Opening, conn: Connection) -> u64 { let id = self.next_id; self.next_id += 1; - let rings = ring::server(opening.region.words()); + let rings = layout::server(opening.region.words()); self.sessions.insert( id, Served { @@ -273,11 +272,14 @@ impl Service { self.holds.release(opening.first); } - /// Put `c` on its session's completion ring. + /// Put `c` on its session's completion ring. [`Self::pull`] leaves room for + /// every answer, so a ring without it is a client whose head went back or + /// past what was posted: the session ends. fn post(session: &mut Served, c: Completion) { - let page = session.region.words(); - session.rings.1.push(page, c.encode()); - session.posted = true; + match session.rings.1.push(session.region.words(), c.encode()) { + Ok(true) => session.posted = true, + Ok(false) | Err(_) => session.closing = true, + } } /// Hand one device answer to the session it is for. @@ -304,11 +306,11 @@ impl Service { /// session's completion ring have room for it. fn pull(&mut self, id: u64) { let s = self.sessions.get_mut(&id).expect("a live session"); - let page = s.region.words(); loop { if s.closing { break; } + let page = s.region.words(); // Room for every answer: what is on the device, and what is posted // and not yet read, never exceeds the completion ring. let Ok(space) = s.rings.1.space(page) else { @@ -316,7 +318,7 @@ impl Service { break; }; let unread = (DEPTH - space) as usize; - if s.state.inflight() + unread >= DEPTH as usize || !self.ctrl.has_room() { + if s.state.inflight().len() + unread >= DEPTH as usize || !self.ctrl.has_room() { break; } let words = match s.rings.0.pop(page) { @@ -329,31 +331,28 @@ impl Service { }; s.requests += 1; match s.state.take(words) { - Taken::Answer(c) => { - s.rings.1.push(page, c.encode()); - s.posted = true; - } + Taken::Answer(c) => Self::post(s, c), Taken::Issue(req) => { - let owner = Owner::Session { session: id, tag: req.tag, write: req.op == Op::Write }; + let write = matches!(req.op, Op::Write { .. }); + let owner = Owner::Session { session: id, tag: req.tag, write }; match req.op { - Op::Read | Op::Write => { - let at = s.device_addr + arena_byte(req.arena) as u64; - let block = s.state.first() + req.lba; - self.ctrl.submit_io(req.op == Op::Write, block, req.blocks, at, owner); + Op::Read { run, lba } | Op::Write { run, lba } => { + let at = s.device_addr + run.span().offset as u64; + let block = s.state.first() + lba; + self.ctrl.submit_io(write, block, run.count(), at, owner); } Op::Flush if self.ctrl.vwc => self.ctrl.submit_flush(owner), // No volatile cache: every write answered is on the // medium already. Op::Flush => { let c = s.state.complete(req.tag, true, &mut self.holds, self.losses); - s.rings.1.push(page, c.expect("just taken").encode()); - s.posted = true; + Self::post(s, c.expect("just taken")); } } } } } - s.rings.0.release(page); + s.rings.0.release(s.region.words()); } /// Publish what each session was answered, and ring its doorbell. @@ -378,7 +377,7 @@ impl Service { let done: Vec = self .sessions .iter() - .filter(|(_, s)| s.closing && s.state.inflight() == 0) + .filter(|(_, s)| s.closing && s.state.inflight().len() == 0) .map(|(id, _)| *id) .collect(); for id in done { diff --git a/userland/blockd/src/region.rs b/userland/blockd/src/region.rs index c4fba33dd0..d268aec9bc 100644 --- a/userland/blockd/src/region.rs +++ b/userland/blockd/src/region.rs @@ -9,8 +9,8 @@ use core::sync::atomic::AtomicU32; use toyos::shm::SharedMemory; use toyos_abi::syscall::SyscallError; -use toyos_blockring::layout::{arena_byte, ARENA_BLOCKS, RING_WORDS, SESSION_BYTES}; -use toyos_blockring::BLOCK_BYTES; +use toyos_blockring::layout::{RING_WORDS, SESSION_BYTES}; +use toyos_blockring::Run; use toyos::volatile::Window; @@ -31,7 +31,7 @@ impl Region { } /// The ring words, as the rings take them. - pub fn words(&self) -> &[AtomicU32] { + pub fn words(&self) -> &[AtomicU32; RING_WORDS] { let base = self.memory.as_ptr(); assert!(base as usize % align_of::() == 0); // SAFETY: the mapping is `SESSION_BYTES` long and lives as long as @@ -39,19 +39,15 @@ impl Region { // `RING_WORDS` words of it are inside it (`layout`'s own assertion); // the base is 2 MiB aligned; and an atomic is the one type that may // alias memory another process writes. - unsafe { core::slice::from_raw_parts(base as *const AtomicU32, RING_WORDS) } + unsafe { &*(base as *const [AtomicU32; RING_WORDS]) } } - /// Arena blocks `first..first + blocks`. - pub fn arena(&self, first: u32, blocks: u32) -> Window { - assert!( - first.checked_add(blocks).is_some_and(|end| end <= ARENA_BLOCKS), - "blockd: arena blocks {first}+{blocks} past the arena" - ); + /// The arena blocks of `run`. + pub fn arena(&self, run: &Run) -> Window { // SAFETY: `SESSION_BYTES` mapped for as long as `self.memory`, and // every window over it is used while the region is held. let whole = unsafe { Window::new(self.memory.as_ptr(), SESSION_BYTES) }; - whole.sub(arena_byte(first), blocks as usize * BLOCK_BYTES) + whole.sub(run.span().offset, run.span().len) } /// A second handle to the region, for a send. diff --git a/userland/blockd/src/session.rs b/userland/blockd/src/session.rs index 2a9aea1d4f..a5b83eed6a 100644 --- a/userland/blockd/src/session.rs +++ b/userland/blockd/src/session.rs @@ -23,10 +23,9 @@ use toyos::poller::{Poller, READABLE}; use toyos_abi::syscall::SyscallError; use toyos_blockring::client::{Client, Outcome, Ticket}; use toyos_blockring::entry::{Completion, Op}; -use toyos_blockring::layout::{ARENA_BLOCKS, MAX_REQUEST_BLOCKS}; -use toyos_blockring::ring::{self, ClientRings}; +use toyos_blockring::layout::{self, ClientRings, ARENA, MAX_REQUEST_BLOCKS}; use toyos_blockring::wire::{self, Opened, Refusal}; -use toyos_blockring::BLOCK_BYTES; +use toyos_blockring::{Run, BLOCK_BYTES}; use crate::region::Region; @@ -79,10 +78,10 @@ struct Arena { impl Arena { fn new() -> Self { - Self { free: vec![true; ARENA_BLOCKS as usize] } + Self { free: vec![true; ARENA.slots() as usize] } } - fn alloc(&mut self, blocks: u32) -> Option { + fn alloc(&mut self, blocks: u32) -> Option { let n = blocks as usize; let mut run = 0; for i in 0..self.free.len() { @@ -90,22 +89,23 @@ impl Arena { if run == n { let first = i + 1 - n; self.free[first..=i].iter_mut().for_each(|b| *b = false); - return Some(first as u32); + return ARENA.run(first as u32, blocks); } } None } - fn release(&mut self, first: u32, blocks: u32) { - for b in &mut self.free[first as usize..(first + blocks) as usize] { - assert!(!*b, "blockd: arena block {first}+{blocks} released twice"); + fn release(&mut self, run: Run) { + let first = run.first() as usize; + for b in &mut self.free[first..first + run.count() as usize] { + assert!(!*b, "blockd: arena run {run:?} released twice"); *b = true; } } } enum Pending { - Read { arena: u32, blocks: u32 }, + Read { run: Run }, Write, Flush, } @@ -135,7 +135,7 @@ impl Session { /// `names` calls `service`. pub fn open(names: Namespace, service: &str, guid: [u8; wire::GUID_BYTES]) -> Result { let region = Region::create().map_err(Error::Kernel)?; - let rings = ring::client(region.words()); + let rings = layout::client(region.words()); let (conn, opened) = handshake(&names, service, guid, ®ion)?; let mut client = Client::new(); client.session_started(); @@ -184,10 +184,10 @@ impl Session { if self.conn.is_none() { return Err(Unsent::Ended); } - let arena = self.arena.alloc(blocks).ok_or(Unsent::ArenaFull)?; - self.region.arena(arena, blocks).copy_in(0, data); + let run = self.arena.alloc(blocks).ok_or(Unsent::ArenaFull)?; + self.region.arena(&run).copy_in(0, data); let ticket = self.ticket(); - self.client.submit(ticket, Op::Write, lba, blocks, arena); + self.client.submit(ticket, Op::Write { run, lba }); self.pending.insert(ticket, Pending::Write); Ok(ticket) } @@ -198,10 +198,10 @@ impl Session { if self.conn.is_none() { return Err(Unsent::Ended); } - let arena = self.arena.alloc(blocks).ok_or(Unsent::ArenaFull)?; + let run = self.arena.alloc(blocks).ok_or(Unsent::ArenaFull)?; let ticket = self.ticket(); - self.client.submit(ticket, Op::Read, lba, blocks, arena); - self.pending.insert(ticket, Pending::Read { arena, blocks }); + self.client.submit(ticket, Op::Read { run, lba }); + self.pending.insert(ticket, Pending::Read { run }); Ok(ticket) } @@ -211,7 +211,7 @@ impl Session { return Err(Unsent::Ended); } let ticket = self.ticket(); - self.client.submit(ticket, Op::Flush, 0, 0, 0); + self.client.submit(ticket, Op::Flush); self.pending.insert(ticket, Pending::Flush); Ok(ticket) } @@ -231,9 +231,10 @@ impl Session { Err(_) => break, } let Some(request) = self.client.next_request() else { break }; - self.rings.0.push(page, request.encode()); + let pushed = self.rings.0.push(page, request.encode()); + assert_eq!(pushed, Ok(true), "blockd: a request past the room just counted"); } - self.peak = self.peak.max(self.client.on_the_wire()); + self.peak = self.peak.max(self.client.on_the_wire().count()); if self.rings.0.publish(page) { // A full pipe is a doorbell already rung; a gone one is a server // that has ended, which the wait finds. @@ -324,9 +325,9 @@ impl Session { for (ticket, outcome) in decided { let pending = self.pending.remove(&ticket).expect("blockd: an answer for no ticket"); let data = match (pending, outcome) { - (Pending::Read { arena, blocks }, Outcome::Done) => { - let mut data = vec![0u8; blocks as usize * BLOCK_BYTES]; - self.region.arena(arena, blocks).copy_out(0, &mut data); + (Pending::Read { run }, Outcome::Done) => { + let mut data = vec![0u8; run.span().len]; + self.region.arena(&run).copy_out(0, &mut data); Some(data) } _ => None, @@ -334,8 +335,8 @@ impl Session { answers.push(Answer { ticket, outcome, data }); } let released: Vec<_> = self.client.take_released().collect(); - for (first, blocks) in released { - self.arena.release(first, blocks); + for run in released { + self.arena.release(run); } } @@ -355,7 +356,7 @@ impl Session { assert!(self.conn.is_none(), "blockd: reconnect while a session is open"); // Before the region goes to the new server: it must find this end's // two indices at zero, as it will set its own. - self.rings = ring::client(self.region.words()); + self.rings = layout::client(self.region.words()); let (conn, opened) = handshake(&self.names, &self.service, self.guid, &self.region)?; if opened != self.opened { return Err(Error::Protocol);