1use core::fmt;
7use std::collections::{BTreeMap, HashMap, VecDeque};
8use std::sync::{Arc, Mutex, MutexGuard, RwLock};
9use std::time::Duration;
10
11use keel_core_api::policy::{
12 BreakerMode, BreakerPolicy, CacheScope, DurationMs, JournalLocation, NondeterminismResponse,
13 Policy, Rate, ResolvedPolicy, RetryPolicy,
14};
15use keel_core_api::{
16 AttemptResult, BreakerState, ENVELOPE_VERSION, ErrorClass, ErrorCode, KeelError, Outcome,
17 OutcomeError, Request,
18};
19use keel_journal::{
20 CacheKey as JournalCacheKey, CallObservation, CallResult, Clock, DiscoveryStore, Journal,
21 ObservedError,
22};
23use serde::{Deserialize, Serialize};
24use serde_json::Value;
25use tokio::time::Instant;
26use tracing::{Instrument, debug, warn};
27
28use crate::events::{CacheStore, EventKind, EventSink, TraceRef};
29use crate::journal_backend::{self, JournalBackend};
30
31#[derive(Debug, Default)]
36struct Breaker {
37 consecutive: u64,
39 outcomes: VecDeque<(Instant, bool)>,
43 open_until: Option<Instant>,
44 opens: u64,
45}
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49enum Admission {
50 Closed,
51 HalfOpen,
53 Rejected,
55}
56
57#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61enum BreakerTransition {
62 None,
64 Opened,
66 Closed,
68}
69
70impl Breaker {
71 fn admit(&self, now: Instant) -> Admission {
72 match self.open_until {
73 Some(until) if now < until => Admission::Rejected,
74 Some(_) => Admission::HalfOpen,
75 None => Admission::Closed,
76 }
77 }
78
79 fn state_at(&self, now: Instant) -> BreakerState {
80 match self.admit(now) {
81 Admission::Rejected => BreakerState::Open,
82 Admission::Closed | Admission::HalfOpen => BreakerState::Closed,
83 }
84 }
85
86 fn on_success(&mut self, now: Instant, config: &BreakerPolicy) -> BreakerTransition {
87 let closed_a_probe = self.open_until.is_some();
91 self.consecutive = 0;
92 self.open_until = None;
93 if closed_a_probe {
94 self.outcomes.clear();
97 return BreakerTransition::Closed;
98 }
99 if let BreakerMode::Rate { window, .. } = config.mode() {
100 self.observe(now, window, false);
101 }
102 BreakerTransition::None
103 }
104
105 fn on_terminal_failure(
106 &mut self,
107 now: Instant,
108 config: &BreakerPolicy,
109 admission: Admission,
110 ) -> BreakerTransition {
111 let should_trip = if admission == Admission::HalfOpen {
112 true } else {
114 match config.mode() {
115 BreakerMode::Count { failures } => {
116 self.consecutive += 1;
117 self.consecutive >= failures.get()
118 }
119 BreakerMode::Rate {
120 window,
121 failure_rate,
122 min_calls,
123 } => {
124 self.observe(now, window, true);
125 self.window_rate_reached(failure_rate, min_calls)
126 }
127 }
128 };
129 if should_trip {
130 self.open_until = Some(now + Duration::from_millis(config.cooldown.0));
131 self.opens += 1;
132 self.consecutive = 0;
133 self.outcomes.clear();
134 BreakerTransition::Opened
135 } else {
136 BreakerTransition::None
137 }
138 }
139
140 fn observe(&mut self, now: Instant, window: DurationMs, failed: bool) {
144 let window = Duration::from_millis(window.0);
145 while let Some(&(at, _)) = self.outcomes.front() {
146 if now.duration_since(at) >= window {
147 self.outcomes.pop_front();
148 } else {
149 break;
150 }
151 }
152 self.outcomes.push_back((now, failed));
153 }
154
155 fn window_rate_reached(&self, failure_rate: f64, min_calls: core::num::NonZeroU32) -> bool {
157 let total = self.outcomes.len();
158 if (total as u64) < u64::from(min_calls.get()) {
159 return false;
160 }
161 let failed = self.outcomes.iter().filter(|&&(_, f)| f).count();
162 #[expect(
163 clippy::cast_precision_loss,
164 reason = "window counts are bounded by the calls observed within one \
165 breaker window — far below f64's 2^53 exact-integer range"
166 )]
167 let rate = failed as f64 / total as f64;
168 rate >= failure_rate
169 }
170}
171
172#[derive(Debug, Default)]
185struct TokenBucket {
186 scaled_tokens: i128,
190 last_refill_ms: u64,
192 primed: bool,
195}
196
197impl TokenBucket {
198 fn plan_admit(&mut self, elapsed_ms: u64, rate: Rate) -> u64 {
202 let limit = i128::from(rate.limit.get());
203 let window = i128::from(rate.window_ms);
204 let capacity = limit * window; if !self.primed {
206 self.primed = true;
207 self.scaled_tokens = capacity;
208 self.last_refill_ms = elapsed_ms;
209 }
210 let elapsed = i128::from(elapsed_ms.saturating_sub(self.last_refill_ms));
213 self.last_refill_ms = self.last_refill_ms.max(elapsed_ms);
214 self.scaled_tokens = capacity.min(
215 self.scaled_tokens
216 .saturating_add(elapsed.saturating_mul(limit)),
217 );
218 self.scaled_tokens -= window;
219 if self.scaled_tokens >= 0 {
220 0
221 } else {
222 let deficit = -self.scaled_tokens;
224 u64::try_from((deficit + limit - 1) / limit).unwrap_or(u64::MAX)
225 }
226 }
227}
228
229#[derive(Debug, Clone, PartialEq, Eq, Hash)]
230struct CacheKey {
231 target: String,
232 args_hash: String,
233}
234
235#[derive(Debug, Clone)]
236struct CacheEntry {
237 expires_at: Instant,
238 payload: Value,
239}
240
241const CACHE_PAYLOAD_SCHEMA: &str = "keel.cache/v1";
246
247#[derive(Serialize)]
250struct CachePayloadRef<'a> {
251 schema: &'a str,
252 payload: &'a Value,
253}
254
255#[derive(Deserialize)]
257struct CachePayloadOwned {
258 schema: String,
259 payload: Value,
260}
261
262fn encode_cache_payload(payload: &Value) -> Result<Vec<u8>, rmp_serde::encode::Error> {
264 rmp_serde::to_vec_named(&CachePayloadRef {
265 schema: CACHE_PAYLOAD_SCHEMA,
266 payload,
267 })
268}
269
270fn decode_cache_payload(bytes: &[u8]) -> Result<Value, String> {
274 let envelope: CachePayloadOwned =
275 rmp_serde::from_slice(bytes).map_err(|e| format!("messagepack decode failed: {e}"))?;
276 if envelope.schema != CACHE_PAYLOAD_SCHEMA {
277 return Err(format!(
278 "unrecognized cache payload schema {:?}",
279 envelope.schema
280 ));
281 }
282 Ok(envelope.payload)
283}
284
285#[derive(Debug)]
290enum CachePlan {
291 None,
292 Memory {
293 key: CacheKey,
294 },
295 Persistent {
296 key: JournalCacheKey,
297 ttl: DurationMs,
298 },
299}
300
301pub trait DiscoveryRecorder: Send + Sync {
305 fn record(&self, observation: &CallObservation) -> keel_journal::Result<()>;
307}
308
309impl<C: Clock> DiscoveryRecorder for DiscoveryStore<C> {
310 fn record(&self, observation: &CallObservation) -> keel_journal::Result<()> {
311 DiscoveryStore::record(self, observation)
312 }
313}
314
315#[derive(Debug, Default)]
316struct TargetMetrics {
317 calls: u64,
318 attempts: u64,
319 retries: u64,
320 successes: u64,
321 failures: u64,
322 cache_hits: u64,
323 throttled: u64,
324}
325
326#[derive(Debug, Serialize)]
329struct TargetReport {
330 attempts: u64,
331 breaker_opens: u64,
332 breaker_state: BreakerState,
333 cache_hits: u64,
334 calls: u64,
335 failures: u64,
336 retries: u64,
337 successes: u64,
338 throttled: u64,
339}
340
341#[derive(Debug, Serialize)]
342struct Report<'a> {
343 v: u32,
344 clock_ms: u64,
345 targets: BTreeMap<&'a str, TargetReport>,
346}
347
348#[derive(Debug, Default)]
349struct State {
350 trace_seq: u64,
351 breakers: HashMap<String, Breaker>,
352 rate_buckets: HashMap<String, TokenBucket>,
353 cache: HashMap<CacheKey, CacheEntry>,
354 metrics: BTreeMap<String, TargetMetrics>,
355}
356
357impl State {
358 fn metrics_for(&mut self, target: &str) -> &mut TargetMetrics {
359 self.metrics.entry(target.to_owned()).or_default()
360 }
361
362 fn breaker_state(&self, target: &str, now: Instant) -> BreakerState {
363 self.breakers
364 .get(target)
365 .map_or(BreakerState::Closed, |b| b.state_at(now))
366 }
367}
368
369#[derive(Debug)]
373struct AttemptOutcome {
374 result: AttemptResult,
375 timed_out_by_layer: bool,
376}
377
378fn terminal_code(
381 retryable: bool,
382 attempt: u32,
383 max_attempts: u32,
384 idempotent: bool,
385) -> Option<ErrorCode> {
386 if !retryable {
387 Some(ErrorCode::NonRetryableError)
388 } else if attempt == max_attempts {
389 Some(ErrorCode::AttemptsExhausted)
390 } else if !idempotent {
391 Some(ErrorCode::NonIdempotentNotRetried)
392 } else {
393 None
394 }
395}
396
397#[derive(Debug, Clone, Copy, PartialEq, Eq)]
400enum PollVerdict {
401 Terminal,
402 Pending,
403 FailOpen,
404}
405
406fn poll_verdict(poll: &keel_core_api::policy::PollPolicy, payload: &Value) -> PollVerdict {
410 use base64::Engine as _;
411 let Some(obj) = payload.as_object() else {
412 return PollVerdict::FailOpen;
413 };
414 let doc_owned;
415 let doc = if obj.contains_key("body_b64")
416 || (obj.contains_key("status") && obj.contains_key("headers"))
417 {
418 let Some(b64) = obj.get("body_b64").and_then(Value::as_str) else {
419 return PollVerdict::FailOpen;
420 };
421 let Ok(bytes) = base64::engine::general_purpose::STANDARD.decode(b64) else {
422 return PollVerdict::FailOpen;
423 };
424 let Ok(parsed) = serde_json::from_slice::<Value>(&bytes) else {
425 return PollVerdict::FailOpen;
426 };
427 if !parsed.is_object() {
428 return PollVerdict::FailOpen;
429 }
430 doc_owned = parsed;
431 doc_owned.as_object().expect("checked object")
432 } else {
433 obj
434 };
435 match doc.get(&poll.until.field) {
436 None => PollVerdict::FailOpen,
437 Some(Value::String(s)) if poll.until.terminal.iter().any(|t| t == s) => {
438 PollVerdict::Terminal
439 }
440 Some(_) => PollVerdict::Pending,
441 }
442}
443
444fn poll_applies(request: &Request) -> bool {
447 request.idempotent && (request.op.starts_with("GET ") || request.op.starts_with("HEAD "))
448}
449
450fn class_str(class: ErrorClass) -> &'static str {
451 match class {
452 ErrorClass::Conn => "conn",
453 ErrorClass::Timeout => "timeout",
454 ErrorClass::Http => "http",
455 ErrorClass::Cancelled => "cancelled",
456 ErrorClass::Other => "other",
457 }
458}
459
460fn breaker_str(state: BreakerState) -> &'static str {
463 match state {
464 BreakerState::Closed => "closed",
465 BreakerState::Open => "open",
466 BreakerState::HalfOpen => "half_open",
467 }
468}
469
470fn record_call_fields(span: &tracing::Span, out: &Outcome) {
475 span.record("trace_id", out.trace_id.as_str());
476 span.record("result", out.result.as_str());
477 if let Some(error) = out.error.as_ref() {
478 span.record("error_code", error.code.as_str());
479 }
480 span.record("attempts", out.attempts);
481 span.record("from_cache", out.from_cache);
482 span.record("throttled", out.throttled);
483 span.record("breaker", breaker_str(out.breaker));
484}
485
486fn emit_breaker_transition(target: &str, transition: BreakerTransition) {
490 match transition {
491 BreakerTransition::Opened => {
492 debug!(target = %target, transition = "opened", "breaker transition");
493 crate::metrics::record_breaker_transition(target, "opened");
494 }
495 BreakerTransition::Closed => {
496 debug!(target = %target, transition = "closed", "breaker transition");
497 crate::metrics::record_breaker_transition(target, "closed");
498 }
499 BreakerTransition::None => {}
500 }
501}
502
503fn warn_inert_breaker_knobs(policy: &Policy) {
508 let defaults = &policy.defaults;
509 let tables = defaults
510 .outbound
511 .iter()
512 .map(|t| (String::from("defaults.outbound"), t))
513 .chain(
514 defaults
515 .llm
516 .iter()
517 .map(|t| (String::from("defaults.llm"), t)),
518 )
519 .chain(
520 policy
521 .target
522 .iter()
523 .map(|(name, t)| (format!("target.\"{name}\""), t)),
524 );
525 for (path, table) in tables {
526 if table
527 .breaker
528 .as_ref()
529 .is_some_and(BreakerPolicy::has_inert_rate_knobs)
530 {
531 warn!(
532 "policy {path}.breaker sets `failures` (count mode) alongside rate-mode knobs \
533 (window/failure_rate/min_calls), which are inert in count mode. Remove \
534 `failures` to select rate mode."
535 );
536 }
537 }
538}
539
540fn with_trace_ref(message: String, trace: Option<&TraceRef>) -> String {
545 match trace {
546 Some(t) => {
547 let sep = if message.ends_with('.') { "" } else { "." };
548 format!("{message}{sep} trace: keel trace {t}")
549 }
550 None => message,
551 }
552}
553
554fn terminal_message(
555 code: ErrorCode,
556 request: &Request,
557 attempt: u32,
558 max_attempts: u32,
559 class: ErrorClass,
560 http_status: Option<u16>,
561 message: &str,
562) -> String {
563 let detail = match http_status {
564 Some(status) => format!("{} {status}", class_str(class)),
565 None => class_str(class).to_owned(),
566 };
567 let text = match code {
568 ErrorCode::Timeout => format!(
569 "{} exceeded its policy timeout on attempt {attempt}/{max_attempts}. {message}",
570 request.op
571 ),
572 ErrorCode::AttemptsExhausted => format!(
573 "{} failed {attempt}/{max_attempts} attempts (last: {detail}). {message}",
574 request.op
575 ),
576 ErrorCode::NonIdempotentNotRetried => format!(
577 "{} failed ({detail}). Not retried: call is not idempotent — observed, not retried. {message}",
578 request.op
579 ),
580 _ => format!(
581 "{} failed ({detail}); error class is not retryable per policy. {message}",
582 request.op
583 ),
584 };
585 text.trim_end().to_owned()
586}
587
588struct TerminalAttemptCtx<'a> {
592 request: &'a Request,
593 attempt: u32,
594 max_attempts: u32,
595 trace: Option<&'a TraceRef>,
596 timed_out_by_layer: bool,
601}
602
603fn terminal_attempt_error(
608 ctx: &TerminalAttemptCtx<'_>,
609 code: ErrorCode,
610 class: ErrorClass,
611 http_status: Option<u16>,
612 message: &str,
613 original: Option<Value>,
614) -> OutcomeError {
615 let code = if ctx.timed_out_by_layer && code != ErrorCode::NonIdempotentNotRetried {
616 ErrorCode::Timeout
617 } else {
618 code
619 };
620 OutcomeError {
621 code,
622 class,
623 http_status,
624 message: with_trace_ref(
625 terminal_message(
626 code,
627 ctx.request,
628 ctx.attempt,
629 ctx.max_attempts,
630 class,
631 http_status,
632 message,
633 ),
634 ctx.trace,
635 ),
636 original,
637 }
638}
639
640pub struct Engine {
646 started: Instant,
647 policy: RwLock<Policy>,
648 state: Mutex<State>,
649 journal: RwLock<JournalSlot>,
654 discovery: Option<Arc<dyn DiscoveryRecorder>>,
656 events: Option<EventSink>,
660}
661
662#[derive(Default)]
667struct JournalSlot {
668 journal: Option<Arc<dyn Journal>>,
669 backend: Option<JournalBackend>,
672}
673
674impl fmt::Debug for Engine {
675 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
676 f.debug_struct("Engine")
678 .field("policy", &self.policy)
679 .field("state", &self.state)
680 .field("journal_attached", &self.current_journal().is_some())
681 .field("discovery_attached", &self.discovery.is_some())
682 .field("events_attached", &self.events.is_some())
683 .finish_non_exhaustive()
684 }
685}
686
687impl Default for Engine {
688 fn default() -> Self {
689 Self::new()
690 }
691}
692
693impl Engine {
694 #[must_use]
695 pub fn new() -> Self {
696 Self {
697 started: Instant::now(),
698 policy: RwLock::new(Policy::default()),
699 state: Mutex::new(State::default()),
700 journal: RwLock::new(JournalSlot::default()),
701 discovery: None,
702 events: EventSink::from_env(),
706 }
707 }
708
709 pub fn attach_journal(&mut self, journal: impl Journal + 'static) -> &mut Self {
715 let slot = self.journal.get_mut().expect("journal lock poisoned");
716 slot.journal = Some(Arc::new(journal));
717 slot.backend = None;
718 self
719 }
720
721 #[must_use]
729 pub fn journal(&self) -> Option<Arc<dyn Journal>> {
730 self.current_journal()
731 }
732
733 fn current_journal(&self) -> Option<Arc<dyn Journal>> {
736 self.journal
737 .read()
738 .expect("journal lock poisoned")
739 .journal
740 .clone()
741 }
742
743 pub fn attach_discovery(&mut self, discovery: impl DiscoveryRecorder + 'static) -> &mut Self {
746 self.discovery = Some(Arc::new(discovery));
747 self
748 }
749
750 pub fn attach_events(&mut self, sink: EventSink) -> &mut Self {
754 self.events = Some(sink);
755 self
756 }
757
758 #[must_use]
761 pub fn events(&self) -> Option<&EventSink> {
762 self.events.as_ref()
763 }
764
765 fn emit_event(&self, kind: impl FnOnce() -> EventKind) {
769 if let Some(sink) = self.events.as_ref() {
770 sink.emit(self.elapsed_ms(), kind());
771 }
772 }
773
774 fn emit_cache(&self, out: &Outcome, target: &str, scope: CacheStore, hit: bool) {
777 self.emit_event(|| {
778 let call = out.trace_id.clone();
779 let target = target.to_owned();
780 if hit {
781 EventKind::CacheHit {
782 call,
783 target,
784 scope,
785 }
786 } else {
787 EventKind::CacheMiss {
788 call,
789 target,
790 scope,
791 }
792 }
793 });
794 }
795
796 pub fn configure(&self, policy_json: &Value) -> Result<(), KeelError> {
802 let policy: Policy =
803 serde_path_to_error::deserialize(policy_json).map_err(|e| KeelError {
804 code: ErrorCode::PolicyInvalid,
805 message: format!("policy invalid at {}: {}", e.path(), e.inner()),
806 })?;
807 if let Some(location) = &policy.journal {
812 self.apply_journal_location(location)?;
813 }
814 if let Some(telemetry) = &policy.telemetry
822 && !telemetry.console
823 {
824 warn!(
825 "policy `telemetry.console = false` is validated but not yet wired: v0.1 always \
826 uses the default local summary. `telemetry.otlp_endpoint` IS honored by front \
827 ends built with the `otel` feature."
828 );
829 }
830 warn_inert_breaker_knobs(&policy);
831 *self.policy.write().expect("policy lock poisoned") = policy;
832 Ok(())
833 }
834
835 #[must_use]
841 pub fn telemetry_otlp_endpoint(&self) -> Option<String> {
842 self.policy
843 .read()
844 .expect("policy lock poisoned")
845 .telemetry
846 .as_ref()
847 .and_then(|t| t.otlp_endpoint.clone())
848 }
849
850 fn apply_journal_location(&self, location: &JournalLocation) -> Result<(), KeelError> {
863 let backend = JournalBackend::select(location);
864 {
865 let slot = self.journal.read().expect("journal lock poisoned");
866 if slot.backend.as_ref() == Some(&backend) {
867 return Ok(()); }
869 }
870 let journal = journal_backend::open(&backend)?;
875 if let JournalBackend::File(path) = &backend {
876 debug!(path = %path.display(), "journal selected by policy");
877 }
878 let mut slot = self.journal.write().expect("journal lock poisoned");
879 slot.journal = Some(journal);
880 slot.backend = Some(backend);
881 Ok(())
882 }
883
884 #[must_use]
891 pub fn idempotency_header(&self, target: &str) -> Option<String> {
892 self.policy
893 .read()
894 .expect("policy lock poisoned")
895 .resolve(target)
896 .idempotency
897 .map(|i| i.header)
898 }
899
900 #[must_use]
904 pub fn nondeterminism_response(&self) -> NondeterminismResponse {
905 self.policy
906 .read()
907 .expect("policy lock poisoned")
908 .flows
909 .as_ref()
910 .map_or(NondeterminismResponse::default(), |f| f.on_nondeterminism)
911 }
912
913 fn state(&self) -> MutexGuard<'_, State> {
914 self.state.lock().expect("state lock poisoned")
915 }
916
917 fn elapsed_ms(&self) -> u64 {
918 u64::try_from(self.started.elapsed().as_millis()).unwrap_or(u64::MAX)
919 }
920
921 pub async fn execute<F>(&self, request: &Request, mut effect: F) -> Outcome
926 where
927 F: AsyncFnMut(u32) -> AttemptResult,
928 {
929 let started = Instant::now();
930 let span = tracing::info_span!(
936 "keel.call",
937 target = %request.target,
938 op = %request.op,
939 trace_id = tracing::field::Empty,
940 result = tracing::field::Empty,
941 error_code = tracing::field::Empty,
942 attempts = tracing::field::Empty,
943 from_cache = tracing::field::Empty,
944 throttled = tracing::field::Empty,
945 breaker = tracing::field::Empty,
946 );
947 let out = self
948 .run_chain(request, &mut effect)
949 .instrument(span.clone())
950 .await;
951 record_call_fields(&span, &out);
952 self.observe(request, &out, started);
953 self.emit_event(|| EventKind::CallEnd {
956 call: out.trace_id.clone(),
957 target: request.target.clone(),
958 result: out.result.clone(),
959 code: out.error.as_ref().map(|e| e.code),
960 attempts: out.attempts,
961 });
962 out
963 }
964
965 async fn run_chain<F>(&self, request: &Request, effect: &mut F) -> Outcome
970 where
971 F: AsyncFnMut(u32) -> AttemptResult,
972 {
973 let target = request.target.as_str();
974 let mut out = self.begin_call(target);
975
976 let trace = self.events.as_ref().map(|sink| {
979 let seq = sink.emit(
980 self.elapsed_ms(),
981 EventKind::CallStart {
982 call: out.trace_id.clone(),
983 target: target.to_owned(),
984 op: request.op.clone(),
985 },
986 );
987 TraceRef {
988 run: sink.run_id().to_owned(),
989 seq,
990 }
991 });
992
993 if request.v != ENVELOPE_VERSION {
994 out.error = Some(OutcomeError {
995 code: ErrorCode::EnvelopeVersion,
996 class: ErrorClass::Other,
997 http_status: None,
998 message: format!("unsupported envelope version {}", request.v),
999 original: None,
1000 });
1001 self.state().metrics_for(target).failures += 1;
1002 return out;
1003 }
1004
1005 let resolved = self
1006 .policy
1007 .read()
1008 .expect("policy lock poisoned")
1009 .resolve(target);
1010
1011 let cache_plan = self.plan_cache(target, &resolved, request);
1014 match &cache_plan {
1015 CachePlan::Memory { key } => {
1016 if self.serve_from_cache(key, &mut out) {
1017 crate::metrics::record_cache_request(target, true);
1018 self.emit_cache(&out, target, CacheStore::Memory, true);
1019 return out;
1020 }
1021 crate::metrics::record_cache_request(target, false);
1022 self.emit_cache(&out, target, CacheStore::Memory, false);
1023 }
1024 CachePlan::Persistent { key, .. } => {
1025 if self.serve_from_persistent(target, key, &mut out) {
1026 crate::metrics::record_cache_request(target, true);
1027 self.emit_cache(&out, target, CacheStore::Persistent, true);
1028 return out;
1029 }
1030 crate::metrics::record_cache_request(target, false);
1031 self.emit_cache(&out, target, CacheStore::Persistent, false);
1032 }
1033 CachePlan::None => {}
1034 }
1035
1036 if let Some(rate) = resolved.rate {
1038 self.throttle(target, rate, &mut out).await;
1039 }
1040
1041 let admission = self.admit(target, &resolved, &mut out, trace.as_ref());
1043 if admission == Admission::Rejected {
1044 return out;
1045 }
1046
1047 let retry = resolved.retry.clone().unwrap_or_else(|| RetryPolicy {
1049 attempts: core::num::NonZeroU32::MIN,
1050 ..RetryPolicy::default()
1051 });
1052 let result = match resolved.poll.as_ref().filter(|_| poll_applies(request)) {
1053 Some(poll) => {
1054 self.run_poll(
1055 request,
1056 &resolved,
1057 &retry,
1058 poll,
1059 effect,
1060 &mut out,
1061 trace.as_ref(),
1062 )
1063 .await
1064 }
1065 None => {
1066 self.run_attempts(request, &resolved, &retry, effect, &mut out, trace.as_ref())
1067 .await
1068 }
1069 };
1070 let memory_key = match &cache_plan {
1073 CachePlan::Memory { key } => Some(key.clone()),
1074 _ => None,
1075 };
1076 self.settle(target, &resolved, admission, memory_key, result, &mut out);
1077
1078 if let CachePlan::Persistent { key, ttl } = &cache_plan
1079 && out.result == "ok"
1080 && let Some(payload) = &out.payload
1081 {
1082 self.write_persistent(target, key, payload, *ttl);
1083 }
1084 out
1085 }
1086
1087 fn plan_cache(&self, target: &str, resolved: &ResolvedPolicy, request: &Request) -> CachePlan {
1091 let (Some(cache), Some(hash)) = (resolved.cache.as_ref(), request.args_hash.as_ref())
1092 else {
1093 return CachePlan::None;
1094 };
1095 let Some(ttl) = cache.ttl else {
1096 return CachePlan::None;
1097 };
1098 match cache.scope {
1099 CacheScope::Persistent if self.current_journal().is_some() => CachePlan::Persistent {
1100 key: JournalCacheKey::new(format!("{target}#{hash}")),
1101 ttl,
1102 },
1103 _ => CachePlan::Memory {
1104 key: CacheKey {
1105 target: target.to_owned(),
1106 args_hash: hash.clone(),
1107 },
1108 },
1109 }
1110 }
1111
1112 fn begin_call(&self, target: &str) -> Outcome {
1114 let mut state = self.state();
1115 state.metrics_for(target).calls += 1;
1116 state.trace_seq += 1;
1117 Outcome {
1118 v: ENVELOPE_VERSION,
1119 result: String::from("error"),
1120 payload: None,
1121 error: None,
1122 attempts: 0,
1123 from_cache: false,
1124 waits_ms: Vec::new(),
1125 throttled: false,
1126 throttle_wait_ms: 0,
1127 breaker: BreakerState::Closed,
1128 trace_id: format!("t-{:06}", state.trace_seq),
1129 }
1130 }
1131
1132 fn serve_from_cache(&self, key: &CacheKey, out: &mut Outcome) -> bool {
1137 let now = Instant::now();
1138 let mut state = self.state();
1139 let payload = match state.cache.get(key) {
1140 Some(entry) if now < entry.expires_at => entry.payload.clone(),
1141 Some(_) => {
1142 state.cache.remove(key);
1145 return false;
1146 }
1147 None => return false,
1148 };
1149 out.result = String::from("ok");
1150 out.payload = Some(payload);
1151 out.from_cache = true;
1152 let metrics = state.metrics_for(&key.target);
1153 metrics.cache_hits += 1;
1154 metrics.successes += 1;
1155 out.breaker = state.breaker_state(&key.target, now);
1156 debug!(target = %key.target, scope = "memory", "cache hit");
1157 true
1158 }
1159
1160 fn serve_from_persistent(
1167 &self,
1168 target: &str,
1169 key: &JournalCacheKey,
1170 out: &mut Outcome,
1171 ) -> bool {
1172 let Some(journal) = self.current_journal() else {
1173 return false;
1174 };
1175 let bytes = match journal.get_cache(key) {
1176 Ok(Some(bytes)) => bytes,
1177 Ok(None) => return false,
1178 Err(error) => {
1179 warn!(target = %target, error = %error, "persistent cache read failed; serving live");
1180 return false;
1181 }
1182 };
1183 let payload = match decode_cache_payload(&bytes) {
1184 Ok(payload) => payload,
1185 Err(reason) => {
1186 warn!(target = %target, reason = %reason, "persistent cache entry undecodable; serving live");
1187 return false;
1188 }
1189 };
1190 let now = Instant::now();
1191 let mut state = self.state();
1192 out.result = String::from("ok");
1193 out.payload = Some(payload);
1194 out.from_cache = true;
1195 let metrics = state.metrics_for(target);
1196 metrics.cache_hits += 1;
1197 metrics.successes += 1;
1198 out.breaker = state.breaker_state(target, now);
1199 debug!(target = %target, scope = "persistent", "cache hit");
1200 true
1201 }
1202
1203 async fn throttle(&self, target: &str, rate: Rate, out: &mut Outcome) {
1205 let wait_ms = {
1206 let elapsed = self.elapsed_ms();
1207 let mut state = self.state();
1208 let bucket = state.rate_buckets.entry(target.to_owned()).or_default();
1209 bucket.plan_admit(elapsed, rate)
1210 };
1211 if wait_ms > 0 {
1212 out.throttled = true;
1213 out.throttle_wait_ms = wait_ms;
1214 self.state().metrics_for(target).throttled += 1;
1215 crate::metrics::record_throttled(target, wait_ms);
1216 self.emit_event(|| EventKind::Throttle {
1218 call: out.trace_id.clone(),
1219 target: target.to_owned(),
1220 wait_ms,
1221 });
1222 tokio::time::sleep(Duration::from_millis(wait_ms)).await;
1223 }
1224 }
1225
1226 fn admit(
1229 &self,
1230 target: &str,
1231 resolved: &ResolvedPolicy,
1232 out: &mut Outcome,
1233 trace: Option<&TraceRef>,
1234 ) -> Admission {
1235 if resolved.breaker.is_none() {
1236 return Admission::Closed;
1237 }
1238 let now = Instant::now();
1239 let admission = {
1240 let mut state = self.state();
1241 let admission = state
1242 .breakers
1243 .entry(target.to_owned())
1244 .or_default()
1245 .admit(now);
1246 if admission == Admission::Rejected {
1247 out.error = Some(OutcomeError {
1248 code: ErrorCode::BreakerOpen,
1249 class: ErrorClass::Other,
1250 http_status: None,
1251 message: with_trace_ref(
1252 format!("breaker OPEN for {target}: failed fast, call not attempted"),
1253 trace,
1254 ),
1255 original: None,
1256 });
1257 out.breaker = BreakerState::Open;
1258 state.metrics_for(target).failures += 1;
1259 }
1260 admission
1261 };
1262 if admission == Admission::Rejected {
1264 self.emit_event(|| EventKind::BreakerReject {
1265 call: out.trace_id.clone(),
1266 target: target.to_owned(),
1267 });
1268 }
1269 if admission == Admission::HalfOpen {
1271 debug!(target = %target, transition = "half_open", "breaker transition");
1272 crate::metrics::record_breaker_transition(target, "half_open");
1273 self.emit_event(|| EventKind::BreakerHalfOpen {
1274 call: out.trace_id.clone(),
1275 target: target.to_owned(),
1276 });
1277 }
1278 admission
1279 }
1280
1281 fn settle(
1284 &self,
1285 target: &str,
1286 resolved: &ResolvedPolicy,
1287 admission: Admission,
1288 cache_key: Option<CacheKey>,
1289 result: Result<Value, OutcomeError>,
1290 out: &mut Outcome,
1291 ) {
1292 let now = Instant::now();
1293 let transition = {
1294 let mut state = self.state();
1295 let transition = match result {
1296 Ok(payload) => {
1297 state.metrics_for(target).successes += 1;
1298 let mut transition = BreakerTransition::None;
1299 if let Some(config) = &resolved.breaker
1300 && let Some(breaker) = state.breakers.get_mut(target)
1301 {
1302 transition = breaker.on_success(now, config);
1303 }
1304 if let (Some(key), Some(cache)) = (cache_key, &resolved.cache)
1305 && let Some(ttl) = cache.ttl
1306 {
1307 state.cache.retain(|_, entry| entry.expires_at > now);
1315 state.cache.insert(
1316 key,
1317 CacheEntry {
1318 expires_at: now + Duration::from_millis(ttl.0),
1319 payload: payload.clone(),
1320 },
1321 );
1322 }
1323 out.result = String::from("ok");
1324 out.payload = Some(payload);
1325 transition
1326 }
1327 Err(error) => {
1328 state.metrics_for(target).failures += 1;
1329 let mut transition = BreakerTransition::None;
1330 if let Some(config) = &resolved.breaker
1331 && let Some(breaker) = state.breakers.get_mut(target)
1332 {
1333 transition = breaker.on_terminal_failure(now, config, admission);
1334 }
1335 out.error = Some(error);
1336 transition
1337 }
1338 };
1339 out.breaker = state.breaker_state(target, now);
1340 transition
1341 };
1342 emit_breaker_transition(target, transition);
1343 match transition {
1344 BreakerTransition::Opened => self.emit_event(|| EventKind::BreakerOpen {
1345 call: out.trace_id.clone(),
1346 target: target.to_owned(),
1347 cooldown_ms: resolved.breaker.as_ref().map_or(0, |b| b.cooldown.0),
1348 }),
1349 BreakerTransition::Closed => self.emit_event(|| EventKind::BreakerClose {
1350 call: out.trace_id.clone(),
1351 target: target.to_owned(),
1352 }),
1353 BreakerTransition::None => {}
1354 }
1355 }
1356
1357 fn write_persistent(
1362 &self,
1363 target: &str,
1364 key: &JournalCacheKey,
1365 payload: &Value,
1366 ttl: DurationMs,
1367 ) {
1368 let Some(journal) = self.current_journal() else {
1369 return;
1370 };
1371 let bytes = match encode_cache_payload(payload) {
1372 Ok(bytes) => bytes,
1373 Err(error) => {
1374 warn!(target = %target, error = %error, "persistent cache encode failed; entry not stored");
1375 return;
1376 }
1377 };
1378 if let Err(error) = journal.put_cache(key, &bytes, Duration::from_millis(ttl.0)) {
1379 warn!(target = %target, error = %error, "persistent cache write failed; entry not stored");
1380 }
1381 }
1382
1383 fn observe(&self, request: &Request, out: &Outcome, started: Instant) {
1388 let Some(discovery) = self.discovery.as_ref() else {
1389 return;
1390 };
1391 let latency_ms = i64::try_from(started.elapsed().as_millis()).unwrap_or(i64::MAX);
1392 let result = if out.from_cache {
1393 CallResult::CacheHit
1394 } else if out.result == "ok" {
1395 CallResult::Success
1396 } else {
1397 CallResult::Failure
1398 };
1399 let error = out.error.as_ref().map(|e| ObservedError {
1400 class: e.class,
1401 http_status: e.http_status,
1402 });
1403 let breaker_opened = out
1404 .error
1405 .as_ref()
1406 .is_some_and(|e| e.code == ErrorCode::BreakerOpen);
1407 let not_retried = out
1410 .error
1411 .as_ref()
1412 .is_some_and(|e| e.code == ErrorCode::NonIdempotentNotRetried);
1413 let wrapped = self
1417 .policy
1418 .read()
1419 .expect("policy lock poisoned")
1420 .target
1421 .contains_key(&request.target);
1422 let observation = CallObservation {
1423 target: request.target.clone(),
1424 result,
1425 attempts: out.attempts,
1426 latency_ms,
1427 throttled: out.throttled,
1428 breaker_opened,
1429 not_retried,
1430 wrapped,
1431 error,
1432 };
1433 if let Err(error) = discovery.record(&observation) {
1434 warn!(target = %request.target, error = %error, "discovery record failed; observation dropped");
1435 }
1436 }
1437
1438 async fn run_one_attempt<F>(
1454 &self,
1455 timeout: Option<DurationMs>,
1456 effect: &mut F,
1457 attempt: u32,
1458 attempt_span: &tracing::Span,
1459 ) -> AttemptOutcome
1460 where
1461 F: AsyncFnMut(u32) -> AttemptResult,
1462 {
1463 match timeout {
1464 Some(limit) => {
1465 match tokio::time::timeout(
1466 Duration::from_millis(limit.0),
1467 effect(attempt).instrument(attempt_span.clone()),
1468 )
1469 .await
1470 {
1471 Ok(result) => AttemptOutcome {
1472 result,
1473 timed_out_by_layer: false,
1474 },
1475 Err(_elapsed) => AttemptOutcome {
1476 result: AttemptResult::Error {
1477 class: ErrorClass::Timeout,
1478 http_status: None,
1479 retry_after_ms: None,
1480 message: format!("no response within {}ms", limit.0),
1481 original: None,
1482 },
1483 timed_out_by_layer: true,
1484 },
1485 }
1486 }
1487 None => AttemptOutcome {
1488 result: effect(attempt).instrument(attempt_span.clone()).await,
1489 timed_out_by_layer: false,
1490 },
1491 }
1492 }
1493
1494 #[expect(clippy::too_many_arguments, reason = "mirrors run_attempts' signature")]
1498 async fn run_poll<F>(
1499 &self,
1500 request: &Request,
1501 resolved: &ResolvedPolicy,
1502 retry: &RetryPolicy,
1503 poll: &keel_core_api::policy::PollPolicy,
1504 effect: &mut F,
1505 out: &mut Outcome,
1506 trace: Option<&TraceRef>,
1507 ) -> Result<Value, OutcomeError>
1508 where
1509 F: AsyncFnMut(u32) -> AttemptResult,
1510 {
1511 let started = Instant::now();
1512 loop {
1513 let payload = self
1514 .run_attempts(request, resolved, retry, effect, out, trace)
1515 .await?;
1516 match poll_verdict(poll, &payload) {
1517 PollVerdict::Terminal => return Ok(payload),
1518 PollVerdict::FailOpen => {
1519 debug!(
1520 target = %request.target, field = %poll.until.field,
1521 "poll fail-open: response not judgeable; returned as-is"
1522 );
1523 return Ok(payload);
1524 }
1525 PollVerdict::Pending => {
1526 let elapsed = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX);
1527 if elapsed.saturating_add(poll.interval.0) > poll.deadline.0 {
1528 return Err(OutcomeError {
1529 code: ErrorCode::PollDeadlineExceeded,
1530 class: ErrorClass::Other,
1531 http_status: None,
1532 message: with_trace_ref(
1533 format!(
1534 "{} poll deadline exceeded: '{}' not terminal after {}ms",
1535 request.op, poll.until.field, poll.deadline.0
1536 ),
1537 trace,
1538 ),
1539 original: None,
1540 });
1541 }
1542 tokio::time::sleep(Duration::from_millis(poll.interval.0)).await;
1543 }
1544 }
1545 }
1546 }
1547
1548 async fn run_attempts<F>(
1551 &self,
1552 request: &Request,
1553 resolved: &ResolvedPolicy,
1554 retry: &RetryPolicy,
1555 effect: &mut F,
1556 out: &mut Outcome,
1557 trace: Option<&TraceRef>,
1558 ) -> Result<Value, OutcomeError>
1559 where
1560 F: AsyncFnMut(u32) -> AttemptResult,
1561 {
1562 let target = request.target.as_str();
1563 let max_attempts = retry.attempts.get();
1564 let attempt_timeout = resolved.timeout.filter(|_| request.idempotent);
1571 for attempt in 1..=max_attempts {
1572 out.attempts += 1;
1573 self.state().metrics_for(target).attempts += 1;
1574 crate::metrics::record_attempt(target);
1575 self.emit_event(|| EventKind::AttemptStart {
1576 call: out.trace_id.clone(),
1577 target: target.to_owned(),
1578 attempt,
1579 });
1580
1581 let attempt_span = tracing::debug_span!(
1585 "keel.attempt",
1586 attempt,
1587 result = tracing::field::Empty,
1588 class = tracing::field::Empty,
1589 http_status = tracing::field::Empty,
1590 wait_ms = tracing::field::Empty,
1591 );
1592
1593 let attempt_outcome = self
1594 .run_one_attempt(attempt_timeout, effect, attempt, &attempt_span)
1595 .await;
1596
1597 match attempt_outcome.result {
1598 AttemptResult::Ok { payload } => {
1599 attempt_span.record("result", "ok");
1600 return Ok(payload);
1601 }
1602 AttemptResult::Error {
1603 class,
1604 http_status,
1605 retry_after_ms,
1606 message,
1607 original,
1608 } => {
1609 attempt_span.record("result", "error");
1610 attempt_span.record("class", class_str(class));
1611 if let Some(status) = http_status {
1612 attempt_span.record("http_status", status);
1613 }
1614 self.emit_event(|| EventKind::AttemptError {
1615 call: out.trace_id.clone(),
1616 target: target.to_owned(),
1617 attempt,
1618 class,
1619 http_status,
1620 });
1621 let retryable = retry.is_retryable(class, http_status);
1622 if let Some(code) =
1623 terminal_code(retryable, attempt, max_attempts, request.idempotent)
1624 {
1625 let ctx = TerminalAttemptCtx {
1626 request,
1627 attempt,
1628 max_attempts,
1629 trace,
1630 timed_out_by_layer: attempt_outcome.timed_out_by_layer,
1631 };
1632 return Err(terminal_attempt_error(
1633 &ctx,
1634 code,
1635 class,
1636 http_status,
1637 &message,
1638 original,
1639 ));
1640 }
1641 let (mut wait, jitter) = retry.schedule.wait_and_jitter(attempt);
1645 if jitter && wait > 0 {
1646 wait = fastrand::u64(wait / 2..=wait);
1647 }
1648 if let Some(server_says) = retry_after_ms {
1649 wait = wait.max(server_says);
1650 }
1651 attempt_span.record("wait_ms", wait);
1652 out.waits_ms.push(wait);
1653 self.state().metrics_for(target).retries += 1;
1654 crate::metrics::record_retry(target, wait);
1655 self.emit_event(|| EventKind::Backoff {
1657 call: out.trace_id.clone(),
1658 target: target.to_owned(),
1659 attempt,
1660 wait_ms: wait,
1661 });
1662 tokio::time::sleep(Duration::from_millis(wait)).await;
1663 }
1664 }
1665 }
1666 unreachable!("loop always returns by the final attempt");
1667 }
1668
1669 pub fn report(&self) -> Value {
1672 let now = Instant::now();
1673 let state = self.state();
1674 let targets = state
1675 .metrics
1676 .iter()
1677 .map(|(name, m)| {
1678 let breaker = state.breakers.get(name);
1679 let row = TargetReport {
1680 attempts: m.attempts,
1681 breaker_opens: breaker.map_or(0, |b| b.opens),
1682 breaker_state: state.breaker_state(name, now),
1683 cache_hits: m.cache_hits,
1684 calls: m.calls,
1685 failures: m.failures,
1686 retries: m.retries,
1687 successes: m.successes,
1688 throttled: m.throttled,
1689 };
1690 (name.as_str(), row)
1691 })
1692 .collect();
1693 serde_json::to_value(Report {
1694 v: 1,
1695 clock_ms: self.elapsed_ms(),
1696 targets,
1697 })
1698 .expect("report serialization is infallible")
1699 }
1700}
1701
1702#[cfg(test)]
1703mod tests {
1704 use super::{
1705 Admission, AttemptResult, Breaker, BreakerTransition, ENVELOPE_VERSION, Engine, Instant,
1706 Request, TokenBucket,
1707 };
1708 use core::num::NonZeroU32;
1709 use core::time::Duration;
1710 use keel_core_api::policy::{BreakerPolicy, DurationMs, Rate};
1711 use serde_json::json;
1712
1713 #[test]
1719 fn breaker_rate_mode_evicts_outcomes_older_than_the_window() {
1720 let config = BreakerPolicy {
1721 failures: None,
1722 cooldown: DurationMs(15_000),
1723 window: Some(DurationMs(10_000)),
1724 failure_rate: Some(0.5),
1725 min_calls: Some(NonZeroU32::new(2).unwrap()),
1726 };
1727 let mut breaker = Breaker::default();
1728 let t0 = Instant::now();
1729
1730 assert_eq!(
1731 breaker.on_terminal_failure(t0, &config, Admission::Closed),
1732 BreakerTransition::None,
1733 "one failure is below min_calls"
1734 );
1735 let t_11s = t0 + Duration::from_secs(11);
1736 assert_eq!(
1737 breaker.on_terminal_failure(t_11s, &config, Admission::Closed),
1738 BreakerTransition::None,
1739 "the stale failure must have aged out of the 10s window"
1740 );
1741 }
1742
1743 #[test]
1747 fn token_bucket_caps_refill_at_burst_capacity() {
1748 let rate = Rate {
1749 limit: core::num::NonZeroU64::new(2).unwrap(),
1750 window_ms: 1000,
1751 };
1752 let mut bucket = TokenBucket::default();
1753
1754 assert_eq!(bucket.plan_admit(0, rate), 0, "burst covers the first call");
1755 assert_eq!(
1756 bucket.plan_admit(0, rate),
1757 0,
1758 "burst covers the second call"
1759 );
1760 assert_eq!(
1761 bucket.plan_admit(0, rate),
1762 500,
1763 "burst drained: paced at window/limit = 500ms"
1764 );
1765
1766 assert_eq!(
1769 bucket.plan_admit(50_000, rate),
1770 0,
1771 "capacity refilled to burst"
1772 );
1773 assert_eq!(
1774 bucket.plan_admit(50_000, rate),
1775 0,
1776 "both burst tokens available"
1777 );
1778 assert_eq!(
1779 bucket.plan_admit(50_000, rate),
1780 500,
1781 "capacity clamp: idle time cannot bank more than `limit` tokens"
1782 );
1783 }
1784
1785 fn req(target: &str, args_hash: &str) -> Request {
1786 Request {
1787 v: ENVELOPE_VERSION,
1788 target: target.to_owned(),
1789 op: format!("GET {target}"),
1790 idempotent: true,
1791 args_hash: Some(args_hash.to_owned()),
1792 }
1793 }
1794
1795 #[tokio::test(start_paused = true)]
1799 async fn in_memory_cache_evicts_expired_entries() {
1800 let engine = Engine::new();
1801 engine
1802 .configure(&json!({
1803 "target": { "api.catalog.internal": { "cache": { "ttl": "60s" } } }
1804 }))
1805 .expect("valid policy");
1806
1807 engine
1808 .execute(&req("api.catalog.internal", "k1"), async |_a| {
1809 AttemptResult::Ok { payload: json!(1) }
1810 })
1811 .await;
1812 assert_eq!(engine.state().cache.len(), 1, "k1 cached");
1813
1814 tokio::time::advance(Duration::from_secs(61)).await;
1816 engine
1817 .execute(&req("api.catalog.internal", "k2"), async |_a| {
1818 AttemptResult::Ok { payload: json!(2) }
1819 })
1820 .await;
1821 assert_eq!(
1822 engine.state().cache.len(),
1823 1,
1824 "expired k1 evicted on write; only the live k2 remains"
1825 );
1826
1827 let out = engine
1830 .execute(&req("api.catalog.internal", "k1"), async |_a| {
1831 AttemptResult::Ok { payload: json!(3) }
1832 })
1833 .await;
1834 assert!(!out.from_cache, "expired/evicted key re-runs live");
1835 }
1836
1837 fn poll_policy_json(interval: &str, deadline: &str) -> serde_json::Value {
1838 json!({ "target": { "api.jobs.internal": {
1839 "retry": { "attempts": 3, "schedule": "fixed(200ms)", "on": ["conn", "timeout", "429", "5xx"] },
1840 "poll": { "interval": interval, "deadline": deadline,
1841 "until": { "field": "status", "terminal": ["completed", "failed"] } }
1842 } } })
1843 }
1844
1845 #[tokio::test(start_paused = true)]
1846 async fn poll_reissues_until_terminal_and_accumulates_attempts() {
1847 let engine = Engine::new();
1848 engine.configure(&poll_policy_json("10s", "90s")).unwrap();
1849 let mut script = vec![
1850 json!({"status": "running"}),
1851 json!({"status": "running"}),
1852 json!({"status": "completed", "result": 42}),
1853 ]
1854 .into_iter();
1855 let out = engine
1856 .execute(&req("api.jobs.internal", "h1"), async |_a| {
1857 AttemptResult::Ok {
1858 payload: script.next().expect("script exhausted"),
1859 }
1860 })
1861 .await;
1862 assert_eq!(out.result, "ok");
1863 assert_eq!(
1864 out.payload,
1865 Some(json!({"status": "completed", "result": 42}))
1866 );
1867 assert_eq!(
1868 out.attempts, 3,
1869 "attempts accumulate across poll iterations"
1870 );
1871 assert!(
1872 out.waits_ms.is_empty(),
1873 "poll intervals are not retry waits"
1874 );
1875 let report = engine.report();
1876 assert_eq!(report["targets"]["api.jobs.internal"]["calls"], 1);
1877 assert_eq!(report["targets"]["api.jobs.internal"]["successes"], 1);
1878 }
1879
1880 #[tokio::test(start_paused = true)]
1881 async fn poll_deadline_is_e016_and_breaker_countable() {
1882 let engine = Engine::new();
1883 let mut policy = poll_policy_json("10s", "25s");
1884 policy["target"]["api.jobs.internal"]["breaker"] = json!({ "failures": 1 });
1885 engine.configure(&policy).unwrap();
1886 let out = engine
1887 .execute(&req("api.jobs.internal", "h1"), async |_a| {
1888 AttemptResult::Ok {
1889 payload: json!({"status": "running"}),
1890 }
1891 })
1892 .await;
1893 assert_eq!(out.result, "error");
1894 let err = out.error.expect("terminal error");
1895 assert_eq!(err.code.as_str(), "KEEL-E016");
1896 assert_eq!(
1897 err.message,
1898 "GET api.jobs.internal poll deadline exceeded: 'status' not terminal after 25000ms"
1899 );
1900 assert_eq!(out.attempts, 3);
1902 assert_eq!(
1903 out.breaker,
1904 keel_core_api::BreakerState::Open,
1905 "breaker-countable"
1906 );
1907 }
1908
1909 #[tokio::test(start_paused = true)]
1910 async fn poll_fails_open_on_missing_field_and_skips_non_get() {
1911 let engine = Engine::new();
1912 engine.configure(&poll_policy_json("10s", "90s")).unwrap();
1913 let out = engine
1915 .execute(&req("api.jobs.internal", "h2"), async |_a| {
1916 AttemptResult::Ok {
1917 payload: json!({"progress": 0.5}),
1918 }
1919 })
1920 .await;
1921 assert_eq!(out.result, "ok");
1922 assert_eq!(out.attempts, 1);
1923 let mut post = req("api.jobs.internal", "h3");
1925 post.op = String::from("POST api.jobs.internal");
1926 let out = engine
1927 .execute(&post, async |_a| AttemptResult::Ok {
1928 payload: json!({"status": "running"}),
1929 })
1930 .await;
1931 assert_eq!(out.attempts, 1);
1932 }
1933
1934 #[tokio::test(start_paused = true)]
1935 async fn poll_judges_http_envelope_bodies() {
1936 use base64::Engine as _;
1937 let engine = Engine::new();
1938 engine.configure(&poll_policy_json("10s", "90s")).unwrap();
1939 let b64 =
1940 |v: &serde_json::Value| base64::engine::general_purpose::STANDARD.encode(v.to_string());
1941 let envelope = |body: &serde_json::Value| {
1942 json!({ "status": 200, "headers": [["content-type", "application/json"]],
1943 "body_b64": b64(body) })
1944 };
1945 let mut script = vec![
1946 envelope(&json!({"status": "running"})),
1947 envelope(&json!({"status": "completed"})),
1948 ]
1949 .into_iter();
1950 let out = engine
1951 .execute(&req("api.jobs.internal", "h4"), async |_a| {
1952 AttemptResult::Ok {
1953 payload: script.next().expect("script exhausted"),
1954 }
1955 })
1956 .await;
1957 assert_eq!(out.result, "ok");
1958 assert_eq!(out.attempts, 2);
1959 let raw = json!({ "status": 200, "headers": [],
1961 "body_b64": base64::engine::general_purpose::STANDARD.encode("not json") });
1962 let out = engine
1963 .execute(&req("api.jobs.internal", "h5"), async |_a| {
1964 AttemptResult::Ok {
1965 payload: raw.clone(),
1966 }
1967 })
1968 .await;
1969 assert_eq!(out.result, "ok");
1970 assert_eq!(out.payload, Some(raw));
1971 assert_eq!(out.attempts, 1);
1972 }
1973}