use crate::{broadcast, cache, stats, track};
use kio::Pollable;
use std::{
collections::{BTreeMap, HashMap},
fmt,
sync::Arc,
sync::atomic::{AtomicU64, Ordering},
task::{Poll, ready},
time::Duration,
};
use rand::RngExt;
use web_async::Lock;
use super::{Requests, WeakCache};
use crate::{
AsPath, Error, Path, PathOwned, PathPrefixes,
coding::{BoundsExceeded, Decode, DecodeError, Encode, EncodeError},
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct Origin {
id: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct InvalidOrigin;
impl fmt::Display for InvalidOrigin {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "local origin id must be non-zero and below 2^62")
}
}
impl std::error::Error for InvalidOrigin {}
impl Origin {
pub(crate) const UNKNOWN: Self = Self { id: 0 };
pub fn new(id: u64) -> Result<Self, InvalidOrigin> {
if id == 0 || id >= 1u64 << 62 {
return Err(InvalidOrigin);
}
Ok(Self { id })
}
pub fn random() -> Self {
let mut rng = rand::rng();
let id = rng.random_range(1..(1u64 << 53));
Self { id }
}
pub fn id(self) -> u64 {
self.id
}
pub fn produce(self) -> Producer {
Info::new(self).produce()
}
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub struct Info {
pub id: Origin,
pub pool: cache::Pool,
pub linger: Duration,
}
impl Default for Info {
fn default() -> Self {
Self {
id: Origin::UNKNOWN,
pool: cache::Pool::default(),
linger: Duration::ZERO,
}
}
}
impl Info {
pub fn new(id: Origin) -> Self {
Self { id, ..Self::default() }
}
pub fn with_pool(mut self, pool: cache::Pool) -> Self {
self.pool = pool;
self
}
pub fn with_linger(mut self, linger: Duration) -> Self {
self.linger = linger;
self
}
pub fn produce(self) -> Producer {
Producer::new(self)
}
}
impl TryFrom<u64> for Origin {
type Error = InvalidOrigin;
fn try_from(id: u64) -> Result<Self, Self::Error> {
Self::new(id)
}
}
impl fmt::Display for Origin {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.id.fmt(f)
}
}
impl<V: Copy> Encode<V> for Origin
where
u64: Encode<V>,
{
fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
self.id.encode(w, version)
}
}
impl<V: Copy> Decode<V> for Origin
where
u64: Decode<V>,
{
fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
let id = u64::decode(r, version)?;
if id >= 1u64 << 62 {
return Err(DecodeError::InvalidValue);
}
Ok(Self { id })
}
}
pub(crate) const MAX_HOPS: usize = 32;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct OriginList(Vec<Origin>);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct TooManyOrigins;
impl fmt::Display for TooManyOrigins {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "too many origins (max {MAX_HOPS})")
}
}
impl std::error::Error for TooManyOrigins {}
impl From<TooManyOrigins> for DecodeError {
fn from(_: TooManyOrigins) -> Self {
DecodeError::BoundsExceeded
}
}
impl OriginList {
pub fn new() -> Self {
Self(Vec::new())
}
pub fn push(&mut self, origin: Origin) -> Result<(), TooManyOrigins> {
if self.0.len() >= MAX_HOPS {
return Err(TooManyOrigins);
}
self.0.push(origin);
Ok(())
}
pub fn replace_first(&mut self, target: Origin, replacement: Origin) -> bool {
for entry in &mut self.0 {
if *entry == target {
*entry = replacement;
return true;
}
}
false
}
pub fn contains(&self, origin: &Origin) -> bool {
self.0.contains(origin)
}
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
pub fn iter(&self) -> std::slice::Iter<'_, Origin> {
self.0.iter()
}
pub fn as_slice(&self) -> &[Origin] {
&self.0
}
}
impl TryFrom<Vec<Origin>> for OriginList {
type Error = TooManyOrigins;
fn try_from(v: Vec<Origin>) -> Result<Self, Self::Error> {
if v.len() > MAX_HOPS {
return Err(TooManyOrigins);
}
Ok(Self(v))
}
}
impl<'a> IntoIterator for &'a OriginList {
type Item = &'a Origin;
type IntoIter = std::slice::Iter<'a, Origin>;
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
impl<V: Copy> Encode<V> for OriginList
where
u64: Encode<V>,
Origin: Encode<V>,
{
fn encode<W: bytes::BufMut>(&self, w: &mut W, version: V) -> Result<(), EncodeError> {
(self.0.len() as u64).encode(w, version)?;
for origin in &self.0 {
origin.encode(w, version)?;
}
Ok(())
}
}
impl<V: Copy> Decode<V> for OriginList
where
u64: Decode<V>,
Origin: Decode<V>,
{
fn decode<R: bytes::Buf>(r: &mut R, version: V) -> Result<Self, DecodeError> {
let count = u64::decode(r, version)? as usize;
if count > MAX_HOPS {
return Err(DecodeError::BoundsExceeded);
}
let mut list = Vec::with_capacity(count);
for _ in 0..count {
list.push(Origin::decode(r, version)?);
}
Ok(Self(list))
}
}
static NEXT_CONSUMER_ID: AtomicU64 = AtomicU64::new(0);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
struct ConsumerId(u64);
impl ConsumerId {
fn new() -> Self {
Self(NEXT_CONSUMER_ID.fetch_add(1, Ordering::Relaxed))
}
}
struct OriginBroadcast {
path: PathOwned,
broadcast: broadcast::Producer,
state: kio::Producer<FrontState>,
announced: bool,
}
fn route_key(name: &Path, hops: &OriginList) -> (usize, u64) {
(hops.len(), fnv_key(name, hops.iter().copied()))
}
fn fnv_key(name: &Path, origins: impl IntoIterator<Item = Origin>) -> u64 {
const SEED: u64 = 0x420C0DECB00B; const FNV_PRIME: u64 = 0x0000_0100_0000_01b3;
let mut hash = SEED;
for &byte in name.as_str().as_bytes() {
hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
}
for origin in origins {
for &byte in &origin.id().to_le_bytes() {
hash = (hash ^ u64::from(byte)).wrapping_mul(FNV_PRIME);
}
}
hash
}
fn route_order(name: &Path, route: &broadcast::Route) -> (bool, u64, usize, u64) {
let (len, hash) = route_key(name, &route.hops);
(!route.announce, route.cost, len, hash)
}
enum PendingUpdate {
Announce(broadcast::Consumer),
Unannounce,
UnannounceAnnounce(broadcast::Consumer),
}
#[derive(Default)]
struct OriginConsumerState {
pending: BTreeMap<PathOwned, PendingUpdate>,
}
impl OriginConsumerState {
fn apply_announce(&mut self, path: PathOwned, broadcast: broadcast::Consumer) {
let new = match self.pending.remove(&path) {
None | Some(PendingUpdate::Announce(_)) => PendingUpdate::Announce(broadcast),
Some(PendingUpdate::Unannounce | PendingUpdate::UnannounceAnnounce(_)) => {
PendingUpdate::UnannounceAnnounce(broadcast)
}
};
self.pending.insert(path, new);
}
fn apply_unannounce(&mut self, path: PathOwned) {
match self.pending.remove(&path) {
Some(PendingUpdate::Announce(_)) => {}
None | Some(PendingUpdate::Unannounce) => {
self.pending.insert(path, PendingUpdate::Unannounce);
}
Some(PendingUpdate::UnannounceAnnounce(_)) => {
self.pending.insert(path, PendingUpdate::Unannounce);
}
}
}
fn take(&mut self) -> Option<OriginAnnounce> {
let path = self.pending.keys().next()?.clone();
let broadcast = match self.pending.remove(&path).unwrap() {
PendingUpdate::Announce(broadcast) => Some(broadcast),
PendingUpdate::Unannounce => None,
PendingUpdate::UnannounceAnnounce(broadcast) => {
self.pending.insert(path.clone(), PendingUpdate::Announce(broadcast));
None
}
};
Some(OriginAnnounce { path, broadcast })
}
}
#[derive(Clone)]
struct AnnounceConsumerNotify {
root: PathOwned,
state: kio::Producer<OriginConsumerState>,
}
impl AnnounceConsumerNotify {
fn announce(&self, path: impl AsPath, broadcast: broadcast::Consumer) {
let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
self.state
.write()
.ok()
.expect("consumer closed")
.apply_announce(path, broadcast);
}
fn unannounce(&self, path: impl AsPath) {
let path = path.as_path().strip_prefix(&self.root).unwrap().to_owned();
self.state.write().ok().expect("consumer closed").apply_unannounce(path);
}
}
struct NotifyNode {
parent: Option<Lock<NotifyNode>>,
consumers: HashMap<ConsumerId, AnnounceConsumerNotify>,
}
impl NotifyNode {
fn new(parent: Option<Lock<NotifyNode>>) -> Self {
Self {
parent,
consumers: HashMap::new(),
}
}
fn announce(&mut self, path: impl AsPath, broadcast: &broadcast::Consumer) {
for consumer in self.consumers.values() {
consumer.announce(path.as_path(), broadcast.clone());
}
if let Some(parent) = &self.parent {
parent.lock().announce(path, broadcast);
}
}
fn unannounce(&mut self, path: impl AsPath) {
for consumer in self.consumers.values() {
consumer.unannounce(path.as_path());
}
if let Some(parent) = &self.parent {
parent.lock().unannounce(path);
}
}
}
struct OriginNode {
broadcast: Option<OriginBroadcast>,
nested: HashMap<String, Lock<OriginNode>>,
notify: Lock<NotifyNode>,
}
impl OriginNode {
fn new(parent: Option<Lock<NotifyNode>>) -> Self {
Self {
broadcast: None,
nested: HashMap::new(),
notify: Lock::new(NotifyNode::new(parent)),
}
}
fn leaf(&mut self, path: &Path) -> Lock<OriginNode> {
let (dir, rest) = path.next_part().expect("leaf called with empty path");
let next = self.entry(dir);
if rest.is_empty() { next } else { next.lock().leaf(&rest) }
}
fn entry(&mut self, dir: &str) -> Lock<OriginNode> {
match self.nested.get(dir) {
Some(next) => next.clone(),
None => {
let next = Lock::new(OriginNode::new(Some(self.notify.clone())));
self.nested.insert(dir.to_string(), next.clone());
next
}
}
}
fn set_announced(&mut self, expect: &kio::Producer<FrontState>, announce: bool) {
let Some(existing) = &mut self.broadcast else { return };
if !existing.state.same_channel(expect) || existing.announced == announce {
return;
}
existing.announced = announce;
let path = existing.path.clone();
let consumer = existing.broadcast.consume();
let mut notify = self.notify.lock();
if announce {
notify.announce(&path, &consumer);
} else {
notify.unannounce(&path);
}
}
fn consume(&mut self, id: ConsumerId, mut notify: AnnounceConsumerNotify) {
self.consume_initial(&mut notify);
self.notify.lock().consumers.insert(id, notify);
}
fn consume_initial(&mut self, notify: &mut AnnounceConsumerNotify) {
if let Some(broadcast) = &self.broadcast
&& broadcast.announced
{
notify.announce(&broadcast.path, broadcast.broadcast.consume());
}
for nested in self.nested.values() {
nested.lock().consume_initial(notify);
}
}
fn consume_broadcast(&self, rest: impl AsPath) -> Option<broadcast::Consumer> {
let rest = rest.as_path();
if let Some((dir, rest)) = rest.next_part() {
let node = self.nested.get(dir)?.lock();
node.consume_broadcast(&rest)
} else {
self.broadcast.as_ref().map(|b| b.broadcast.consume())
}
}
fn unconsume(&mut self, id: ConsumerId) {
self.notify.lock().consumers.remove(&id).expect("consumer not found");
if self.is_empty() {
}
}
fn remove(&mut self, expect: &kio::Producer<FrontState>, relative: impl AsPath) {
let relative = relative.as_path();
if let Some((dir, relative)) = relative.next_part() {
let Some(nested) = self.nested.get(dir) else { return };
let nested = nested.clone();
let mut locked = nested.lock();
locked.remove(expect, &relative);
if locked.is_empty() {
drop(locked);
self.nested.remove(dir);
}
} else if let Some(existing) = &self.broadcast
&& existing.state.same_channel(expect)
{
let existing = self.broadcast.take().expect("checked above");
if existing.announced {
self.notify.lock().unannounce(&existing.path);
}
}
}
fn is_empty(&self) -> bool {
self.broadcast.is_none() && self.nested.is_empty() && self.notify.lock().consumers.is_empty()
}
}
#[derive(Clone)]
struct OriginNodes {
nodes: Vec<(PathOwned, Lock<OriginNode>)>,
}
impl OriginNodes {
pub fn select(&self, prefixes: &PathPrefixes) -> Option<Self> {
let mut roots = Vec::new();
for (root, state) in &self.nodes {
for prefix in prefixes {
if root.has_prefix(prefix) {
roots.push((root.to_owned(), state.clone()));
continue;
}
if let Some(suffix) = prefix.strip_prefix(root) {
let nested = state.lock().leaf(&suffix);
roots.push((prefix.to_owned(), nested));
}
}
}
if roots.is_empty() {
None
} else {
Some(Self { nodes: roots })
}
}
pub fn root(&self, new_root: impl AsPath) -> Option<Self> {
let new_root = new_root.as_path();
let mut roots = Vec::new();
if new_root.is_empty() {
return Some(self.clone());
}
for (root, state) in &self.nodes {
if let Some(suffix) = root.strip_prefix(&new_root) {
roots.push((suffix.to_owned(), state.clone()));
} else if let Some(suffix) = new_root.strip_prefix(root) {
let nested = state.lock().leaf(&suffix);
roots.push(("".into(), nested));
}
}
if roots.is_empty() {
None
} else {
Some(Self { nodes: roots })
}
}
pub fn get(&self, path: impl AsPath) -> Option<(Lock<OriginNode>, PathOwned)> {
let path = path.as_path();
for (root, state) in &self.nodes {
if let Some(suffix) = path.strip_prefix(root) {
return Some((state.clone(), suffix.to_owned()));
}
}
None
}
}
impl Default for OriginNodes {
fn default() -> Self {
Self {
nodes: vec![("".into(), Lock::new(OriginNode::new(None)))],
}
}
}
#[derive(Clone)]
pub struct OriginAnnounce {
pub path: PathOwned,
pub broadcast: Option<broadcast::Consumer>,
}
#[derive(Clone)]
pub struct Producer {
info: Origin,
nodes: OriginNodes,
root: PathOwned,
dynamic: kio::Shared<OriginDynamicState>,
pool: cache::Pool,
linger: Duration,
stats: stats::Session,
}
impl std::ops::Deref for Producer {
type Target = Origin;
fn deref(&self) -> &Self::Target {
&self.info
}
}
impl Producer {
pub fn new(info: Info) -> Self {
Self {
info: info.id,
nodes: OriginNodes::default(),
root: PathOwned::default(),
dynamic: kio::Shared::default(),
pool: info.pool,
linger: info.linger,
stats: stats::Session::default(),
}
}
pub fn with_stats(mut self, session: stats::Session) -> Self {
self.stats = session;
self
}
pub fn with_linger(mut self, linger: Duration) -> Self {
self.linger = linger;
self
}
pub fn info(&self) -> Info {
Info {
id: self.info,
pool: self.pool.clone(),
linger: self.linger,
}
}
pub(crate) fn empty(info: Origin) -> Self {
Self {
info,
nodes: OriginNodes { nodes: Vec::new() },
root: PathOwned::default(),
dynamic: kio::Shared::default(),
pool: cache::Pool::default(),
linger: Duration::ZERO,
stats: stats::Session::default(),
}
}
pub fn create_broadcast(&self, path: impl AsPath, route: broadcast::Route) -> Result<broadcast::Producer, Error> {
let path = path.as_path();
debug_assert!(
!route.hops.contains(&self.info),
"create_broadcast called with a looping hop chain",
);
let (node, rest) = self.nodes.get(&path).ok_or(Error::Unauthorized)?;
let full = self.root.join(&path).to_owned();
if full.parts().count() > Path::MAX_PARTS {
return Err(BoundsExceeded.into());
}
let ingress = self.stats.ingress(&full);
let mut source = broadcast::Info { origin: self.info() }
.produce()
.with_stats(ingress.clone());
source.set_route(route).expect("fresh producer");
web_async::spawn(run_source(self.info(), node, full, rest, source.consume(), ingress));
Ok(source)
}
pub fn scope(&self, prefixes: &[Path]) -> Option<Producer> {
let prefixes = PathPrefixes::new(prefixes);
Some(Producer {
info: self.info,
nodes: self.nodes.select(&prefixes)?,
root: self.root.clone(),
dynamic: self.dynamic.clone(),
pool: self.pool.clone(),
linger: self.linger,
stats: self.stats.clone(),
})
}
pub fn dynamic(&self) -> Dynamic {
Dynamic::new(self.info, self.root.clone(), self.dynamic.clone())
}
pub fn consume(&self) -> Consumer {
Consumer::new(
self.info,
self.root.clone(),
self.nodes.clone(),
self.dynamic.clone(),
stats::Session::default(),
)
}
pub fn announces(&self) -> AnnounceProducer {
AnnounceProducer::new(self.root.clone(), self.nodes.clone())
}
pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
let prefix = prefix.as_path();
Some(Self {
info: self.info,
root: self.root.join(&prefix).to_owned(),
nodes: self.nodes.root(&prefix)?,
dynamic: self.dynamic.clone(),
pool: self.pool.clone(),
linger: self.linger,
stats: self.stats.clone(),
})
}
pub fn root(&self) -> &Path<'_> {
&self.root
}
pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
self.nodes.nodes.iter().map(|(root, _)| root)
}
pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
self.root.join(path)
}
}
const MAX_TRACK_RETRIES: u32 = 3;
struct FrontRoute {
id: u64,
route: broadcast::Route,
source: broadcast::Consumer,
}
struct FrontState {
path: PathOwned,
self_origin: Origin,
next_route: u64,
routes: Vec<FrontRoute>,
active: Option<u64>,
linger: Duration,
closed: bool,
}
impl FrontState {
fn best_route(&self) -> Option<u64> {
self.routes
.iter()
.min_by_key(|r| route_order(&self.path.as_path(), &r.route))
.map(|r| r.id)
}
fn reselect(&mut self, carrying: bool) {
let best = self.best_route();
if carrying
&& let (Some(best_id), Some(cur_id)) = (best, self.active)
&& best_id != cur_id
&& let Some(candidate) = self.routes.iter().find(|r| r.id == best_id)
&& let Some(incumbent) = self.routes.iter().find(|r| r.id == cur_id)
&& incumbent.route.announce
&& candidate.route.cost < incumbent.route.cost
&& candidate.route.advertised == 0
&& candidate.route.hops.len() >= 2
&& !self.handover_allowed(&candidate.route)
{
return;
}
self.active = best;
}
fn handover_allowed(&self, route: &broadcast::Route) -> bool {
let name = self.path.as_path();
match route.hops.iter().last() {
Some(peer) => fnv_key(&name, [*peer]) < fnv_key(&name, [self.self_origin]),
None => true,
}
}
fn active_route(&self) -> Option<broadcast::Route> {
let id = self.active?;
self.routes.iter().find(|r| r.id == id).map(|r| r.route.clone())
}
}
fn sync_front(state: &kio::Producer<FrontState>, broadcast: &broadcast::Producer, leaf: &Lock<OriginNode>) {
let mut leaf_guard = leaf.lock();
let advert = state.read().active_route();
if let Some(advert) = advert {
let announce = advert.announce;
let _ = broadcast.clone().set_route(advert);
leaf_guard.set_announced(state, announce);
}
}
fn detach_source(
state: &kio::Producer<FrontState>,
broadcast: &broadcast::Producer,
leaf: &Lock<OriginNode>,
id: u64,
graceful: bool,
) {
let close = {
let carrying = broadcast.demand().is_used();
let Ok(mut s) = state.write() else { return };
let Some(pos) = s.routes.iter().position(|r| r.id == id) else {
return;
};
s.routes.remove(pos);
s.reselect(carrying);
if s.routes.is_empty() && !s.closed && (graceful || s.linger.is_zero()) {
s.closed = true;
true
} else {
false
}
};
if close {
broadcast.abort_spliced(Error::Dropped);
}
sync_front(state, broadcast, leaf);
}
async fn run_source(
origin: Info,
node: Lock<OriginNode>,
full: PathOwned,
rest: PathOwned,
mut source: broadcast::Consumer,
ingress: stats::Scope,
) {
let Ok(route) = source.route_changed().await else {
return;
};
let mut announce = route.announce.then(|| ingress.announce());
let leaf = if rest.is_empty() {
node.clone()
} else {
node.lock().leaf(&rest)
};
let (state, broadcast, id) = attach_source(&origin, &node, &leaf, &full, &rest, &source, route);
loop {
match source.route_changed().await {
Ok(route) => {
let announced = route.announce;
{
let carrying = broadcast.demand().is_used();
let Ok(mut s) = state.write() else { return };
let Some(entry) = s.routes.iter_mut().find(|r| r.id == id) else {
return;
};
if entry.route == route {
continue;
}
entry.route = route;
s.reselect(carrying);
}
match (announced, announce.is_some()) {
(true, false) => announce = Some(ingress.announce()),
(false, true) => announce = None,
_ => {}
}
sync_front(&state, &broadcast, &leaf);
}
Err(_) => {
detach_source(&state, &broadcast, &leaf, id, source.is_finished());
return;
}
}
}
}
fn attach_source(
origin: &Info,
node: &Lock<OriginNode>,
leaf: &Lock<OriginNode>,
full: &PathOwned,
rest: &PathOwned,
source: &broadcast::Consumer,
route: broadcast::Route,
) -> (kio::Producer<FrontState>, broadcast::Producer, u64) {
let mut leaf_guard = leaf.lock();
if let Some(existing) = &leaf_guard.broadcast {
let mut joined = None;
let carrying = existing.broadcast.demand().is_used();
if let Ok(mut s) = existing.state.write()
&& !s.closed
{
let id = s.next_route;
s.next_route += 1;
s.routes.push(FrontRoute {
id,
route: route.clone(),
source: source.clone(),
});
s.reselect(carrying);
joined = Some(id);
}
if let Some(id) = joined {
let state = existing.state.clone();
let broadcast = existing.broadcast.clone();
drop(leaf_guard);
sync_front(&state, &broadcast, leaf);
return (state, broadcast, id);
}
}
let announce = route.announce;
let broadcast = broadcast::Producer::new_spliced(broadcast::Info { origin: origin.clone() });
let _ = broadcast.clone().set_route(route.clone());
let state = kio::Producer::new(FrontState {
path: full.clone(),
self_origin: origin.id,
next_route: 1,
routes: vec![FrontRoute {
id: 0,
route,
source: source.clone(),
}],
active: Some(0),
linger: origin.linger,
closed: false,
});
if let Some(stale) = leaf_guard.broadcast.take()
&& stale.announced
{
leaf_guard.notify.lock().unannounce(&stale.path);
}
let entry = OriginBroadcast {
path: full.clone(),
broadcast: broadcast.clone(),
state: state.clone(),
announced: announce,
};
if entry.announced {
leaf_guard.notify.lock().announce(full, &broadcast.consume());
}
leaf_guard.broadcast = Some(entry);
drop(leaf_guard);
web_async::spawn(run_front(state.clone(), broadcast.clone(), node.clone(), rest.clone()));
(state, broadcast, 0)
}
async fn run_front(
state: kio::Producer<FrontState>,
mut broadcast: broadcast::Producer,
node: Lock<OriginNode>,
rest: PathOwned,
) {
enum Step {
Serve(Arc<str>, super::resume::Producer),
Changed,
Expired,
Closed,
}
let linger = state.read().linger;
let mut deadline: Option<web_async::time::Instant> = None;
loop {
let empty = {
let s = state.read();
!s.closed && s.routes.is_empty()
};
deadline = match (empty, deadline) {
(true, None) => web_async::time::Instant::now().checked_add(linger),
(true, at) => at,
(false, _) => None,
};
let step = {
let mut sleep = std::pin::pin!(async {
match deadline {
Some(at) => {
web_async::time::sleep(at.saturating_duration_since(web_async::time::Instant::now())).await
}
None => std::future::pending().await,
}
});
let mut fired = false;
kio::wait(|waiter| {
if let Poll::Ready((name, resume)) = broadcast.poll_spliced_assigned(waiter) {
return Poll::Ready(Step::Serve(name, resume));
}
match state.poll(waiter, |s| {
if s.closed || s.routes.is_empty() != empty {
Poll::Ready(())
} else {
Poll::Pending
}
}) {
Poll::Ready(Ok(guard)) => {
return Poll::Ready(if guard.closed { Step::Closed } else { Step::Changed });
}
Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
Poll::Pending => {}
}
if deadline.is_some() && !fired && waiter.poll_future(sleep.as_mut()).is_ready() {
fired = true;
}
match fired {
true => Poll::Ready(Step::Expired),
false => Poll::Pending,
}
})
.await
};
match step {
Step::Serve(name, resume) => {
web_async::spawn(serve_track(state.clone(), name, resume));
}
Step::Changed => {}
Step::Expired => {
let close = {
let Ok(mut s) = state.write() else { break };
if !s.closed && s.routes.is_empty() {
s.closed = true;
true
} else {
false
}
};
if close {
break;
}
}
Step::Closed => break,
}
}
broadcast.abort_spliced(Error::Dropped);
broadcast.finish();
node.lock().remove(&state, &rest);
}
async fn serve_track(state: kio::Producer<FrontState>, name: Arc<str>, mut resume: super::resume::Producer) {
enum Step {
Closed,
Splice(u64, broadcast::Consumer),
Complete,
Failed,
}
let mut fails = 0u32;
let mut serving: Option<(u64, track::Consumer)> = None;
let mut dead: Option<u64> = None;
loop {
let serving_id = serving.as_ref().map(|(id, _)| *id);
let step = kio::wait(|waiter| {
match state.poll(waiter, |s| {
if s.closed || matches!(s.active, Some(active) if Some(active) != serving_id && Some(active) != dead) {
Poll::Ready(())
} else {
Poll::Pending
}
}) {
Poll::Ready(Ok(guard)) => {
if guard.closed {
return Poll::Ready(Step::Closed);
}
let active = guard.active.expect("predicate guaranteed an active source");
let source = guard
.routes
.iter()
.find(|r| r.id == active)
.expect("active source in table")
.source
.clone();
return Poll::Ready(Step::Splice(active, source));
}
Poll::Ready(Err(_)) => return Poll::Ready(Step::Closed),
Poll::Pending => {}
}
if let Some((_, track)) = &serving
&& let Poll::Ready(result) = track.poll_complete(waiter)
{
return Poll::Ready(match result {
Ok(()) => Step::Complete,
Err(_) => Step::Failed,
});
}
Poll::Pending
})
.await;
match step {
Step::Closed => return,
Step::Complete => {
let _ = resume.finish();
return;
}
Step::Failed => {
serving = None;
}
Step::Splice(id, source) => {
let attempt = match source.track(&name) {
Ok(track) => {
let query = track.info().into_inner();
let info = kio::wait(|waiter| {
if let Poll::Ready(result) = query.poll(waiter) {
return Poll::Ready(Some(result));
}
match state.poll(waiter, |s| {
if s.closed || s.active != Some(id) {
Poll::Ready(())
} else {
Poll::Pending
}
}) {
Poll::Ready(_) => Poll::Ready(None),
Poll::Pending => Poll::Pending,
}
})
.await;
match info {
None => continue,
Some(Ok(_)) => match track.poll_complete(&kio::Waiter::noop()) {
Poll::Ready(Err(err)) => Err(err),
_ => Ok(track),
},
Some(Err(err)) => Err(err),
}
}
Err(err) => Err(err),
};
match attempt {
Ok(track) => {
if resume.takeover(&track).is_err() {
return;
}
fails = 0;
dead = None;
serving = Some((id, track));
}
Err(_) if source.is_closing() => {
dead = Some(id);
serving = None;
}
Err(err) => {
fails += 1;
if fails >= MAX_TRACK_RETRIES {
tracing::debug!(name = %name, %err, "aborting unservable track");
let _ = resume.abort(Error::Unroutable);
return;
}
serving = None;
}
}
}
}
}
}
#[derive(Default)]
struct OriginDynamicState {
requests: Requests<PathOwned, kio::Producer<PendingBroadcast>>,
served: WeakCache<PathOwned, broadcast::WeakConsumer>,
}
#[derive(Default)]
struct PendingBroadcast {
resolved: Option<Result<broadcast::Consumer, Error>>,
}
pub struct Dynamic {
info: Origin,
root: PathOwned,
state: kio::Shared<OriginDynamicState>,
}
impl Clone for Dynamic {
fn clone(&self) -> Self {
self.state.lock().requests.add_handler();
Self {
info: self.info,
root: self.root.clone(),
state: self.state.clone(),
}
}
}
impl Dynamic {
fn new(info: Origin, root: PathOwned, state: kio::Shared<OriginDynamicState>) -> Self {
state.lock().requests.add_handler();
Self { info, root, state }
}
pub fn info(&self) -> &Origin {
&self.info
}
pub fn poll_requested_broadcast(&mut self, waiter: &kio::Waiter) -> Poll<Result<Request, Error>> {
let mut state = ready!(self.state.poll(waiter, |state| {
if state.requests.has_queued() {
Poll::Ready(())
} else {
Poll::Pending
}
}));
let path = state.requests.pop().expect("predicate guaranteed a request");
let producer = state.requests.get(&path).expect("popped key must be pending").clone();
Poll::Ready(Ok(Request {
path,
producer,
state: self.state.clone(),
}))
}
pub async fn requested_broadcast(&mut self) -> Result<Request, Error> {
kio::wait(|waiter| self.poll_requested_broadcast(waiter)).await
}
pub fn root(&self) -> &Path<'_> {
&self.root
}
}
impl Drop for Dynamic {
fn drop(&mut self) {
let mut state = self.state.lock();
if state.requests.remove_handler() {
state.requests.drain_queued();
}
}
}
pub struct Request {
path: PathOwned,
producer: kio::Producer<PendingBroadcast>,
state: kio::Shared<OriginDynamicState>,
}
impl Request {
pub fn path(&self) -> &Path<'_> {
&self.path
}
pub fn accept(self, broadcast: impl Consume<broadcast::Consumer>) {
let broadcast = broadcast.consume();
let resolved = {
let mut state = self.state.lock();
let existing = state.served.insert(self.path.clone(), broadcast.weak());
state
.requests
.remove_if(&self.path, |producer| producer.same_channel(&self.producer));
existing.map(|weak| weak.consume()).unwrap_or(broadcast)
};
if let Ok(mut pending) = self.producer.write() {
pending.resolved = Some(Ok(resolved));
}
}
pub fn reject(self, err: Error) {
self.state
.lock()
.requests
.remove_if(&self.path, |producer| producer.same_channel(&self.producer));
if let Ok(mut state) = self.producer.write() {
state.resolved = Some(Err(err));
}
}
}
impl Drop for Request {
fn drop(&mut self) {
self.state
.lock()
.requests
.remove_if(&self.path, |producer| producer.same_channel(&self.producer));
}
}
pub struct Requesting {
inner: RequestState,
stats: stats::Scope,
}
enum RequestState {
Ready(broadcast::Consumer),
Failed(Error),
Pending(kio::Consumer<PendingBroadcast>),
}
impl Requesting {
fn ready(broadcast: broadcast::Consumer) -> Self {
Self {
inner: RequestState::Ready(broadcast),
stats: stats::Scope::default(),
}
}
fn failed(error: Error) -> Self {
Self {
inner: RequestState::Failed(error),
stats: stats::Scope::default(),
}
}
fn pending(consumer: kio::Consumer<PendingBroadcast>) -> Self {
Self {
inner: RequestState::Pending(consumer),
stats: stats::Scope::default(),
}
}
fn with_stats(mut self, scope: stats::Scope) -> Self {
self.stats = scope;
self
}
pub fn poll_ok(&self, waiter: &kio::Waiter) -> Poll<Result<broadcast::Consumer, Error>> {
match &self.inner {
RequestState::Ready(broadcast) => Poll::Ready(Ok(broadcast.clone().with_stats(self.stats.clone()))),
RequestState::Failed(error) => Poll::Ready(Err(error.clone())),
RequestState::Pending(consumer) => Poll::Ready(
match ready!(consumer.poll(waiter, |state| match &state.resolved {
Some(result) => Poll::Ready(result.clone()),
None => Poll::Pending,
})) {
Ok(result) => result.map(|broadcast| broadcast.with_stats(self.stats.clone())),
Err(_closed) => Err(Error::Unroutable),
},
),
}
}
}
impl kio::Pollable for Requesting {
type Output = Result<broadcast::Consumer, Error>;
fn poll(&self, waiter: &kio::Waiter) -> Poll<Self::Output> {
self.poll_ok(waiter)
}
}
pub trait Consume<T> {
fn consume(&self) -> T;
}
impl<T, U: Consume<T>> Consume<T> for &U {
fn consume(&self) -> T {
(**self).consume()
}
}
impl Consume<Consumer> for Producer {
fn consume(&self) -> Consumer {
Consumer::new(
self.info,
self.root.clone(),
self.nodes.clone(),
self.dynamic.clone(),
stats::Session::default(),
)
}
}
impl Consume<Consumer> for Consumer {
fn consume(&self) -> Consumer {
self.clone()
}
}
impl Consume<broadcast::Consumer> for broadcast::Producer {
fn consume(&self) -> broadcast::Consumer {
self.consume()
}
}
impl Consume<broadcast::Consumer> for broadcast::Consumer {
fn consume(&self) -> broadcast::Consumer {
self.clone()
}
}
impl Consume<track::Consumer> for track::Producer {
fn consume(&self) -> track::Consumer {
self.consume()
}
}
impl Consume<track::Consumer> for track::Consumer {
fn consume(&self) -> track::Consumer {
self.clone()
}
}
#[derive(Clone)]
pub struct Consumer {
info: Origin,
nodes: OriginNodes,
root: PathOwned,
dynamic: kio::Shared<OriginDynamicState>,
stats: stats::Session,
}
impl std::ops::Deref for Consumer {
type Target = Origin;
fn deref(&self) -> &Self::Target {
&self.info
}
}
impl Consumer {
fn new(
info: Origin,
root: PathOwned,
nodes: OriginNodes,
dynamic: kio::Shared<OriginDynamicState>,
stats: stats::Session,
) -> Self {
Self {
info,
nodes,
root,
dynamic,
stats,
}
}
pub fn with_stats(mut self, session: stats::Session) -> Self {
self.stats = session;
self
}
fn untagged(&self) -> Self {
Self {
stats: stats::Session::default(),
..self.clone()
}
}
pub(crate) fn empty(&self) -> Self {
Self {
info: self.info,
nodes: OriginNodes { nodes: Vec::new() },
root: self.root.clone(),
dynamic: self.dynamic.clone(),
stats: self.stats.clone(),
}
}
pub fn announced(&self) -> AnnounceConsumer {
AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), self.stats.clone())
}
pub fn consume(&self) -> Self {
self.clone()
}
fn get_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
let path = path.as_path();
let (root, rest) = self.nodes.get(&path)?;
let state = root.lock();
state.consume_broadcast(&rest)
}
pub async fn announced_broadcast(&self, path: impl AsPath) -> Option<broadcast::Consumer> {
let path = path.as_path();
let consumer = self.scope(std::slice::from_ref(&path))?;
if !consumer.allowed().any(|allowed| path.has_prefix(allowed)) {
return None;
}
let mut announced = consumer.untagged().announced();
let scope = self.stats.egress(self.root.join(&path).to_owned());
loop {
let OriginAnnounce {
path: announced_path,
broadcast,
} = announced.next().await?;
if announced_path.as_path() == path
&& let Some(broadcast) = broadcast
{
return Some(broadcast.with_stats(scope));
}
}
}
pub fn scope(&self, prefixes: &[Path]) -> Option<Consumer> {
let prefixes = PathPrefixes::new(prefixes);
Some(Consumer::new(
self.info,
self.root.clone(),
self.nodes.select(&prefixes)?,
self.dynamic.clone(),
self.stats.clone(),
))
}
pub fn request_broadcast(&self, path: impl AsPath) -> kio::Pending<Requesting> {
let path = path.as_path();
let absolute = self.root.join(&path).to_owned();
let scope = self.stats.egress(&absolute);
if let Some(broadcast) = self.get_broadcast(&path) {
return kio::Pending::new(Requesting::ready(broadcast).with_stats(scope));
}
let mut state = self.dynamic.lock();
if let Some(weak) = state.served.get(&absolute) {
return kio::Pending::new(Requesting::ready(weak.consume()).with_stats(scope));
}
let consumer = if let Some(producer) = state.requests.join(&absolute) {
producer.consume()
} else {
let producer = kio::Producer::<PendingBroadcast>::default();
let consumer = producer.consume();
if state.requests.insert(absolute, producer).is_err() {
return kio::Pending::new(Requesting::failed(Error::Unroutable));
}
consumer
};
kio::Pending::new(Requesting::pending(consumer).with_stats(scope))
}
pub fn with_root(&self, prefix: impl AsPath) -> Option<Self> {
let prefix = prefix.as_path();
Some(Self::new(
self.info,
self.root.join(&prefix).to_owned(),
self.nodes.root(&prefix)?,
self.dynamic.clone(),
self.stats.clone(),
))
}
pub fn root(&self) -> &Path<'_> {
&self.root
}
pub fn allowed(&self) -> impl Iterator<Item = &Path<'_>> {
self.nodes.nodes.iter().map(|(root, _)| root)
}
pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
self.root.join(path)
}
}
#[derive(Clone)]
pub struct AnnounceProducer {
nodes: OriginNodes,
root: PathOwned,
}
impl AnnounceProducer {
fn new(root: PathOwned, nodes: OriginNodes) -> Self {
Self { nodes, root }
}
pub fn consume(&self) -> AnnounceConsumer {
AnnounceConsumer::new(self.root.clone(), self.nodes.clone(), stats::Session::default())
}
pub fn root(&self) -> &Path<'_> {
&self.root
}
}
pub struct AnnounceConsumer {
id: ConsumerId,
nodes: OriginNodes,
root: PathOwned,
state: kio::Producer<OriginConsumerState>,
stats: stats::Session,
guards: HashMap<PathOwned, stats::Announce>,
}
impl AnnounceConsumer {
fn new(root: PathOwned, nodes: OriginNodes, stats: stats::Session) -> Self {
let state = kio::Producer::<OriginConsumerState>::default();
let id = ConsumerId::new();
for (_, node) in &nodes.nodes {
let notify = AnnounceConsumerNotify {
root: root.clone(),
state: state.clone(),
};
node.lock().consume(id, notify);
}
Self {
id,
nodes,
root,
state,
stats,
guards: HashMap::new(),
}
}
fn attribute(&mut self, update: OriginAnnounce) -> OriginAnnounce {
let OriginAnnounce { path, broadcast } = update;
let absolute = self.root.join(&path).to_owned();
match broadcast {
Some(broadcast) => {
let scope = self.stats.egress(&absolute);
self.guards.entry(absolute).or_insert_with(|| scope.announce());
OriginAnnounce {
path,
broadcast: Some(broadcast.with_stats(scope)),
}
}
None => {
self.guards.remove(&absolute);
OriginAnnounce { path, broadcast: None }
}
}
}
pub async fn next(&mut self) -> Option<OriginAnnounce> {
kio::wait(|waiter| self.poll_next(waiter)).await
}
pub fn poll_next(&mut self, waiter: &kio::Waiter) -> Poll<Option<OriginAnnounce>> {
let update = {
let mut state = match ready!(self.state.poll(waiter, |state| {
if state.pending.is_empty() {
Poll::Pending
} else {
Poll::Ready(())
}
})) {
Ok(state) => state,
Err(_) => return Poll::Ready(None),
};
state.take().expect("predicate guaranteed an update")
};
Poll::Ready(Some(self.attribute(update)))
}
pub fn try_next(&mut self) -> Option<OriginAnnounce> {
let update = self.state.write().ok()?.take()?;
Some(self.attribute(update))
}
pub fn is_closed(&self) -> bool {
self.state.write().is_err()
}
pub fn root(&self) -> &Path<'_> {
&self.root
}
pub fn absolute(&self, path: impl AsPath) -> Path<'_> {
self.root.join(path)
}
}
impl Drop for AnnounceConsumer {
fn drop(&mut self) {
for (_, root) in &self.nodes.nodes {
root.lock().unconsume(self.id);
}
}
}
#[cfg(test)]
use futures::FutureExt;
#[cfg(test)]
#[allow(missing_docs)] impl AnnounceConsumer {
pub fn assert_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
let expected = expected.as_path();
let announce = self.next().now_or_never().expect("next blocked").expect("no next");
assert_eq!(announce.path, expected, "wrong path");
let announced = announce.broadcast.expect("should be an active announce");
assert!(announced.is_clone(broadcast), "should be the same broadcast");
}
pub fn assert_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
let expected = expected.as_path();
let announce = self.next().now_or_never().expect("next blocked").expect("no next");
assert_eq!(announce.path, expected, "wrong path");
announce.broadcast.expect("should be an active announce")
}
pub fn assert_try_next(&mut self, expected: impl AsPath, broadcast: &broadcast::Consumer) {
let expected = expected.as_path();
let announce = self.try_next().expect("no next");
assert_eq!(announce.path, expected, "wrong path");
let announced = announce.broadcast.expect("should be an active announce");
assert!(announced.is_clone(broadcast), "should be the same broadcast");
}
pub fn assert_try_next_some(&mut self, expected: impl AsPath) -> broadcast::Consumer {
let expected = expected.as_path();
let announce = self.try_next().expect("no next");
assert_eq!(announce.path, expected, "wrong path");
announce.broadcast.expect("should be an active announce")
}
pub fn assert_next_none(&mut self, expected: impl AsPath) {
let expected = expected.as_path();
let announce = self.next().now_or_never().expect("next blocked").expect("no next");
assert_eq!(announce.path, expected, "wrong path");
assert!(announce.broadcast.is_none(), "should be unannounced");
}
pub fn assert_next_wait(&mut self) {
if let Some(res) = self.next().now_or_never() {
panic!("next should block: got {:?}", res.map(|a| a.path));
}
}
}
#[cfg(test)]
mod tests {
use crate::coding::Decode;
use crate::group;
use super::*;
fn announce() -> broadcast::Route {
broadcast::Route::new().with_announce(true)
}
fn origin_keyed(name: &str, peer: Origin, above: bool) -> Origin {
let name = Path::new(name);
let peer_key = fnv_key(&name, [peer]);
(100u64..)
.map(|id| Origin::new(id).unwrap())
.find(|origin| (fnv_key(&name, [*origin]) > peer_key) == above)
.unwrap()
}
fn front_state(self_origin: Origin, routes: Vec<broadcast::Route>) -> FrontState {
let source = broadcast::Info::new().produce().consume();
FrontState {
path: Path::new("test").to_owned(),
self_origin,
next_route: routes.len() as u64,
routes: routes
.into_iter()
.enumerate()
.map(|(id, route)| FrontRoute {
id: id as u64,
route,
source: source.clone(),
})
.collect(),
active: Some(0),
linger: Duration::ZERO,
closed: false,
}
}
fn sibling_route(peer: Origin) -> broadcast::Route {
let hops = OriginList::try_from(vec![Origin::new(90).unwrap(), peer]).unwrap();
announce().with_hops(hops)
}
fn upstream_route(cost: u64) -> broadcast::Route {
let hops = OriginList::try_from(vec![Origin::new(90).unwrap()]).unwrap();
announce().with_hops(hops).with_cost(cost)
}
#[test]
fn test_carrying_gate_keys() {
let peer = Origin::new(3).unwrap();
let mut lost = front_state(
origin_keyed("test", peer, false),
vec![upstream_route(10), sibling_route(peer)],
);
lost.reselect(true);
assert_eq!(
lost.active,
Some(0),
"carrying front re-parented onto a higher-keyed peer"
);
lost.reselect(false);
assert_eq!(lost.active, Some(1), "idle front must take the cheaper route");
let mut won = front_state(
origin_keyed("test", peer, true),
vec![upstream_route(10), sibling_route(peer)],
);
won.reselect(true);
assert_eq!(won.active, Some(1), "carrying front must follow a lower-keyed peer");
}
#[test]
fn test_carrying_gate_symmetric_race() {
let a = Origin::new(1).unwrap();
let b = Origin::new(2).unwrap();
let mut a_view = front_state(a, vec![upstream_route(10), sibling_route(b)]);
let mut b_view = front_state(b, vec![upstream_route(10), sibling_route(a)]);
a_view.reselect(true);
b_view.reselect(true);
let a_moved = a_view.active == Some(1);
let b_moved = b_view.active == Some(1);
assert!(
a_moved != b_moved,
"exactly one side must re-parent (a: {a_moved}, b: {b_moved})"
);
}
#[test]
fn test_carrying_switches_to_benign_routes() {
let peer = Origin::new(3).unwrap();
let lost = origin_keyed("test", peer, false);
let mut forwarder = sibling_route(peer).with_cost(4);
forwarder.advertised = 4;
let mut state = front_state(lost, vec![upstream_route(10), forwarder]);
state.reselect(true);
assert_eq!(
state.active,
Some(1),
"a cheaper forwarder path must win while carrying"
);
let direct = announce().with_hops(OriginList::try_from(vec![peer]).unwrap());
let mut state = front_state(lost, vec![upstream_route(10), direct]);
state.reselect(true);
assert_eq!(
state.active,
Some(1),
"a direct publisher route must win while carrying"
);
}
#[test]
fn test_carrying_gate_ignores_unannounced_incumbent() {
let peer = Origin::new(3).unwrap();
let unannounced = upstream_route(10).with_announce(false);
let mut state = front_state(
origin_keyed("test", peer, false),
vec![unannounced, sibling_route(peer)],
);
state.reselect(true);
assert_eq!(
state.active,
Some(1),
"an unannounced incumbent must always be displaced"
);
}
async fn settle() {
tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
}
async fn accept_track(dynamic: &mut broadcast::Dynamic, name: &str) -> track::Producer {
let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
.await
.expect("timed out waiting for a track request")
.expect("source closed");
assert_eq!(request.name(), name, "unexpected track dispatched");
request.accept(None)
}
#[tokio::test]
async fn test_stats_tagged_end_to_end() {
use crate::Timestamp;
use crate::stats::{Config, Registry, Tier};
use bytes::Bytes;
tokio::time::pause();
let registry = Registry::new(Config::new());
let ctx = registry.tier(Tier::default()).session("acme");
let origin = Origin::random().produce();
let ingress = origin.clone().with_stats(ctx.clone());
let egress = origin.consume().with_stats(ctx.clone());
let mut announced = egress.announced();
let source = ingress.create_broadcast("demo", announce()).unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
let update = announced.next().await.unwrap();
assert_eq!(update.path.as_str(), "demo");
let broadcast = update.broadcast.unwrap();
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer = accept_track(&mut dynamic, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
let mut group = producer.append_group().unwrap();
group
.write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
.unwrap();
group
.write_frame(Timestamp::ZERO, Bytes::from_static(b"world"))
.unwrap();
group.finish().unwrap();
let mut group_c = sub.recv_group().await.unwrap().unwrap();
let mut frames = 0;
while let Some(frame) = group_c.read_frame().await.unwrap() {
assert_eq!(frame.payload.len(), 5);
frames += 1;
}
assert_eq!(frames, 2);
settle().await;
let report = registry.report();
let entry = report
.traffic
.iter()
.find(|e| e.path.as_str() == "demo")
.expect("demo tracked");
let path_len = "demo".len() as u64;
let egress = &entry.publisher;
assert_eq!(egress.announced, 1, "one egress announce");
assert_eq!(egress.announced_bytes, path_len);
assert_eq!(egress.subscriptions, 1, "one egress subscription");
assert_eq!(egress.broadcasts, 1, "one viewer");
assert_eq!(egress.groups, 1);
assert_eq!(egress.frames, 2);
assert_eq!(egress.bytes, 10);
assert_eq!(egress.fetches, 0);
let ingress = &entry.subscriber;
assert_eq!(ingress.announced, 1, "one ingress announce");
assert_eq!(ingress.announced_bytes, path_len);
assert_eq!(ingress.subscriptions, 1, "one ingress track");
assert_eq!(ingress.broadcasts, 0, "ingress has no viewer refcount");
assert_eq!(ingress.groups, 1);
assert_eq!(ingress.frames, 2);
assert_eq!(ingress.bytes, 10);
let fetched = broadcast.track("video").unwrap().fetch_group(0, None).await.unwrap();
let _ = fetched;
settle().await;
let report = registry.report();
let entry = report.traffic.iter().find(|e| e.path.as_str() == "demo").unwrap();
assert_eq!(entry.publisher.fetches, 1, "one fetch");
assert_eq!(entry.publisher.subscriptions, 1, "fetch does not bump subscriptions");
assert_eq!(entry.publisher.broadcasts, 1, "fetch does not bump the viewer refcount");
assert_eq!(entry.subscriber.fetches, 0, "ingress cannot fetch");
}
#[tokio::test]
async fn test_stats_read_frame_counts_once() {
use crate::Timestamp;
use crate::stats::{Config, Registry, Tier};
use bytes::Bytes;
tokio::time::pause();
let registry = Registry::new(Config::new());
let ctx = registry.tier(Tier::default()).session("acme");
let origin = Origin::random().produce();
let ingress = origin.clone().with_stats(ctx.clone());
let egress = origin.consume().with_stats(ctx.clone());
let mut announced = egress.announced();
let source = ingress.create_broadcast("demo", announce()).unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
let broadcast = announced.next().await.unwrap().broadcast.unwrap();
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer = accept_track(&mut dynamic, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
producer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
.unwrap();
let frame = sub.read_frame().await.unwrap().expect("frame");
assert_eq!(frame.payload.len(), 5);
settle().await;
let report = registry.report();
let entry = report
.traffic
.iter()
.find(|e| e.path.as_str() == "demo")
.expect("demo tracked");
assert_eq!(entry.publisher.groups, 1, "one group, counted once");
assert_eq!(entry.publisher.frames, 1, "one frame, counted once");
assert_eq!(
entry.publisher.bytes, 5,
"payload counted once, not zero and not doubled"
);
}
#[tokio::test]
async fn test_stats_datagrams_counted_both_sides() {
use crate::Timestamp;
use crate::stats::{Config, Registry, Tier};
tokio::time::pause();
let registry = Registry::new(Config::new());
let ctx = registry.tier(Tier::default()).session("acme");
let origin = Origin::random().produce();
let ingress = origin.clone().with_stats(ctx.clone());
let egress = origin.consume().with_stats(ctx.clone());
let mut announced = egress.announced();
let source = ingress.create_broadcast("demo", announce()).unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
let broadcast = announced.next().await.unwrap().broadcast.unwrap();
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer = accept_track(&mut dynamic, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
producer.append_datagram(Timestamp::ZERO, &b"hello"[..]).unwrap();
let datagram = sub.recv_datagram().await.unwrap().expect("datagram");
assert_eq!(&datagram.payload[..], b"hello");
settle().await;
let report = registry.report();
let entry = report
.traffic
.iter()
.find(|e| e.path.as_str() == "demo")
.expect("demo tracked");
for (side, traffic) in [("egress", &entry.publisher), ("ingress", &entry.subscriber)] {
assert_eq!(traffic.datagrams, 1, "{side}: one datagram");
assert_eq!(traffic.groups, 1, "{side}: counted as its single-frame group");
assert_eq!(traffic.frames, 1, "{side}: one frame");
assert_eq!(traffic.bytes, 5, "{side}: payload counted once");
}
}
#[test]
fn origin_rejects_reserved_ids() {
assert!(Origin::new(0).is_err());
assert!(Origin::new(1u64 << 62).is_err());
assert_eq!(Origin::new(1).unwrap().id(), 1);
let mut zero = [0u8].as_slice();
assert_eq!(
Origin::decode(&mut zero, crate::lite::Version::Lite05).unwrap(),
Origin::UNKNOWN
);
}
#[test]
fn origin_list_push_fails_at_limit() {
let mut list = OriginList::new();
for _ in 0..MAX_HOPS {
list.push(Origin::random()).unwrap();
}
assert_eq!(list.len(), MAX_HOPS);
assert_eq!(list.push(Origin::random()), Err(TooManyOrigins));
}
#[test]
fn origin_list_replace_first() {
let mut list = OriginList::new();
for _ in 0..3 {
list.push(Origin::UNKNOWN).unwrap();
}
assert!(list.replace_first(Origin::UNKNOWN, Origin::new(7).unwrap()));
assert_eq!(
list.as_slice(),
&[Origin::new(7).unwrap(), Origin::UNKNOWN, Origin::UNKNOWN]
);
assert!(!list.replace_first(Origin::new(99).unwrap(), Origin::new(8).unwrap()));
assert_eq!(list.len(), 3);
}
#[test]
fn origin_list_try_from_vec_enforces_limit() {
let under: Vec<Origin> = (0..MAX_HOPS).map(|_| Origin::random()).collect();
assert!(OriginList::try_from(under).is_ok());
let over: Vec<Origin> = (0..MAX_HOPS + 1).map(|_| Origin::random()).collect();
assert_eq!(OriginList::try_from(over), Err(TooManyOrigins));
}
#[tokio::test]
async fn test_announce() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut consumer1 = origin.consume().announced();
consumer1.assert_next_wait();
let mut broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
settle().await;
consumer1.assert_next_some("test1");
consumer1.assert_next_wait();
let mut consumer2 = origin.consume().announced();
let mut broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
settle().await;
consumer1.assert_next_some("test2");
consumer1.assert_next_wait();
consumer2.assert_next_some("test1");
consumer2.assert_next_some("test2");
consumer2.assert_next_wait();
broadcast1.finish();
settle().await;
consumer1.assert_next_none("test1");
consumer2.assert_next_none("test1");
consumer1.assert_next_wait();
consumer2.assert_next_wait();
let mut consumer3 = origin.consume().announced();
consumer3.assert_next_some("test2");
consumer3.assert_next_wait();
broadcast2.finish();
settle().await;
consumer1.assert_next_none("test2");
consumer2.assert_next_none("test2");
consumer3.assert_next_none("test2");
}
#[tokio::test]
async fn test_duplicate() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
let mut broadcast3 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
assert!(consumer.get_broadcast("test").is_some());
announced.assert_next_some("test");
announced.assert_next_wait();
broadcast2.finish();
settle().await;
assert!(consumer.get_broadcast("test").is_some());
announced.assert_next_wait();
broadcast1.finish();
settle().await;
assert!(consumer.get_broadcast("test").is_some());
announced.assert_next_wait();
broadcast3.finish();
settle().await;
assert!(consumer.get_broadcast("test").is_none());
announced.assert_next_none("test");
announced.assert_next_wait();
}
#[tokio::test]
async fn test_route_failover() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
let mut dynamic_a = source_a.dynamic();
settle().await;
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
let mut dynamic_b = source_b.dynamic();
settle().await;
settle().await;
announced.assert_next_wait();
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer = accept_track(&mut dynamic_a, "video").await;
settle().await;
dynamic_b.assert_no_request();
let mut sub = subscribing.await.unwrap();
sub.assert_no_group();
assert_eq!(producer.subscription().unwrap().group_start, None);
producer.append_group().unwrap();
producer.append_group().unwrap();
assert_eq!(sub.assert_group().sequence, 0);
assert_eq!(sub.assert_group().sequence, 1);
producer.abort(Error::Dropped).unwrap();
source_a.abort(Error::Dropped).unwrap();
drop(dynamic_a);
settle().await;
announced.assert_next_wait();
let mut producer = accept_track(&mut dynamic_b, "video").await;
settle().await;
sub.assert_no_group();
assert_eq!(producer.subscription().unwrap().group_start, Some(2));
producer.create_group(group::Info { sequence: 1 }).unwrap();
producer.create_group(group::Info { sequence: 2 }).unwrap();
assert_eq!(sub.assert_group().sequence, 2, "groups below the boundary are filtered");
sub.assert_not_closed();
}
#[tokio::test]
async fn test_broadcast_route_watch() {
let mut producer = broadcast::Info::new().produce();
let mut consumer = producer.consume();
assert_eq!(consumer.route_changed().await.unwrap(), broadcast::Route::default());
producer.set_route(broadcast::Route::default()).unwrap();
assert!(consumer.route_changed().now_or_never().is_none());
let mut hops = OriginList::new();
hops.push(Origin::new(7).unwrap()).unwrap();
let route = broadcast::Route::new().with_hops(hops).with_cost(3);
producer.set_route(route.clone()).unwrap();
assert_eq!(consumer.route_changed().await.unwrap(), route);
let mut fresh = producer.consume();
assert_eq!(fresh.route_changed().await.unwrap(), route);
drop(producer);
assert!(matches!(consumer.route_changed().await.unwrap_err(), Error::Dropped));
}
#[tokio::test]
async fn test_route_cost_update() {
tokio::time::pause();
let origin = Info::new(origin_keyed("test", Origin::new(3).unwrap(), true)).produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
let mut source_a = origin
.create_broadcast("test", announce().with_hops(hops_a.clone()))
.unwrap();
let mut dynamic_a = source_a.dynamic();
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
let mut watch = broadcast.clone();
assert_eq!(watch.route_changed().await.unwrap().hops, hops_a);
let mut source_b = origin
.create_broadcast("test", announce().with_hops(hops_b.clone()))
.unwrap();
let mut dynamic_b = source_b.dynamic();
settle().await;
assert!(
watch.route_changed().now_or_never().is_none(),
"a losing standby must not change the advertised route"
);
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer = accept_track(&mut dynamic_a, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
producer.append_group().unwrap();
assert_eq!(sub.assert_group().sequence, 0);
source_a
.set_route(announce().with_hops(hops_a.clone()).with_cost(10))
.unwrap();
settle().await;
assert_eq!(watch.route_changed().await.unwrap().hops, hops_b);
announced.assert_next_wait();
let mut producer_b = accept_track(&mut dynamic_b, "video").await;
settle().await;
sub.assert_no_group();
assert_eq!(producer_b.subscription().unwrap().group_start, Some(1));
producer_b.create_group(group::Info { sequence: 1 }).unwrap();
assert_eq!(sub.assert_group().sequence, 1);
sub.assert_not_closed();
source_b
.set_route(announce().with_hops(hops_b.clone()).with_cost(5))
.unwrap();
settle().await;
let advertised = watch.route_changed().await.unwrap();
assert_eq!(advertised.hops, hops_b);
assert_eq!(advertised.cost, 5);
announced.assert_next_wait();
}
#[tokio::test]
async fn test_completed_track_survives_route_churn() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let hops_a = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let hops_b = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
let source_a = origin.create_broadcast("test", announce().with_hops(hops_a)).unwrap();
let mut dynamic_a = source_a.dynamic();
settle().await;
let source_b = origin.create_broadcast("test", announce().with_hops(hops_b)).unwrap();
let mut dynamic_b = source_b.dynamic();
settle().await;
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer = accept_track(&mut dynamic_a, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
producer.append_group().unwrap();
assert_eq!(sub.assert_group().sequence, 0);
producer.finish().unwrap();
drop(producer);
settle().await;
sub.assert_closed();
source_a.abort(Error::Dropped).unwrap();
drop(dynamic_a);
settle().await;
dynamic_b.assert_no_request();
let mut late = broadcast.track("video").unwrap().subscribe(None).await.unwrap();
late.assert_closed();
}
#[tokio::test]
async fn test_serve_resets_retry_budget() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
let subscribing = broadcast.track("video").unwrap().subscribe(None);
for _ in 0..2 * MAX_TRACK_RETRIES {
let request = tokio::time::timeout(std::time::Duration::from_secs(1), dynamic.requested_track())
.await
.expect("timed out waiting for a retry")
.unwrap();
request.reject(Error::NotFound);
let producer = accept_track(&mut dynamic, "video").await;
settle().await;
drop(producer);
}
let _producer = accept_track(&mut dynamic, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
sub.assert_not_closed();
}
#[tokio::test]
async fn test_route_handover() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops_long = OriginList::try_from(vec![Origin::new(2).unwrap(), Origin::new(3).unwrap()]).unwrap();
let hops_short = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let source_a = origin
.create_broadcast("test", announce().with_hops(hops_long))
.unwrap();
let mut dynamic_a = source_a.dynamic();
settle().await;
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer_a = accept_track(&mut dynamic_a, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
producer_a.append_group().unwrap();
producer_a.append_group().unwrap();
assert_eq!(sub.assert_group().sequence, 0);
assert_eq!(sub.assert_group().sequence, 1);
let source_b = origin
.create_broadcast("test", announce().with_hops(hops_short))
.unwrap();
let mut dynamic_b = source_b.dynamic();
settle().await;
settle().await;
announced.assert_next_wait();
let mut producer_b = accept_track(&mut dynamic_b, "video").await;
settle().await;
sub.assert_no_group();
assert_eq!(producer_a.subscription().unwrap().group_end, Some(1));
assert_eq!(producer_b.subscription().unwrap().group_start, Some(2));
producer_a.create_group(group::Info { sequence: 2 }).unwrap();
producer_b.create_group(group::Info { sequence: 2 }).unwrap();
producer_b.create_group(group::Info { sequence: 3 }).unwrap();
assert_eq!(sub.assert_group().sequence, 2);
assert_eq!(sub.assert_group().sequence, 3);
sub.assert_no_group();
sub.assert_not_closed();
}
#[tokio::test(start_paused = true)]
async fn test_route_unannounce_immediate() {
let origin = Origin::random().produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let mut source = origin
.create_broadcast("test", announce().with_hops(hops.clone()))
.unwrap();
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
source.finish();
settle().await;
announced.assert_next_none("test");
let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
settle().await;
let fresh = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
assert!(
!fresh.is_clone(&broadcast),
"re-create must not splice the old broadcast"
);
}
#[tokio::test(start_paused = true)]
async fn test_route_detach_immediate() {
let origin = Origin::random().produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let source = origin
.create_broadcast("test", announce().with_hops(hops.clone()))
.unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let producer = accept_track(&mut dynamic, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
drop(producer);
source.abort(Error::Dropped).unwrap();
drop(dynamic);
settle().await;
announced.assert_next_none("test");
sub.assert_error();
let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
settle().await;
settle().await;
let fresh = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
assert!(
!fresh.is_clone(&broadcast),
"re-create must not splice the old broadcast"
);
}
#[tokio::test(start_paused = true)]
async fn test_linger_reconnect_splices() {
let origin = Info::new(Origin::random())
.with_linger(Duration::from_secs(5))
.produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let source = origin
.create_broadcast("test", announce().with_hops(hops.clone()))
.unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let mut producer = accept_track(&mut dynamic, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
producer.append_group().unwrap();
producer.append_group().unwrap();
assert_eq!(sub.assert_group().sequence, 0);
assert_eq!(sub.assert_group().sequence, 1);
drop(producer);
source.abort(Error::Dropped).unwrap();
drop(dynamic);
settle().await;
announced.assert_next_wait();
sub.assert_no_group();
sub.assert_not_closed();
let during = consumer.request_broadcast("test").await.unwrap();
assert!(during.is_clone(&broadcast), "the lingering broadcast still resolves");
let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
announced.assert_next_wait();
let again = consumer.request_broadcast("test").await.unwrap();
assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
let mut producer = accept_track(&mut dynamic, "video").await;
settle().await;
sub.assert_no_group();
assert_eq!(producer.subscription().unwrap().group_start, Some(2));
producer.create_group(group::Info { sequence: 2 }).unwrap();
assert_eq!(sub.assert_group().sequence, 2);
sub.assert_not_closed();
}
#[tokio::test(start_paused = true)]
async fn test_linger_expiry_closes() {
let origin = Info::new(Origin::random())
.with_linger(Duration::from_secs(5))
.produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let source = origin
.create_broadcast("test", announce().with_hops(hops.clone()))
.unwrap();
let mut dynamic = source.dynamic();
settle().await;
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
let subscribing = broadcast.track("video").unwrap().subscribe(None);
let producer = accept_track(&mut dynamic, "video").await;
settle().await;
let mut sub = subscribing.await.unwrap();
drop(producer);
source.abort(Error::Dropped).unwrap();
drop(dynamic);
settle().await;
announced.assert_next_wait();
tokio::time::sleep(std::time::Duration::from_secs(6)).await;
settle().await;
announced.assert_next_none("test");
sub.assert_error();
let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
settle().await;
settle().await;
let fresh = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
assert!(
!fresh.is_clone(&broadcast),
"a late re-create must not splice the expired broadcast"
);
}
#[tokio::test(start_paused = true)]
async fn test_linger_forever() {
let origin = Info::new(Origin::random()).with_linger(Duration::MAX).produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let source = origin
.create_broadcast("test", announce().with_hops(hops.clone()))
.unwrap();
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
source.abort(Error::Dropped).unwrap();
settle().await;
tokio::time::sleep(std::time::Duration::from_secs(60 * 60 * 24 * 3)).await;
announced.assert_next_wait();
let source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
settle().await;
settle().await;
let again = consumer.request_broadcast("test").await.unwrap();
assert!(again.is_clone(&broadcast), "the reconnect must splice, not replace");
drop(source);
}
#[tokio::test(start_paused = true)]
async fn test_linger_skipped_on_finish() {
let origin = Info::new(Origin::random())
.with_linger(Duration::from_secs(5))
.produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let mut source = origin
.create_broadcast("test", announce().with_hops(hops.clone()))
.unwrap();
settle().await;
let broadcast = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
source.finish();
settle().await;
announced.assert_next_none("test");
let _source = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
settle().await;
let fresh = consumer.request_broadcast("test").await.unwrap();
announced.assert_next_some("test");
assert!(
!fresh.is_clone(&broadcast),
"a finish must not leave a lingering broadcast to splice into"
);
}
#[tokio::test]
async fn test_announce_toggle() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let mut source = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
settle().await;
announced.assert_next_wait();
let broadcast = consumer
.get_broadcast("test")
.expect("offline broadcast is still routable");
assert!(!broadcast.route().announce);
let requested = consumer.request_broadcast("test").await.unwrap();
assert!(requested.is_clone(&broadcast));
source.set_route(announce()).unwrap();
settle().await;
let face = announced.assert_next_some("test");
assert!(face.is_clone(&broadcast));
let mut fresh = origin.consume().announced();
fresh.assert_next_some("test");
fresh.assert_next_wait();
source.set_route(broadcast::Route::new()).unwrap();
settle().await;
announced.assert_next_none("test");
assert!(consumer.get_broadcast("test").is_some());
let mut fresh = origin.consume().announced();
fresh.assert_next_wait();
source.finish();
settle().await;
assert!(consumer.get_broadcast("test").is_none());
}
#[tokio::test]
async fn test_announce_beats_offline() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let mut announced = consumer.announced();
let _offline = origin.create_broadcast("test", broadcast::Route::new()).unwrap();
settle().await;
announced.assert_next_wait();
let mut announced_source = origin.create_broadcast("test", announce().with_cost(10)).unwrap();
settle().await;
announced.assert_next_some("test");
let face = consumer.get_broadcast("test").unwrap();
assert!(face.route().announce);
assert_eq!(face.route().cost, 10);
announced_source.finish();
settle().await;
announced.assert_next_none("test");
assert!(consumer.get_broadcast("test").is_some());
}
#[tokio::test]
async fn test_better_source_no_churn() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut announced = origin.consume().announced();
let hops = OriginList::try_from(vec![Origin::new(1).unwrap()]).unwrap();
let _a = origin.create_broadcast("test", announce().with_hops(hops)).unwrap();
settle().await;
let face = announced.assert_next_some("test");
let _b = origin.create_broadcast("test", announce()).unwrap();
settle().await;
announced.assert_next_wait();
let current = origin.consume().get_broadcast("test").unwrap();
assert!(current.is_clone(&face), "the broadcast identity must not change");
assert!(current.route().hops.is_empty());
}
#[tokio::test]
async fn test_duplicate_reverse() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
assert!(origin.consume().get_broadcast("test").is_some());
broadcast2.finish();
settle().await;
assert!(origin.consume().get_broadcast("test").is_some());
broadcast1.finish();
settle().await;
assert!(origin.consume().get_broadcast("test").is_none());
}
#[tokio::test]
async fn test_deterministic_tiebreak() {
tokio::time::pause();
fn hops(ids: &[u64]) -> OriginList {
OriginList::try_from(
ids.iter()
.copied()
.map(|id| Origin::new(id).unwrap())
.collect::<Vec<_>>(),
)
.unwrap()
}
async fn winner(first: &[u64], second: &[u64]) -> OriginList {
let origin = Origin::random().produce();
let _a = origin
.create_broadcast("test", announce().with_hops(hops(first)))
.unwrap();
let _b = origin
.create_broadcast("test", announce().with_hops(hops(second)))
.unwrap();
settle().await;
origin.consume().get_broadcast("test").unwrap().route().hops
}
let forward = winner(&[10, 20], &[30, 40]).await;
let reverse = winner(&[30, 40], &[10, 20]).await;
assert_eq!(forward, reverse, "tie-break must not depend on publish order");
assert_eq!(winner(&[10, 20], &[30]).await.len(), 1);
assert_eq!(winner(&[30], &[10, 20]).await.len(), 1);
}
#[tokio::test]
async fn test_many_announces() {
let origin = Origin::random().produce();
let mut consumer = origin.consume().announced();
let mut broadcasts = Vec::new();
for i in 0..256 {
broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
settle().await;
}
for i in 0..256 {
consumer.assert_next_some(format!("test{i:03}"));
}
consumer.assert_next_wait();
}
#[tokio::test]
async fn test_many_announces_try() {
let origin = Origin::random().produce();
let mut consumer = origin.consume().announced();
let mut broadcasts = Vec::new();
for i in 0..256 {
broadcasts.push(origin.create_broadcast(format!("test{i:03}"), announce()).unwrap());
settle().await;
}
for i in 0..256 {
consumer.assert_try_next_some(format!("test{i:03}"));
}
}
#[tokio::test]
async fn test_with_root_basic() {
let origin = Origin::random().produce();
let foo_producer = origin.with_root("foo").expect("should create root");
assert_eq!(foo_producer.root().as_str(), "foo");
let mut consumer = origin.consume().announced();
let _broadcast = foo_producer
.create_broadcast("bar/baz", announce())
.expect("publish allowed");
settle().await;
consumer.assert_next_some("foo/bar/baz");
let mut foo_consumer = foo_producer.consume().announced();
foo_consumer.assert_next_some("bar/baz");
}
#[tokio::test]
async fn test_with_root_nested() {
let origin = Origin::random().produce();
let foo_producer = origin.with_root("foo").expect("should create foo root");
let foo_bar_producer = foo_producer.with_root("bar").expect("should create bar root");
assert_eq!(foo_bar_producer.root().as_str(), "foo/bar");
let mut consumer = origin.consume().announced();
let _broadcast = foo_bar_producer
.create_broadcast("baz", announce())
.expect("publish allowed");
settle().await;
consumer.assert_next_some("foo/bar/baz");
let mut foo_bar_consumer = foo_bar_producer.consume().announced();
foo_bar_consumer.assert_next_some("baz");
}
#[tokio::test]
async fn test_publish_scope_allows() {
let origin = Origin::random().produce();
let limited_producer = origin
.scope(&["allowed/path1".into(), "allowed/path2".into()])
.expect("should create limited producer");
let _broadcast = limited_producer
.create_broadcast("allowed/path1", announce())
.expect("publish allowed");
let _keep2 = limited_producer
.create_broadcast("allowed/path1/nested", announce())
.expect("publish allowed");
let _keep3 = limited_producer
.create_broadcast("allowed/path2", announce())
.expect("publish allowed");
settle().await;
assert!(limited_producer.create_broadcast("notallowed", announce()).is_err());
assert!(limited_producer.create_broadcast("allowed", announce()).is_err()); assert!(limited_producer.create_broadcast("other/path", announce()).is_err());
}
#[tokio::test]
async fn test_publish_max_parts() {
let origin = Origin::random().produce();
let at_limit = (0..Path::MAX_PARTS)
.map(|i| i.to_string())
.collect::<Vec<_>>()
.join("/");
let _broadcast = origin
.create_broadcast(at_limit.as_str(), announce())
.expect("publish allowed");
settle().await;
let too_deep = format!("{at_limit}/extra");
assert!(origin.create_broadcast(too_deep.as_str(), announce()).is_err());
let rooted = origin.with_root("root").expect("wildcard allows any root");
assert!(rooted.create_broadcast(at_limit.as_str(), announce()).is_err());
}
#[tokio::test]
async fn test_publish_scope_empty() {
let origin = Origin::random().produce();
assert!(origin.scope(&[]).is_none());
}
#[tokio::test]
async fn test_consume_scope_filters() {
let origin = Origin::random().produce();
let mut consumer = origin.consume().announced();
let _broadcast1 = origin.create_broadcast("allowed", announce()).unwrap();
let _broadcast2 = origin.create_broadcast("allowed/nested", announce()).unwrap();
let _broadcast3 = origin.create_broadcast("notallowed", announce()).unwrap();
settle().await;
let mut limited_consumer = origin
.consume()
.scope(&["allowed".into()])
.expect("should create limited consumer")
.announced();
limited_consumer.assert_next_some("allowed");
limited_consumer.assert_next_some("allowed/nested");
limited_consumer.assert_next_wait();
consumer.assert_next_some("allowed");
consumer.assert_next_some("allowed/nested");
consumer.assert_next_some("notallowed");
}
#[tokio::test]
async fn test_consume_scope_multiple_prefixes() {
let origin = Origin::random().produce();
let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
settle().await;
let mut limited_consumer = origin
.consume()
.scope(&["foo".into(), "bar".into()])
.expect("should create limited consumer")
.announced();
limited_consumer.assert_next_some("bar/test");
limited_consumer.assert_next_some("foo/test");
limited_consumer.assert_next_wait(); }
#[tokio::test]
async fn test_with_root_and_publish_scope() {
let origin = Origin::random().produce();
let foo_producer = origin.with_root("foo").expect("should create foo root");
let limited_producer = foo_producer
.scope(&["bar".into(), "goop/pee".into()])
.expect("should create limited producer");
let mut consumer = origin.consume().announced();
let _broadcast = limited_producer
.create_broadcast("bar", announce())
.expect("publish allowed");
let _keep2 = limited_producer
.create_broadcast("bar/nested", announce())
.expect("publish allowed");
let _keep3 = limited_producer
.create_broadcast("goop/pee", announce())
.expect("publish allowed");
let _keep4 = limited_producer
.create_broadcast("goop/pee/nested", announce())
.expect("publish allowed");
settle().await;
assert!(limited_producer.create_broadcast("baz", announce()).is_err());
assert!(limited_producer.create_broadcast("goop", announce()).is_err()); assert!(limited_producer.create_broadcast("goop/other", announce()).is_err());
consumer.assert_next_some("foo/bar");
consumer.assert_next_some("foo/bar/nested");
consumer.assert_next_some("foo/goop/pee");
consumer.assert_next_some("foo/goop/pee/nested");
}
#[tokio::test]
async fn test_with_root_and_consume_scope() {
let origin = Origin::random().produce();
let _broadcast1 = origin.create_broadcast("foo/bar/test", announce()).unwrap();
let _broadcast2 = origin.create_broadcast("foo/goop/pee/test", announce()).unwrap();
let _broadcast3 = origin.create_broadcast("foo/other/test", announce()).unwrap();
settle().await;
let foo_producer = origin.with_root("foo").expect("should create foo root");
let mut limited_consumer = foo_producer
.consume()
.scope(&["bar".into(), "goop/pee".into()])
.expect("should create limited consumer")
.announced();
limited_consumer.assert_next_some("bar/test");
limited_consumer.assert_next_some("goop/pee/test");
limited_consumer.assert_next_wait(); }
#[tokio::test]
async fn test_with_root_unauthorized() {
let origin = Origin::random().produce();
let limited_producer = origin
.scope(&["allowed".into()])
.expect("should create limited producer");
assert!(limited_producer.with_root("notallowed").is_none());
let allowed_root = limited_producer
.with_root("allowed")
.expect("should create allowed root");
assert_eq!(allowed_root.root().as_str(), "allowed");
}
#[tokio::test]
async fn test_wildcard_permission() {
let origin = Origin::random().produce();
let root_producer = origin.clone();
let _broadcast = root_producer
.create_broadcast("any/path", announce())
.expect("publish allowed");
let _keep2 = root_producer
.create_broadcast("other/path", announce())
.expect("publish allowed");
settle().await;
let foo_producer = root_producer.with_root("foo").expect("should create any root");
assert_eq!(foo_producer.root().as_str(), "foo");
}
#[tokio::test]
async fn test_consume_broadcast_with_permissions() {
let origin = Origin::random().produce();
let _broadcast1 = origin.create_broadcast("allowed/test", announce()).unwrap();
let _broadcast2 = origin.create_broadcast("notallowed/test", announce()).unwrap();
settle().await;
let limited_consumer = origin
.consume()
.scope(&["allowed".into()])
.expect("should create limited consumer");
let result = limited_consumer.get_broadcast("allowed/test");
assert!(result.is_some());
assert!(
result
.unwrap()
.is_clone(&origin.consume().get_broadcast("allowed/test").unwrap())
);
assert!(limited_consumer.get_broadcast("notallowed/test").is_none());
let consumer = origin.consume();
assert!(consumer.get_broadcast("allowed/test").is_some());
assert!(consumer.get_broadcast("notallowed/test").is_some());
}
#[tokio::test]
async fn test_nested_paths_with_permissions() {
let origin = Origin::random().produce();
let limited_producer = origin.scope(&["a/b/c".into()]).expect("should create limited producer");
let _broadcast = limited_producer
.create_broadcast("a/b/c", announce())
.expect("publish allowed");
let _keep2 = limited_producer
.create_broadcast("a/b/c/d", announce())
.expect("publish allowed");
let _keep3 = limited_producer
.create_broadcast("a/b/c/d/e", announce())
.expect("publish allowed");
settle().await;
assert!(limited_producer.create_broadcast("a", announce()).is_err());
assert!(limited_producer.create_broadcast("a/b", announce()).is_err());
assert!(limited_producer.create_broadcast("a/b/other", announce()).is_err());
}
#[tokio::test]
async fn test_multiple_consumers_with_different_permissions() {
let origin = Origin::random().produce();
let _broadcast1 = origin.create_broadcast("foo/test", announce()).unwrap();
let _broadcast2 = origin.create_broadcast("bar/test", announce()).unwrap();
let _broadcast3 = origin.create_broadcast("baz/test", announce()).unwrap();
settle().await;
let mut foo_consumer = origin
.consume()
.scope(&["foo".into()])
.expect("should create foo consumer")
.announced();
let mut bar_consumer = origin
.consume()
.scope(&["bar".into()])
.expect("should create bar consumer")
.announced();
let mut foobar_consumer = origin
.consume()
.scope(&["foo".into(), "bar".into()])
.expect("should create foobar consumer")
.announced();
foo_consumer.assert_next_some("foo/test");
foo_consumer.assert_next_wait();
bar_consumer.assert_next_some("bar/test");
bar_consumer.assert_next_wait();
foobar_consumer.assert_next_some("bar/test");
foobar_consumer.assert_next_some("foo/test");
foobar_consumer.assert_next_wait();
}
#[tokio::test]
async fn test_select_with_empty_prefix() {
let origin = Origin::random().produce();
let demo_producer = origin.with_root("demo").expect("should create demo root");
let limited_producer = demo_producer
.scope(&["worm-node".into(), "foobar".into()])
.expect("should create limited producer");
let _broadcast1 = limited_producer
.create_broadcast("worm-node/test", announce())
.expect("publish allowed");
let _broadcast2 = limited_producer
.create_broadcast("foobar/test", announce())
.expect("publish allowed");
settle().await;
let mut consumer = limited_producer
.consume()
.scope(&["".into()])
.expect("should create consumer with empty prefix")
.announced();
let a1 = consumer.try_next().expect("expected first announcement");
let a2 = consumer.try_next().expect("expected second announcement");
consumer.assert_next_wait();
let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
paths.sort();
assert_eq!(paths, ["foobar/test", "worm-node/test"]);
}
#[tokio::test]
async fn test_select_narrowing_scope() {
let origin = Origin::random().produce();
let demo_producer = origin.with_root("demo").expect("should create demo root");
let limited_producer = demo_producer
.scope(&["worm-node".into(), "foobar".into()])
.expect("should create limited producer");
let _broadcast1 = limited_producer
.create_broadcast("worm-node", announce())
.expect("publish allowed");
let _broadcast2 = limited_producer
.create_broadcast("worm-node/foo", announce())
.expect("publish allowed");
let _broadcast3 = limited_producer
.create_broadcast("foobar/bar", announce())
.expect("publish allowed");
settle().await;
let mut worm_consumer = limited_producer
.consume()
.scope(&["worm-node".into()])
.expect("should create worm-node consumer")
.announced();
worm_consumer.assert_next_some("worm-node");
worm_consumer.assert_next_some("worm-node/foo");
worm_consumer.assert_next_wait();
let mut foo_consumer = limited_producer
.consume()
.scope(&["worm-node/foo".into()])
.expect("should create worm-node/foo consumer")
.announced();
foo_consumer.assert_next_some("worm-node/foo");
foo_consumer.assert_next_wait(); }
#[tokio::test]
async fn test_select_multiple_roots_with_empty_prefix() {
let origin = Origin::random().produce();
let limited_producer = origin
.scope(&["app1".into(), "app2".into(), "shared".into()])
.expect("should create limited producer");
let _broadcast1 = limited_producer
.create_broadcast("app1/data", announce())
.expect("publish allowed");
let _broadcast2 = limited_producer
.create_broadcast("app2/config", announce())
.expect("publish allowed");
let _broadcast3 = limited_producer
.create_broadcast("shared/resource", announce())
.expect("publish allowed");
settle().await;
let mut consumer = limited_producer
.consume()
.scope(&["".into()])
.expect("should create consumer with empty prefix")
.announced();
consumer.assert_next_some("app1/data");
consumer.assert_next_some("app2/config");
consumer.assert_next_some("shared/resource");
consumer.assert_next_wait();
}
#[tokio::test]
async fn test_publish_scope_with_empty_prefix() {
let origin = Origin::random().produce();
let limited_producer = origin
.scope(&["services/api".into(), "services/web".into()])
.expect("should create limited producer");
let same_producer = limited_producer
.scope(&["".into()])
.expect("should create producer with empty prefix");
let _broadcast = same_producer
.create_broadcast("services/api", announce())
.expect("publish allowed");
let _keep2 = same_producer
.create_broadcast("services/web", announce())
.expect("publish allowed");
assert!(same_producer.create_broadcast("services/db", announce()).is_err());
assert!(same_producer.create_broadcast("other", announce()).is_err());
}
#[tokio::test]
async fn test_select_narrowing_to_deeper_path() {
let origin = Origin::random().produce();
let limited_producer = origin.scope(&["org".into()]).expect("should create limited producer");
let _broadcast1 = limited_producer
.create_broadcast("org/team1/project1", announce())
.expect("publish allowed");
let _broadcast2 = limited_producer
.create_broadcast("org/team1/project2", announce())
.expect("publish allowed");
let _broadcast3 = limited_producer
.create_broadcast("org/team2/project1", announce())
.expect("publish allowed");
settle().await;
let mut team2_consumer = limited_producer
.consume()
.scope(&["org/team2".into()])
.expect("should create team2 consumer")
.announced();
team2_consumer.assert_next_some("org/team2/project1");
team2_consumer.assert_next_wait();
let mut project1_consumer = limited_producer
.consume()
.scope(&["org/team1/project1".into()])
.expect("should create project1 consumer")
.announced();
project1_consumer.assert_next_some("org/team1/project1");
project1_consumer.assert_next_wait();
}
#[tokio::test]
async fn test_select_with_non_matching_prefix() {
let origin = Origin::random().produce();
let limited_producer = origin
.scope(&["allowed/path".into()])
.expect("should create limited producer");
assert!(limited_producer.consume().scope(&["different/path".into()]).is_none());
assert!(limited_producer.scope(&["other/path".into()]).is_none());
}
#[tokio::test]
async fn test_with_root_trailing_slash_consumer() {
let origin = Origin::random().produce();
let prefix = "some_prefix/".to_string();
let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
let _b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
settle().await;
consumer.assert_next_some("test");
}
#[tokio::test]
async fn test_with_root_trailing_slash_producer() {
let origin = Origin::random().produce();
let prefix = "some_prefix/".to_string();
let rooted = origin.with_root(prefix).unwrap();
let _b = rooted.create_broadcast("test", announce()).unwrap();
settle().await;
let mut consumer = rooted.consume().announced();
consumer.assert_next_some("test");
}
#[tokio::test]
async fn test_with_root_trailing_slash_unannounce() {
tokio::time::pause();
let origin = Origin::random().produce();
let prefix = "some_prefix/".to_string();
let mut consumer = origin.consume().with_root(prefix).unwrap().announced();
let mut b = origin.create_broadcast("some_prefix/test", announce()).unwrap();
settle().await;
consumer.assert_next_some("test");
b.finish();
settle().await;
consumer.assert_next_none("test");
}
#[tokio::test]
async fn test_select_maintains_access_with_wider_prefix() {
let origin = Origin::random().produce();
let demo_producer = origin.with_root("demo").expect("should create demo root");
let user_producer = demo_producer
.scope(&["worm-node".into(), "foobar".into()])
.expect("should create user producer");
let _broadcast1 = user_producer
.create_broadcast("worm-node/data", announce())
.expect("publish allowed");
let _broadcast2 = user_producer
.create_broadcast("foobar", announce())
.expect("publish allowed");
settle().await;
let mut consumer = user_producer
.consume()
.scope(&["".into()])
.expect("scope with empty prefix should not fail when user has specific permissions")
.announced();
let a1 = consumer.try_next().expect("expected first announcement");
let a2 = consumer.try_next().expect("expected second announcement");
consumer.assert_next_wait();
let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
paths.sort();
assert_eq!(paths, ["foobar", "worm-node/data"]);
let mut narrow_consumer = user_producer
.consume()
.scope(&["worm-node".into()])
.expect("should be able to narrow scope to worm-node")
.announced();
narrow_consumer.assert_next_some("worm-node/data");
narrow_consumer.assert_next_wait(); }
#[tokio::test]
async fn test_duplicate_prefixes_deduped() {
let origin = Origin::random().produce();
let producer = origin
.scope(&["demo".into(), "demo".into()])
.expect("should create producer");
let _broadcast = producer
.create_broadcast("demo/stream", announce())
.expect("publish allowed");
settle().await;
let mut consumer = producer.consume().announced();
consumer.assert_next_some("demo/stream");
consumer.assert_next_wait();
}
#[tokio::test]
async fn test_overlapping_prefixes_deduped() {
let origin = Origin::random().produce();
let producer = origin
.scope(&["demo".into(), "demo/foo".into()])
.expect("should create producer");
let _broadcast = producer
.create_broadcast("demo/bar/stream", announce())
.expect("publish allowed");
settle().await;
let mut consumer = producer.consume().announced();
consumer.assert_next_some("demo/bar/stream");
consumer.assert_next_wait();
}
#[tokio::test]
async fn test_overlapping_prefixes_no_duplicate_announcements() {
let origin = Origin::random().produce();
let producer = origin
.scope(&["demo".into(), "demo/foo".into()])
.expect("should create producer");
let _broadcast = producer
.create_broadcast("demo/foo/stream", announce())
.expect("publish allowed");
settle().await;
let mut consumer = producer.consume().announced();
consumer.assert_next_some("demo/foo/stream");
consumer.assert_next_wait();
}
#[tokio::test]
async fn test_allowed_returns_deduped_prefixes() {
let origin = Origin::random().produce();
let producer = origin
.scope(&["demo".into(), "demo/foo".into(), "anon".into()])
.expect("should create producer");
let allowed: Vec<_> = producer.allowed().collect();
assert_eq!(allowed.len(), 2, "demo/foo should be subsumed by demo");
}
#[tokio::test]
async fn test_announced_broadcast_already_announced() {
let origin = Origin::random().produce();
let _broadcast = origin.create_broadcast("test", announce()).unwrap();
settle().await;
let consumer = origin.consume();
let result = consumer.announced_broadcast("test").await.expect("should find it");
assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
}
#[tokio::test]
async fn test_announced_broadcast_delayed() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let wait = tokio::spawn({
let consumer = consumer.clone();
async move { consumer.announced_broadcast("test").await }
});
tokio::task::yield_now().await;
let _broadcast = origin.create_broadcast("test", announce()).unwrap();
settle().await;
let result = wait.await.unwrap().expect("should find it");
assert!(result.is_clone(&consumer.get_broadcast("test").unwrap()));
}
#[tokio::test]
async fn test_announced_broadcast_ignores_unrelated_paths() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let wait = tokio::spawn({
let consumer = consumer.clone();
async move { consumer.announced_broadcast("target").await }
});
tokio::task::yield_now().await;
let _other = origin.create_broadcast("other", announce()).unwrap();
settle().await;
tokio::task::yield_now().await;
assert!(!wait.is_finished(), "must not resolve on unrelated path");
let _target = origin.create_broadcast("target", announce()).unwrap();
settle().await;
let result = wait.await.unwrap().expect("should find target");
assert!(result.is_clone(&consumer.get_broadcast("target").unwrap()));
}
#[tokio::test]
async fn test_announced_broadcast_skips_nested_paths() {
tokio::time::pause();
let origin = Origin::random().produce();
let consumer = origin.consume();
let wait = tokio::spawn({
let consumer = consumer.clone();
async move { consumer.announced_broadcast("foo").await }
});
tokio::task::yield_now().await;
let _nested = origin.create_broadcast("foo/bar", announce()).unwrap();
settle().await;
tokio::task::yield_now().await;
assert!(!wait.is_finished(), "must not resolve on a nested path");
let _exact = origin.create_broadcast("foo", announce()).unwrap();
settle().await;
let result = wait.await.unwrap().expect("should find foo exactly");
assert!(result.is_clone(&consumer.get_broadcast("foo").unwrap()));
}
#[tokio::test]
async fn test_announced_broadcast_disallowed() {
let origin = Origin::random().produce();
let limited = origin
.consume()
.scope(&["allowed".into()])
.expect("should create limited");
assert!(limited.announced_broadcast("notallowed").await.is_none());
}
#[tokio::test]
async fn test_announced_broadcast_scope_too_narrow() {
let origin = Origin::random().produce();
let limited = origin
.consume()
.scope(&["foo/specific".into()])
.expect("should create limited");
let result = limited
.announced_broadcast("foo")
.now_or_never()
.expect("must not block");
assert!(result.is_none());
}
#[tokio::test]
async fn test_coalesce_announce_then_unannounce() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut announced = origin.consume().announced();
let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
settle().await;
broadcast.finish();
settle().await;
announced.assert_next_wait();
}
#[tokio::test]
async fn test_coalesce_announce_unannounce_announce() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut announced = origin.consume().announced();
let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
broadcast1.finish();
settle().await;
let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
announced.assert_next_some("test");
announced.assert_next_wait();
}
#[tokio::test]
async fn test_coalesce_unannounce_announce_preserved() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
let mut announced = origin.consume().announced();
announced.assert_next_some("test");
broadcast1.finish();
settle().await;
let _broadcast2 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
announced.assert_next_none("test");
announced.assert_next_some("test");
announced.assert_next_wait();
}
#[tokio::test]
async fn test_coalesce_unannounce_announce_unannounce() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut broadcast1 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
let mut announced = origin.consume().announced();
announced.assert_next_some("test");
broadcast1.finish();
settle().await;
let mut broadcast2 = origin.create_broadcast("test", announce()).unwrap();
settle().await;
broadcast2.finish();
settle().await;
announced.assert_next_none("test");
announced.assert_next_wait();
}
#[tokio::test]
async fn test_coalesce_churn_bounded() {
tokio::time::pause();
let origin = Origin::random().produce();
let mut announced = origin.consume().announced();
for _ in 0..1000 {
let mut broadcast = origin.create_broadcast("test", announce()).unwrap();
settle().await;
broadcast.finish();
}
settle().await;
let mut collected = Vec::new();
while let Some(update) = announced.try_next() {
collected.push(update);
}
assert!(
collected.len() <= 1,
"expected at most one pending update, got {}",
collected.len()
);
assert!(
collected.iter().all(|a| a.path == Path::new("test")),
"unexpected path in pending updates",
);
}
#[tokio::test]
async fn test_consumer_clone_is_side_effect_free() {
let origin = Origin::random().produce();
let _broadcast1 = origin.create_broadcast("test1", announce()).unwrap();
let _broadcast2 = origin.create_broadcast("test2", announce()).unwrap();
settle().await;
let consumer = origin.consume();
let mut announced = consumer.announced();
for _ in 0..16 {
let cloned = consumer.clone();
assert!(cloned.get_broadcast("test1").is_some());
assert!(cloned.get_broadcast("test2").is_some());
}
let a1 = announced.try_next().expect("first announcement");
let a2 = announced.try_next().expect("second announcement");
announced.assert_next_wait();
let mut paths: Vec<_> = [&a1, &a2].iter().map(|a| a.path.to_string()).collect();
paths.sort();
assert_eq!(paths, ["test1", "test2"]);
let mut fresh = consumer.announced();
let b1 = fresh.try_next().expect("backlog: first");
let b2 = fresh.try_next().expect("backlog: second");
fresh.assert_next_wait();
let mut paths: Vec<_> = [&b1, &b2].iter().map(|a| a.path.to_string()).collect();
paths.sort();
assert_eq!(paths, ["test1", "test2"]);
}
#[tokio::test]
async fn dynamic_request_unroutable_without_handler() {
let origin = Origin::random().produce();
let consumer = origin.consume();
assert!(matches!(
consumer.request_broadcast("missing").await,
Err(Error::Unroutable)
));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_served_not_announced() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let mut announced = origin.consume().announced();
announced.assert_next_wait();
let served = broadcast::Info::new().produce();
let request_fut = consumer.request_broadcast("fallback");
let mut served_dynamic = served.dynamic();
let request = dynamic.requested_broadcast().await.unwrap();
assert_eq!(request.path(), &Path::new("fallback"));
request.accept(&served);
let broadcast = request_fut.await.unwrap();
assert!(broadcast.is_clone(&served.consume()));
let track_fut = broadcast.track("video").unwrap().subscribe(None);
let mut producer = served_dynamic.requested_track().await.unwrap().accept(None);
let mut track = track_fut.await.unwrap();
producer.append_group().unwrap();
track.assert_group();
announced.assert_next_wait();
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_coalesces() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let f1 = consumer.request_broadcast("dup");
let f2 = consumer.request_broadcast("dup");
let request = dynamic.requested_broadcast().await.unwrap();
assert_eq!(request.path(), &Path::new("dup"));
assert!(
dynamic.requested_broadcast().now_or_never().is_none(),
"a coalesced request must not be served twice"
);
let served = broadcast::Info::new().produce();
request.accept(&served);
assert!(f1.await.unwrap().is_clone(&served.consume()));
assert!(f2.await.unwrap().is_clone(&served.consume()));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_dedups_served() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let request_fut = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
let served = broadcast::Info::new().produce();
request.accept(&served);
let first = request_fut.await.unwrap();
assert!(first.is_clone(&served.consume()));
let second = consumer.request_broadcast("fallback").await.unwrap();
assert!(second.is_clone(&served.consume()));
assert!(
dynamic.requested_broadcast().now_or_never().is_none(),
"a still-live served broadcast must not be re-requested from the handler"
);
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_reserves_after_close() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let request_fut = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
let served = broadcast::Info::new().produce();
request.accept(&served);
request_fut.await.unwrap();
drop(served);
let request_fut = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
assert_eq!(request.path(), &Path::new("fallback"));
let served = broadcast::Info::new().produce();
request.accept(&served);
assert!(request_fut.await.unwrap().is_clone(&served.consume()));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_served_cache_bounded() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
for i in 0..100 {
let path = format!("one-shot/{i}");
let request_fut = consumer.request_broadcast(&path);
let request = dynamic.requested_broadcast().await.unwrap();
let served = broadcast::Info::new().produce();
request.accept(&served);
request_fut.await.unwrap();
drop(served);
}
assert!(
origin.dynamic.read().served.len() <= 4,
"stale served entries must be reclaimed, not accumulate per distinct path: {}",
origin.dynamic.read().served.len()
);
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_coalesces_after_handoff() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let f1 = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
let f2 = consumer.request_broadcast("fallback");
assert!(
dynamic.requested_broadcast().now_or_never().is_none(),
"a repeat request during hand-off must coalesce, not re-queue"
);
let served = broadcast::Info::new().produce();
request.accept(&served);
assert!(f1.await.unwrap().is_clone(&served.consume()));
assert!(f2.await.unwrap().is_clone(&served.consume()));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_dropped_after_handoff() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let f1 = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
let f2 = consumer.request_broadcast("fallback");
drop(request);
assert!(matches!(f1.await, Err(Error::Unroutable)));
assert!(matches!(f2.await, Err(Error::Unroutable)));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_rejected() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let request_fut = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
request.reject(Error::Cancel);
assert!(matches!(request_fut.await, Err(Error::Cancel)));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_rerequest_after_reject() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let f1 = consumer.request_broadcast("fallback");
dynamic.requested_broadcast().await.unwrap().reject(Error::Unroutable);
assert!(matches!(f1.await, Err(Error::Unroutable)));
let served = broadcast::Info::new().produce();
let f2 = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
assert_eq!(request.path(), &Path::new("fallback"));
request.accept(&served);
assert!(f2.await.unwrap().is_clone(&served.consume()));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_handler_dropped() {
let origin = Origin::random().produce();
let dynamic = origin.dynamic();
let consumer = origin.consume();
let request_fut = consumer.request_broadcast("fallback");
drop(dynamic);
assert!(matches!(request_fut.await, Err(Error::Unroutable)));
assert!(matches!(
consumer.request_broadcast("again").await,
Err(Error::Unroutable)
));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_accept_after_handler_dropped() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let request_fut = consumer.request_broadcast("fallback");
let request = dynamic.requested_broadcast().await.unwrap();
drop(dynamic);
let served = broadcast::Info::new().produce();
request.accept(&served);
assert!(request_fut.await.unwrap().is_clone(&served.consume()));
}
#[tokio::test(start_paused = true)]
async fn dynamic_request_prefers_announced() {
let origin = Origin::random().produce();
let mut dynamic = origin.dynamic();
let consumer = origin.consume();
let _broadcast = origin.create_broadcast("live", announce()).unwrap();
settle().await;
let got = consumer.request_broadcast("live").await.unwrap();
assert!(
got.is_clone(&consumer.get_broadcast("live").unwrap()),
"should return the published broadcast"
);
assert!(
dynamic.requested_broadcast().now_or_never().is_none(),
"a published path must not queue a fallback request"
);
}
#[tokio::test(start_paused = true)]
async fn dynamic_clone_keeps_alive() {
let origin = Origin::random().produce();
let dynamic = origin.dynamic();
let consumer = origin.consume();
drop(dynamic.clone());
let request_fut = consumer.request_broadcast("fallback");
assert!(
request_fut.now_or_never().is_none(),
"request should stay pending until served"
);
}
}