use a2a_protocol_types::events::StreamResponse;
use a2a_protocol_types::task::TaskId;
use crate::handler::limits::HandlerLimits;
use crate::metrics::{push_outcome, Metrics};
use crate::push::{PushConfigStore, PushSender};
fn timeout_outcome(sender: &dyn PushSender, limits: &HandlerLimits) -> &'static str {
if sender
.max_delivery_duration()
.is_some_and(|wanted| wanted > limits.push_delivery_timeout)
{
push_outcome::TIMEOUT_TRUNCATED
} else {
push_outcome::TIMEOUT
}
}
pub(super) async fn deliver_push_bg(
task_id: &TaskId,
event: &StreamResponse,
push_config_store: &dyn PushConfigStore,
push_sender: Option<&dyn PushSender>,
limits: &HandlerLimits,
metrics: &dyn Metrics,
) {
let Some(sender) = push_sender else {
return;
};
let Ok(configs) = push_config_store.list(task_id.as_ref()).await else {
return;
};
let max_total_push_time = std::time::Duration::from_secs(30);
let deadline = tokio::time::Instant::now() + max_total_push_time;
for (reached, config) in configs.iter().enumerate() {
if tokio::time::Instant::now() >= deadline {
trace_warn!(
task_id = %task_id,
remaining_configs = configs.len() - reached,
"push delivery deadline exceeded; skipping remaining configs"
);
for _ in reached..configs.len() {
metrics.on_push_delivery(push_outcome::SKIPPED);
}
break;
}
let result = tokio::time::timeout(
limits.push_delivery_timeout,
sender.send(&config.url, event, config),
)
.await;
match result {
Ok(Err(_err)) => {
trace_warn!(
task_id = %task_id,
url = %config.url,
error = %_err,
"push notification delivery failed (background)"
);
metrics.on_push_delivery(push_outcome::FAILED);
}
Err(_) => {
let outcome = timeout_outcome(sender, limits);
trace_warn!(
task_id = %task_id,
url = %config.url,
outcome,
"push notification delivery timed out (background)"
);
metrics.on_push_delivery(outcome);
}
Ok(Ok(())) => metrics.on_push_delivery(push_outcome::DELIVERED),
}
}
}
#[cfg(test)]
mod tests;