use crate::context::Root;
use crate::fiber::spawn_state::EffectiveApplyInput;
use crate::fiber::{DependencyEdge, Fiber, RestartError};
use crate::service::ServiceOccurrenceId;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct SemanticTarget {
input: EffectiveApplyInput,
assignments: Vec<(DependencyEdge, Option<ServiceOccurrenceId>)>,
}
impl SemanticTarget {
pub(super) fn new(
input: EffectiveApplyInput,
assignments: Vec<(DependencyEdge, Option<ServiceOccurrenceId>)>,
) -> Self {
Self { input, assignments }
}
pub(crate) fn has_missing(&self) -> bool {
self.assignments
.iter()
.any(|(_, publication)| publication.is_none())
}
}
const INERTIA_IDLE: u64 = 0;
const INERTIA_ACTIVE: u64 = 1;
const INERTIA_RELEASING: u64 = 2;
pub(crate) struct InertiaSlot {
inertia: AtomicU64,
target: parking_lot::Mutex<Option<SemanticTarget>>,
recheck_committed: AtomicU64,
recheck_settled: AtomicU64,
done: tokio::sync::Notify,
creation_owned: std::sync::atomic::AtomicBool,
}
pub(crate) enum InitialOutcome {
Active,
Pending,
Failed(super::PluginFailure),
Interrupted,
}
impl InertiaSlot {
pub(crate) fn new() -> Self {
Self {
inertia: AtomicU64::new(INERTIA_IDLE),
target: parking_lot::Mutex::new(None),
recheck_committed: AtomicU64::new(0),
recheck_settled: AtomicU64::new(0),
done: tokio::sync::Notify::new(),
creation_owned: std::sync::atomic::AtomicBool::new(false),
}
}
fn set_creation_owned(&self, owned: bool) {
self.creation_owned.store(owned, Ordering::SeqCst);
}
pub(crate) fn take_creation_ownership(&self) -> bool {
self.creation_owned.swap(false, Ordering::SeqCst)
}
fn current_target(&self) -> SemanticTarget {
self.target
.lock()
.clone()
.expect("an active convergence pass has a recorded SemanticTarget")
}
fn record_target(&self, target: SemanticTarget) {
*self.target.lock() = Some(target);
}
pub(crate) fn pin_inactive(&self) {
*self.target.lock() = None;
}
fn mark_pending_if_untouched(&self, pending: &SemanticTarget) {
let mut stored = self.target.lock();
if stored.is_none() {
*stored = Some(pending.clone());
}
}
pub(crate) async fn claim(&self) {
loop {
let notified = self.done.notified();
tokio::pin!(notified);
notified.as_mut().enable();
if self
.inertia
.compare_exchange(
INERTIA_IDLE,
INERTIA_ACTIVE,
Ordering::SeqCst,
Ordering::SeqCst,
)
.is_ok()
{
return;
}
notified.await;
}
}
pub(crate) fn commit_recheck(&self) {
self.recheck_committed.fetch_add(1, Ordering::SeqCst);
}
pub(crate) fn has_committed_recheck(&self) -> bool {
self.recheck_committed.load(Ordering::SeqCst) != self.recheck_settled.load(Ordering::SeqCst)
}
fn recheck_revision(&self) -> u64 {
self.recheck_committed.load(Ordering::SeqCst)
}
pub(crate) async fn wait_idle(&self) {
loop {
let notified = self.done.notified();
tokio::pin!(notified);
notified.as_mut().enable();
if self.inertia.load(Ordering::SeqCst) == INERTIA_IDLE {
return;
}
notified.await;
}
}
pub(crate) fn abandon(&self) {
self.inertia.store(INERTIA_IDLE, Ordering::SeqCst);
self.done.notify_waiters();
}
pub(crate) fn is_idle(&self) -> bool {
self.inertia.load(Ordering::SeqCst) == INERTIA_IDLE
}
fn try_claim(&self) -> bool {
self.inertia
.compare_exchange(
INERTIA_IDLE,
INERTIA_ACTIVE,
Ordering::SeqCst,
Ordering::SeqCst,
)
.is_ok()
}
fn exit_recheck(
&self,
fiber: &Arc<Fiber>,
root: &Arc<Root>,
settled: &SemanticTarget,
mut revision: u64,
) -> bool {
loop {
self.inertia.store(INERTIA_RELEASING, Ordering::SeqCst);
let live = fiber.compute_target(root);
let observed_revision = self.recheck_revision();
if live != *settled {
self.record_target(live);
self.inertia.store(INERTIA_ACTIVE, Ordering::SeqCst);
return true;
}
if observed_revision != revision {
revision = observed_revision;
continue;
}
self.recheck_settled.store(revision, Ordering::SeqCst);
match self.inertia.compare_exchange(
INERTIA_RELEASING,
INERTIA_IDLE,
Ordering::SeqCst,
Ordering::SeqCst,
) {
Ok(_) => {
self.done.notify_waiters();
return false;
}
Err(INERTIA_ACTIVE) => {
revision = self.recheck_revision();
}
Err(other) => panic!("invalid inertia release state {other}"),
}
}
}
pub(crate) fn kick(&self, fiber: &Arc<Fiber>, root: &Arc<Root>) {
if !fiber.is_alive() || !self.has_committed_recheck() {
return;
}
loop {
match self.inertia.load(Ordering::SeqCst) {
INERTIA_ACTIVE => return,
INERTIA_RELEASING => {
if self
.inertia
.compare_exchange(
INERTIA_RELEASING,
INERTIA_ACTIVE,
Ordering::SeqCst,
Ordering::SeqCst,
)
.is_ok()
{
return;
}
}
INERTIA_IDLE => {
let Ok(handle) = tokio::runtime::Handle::try_current() else {
return;
};
if self
.inertia
.compare_exchange(
INERTIA_IDLE,
INERTIA_ACTIVE,
Ordering::SeqCst,
Ordering::SeqCst,
)
.is_ok()
{
let revision = self.recheck_revision();
let settled = self.target.lock().clone();
if let Some(settled) = settled {
if self.exit_recheck(fiber, root, &settled, revision) {
drop(Self::spawn_convergence_pass(&handle, fiber, root));
}
} else {
self.record_target(fiber.compute_target(root));
drop(Self::spawn_convergence_pass(&handle, fiber, root));
}
return;
}
}
other => panic!("invalid inertia state {other}"),
}
}
}
fn spawn_convergence_pass(
handle: &tokio::runtime::Handle,
fiber: &Arc<Fiber>,
root: &Arc<Root>,
) -> tokio::task::JoinHandle<()> {
let logger = root.logger.logger_for_fiber(&fiber.name);
let context = format!("fiber={:?}", fiber.name);
let fiber = fiber.clone();
let root = root.clone();
crate::framework_task::spawn(handle, logger, "fiber convergence", context, async move {
fiber.slot.convergence_pass(&fiber, &root).await
})
}
pub(crate) async fn initial_spawn_pass(
&self,
fiber: &Arc<Fiber>,
root: &Arc<Root>,
) -> InitialOutcome {
if !self.try_claim() {
self.claim().await;
}
if !fiber.is_alive() {
fiber.force_unlink();
self.abandon();
return InitialOutcome::Interrupted;
}
self.set_creation_owned(true);
let outcome = loop {
let revision = self.recheck_revision();
let target = fiber.compute_target(root);
let inactive = target.has_missing();
fiber.settle_once(&target, true).await;
if inactive {
self.mark_pending_if_untouched(&target);
if self.exit_recheck(fiber, root, &target, revision) {
continue;
}
break InitialOutcome::Pending;
}
self.record_target(target.clone());
if fiber.state() == crate::fiber::FiberState::Failed {
let failure = fiber
.take_error()
.map(super::spawn::into_owned)
.unwrap_or_else(|| {
super::PluginFailure::returned(
"plugin failed without a parked failure".to_owned(),
)
});
fiber.creation_teardown().await;
break InitialOutcome::Failed(failure);
}
if self.exit_recheck(fiber, root, &target, revision) {
continue;
}
break InitialOutcome::Active;
};
self.set_creation_owned(false);
if !self.is_idle() {
self.abandon();
}
outcome
}
async fn convergence_pass(&self, fiber: &Arc<Fiber>, root: &Arc<Root>) {
loop {
let revision = self.recheck_revision();
let target = self.current_target();
fiber.settle_once(&target, false).await;
let live = fiber.compute_target(root);
if live != target {
self.record_target(live);
continue;
}
if self.exit_recheck(fiber, root, &target, revision) {
continue; }
return;
}
}
pub(crate) async fn restart_pass(
&self,
fiber: &Arc<Fiber>,
) -> std::result::Result<(), RestartError> {
self.claim().await;
if !fiber.is_alive() {
self.abandon();
return Err(RestartError::Closed);
}
let root = fiber
.spawn_state
.required_root()
.expect("live non-root Fiber has installed spawn state");
let revision = self.recheck_revision();
let target = fiber.compute_target(&root);
self.record_target(target.clone());
fiber.generation_replacing.store(true, Ordering::SeqCst);
let (tx, rx) = tokio::sync::oneshot::channel();
let fiber = fiber.clone();
crate::effect::detach(async move {
let result = fiber
.slot
.restart_committed(&fiber, &root, target, revision)
.await;
let _ = tx.send(result);
});
rx.await
.expect("framework-owned restart transaction always publishes completion")
}
async fn restart_committed(
&self,
fiber: &Arc<Fiber>,
root: &Arc<Root>,
mut target: SemanticTarget,
mut revision: u64,
) -> std::result::Result<(), RestartError> {
loop {
fiber.settle_once(&target, false).await;
let result = match fiber.error.lock().clone() {
Some(failure) => Err(RestartError::Apply(super::spawn::clone_owned(&failure))),
None => Ok(()),
};
if !self.exit_recheck(fiber, root, &target, revision) {
return result;
}
revision = self.recheck_revision();
target = fiber.compute_target(root);
self.record_target(target.clone());
}
}
pub(crate) async fn update_pass(
&self,
fiber: &Arc<Fiber>,
change: crate::PreparedChange,
) -> std::result::Result<crate::fiber::FiberState, crate::fiber::UpdateError> {
self.claim().await;
if !fiber.is_alive() {
self.abandon();
return Err(crate::fiber::UpdateError::AdmissionLost);
}
let root = fiber
.spawn_state
.required_root()
.expect("live Fiber has spawn state");
fiber.generation_replacing.store(true, Ordering::SeqCst);
fiber
.spawn_state
.commit_change(change)
.expect("precommit contract check and slot ownership make commit infallible");
let revision = self.recheck_revision();
let target = fiber.compute_target(&root);
self.record_target(target.clone());
let (tx, rx) = tokio::sync::oneshot::channel();
let fiber = fiber.clone();
crate::effect::detach(async move {
let result = fiber
.slot
.update_committed(&fiber, &root, target, revision)
.await;
let _ = tx.send(result);
});
rx.await
.expect("framework-owned update transaction publishes completion")
}
async fn update_committed(
&self,
fiber: &Arc<Fiber>,
root: &Arc<Root>,
mut target: SemanticTarget,
mut revision: u64,
) -> std::result::Result<crate::fiber::FiberState, crate::fiber::UpdateError> {
loop {
fiber.settle_once(&target, false).await;
let result = match fiber.error.lock().clone() {
Some(failure) => Err(crate::fiber::UpdateError::Apply(super::spawn::clone_owned(
&failure,
))),
None => Ok(fiber.state()),
};
if !self.exit_recheck(fiber, root, &target, revision) {
return result;
}
revision = self.recheck_revision();
target = fiber.compute_target(root);
self.record_target(target.clone());
}
}
}
#[cfg(test)]
mod tests {
use crate::context::Context;
use crate::context::RealmKey;
use crate::deadline::bounded;
use crate::fiber::{Fiber, FiberHandle, FiberState};
use crate::logger::{BufferExporter, Level};
use crate::service::Service;
use std::future::Future;
use std::sync::Arc;
use std::task::{Context as TaskContext, Poll, Waker};
fn fiber_requiring(ctx: &Context, inject: &[&str]) -> Arc<Fiber> {
let edges = inject
.iter()
.map(|name| crate::fiber::DependencyEdge::new((*name).to_string(), RealmKey::DEFAULT))
.collect();
let fiber = Fiber::new_with_edges("test", edges);
let prepared = crate::PreparedPlugin::from_input(
CountApplies {
applies: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
fail: false,
},
(),
);
let crate::PreparedPlugin {
plugin,
name,
inject,
contract,
} = prepared;
fiber
.spawn_state
.install(
plugin,
name,
inject,
ctx.clone(),
crate::registry::PluginKey::Typed(contract),
ctx.clone(),
)
.unwrap();
ctx.root.deps.register(&fiber);
fiber
}
#[derive(Debug, thiserror::Error)]
#[error("test apply failure")]
struct TestApplyError;
struct RetryWakeService;
impl crate::Service for RetryWakeService {
const NAME: &'static str = "phase2-retry-wake-service";
}
struct CountApplies {
applies: Arc<std::sync::atomic::AtomicUsize>,
fail: bool,
}
impl crate::Plugin for CountApplies {
type Config = ();
type Input = ();
type PrepareError = std::convert::Infallible;
type ApplyError = TestApplyError;
fn prepare(&self, _config: ()) -> std::result::Result<(), std::convert::Infallible> {
Ok(())
}
async fn apply(
&self,
_ctx: Context,
_prepared: &(),
) -> std::result::Result<(), TestApplyError> {
self.applies
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if self.fail {
return Err(TestApplyError);
}
Ok(())
}
}
fn fiber_in_creation(
ctx: &Context,
applies: Arc<std::sync::atomic::AtomicUsize>,
fail: bool,
) -> Arc<Fiber> {
let fiber = Fiber::new("test");
fiber
.creation_pending
.store(true, std::sync::atomic::Ordering::SeqCst);
let crate::PreparedPlugin {
plugin,
name,
inject,
contract,
} = crate::PreparedPlugin::from_input(CountApplies { applies, fail }, ());
fiber
.spawn_state
.install(
plugin,
name,
inject,
ctx.clone(),
crate::registry::PluginKey::Typed(contract),
ctx.clone(),
)
.unwrap();
fiber
}
#[tokio::test]
async fn a_convergence_winner_leaves_the_first_apply_to_the_initial_pass() {
let ctx = Context::new();
let root = ctx.root.clone();
let applies = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let fiber = fiber_in_creation(&ctx, applies.clone(), false);
assert!(fiber.slot.try_claim());
let target = fiber.compute_target(&root);
fiber.settle_once(&target, false).await;
assert_eq!(
applies.load(std::sync::atomic::Ordering::SeqCst),
0,
"a convergence pass never runs a creation-pending first apply"
);
assert_eq!(fiber.state(), FiberState::Pending);
fiber.slot.abandon();
let outcome = fiber.slot.initial_spawn_pass(&fiber, &root).await;
assert!(matches!(outcome, super::InitialOutcome::Active));
assert_eq!(
applies.load(std::sync::atomic::Ordering::SeqCst),
1,
"the creation applied exactly once (LF-01's single transition)"
);
}
#[tokio::test]
async fn a_first_apply_failure_is_the_creations_initial_apply_ending() {
let ctx = Context::new();
let root = ctx.root.clone();
let applies = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let fiber = fiber_in_creation(&ctx, applies.clone(), true);
assert!(fiber.slot.try_claim());
let target = fiber.compute_target(&root);
fiber.settle_once(&target, false).await;
fiber.slot.abandon();
let outcome = fiber.slot.initial_spawn_pass(&fiber, &root).await;
assert!(
matches!(outcome, super::InitialOutcome::Failed(_)),
"the first apply's failure is the creation's own ending"
);
assert_eq!(
applies.load(std::sync::atomic::Ordering::SeqCst),
1,
"no second apply ran"
);
assert!(!fiber.is_alive(), "the attempted fiber was torn down");
}
#[tokio::test]
async fn detached_convergence_invariant_panic_is_reported_and_remains_panicked() {
let ctx = Context::new();
let buffer = Arc::new(BufferExporter::new(8, Level::Debug).unwrap());
let _registration = ctx.add_exporter(buffer.clone()).unwrap();
let fiber = fiber_requiring(&ctx, &[]);
assert!(fiber.slot.try_claim());
let join = super::InertiaSlot::spawn_convergence_pass(
&tokio::runtime::Handle::current(),
&fiber,
&ctx.root,
);
let error = join
.await
.expect_err("the invariant panic still ends the task");
assert!(error.is_panic());
assert!(
!fiber.slot.is_idle(),
"framework panic reporting must not synthesize a recovery transition"
);
let records = buffer.snapshot();
assert_eq!(
records.len(),
1,
"one structured report per framework panic"
);
assert_eq!(records[0].level(), Level::Error);
assert_eq!(records[0].channel(), "test");
assert!(
records[0]
.text()
.contains("framework task fiber convergence panicked")
);
assert!(records[0].text().contains("fiber=\"test\""));
assert!(
records[0]
.text()
.contains("an active convergence pass has a recorded SemanticTarget")
);
fiber.slot.abandon();
}
#[tokio::test]
async fn a_convergence_applies_a_never_applied_fiber_after_handoff() {
let ctx = Context::new();
let root = ctx.root.clone();
let applies = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let fiber = fiber_in_creation(&ctx, applies.clone(), false);
fiber
.creation_pending
.store(false, std::sync::atomic::Ordering::SeqCst);
assert!(fiber.slot.try_claim());
let target = fiber.compute_target(&root);
fiber.settle_once(&target, false).await;
fiber.slot.abandon();
assert_eq!(
applies.load(std::sync::atomic::Ordering::SeqCst),
1,
"post-handoff convergence runs the first apply normally"
);
assert_eq!(fiber.state(), FiberState::Active);
}
#[test]
fn semantic_target_uses_fiber_edges_and_exact_publication_assignments() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &["counter"]);
let missing = fiber.compute_target(&ctx.root);
assert!(missing.has_missing());
struct Counter;
impl crate::Service for Counter {
const NAME: &'static str = "counter";
}
let publication = ctx.provide(Arc::new(Counter)).unwrap();
let present = fiber.compute_target(&ctx.root);
assert!(!present.has_missing());
assert_ne!(missing, present);
publication.set(Arc::new(Counter)).unwrap();
assert_eq!(
present,
fiber.compute_target(&ctx.root),
"same publication payload mutation is target-neutral"
);
}
#[test]
fn target_neutral_recheck_during_release_is_acknowledged_without_another_pass() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &[]);
assert!(fiber.slot.try_claim());
let target = fiber.compute_target(&ctx.root);
let revision = fiber.slot.recheck_revision();
fiber.slot.commit_recheck();
assert!(
!fiber
.slot
.exit_recheck(&fiber, &ctx.root, &target, revision),
"mechanism-only recheck revisions are not semantic target drift"
);
assert!(fiber.slot.is_idle());
assert!(!fiber.slot.has_committed_recheck());
}
#[tokio::test]
async fn target_neutral_idle_recheck_does_not_reapply_the_settled_target() {
let ctx = Context::new();
let root = ctx.root.clone();
let applies = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let fiber = fiber_in_creation(&ctx, applies.clone(), false);
let outcome = fiber.slot.initial_spawn_pass(&fiber, &root).await;
assert!(matches!(outcome, super::InitialOutcome::Active));
fiber
.creation_pending
.store(false, std::sync::atomic::Ordering::SeqCst);
assert_eq!(applies.load(std::sync::atomic::Ordering::SeqCst), 1);
fiber.slot.commit_recheck();
fiber.slot.kick(&fiber, &root);
fiber.slot.wait_idle().await;
assert_eq!(
applies.load(std::sync::atomic::Ordering::SeqCst),
1,
"same SemanticTarget rechecks must not trigger another apply"
);
assert!(!fiber.slot.has_committed_recheck());
}
#[test]
fn index_state_is_not_needed_to_reconstruct_pending_diagnostics_or_target() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &["counter"]);
ctx.root.deps.disable_for_test();
assert_eq!(fiber.missing_from(&ctx.root), vec!["counter".to_owned()]);
assert!(fiber.compute_target(&ctx.root).has_missing());
}
#[test]
fn off_runtime_recheck_claims_releasing_slot_for_the_existing_holder() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &["counter"]);
fiber.slot.inertia.store(
super::INERTIA_RELEASING,
std::sync::atomic::Ordering::SeqCst,
);
fiber.slot.commit_recheck();
fiber.slot.kick(&fiber, &ctx.root);
assert_eq!(
fiber.slot.inertia.load(std::sync::atomic::Ordering::SeqCst),
super::INERTIA_ACTIVE
);
assert!(fiber.slot.has_committed_recheck());
fiber.slot.abandon();
}
#[tokio::test]
async fn ready_rechecks_after_an_old_wake_before_returning() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &[RetryWakeService::NAME]);
let outcome = fiber.slot.initial_spawn_pass(&fiber, &ctx.root).await;
assert!(matches!(outcome, super::InitialOutcome::Pending));
assert_eq!(fiber.state(), FiberState::Pending);
fiber.slot.claim().await;
let handle = FiberHandle::new(fiber.clone());
let mut ready = Box::pin(handle.ready());
let mut task_cx = TaskContext::from_waker(Waker::noop());
assert!(matches!(ready.as_mut().poll(&mut task_cx), Poll::Pending));
fiber.slot.abandon();
let publication = std::thread::spawn({
let root = ctx.clone();
move || root.provide(Arc::new(RetryWakeService)).unwrap()
})
.join()
.unwrap();
assert!(fiber.slot.is_idle());
assert!(fiber.slot.has_committed_recheck());
assert_eq!(fiber.state(), FiberState::Pending);
assert!(matches!(ready.as_mut().poll(&mut task_cx), Poll::Pending));
drop(ready);
assert_eq!(
bounded(2000, handle.ready())
.await
.expect("replacement ready observes framework-owned convergence")
.unwrap(),
FiberState::Active
);
assert!(!fiber.slot.has_committed_recheck());
drop(publication);
}
#[tokio::test]
async fn cancelling_one_ready_waiter_does_not_block_another() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &[]);
let outcome = fiber.slot.initial_spawn_pass(&fiber, &ctx.root).await;
assert!(matches!(outcome, super::InitialOutcome::Active));
fiber.slot.claim().await;
let first_handle = FiberHandle::new(fiber.clone());
let second_handle = FiberHandle::new(fiber.clone());
let mut first = Box::pin(first_handle.ready());
let second = second_handle.ready();
tokio::pin!(second);
let mut task_cx = TaskContext::from_waker(Waker::noop());
assert!(matches!(first.as_mut().poll(&mut task_cx), Poll::Pending));
assert!(matches!(second.as_mut().poll(&mut task_cx), Poll::Pending));
drop(first);
fiber.slot.abandon();
assert!(matches!(
second.as_mut().poll(&mut task_cx),
Poll::Ready(Ok(FiberState::Active))
));
}
#[tokio::test]
async fn cancelling_one_claim_waiter_does_not_consume_release_for_another() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &[]);
fiber.slot.claim().await;
let mut first = Box::pin(fiber.slot.claim());
let second = fiber.slot.claim();
tokio::pin!(second);
let mut task_cx = TaskContext::from_waker(Waker::noop());
assert!(matches!(first.as_mut().poll(&mut task_cx), Poll::Pending));
assert!(matches!(second.as_mut().poll(&mut task_cx), Poll::Pending));
drop(first);
fiber.slot.abandon();
assert!(matches!(
second.as_mut().poll(&mut task_cx),
Poll::Ready(())
));
assert!(
!fiber.slot.is_idle(),
"surviving claimant owns the released slot"
);
fiber.slot.abandon();
}
#[test]
fn notify_waiters_between_future_creation_and_first_poll_is_observed() {
let notify = tokio::sync::Notify::new();
let notified = notify.notified();
tokio::pin!(notified);
notify.notify_waiters();
assert!(
notified.as_mut().enable(),
"broadcast generation survives until the first poll"
);
}
#[test]
fn notify_waiters_before_future_creation_leaves_no_later_permit() {
let notify = tokio::sync::Notify::new();
notify.notify_waiters();
let notified = notify.notified();
tokio::pin!(notified);
assert!(
!notified.as_mut().enable(),
"notify_waiters is a broadcast generation, not a stored later permit"
);
}
#[test]
fn one_notify_waiters_broadcast_reaches_all_precreated_futures() {
let notify = tokio::sync::Notify::new();
let first = notify.notified();
let second = notify.notified();
tokio::pin!(first);
tokio::pin!(second);
notify.notify_waiters();
assert!(
first.as_mut().enable(),
"one broadcast reaches the first pre-created Notified future"
);
assert!(
second.as_mut().enable(),
"one broadcast reaches the second pre-created Notified future"
);
}
#[test]
fn cancelling_one_precreated_waiter_does_not_consume_broadcast_for_another() {
let notify = tokio::sync::Notify::new();
let cancelled = Box::pin(notify.notified());
let survivor = notify.notified();
tokio::pin!(survivor);
drop(cancelled);
notify.notify_waiters();
assert!(
survivor.as_mut().enable(),
"cancelled waiter cannot consume another waiter's broadcast"
);
}
#[test]
fn enable_registers_multiple_waiters_for_single_permit_selection() {
let notify = tokio::sync::Notify::new();
let first = notify.notified();
let second = notify.notified();
tokio::pin!(first);
tokio::pin!(second);
assert!(!first.as_mut().enable());
assert!(!second.as_mut().enable());
notify.notify_one();
let ready_after_one =
usize::from(first.as_mut().enable()) + usize::from(second.as_mut().enable());
assert_eq!(
ready_after_one, 1,
"one notify_one permit selects exactly one enabled waiter"
);
notify.notify_one();
let ready_after_two =
usize::from(first.as_mut().enable()) + usize::from(second.as_mut().enable());
assert_eq!(
ready_after_two, 2,
"a second permit releases the other enabled waiter"
);
}
#[tokio::test]
async fn abandon_wakes_waiters() {
let ctx = Context::new();
let fiber = fiber_requiring(&ctx, &[]);
fiber.slot.claim().await;
let waiter = tokio::spawn({
let fiber = fiber.clone();
async move { fiber.slot.wait_idle().await }
});
fiber.slot.abandon();
bounded(2000, waiter)
.await
.expect("every IDLE store owes a notify_waiters")
.unwrap();
}
#[tokio::test]
async fn initial_spawn_pass_stays_pending_and_releases_when_services_missing() {
let ctx = Context::new();
let root = ctx.root.clone();
let fiber = fiber_requiring(&ctx, &["counter"]);
let outcome = fiber.slot.initial_spawn_pass(&fiber, &root).await;
assert!(
matches!(outcome, super::InitialOutcome::Pending),
"missing requirements end the pass as stable Pending"
);
assert_eq!(fiber.state(), FiberState::Pending);
assert!(fiber.slot.current_target().has_missing());
assert!(fiber.slot.is_idle());
}
#[tokio::test]
async fn initial_spawn_pass_on_an_invalidated_fiber_reports_interrupted() {
let ctx = Context::new();
let root = ctx.root.clone();
let fiber = fiber_requiring(&ctx, &[]);
root.registry
.attach_fiber(crate::registry::PluginKey::Anonymous(1), fiber.clone());
fiber.dispose().await;
let outcome = fiber.slot.initial_spawn_pass(&fiber, &root).await;
assert!(
matches!(outcome, super::InitialOutcome::Interrupted),
"invalidation before the pass reports Interrupted"
);
assert!(fiber.slot.is_idle(), "the arbiter is released");
assert_eq!(
root.registry.snapshot_fibers().len(),
0,
"no resident attempted Fiber remains observable"
);
}
#[test]
fn changed_committed_effective_input_changes_the_semantic_target() {
let ctx = Context::new();
let root = ctx.root.clone();
let applies = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let fiber = fiber_in_creation(&ctx, applies, false);
let before = fiber.compute_target(&root);
assert_eq!(
before,
fiber.compute_target(&root),
"plain restart with the same committed input is target-neutral"
);
fiber
.spawn_state
.commit_change(crate::PreparedChange::from_input::<CountApplies>(()))
.unwrap();
let after = fiber.compute_target(&root);
assert_ne!(
before, after,
"a committed apply-input change is target drift"
);
}
#[test]
fn normalized_declaration_order_produces_the_same_exact_edges() {
let a = crate::InjectSpec::none().require("beta").require("alpha");
let b = crate::InjectSpec::none().require("alpha").require("beta");
let edges = |spec: &crate::InjectSpec| {
spec.entries()
.map(|entry| {
crate::fiber::DependencyEdge::new(entry.name.clone(), RealmKey::DEFAULT)
})
.collect::<Vec<_>>()
};
assert_eq!(edges(&a), edges(&b));
}
}