use std::any::TypeId;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use parking_lot::Mutex;
use crate::context::Context;
use crate::effect::Disposable;
use crate::fiber::{Fiber, FiberState};
use crate::registry::{Plugin, RegistryService};
use crate::service::{CordisError, Service, ServiceInitFuture};
use crate::{FiberId, ReflectService};
#[derive(Debug)]
struct MtU1(pub u64);
impl Service for MtU1 {}
#[derive(Debug)]
struct MtU2(pub u64);
impl Service for MtU2 {}
#[derive(Debug)]
struct MtI1(pub u64);
impl Service for MtI1 {}
#[derive(Debug)]
struct MtCVal {
i1: u64,
u1: u64,
ready: bool,
}
impl Service for MtCVal {
fn check(&self) -> bool {
self.ready
}
}
struct MtIndepPlugin;
impl Plugin for MtIndepPlugin {
type Config = ();
type Provides = MtI1;
fn apply(&self, _ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtI1>, CordisError> {
Ok(Arc::new(MtI1(7)))
}
}
struct MtConsumerPlugin;
impl Plugin for MtConsumerPlugin {
type Config = ();
type Provides = MtCVal;
fn apply(&self, ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtCVal>, CordisError> {
let i1 = ctx.get::<MtI1>().map(|v| v.0);
let u1 = ctx.get::<MtU1>().map(|v| v.0);
Ok(Arc::new(MtCVal {
i1: i1.unwrap_or(0),
u1: u1.unwrap_or(0),
ready: i1.is_some() && u1.is_some(),
}))
}
}
#[derive(Debug)]
struct MtP1(pub u32);
impl Service for MtP1 {}
#[derive(Debug)]
struct MtP2(pub u32);
impl Service for MtP2 {}
#[derive(Debug)]
struct MtPair {
a: u32,
b: u32,
ready: bool,
}
impl Service for MtPair {
fn check(&self) -> bool {
self.ready
}
}
struct MtP1Plugin;
impl Plugin for MtP1Plugin {
type Config = ();
type Provides = MtP1;
fn apply(&self, _ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtP1>, CordisError> {
Ok(Arc::new(MtP1(7)))
}
}
struct MtP2Plugin;
impl Plugin for MtP2Plugin {
type Config = ();
type Provides = MtP2;
fn apply(&self, _ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtP2>, CordisError> {
Ok(Arc::new(MtP2(11)))
}
}
struct MtPairPlugin;
impl Plugin for MtPairPlugin {
type Config = ();
type Provides = MtPair;
fn apply(&self, ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtPair>, CordisError> {
let a = ctx.get::<MtP1>().map(|v| v.0);
let b = ctx.get::<MtP2>().map(|v| v.0);
Ok(Arc::new(MtPair {
a: a.unwrap_or(0),
b: b.unwrap_or(0),
ready: a.is_some() && b.is_some(),
}))
}
}
#[derive(Debug)]
struct MtProv(pub u32);
impl Service for MtProv {}
#[derive(Debug)]
struct MtDerived {
src: u32,
ready: bool,
}
impl Service for MtDerived {
fn check(&self) -> bool {
self.ready
}
}
struct MtDependentPlugin;
impl Plugin for MtDependentPlugin {
type Config = ();
type Provides = MtDerived;
fn apply(&self, ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtDerived>, CordisError> {
match ctx.get::<MtProv>() {
Some(v) => Ok(Arc::new(MtDerived {
src: v.0,
ready: true,
})),
None => Ok(Arc::new(MtDerived {
src: 0,
ready: false,
})),
}
}
}
struct MtFailPlugin;
impl Plugin for MtFailPlugin {
type Config = ();
type Provides = MtProv;
fn apply(&self, _ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtProv>, CordisError> {
Err(CordisError::Configuration("factory exploded".into()))
}
}
struct MtRevivePlugin;
impl Plugin for MtRevivePlugin {
type Config = ();
type Provides = MtProv;
fn apply(&self, _ctx: &Arc<Context>, _config: ()) -> Result<Arc<MtProv>, CordisError> {
Ok(Arc::new(MtProv(9)))
}
}
#[derive(Debug)]
struct MissingProbe;
impl Service for MissingProbe {}
#[derive(Debug)]
struct MtMark1(pub u64);
impl Service for MtMark1 {}
#[derive(Debug)]
struct MtMark2(pub u64);
impl Service for MtMark2 {}
struct MtEffA {
log: Arc<Mutex<Vec<&'static str>>>,
count: Arc<AtomicUsize>,
}
impl Service for MtEffA {
fn init(&self, _ctx: &Arc<Context>) -> ServiceInitFuture<'_> {
let log = self.log.clone();
let count = self.count.clone();
Box::pin(async move {
Ok(Some(Box::new(move || {
log.lock().push("A");
count.fetch_add(1, Ordering::SeqCst);
}) as Box<dyn Disposable>))
})
}
}
struct MtEffB {
log: Arc<Mutex<Vec<&'static str>>>,
count: Arc<AtomicUsize>,
}
impl Service for MtEffB {
fn init(&self, _ctx: &Arc<Context>) -> ServiceInitFuture<'_> {
let log = self.log.clone();
let count = self.count.clone();
Box::pin(async move {
Ok(Some(Box::new(move || {
log.lock().push("B");
count.fetch_add(1, Ordering::SeqCst);
}) as Box<dyn Disposable>))
})
}
}
fn base_root() -> (Arc<Context>, Arc<RegistryService>, Arc<ReflectService>) {
let ctx = Context::new_root();
let reg = ctx.provide(RegistryService::new());
let reflect = ctx.provide(ReflectService::new());
(ctx, reg, reflect)
}
async fn drain_spawned() {
for _ in 0..8 {
tokio::task::yield_now().await;
}
}
fn assert_quiescent(fibers: &[(String, Arc<Fiber>)], ctx: &Arc<Context>) -> Result<(), String> {
for (name, fiber) in fibers {
match fiber.state() {
FiberState::Loading | FiberState::Reloading | FiberState::Unloading { .. } => {
return Err(format!(
"fiber '{name}' rests in transitional state {:?}",
fiber.state()
));
}
FiberState::Failed { .. }
| FiberState::Inactive { .. }
| FiberState::Pending => {}
FiberState::Active { .. } => {
for tid in fiber.injected_type_ids() {
if !ctx.is_available(tid) {
return Err(format!(
"fiber '{name}' is Active but injected dependency {tid:?} is unavailable"
));
}
}
}
}
}
Ok(())
}
struct MtRng(u64);
impl MtRng {
const SEED: u64 = 0x9E37_79B9_7F4A_7C15;
fn new() -> Self {
MtRng(Self::SEED)
}
fn next(&mut self) -> u64 {
let mut x = self.0;
x ^= x >> 12;
x ^= x << 25;
x ^= x >> 27;
self.0 = x;
x.wrapping_mul(0x0254_5F49_14BF_49C9)
}
fn below(&mut self, n: u64) -> u64 {
self.next() % n.max(1)
}
}
pub async fn quiescence_after_every_op() -> Result<(), String> {
let (ctx, reg, reflect) = base_root();
let mut rng = MtRng::new();
let mut fibers: Vec<(String, Arc<Fiber>)> = vec![("root".to_string(), ctx.fiber())];
let mut newest: Option<usize> = None;
for step in 0..24u32 {
match rng.below(8) {
0 => {
let v = rng.next() % 1000;
ctx.provide(MtU1(v));
}
1 => {
let v = rng.next() % 1000;
ctx.provide(MtU2(v));
}
2 => {
reflect.notify_with_ctx(TypeId::of::<MtU1>(), &ctx).await;
}
3 => {
reflect.notify_with_ctx(TypeId::of::<MtI1>(), &ctx).await;
}
4 => match reg.register(&ctx, MtIndepPlugin, ()) {
Ok(fid) => {
let fiber = reg.get_fiber(fid).expect("registered fiber tracked");
fibers.push((format!("indep-{fid}"), fiber));
newest = Some(fibers.len() - 1);
}
Err(e) => {
if !e.to_string().contains("duplicate provider") {
return Err(format!("step {step}: unexpected register error: {e}"));
}
}
},
5 => match reg.register(&ctx, MtConsumerPlugin, ()) {
Ok(fid) => {
let fiber = reg.get_fiber(fid).expect("registered fiber tracked");
fiber.declare_inject::<MtI1>();
fiber.declare_inject::<MtU1>();
fiber.refresh(&ctx).await;
fibers.push((format!("consumer-{fid}"), fiber));
newest = Some(fibers.len() - 1);
}
Err(e) => {
if !e.to_string().contains("duplicate provider") {
return Err(format!("step {step}: unexpected register error: {e}"));
}
}
},
6 => match ctx.remove::<MtU1>() {
Ok(_) => {}
Err(e) => {
if !e.to_string().contains("guarded withdrawal") {
return Err(format!("step {step}: unexpected removal error: {e}"));
}
}
},
_ => {
if let Some(idx) = newest.take() {
let _ = fibers[idx].1.dispose().await;
}
}
}
drain_spawned().await;
assert_quiescent(&fibers, &ctx)
.map_err(|e| format!("quiescence broken after step {step}: {e}"))?;
}
if ctx.get::<MtI1>().is_none() {
let fid = reg
.register(&ctx, MtIndepPlugin, ())
.map_err(|e| format!("convergence register failed: {e}"))?;
fibers.push((format!("indep-{fid}"), reg.get_fiber(fid).unwrap()));
}
let consumer_idx = match fibers.iter().rposition(|(n, _)| n.starts_with("consumer")) {
Some(i) => i,
None => {
let fid = reg
.register(&ctx, MtConsumerPlugin, ())
.map_err(|e| format!("convergence register failed: {e}"))?;
let fiber = reg.get_fiber(fid).unwrap();
fiber.declare_inject::<MtI1>();
fiber.declare_inject::<MtU1>();
fiber.refresh(&ctx).await;
fibers.push((format!("consumer-{fid}"), fiber));
fibers.len() - 1
}
};
if ctx.get::<MtU1>().is_none() {
ctx.provide(MtU1(4242));
}
reflect.notify_with_ctx(TypeId::of::<MtI1>(), &ctx).await;
reflect.notify_with_ctx(TypeId::of::<MtU1>(), &ctx).await;
drain_spawned().await;
assert_quiescent(&fibers, &ctx).map_err(|e| format!("final quiescence: {e}"))?;
let consumer = fibers[consumer_idx].1.clone();
consumer.refresh(&ctx).await;
match consumer.state() {
FiberState::Active { .. } => {}
other => return Err(format!("consumer did not converge to Active: {other:?}")),
}
let fresh = consumer.compute_epoch(&ctx);
if consumer.epoch() != fresh {
return Err(format!(
"stale epoch '{}' != freshly computed '{}'",
consumer.epoch(),
fresh
));
}
let cval = ctx
.get::<MtCVal>()
.ok_or("consumer projection missing after convergence")?;
let u1 = ctx.get::<MtU1>().ok_or("MtU1 missing after convergence")?;
if !cval.ready || cval.i1 != 7 || cval.u1 != u1.0 {
return Err(format!("consumer view diverged from store: {cval:?}"));
}
Ok(())
}
pub async fn order_confluence_of_registrations() -> Result<(), String> {
let (ca, ra, _fa) = base_root();
ra.register(&ca, MtP1Plugin, ())
.map_err(|e| format!("A/P1: {e}"))?;
ra.register(&ca, MtP2Plugin, ())
.map_err(|e| format!("A/P2: {e}"))?;
let fid_a = ra
.register(&ca, MtPairPlugin, ())
.map_err(|e| format!("A/C: {e}"))?;
let fib_a = ra.get_fiber(fid_a).expect("tracked");
fib_a.declare_inject::<MtP1>();
fib_a.declare_inject::<MtP2>();
fib_a.refresh(&ca).await;
drain_spawned().await;
match fib_a.state() {
FiberState::Active { .. } => {}
other => return Err(format!("A: consumer not Active: {other:?}")),
}
let epoch_a = fib_a.epoch();
let pair_a = ca.get::<MtPair>().ok_or("A: pair value missing")?;
if !(pair_a.ready && pair_a.a == 7 && pair_a.b == 11) {
return Err(format!("A: wrong projection {pair_a:?}"));
}
let (cb, rb, fb) = base_root();
let fid_b = rb
.register(&cb, MtPairPlugin, ())
.map_err(|e| format!("B/C: {e}"))?;
let fib_b = rb.get_fiber(fid_b).expect("tracked");
fib_b.declare_inject::<MtP1>();
fib_b.declare_inject::<MtP2>();
fib_b.refresh(&cb).await;
match fib_b.state() {
FiberState::Inactive { .. } => {}
other => {
return Err(format!(
"B: consumer should be Inactive pre-providers: {other:?}"
))
}
}
rb.register(&cb, MtP1Plugin, ())
.map_err(|e| format!("B/P1: {e}"))?;
fb.notify_with_ctx(TypeId::of::<MtP1>(), &cb).await;
drain_spawned().await;
match fib_b.state() {
FiberState::Inactive { .. } => {}
other => {
return Err(format!(
"B: consumer should stay Inactive with one provider: {other:?}"
))
}
}
rb.register(&cb, MtP2Plugin, ())
.map_err(|e| format!("B/P2: {e}"))?;
fb.notify_with_ctx(TypeId::of::<MtP2>(), &cb).await;
drain_spawned().await;
match fib_b.state() {
FiberState::Active { .. } => {}
other => {
return Err(format!(
"B: consumer not Active after both providers: {other:?}"
))
}
}
let epoch_b = fib_b.epoch();
let pair_b = cb.get::<MtPair>().ok_or("B: pair value missing")?;
if epoch_a != epoch_b {
return Err(format!(
"confluence violated: epoch A '{epoch_a}' != epoch B '{epoch_b}'"
));
}
if pair_b.a != pair_a.a || pair_b.b != pair_a.b || pair_b.ready != pair_a.ready {
return Err(format!(
"confluence violated: projections differ ({:?} vs {:?})",
*pair_a, *pair_b
));
}
Ok(())
}
pub async fn dependent_never_active_without_provider() -> Result<(), String> {
let (ctx, reg, reflect) = base_root();
let fid = reg
.register(&ctx, MtDependentPlugin, ())
.map_err(|e| format!("register dependent: {e}"))?;
let dep = reg.get_fiber(fid).ok_or("dependent fiber not tracked")?;
dep.declare_inject::<MtProv>();
dep.refresh(&ctx).await;
match dep.state() {
FiberState::Inactive { .. } => {}
other => return Err(format!("pre-provider state should be Inactive: {other:?}")),
}
let raw = Arc::new(Fiber::new());
raw.declare_inject::<MtProv>();
match raw.state() {
FiberState::Inactive { error: None } => {}
other => {
return Err(format!(
"fresh raw fiber should be Inactive{{error:None}}: {other:?}"
))
}
}
raw.refresh(&ctx).await;
match raw.state() {
FiberState::Inactive { error: Some(_) } => {}
other => {
return Err(format!(
"refreshed raw fiber should report the missing dep: {other:?}"
))
}
}
ctx.provide(MtProv(1));
reflect.notify_with_ctx(TypeId::of::<MtProv>(), &ctx).await;
drain_spawned().await;
match dep.state() {
FiberState::Active { .. } => {}
other => return Err(format!("dependent should activate on provide: {other:?}")),
}
let d1 = ctx.get::<MtDerived>().ok_or("derived value missing")?;
if d1.src != 1 {
return Err(format!("projection should see v1, got {}", d1.src));
}
let epoch_v1 = dep.epoch();
raw.refresh(&ctx).await;
if !matches!(raw.state(), FiberState::Active { .. }) {
return Err(format!("raw fiber should activate: {:?}", raw.state()));
}
ctx.provide(MtProv(2));
reflect.notify_with_ctx(TypeId::of::<MtProv>(), &ctx).await;
drain_spawned().await;
match dep.state() {
FiberState::Active { .. } => {}
other => {
return Err(format!(
"dependent should survive provider swap Active: {other:?}"
))
}
}
let d2 = ctx
.get::<MtDerived>()
.ok_or("derived value missing after swap")?;
if d2.src != 2 {
return Err(format!("projection should see v2, got {}", d2.src));
}
let epoch_v2 = dep.epoch();
if epoch_v1 == epoch_v2 {
return Err(format!(
"epoch must change across a provider swap (both were '{epoch_v1}')"
));
}
let err = match ctx.remove::<MtProv>() {
Err(e) => e,
Ok(_) => {
return Err("guarded withdrawal must refuse removal of a consumed provider".to_string())
}
};
if !err.to_string().contains("guarded withdrawal") {
return Err(format!("unexpected refusal reason: {err}"));
}
let _ = dep.dispose().await;
if ctx.get::<MtDerived>().is_some() {
return Err("disposal must retract the dependent's projection".to_string());
}
let removed = ctx
.remove::<MtProv>()
.map_err(|e| format!("removal after disposal should pass: {e}"))?;
if removed.as_ref().map(|v| v.0) != Some(2) {
return Err("removed provider should be the v2 instance".to_string());
}
reflect.notify_with_ctx(TypeId::of::<MtProv>(), &ctx).await;
drain_spawned().await;
raw.refresh(&ctx).await;
match raw.state() {
FiberState::Inactive { .. } => {}
other => {
return Err(format!(
"raw fiber must deactivate after provider removal: {other:?}"
))
}
}
let fail_err = reg
.register(&ctx, MtFailPlugin, ())
.expect_err("failing factory must be refused");
if !fail_err.to_string().contains("factory exploded") {
return Err(format!("unexpected failure reason: {fail_err}"));
}
let fail_fid = next_fid_after(®, fid)?;
let failed = reg
.get_fiber(fail_fid)
.ok_or("failed fiber must stay inspectable via get_fiber")?;
match failed.state() {
FiberState::Failed { error } => {
if !error.as_deref().unwrap_or("").contains("factory exploded") {
return Err("Failed state must carry the factory error".to_string());
}
}
other => return Err(format!("expected Failed{{error}}, got {other:?}")),
}
if ctx.get::<MtProv>().is_some() {
return Err("failed factory must not provide an instance".to_string());
}
match dep.state() {
FiberState::Inactive { .. } => {}
other => {
return Err(format!(
"dependent must rest Inactive while its key has no live provider: {other:?}"
))
}
}
let revive_fid = reg
.register(&ctx, MtRevivePlugin, ())
.map_err(|e| format!("re-register after failure must succeed: {e}"))?;
if revive_fid == fail_fid {
return Err("re-registration must allocate a fresh fiber id".to_string());
}
let revived = reg
.get_fiber(revive_fid)
.ok_or("revived fiber must be tracked")?;
match revived.state() {
FiberState::Active { .. } => {}
other => return Err(format!("revived registration should be Active: {other:?}")),
}
if !matches!(
reg.get_fiber(fail_fid)
.ok_or("failed fiber vanished")?
.state(),
FiberState::Failed { .. }
) {
return Err("failed fiber must remain terminal Failed".to_string());
}
let late_fid = reg
.register(&ctx, MtDependentPlugin, ())
.map_err(|e| format!("dependent registration on revived key failed: {e}"))?;
let late = reg
.get_fiber(late_fid)
.ok_or("late dependent fiber not tracked")?;
match late.state() {
FiberState::Active { .. } if late.epoch() == ":" => {}
other => {
return Err(format!(
"healthy dependent over a present provider must register Active \
with an inject-less epoch: {other:?}"
))
}
}
late.declare_inject::<MtProv>();
match late.state() {
FiberState::Active { .. } => {}
other => {
return Err(format!(
"satisfied declaration must keep the fiber Active eagerly: {other:?}"
))
}
}
if !late.epoch().contains("MtProv") {
return Err(format!(
"declaration must fold the provider into the epoch eagerly: {}",
late.epoch()
));
}
let d3 = ctx
.get::<MtDerived>()
.ok_or("derived value missing after revival")?;
if d3.src != 9 {
return Err(format!(
"projection should observe revived v9, got {}",
d3.src
));
}
let before = revived.epoch();
revived.declare_inject::<MtProv>();
let after = revived.epoch();
if !matches!(revived.state(), FiberState::Active { .. }) {
return Err(format!(
"satisfied declaration must keep the fiber Active: {:?}",
revived.state()
));
}
if after == before {
return Err("satisfied declaration must fold into the epoch immediately".to_string());
}
revived.declare_inject::<MissingProbe>();
match revived.state() {
FiberState::Inactive { error: Some(note) } => {
if !note.contains("missing or inactive dependency") {
return Err(format!("unexpected deactivation note: {note}"));
}
}
other => {
return Err(format!(
"unsatisfied declaration must deactivate the fiber eagerly: {other:?}"
))
}
}
if ctx.get::<MtProv>().is_some() {
return Err("deactivated fiber's projection must be withdrawn".to_string());
}
Ok(())
}
fn next_fid_after(reg: &RegistryService, prev_max: FiberId) -> Result<FiberId, String> {
for fid in (prev_max + 1)..=(prev_max + 64) {
if reg.get_fiber(fid).is_some() {
return Ok(fid);
}
}
Err("no tracked fiber found above the previous max id".to_string())
}
pub async fn lifo_dispose_restores_store() -> Result<(), String> {
let ctx = Context::new_root();
let pre = ctx.snapshot_len();
let log: Arc<Mutex<Vec<&'static str>>> = Arc::new(Mutex::new(Vec::new()));
let count = Arc::new(AtomicUsize::new(0));
ctx.provide(MtMark1(1));
ctx.provide(MtMark2(2));
ctx.plugin(MtEffA {
log: log.clone(),
count: count.clone(),
})
.await
.map_err(|e| format!("plugin A: {e}"))?;
ctx.plugin(MtEffB {
log: log.clone(),
count: count.clone(),
})
.await
.map_err(|e| format!("plugin B: {e}"))?;
if !matches!(ctx.fiber().state(), FiberState::Active { .. }) {
return Err(format!(
"root fiber should be Active after plugins: {:?}",
ctx.fiber().state()
));
}
let _ = ctx.fiber().dispose().await;
if count.load(Ordering::SeqCst) != 2 {
return Err(format!(
"each disposable must run exactly once, counter={}",
count.load(Ordering::SeqCst)
));
}
let order = log.lock().clone();
if order != ["B", "A"] {
return Err(format!(
"LIFO order violated: {order:?} (expected [\"B\", \"A\"])"
));
}
if ctx.get::<MtMark1>().is_some()
|| ctx.get::<MtMark2>().is_some()
|| ctx.get::<MtEffA>().is_some()
|| ctx.get::<MtEffB>().is_some()
{
return Err("store not fully restored after disposal".to_string());
}
if ctx.snapshot_len() != pre {
return Err(format!(
"store length {} after disposal, expected {pre}",
ctx.snapshot_len()
));
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn metatheory_quiescence_after_every_op() {
quiescence_after_every_op()
.await
.expect("Thm 66 quiescence property");
}
#[tokio::test]
async fn metatheory_order_confluence_of_registrations() {
order_confluence_of_registrations()
.await
.expect("Cor 21 / Thm 73 confluence property");
}
#[tokio::test]
async fn metatheory_dependent_never_active_without_provider() {
dependent_never_active_without_provider()
.await
.expect("spatial reactive invariant");
}
#[tokio::test]
async fn metatheory_lifo_dispose_restores_store() {
lifo_dispose_restores_store()
.await
.expect("Thm 16 LIFO disposal property");
}
}
#[cfg(test)]
mod version_conformance {
use super::*;
#[derive(Debug)]
struct VcPeer(pub u64);
impl Service for VcPeer {}
#[derive(Debug)]
struct VcDependent {
peer: Option<u64>,
}
impl Service for VcDependent {
fn check(&self) -> bool {
self.peer.is_some()
}
}
struct VcConsumerPlugin;
impl Plugin for VcConsumerPlugin {
type Config = ();
type Provides = VcDependent;
fn apply(&self, ctx: &Arc<Context>, _config: ()) -> Result<Arc<VcDependent>, CordisError> {
Ok(Arc::new(VcDependent {
peer: ctx.get::<VcPeer>().map(|p| p.0),
}))
}
}
fn requirement(major: u64, floor: u64) -> u64 {
major * crate::Context::VERSION_MAJOR_SCALE + floor
}
#[tokio::test]
async fn versioned_provide_satisfies_compatible_constraint() {
let (ctx, reg, _reflect) = base_root();
ctx.provide_versioned(VcPeer(100_001), 100_001);
let fid = reg
.register(&ctx, VcConsumerPlugin, ())
.expect("consumer registration");
let fiber = reg.get_fiber(fid).expect("tracked consumer");
fiber.declare_inject_versioned::<VcPeer>(Some(requirement(1, 1)));
fiber.refresh(&ctx).await;
match fiber.state() {
FiberState::Active { .. } => {}
other => panic!("compatible provider must activate the consumer, got {other:?}"),
}
let dep = ctx.get::<VcDependent>().expect("projection provided");
assert_eq!(dep.peer, Some(100_001));
}
#[tokio::test]
async fn version_mismatch_keeps_dependent_inactive_until_compatible_upgrade() {
let (ctx, reg, reflect) = base_root();
ctx.provide_versioned(VcPeer(200_001), 200_001);
let fid = reg
.register(&ctx, VcConsumerPlugin, ())
.expect("consumer registration");
let fiber = reg.get_fiber(fid).expect("tracked consumer");
fiber.declare_inject_versioned::<VcPeer>(Some(requirement(1, 1)));
fiber.refresh(&ctx).await;
assert!(matches!(fiber.state(), FiberState::Inactive { .. }));
assert!(ctx.get::<VcDependent>().is_none());
let old = ctx
.remove::<VcPeer>()
.expect("withdrawal allowed while consumer is Inactive")
.expect("provider present");
assert_eq!(old.0, 200_001);
ctx.provide_versioned(VcPeer(100_005), 100_005);
reflect.notify_with_ctx(TypeId::of::<VcPeer>(), &ctx).await;
drain_spawned().await;
match fiber.state() {
FiberState::Active { .. } => {}
other => panic!("compatible upgrade must reactivate the consumer, got {other:?}"),
}
let dep = ctx.get::<VcDependent>().expect("projection re-provided");
assert_eq!(dep.peer, Some(100_005));
}
#[tokio::test]
async fn legacy_provide_defaults_zero_and_satisfies_unconstrained_only() {
let (ctx, reg, _reflect) = base_root();
ctx.provide(VcPeer(42));
let fid = reg
.register(&ctx, VcConsumerPlugin, ())
.expect("consumer registration");
let fiber = reg.get_fiber(fid).expect("tracked consumer");
fiber.declare_inject_versioned::<VcPeer>(None);
fiber.refresh(&ctx).await;
assert!(matches!(fiber.state(), FiberState::Active { .. }));
assert_eq!(ctx.get::<VcDependent>().unwrap().peer, Some(42));
fiber.declare_inject_versioned::<VcPeer>(Some(requirement(1, 0)));
fiber.refresh(&ctx).await;
assert!(matches!(fiber.state(), FiberState::Inactive { .. }));
assert_eq!(ctx.provider_version(TypeId::of::<VcPeer>()), 0);
}
}