
Schedule.pudu
Pudu52 lines2.2 KB
1/** @Log.Domain.Schedule.Module — batch retry timing after sink failures */2module PuduLangLog.Domain.Schedule34import Std.Math as Math5import Std.Option as Option67/** @Log.Domain.Schedule.State — failures since the last successful batch */8export type State = { buffering: Int, retryLimit: Int, failures: Int, dropped: Int, firstFailure: Option[Int] }910/** @Log.Domain.Schedule.Verdict — what to do with a failed batch */11export type Verdict = { state: State, dropBatch: Bool, dropQueue: Bool }121314const MINIMUM_BACKOFF: Int = 5000151617const MAXIMUM_BACKOFF: Int = 60000181920const DROPS_BEFORE_QUEUE: Int = 1021222324export fn create(buffering: Int, retryLimit: Int) -> State {25 State{buffering: buffering, retryLimit: retryLimit, failures: 0, dropped: 0, firstFailure: None}26}272829export fn succeeded(state: &State) -> State { create(state.buffering, state.retryLimit) }303132export fn failed(state: &State, now: Int) -> Verdict {33 let first = Option.unwrapOr(state.firstFailure, now)34 let counted = State{..*state, failures: state.failures + 1, firstFailure: Some(first)}35 let dropBatch = now - first + interval(&counted) >= state.retryLimit36 let dropped = if dropBatch { counted.dropped + 1 } else { counted.dropped }37 Verdict{state: State{..counted, dropped: dropped}, dropBatch: dropBatch, dropQueue: dropped >= DROPS_BEFORE_QUEUE}38}39404142export fn interval(state: &State) -> Int {43 if state.failures <= 1 { return state.buffering }44 var backoff = Math.max(state.buffering, MINIMUM_BACKOFF)45 var doublings = state.failures - 146 while doublings > 0 {47 backoff = Math.min(MAXIMUM_BACKOFF, backoff * 2)48 doublings = doublings - 149 }50 Math.max(state.buffering, backoff)51}52