use core::time::Duration;
use anyhow::Context as _;
use async_nats::jetstream;
use crate::partition::PARTITIONS;
pub mod jobs;
pub mod output;
pub mod raw;
pub mod results;
pub use jobs::*;
pub use output::*;
pub use raw::*;
pub use results::*;
pub const DUPLICATE_WINDOW: Duration = Duration::from_secs(2 * 60);
pub fn duplicate_window(max_age: Duration) -> Duration {
if max_age.is_zero() {
DUPLICATE_WINDOW
} else {
DUPLICATE_WINDOW.min(max_age)
}
}
pub fn partition_of_subject(subject: &str) -> Option<u16> {
let token = subject.rsplit_once(".p.")?.1;
if token.is_empty() || !token.bytes().all(|byte| byte.is_ascii_digit()) {
return None;
}
let partition: u16 = token.parse().ok()?;
(u64::from(partition) < PARTITIONS).then_some(partition)
}
async fn create_or_update_stream(
context: &jetstream::Context,
config: jetstream::stream::Config,
) -> anyhow::Result<jetstream::stream::Stream> {
let name = config.name.clone();
match context.update_stream(&config).await {
Ok(_) => context
.get_stream(&name)
.await
.with_context(|| format!("could not open stream {name}")),
Err(_) => context
.create_stream(config)
.await
.with_context(|| format!("could not create stream {name}")),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn duplicate_window_fits_inside_max_age() {
assert_eq!(
duplicate_window(Duration::from_secs(60)),
Duration::from_secs(60)
);
assert_eq!(duplicate_window(Duration::from_secs(600)), DUPLICATE_WINDOW);
assert_eq!(duplicate_window(Duration::ZERO), DUPLICATE_WINDOW);
}
#[test]
fn partition_of_subject_reads_every_partitioned_plane() {
for partition in [0u64, 1, 485, PARTITIONS - 1] {
let expect = Some(partition as u16);
assert_eq!(partition_of_subject(&raw::raw_subject(partition)), expect);
assert_eq!(
partition_of_subject(&results::result_subject(partition)),
expect
);
assert_eq!(
partition_of_subject(&output::output_subject(partition)),
expect
);
}
}
#[test]
fn partition_of_subject_rejects_the_unaddressable() {
let rejected = [
"solve.v1.g.europe.r.syd.q.0", "events.raw.p", "events.raw.p.", "events.raw.p.1024", "events.raw.p.99999", "events.raw.p.3.oops", "events.raw.p.-1", "events.raw.p.+1", "events.raw.p.0x1", "", ];
for subject in rejected {
assert_eq!(
partition_of_subject(subject),
None,
"{subject:?} should not resolve to a partition"
);
}
}
}