use super::*;
#[tokio::test(start_paused = true)]
async fn configs_the_deadline_never_reaches_are_counted_as_skipped() {
use std::sync::atomic::Ordering;
use std::time::Duration;
let task_id = TaskId::new("t-skip");
let event = make_status_event("t-skip", TaskState::Working);
let store = store_with_configs("t-skip", 100).await;
let metrics = CountingMetrics::default();
let limits = HandlerLimits::default().with_push_delivery_timeout(Duration::from_secs(5));
deliver_push_bg(
&task_id,
&event,
&store,
Some(&SlowPushSender),
&limits,
&metrics,
)
.await;
let delivered = metrics.delivered.load(Ordering::Relaxed);
let skipped = metrics.skipped.load(Ordering::Relaxed);
assert_eq!(
delivered + skipped,
100,
"every config must be accounted for: {delivered} delivered, {skipped} skipped"
);
assert!(
skipped > 0,
"the 30s budget cannot reach 100 configs at 5s each; nothing was reported skipped"
);
assert_eq!(
delivered, 6,
"budget / per-delivery cost is the reachable count"
);
assert_eq!(skipped, 94, "the rest are skipped, and are now counted");
assert_eq!(
metrics.other.load(Ordering::Relaxed),
0,
"no delivery failed or timed out in this scenario"
);
}
struct WantsLongerThanAllowed;
impl crate::push::PushSender for WantsLongerThanAllowed {
fn send<'a>(
&'a self,
_url: &'a str,
_event: &'a StreamResponse,
_config: &'a TaskPushNotificationConfig,
) -> Pin<Box<dyn Future<Output = a2a_protocol_types::error::A2aResult<()>> + Send + 'a>> {
Box::pin(async {
tokio::time::sleep(std::time::Duration::from_secs(93)).await;
Ok(())
})
}
fn max_delivery_duration(&self) -> Option<std::time::Duration> {
Some(std::time::Duration::from_secs(93))
}
}
struct SaysNothing;
impl crate::push::PushSender for SaysNothing {
fn send<'a>(
&'a self,
_url: &'a str,
_event: &'a StreamResponse,
_config: &'a TaskPushNotificationConfig,
) -> Pin<Box<dyn Future<Output = a2a_protocol_types::error::A2aResult<()>> + Send + 'a>> {
Box::pin(async {
tokio::time::sleep(std::time::Duration::from_secs(93)).await;
Ok(())
})
}
}
#[tokio::test(start_paused = true)]
async fn a_sender_cut_short_by_the_handler_bound_is_reported_as_truncated() {
use std::sync::atomic::Ordering;
use std::time::Duration;
let task_id = TaskId::new("t-trunc");
let event = make_status_event("t-trunc", TaskState::Working);
let store = store_with_configs("t-trunc", 1).await;
let limits = HandlerLimits::default().with_push_delivery_timeout(Duration::from_secs(5));
let declared = CountingMetrics::default();
deliver_push_bg(
&task_id,
&event,
&store,
Some(&WantsLongerThanAllowed),
&limits,
&declared,
)
.await;
assert_eq!(
declared.truncated.load(Ordering::Relaxed),
1,
"a sender that declared 93s against a 5s bound was cut short"
);
assert_eq!(
declared.timed_out.load(Ordering::Relaxed),
0,
"and must not also be reported as an ordinary timeout"
);
let silent = CountingMetrics::default();
deliver_push_bg(
&task_id,
&event,
&store,
Some(&SaysNothing),
&limits,
&silent,
)
.await;
assert_eq!(
silent.timed_out.load(Ordering::Relaxed),
1,
"a sender that says nothing about its schedule is a plain timeout"
);
assert_eq!(
silent.truncated.load(Ordering::Relaxed),
0,
"nothing is known to have been truncated, so nothing is claimed"
);
}