style: cargo fmt sweep under rustc 1.97.1 (toolchain reformat, no semantic change)

This commit is contained in:
smarm-agent
2026-08-13 05:56:49 +00:00
parent 1262cc30e3
commit 95306c7f60
71 changed files with 1956 additions and 763 deletions
+1 -1
View File
@@ -99,7 +99,7 @@ pub extern "C-unwind" fn trampoline() {
};
let outcome = match panic::catch_unwind(panic::AssertUnwindSafe(b)) {
Ok(()) => Outcome::Exit,
Ok(()) => Outcome::Exit,
Err(payload) => {
if payload.is::<StopSentinel>() {
Outcome::Stopped
+20 -10
View File
@@ -342,8 +342,7 @@ mod inner {
// Count the loss in would-be delta terms so the audit's columns
// compare directly against `injected_cycles`.
DISCARD_OVERMAX_N.fetch_add(1, Ordering::Relaxed);
DISCARD_OVERMAX_CYCLES
.fetch_add(interval.saturating_mul(pct) / 100, Ordering::Relaxed);
DISCARD_OVERMAX_CYCLES.fetch_add(interval.saturating_mul(pct) / 100, Ordering::Relaxed);
return;
}
let delta = interval.saturating_mul(pct) / 100;
@@ -419,8 +418,7 @@ mod inner {
let gap = preempt::rdtsc()
.saturating_sub(desched_tsc)
.min(MAX_SAMPLE_CYCLES);
OFFCPU_IN_SITE_CYCLES
.fetch_add(gap.saturating_mul(pct) / 100, Ordering::Relaxed);
OFFCPU_IN_SITE_CYCLES.fetch_add(gap.saturating_mul(pct) / 100, Ordering::Relaxed);
OFFCPU_IN_SITE_N.fetch_add(1, Ordering::Relaxed);
}
}
@@ -533,7 +531,9 @@ mod inner {
park_forgiven_cycles: self
.park_forgiven_cycles
.saturating_sub(before.park_forgiven_cycles),
drop_park_cycles: self.drop_park_cycles.saturating_sub(before.drop_park_cycles),
drop_park_cycles: self
.drop_park_cycles
.saturating_sub(before.drop_park_cycles),
drop_park_n: self.drop_park_n.saturating_sub(before.drop_park_n),
drop_yield_cycles: self
.drop_yield_cycles
@@ -542,12 +542,18 @@ mod inner {
discard_overmax_cycles: self
.discard_overmax_cycles
.saturating_sub(before.discard_overmax_cycles),
discard_overmax_n: self.discard_overmax_n.saturating_sub(before.discard_overmax_n),
discard_unarmed_n: self.discard_unarmed_n.saturating_sub(before.discard_unarmed_n),
discard_overmax_n: self
.discard_overmax_n
.saturating_sub(before.discard_overmax_n),
discard_unarmed_n: self
.discard_unarmed_n
.saturating_sub(before.discard_unarmed_n),
offcpu_in_site_cycles: self
.offcpu_in_site_cycles
.saturating_sub(before.offcpu_in_site_cycles),
offcpu_in_site_n: self.offcpu_in_site_n.saturating_sub(before.offcpu_in_site_n),
offcpu_in_site_n: self
.offcpu_in_site_n
.saturating_sub(before.offcpu_in_site_n),
}
}
}
@@ -795,7 +801,9 @@ mod inner {
let cell = results
.iter()
.find(|r| r.site == site && r.speedup_pct == speedup_pct)?;
let base = results.iter().find(|r| r.site == site && r.speedup_pct == 0)?;
let base = results
.iter()
.find(|r| r.site == site && r.speedup_pct == 0)?;
let rate = normalized_rate(cell, point)?;
let b = normalized_rate(base, point)?;
if b <= 0.0 {
@@ -946,7 +954,9 @@ macro_rules! progress {
macro_rules! causal_site {
($name:literal) => {{
static __SMARM_SITE: ::std::sync::OnceLock<u32> = ::std::sync::OnceLock::new();
$crate::causal::SiteGuard::enter(*__SMARM_SITE.get_or_init(|| $crate::causal::site_id($name)))
$crate::causal::SiteGuard::enter(
*__SMARM_SITE.get_or_init(|| $crate::causal::site_id($name)),
)
}};
}
+51 -27
View File
@@ -104,7 +104,12 @@ pub fn channel<T>() -> (Sender<T>, Receiver<T>) {
senders: 1,
receiver_alive: true,
}));
(Sender { inner: inner.clone() }, Receiver { inner })
(
Sender {
inner: inner.clone(),
},
Receiver { inner },
)
}
struct Inner<T> {
@@ -178,7 +183,9 @@ impl std::error::Error for RecvTimeoutError {}
impl<T> Clone for Sender<T> {
fn clone(&self) -> Self {
self.inner.lock().senders += 1;
Sender { inner: self.inner.clone() }
Sender {
inner: self.inner.clone(),
}
}
}
@@ -248,10 +255,18 @@ impl<T> Sender<T> {
g.parked_receiver.take()
};
if let Some((pid, epoch)) = unpark {
crate::te!(crate::trace::Event::Send { sender: crate::actor::current_pid().unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)), receiver: Some(pid) });
crate::te!(crate::trace::Event::Send {
sender: crate::actor::current_pid()
.unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)),
receiver: Some(pid)
});
crate::scheduler::unpark_at(pid, epoch);
} else {
crate::te!(crate::trace::Event::Send { sender: crate::actor::current_pid().unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)), receiver: None });
crate::te!(crate::trace::Event::Send {
sender: crate::actor::current_pid()
.unwrap_or(crate::pid::Pid::new(u32::MAX, u32::MAX)),
receiver: None
});
}
Ok(())
}
@@ -290,10 +305,12 @@ impl<T> Receiver<T> {
// Release the lock before parking: the unparker will need it.
crate::scheduler::park_current();
// Woken up. Record it before looping to check the queue.
crate::te!(crate::trace::Event::RecvWake(match crate::actor::current_pid() {
Some(p) => p,
None => panic!("smarm: RecvWake outside an actor (core corrupt)"),
}));
crate::te!(crate::trace::Event::RecvWake(
match crate::actor::current_pid() {
Some(p) => p,
None => panic!("smarm: RecvWake outside an actor (core corrupt)"),
}
));
}
}
@@ -347,10 +364,12 @@ impl<T> Receiver<T> {
crate::scheduler::insert_wait_timer(deadline, me, target, epoch);
crate::scheduler::park_current();
crate::te!(crate::trace::Event::RecvWake(match crate::actor::current_pid() {
Some(p) => p,
None => panic!("smarm: RecvWake outside an actor (core corrupt)"),
}));
crate::te!(crate::trace::Event::RecvWake(
match crate::actor::current_pid() {
Some(p) => p,
None => panic!("smarm: RecvWake outside an actor (core corrupt)"),
}
));
let mut g = self.inner.lock();
if let Some(v) = g.queue.pop_front() {
crate::preempt::note_message_received();
@@ -412,10 +431,12 @@ impl<T> Receiver<T> {
}
// Release the lock before parking: the unparker will need it.
crate::scheduler::park_current();
crate::te!(crate::trace::Event::RecvWake(match crate::actor::current_pid() {
Some(p) => p,
None => panic!("smarm: RecvWake outside an actor (core corrupt)"),
}));
crate::te!(crate::trace::Event::RecvWake(
match crate::actor::current_pid() {
Some(p) => p,
None => panic!("smarm: RecvWake outside an actor (core corrupt)"),
}
));
}
}
@@ -616,7 +637,12 @@ pub fn try_select(arms: &[&dyn Selectable]) -> std::io::Result<usize> {
// Channel-only selects skip all of it: `eager` is false, the guard
// is disarmed, and the loser-arm self-cleaning story is unchanged.
let eager = arms.iter().any(|a| a.sel_eager_cleanup());
let mut guard = UnregisterGuard { arms, me, epoch, armed: eager };
let mut guard = UnregisterGuard {
arms,
me,
epoch,
armed: eager,
};
crate::scheduler::park_current();
@@ -687,11 +713,7 @@ impl Drop for UnregisterGuard<'_> {
// unregistered eagerly so none are left dangling. `Err` = an arm failed to
// register; same unwind (earlier fd arms unregistered, wait retired).
// `Ok(None)` = every arm registered successfully; the caller parks.
fn register_arms(
me: Pid,
epoch: u32,
arms: &[&dyn Selectable],
) -> std::io::Result<Option<usize>> {
fn register_arms(me: Pid, epoch: u32, arms: &[&dyn Selectable]) -> std::io::Result<Option<usize>> {
for (i, arm) in arms.iter().enumerate() {
let registered = match arm.sel_register(me, epoch) {
Ok(r) => r,
@@ -736,10 +758,7 @@ impl crate::timer::TimerTarget for SelectTimeout {
/// Panics if `arms` is empty, if called outside an actor, or if an fd arm
/// fails to register (see [`try_select_timeout`] for the fallible form; a
/// channel-only select can never fail).
pub fn select_timeout(
arms: &[&dyn Selectable],
timeout: std::time::Duration,
) -> Option<usize> {
pub fn select_timeout(arms: &[&dyn Selectable], timeout: std::time::Duration) -> Option<usize> {
match try_select_timeout(arms, timeout) {
Ok(r) => r,
Err(e) => panic!(
@@ -776,7 +795,12 @@ pub fn try_select_timeout(
// would leave those fds unusable until a kernel event happened to
// clear them.
let eager = arms.iter().any(|a| a.sel_eager_cleanup());
let mut guard = UnregisterGuard { arms, me, epoch, armed: eager };
let mut guard = UnregisterGuard {
arms,
me,
epoch,
armed: eager,
};
crate::scheduler::park_current();
+26 -11
View File
@@ -16,10 +16,18 @@ thread_local! {
static ACTOR_SP: Cell<usize> = const { Cell::new(0) };
}
fn get_scheduler_sp() -> usize { SCHEDULER_SP.with(|c| c.get()) }
fn set_scheduler_sp(v: usize) { SCHEDULER_SP.with(|c| c.set(v)) }
pub fn get_actor_sp() -> usize { ACTOR_SP.with(|c| c.get()) }
pub fn set_actor_sp(v: usize) { ACTOR_SP.with(|c| c.set(v)) }
fn get_scheduler_sp() -> usize {
SCHEDULER_SP.with(|c| c.get())
}
fn set_scheduler_sp(v: usize) {
SCHEDULER_SP.with(|c| c.set(v))
}
pub fn get_actor_sp() -> usize {
ACTOR_SP.with(|c| c.get())
}
pub fn set_actor_sp(v: usize) {
ACTOR_SP.with(|c| c.set(v))
}
// ---------------------------------------------------------------------------
// Initial stack layout
@@ -49,13 +57,20 @@ pub fn set_actor_sp(v: usize) { ACTOR_SP.with(|c| c.set(v)) }
pub fn init_actor_stack(top: *mut u8, entry: extern "C-unwind" fn()) -> usize {
unsafe {
let mut sp = (top as usize & !15) - 8;
sp -= 8; (sp as *mut usize).write(entry as usize); // ret target
sp -= 8; (sp as *mut usize).write(0); // rbx
sp -= 8; (sp as *mut usize).write(0); // rbp
sp -= 8; (sp as *mut usize).write(0); // r12
sp -= 8; (sp as *mut usize).write(0); // r13
sp -= 8; (sp as *mut usize).write(0); // r14
sp -= 8; (sp as *mut usize).write(0); // r15
sp -= 8;
(sp as *mut usize).write(entry as usize); // ret target
sp -= 8;
(sp as *mut usize).write(0); // rbx
sp -= 8;
(sp as *mut usize).write(0); // rbp
sp -= 8;
(sp as *mut usize).write(0); // r12
sp -= 8;
(sp as *mut usize).write(0); // r13
sp -= 8;
(sp as *mut usize).write(0); // r14
sp -= 8;
(sp as *mut usize).write(0); // r15
sp
}
}
+78 -34
View File
@@ -178,7 +178,9 @@
//! from any handler via [`Watcher::watch`]) because monitors are inherently
//! created at runtime. The idle window is set once, in `init`.
use crate::channel::{channel, select, select_timeout, Receiver, RecvTimeoutError, Selectable, Sender};
use crate::channel::{
channel, select, select_timeout, Receiver, RecvTimeoutError, Selectable, Sender,
};
use crate::monitor::{demonitor, monitor, Down, Monitor};
use crate::pid::Pid;
use crate::registry::{register_with, resolve_named_sender, RegisterError};
@@ -273,7 +275,10 @@ pub struct GenServerRef<G: GenServer> {
impl<G: GenServer> Clone for GenServerRef<G> {
fn clone(&self) -> Self {
GenServerRef { tx: self.tx.clone(), pid: self.pid }
GenServerRef {
tx: self.tx.clone(),
pid: self.pid,
}
}
}
@@ -412,7 +417,9 @@ impl<G: GenServer> GenServerCtx<G> {
/// A clonable handle to the loop's monitor intake. Store it in the state
/// during `init` to watch monitors from later handlers.
pub fn watcher(&self) -> Watcher<G> {
Watcher { tx: self.sys_tx.clone() }
Watcher {
tx: self.sys_tx.clone(),
}
}
/// Shorthand for `ctx.watcher().watch(m)` when watching during `init`.
@@ -426,7 +433,10 @@ impl<G: GenServer> GenServerCtx<G> {
/// [`tick_every`](TimerHandle::tick_every) /
/// [`cancel`](TimerHandle::cancel) from any later handler.
pub fn timer(&self) -> TimerHandle<G> {
TimerHandle { sys_tx: self.sys_tx.clone(), reg: self.reg.clone() }
TimerHandle {
sys_tx: self.sys_tx.clone(),
reg: self.reg.clone(),
}
}
/// Set a quiet-period window: if the loop goes `after` without dispatching
@@ -518,7 +528,10 @@ pub struct TimerHandle<G: GenServer> {
// Manual Clone for the same reason as `Watcher`: no `G: Clone` needed.
impl<G: GenServer> Clone for TimerHandle<G> {
fn clone(&self) -> Self {
TimerHandle { sys_tx: self.sys_tx.clone(), reg: self.reg.clone() }
TimerHandle {
sys_tx: self.sys_tx.clone(),
reg: self.reg.clone(),
}
}
}
@@ -569,8 +582,18 @@ impl<G: GenServer> TimerHandle<G> {
// First instance fires after `every`; the payload is produced loop-side
// from `make` on fire, so the tick carries only the stable id.
let sub = send_after_to(every, self.sys_tx.clone(), Sys::Tick(local));
reg.periodics.insert(local, Periodic { every, live: sub, make });
debug_assert!(reg.rearm_tx.is_some(), "rearm_tx must be Some while periodics is non-empty");
reg.periodics.insert(
local,
Periodic {
every,
live: sub,
make,
},
);
debug_assert!(
reg.rearm_tx.is_some(),
"rearm_tx must be Some while periodics is non-empty"
);
local
}
@@ -617,7 +640,9 @@ pub struct Watcher<G: GenServer> {
// regardless of the server type (it clones only the inner sender).
impl<G: GenServer> Clone for Watcher<G> {
fn clone(&self) -> Self {
Watcher { tx: self.tx.clone() }
Watcher {
tx: self.tx.clone(),
}
}
}
@@ -691,7 +716,10 @@ impl<G: GenServer> GenServerBuilder<G> {
/// live server). Consumes the builder, carrying its `with_info` / `under`
/// configuration through.
pub fn named(self, name: GenServerName<G>) -> NamedGenServerBuilder<G> {
NamedGenServerBuilder { builder: self, name: name.as_str() }
NamedGenServerBuilder {
builder: self,
name: name.as_str(),
}
}
/// Private shared body behind [`start`](Self::start) and
@@ -700,18 +728,24 @@ impl<G: GenServer> GenServerBuilder<G> {
/// under the name before returning.
fn spawn_server(self) -> GenServerRef<G> {
let (tx, rx) = channel::<Envelope<G>>();
let GenServerBuilder { state, infos, supervisor, stack_opts } = self;
let GenServerBuilder {
state,
infos,
supervisor,
stack_opts,
} = self;
let handle = match supervisor {
Some(sup) => {
crate::scheduler::spawn_under_with(sup, stack_opts, move || {
server_loop::<G>(rx, state, infos)
})
}
None => crate::scheduler::spawn_with(stack_opts, move || {
Some(sup) => crate::scheduler::spawn_under_with(sup, stack_opts, move || {
server_loop::<G>(rx, state, infos)
}),
None => {
crate::scheduler::spawn_with(stack_opts, move || server_loop::<G>(rx, state, infos))
}
};
GenServerRef { tx, pid: handle.pid() }
GenServerRef {
tx,
pid: handle.pid(),
}
}
}
@@ -739,7 +773,10 @@ impl<G> GenServerName<G> {
/// associated constants at call sites.
#[inline]
pub const fn new(name: &'static str) -> Self {
Self { name, _marker: PhantomData }
Self {
name,
_marker: PhantomData,
}
}
/// The underlying registry key.
@@ -933,7 +970,11 @@ fn server_loop<G: GenServer>(
// Bind the ctx so the idle window set during init can be read back, then
// drop it — that drops the loop's own Sys sender, so a state that cloned no
// Watcher/TimerHandle lets the arm auto-close (the unused-ctx behaviour).
let ctx = GenServerCtx { sys_tx, reg: reg.clone(), idle: Cell::new(None) };
let ctx = GenServerCtx {
sys_tx,
reg: reg.clone(),
idle: Cell::new(None),
};
guard.0.init(&ctx);
let idle = ctx.idle.get();
drop(ctx);
@@ -983,13 +1024,12 @@ fn server_loop<G: GenServer>(
// inbox. The slice is rebuilt each iteration because the monitor
// and info sets shrink/grow. Mirrors the fast-path inbox park
// above; keep them in sync.
let nd = monitors.len(); // monitor band: [0, nd)
let nw = sys_open as usize; // system arm: [nd, nd+nw)
// info band: [nd+nw, nd+nw+ni)
// inbox arm: [nd+nw+ni]
let nd = monitors.len(); // monitor band: [0, nd)
let nw = sys_open as usize; // system arm: [nd, nd+nw)
// info band: [nd+nw, nd+nw+ni)
// inbox arm: [nd+nw+ni]
let sel = {
let mut arms: Vec<&dyn Selectable> =
Vec::with_capacity(nd + nw + infos.len() + 1);
let mut arms: Vec<&dyn Selectable> = Vec::with_capacity(nd + nw + infos.len() + 1);
for m in &monitors {
arms.push(&m.rx);
}
@@ -1001,9 +1041,7 @@ fn server_loop<G: GenServer>(
}
arms.push(&rx);
match idle_deadline {
Some(dl) => {
select_timeout(&arms, dl.saturating_duration_since(Instant::now()))
}
Some(dl) => select_timeout(&arms, dl.saturating_duration_since(Instant::now())),
None => Some(select(&arms)),
}
};
@@ -1034,8 +1072,12 @@ fn server_loop<G: GenServer>(
// live set tracks only still-pending timers, then
// dispatch.
match reg.lock() {
Ok(mut g) => { g.oneshots.remove(&id); }
Err(e) => panic!("smarm: gen_server reg lock poisoned (core corrupt): {e}"),
Ok(mut g) => {
g.oneshots.remove(&id);
}
Err(e) => {
panic!("smarm: gen_server reg lock poisoned (core corrupt): {e}")
}
}
guard.0.handle_timer(msg);
reset_idle(&mut idle_deadline);
@@ -1049,7 +1091,9 @@ fn server_loop<G: GenServer>(
let msg = {
let mut g = match reg.lock() {
Ok(g) => g,
Err(e) => panic!("smarm: gen_server reg lock poisoned (core corrupt): {e}"),
Err(e) => panic!(
"smarm: gen_server reg lock poisoned (core corrupt): {e}"
),
};
let r = &mut *g;
if let Some(p) = r.periodics.get_mut(&id) {
@@ -1057,9 +1101,9 @@ fn server_loop<G: GenServer>(
let msg = (p.make)();
let tx = match r.rearm_tx.as_ref() {
Some(tx) => tx.clone(),
None => panic!(
"smarm: live periodic without rearm_tx (logic bug)"
),
None => {
panic!("smarm: live periodic without rearm_tx (logic bug)")
}
};
p.live = send_after_to(every, tx, Sys::Tick(id));
Some(msg)
+25 -10
View File
@@ -219,7 +219,11 @@ struct Timers {
impl Timers {
fn new() -> Self {
Timers { next_local: 0, state: None, named: HashMap::new() }
Timers {
next_local: 0,
state: None,
named: HashMap::new(),
}
}
fn mint(&mut self) -> u64 {
@@ -248,7 +252,11 @@ pub struct Cx<Ev> {
impl<Ev> Cx<Ev> {
fn new(sys_tx: Sender<Sys>, reg: Arc<Mutex<Timers>>) -> Self {
Cx { sys_tx, reg, _ev: PhantomData }
Cx {
sys_tx,
reg,
_ev: PhantomData,
}
}
/// Arm the **state timeout**: fire a `state_timeout` event after `after` in
@@ -387,7 +395,10 @@ pub struct GenStatemRef<M: Machine> {
impl<M: Machine> Clone for GenStatemRef<M> {
fn clone(&self) -> Self {
GenStatemRef { tx: self.tx.clone(), pid: self.pid }
GenStatemRef {
tx: self.tx.clone(),
pid: self.pid,
}
}
}
@@ -443,13 +454,13 @@ pub fn spawn<M: Machine>(machine: M) -> GenStatemRef<M> {
/// a `_with` variant like the scheduler's own spawns.
///
/// Panics if called outside `Runtime::run()`.
pub fn spawn_with<M: Machine>(
opts: crate::scheduler::SpawnOpts,
machine: M,
) -> GenStatemRef<M> {
pub fn spawn_with<M: Machine>(opts: crate::scheduler::SpawnOpts, machine: M) -> GenStatemRef<M> {
let (tx, rx) = channel::<M::Ev>();
let handle = crate::scheduler::spawn_with(opts, move || statem_loop(rx, machine));
GenStatemRef { tx, pid: handle.pid() }
GenStatemRef {
tx,
pid: handle.pid(),
}
}
/// The machine actor body: `on_start`, then one `handle` per event until the
@@ -488,7 +499,9 @@ fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, mut machine: M) {
Sys::StateTimeout(local) => {
let mut t = match reg.lock() {
Ok(g) => g,
Err(e) => panic!("smarm: gen_statem reg lock poisoned (core corrupt): {e}"),
Err(e) => panic!(
"smarm: gen_statem reg lock poisoned (core corrupt): {e}"
),
};
match t.state {
Some((live, _)) if live == local => {
@@ -501,7 +514,9 @@ fn statem_loop<M: Machine>(rx: Receiver<M::Ev>, mut machine: M) {
Sys::Timeout(name, local) => {
let mut t = match reg.lock() {
Ok(g) => g,
Err(e) => panic!("smarm: gen_statem reg lock poisoned (core corrupt): {e}"),
Err(e) => panic!(
"smarm: gen_statem reg lock poisoned (core corrupt): {e}"
),
};
match t.named.get(name) {
Some(&(live, _)) if live == local => {
+17 -4
View File
@@ -239,7 +239,10 @@ pub fn snapshot() -> RuntimeSnapshot {
actors.push(info);
}
}
RuntimeSnapshot { format_version: SNAPSHOT_FORMAT_VERSION, actors }
RuntimeSnapshot {
format_version: SNAPSHOT_FORMAT_VERSION,
actors,
}
})
}
@@ -387,7 +390,10 @@ pub fn tree() -> RuntimeTree {
/// want to inspect again) and want the tree view of it without re-reading
/// the runtime.
pub fn tree_from(snap: RuntimeSnapshot) -> RuntimeTree {
let RuntimeSnapshot { format_version, actors } = snap;
let RuntimeSnapshot {
format_version,
actors,
} = snap;
let mut index_of: HashMap<Pid, usize> = HashMap::with_capacity(actors.len());
for (i, a) in actors.iter().enumerate() {
@@ -418,7 +424,10 @@ pub fn tree_from(snap: RuntimeSnapshot) -> RuntimeTree {
.into_iter()
.filter_map(|i| build_node(i, &children_of, &orphaned, &mut slots))
.collect();
RuntimeTree { format_version, roots: root_nodes }
RuntimeTree {
format_version,
roots: root_nodes,
}
}
fn build_node(
@@ -436,5 +445,9 @@ fn build_node(
.collect()
})
.unwrap_or_default();
Some(TreeNode { info, orphaned: orphaned[i], children })
Some(TreeNode {
info,
orphaned: orphaned[i],
children,
})
}
+4 -17
View File
@@ -136,7 +136,6 @@ pub struct IoThread {
waiters: Waiters,
// ----- Epoll machinery -----
/// The epollfd, owned by `IoThread`. Callable cross-thread via
/// `epoll_ctl` per the man page.
epollfd: RawFd,
@@ -147,7 +146,6 @@ pub struct IoThread {
shutdown_write: RawFd,
// ----- Threads -----
pool_thread: Option<OsJoinHandle<()>>,
epoll_thread: Option<OsJoinHandle<()>>,
}
@@ -284,9 +282,8 @@ impl IoThread {
events,
u64: fd as u64,
};
let r = unsafe {
libc::epoll_ctl(self.epollfd, libc::EPOLL_CTL_ADD, fd, &mut ev as *mut _)
};
let r =
unsafe { libc::epoll_ctl(self.epollfd, libc::EPOLL_CTL_ADD, fd, &mut ev as *mut _) };
if r < 0 {
return Err(io::Error::last_os_error());
}
@@ -398,12 +395,7 @@ fn epoll_loop(epollfd: RawFd, waiters: Waiters, rt: Weak<RuntimeInner>) {
loop {
let n = unsafe {
libc::epoll_wait(
epollfd,
events.as_mut_ptr(),
MAX_EVENTS as libc::c_int,
-1,
)
libc::epoll_wait(epollfd, events.as_mut_ptr(), MAX_EVENTS as libc::c_int, -1)
};
if n < 0 {
@@ -438,12 +430,7 @@ fn epoll_loop(epollfd: RawFd, waiters: Waiters, rt: Weak<RuntimeInner>) {
let entry = w.remove(&fd);
if entry.is_some() {
unsafe {
libc::epoll_ctl(
epollfd,
libc::EPOLL_CTL_DEL,
fd,
std::ptr::null_mut(),
);
libc::epoll_ctl(epollfd, libc::EPOLL_CTL_DEL, fd, std::ptr::null_mut());
}
}
entry
+34 -30
View File
@@ -11,36 +11,36 @@
//!
//! See `LOOM.md` for the design intent and the deferred-for-later list.
pub mod stack;
pub(crate) mod signal;
pub mod context;
pub mod preempt;
pub mod pid;
pub mod actor;
pub mod causal;
pub mod channel;
pub mod scheduler;
pub mod supervisor;
pub mod timer;
pub mod io;
pub mod mutex;
pub mod monitor;
pub mod registry;
pub mod pg;
pub mod link;
pub mod context;
pub mod gen_server;
pub mod gen_statem;
pub mod introspect;
pub mod io;
pub mod link;
pub mod monitor;
pub mod mutex;
#[cfg(feature = "observer")]
pub mod observer;
pub mod runtime;
pub(crate) mod park;
pub mod pg;
pub mod pid;
pub mod preempt;
pub(crate) mod raw_mutex;
pub(crate) mod slot_state;
pub(crate) mod sync_shim;
pub mod registry;
#[doc(hidden)] // pub only so benches/rq_micro.rs can drive the raw structures
pub mod run_queue;
pub mod runtime;
pub mod scheduler;
pub(crate) mod signal;
pub(crate) mod slot_state;
pub mod stack;
pub mod supervisor;
pub(crate) mod sync_shim;
pub mod timer;
pub mod trace;
pub mod causal;
// ---------------------------------------------------------------------------
// Global allocator
@@ -59,23 +59,28 @@ pub use channel::{
};
pub use gen_server::{
call, cast, shutdown, whereis_server, CallError, CallTimeoutError, CastError, GenServer,
NamedGenServerBuilder, GenServerBuilder, GenServerCtx, GenServerName, GenServerRef, TimerHandle, Watcher,
GenServerBuilder, GenServerCtx, GenServerName, GenServerRef, NamedGenServerBuilder,
TimerHandle, Watcher,
};
pub use gen_statem::{
CallError as GenStatemCallError, Cx, Machine, Reply, Resolution, SendError as GenStatemSendError,
GenStatemRef,
CallError as GenStatemCallError, Cx, GenStatemRef, Machine, Reply, Resolution,
SendError as GenStatemSendError,
};
pub use introspect::{StackInfo,
pub use introspect::{
actor_info, snapshot, tree, tree_from, ActorInfo, ActorState, RuntimeSnapshot, RuntimeTree,
TreeNode, SNAPSHOT_FORMAT_VERSION,
StackInfo, TreeNode, SNAPSHOT_FORMAT_VERSION,
};
pub use link::{link, trap_exit, unlink, ExitSignal};
pub use monitor::{
demonitor, mark_watchable, monitor, terminal_reason, Down, DownReason, Monitor, MonitorId,
};
pub use mutex::{LockTimeout, Mutex, MutexGuard};
#[cfg(feature = "observer")]
pub use observer::{ObserverReply, ObserverRequest};
pub use link::{link, trap_exit, unlink, ExitSignal};
pub use monitor::{demonitor, mark_watchable, monitor, terminal_reason, Down, DownReason, Monitor, MonitorId};
pub use mutex::{LockTimeout, Mutex, MutexGuard};
pub use pg::{
dispatch, join, leave, members, members_as, pick, pick_as, Incarnation, Member, NodeId,
};
pub use pid::{Addressable, Erased, Name, Pid, RawPid};
pub use pg::{dispatch, join, leave, members, members_as, pick, pick_as, Incarnation, Member, NodeId};
pub use registry::{
install, lookup_as, register, resolve_name, send, send_dyn, send_to, unregister, whereis,
NameResolution, RegisterError, SendError,
@@ -83,9 +88,8 @@ pub use registry::{
pub use runtime::{init, Config, Runtime};
pub use scheduler::{
block_on_io, cancel_timer, request_stop, run, self_pid, send_after, send_after_named,
send_after_named_wall, send_after_wall, sleep, sleep_wall,
spawn, spawn_addr, spawn_addr_with, spawn_under, spawn_under_with, spawn_with,
wait_readable, wait_readable_timeout, wait_writable,
send_after_named_wall, send_after_wall, sleep, sleep_wall, spawn, spawn_addr, spawn_addr_with,
spawn_under, spawn_under_with, spawn_with, wait_readable, wait_readable_timeout, wait_writable,
wait_writable_timeout, yield_now, FdArm, JoinError, JoinHandle, SpawnOpts,
};
pub use supervisor::{ChildSpec, OneForOne, Restart, Signal, Strategy};
+4 -1
View File
@@ -157,7 +157,10 @@ pub fn link<A>(target: Pid<A>) {
});
match my_trap {
Some(tx) => {
let _ = tx.send(ExitSignal { from: target, reason: DownReason::NoProc });
let _ = tx.send(ExitSignal {
from: target,
reason: DownReason::NoProc,
});
}
None => request_stop(me),
}
+4 -1
View File
@@ -171,7 +171,10 @@ pub fn monitor<A>(target: Pid<A>) -> Monitor {
});
if !registered {
let _ = tx.send(Down { pid: target, reason: DownReason::NoProc });
let _ = tx.send(Down {
pid: target,
reason: DownReason::NoProc,
});
}
Monitor { id, target, rx }
+29 -10
View File
@@ -158,7 +158,11 @@ impl TimerTarget for MutexCore {
if st.holder == Some(pid) {
return;
}
match st.waiters.iter().position(|w| w.pid == pid && w.epoch == epoch) {
match st
.waiters
.iter()
.position(|w| w.pid == pid && w.epoch == epoch)
{
Some(pos) => {
st.waiters.remove(pos);
true
@@ -246,7 +250,10 @@ impl<T> Mutex<T> {
Some(v) => v,
None => panic!("smarm: Mutex value missing on free fast path (core corrupt)"),
};
return Ok(MutexGuard { mutex: self, value: Some(value) });
return Ok(MutexGuard {
mutex: self,
value: Some(value),
});
}
}
@@ -287,7 +294,10 @@ impl<T> Mutex<T> {
Some(v) => v,
None => panic!("smarm: Mutex value missing after grant (core corrupt)"),
};
Ok(MutexGuard { mutex: self, value: Some(value) })
Ok(MutexGuard {
mutex: self,
value: Some(value),
})
} else {
Err(LockTimeout)
}
@@ -315,7 +325,10 @@ impl<T> Mutex<T> {
Some(v) => v,
None => panic!("smarm: Mutex value missing on try_lock free path (core corrupt)"),
};
Some(MutexGuard { mutex: self, value: Some(value) })
Some(MutexGuard {
mutex: self,
value: Some(value),
})
}
/// Blocking fallback used when called outside the smarm runtime.
@@ -329,10 +342,15 @@ impl<T> Mutex<T> {
Ok(mut g) => g.take(),
Err(e) => panic!("smarm: mutex value lock poisoned (core corrupt): {e}"),
};
if let Some(v) = v { break v; }
if let Some(v) = v {
break v;
}
std::thread::yield_now();
};
Ok(MutexGuard { mutex: self, value: Some(value) })
Ok(MutexGuard {
mutex: self,
value: Some(value),
})
}
}
@@ -342,7 +360,10 @@ impl<T> Clone for Mutex<T> {
/// lock and one protected value; locking through any clone excludes
/// every other clone.
fn clone(&self) -> Self {
Self { core: self.core.clone(), value: self.value.clone() }
Self {
core: self.core.clone(),
value: self.value.clone(),
}
}
}
@@ -388,9 +409,7 @@ impl<T: std::fmt::Debug> std::fmt::Debug for MutexGuard<'_, T> {
Some(v) => v,
None => panic!("smarm: MutexGuard value missing (core corrupt)"),
};
f.debug_tuple("MutexGuard")
.field(value)
.finish()
f.debug_tuple("MutexGuard").field(value).finish()
}
}
+38 -15
View File
@@ -105,7 +105,9 @@ mod parker {
impl Parker {
pub(super) fn new() -> Self {
Self { state: AtomicU32::new(EMPTY) }
Self {
state: AtomicU32::new(EMPTY),
}
}
/// Returns `true` = woken (permit consumed), `false` = timed out.
@@ -211,7 +213,10 @@ mod parker {
impl Parker {
pub(super) fn new() -> Self {
Self { permit: Mutex::new(false), cv: Condvar::new() }
Self {
permit: Mutex::new(false),
cv: Condvar::new(),
}
}
/// Returns `true` = woken (permit consumed), `false` = timed out.
@@ -395,12 +400,10 @@ impl Coordinator {
// The CAS is the exactly-one guarantee: whoever clears the bit
// owns the wake; a racing wake_one retries on the observed value
// (coherence: a failed CAS can never read older than `mask`).
match self.idle.compare_exchange(
mask,
mask & !bit,
Ordering::AcqRel,
Ordering::Acquire,
) {
match self
.idle
.compare_exchange(mask, mask & !bit, Ordering::AcqRel, Ordering::Acquire)
{
Ok(_) => {
self.parkers[id].unpark();
return true;
@@ -467,7 +470,8 @@ impl Coordinator {
// with the deadline still NO_DEADLINE compares `new < MAX` = true
// and over-wakes — the benign direction. (Under the mandated timer
// serialization this interleaving cannot occur anyway.)
self.tk_armed.store(self.deadline_nanos(deadline), Ordering::SeqCst);
self.tk_armed
.store(self.deadline_nanos(deadline), Ordering::SeqCst);
true
}
@@ -578,7 +582,10 @@ mod tests {
let t0 = Instant::now();
let r = c.park(0, None, || false);
assert_eq!(r, ParkResult::Woken);
assert!(t0.elapsed() < Duration::from_millis(100), "park blocked despite permit");
assert!(
t0.elapsed() < Duration::from_millis(100),
"park blocked despite permit"
);
}
#[test]
@@ -632,14 +639,20 @@ mod tests {
// Wait until all four are published idle.
let t0 = Instant::now();
while c.idle_mask().count_ones() != N as u32 {
assert!(t0.elapsed() < Duration::from_secs(5), "threads never parked");
assert!(
t0.elapsed() < Duration::from_secs(5),
"threads never parked"
);
std::thread::yield_now();
}
assert!(c.wake_one());
// Exactly one wakes; give the others a beat to (incorrectly) wake.
let t0 = Instant::now();
while woken.load(O::SeqCst) == 0 {
assert!(t0.elapsed() < Duration::from_secs(5), "wake_one woke nobody");
assert!(
t0.elapsed() < Duration::from_secs(5),
"wake_one woke nobody"
);
std::thread::yield_now();
}
std::thread::sleep(Duration::from_millis(100));
@@ -706,7 +719,10 @@ mod tests {
assert_eq!(c.armed_deadline_nanos(), c.deadline_nanos(d2));
c.disarm_timer(0);
assert_eq!(c.armed_deadline_nanos(), NO_DEADLINE);
assert!(c.try_arm_timer(1, d1), "role must be re-takeable after disarm");
assert!(
c.try_arm_timer(1, d1),
"role must be re-takeable after disarm"
);
c.disarm_timer(1);
}
@@ -721,7 +737,10 @@ mod tests {
let t0 = Instant::now();
let r = c.park(0, Some(far), || false);
assert_eq!(r, ParkResult::Woken, "re-arm wake lost");
assert!(t0.elapsed() < Duration::from_secs(5), "slept toward the stale deadline");
assert!(
t0.elapsed() < Duration::from_secs(5),
"slept toward the stale deadline"
);
c.disarm_timer(0);
}
@@ -796,7 +815,11 @@ mod tests {
assert!(c.try_arm_timer(0, far));
c.note_deadline(near);
let r = c.park(0, Some(far), || false);
assert_eq!(r, ParkResult::Woken, "re-arm wake lost through note_deadline");
assert_eq!(
r,
ParkResult::Woken,
"re-arm wake lost through note_deadline"
);
c.disarm_timer(0);
}
}
+80 -17
View File
@@ -221,7 +221,9 @@ pub(crate) struct ProcessGroups {
impl ProcessGroups {
pub(crate) fn new() -> Self {
Self { groups: HashMap::new() }
Self {
groups: HashMap::new(),
}
}
/// Insert `ms` into `group`. Idempotent on the *member*: if the member is
@@ -323,20 +325,33 @@ impl ProcessGroups {
fn members_where(&self, group: &str, mut is_live: impl FnMut(Pid) -> bool) -> Vec<Pid> {
self.groups
.get(group)
.map(|v| v.iter().map(|e| e.member.pid).filter(|&p| is_live(p)).collect())
.map(|v| {
v.iter()
.map(|e| e.member.pid)
.filter(|&p| is_live(p))
.collect()
})
.unwrap_or_default()
}
/// The first live member of `group` in insertion order — stateless
/// first-live `pick`, with the same read-path backstop as `members_where`.
fn first_member_where(&self, group: &str, mut is_live: impl FnMut(Pid) -> bool) -> Option<Pid> {
self.groups.get(group)?.iter().map(|e| e.member.pid).find(|&p| is_live(p))
self.groups
.get(group)?
.iter()
.map(|e| e.member.pid)
.find(|&p| is_live(p))
}
}
/// Build the full member identity for `pid` from runtime identity.
fn member_for(inner: &crate::runtime::RuntimeInner, pid: Pid) -> Member {
Member { node: inner.node_id, incarnation: inner.incarnation, pid }
Member {
node: inner.node_id,
incarnation: inner.incarnation,
pid,
}
}
/// Is `pid` a live actor right now? Generation-checked atomic slot-word read,
@@ -367,7 +382,10 @@ pub fn join<A>(group: impl Into<String>, pid: Pid<A>) -> bool {
let mon = monitor(pid);
let (rejected, reaped) = with_runtime(|inner| {
let ms = Membership { member: member_for(inner, pid), monitor: mon };
let ms = Membership {
member: member_for(inner, pid),
monitor: mon,
};
let mut pg = inner.process_groups.lock();
let reaped = pg.reap_group(&group);
let rejected = pg.join(&group, ms);
@@ -507,7 +525,11 @@ mod tests {
let (tx, rx) = channel::<Down>();
let ms = Membership {
member: member(index, generation),
monitor: Monitor { id: MonitorId(0), target: pid, rx },
monitor: Monitor {
id: MonitorId(0),
target: pid,
rx,
},
};
(ms, tx)
}
@@ -518,7 +540,10 @@ mod tests {
let (a, _ta) = synth(1, 0);
let (b, _tb) = synth(1, 0);
assert!(pg.join("workers", a).is_none(), "first join inserts");
assert!(pg.join("workers", b).is_some(), "second identical join is handed back");
assert!(
pg.join("workers", b).is_some(),
"second identical join is handed back"
);
assert_eq!(pg.members_of("workers"), vec![member(1, 0)]);
}
@@ -542,7 +567,10 @@ mod tests {
let (a, _ta) = synth(1, 0);
let (b, _tb) = synth(1, 1);
assert!(pg.join("g", a).is_none());
assert!(pg.join("g", b).is_none(), "different generation is a distinct member");
assert!(
pg.join("g", b).is_none(),
"different generation is a distinct member"
);
assert_eq!(pg.members_of("g"), vec![member(1, 0), member(1, 1)]);
}
@@ -555,16 +583,27 @@ mod tests {
pg.join("g", b);
assert!(pg.leave("g", member(1, 0)).is_some());
assert_eq!(pg.members_of("g"), vec![member(2, 0)]);
assert!(pg.leave("g", member(1, 0)).is_none(), "second leave finds nothing");
assert!(
pg.leave("g", member(1, 0)).is_none(),
"second leave finds nothing"
);
assert!(pg.leave("g", member(2, 0)).is_some());
assert!(pg.members_of("g").is_empty(), "group is now empty");
assert!(pg.leave("never", member(9, 0)).is_none(), "leaving an unknown group is a no-op");
assert!(
pg.leave("never", member(9, 0)).is_none(),
"leaving an unknown group is a no-op"
);
}
#[test]
fn remove_where_sweeps_every_group() {
let mut pg = ProcessGroups::new();
for (g, (m, _t)) in [("a", synth(1, 0)), ("a", synth(2, 0)), ("b", synth(1, 0)), ("c", synth(3, 0))] {
for (g, (m, _t)) in [
("a", synth(1, 0)),
("a", synth(2, 0)),
("b", synth(1, 0)),
("c", synth(3, 0)),
] {
pg.join(g, m);
}
// Death of pid index 1 (any generation) evicts it everywhere.
@@ -582,8 +621,16 @@ mod tests {
let pid = Pid::new(1, 0);
let (tx, rx) = channel::<Down>();
let dead = Membership {
member: Member { node: DEFAULT_NODE_ID, incarnation: Incarnation::new(7), pid },
monitor: Monitor { id: MonitorId(0), target: pid, rx },
member: Member {
node: DEFAULT_NODE_ID,
incarnation: Incarnation::new(7),
pid,
},
monitor: Monitor {
id: MonitorId(0),
target: pid,
rx,
},
};
let _keep = tx;
let (live, _tl) = synth(2, 0);
@@ -614,9 +661,17 @@ mod tests {
pg.join("b", b1);
// pid 1 dies: its group-a monitor receives a Down. Its group-b monitor
// has not — reap must still sweep pid 1 out of b by the pid predicate.
ta1.send(Down { pid: Pid::new(1, 0), reason: DownReason::Exit }).unwrap();
ta1.send(Down {
pid: Pid::new(1, 0),
reason: DownReason::Exit,
})
.unwrap();
let evicted = pg.reap_group("a");
assert_eq!(evicted.len(), 2, "pid 1's memberships in both a and b are evicted");
assert_eq!(
evicted.len(),
2,
"pid 1's memberships in both a and b are evicted"
);
assert_eq!(pg.members_of("a"), vec![member(2, 0)]);
assert!(pg.members_of("b").is_empty(), "swept from b too; pruned");
}
@@ -646,8 +701,16 @@ mod tests {
let dead = Pid::new(1, 0);
let oracle = |pid: Pid| pid != dead;
assert_eq!(pg.members_where("g", oracle), vec![Pid::new(2, 0)], "dead pid filtered from read");
assert_eq!(pg.first_member_where("g", oracle), Some(Pid::new(2, 0)), "pick skips the dead first member");
assert_eq!(
pg.members_where("g", oracle),
vec![Pid::new(2, 0)],
"dead pid filtered from read"
);
assert_eq!(
pg.first_member_where("g", oracle),
Some(Pid::new(2, 0)),
"pick skips the dead first member"
);
// Backstop does not evict — that stays the monitor's job; raw storage
// still holds both until reap runs.
+12 -3
View File
@@ -79,7 +79,10 @@ impl Pid<Erased> {
/// here; typing happens at typed-actor boundaries via [`Pid::from_raw`].
#[inline]
pub const fn new(index: u32, generation: u32) -> Self {
Self { raw: RawPid::new(index, generation), _marker: PhantomData }
Self {
raw: RawPid::new(index, generation),
_marker: PhantomData,
}
}
}
@@ -90,7 +93,10 @@ impl<A> Pid<A> {
/// resolution paths.
#[inline]
pub(crate) const fn from_raw(raw: RawPid) -> Self {
Self { raw, _marker: PhantomData }
Self {
raw,
_marker: PhantomData,
}
}
/// The raw identity, dropping the actor type — the key for identity-only
@@ -192,7 +198,10 @@ impl<M> Name<M> {
/// associated constants at call sites.
#[inline]
pub const fn new(name: &'static str) -> Self {
Self { name, _marker: PhantomData }
Self {
name,
_marker: PhantomData,
}
}
/// The underlying registry key.
+4 -1
View File
@@ -166,7 +166,10 @@ impl<T> RawMutex<T> {
{
self.lock_slow();
}
RawMutexGuard { m: self, prev_preempt }
RawMutexGuard {
m: self,
prev_preempt,
}
}
#[cold]
+56 -15
View File
@@ -1,4 +1,3 @@
//! Give an actor a name so other actors can find it and message it.
//!
//! Without the registry, the only way to reach an actor is to already be
@@ -196,7 +195,9 @@ impl<M> std::fmt::Display for SendError<M> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SendError::Unresolved(_) => write!(f, "no live actor registered under that name"),
SendError::Dead(_) => write!(f, "the addressed actor is no longer the live incarnation"),
SendError::Dead(_) => {
write!(f, "the addressed actor is no longer the live incarnation")
}
SendError::NoChannel(_) => write!(f, "actor has no channel for this message type"),
SendError::Closed(_) => write!(f, "the actor's channel for this type is closed"),
SendError::NoMember(_) => write!(f, "no live member in the process group"),
@@ -243,7 +244,10 @@ struct Mailbox {
impl Mailbox {
fn new(pid: Pid) -> Self {
Self { pid, channels: HashMap::new() }
Self {
pid,
channels: HashMap::new(),
}
}
/// Clone the `Sender<M>` for this actor, if it has one. Called **under the
@@ -294,7 +298,10 @@ pub(crate) struct Registry {
impl Registry {
pub(crate) fn new() -> Self {
Self { by_index: HashMap::new(), by_name: HashMap::new() }
Self {
by_index: HashMap::new(),
by_name: HashMap::new(),
}
}
/// Drop a dead holder's artifacts: every name bound to it, and its
@@ -303,7 +310,11 @@ impl Registry {
/// wholesale on pid mismatch) and is left untouched.
fn prune_holder(&mut self, holder: Pid) {
self.by_name.retain(|_, p| *p != holder);
if self.by_index.get(&holder.index()).is_some_and(|mb| mb.pid == holder) {
if self
.by_index
.get(&holder.index())
.is_some_and(|mb| mb.pid == holder)
{
self.by_index.remove(&holder.index());
}
}
@@ -351,7 +362,11 @@ impl Registry {
.iter()
.filter_map(|(&n, &p)| (p == mb.pid).then_some(n))
.collect();
Some(MailboxInfo { pid: mb.pid, names, depth: depth.min(u32::MAX as usize) as u32 })
Some(MailboxInfo {
pid: mb.pid,
names,
depth: depth.min(u32::MAX as usize) as u32,
})
}
}
@@ -428,13 +443,19 @@ pub(crate) fn register_with<M: Send + 'static>(
/// index from a dead prior incarnation (pid mismatch) is replaced wholesale.
/// Caller holds the registry lock and has established that `me` is live.
fn publish_channel<M: Send + 'static>(reg: &mut Registry, me: Pid, tx: Sender<M>) {
let mb = reg.by_index.entry(me.index()).or_insert_with(|| Mailbox::new(me));
let mb = reg
.by_index
.entry(me.index())
.or_insert_with(|| Mailbox::new(me));
if mb.pid != me {
*mb = Mailbox::new(me);
}
mb.channels.insert(
TypeId::of::<M>(),
Channel { sender: Box::new(tx), msg_type: type_name::<M>() },
Channel {
sender: Box::new(tx),
msg_type: type_name::<M>(),
},
);
}
@@ -473,7 +494,10 @@ pub fn install<A: Addressable>(tx: Sender<A::Msg>) -> Pid<A> {
pub(crate) fn install_for<M: Send + 'static>(pid: Pid, tx: Sender<M>) {
with_runtime(|inner| {
let mut reg = inner.registry.lock();
debug_assert!(live(inner, pid), "install_for: pid must be a freshly spawned, live actor");
debug_assert!(
live(inner, pid),
"install_for: pid must be a freshly spawned, live actor"
);
publish_channel::<M>(&mut reg, pid, tx);
});
}
@@ -573,7 +597,10 @@ pub(crate) fn resolve_named_sender<M: Send + 'static>(name: &str) -> Option<(Pid
}
// A live holder's mailbox is its own (publish replaces wholesale on
// pid mismatch, and one live actor per slot), so index lookup is safe.
let tx = reg.by_index.get(&pid.index()).and_then(Mailbox::clone_sender::<M>)?;
let tx = reg
.by_index
.get(&pid.index())
.and_then(Mailbox::clone_sender::<M>)?;
Some((pid, tx))
})
}
@@ -586,7 +613,11 @@ pub fn unregister(name: &str) -> Option<Pid> {
with_runtime(|inner| {
let mut reg = inner.registry.lock();
let pid = reg.by_name.remove(name)?;
if live(inner, pid) { Some(pid) } else { None }
if live(inner, pid) {
Some(pid)
} else {
None
}
})
}
@@ -618,12 +649,17 @@ pub fn send<M: Send + 'static>(name: Name<M>, msg: M) -> Result<(), SendError<M>
reg.prune_holder(pid);
return Err(SendError::Unresolved(msg));
}
match reg.by_index.get(&pid.index()).and_then(Mailbox::clone_sender::<M>) {
match reg
.by_index
.get(&pid.index())
.and_then(Mailbox::clone_sender::<M>)
{
Some(tx) => tx,
None => return Err(SendError::NoChannel(msg)),
}
};
tx.send(msg).map_err(|crate::channel::SendError(m)| SendError::Closed(m))
tx.send(msg)
.map_err(|crate::channel::SendError(m)| SendError::Closed(m))
})
}
@@ -647,7 +683,11 @@ fn send_to_pid<M: Send + 'static>(
match reg.by_index.get(&pid.index()).map(|m| m.pid) {
// Exact incarnation, still alive: its `M` channel, or NoChannel.
Some(stored) if stored == pid && live(inner, pid) => {
match reg.by_index.get(&pid.index()).and_then(Mailbox::clone_sender::<M>) {
match reg
.by_index
.get(&pid.index())
.and_then(Mailbox::clone_sender::<M>)
{
Some(tx) => tx,
None => return Err(SendError::NoChannel(msg)),
}
@@ -662,7 +702,8 @@ fn send_to_pid<M: Send + 'static>(
_ => return Err(SendError::Dead(msg)),
}
};
tx.send(msg).map_err(|crate::channel::SendError(m)| SendError::Closed(m))
tx.send(msg)
.map_err(|crate::channel::SendError(m)| SendError::Closed(m))
}
/// Deliver `msg` directly to the exact actor identified by `pid`. Unlike
+28 -5
View File
@@ -222,7 +222,10 @@ impl MpmcRing {
if diff == 0 {
// Our turn: claim the position.
match self.enqueue_pos.0.compare_exchange_weak(
pos, pos + 1, Ordering::Relaxed, Ordering::Relaxed,
pos,
pos + 1,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => {
// SAFETY: the claim gives us exclusive write access
@@ -250,7 +253,10 @@ impl MpmcRing {
let diff = seq as isize - (pos + 1) as isize;
if diff == 0 {
match self.dequeue_pos.0.compare_exchange_weak(
pos, pos + 1, Ordering::Relaxed, Ordering::Relaxed,
pos,
pos + 1,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => {
// SAFETY: the claim gives us exclusive read access;
@@ -464,19 +470,36 @@ mod tests {
let popped = popped.lock().unwrap();
assert_eq!(popped.len(), total, "count mismatch");
let set: HashSet<u64> = popped.iter().map(|p| ((p.index() as u64) << 32) | p.generation() as u64).collect();
let set: HashSet<u64> = popped
.iter()
.map(|p| ((p.index() as u64) << 32) | p.generation() as u64)
.collect();
assert_eq!(set.len(), total, "duplicate or lost element");
assert_eq!(pop(&q), None);
}
#[test]
fn mpmc_exactly_once_contended() {
exactly_once(MpmcRing::new(8, 4096), |q, p| q.push(p), |q| q.pop(), 4, 4, 1000);
exactly_once(
MpmcRing::new(8, 4096),
|q, p| q.push(p),
|q| q.pop(),
4,
4,
1000,
);
}
#[test]
fn striped_exactly_once_contended() {
exactly_once(StripedRing::new(8, 4096), |q, p| q.push(p), |q| q.pop(), 4, 4, 1000);
exactly_once(
StripedRing::new(8, 4096),
|q, p| q.push(p),
|q| q.pop(),
4,
4,
1000,
);
}
#[test]
+79 -31
View File
@@ -112,10 +112,11 @@
//! deadline, so an expiry wakes one scheduler, not a herd.
use crate::actor::{
clear_current_pid, is_actor_done, reset_actor_done, set_current_actor_box,
set_current_pid, take_last_outcome, Actor, Outcome,
clear_current_pid, is_actor_done, reset_actor_done, set_current_actor_box, set_current_pid,
take_last_outcome, Actor, Outcome,
};
use crate::channel::Sender;
use crate::context::{get_actor_sp, set_actor_sp, switch_to_actor};
use crate::io::IoThread;
use crate::monitor::{Down, DownReason, MonitorId};
use crate::pid::Pid;
@@ -124,11 +125,8 @@ use crate::raw_mutex::RawMutex;
use crate::slot_state::{StateWord, Status, Unpark};
use crate::supervisor::Signal;
use crate::timer::Timers;
use crate::context::{get_actor_sp, set_actor_sp, switch_to_actor};
use std::sync::atomic::{
AtomicBool, AtomicPtr, AtomicU32, AtomicU64, AtomicUsize, Ordering,
};
use std::sync::atomic::{AtomicBool, AtomicPtr, AtomicU32, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
@@ -175,7 +173,9 @@ impl Config {
pub fn exact(n: usize) -> Self {
assert!(n >= 1, "scheduler thread count must be ≥ 1");
Self {
min: n, max: n, exact: Some(n),
min: n,
max: n,
exact: Some(n),
alloc_interval: crate::preempt::DEFAULT_ALLOC_INTERVAL,
timeslice_cycles: crate::preempt::DEFAULT_TIMESLICE_CYCLES,
stack_pool_cap: n * 4,
@@ -196,7 +196,9 @@ impl Config {
assert!(e >= 1, "exact must be ≥ 1");
}
Self {
min, max, exact,
min,
max,
exact,
alloc_interval: crate::preempt::DEFAULT_ALLOC_INTERVAL,
timeslice_cycles: crate::preempt::DEFAULT_TIMESLICE_CYCLES,
stack_pool_cap: max * 4,
@@ -322,7 +324,9 @@ impl Default for Config {
.map(|n| n.get())
.unwrap_or(1);
Self {
min: 1, max: avail, exact: None,
min: 1,
max: avail,
exact: None,
alloc_interval: crate::preempt::DEFAULT_ALLOC_INTERVAL,
timeslice_cycles: crate::preempt::DEFAULT_TIMESLICE_CYCLES,
stack_pool_cap: avail * 4,
@@ -376,7 +380,9 @@ pub struct RuntimeStats {
impl RuntimeStats {
/// Sum of run queue lengths across all scheduler threads.
pub fn total_run_queue_len(&self) -> u64 {
self.inner.stats.iter()
self.inner
.stats
.iter()
.map(|s| s.run_queue_len.load(Ordering::Relaxed))
.sum()
}
@@ -400,7 +406,9 @@ impl RuntimeStats {
/// scheduler threads. Counters are reset at the start of each `run()`,
/// so after a run this reads that run's total.
pub fn slot_hits(&self) -> u64 {
self.inner.stats.iter()
self.inner
.stats
.iter()
.map(|s| s.slot_hits.load(Ordering::Relaxed))
.sum()
}
@@ -408,7 +416,9 @@ impl RuntimeStats {
/// RFC 005: total slot occupants displaced to the shared queue, summed
/// across scheduler threads. Reset at the start of each `run()`.
pub fn slot_displacements(&self) -> u64 {
self.inner.stats.iter()
self.inner
.stats
.iter()
.map(|s| s.slot_displacements.load(Ordering::Relaxed))
.sum()
}
@@ -713,7 +723,8 @@ impl Slot {
#[inline]
pub(crate) fn record_message(&self) {
let v = self.messages_received.load(Ordering::Relaxed);
self.messages_received.store(v.wrapping_add(1), Ordering::Relaxed);
self.messages_received
.store(v.wrapping_add(1), Ordering::Relaxed);
}
/// Read the received-message tally (Relaxed; cross-thread snapshot read).
@@ -732,7 +743,8 @@ impl Slot {
#[inline]
pub(crate) fn add_budget(&self, cycles: u64) {
let v = self.budget_cycles.load(Ordering::Relaxed);
self.budget_cycles.store(v.wrapping_add(cycles), Ordering::Relaxed);
self.budget_cycles
.store(v.wrapping_add(cycles), Ordering::Relaxed);
}
/// Read the accumulated budget cycles (Relaxed). Always 0 unless the
@@ -871,7 +883,6 @@ impl Slot {
Some(*unsafe { Box::from_raw(raw) })
}
}
}
// ---------------------------------------------------------------------------
@@ -1270,11 +1281,13 @@ impl Runtime {
// Re-initialise shared state for this run.
assert_eq!(
self.inner.run_queue.len(), 0,
self.inner.run_queue.len(),
0,
"run() called while previous run still active"
);
debug_assert_eq!(
self.inner.live_actors.load(Ordering::Acquire), 0,
self.inner.live_actors.load(Ordering::Acquire),
0,
"run() called while previous run still active"
);
// RFC 018: the IO producers reach the runtime (slot table + unpark)
@@ -1421,7 +1434,9 @@ impl Runtime {
/// Snapshot of runtime statistics for introspection / tests.
pub fn stats(&self) -> RuntimeStats {
RuntimeStats { inner: self.inner.clone() }
RuntimeStats {
inner: self.inner.clone(),
}
}
}
@@ -1456,7 +1471,10 @@ thread_local! {
}
#[derive(Copy, Clone)]
pub(crate) enum YieldIntent { Yield, Park }
pub(crate) enum YieldIntent {
Yield,
Park,
}
pub(crate) fn set_yield_intent(i: YieldIntent) {
YIELD_INTENT.with(|c| c.set(i));
@@ -1491,7 +1509,10 @@ pub const ROOT_PID: Pid = Pid::new(u32::MAX, u32::MAX);
/// for writes after reclaim — so an over-eager mark costs a cancel-write,
/// never data.
fn maybe_shrink_stack(slot: &Slot) {
let parks = slot.parks_since_shrink.load(Ordering::Relaxed).saturating_add(1);
let parks = slot
.parks_since_shrink
.load(Ordering::Relaxed)
.saturating_add(1);
slot.parks_since_shrink.store(parks, Ordering::Relaxed);
let sp = slot.sp.load(Ordering::Relaxed);
@@ -1590,12 +1611,19 @@ pub(crate) fn install_actor(
// publish below.
let (diag_reserve, diag_guard) = stack.shape();
let diag_top = stack.top() as usize;
slot.stop_ptr.store(Arc::as_ptr(&stop) as *mut _, Ordering::Release);
slot.stop_ptr
.store(Arc::as_ptr(&stop) as *mut _, Ordering::Release);
{
let mut cold = slot.cold.lock();
debug_assert!(cold.actor.is_none(), "install over live actor");
debug_assert!(cold.waiters.is_empty() && cold.monitors.is_empty() && cold.links.is_empty());
cold.actor = Some(Actor { pid, stack, supervisor, stop, trap: None });
cold.actor = Some(Actor {
pid,
stack,
supervisor,
stop,
trap: None,
});
cold.outstanding_handles = 1;
cold.outcome = None;
cold.pending_io_result = None;
@@ -1607,9 +1635,11 @@ pub(crate) fn install_actor(
slot.parks_since_shrink.store(0, Ordering::Relaxed);
slot.shrink_count.store(0, Ordering::Relaxed);
slot.diag_stack_top.store(diag_top, Ordering::Relaxed);
slot.diag_stack_reserve.store(diag_reserve, Ordering::Relaxed);
slot.diag_stack_reserve
.store(diag_reserve, Ordering::Relaxed);
slot.diag_stack_guard.store(diag_guard, Ordering::Relaxed);
slot.diag_pid.store(((idx as u64) << 32) | gen as u64, Ordering::Relaxed);
slot.diag_pid
.store(((idx as u64) << 32) | gen as u64, Ordering::Relaxed);
slot.store_closure(closure);
slot.reset_counters();
inner.live_actors.fetch_add(1, Ordering::Relaxed);
@@ -1618,7 +1648,10 @@ pub(crate) fn install_actor(
// Release store orders everything above before any Acquire reader.
slot.word.publish_queued(gen);
inner.enqueue(pid);
crate::te!(crate::trace::Event::Spawn { parent: supervisor, child: pid });
crate::te!(crate::trace::Event::Spawn {
parent: supervisor,
child: pid
});
pid
}
@@ -1635,14 +1668,19 @@ pub(crate) fn install_actor(
/// is released — a last-sender drop can unpark a receiver, which takes the
/// run-queue mutex; legal under a cold lock, but pointless to nest.
pub(crate) fn reclaim_slot(inner: &RuntimeInner, pid: Pid) {
let Some(slot) = inner.slot_at(pid) else { return };
let Some(slot) = inner.slot_at(pid) else {
return;
};
let dropped_outside;
{
let mut cold = slot.cold.lock();
if slot.status_for(pid) != Status::Done || cold.outstanding_handles != 0 {
return; // already reclaimed, or not yet eligible
}
debug_assert!(cold.actor.is_none(), "reclaiming a slot that still owns an actor");
debug_assert!(
cold.actor.is_none(),
"reclaiming a slot that still owns an actor"
);
dropped_outside = (
cold.outcome.take(),
cold.supervisor_channel.take(),
@@ -1741,7 +1779,10 @@ fn finalize_actor(inner: &Arc<RuntimeInner>, pid: Pid, outcome: Outcome) {
// Notify monitors. Sent outside any slot lock: `send` may unpark a parked
// receiver, which takes the run-queue mutex.
for (_, m) in monitors {
let _ = m.send(Down { pid, reason: down_reason });
let _ = m.send(Down {
pid,
reason: down_reason,
});
}
// Walk linked peers ONE AT A TIME (cold locks are leaves). For every
@@ -1773,7 +1814,10 @@ fn finalize_actor(inner: &Arc<RuntimeInner>, pid: Pid, outcome: Outcome) {
};
match trap {
Some(Some(tx)) => {
let _ = tx.send(crate::link::ExitSignal { from: pid, reason: down_reason });
let _ = tx.send(crate::link::ExitSignal {
from: pid,
reason: down_reason,
});
}
Some(None) => crate::scheduler::request_stop(peer),
None => {}
@@ -1929,7 +1973,9 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
let io_out = inner.io_outstanding.load(Ordering::Acquire)
+ inner.io_fd_waiters.load(Ordering::Acquire);
stats.run_queue_len.store(inner.run_queue.len(), Ordering::Relaxed);
stats
.run_queue_len
.store(inner.run_queue.len(), Ordering::Relaxed);
let pop = match inner.run_queue.pop() {
Some(pid) => Pop::Got(pid),
None => {
@@ -2081,7 +2127,9 @@ fn schedule_loop(inner: &Arc<RuntimeInner>, slot_idx: usize) {
}
// Update per-thread stats: record who's on-CPU.
stats.current_pid_index.store(pid.index(), Ordering::Relaxed);
stats
.current_pid_index
.store(pid.index(), Ordering::Relaxed);
set_actor_sp(sp);
set_current_pid(pid);
+83 -78
View File
@@ -68,9 +68,7 @@
use crate::actor::current_pid;
use crate::channel::Sender;
use crate::pid::{Name, Pid};
use crate::runtime::{
self, RuntimeInner, YieldIntent, RUNTIME,
};
use crate::runtime::{self, RuntimeInner, YieldIntent, RUNTIME};
use crate::supervisor::Signal;
use std::sync::atomic::Ordering;
use std::sync::Arc;
@@ -152,7 +150,9 @@ pub struct JoinHandle {
impl JoinHandle {
/// The identity of the actor this handle refers to.
pub fn pid(&self) -> Pid { self.pid }
pub fn pid(&self) -> Pid {
self.pid
}
/// Block the calling actor until the spawned actor finishes, then
/// report how it finished: `Ok(())` if it returned normally or stopped
@@ -182,12 +182,10 @@ impl JoinHandle {
crate::slot_state::Status::Stale => {
panic!("join: target slot has been reused")
}
crate::slot_state::Status::Done => {
Some(match cold.outcome.take() {
Some(outcome) => outcome,
None => panic!("Done slot must have outcome"),
})
}
crate::slot_state::Status::Done => Some(match cold.outcome.take() {
Some(outcome) => outcome,
None => panic!("Done slot must have outcome"),
}),
crate::slot_state::Status::Live => {
// begin_wait is lock-free, legal under the cold lock;
// registering under it makes the epoch atomic with
@@ -227,8 +225,7 @@ impl JoinHandle {
match slot.status_for(self.pid) {
crate::slot_state::Status::Stale => false,
status => {
cold.outstanding_handles =
cold.outstanding_handles.saturating_sub(1);
cold.outstanding_handles = cold.outstanding_handles.saturating_sub(1);
cold.outstanding_handles == 0
&& status == crate::slot_state::Status::Done
}
@@ -305,9 +302,7 @@ pub fn spawn(f: impl FnOnce() + Send + 'static) -> JoinHandle {
/// [`spawn`] with per-actor stack shape overrides (RFC 019).
pub fn spawn_with(opts: SpawnOpts, f: impl FnOnce() + Send + 'static) -> JoinHandle {
let parent = current_pid().unwrap_or_else(|| {
with_runtime(|_| crate::runtime::ROOT_PID)
});
let parent = current_pid().unwrap_or_else(|| with_runtime(|_| crate::runtime::ROOT_PID));
spawn_under_with(parent, opts, f)
}
@@ -339,7 +334,10 @@ pub fn spawn_under_with<A>(
crate::runtime::install_actor(inner, idx, sp, stack, supervisor, closure)
});
JoinHandle { pid, consumed: false }
JoinHandle {
pid,
consumed: false,
}
}
/// Spawn an actor that other actors can message directly by its [`Pid<A>`],
@@ -566,11 +564,9 @@ pub fn sleep(duration: std::time::Duration) {
let _np = NoPreempt::enter();
let epoch = begin_wait();
let deadline = crate::timer::deadline_from_now(duration);
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert_sleep(deadline, me, epoch),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert_sleep(deadline, me, epoch),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
});
park_current();
}
@@ -588,11 +584,9 @@ pub fn sleep_wall(duration: std::time::Duration) {
let _np = NoPreempt::enter();
let epoch = begin_wait();
let deadline = crate::timer::deadline_from_now(duration);
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert_sleep_wall(deadline, me, epoch),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert_sleep_wall(deadline, me, epoch),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
});
park_current();
}
@@ -608,15 +602,13 @@ pub fn insert_wait_timer(
target: std::sync::Arc<dyn crate::timer::TimerTarget>,
epoch: u32,
) {
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert(
deadline,
pid,
crate::timer::Reason::WaitTimeout { target, epoch },
),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert(
deadline,
pid,
crate::timer::Reason::WaitTimeout { target, epoch },
),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
});
}
@@ -646,11 +638,9 @@ pub fn send_after<A: crate::pid::Addressable>(
let fire = Box::new(move || {
let _ = crate::registry::send_to(dest, msg);
});
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert_send(deadline, dest.erase(), fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert_send(deadline, dest.erase(), fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
})
}
@@ -670,11 +660,9 @@ pub fn send_after_named<M: Send + 'static>(
let fire = Box::new(move || {
let _ = crate::registry::send(dest, msg);
});
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert_send(deadline, armed_by, fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert_send(deadline, armed_by, fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
})
}
@@ -693,11 +681,9 @@ pub fn send_after_wall<A: crate::pid::Addressable>(
let fire = Box::new(move || {
let _ = crate::registry::send_to(dest, msg);
});
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert_send_wall(deadline, dest.erase(), fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert_send_wall(deadline, dest.erase(), fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
})
}
@@ -714,11 +700,9 @@ pub fn send_after_named_wall<M: Send + 'static>(
let fire = Box::new(move || {
let _ = crate::registry::send(dest, msg);
});
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert_send_wall(deadline, armed_by, fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert_send_wall(deadline, armed_by, fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
})
}
@@ -738,11 +722,9 @@ pub(crate) fn send_after_to<T: Send + 'static>(
let fire = Box::new(move || {
let _ = tx.send(msg);
});
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.insert_send(deadline, armed_by, fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.insert_send(deadline, armed_by, fire),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
})
}
@@ -750,11 +732,9 @@ pub(crate) fn send_after_to<T: Send + 'static>(
/// it fires. Returns `true` if the timer was still pending and delivery is
/// now prevented, `false` if it had already fired or was already cancelled.
pub fn cancel_timer(id: crate::timer::TimerId) -> bool {
with_runtime(|inner| {
match inner.timers.lock() {
Ok(mut timers) => timers.cancel(id),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
}
with_runtime(|inner| match inner.timers.lock() {
Ok(mut timers) => timers.cancel(id),
Err(e) => panic!("smarm: timers lock poisoned (core corrupt): {e}"),
})
}
@@ -820,7 +800,8 @@ where
};
let mut cold = slot.cold.lock();
debug_assert_eq!(
slot.generation(), me.generation(),
slot.generation(),
me.generation(),
"block_on_io: own slot reused mid-park"
);
match cold.pending_io_result.take() {
@@ -946,12 +927,20 @@ pub struct FdArm {
impl FdArm {
/// An arm that becomes ready when `fd` is readable.
pub fn readable(fd: std::os::fd::RawFd) -> Self {
FdArm { fd, readable: true, writable: false }
FdArm {
fd,
readable: true,
writable: false,
}
}
/// An arm that becomes ready when `fd` is writable.
pub fn writable(fd: std::os::fd::RawFd) -> Self {
FdArm { fd, readable: false, writable: true }
FdArm {
fd,
readable: false,
writable: true,
}
}
}
@@ -976,9 +965,7 @@ impl crate::channel::Selectable for FdArm {
match io.as_mut() {
Some(io) => {
inner.io_fd_waiters.fetch_add(1, Ordering::AcqRel);
let r = io.epoll_register(
self.fd, pid, epoch, self.readable, self.writable,
);
let r = io.epoll_register(self.fd, pid, epoch, self.readable, self.writable);
if r.is_err() {
inner.io_fd_waiters.fetch_sub(1, Ordering::AcqRel);
}
@@ -1033,7 +1020,11 @@ fn poll_events(fd: std::os::fd::RawFd, readable: bool, writable: bool) -> std::i
if writable {
events |= libc::POLLOUT;
}
let mut pfd = libc::pollfd { fd, events, revents: 0 };
let mut pfd = libc::pollfd {
fd,
events,
revents: 0,
};
loop {
let r = unsafe { libc::poll(&mut pfd, 1, 0) };
if r < 0 {
@@ -1081,7 +1072,11 @@ pub fn wait_writable_timeout(
pub fn read(fd: std::os::fd::RawFd, buf: &mut [u8]) -> std::io::Result<usize> {
wait_readable(fd)?;
let n = unsafe { libc::read(fd, buf.as_mut_ptr() as *mut _, buf.len()) };
if n < 0 { Err(std::io::Error::last_os_error()) } else { Ok(n as usize) }
if n < 0 {
Err(std::io::Error::last_os_error())
} else {
Ok(n as usize)
}
}
/// Convenience wrapper: park until `fd` is writable, then perform the
@@ -1090,7 +1085,11 @@ pub fn read(fd: std::os::fd::RawFd, buf: &mut [u8]) -> std::io::Result<usize> {
pub fn write(fd: std::os::fd::RawFd, buf: &[u8]) -> std::io::Result<usize> {
wait_writable(fd)?;
let n = unsafe { libc::write(fd, buf.as_ptr() as *const _, buf.len()) };
if n < 0 { Err(std::io::Error::last_os_error()) } else { Ok(n as usize) }
if n < 0 {
Err(std::io::Error::last_os_error())
} else {
Ok(n as usize)
}
}
// ---------------------------------------------------------------------------
@@ -1099,12 +1098,15 @@ pub fn write(fd: std::os::fd::RawFd, buf: &[u8]) -> std::io::Result<usize> {
pub fn register_supervisor_channel(pid: Pid, sender: Sender<Signal>) {
with_runtime(|inner| {
let slot = inner.slot_at(pid)
let slot = inner
.slot_at(pid)
.unwrap_or_else(|| panic!("register_supervisor_channel: pid {:?} not found", pid));
let mut cold = slot.cold.lock();
assert_eq!(
slot.generation(), pid.generation(),
"register_supervisor_channel: pid {:?} not found", pid
slot.generation(),
pid.generation(),
"register_supervisor_channel: pid {:?} not found",
pid
);
cold.supervisor_channel = Some(sender);
});
@@ -1177,6 +1179,9 @@ mod send_after_to_tests {
crate::sleep(Duration::from_millis(30));
r2.store(true, Ordering::SeqCst);
});
assert!(reached.load(Ordering::SeqCst), "runtime survived the dead-channel fire");
assert!(
reached.load(Ordering::SeqCst),
"runtime survived the dead-channel fire"
);
}
}
+12 -3
View File
@@ -214,7 +214,10 @@ struct Buf {
impl Buf {
fn new() -> Self {
Buf { b: [0; 320], len: 0 }
Buf {
b: [0; 320],
len: 0,
}
}
fn s(&mut self, s: &str) {
for &c in s.as_bytes() {
@@ -274,8 +277,14 @@ mod tests {
#[test]
fn inside_guard_both_edges() {
assert_eq!(classify(GUARD_LO, TOP, RESERVE, GUARD), FaultClass::Guard);
assert_eq!(classify(GUARD_HI - 1, TOP, RESERVE, GUARD), FaultClass::Guard);
assert_eq!(classify(GUARD_LO + GUARD / 2, TOP, RESERVE, GUARD), FaultClass::Guard);
assert_eq!(
classify(GUARD_HI - 1, TOP, RESERVE, GUARD),
FaultClass::Guard
);
assert_eq!(
classify(GUARD_LO + GUARD / 2, TOP, RESERVE, GUARD),
FaultClass::Guard
);
}
#[test]
+9 -9
View File
@@ -188,8 +188,7 @@ impl StateWord {
loop {
let w = self.load();
debug_assert!(
matches!(word_state(w), ST_RUNNING | ST_RUNNING_NOTIFIED)
&& word_gen(w) == gen,
matches!(word_state(w), ST_RUNNING | ST_RUNNING_NOTIFIED) && word_gen(w) == gen,
"yield return from invalid word {w:#x}"
);
if self
@@ -247,8 +246,7 @@ impl StateWord {
loop {
let w = self.load();
debug_assert!(
matches!(word_state(w), ST_RUNNING | ST_RUNNING_NOTIFIED)
&& word_gen(w) == gen,
matches!(word_state(w), ST_RUNNING | ST_RUNNING_NOTIFIED) && word_gen(w) == gen,
"begin_wait from invalid word {w:#x}"
);
let next = word_epoch(w).wrapping_add(1) & EPOCH_MASK;
@@ -342,8 +340,7 @@ impl StateWord {
loop {
let w = self.load();
debug_assert!(
matches!(word_state(w), ST_RUNNING | ST_RUNNING_NOTIFIED)
&& word_gen(w) == gen,
matches!(word_state(w), ST_RUNNING | ST_RUNNING_NOTIFIED) && word_gen(w) == gen,
"clear_notify from invalid word {w:#x}"
);
if word_state(w) != ST_RUNNING_NOTIFIED {
@@ -372,8 +369,7 @@ impl StateWord {
pub(crate) fn set_done(&self, gen: u32) {
let prev = self.0.swap(pack(gen, 0, ST_DONE), Ordering::AcqRel);
debug_assert!(
matches!(word_state(prev), ST_RUNNING | ST_RUNNING_NOTIFIED)
&& word_gen(prev) == gen,
matches!(word_state(prev), ST_RUNNING | ST_RUNNING_NOTIFIED) && word_gen(prev) == gen,
"finalize from invalid word {prev:#x}"
);
}
@@ -538,7 +534,11 @@ mod loom_tests {
// not a pending notification.
assert!(word.try_claim(0));
assert_eq!(word.unpark(0, Some(epoch)), Unpark::Noop);
assert_eq!(word_state(word.load()), ST_RUNNING, "stale epoch notified a live run");
assert_eq!(
word_state(word.load()),
ST_RUNNING,
"stale epoch notified a live run"
);
});
}
+12 -5
View File
@@ -57,16 +57,19 @@ impl Stack {
}
let base = base as *mut u8;
let ret = unsafe {
libc::mprotect(base as *mut libc::c_void, guard_size, libc::PROT_NONE)
};
let ret = unsafe { libc::mprotect(base as *mut libc::c_void, guard_size, libc::PROT_NONE) };
if ret != 0 {
let err = io::Error::last_os_error();
unsafe { libc::munmap(base as *mut libc::c_void, total_size) };
return Err(err);
}
Ok(Self { base, total_size, stack_size, guard_size })
Ok(Self {
base,
total_size,
stack_size,
guard_size,
})
}
/// 16-byte-aligned top of the usable region.
@@ -177,7 +180,11 @@ pub(crate) fn shrink_range(hwm: usize, sp: usize, page: usize) -> Option<(usize,
/// with `stack_size` page-rounded by `Stack::new` the result is always
/// page-aligned. Checked math: `retain ≥ stack_size` (notably the default
/// 64 KiB reserve with the 64 KiB RETAIN) and overflow collapse to `None`.
pub(crate) fn retain_range(stack_size: usize, retain: usize, page: usize) -> Option<(usize, usize)> {
pub(crate) fn retain_range(
stack_size: usize,
retain: usize,
page: usize,
) -> Option<(usize, usize)> {
debug_assert!(page.is_power_of_two());
let retain = retain.checked_add(page - 1)? & !(page - 1); // page_up(retain)
let len = stack_size.checked_sub(retain)?;
+8 -6
View File
@@ -123,9 +123,9 @@ pub enum Signal {
impl std::fmt::Debug for Signal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Signal::Exit(pid) => write!(f, "Signal::Exit({:?})", pid),
Signal::Exit(pid) => write!(f, "Signal::Exit({:?})", pid),
Signal::Panic(pid, _) => write!(f, "Signal::Panic({:?}, ..)", pid),
Signal::Stopped(pid) => write!(f, "Signal::Stopped({:?})", pid),
Signal::Stopped(pid) => write!(f, "Signal::Stopped({:?})", pid),
}
}
}
@@ -133,9 +133,9 @@ impl std::fmt::Debug for Signal {
impl Signal {
pub fn pid(&self) -> Pid {
match self {
Signal::Exit(p) => *p,
Signal::Exit(p) => *p,
Signal::Panic(p, _) => *p,
Signal::Stopped(p) => *p,
Signal::Stopped(p) => *p,
}
}
}
@@ -169,7 +169,10 @@ pub struct ChildSpec {
impl ChildSpec {
pub fn new(restart: Restart, start: impl Fn() + Send + Sync + 'static) -> Self {
Self { start: Arc::new(start), restart }
Self {
start: Arc::new(start),
restart,
}
}
}
@@ -392,4 +395,3 @@ impl OneForOne {
}
}
}
+3 -1
View File
@@ -129,7 +129,9 @@ impl Ord for Entry {
// Earlier deadline first; ties broken by insertion order so the
// ordering is total. `Reason` and `Pid` deliberately don't
// participate.
self.deadline.cmp(&other.deadline).then_with(|| self.seq.cmp(&other.seq))
self.deadline
.cmp(&other.deadline)
.then_with(|| self.seq.cmp(&other.seq))
}
}
+44 -28
View File
@@ -16,13 +16,17 @@
#[cfg(feature = "smarm-trace")]
#[macro_export]
macro_rules! te {
($kind:expr) => { $crate::trace::record($kind) };
($kind:expr) => {
$crate::trace::record($kind)
};
}
#[cfg(not(feature = "smarm-trace"))]
#[macro_export]
macro_rules! te {
($kind:expr) => { () };
($kind:expr) => {
()
};
}
#[cfg(feature = "smarm-trace")]
@@ -68,8 +72,8 @@ mod inner {
// -----------------------------------------------------------------------
struct Record {
nanos: u64, // ns since open()
tid: u64, // OS thread id
nanos: u64, // ns since open()
tid: u64, // OS thread id
event: Event,
}
@@ -84,8 +88,8 @@ mod inner {
// -----------------------------------------------------------------------
struct Global {
sender: mpsc::Sender<Msg>,
start: Instant,
sender: mpsc::Sender<Msg>,
start: Instant,
}
static GLOBAL: Mutex<Option<Global>> = Mutex::new(None);
@@ -95,7 +99,7 @@ mod inner {
// The start Instant is copied alongside it — also one mutex hit per thread.
// record() never touches GLOBAL after that.
struct LocalState {
tx: mpsc::Sender<Msg>,
tx: mpsc::Sender<Msg>,
start: Instant,
}
@@ -109,8 +113,8 @@ mod inner {
// -----------------------------------------------------------------------
pub fn open() {
let path = std::env::var("SMARM_TRACE_FILE")
.unwrap_or_else(|_| "smarm_trace.json".to_owned());
let path =
std::env::var("SMARM_TRACE_FILE").unwrap_or_else(|_| "smarm_trace.json".to_owned());
let (tx, rx) = mpsc::channel::<Msg>();
let start = Instant::now();
@@ -164,8 +168,11 @@ mod inner {
// which would try to re-acquire inner.shared (already held at many
// te!() call sites) -> deadlock. Guard at the very top, before any
// allocation-capable call.
let was_enabled = crate::preempt::PREEMPTION_ENABLED
.with(|e| { let v = e.get(); e.set(false); v });
let was_enabled = crate::preempt::PREEMPTION_ENABLED.with(|e| {
let v = e.get();
e.set(false);
v
});
LOCAL_STATE.with(|cell| {
let mut opt = cell.borrow_mut();
@@ -182,7 +189,7 @@ mod inner {
}
if let Some(ls) = opt.as_ref() {
let nanos = ls.start.elapsed().as_nanos() as u64;
let tid = os_tid();
let tid = os_tid();
let _ = ls.tx.send(Msg::Event(Record { nanos, tid, event }));
}
});
@@ -197,7 +204,10 @@ mod inner {
fn drain_thread(rx: mpsc::Receiver<Msg>, path: &str) {
let f = match std::fs::File::create(path) {
Ok(f) => f,
Err(e) => { eprintln!("[smarm-trace] create failed: {}", e); return; }
Err(e) => {
eprintln!("[smarm-trace] create failed: {}", e);
return;
}
};
let mut w = std::io::BufWriter::new(f);
let _ = writeln!(w, "{{\"traceEvents\":[");
@@ -210,7 +220,9 @@ mod inner {
Ok(Msg::Event(r)) => {
let (name, actor_idx) = chrome_fields(&r.event);
let ts_us = r.nanos as f64 / 1000.0;
if !first { let _ = w.write_all(b",\n"); }
if !first {
let _ = w.write_all(b",\n");
}
first = false;
let _ = write!(w,
"{{\"ph\":\"i\",\"ts\":{:.3},\"pid\":{},\"tid\":{},\"name\":{:?},\"s\":\"g\"}}",
@@ -234,27 +246,31 @@ mod inner {
fn chrome_fields(ev: &Event) -> (String, u32) {
match ev {
Event::Spawn { parent, child } =>
(format!("spawn c={}", child.index()), parent.index()),
Event::Resume(p) => ("resume".into(), p.index()),
Event::Yield(p) => ("yield".into(), p.index()),
Event::Park(p) => ("park".into(), p.index()),
Event::Done(p) => ("done".into(), p.index()),
Event::UnparkDirect(p) => ("unpark_direct".into(), p.index()),
Event::UnparkDeferred(p) => ("unpark_deferred".into(), p.index()),
Event::Spawn { parent, child } => {
(format!("spawn c={}", child.index()), parent.index())
}
Event::Resume(p) => ("resume".into(), p.index()),
Event::Yield(p) => ("yield".into(), p.index()),
Event::Park(p) => ("park".into(), p.index()),
Event::Done(p) => ("done".into(), p.index()),
Event::UnparkDirect(p) => ("unpark_direct".into(), p.index()),
Event::UnparkDeferred(p) => ("unpark_deferred".into(), p.index()),
Event::UnparkFlagConsumed(p) => ("unpark_flag_consumed".into(), p.index()),
Event::Send { sender, receiver } => (
format!("send rx={}", receiver
.map(|p| p.index().to_string())
.unwrap_or_else(|| "none".into())),
format!(
"send rx={}",
receiver
.map(|p| p.index().to_string())
.unwrap_or_else(|| "none".into())
),
sender.index(),
),
Event::RecvPark(p) => ("recv_park".into(), p.index()),
Event::RecvWake(p) => ("recv_wake".into(), p.index()),
Event::Enqueue(p) => ("enqueue".into(), p.index()),
Event::Dequeue(p) => ("dequeue".into(), p.index()),
Event::Enqueue(p) => ("enqueue".into(), p.index()),
Event::Dequeue(p) => ("dequeue".into(), p.index()),
Event::SlotPush(p) => ("slot_push".into(), p.index()),
Event::SlotPop(p) => ("slot_pop".into(), p.index()),
Event::SlotPop(p) => ("slot_pop".into(), p.index()),
}
}