diff --git a/crates/didbot-serve/src/subscribe.rs b/crates/didbot-serve/src/subscribe.rs index 825d4939..e6d7c469 100644 --- a/crates/didbot-serve/src/subscribe.rs +++ b/crates/didbot-serve/src/subscribe.rs @@ -365,6 +365,25 @@ impl Repos { self.inner.sequence.latest() } + /// How many connections are attached to the live channel right now. + /// + /// The per-connection state this producer holds is exactly one broadcast + /// receiver, so this figure is the whole of it: a consumer that hung up + /// and left its receiver behind would show up here and nowhere else. + /// Observability rather than a bound — nothing refuses a subscriber for + /// being the n-th — and the number a test watches to see a departed + /// consumer actually freed. + #[must_use] + pub fn subscribers(&self) -> usize { + self.inner.tx.receiver_count() + } + + /// How many frames the replay buffer holds at most. + #[must_use] + pub fn capacity(&self) -> usize { + self.inner.capacity + } + /// Takes the ring lock, recovering from poisoning. /// /// The same posture [`crate::firehose`] takes: every operation under it is @@ -987,6 +1006,98 @@ mod tests { assert_eq!(seqs, vec![4u64, 5, 6]); } + /// The load-bearing one for a consumer that stops reading: it is *ended*, + /// not silently fast-forwarded. A stream that resumed after a gap would + /// leave a relay believing it had applied every commit when it had not, + /// and the gap would never be visible again. + #[tokio::test] + async fn a_stream_that_fell_behind_ends_rather_than_skipping_to_the_present() { + let repos = Repos::new(Arc::new(Sequence::in_memory()), 4); + let mut stream = repos + .subscribe(None, Exclusions::default()) + .expect("a live consumer"); + // Nothing is read while these are published, so the live channel + // overruns and the frames at the front of it are gone. + for _ in 0..12 { + repos.publish_commit(&commit("3lqaaa", None)); + } + assert_eq!( + stream.next().await, + None, + "a consumer that fell behind was handed a frame with a hole in front of it" + ); + assert!( + stream.lagged() > 0, + "the stream ended without recording that anything was dropped" + ); + } + + /// The writer's side of the same fact: publishing does not wait on a + /// subscriber that is not reading. `Repos` is driven from inside a write, + /// so a `publish` that blocked here would be a way for any consumer to + /// halt every write on the server. + #[tokio::test] + async fn publishing_does_not_wait_on_a_subscriber_that_never_reads() { + let repos = Repos::new(Arc::new(Sequence::in_memory()), 4); + let _idle = repos + .subscribe(None, Exclusions::default()) + .expect("a live consumer"); + let published = tokio::task::spawn_blocking({ + let repos = repos.clone(); + move || { + // Far past the channel's depth, from a thread with no + // reactor under it: if `publish` waited on the reader, this + // never returns. + (0..64) + .map(|_| repos.publish_commit(&commit("3lqaaa", None)).seq) + .collect::>() + } + }); + let published = tokio::time::timeout(std::time::Duration::from_secs(10), published) + .await + .expect("publishing blocked on a subscriber that was not reading") + .expect("the publishing thread did not panic"); + assert_eq!( + published, + (1..=64).collect::>(), + "every write was numbered, in order, with the reader stalled" + ); + assert_eq!(repos.all().len(), 4, "and the buffer stayed at its bound"); + } + + /// The exact edge of the replay window. A cursor naming the frame just + /// before the oldest one held is still reachable — the next frame the + /// consumer wants is the oldest one — and one number below that is not. + #[tokio::test] + async fn the_oldest_frame_held_is_reachable_and_the_one_before_it_is_not() { + let repos = Repos::new(Arc::new(Sequence::in_memory()), 3); + for _ in 0..6 { + repos.publish_commit(&commit("3lqaaa", None)); + } + let oldest = repos.oldest(); + assert_eq!(oldest, 4, "the buffer holds 4, 5 and 6"); + + let mut reachable = repos + .subscribe(Some(oldest as i64 - 1), Exclusions::default()) + .expect("not a future cursor"); + assert_eq!( + seq_of(&reachable.next().await.expect("a frame")), + oldest, + "a cursor one below the oldest frame held must replay it, not admit a gap" + ); + + let mut lost = repos + .subscribe(Some(oldest as i64 - 2), Exclusions::default()) + .expect("not a future cursor"); + let (header, body) = decode(&lost.next().await.expect("a frame")); + assert_eq!(field(&header, "t"), &Value::String("#info".to_owned())); + assert_eq!( + field(&body, "name"), + &Value::String(OUTDATED_CURSOR.to_owned()), + "a cursor two below the oldest frame held has a hole and must be told so" + ); + } + /// The load-bearing one: a consumer that resumes from its last cursor sees /// every frame after it exactly once, with nothing repeated and nothing /// skipped, including frames published while it was reconnecting.