use std::sync::Arc;
use zygo_core::supervisor::Response as Reply;
use super::Api;
#[derive(Default)]
pub(super) struct Usage {
totals: std::collections::BTreeMap<String, crate::cmd::otlp::TenantUsage>,
queued: std::collections::VecDeque<zygo_core::pool::Usage>,
dropped: u64,
}
const USAGE_QUEUE: usize = 10_000;
const USAGE_BATCH: usize = 256;
impl Usage {
fn record(&mut self, usage: zygo_core::pool::Usage) {
let totals = self.totals.entry(usage.tenant.clone()).or_default();
totals.tenant = usage.tenant.clone();
totals.requests += 1;
if usage.outcome != "ok" {
totals.failures += 1;
}
totals.cpu_ms += usage.cpu_ms;
totals.wall_ms += usage.wall_ms;
*totals.by_outcome.entry(usage.outcome.clone()).or_default() += 1;
if self.queued.len() >= USAGE_QUEUE {
self.queued.pop_front();
self.dropped += 1;
}
self.queued.push_back(usage);
}
pub(super) fn snapshot(&self) -> Vec<crate::cmd::otlp::TenantUsage> {
self.totals.values().cloned().collect()
}
fn take_batch(&mut self) -> Vec<zygo_core::pool::Usage> {
self.queued
.drain(..USAGE_BATCH.min(self.queued.len()))
.collect()
}
fn return_batch(&mut self, batch: Vec<zygo_core::pool::Usage>) {
for usage in batch.into_iter().rev() {
if self.queued.len() >= USAGE_QUEUE {
self.dropped += 1;
continue;
}
self.queued.push_front(usage);
}
}
}
pub(super) async fn deliver_usage(api: Arc<Api>, url: reqwest::Url, interval: std::time::Duration) {
let client = match reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(30))
.build()
{
Ok(client) => client,
Err(e) => {
tracing::error!("usage: cannot build an HTTP client: {e}");
return;
}
};
let mut ticker = tokio::time::interval(interval);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut failing = false;
loop {
ticker.tick().await;
loop {
let batch = api.usage.lock().expect("usage").take_batch();
if batch.is_empty() {
break;
}
let body = serde_json::json!({ "events": batch });
match client.post(url.clone()).json(&body).send().await {
Ok(response) if response.status().is_success() => {
if failing {
tracing::info!("usage: the webhook is answering again");
failing = false;
}
}
outcome => {
api.usage.lock().expect("usage").return_batch(batch);
if !failing {
failing = true;
match outcome {
Ok(r) => tracing::warn!("usage: the webhook answered {}", r.status()),
Err(e) => tracing::warn!("usage: the webhook is unreachable: {e}"),
}
}
break;
}
}
}
let dropped = {
let mut usage = api.usage.lock().expect("usage");
std::mem::take(&mut usage.dropped)
};
if dropped > 0 {
tracing::warn!(
"usage: dropped {dropped} events; the queue holds {USAGE_QUEUE} and the \
webhook is behind"
);
}
}
}
pub(super) fn count_usage(api: &Api, reply: &Reply) {
if let Reply::Executed { outcome } = reply {
api.usage
.lock()
.expect("usage")
.record(zygo_core::pool::Usage::from(outcome.as_ref()));
}
}