import { StreamDataType } from '../../../types/stream'; import { DuressLevel } from '../../../types/quorum'; export interface DuressSnapshot { level: DuressLevel; score_ms: number; throttle_ms: number; per_type: Record; } /** * Adaptive engine duress detection via processing latency. * * ## Why this exists * * Prior fixes responded to queue *depth* (a symptom) — doubling reservation * timeouts and halving batch sizes when the stream backed up. A deep queue * doesn't necessarily mean duress (it could be a burst of external triggers), * and a shallow queue doesn't necessarily mean health. This module responds * to the *cause*: actual processing latency per message type. * * ## How it works * * Each engine router tracks an exponential moving average (EMA) of how long * each canonical message type (transition, timehook, webhook, worker response, * etc.) takes to process. When healthy, these are sub-50ms. When the max EMA * crosses configurable thresholds (200ms → mild, 1s → moderate, 5s → severe), * the manager computes a proportional throttle delay that the ThrottleManager * applies as a floor on engine consumption rate. * * ## Hysteresis (asymmetric by design) * * Escalation is immediate — if the engine suddenly enters duress, the throttle * kicks in on the next evaluation. De-escalation requires `HYSTERESIS_COUNT` * (default 3) consecutive improving evaluations before dropping a level. This * prevents oscillation: throttle → drain → un-throttle → refill → throttle. * The EMA already smooths individual outliers; hysteresis gates the recovery * path specifically. * * ## Quorum coordination * * When a router detects a level change (or remains in duress), it broadcasts * a `'duress'` message via the quorum. Peers adopt the signal only if it's * worse than their local state, so the mesh converges on the worst-case * throttle without coordination. * * ## What this does NOT do * * External messages (triggers, signalIn/webhooks from the outside world) are * never throttled. They always enter `engine_streams`. Only the engine * routers' pull rate slows down, giving the system breathing room. */ export declare class DuressManager { private emas; private sampleCounts; private currentLevel; private belowThresholdCount; private duressThrottle; private lastBroadcastAt; private lastBroadcastLevel; /** * Record a processing duration for a message type. * Updates the exponential moving average for that type. */ recordLatency(type: StreamDataType, durationMs: number): void; /** * Evaluate duress state from current EMAs. * Returns a snapshot with level, score, recommended throttle, * and per-type latencies. */ evaluate(): DuressSnapshot; getDuressThrottle(): number; getCurrentLevel(): DuressLevel; /** * Apply a duress snapshot received from another engine via quorum. * Adopts the remote signal only if it indicates worse duress than local. */ applyRemoteDuress(throttleMs: number, level: DuressLevel): void; /** * Whether a quorum broadcast is warranted. * Rate-limited and only fires when level changes or duress is active. */ shouldBroadcast(): boolean; markBroadcast(): void; /** * Returns a snapshot for inclusion in quorum rollcall profiles. */ getSnapshot(): DuressSnapshot; private scoreToLevel; private scoreToThrottle; private lerp; private levelOrdinal; }