use std::sync::Arc;
use crate::core::{RunId, Secret, StoreError};
use crate::push::{DueBatch, PushConfig, PushNamespace, PushRegistration, PushStore};
#[derive(Debug)]
pub struct DefaultPaged(pub Arc<dyn PushStore>);
#[async_trait::async_trait]
impl PushStore for DefaultPaged {
async fn put(&self, config: &PushConfig, next_seq: u64) -> Result<(), StoreError> {
self.0.put(config, next_seq).await
}
async fn get(&self, task: RunId, id: &str) -> Result<Option<PushConfig>, StoreError> {
self.0.get(task, id).await
}
async fn list(&self, task: RunId) -> Result<Vec<PushConfig>, StoreError> {
self.0.list(task).await
}
async fn due(&self, at: u64, limit: usize) -> Result<Vec<PushRegistration>, StoreError> {
self.0.due(at, limit).await
}
async fn advance(&self, task: RunId, id: &str, next_seq: u64) -> Result<(), StoreError> {
self.0.advance(task, id, next_seq).await
}
async fn retry(
&self,
task: RunId,
id: &str,
next_attempt_at: u64,
error: &str,
) -> Result<(), StoreError> {
self.0.retry(task, id, next_attempt_at, error).await
}
async fn park(&self, task: RunId, id: &str, error: &str) -> Result<(), StoreError> {
self.0.park(task, id, error).await
}
async fn parked(&self, limit: usize) -> Result<Vec<PushRegistration>, StoreError> {
self.0.parked(limit).await
}
async fn unpark(&self, task: RunId, id: &str, at: u64) -> Result<bool, StoreError> {
self.0.unpark(task, id, at).await
}
async fn delete(&self, task: RunId, id: &str) -> Result<(), StoreError> {
self.0.delete(task, id).await
}
}
fn config(task: RunId, id: &str) -> PushConfig {
PushConfig {
id: id.to_owned(),
task,
url: "https://hooks.acme.example/a2a".to_owned(),
token: Some(Secret::new("opaque")),
authentication: None,
}
}
fn shape(rows: &[PushRegistration]) -> Vec<(RunId, String, u64, u32, u64)> {
rows.iter()
.map(|registration| {
(
registration.config.task,
registration.config.id.clone(),
registration.next_seq,
registration.attempts,
registration.next_attempt_at,
)
})
.collect()
}
async fn both(
store: &Arc<dyn PushStore>,
at: u64,
limit: usize,
namespace: PushNamespace,
) -> (DueBatch, DueBatch) {
let native = store.due_in(at, limit, namespace).await.expect("native");
let paged = DefaultPaged(Arc::clone(store))
.due_in(at, limit, namespace)
.await
.expect("default");
(native, paged)
}
pub async fn pin_due_in_against_the_default(store: Arc<dyn PushStore>) {
let (t1, t2) = (RunId::generate(), RunId::generate());
for (task, id, retry_at) in [
(t1, "hook-a", None),
(t1, "hook-b", Some(5)),
(t2, "hook-c", Some(10)),
(t2, "hook-late", Some(200)),
(t1, "operator:bus", None),
(t2, "operator:audit", Some(7)),
(t2, "operator:late", Some(100)),
] {
store.put(&config(task, id), 1).await.expect("put");
if let Some(at) = retry_at {
store.retry(task, id, at, "staggered").await.expect("retry");
}
}
for namespace in [PushNamespace::Caller, PushNamespace::Operator] {
let (native, paged) = both(&store, 20, 10, namespace).await;
assert_eq!(
shape(&native.rows),
shape(&paged.rows),
"the native due_in and the trait default disagree about which rows \
{namespace:?} owns"
);
assert_eq!(
native.unserved, paged.unserved,
"the native due_in and the trait default disagree about the \
foreign backlog visible past an exhausted read"
);
let (native, paged) = both(&store, 20, 1, namespace).await;
assert_eq!(
shape(&native.rows),
shape(&paged.rows),
"under a truncating limit the native due_in serves different rows \
than the default would"
);
assert!(
native.unserved >= paged.unserved,
"the native override reported a smaller foreign backlog than the \
default saw, so the lower bound shrank ({} < {})",
native.unserved,
paged.unserved
);
let (native, paged) = both(&store, 0, 10, namespace).await;
assert_eq!(shape(&native.rows), shape(&paged.rows));
assert_eq!(
native.rows.len(),
1,
"the due instant stopped filtering: {namespace:?} saw rows whose \
retry is still in the future"
);
assert_eq!(native.unserved, paged.unserved);
assert_eq!(native.unserved, 1);
}
}