use std::{
sync::{
Arc, Mutex, PoisonError,
atomic::{AtomicUsize, Ordering},
mpsc,
},
thread::JoinHandle,
time::Duration,
};
use polyc_state::{feed::JournalFeed, revision::JournalPosition};
use crate::feed::{
executor::FeedExecutor,
service::{ChunkSender, MAX_CONCURRENT_SUBSCRIPTIONS, Pump, StepOutcome, TAIL_POLL},
};
struct Parked {
pump: Pump,
sender: ChunkSender,
known: JournalPosition,
idle: Duration,
}
#[derive(Clone)]
pub(crate) struct Registry(Arc<Mutex<Vec<Parked>>>);
impl Registry {
fn new() -> Self {
Self(Arc::new(Mutex::new(Vec::new())))
}
pub(crate) fn park(self, pump: Pump, sender: ChunkSender) {
let known = pump.head_hint();
self.0
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push(Parked {
pump,
sender,
known,
idle: Duration::ZERO,
});
}
fn drain(&self) -> Vec<Parked> {
std::mem::take(&mut *self.0.lock().unwrap_or_else(PoisonError::into_inner))
}
fn restore(&self, entries: Vec<Parked>) {
if entries.is_empty() {
return;
}
self.0
.lock()
.unwrap_or_else(PoisonError::into_inner)
.extend(entries);
}
}
pub(crate) fn drive(mut pump: Pump, mut sender: ChunkSender, registry: Registry) {
loop {
match pump.step(&mut sender) {
StepOutcome::Sent => {}
StepOutcome::CaughtUp => {
registry.park(pump, sender);
return;
}
StepOutcome::Ended => return,
}
}
}
pub struct Tailer {
registry: Registry,
admitted: Arc<AtomicUsize>,
budget: usize,
stop: Option<mpsc::Sender<()>>,
thread: Option<JoinHandle<()>>,
}
impl std::fmt::Debug for Tailer {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Tailer").finish_non_exhaustive()
}
}
impl Tailer {
pub fn start_production(
feed: Arc<dyn JournalFeed>,
executor: Arc<FeedExecutor>,
) -> std::io::Result<Self> {
Self::start(feed, executor, MAX_CONCURRENT_SUBSCRIPTIONS)
}
pub fn start(
feed: Arc<dyn JournalFeed>,
executor: Arc<FeedExecutor>,
budget: usize,
) -> std::io::Result<Self> {
let registry = Registry::new();
let (stop, stopped) = mpsc::channel();
let loop_registry = registry.clone();
let thread = std::thread::Builder::new()
.name("polychrome-feed-tailer".to_owned())
.spawn(move || run(&feed, &loop_registry, &executor, &stopped))?;
Ok(Self {
registry,
admitted: Arc::new(AtomicUsize::new(0)),
budget,
stop: Some(stop),
thread: Some(thread),
})
}
pub(crate) fn registry(&self) -> Registry {
self.registry.clone()
}
#[must_use]
pub(crate) fn try_admit(&self) -> Option<SubscriptionPermit> {
self.admitted
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |count| {
(count < self.budget).then_some(count + 1)
})
.ok()
.map(|_| SubscriptionPermit(Arc::clone(&self.admitted)))
}
}
pub(crate) struct SubscriptionPermit(Arc<AtomicUsize>);
impl Drop for SubscriptionPermit {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::AcqRel);
}
}
#[cfg(test)]
impl SubscriptionPermit {
pub(crate) fn for_test() -> Self {
Self(Arc::new(AtomicUsize::new(1)))
}
}
impl Drop for Tailer {
fn drop(&mut self) {
if let Some(stop) = self.stop.take() {
let _ = stop.send(());
}
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
const HEADS_BATCH: usize = 128;
fn tick(
feed: &Arc<dyn JournalFeed>,
registry: &Registry,
executor: &Arc<FeedExecutor>,
elapsed: Duration,
) {
let parked = registry.drain();
if parked.is_empty() {
return;
}
let mut live = Vec::with_capacity(parked.len());
for mut entry in parked {
if entry.sender.is_closed() {
continue;
}
if entry.pump.is_draining() {
entry.pump.drain_now(&mut entry.sender);
continue;
}
live.push(entry);
}
if live.is_empty() {
return;
}
let mut still_parked = Vec::with_capacity(live.len());
while !live.is_empty() {
let batch_len = live.len().min(HEADS_BATCH);
let batch: Vec<Parked> = live.drain(..batch_len).collect();
let sources: Vec<_> = batch
.iter()
.map(|entry| entry.pump.source().clone())
.collect();
let heads = feed.heads(&sources);
debug_assert_eq!(
heads.len(),
batch.len(),
"JournalFeed::heads must answer one entry per source asked about, in the same order"
);
for (entry, (_, head)) in batch.into_iter().zip(heads) {
let Parked {
pump,
mut sender,
known,
idle,
} = entry;
let moved = head != Ok(known);
if moved {
let registry = registry.clone();
executor.spawn_pump(move || drive(pump, sender, registry));
continue;
}
let idle = idle.saturating_add(elapsed);
if idle >= pump.max_tail_idle() {
pump.pause_now(&mut sender);
continue;
}
still_parked.push(Parked {
pump,
sender,
known,
idle,
});
}
}
registry.restore(still_parked);
}
fn run(
feed: &Arc<dyn JournalFeed>,
registry: &Registry,
executor: &Arc<FeedExecutor>,
stopped: &mpsc::Receiver<()>,
) {
let mut previous = std::time::Instant::now();
loop {
match stopped.recv_timeout(TAIL_POLL) {
Ok(()) | Err(mpsc::RecvTimeoutError::Disconnected) => return,
Err(mpsc::RecvTimeoutError::Timeout) => {}
}
let now = std::time::Instant::now();
let elapsed = now.duration_since(previous);
previous = now;
tick(feed, registry, executor, elapsed);
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use polyc_state::memory::MemoryFeed;
use super::{FeedExecutor, Tailer};
fn tailer_over(budget: usize) -> Tailer {
let feed = Arc::new(MemoryFeed::new());
let executor = Arc::new(FeedExecutor::start(2).expect("start the executor"));
Tailer::start(feed, executor, budget).expect("start the tailer")
}
#[test]
fn admits_up_to_the_budget_then_refuses() {
let tailer = tailer_over(2);
let first = tailer.try_admit();
assert!(first.is_some());
let second = tailer.try_admit();
assert!(second.is_some());
assert!(
tailer.try_admit().is_none(),
"a third subscription past a budget of two must be refused"
);
drop(first);
drop(second);
}
#[test]
fn dropping_a_permit_frees_its_slot() {
let tailer = tailer_over(1);
let permit = tailer.try_admit();
assert!(permit.is_some());
assert!(
tailer.try_admit().is_none(),
"already at the one-subscription budget"
);
drop(permit);
assert!(
tailer.try_admit().is_some(),
"the freed slot must be admittable again"
);
}
}
#[cfg(test)]
mod proofs {
use std::{
sync::{
Arc, Mutex as StdMutex,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
time::Duration,
};
use polyc_state::{
context::CallContext,
error::StateError,
feed::{
AcknowledgeProjectorCursor, CompactFeedPrefix, CreateSnapshot, FeedChunk,
FeedCompaction, FeedCursor, FeedRetention, FeedSnapshot, GetFeedRetention,
GetProjectorStatus, JournalFeed, ListProjectors, ProjectorListing, ProjectorStatus,
RegisterProjector, SubscribeCommits,
},
id::PartitionId,
receipt::Receipt,
revision::{JournalPosition, JournalSource, PartitionIncarnation},
stream::StreamEnd,
};
use super::{FeedExecutor, Registry, TAIL_POLL, tick};
use crate::feed::service::{FeedStreamTuning, Pump, PumpedChunk};
fn source(n: usize) -> JournalSource {
JournalSource::new(
PartitionId::new(format!("tailer-proof-{n}")),
PartitionIncarnation::from_bytes([1; PartitionIncarnation::LEN]),
)
}
struct FakeFeed {
heads_calls: AtomicUsize,
commits_calls: AtomicUsize,
positions: StdMutex<Vec<JournalPosition>>,
commit_delay: StdMutex<Vec<Duration>>,
completed: StdMutex<Vec<bool>>,
snapshot_at_delay: StdMutex<Option<Vec<bool>>>,
}
impl FakeFeed {
fn new(sources: usize) -> Self {
Self {
heads_calls: AtomicUsize::new(0),
commits_calls: AtomicUsize::new(0),
positions: StdMutex::new(vec![JournalPosition::ORIGIN; sources]),
commit_delay: StdMutex::new(vec![Duration::ZERO; sources]),
completed: StdMutex::new(vec![false; sources]),
snapshot_at_delay: StdMutex::new(None),
}
}
fn snapshot_at_delay_within(&self, timeout: Duration) -> Vec<bool> {
let deadline = std::time::Instant::now() + timeout;
loop {
let snapshot = self.snapshot_at_delay.lock().unwrap().clone();
if let Some(snapshot) = snapshot {
return snapshot;
}
assert!(
std::time::Instant::now() < deadline,
"the deliberately delayed source never completed within the hang-guard"
);
std::thread::sleep(Duration::from_millis(2));
}
}
fn advance(&self, n: usize, position: JournalPosition) {
self.positions.lock().unwrap()[n] = position;
}
fn delay(&self, n: usize, delay: Duration) {
self.commit_delay.lock().unwrap()[n] = delay;
}
fn index_of(source: &JournalSource) -> usize {
source
.partition()
.as_str()
.strip_prefix("tailer-proof-")
.expect("every source this fake serves is one it minted")
.parse()
.expect("the index is the whole suffix")
}
}
impl JournalFeed for FakeFeed {
fn create_snapshot(
&self,
_command: CreateSnapshot,
_context: &CallContext,
) -> Result<FeedSnapshot, StateError> {
unimplemented!("not exercised by a tail-loop proof")
}
fn commits(
&self,
request: SubscribeCommits,
_context: &CallContext,
) -> Result<FeedChunk, StateError> {
self.commits_calls.fetch_add(1, Ordering::SeqCst);
let index = Self::index_of(request.source());
let delay = self.commit_delay.lock().unwrap()[index];
if delay > Duration::ZERO {
std::thread::sleep(delay);
}
let mut completed = self.completed.lock().unwrap();
completed[index] = true;
if delay > Duration::ZERO {
*self.snapshot_at_delay.lock().unwrap() = Some(completed.clone());
}
drop(completed);
Ok(FeedChunk::new(
Vec::new(),
FeedCursor::origin(request.source().clone()),
StreamEnd::More,
))
}
fn heads(
&self,
sources: &[JournalSource],
) -> Vec<(JournalSource, Result<JournalPosition, StateError>)> {
self.heads_calls.fetch_add(1, Ordering::SeqCst);
let positions = self.positions.lock().unwrap();
sources
.iter()
.map(|source| {
let index = Self::index_of(source);
(source.clone(), Ok(positions[index]))
})
.collect()
}
fn register_projector(
&self,
_command: RegisterProjector,
_context: &CallContext,
) -> Result<Receipt, StateError> {
unimplemented!("not exercised by a tail-loop proof")
}
fn acknowledge(
&self,
_command: AcknowledgeProjectorCursor,
_context: &CallContext,
) -> Result<Receipt, StateError> {
unimplemented!("not exercised by a tail-loop proof")
}
fn projector_status(
&self,
_request: GetProjectorStatus,
_context: &CallContext,
) -> Result<Option<ProjectorStatus>, StateError> {
unimplemented!("not exercised by a tail-loop proof")
}
fn projectors(
&self,
_request: ListProjectors,
_context: &CallContext,
) -> Result<ProjectorListing, StateError> {
unimplemented!("not exercised by a tail-loop proof")
}
fn retention(
&self,
_request: GetFeedRetention,
_context: &CallContext,
) -> Result<FeedRetention, StateError> {
unimplemented!("not exercised by a tail-loop proof")
}
fn compact(
&self,
_command: CompactFeedPrefix,
_context: &CallContext,
) -> Result<FeedCompaction, StateError> {
unimplemented!("not exercised by a tail-loop proof")
}
}
fn park(
registry: &Registry,
feed: &Arc<dyn JournalFeed>,
draining: &Arc<AtomicBool>,
n: usize,
) -> futures::channel::mpsc::Receiver<PumpedChunk> {
let (sender, receiver) = futures::channel::mpsc::channel(8);
let pump = Pump::for_proof(
Arc::clone(feed),
Arc::clone(draining),
source(n),
FeedStreamTuning::new(),
);
registry.clone().park(pump, sender);
receiver
}
fn recv_within(
receiver: &mut futures::channel::mpsc::Receiver<PumpedChunk>,
timeout: Duration,
) -> bool {
let deadline = std::time::Instant::now() + timeout;
loop {
if receiver.try_recv().is_ok() {
return true;
}
if std::time::Instant::now() >= deadline {
return false;
}
std::thread::sleep(Duration::from_millis(2));
}
}
#[test]
fn a_caught_up_partition_costs_no_durable_read() {
const PARKED: usize = 1_000;
const TICKS: usize = 20;
let feed = Arc::new(FakeFeed::new(PARKED));
let dyn_feed: Arc<dyn JournalFeed> = feed.clone();
let executor = Arc::new(FeedExecutor::start(4).expect("start the executor"));
let registry = Registry::new();
let draining = Arc::new(AtomicBool::new(false));
let _receivers: Vec<_> = (0..PARKED)
.map(|n| park(®istry, &dyn_feed, &draining, n))
.collect();
for _ in 0..TICKS {
tick(&dyn_feed, ®istry, &executor, TAIL_POLL);
}
assert_eq!(
feed.commits_calls.load(Ordering::SeqCst),
0,
"nothing moved; commits must never run"
);
assert_eq!(
feed.heads_calls.load(Ordering::SeqCst),
TICKS * PARKED.div_ceil(super::HEADS_BATCH),
"heads() batches by HEADS_BATCH, not once per parked subscription and not once for \
all 1,000 at once"
);
}
#[test]
fn only_the_moved_source_is_dispatched() {
const PARKED: usize = 1_000;
const MOVED: usize = 507;
let feed = Arc::new(FakeFeed::new(PARKED));
let dyn_feed: Arc<dyn JournalFeed> = feed.clone();
let executor = Arc::new(FeedExecutor::start(4).expect("start the executor"));
let registry = Registry::new();
let draining = Arc::new(AtomicBool::new(false));
let _receivers: Vec<_> = (0..PARKED)
.map(|n| park(®istry, &dyn_feed, &draining, n))
.collect();
tick(&dyn_feed, ®istry, &executor, TAIL_POLL);
assert_eq!(
feed.commits_calls.load(Ordering::SeqCst),
0,
"nothing moved on the first tick"
);
feed.advance(MOVED, JournalPosition::new(1));
tick(&dyn_feed, ®istry, &executor, TAIL_POLL);
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while feed.commits_calls.load(Ordering::SeqCst) == 0 && std::time::Instant::now() < deadline
{
std::thread::sleep(Duration::from_millis(2));
}
assert_eq!(
feed.commits_calls.load(Ordering::SeqCst),
1,
"only the one moved source's pump is dispatched"
);
}
#[test]
fn a_drain_ends_every_parked_stream_in_one_tick() {
const PARKED: usize = 1_000;
let feed = Arc::new(FakeFeed::new(PARKED));
let dyn_feed: Arc<dyn JournalFeed> = feed.clone();
let executor = Arc::new(FeedExecutor::start(4).expect("start the executor"));
let registry = Registry::new();
let draining = Arc::new(AtomicBool::new(false));
let receivers: Vec<_> = (0..PARKED)
.map(|n| park(®istry, &dyn_feed, &draining, n))
.collect();
draining.store(true, Ordering::SeqCst);
tick(&dyn_feed, ®istry, &executor, TAIL_POLL);
assert_eq!(
registry.drain().len(),
0,
"every parked entry left the registry in one tick"
);
assert_eq!(
feed.commits_calls.load(Ordering::SeqCst),
0,
"a drain costs no durable read"
);
for mut receiver in receivers {
assert!(
recv_within(&mut receiver, Duration::from_secs(1)),
"every parked subscription receives its drain marker"
);
}
}
#[test]
fn a_slow_source_does_not_delay_the_others() {
const PARKED: usize = 10;
const SLOW: usize = 3;
const SLOW_DELAY: Duration = Duration::from_millis(2_000);
const SAFETY_TIMEOUT: Duration = Duration::from_secs(10);
let feed = Arc::new(FakeFeed::new(PARKED));
let dyn_feed: Arc<dyn JournalFeed> = feed.clone();
let executor = Arc::new(FeedExecutor::start(PARKED + 2).expect("start the executor"));
let registry = Registry::new();
let draining = Arc::new(AtomicBool::new(false));
let _receivers: Vec<_> = (0..PARKED)
.map(|n| {
feed.advance(n, JournalPosition::new(1));
park(®istry, &dyn_feed, &draining, n)
})
.collect();
feed.delay(SLOW, SLOW_DELAY);
std::thread::spawn({
let dyn_feed = Arc::clone(&dyn_feed);
move || tick(&dyn_feed, ®istry, &executor, TAIL_POLL)
});
let snapshot = feed.snapshot_at_delay_within(SAFETY_TIMEOUT);
for (n, completed) in snapshot.iter().enumerate().take(PARKED) {
if n == SLOW {
continue;
}
assert!(
*completed,
"source {n} had not completed by the instant the slow source's own \
{SLOW_DELAY:?} delay did, so it must have been queued behind it \
instead of running concurrently"
);
}
}
}