use std::time::{Duration, Instant};
use crate::channels::whatsapp::rate_limit::{
Decision, LimiterState, QueuedSend, WhatsappRateLimiter, admit, drain_ready,
};
use crate::config::types::WaRateLimitConfig;
fn cfg(per_minute: u32, daily_cap: u32) -> WaRateLimitConfig {
WaRateLimitConfig {
messages_per_minute: per_minute,
daily_cap,
}
}
#[test]
fn burst_of_fifty_paces_without_drops() {
let cfg = WaRateLimitConfig::default(); let mut state = LimiterState::default();
let t0 = Instant::now();
let mut now = t0;
let mut delivered = 0;
let mut paced = 0;
while delivered < 50 {
match admit(&mut state, &cfg, now, false) {
Decision::SendNow => delivered += 1,
Decision::Pace(wait) => {
assert!(wait > Duration::ZERO, "pace must actually wait");
assert!(
wait <= Duration::from_secs(2),
"refill gap at 30/min is at most 2s, got {:?}",
wait
);
now += wait;
paced += 1;
}
Decision::Queue { .. } => panic!("burst of 50 must never hit the 800/day cap"),
}
}
assert_eq!(delivered, 50);
assert_eq!(
paced, 20,
"bucket holds 30; sends 31-50 pace one refill gap each"
);
assert!(
now - t0 >= Duration::from_secs(38),
"20 waits of ~2s must stretch the burst to ~40s, got {:?}",
now - t0
);
}
#[test]
fn daily_cap_queues_with_exactly_one_alert() {
let cfg = cfg(0, 5); let mut state = LimiterState::default();
let now = Instant::now();
for _ in 0..5 {
assert_eq!(admit(&mut state, &cfg, now, false), Decision::SendNow);
}
assert_eq!(
admit(&mut state, &cfg, now, false),
Decision::Queue { alert: true },
"first over-cap send raises the one alert"
);
assert_eq!(
admit(&mut state, &cfg, now, false),
Decision::Queue { alert: false }
);
assert_eq!(
admit(&mut state, &cfg, now, false),
Decision::Queue { alert: false }
);
assert!(state.saturated());
}
#[test]
fn window_slide_resets_episode_and_drains_fifo() {
let cfg = cfg(0, 2);
let mut state = LimiterState::default();
let t0 = Instant::now();
assert_eq!(admit(&mut state, &cfg, t0, false), Decision::SendNow);
assert_eq!(admit(&mut state, &cfg, t0, false), Decision::SendNow);
assert_eq!(
admit(&mut state, &cfg, t0, false),
Decision::Queue { alert: true }
);
state.queue.push_back(QueuedSend {
jid: "a@s.whatsapp.net".into(),
text: "first".into(),
});
state.queue.push_back(QueuedSend {
jid: "b@s.whatsapp.net".into(),
text: "second".into(),
});
assert!(state.saturated(), "episode in flight before the slide");
let later = t0 + Duration::from_secs(24 * 3600 + 1);
let drained = drain_ready(&mut state, &cfg, later);
assert_eq!(
drained.len(),
2,
"both queued sends flush once the window empties"
);
assert_eq!(drained[0].text, "first");
assert_eq!(drained[1].text, "second");
assert!(!state.saturated(), "episode cleared by the slide");
assert_eq!(
admit(&mut state, &cfg, later, false),
Decision::Queue { alert: true }
);
}
#[test]
fn owner_bypasses_a_saturated_budget() {
let cfg = cfg(1, 1);
let mut state = LimiterState::default();
let now = Instant::now();
assert_eq!(admit(&mut state, &cfg, now, false), Decision::SendNow);
assert_eq!(
admit(&mut state, &cfg, now, false),
Decision::Queue { alert: true }
);
for _ in 0..5 {
assert_eq!(admit(&mut state, &cfg, now, true), Decision::SendNow);
}
assert_eq!(state.queue_len(), 0);
assert_eq!(
admit(&mut state, &cfg, now, false),
Decision::Queue { alert: false }
);
}
#[test]
fn zero_knobs_disable_limiting() {
let cfg = cfg(0, 0);
let mut state = LimiterState::default();
let now = Instant::now();
for _ in 0..1000 {
assert_eq!(admit(&mut state, &cfg, now, false), Decision::SendNow);
}
}
#[test]
fn drain_ready_respects_budget() {
let cfg = cfg(0, 2);
let mut state = LimiterState::default();
for i in 0..3 {
state.queue.push_back(QueuedSend {
jid: "x@s.whatsapp.net".into(),
text: format!("m{i}"),
});
}
let drained = drain_ready(&mut state, &cfg, Instant::now());
assert_eq!(drained.len(), 2, "cap 2 allows two flushes");
assert_eq!(drained[0].text, "m0");
assert_eq!(drained[1].text, "m1");
assert_eq!(state.queue_len(), 1, "m2 waits for the next window slide");
}
#[test]
fn config_serde_defaults_and_overrides() {
let empty: WaRateLimitConfig = toml::from_str("").unwrap();
assert_eq!(empty.messages_per_minute, 30);
assert_eq!(empty.daily_cap, 800);
let partial: WaRateLimitConfig = toml::from_str("messages_per_minute = 10").unwrap();
assert_eq!(partial.messages_per_minute, 10);
assert_eq!(partial.daily_cap, 800);
let full: WaRateLimitConfig =
toml::from_str("messages_per_minute = 5\ndaily_cap = 100").unwrap();
assert_eq!(
full,
WaRateLimitConfig {
messages_per_minute: 5,
daily_cap: 100
}
);
let wa: crate::config::types::WhatsAppConfig =
toml::from_str("enabled = true\n[rate_limit]\nmessages_per_minute = 7").unwrap();
assert_eq!(wa.rate_limit.messages_per_minute, 7);
assert_eq!(wa.rate_limit.daily_cap, 800);
}
#[tokio::test]
async fn gate_queues_and_reports_position() {
use crate::channels::whatsapp::rate_limit::{GateOutcome, WhatsappRateLimiter};
let limiter = WhatsappRateLimiter::new();
let cfg = cfg(0, 1);
let jid = "25512345678@s.whatsapp.net";
assert_eq!(
limiter.gate(&cfg, jid, "first", false).await,
GateOutcome::SendNow
);
assert_eq!(
limiter.gate(&cfg, jid, "second", false).await,
GateOutcome::Queued { position: 1 }
);
assert_eq!(limiter.queue_len().await, 1);
assert!(
limiter.alert_due(),
"first queue of the episode owes an alert"
);
limiter.mark_alert_sent();
assert!(!limiter.alert_due(), "exactly once: alert delivered");
assert_eq!(
limiter.gate(&cfg, jid, "to-owner", true).await,
GateOutcome::SendNow
);
}
#[tokio::test]
async fn ephemeral_gate_drops_instead_of_queueing() {
let limiter = WhatsappRateLimiter::new();
let cfg = cfg(0, 1);
assert!(
limiter.gate_ephemeral(&cfg, false).await,
"first send within cap"
);
assert!(
!limiter.gate_ephemeral(&cfg, false).await,
"saturated: ephemeral sends drop"
);
assert_eq!(limiter.queue_len().await, 0, "drops must not queue");
assert!(
limiter.alert_due(),
"saturation still arms the one-time alert"
);
assert!(limiter.gate_ephemeral(&cfg, true).await, "owner bypasses");
}
#[test]
fn agent_output_paths_gate_through_the_limiter() {
const TOOL: &str = include_str!("../brain/tools/whatsapp_send.rs");
const HANDLER: &str = include_str!("../channels/whatsapp/handler.rs");
const RESUME: &str = include_str!("../channels/whatsapp/resume.rs");
const AGENT: &str = include_str!("../channels/whatsapp/agent.rs");
assert_eq!(
TOOL.matches(".gate(&rl_cfg").count(),
2,
"tool send and reply arms must gate per chunk"
);
assert!(
HANDLER.contains(".gate_ephemeral(&rl_cfg_c"),
"streaming intermediates must use the ephemeral gate"
);
assert!(
HANDLER.contains(".gate(&wa_cfg.rate_limit"),
"final-text chunks must gate"
);
assert!(
RESUME.contains(".gate(&wa_cfg.rate_limit"),
"bg-resume send must gate"
);
assert!(
AGENT.contains("spawn_drainer("),
"agent start must spawn the drainer exactly once"
);
}