use core::time::Duration;
use anyhow::Context as _;
use async_nats::jetstream::{
self,
consumer::{AckPolicy, PullConsumer, pull},
stream::{Config, RetentionPolicy, StorageType},
};
use super::{create_or_update_stream, duplicate_window};
use crate::protocol::ids::{GraphVersion, Lane, RegionId};
pub const JOB_PREFIX: &str = "solve.v1";
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct JobsConfig {
pub max_age: Duration,
pub max_deliver: i64,
pub ack_wait: Duration,
pub max_ack_pending: i64,
}
impl Default for JobsConfig {
fn default() -> Self {
Self {
max_age: Duration::ZERO,
max_deliver: -1,
ack_wait: Duration::from_secs(30),
max_ack_pending: 64,
}
}
}
pub fn job_subject(graph: &GraphVersion, region: &RegionId, lane: Lane) -> String {
format!("{JOB_PREFIX}.g.{graph}.r.{region}.q.{}", lane.0)
}
pub fn job_stream_name(region: &RegionId) -> String {
format!("SOLVE-JOBS-{region}")
}
pub fn job_stream_subjects(region: &RegionId) -> String {
format!("{JOB_PREFIX}.g.*.r.{region}.q.*")
}
pub fn job_consumer_name(graph: &GraphVersion, region: &RegionId) -> String {
format!("matchers-g{graph}-r{region}")
}
pub fn job_consumer_filter(graph: &GraphVersion, region: &RegionId) -> String {
format!("{JOB_PREFIX}.g.{graph}.r.{region}.q.>")
}
pub fn parse_job_subject(subject: &str) -> Option<(GraphVersion, RegionId, Lane)> {
let parts: [&str; 8] = subject.split('.').collect::<Vec<_>>().try_into().ok()?;
match parts {
["solve", "v1", "g", graph, "r", region, "q", lane] => {
let graph = GraphVersion::new(graph).ok()?;
let region = RegionId::new(region).ok()?;
let lane = lane.parse::<u8>().ok().map(Lane)?;
Some((graph, region, lane))
}
_ => None,
}
}
fn job_stream_config(region: &RegionId, config: &JobsConfig) -> Config {
Config {
name: job_stream_name(region),
subjects: vec![job_stream_subjects(region)],
retention: RetentionPolicy::WorkQueue,
storage: StorageType::File,
max_age: config.max_age,
duplicate_window: duplicate_window(config.max_age),
..Default::default()
}
}
pub async fn ensure_job_stream(
context: &jetstream::Context,
region: &RegionId,
config: &JobsConfig,
) -> anyhow::Result<jetstream::stream::Stream> {
create_or_update_stream(context, job_stream_config(region, config))
.await
.with_context(|| format!("could not reconcile job stream for region {region}"))
}
pub async fn open_job_stream(
context: &jetstream::Context,
region: &RegionId,
config: &JobsConfig,
) -> anyhow::Result<jetstream::stream::Stream> {
context
.get_or_create_stream(job_stream_config(region, config))
.await
.with_context(|| format!("could not open job stream for region {region}"))
}
pub async fn job_consumer(
stream: &jetstream::stream::Stream,
graph: &GraphVersion,
region: &RegionId,
config: &JobsConfig,
) -> anyhow::Result<PullConsumer> {
let name = job_consumer_name(graph, region);
stream
.get_or_create_consumer(
&name,
pull::Config {
durable_name: Some(name.clone()),
filter_subject: job_consumer_filter(graph, region),
ack_policy: AckPolicy::Explicit,
max_deliver: config.max_deliver,
ack_wait: config.ack_wait,
max_ack_pending: config.max_ack_pending,
..Default::default()
},
)
.await
.with_context(|| format!("could not create job consumer {name}"))
}
#[cfg(test)]
mod tests {
use super::*;
fn graph() -> GraphVersion {
GraphVersion::new("europe-2026-09").unwrap()
}
fn region() -> RegionId {
RegionId::new("syd").unwrap()
}
#[test]
fn subject_and_names_format() {
assert_eq!(
job_subject(&graph(), ®ion(), Lane(0)),
"solve.v1.g.europe-2026-09.r.syd.q.0"
);
assert_eq!(
job_subject(&graph(), ®ion(), Lane(7)),
"solve.v1.g.europe-2026-09.r.syd.q.7"
);
assert_eq!(job_stream_name(®ion()), "SOLVE-JOBS-syd");
assert_eq!(job_stream_subjects(®ion()), "solve.v1.g.*.r.syd.q.*");
assert_eq!(
job_consumer_name(&graph(), ®ion()),
"matchers-geurope-2026-09-rsyd"
);
assert_eq!(
job_consumer_filter(&graph(), ®ion()),
"solve.v1.g.europe-2026-09.r.syd.q.>"
);
}
#[test]
fn defaults_retain_unanswered_work_for_redelivery() {
let cfg = JobsConfig::default();
assert_eq!(cfg.max_age, Duration::ZERO);
assert!(cfg.max_deliver < 0, "negative means unlimited delivery");
assert!(cfg.max_ack_pending > 0, "claiming remains capacity bounded");
}
#[test]
fn subject_round_trips_through_parser() {
for lane in [0u8, 1, 9, 255] {
let subject = job_subject(&graph(), ®ion(), Lane(lane));
assert_eq!(
parse_job_subject(&subject),
Some((graph(), region(), Lane(lane)))
);
}
}
#[test]
fn parser_rejects_the_malformed() {
let rejected = [
"events.raw.p.3", "solve.v1.g.europe.r.syd.q.>", "solve.v1.g.europe.r.syd.q.-1", "solve.v1.g.europe.r.syd.q.256", "solve.v1.g.europe.r.syd.q", "solve.v1.g.europe.r.syd.q.0.extra", "solve.v2.g.europe.r.syd.q.0", "solve.v1.x.europe.r.syd.q.0", "", ];
for subject in rejected {
assert_eq!(
parse_job_subject(subject),
None,
"{subject:?} should not parse"
);
}
}
#[test]
fn default_delivery_window_matches_the_regional_capacity_envelope() {
assert_eq!(JobsConfig::default().max_ack_pending, 64);
}
}