use core::time::Duration;
use routers_network::Entry;
use serde::de::DeserializeOwned;
use tracing::{debug, warn};
use crate::bus::adapter::{AckHandle, Source};
use crate::lifecycle::Shutdown;
use crate::materializer::sink::{Applied, Sink};
use crate::metrics::Metrics;
use crate::protocol::output::CommittedOutput;
const NAK_BACKOFF: Duration = Duration::from_secs(1);
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Stats {
pub applied: u64,
pub duplicates: u64,
pub retracted: u64,
pub poison: u64,
pub errors: u64,
}
impl Stats {
fn record(&mut self, applied: Applied) {
match applied {
Applied::Duplicate => self.duplicates += 1,
Applied::Retracted { .. } => self.retracted += 1,
Applied::Inserted { .. }
| Applied::Superseded { .. }
| Applied::SegmentOpened
| Applied::Terminal => self.applied += 1,
}
}
}
fn applied_kind(applied: &Applied) -> &'static str {
match applied {
Applied::Inserted { .. } => "inserted",
Applied::Superseded { .. } => "superseded",
Applied::Duplicate => "duplicate",
Applied::Retracted { .. } => "retracted",
Applied::SegmentOpened => "segment_opened",
Applied::Terminal => "terminal",
}
}
pub async fn run<E, S, Src>(source: Src, sink: S, shutdown: Shutdown) -> anyhow::Result<Stats>
where
E: Entry + DeserializeOwned,
S: Sink<E>,
Src: Source<CommittedOutput<E>>,
{
run_with_metrics(source, sink, shutdown, &Metrics::noop()).await
}
pub async fn run_with_metrics<E, S, Src>(
mut source: Src,
sink: S,
shutdown: Shutdown,
metrics: &Metrics,
) -> anyhow::Result<Stats>
where
E: Entry + DeserializeOwned,
S: Sink<E>,
Src: Source<CommittedOutput<E>>,
{
let mut stats = Stats::default();
loop {
tokio::select! {
biased;
() = shutdown.triggered() => break,
next = source.next() => {
let delivery = match next {
None => break,
Some(Ok(delivery)) => delivery,
Some(Err(err)) => {
warn!(error = %err, "committed-output source failed; continuing");
stats.poison += 1;
continue;
}
};
match sink.apply(&delivery.item).await {
Ok(applied) => {
metrics.materialized(applied_kind(&applied));
stats.record(applied);
if let Err(err) = delivery.handle.ack().await {
warn!(error = %err, "could not acknowledge applied output");
} else {
debug!(?applied, "applied committed output");
}
}
Err(err) => {
warn!(error = %err, "sink failed to apply output; naking");
stats.errors += 1;
if let Err(nak_err) =
delivery.handle.nak(Some(NAK_BACKOFF)).await
{
warn!(error = %nak_err, "could not nak failed output");
}
}
}
}
}
}
Ok(stats)
}
#[cfg(test)]
mod tests {
use super::*;
use geo::Point;
use routers_network::mock::MockEntryId;
use routers_network::{DirectionAwareEdgeId, Edge};
use crate::bus::memory::MemoryBus;
use crate::event::{MatchedDiff, MatchedLayer, VehicleId};
use crate::lifecycle::{DrainReason, Shutdown};
use crate::materializer::memory::MemorySink;
use crate::materializer::sink::Applied;
use crate::protocol::ids::headers;
use crate::protocol::ids::{JobId, ObservationId, Revision, SegmentId};
use crate::protocol::output::{CommittedOutput, OutputKind};
use async_nats::HeaderMap;
type E = MockEntryId;
fn matched_job(job: u128, vehicle: u64, revision: u64, timestamp: i64) -> CommittedOutput<E> {
let diff = MatchedDiff {
revision,
downgraded: false,
layers: vec![MatchedLayer {
timestamp,
edge: Edge {
source: MockEntryId(1),
target: MockEntryId(2),
weight: 1,
id: DirectionAwareEdgeId::new(MockEntryId(3)),
},
position: Point::new(0.0, 0.0),
path: Vec::new(),
}],
};
CommittedOutput::new(
JobId(job),
VehicleId(vehicle),
ObservationId {
partition: 0,
sequence: revision,
},
Revision(revision),
SegmentId(1),
OutputKind::Matched {
diff,
finalized_through: None,
},
)
}
fn matched(vehicle: u64, revision: u64, timestamp: i64) -> CommittedOutput<E> {
matched_job(u128::from(revision), vehicle, revision, timestamp)
}
const FILTER: &str = "events.matched.v1.p.>";
async fn publish(bus: &MemoryBus, output: &CommittedOutput<E>) {
use crate::bus::adapter::Publisher;
bus.publisher::<CommittedOutput<E>>()
.publish(
&output.subject(),
&output.msg_id(),
HeaderMap::new(),
output,
)
.await
.expect("publish");
}
#[tokio::test]
async fn processes_and_acks_three_outputs() {
let bus = MemoryBus::new();
let source = bus.source::<CommittedOutput<E>>(FILTER);
let sink = MemorySink::<E>::new();
for vehicle in 1..=3 {
publish(&bus, &matched_job(u128::from(vehicle), vehicle, 5, 100)).await;
}
bus.close();
let stats = run::<E, _, _>(source, sink.clone(), Shutdown::new())
.await
.expect("run");
assert_eq!(stats.applied, 3);
assert_eq!(stats.duplicates, 0);
assert_eq!(bus.acked_count(FILTER), 3, "every output was acked");
assert_eq!(sink.len(), 3, "every vehicle materialised");
}
#[tokio::test]
async fn an_equal_revision_output_counts_as_a_duplicate() {
let bus = MemoryBus::new();
let source = bus.source::<CommittedOutput<E>>(FILTER);
let sink = MemorySink::<E>::new();
publish(&bus, &matched_job(1, 1, 5, 100)).await;
publish(&bus, &matched_job(2, 1, 5, 100)).await;
bus.close();
let stats = run::<E, _, _>(source, sink, Shutdown::new())
.await
.expect("run");
assert_eq!(stats.applied, 1);
assert_eq!(stats.duplicates, 1);
assert_eq!(bus.acked_count(FILTER), 2, "both were acked");
}
#[derive(Clone)]
struct FailingSink {
shutdown: Shutdown,
}
#[derive(Debug, thiserror::Error)]
#[error("sink is down")]
struct SinkDown;
impl Sink<E> for FailingSink {
type Error = SinkDown;
async fn apply(&self, _output: &CommittedOutput<E>) -> Result<Applied, SinkDown> {
self.shutdown.trigger(DrainReason::Operator);
Err(SinkDown)
}
}
#[tokio::test]
async fn a_failing_sink_naks_and_does_not_ack() {
let bus = MemoryBus::new();
let source = bus.source::<CommittedOutput<E>>(FILTER);
let shutdown = Shutdown::new();
let sink = FailingSink {
shutdown: shutdown.clone(),
};
publish(&bus, &matched(1, 5, 100)).await;
let stats = run::<E, _, _>(source, sink, shutdown).await.expect("run");
assert_eq!(stats.errors, 1, "the failed apply was counted");
assert_eq!(stats.applied, 0);
assert_eq!(bus.acked_count(FILTER), 0, "nothing was acked");
assert_eq!(
bus.nak_delays(),
vec![Some(NAK_BACKOFF)],
"the output was naked with the backoff",
);
}
#[tokio::test]
async fn shutdown_ends_the_loop() {
let bus = MemoryBus::new();
let source = bus.source::<CommittedOutput<E>>(FILTER);
let sink = MemorySink::<E>::new();
let shutdown = Shutdown::new();
let handle = {
let shutdown = shutdown.clone();
tokio::spawn(async move { run::<E, _, _>(source, sink, shutdown).await })
};
shutdown.trigger(DrainReason::Signal);
let stats = handle.await.expect("join").expect("run");
assert_eq!(stats, Stats::default(), "no work was done");
}
#[tokio::test]
async fn malformed_payload_is_acked_by_the_transport_before_materializing() {
use crate::bus::adapter::Publisher;
let bus = MemoryBus::new();
let source = bus.source::<CommittedOutput<E>>(FILTER);
let sink = MemorySink::<E>::new();
let mut poison_headers = HeaderMap::new();
headers::stamp_schema(&mut poison_headers);
bus.publisher::<CommittedOutput<E>>()
.publish_bytes(
&crate::topology::output::output_subject(0),
"poison-1",
poison_headers,
b"not a committed output",
)
.await
.expect("publish bytes");
publish(&bus, &matched(1, 5, 100)).await;
bus.close();
let stats = run::<E, _, _>(source, sink.clone(), Shutdown::new())
.await
.expect("run");
assert_eq!(
stats.poison, 0,
"transport poison does not leak into the reducer"
);
assert_eq!(stats.applied, 1, "the good output still applied");
assert_eq!(sink.len(), 1);
assert_eq!(
bus.acked_count(FILTER),
2,
"the poison and good output were retired"
);
}
}