///|
/// The optional delivery engine: deduplication, rate limiting and retry with
/// backoff around the plain `deliver` function. Pure and testable — every
/// timestamp comes from the caller, and outcomes are returned as values
/// instead of being logged, so the host application decides what to record.
///|
struct DedupeEntry {
first_ms : UInt64
count : Int
}
///|
struct Queued {
message : Message
/// Owner key for per-group queue accounting (e.g. one project's events).
owner : String
targets : Array[Target]
timeout_ms : Int
}
///|
const MAX_ATTEMPTS_DEFAULT : Int = 3
///|
pub struct Engine {
ctx : Context
max_attempts : Int
mut queue : Array[Queued]
mut dropped : Int
/// Owner → (window started ms, emitted in window). Each owner gets its own
/// rolling minute budget, so one project's throttle never starves another.
rate_window : Map[String, (UInt64, Int)]
dedupe : Map[String, DedupeEntry]
mut next_attempt_ms : UInt64
mut attempts : Int
}
///|
/// Create an engine. `ctx` brands the outgoing messages (footer and EHLO
/// name); `max_attempts` is how many delivery passes one message gets before
/// it is dropped.
pub fn Engine::new(
ctx : Context,
max_attempts? : Int = MAX_ATTEMPTS_DEFAULT,
) -> Engine {
{
ctx,
max_attempts,
queue: [],
dropped: 0,
rate_window: Map([]),
dedupe: Map([]),
next_attempt_ms: 0UL,
attempts: 0,
}
}
///|
/// What happened to one emit call.
pub(all) enum EmitOutcome {
Queued
/// Same condition inside the dedupe window; `count` is how many times the
/// key has been seen in the current window.
Deduped(count~ : Int)
RateLimited
QueueFull
} derive(Debug, Eq)
///|
pub extend EmitOutcome with Eq::{not_equal, equal}
///|
pub extend EmitOutcome with @debug.Debug::{to_repr}
///|
/// Queue one message. Deduplication, rate limiting and queue accounting are
/// all driven by the parameters — the engine keeps no policy of its own.
///
/// - `key` deduplicates identical conditions inside `dedupe_window_s`.
/// - `owner` + `queue_limit` keep a burst from one group from starving others:
/// once `queue_limit` messages of `owner` are queued, further ones drop.
/// - `rate_per_minute` caps emissions per rolling 60s window (0 = unlimited).
pub fn Engine::emit(
self : Engine,
now : UInt64,
key : String,
owner : String,
msg : Message,
targets : Array[Target],
timeout_ms : Int,
dedupe_window_s? : Int = 0,
rate_per_minute? : Int = 0,
queue_limit? : Int = 0,
) -> EmitOutcome {
if targets.length() == 0 {
return EmitOutcome::Queued
}
let mut message = msg
if dedupe_window_s > 0 {
let window = (dedupe_window_s * 1000).to_uint64()
match self.dedupe.get(key) {
Some(entry) =>
if now >= entry.first_ms && now - entry.first_ms < window {
// Same condition inside the window: count it, stay silent.
self.dedupe[key] = {
first_ms: entry.first_ms,
count: entry.count + 1,
}
return Deduped(count=entry.count + 1)
} else {
// The window expired: send again and say how quiet it has been.
if entry.count > 1 {
message = {
..message,
body: message.body +
"\n(上一个窗口内同样的事件共 \{entry.count} 次)",
}
}
self.dedupe[key] = { first_ms: now, count: 1, }
}
None => self.dedupe[key] = { first_ms: now, count: 1, }
}
}
if rate_per_minute > 0 {
// Per-owner rolling minute window: one project's throttle never counts
// another project's emissions.
let (started, sent) = match self.rate_window.get(owner) {
Some(window) => window
None => (0UL, 0)
}
let (started, sent) = if now >= started && now - started >= 60000UL {
(now, 0)
} else {
(started, sent)
}
if sent >= rate_per_minute {
self.dropped += 1
self.rate_window[owner] = (started, sent)
return RateLimited
}
self.rate_window[owner] = (started, sent + 1)
}
let queued = self.queue.filter(fn(notice) { notice.owner == owner }).length()
if queue_limit > 0 && queued >= queue_limit {
// Keep the first events of a burst: they carry the root cause.
self.dropped += 1
return QueueFull
}
self.queue.push({ message, owner, targets, timeout_ms, })
EmitOutcome::Queued
}
///|
/// What one pump pass did with the head of the queue.
pub(all) enum PumpOutcome {
/// Queue empty or waiting out the retry backoff.
Idle
Delivered
/// Some targets got it, some failed; retrying would duplicate for the
/// ones that succeeded, so the head is dropped either way.
Partial
/// Every target failed; `attempt` counts this failure. The message is
/// retried after a backoff until `max_attempts` is reached.
Failed(attempt~ : Int)
/// Every target failed and the attempt budget is spent; the head dropped.
Exhausted
} derive(Debug, Eq)
///|
pub extend PumpOutcome with Eq::{not_equal, equal}
///|
pub extend PumpOutcome with @debug.Debug::{to_repr}
///|
pub(all) struct PumpReport {
remaining : Int
outcome : PumpOutcome
/// Owner of the attempted message, so the host application can attribute
/// the outcome in its own logs. None when idle.
owner : String?
} derive(Debug)
///|
pub extend PumpReport with Debug::{to_repr}
///|
/// One pass of the engine: deliver at most one queued message.
pub async fn Engine::pump(self : Engine, now : UInt64) -> PumpReport {
if self.queue.length() == 0 || now < self.next_attempt_ms {
return { remaining: self.queue.length(), outcome: Idle, owner: None, }
}
let notice = self.queue[0]
let results = deliver_all(
notice.targets,
notice.message,
notice.timeout_ms,
self.ctx,
)
let delivered = results.filter(fn(pair) {
let (_, failure) = pair
failure is None
})
let failures = results.filter(fn(pair) {
let (_, failure) = pair
failure is Some(_)
})
let outcome = if failures.length() == 0 {
self.drop_head()
self.attempts = 0
self.next_attempt_ms = 0UL
Delivered
} else if delivered.length() > 0 {
// Somebody got it; retrying now would only duplicate it for them.
self.drop_head()
self.attempts = 0
self.next_attempt_ms = 0UL
Partial
} else {
self.attempts += 1
if self.attempts >= self.max_attempts {
self.drop_head()
self.attempts = 0
self.next_attempt_ms = 0UL
Exhausted
} else {
// Backoff grows linearly: 5s after the first failure, 15s after the
// second — long enough to clear a blip, short enough for an alert.
self.next_attempt_ms = now +
(5000 * (1 + 2 * (self.attempts - 1))).to_uint64()
Failed(attempt=self.attempts)
}
}
{ remaining: self.queue.length(), outcome, owner: Some(notice.owner), }
}
///|
fn Engine::drop_head(self : Engine) -> Unit {
self.queue = self.queue[1:].to_owned()
}
///|
/// Messages dropped since the engine was created: queue overflow plus rate
/// limiting.
pub fn Engine::dropped(self : Engine) -> Int {
self.dropped
}
///|
/// Number of messages still queued.
pub fn Engine::queued(self : Engine) -> Int {
self.queue.length()
}