Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-log

Structured event logging for Pudu: message templates, enrichment, filtering, formatting, and sinks

0.1.0Apache-2.01

InstallClose

Batching.pudu

Pudu161 lines7.1 KB

GitHub ↗
1/** @Log.Sinks.Batching.Sink — groups events and emits them in batches */2module PuduLangLog.Sinks.Batching34import Std.Concurrent as Concurrent5import Std.List as List6import Std.Sync as Sync7import Std.Time as Time8import PuduLangLog.Domain.Schedule as Schedule9import PuduLangLog as Log10import PuduLangLog.SelfLog as SelfLog11import PuduLangLog.Sink as Sink1213/** @Log.Sinks.Batching.Options — batch size, timing, queue, and retries */14export type Options = {15  batchSizeLimit: Int,16  bufferingTimeLimit: Int,17  eagerlyEmitFirstEvent: Bool,18  queueLimit: Option[Int],19  retryTimeLimit: Int,20  selfLog: SelfLog.SelfLog21}2223/** @Log.Sinks.Batching.Target — where a batch of events is written */24export type Target = { emitBatch: fn(&Array[Log.Event]) -> Result[(), Str], onEmptyBatch: fn() -> () }2526/** @Log.Sinks.Batching.State — the queue, the batch being retried, and the schedule */27type State = { queue: Array[Log.Event], retrying: Array[Log.Event], schedule: Schedule.State, eager: Bool, closed: Bool }2829/** @Log.Sinks.Batching.Context — what the worker and the sink share */30type Context = { target: Target, options: Options, state: Sync.Cell[State], lock: Sync.Mutex, listener: Sync.Cell[Option[fn(&Sink.Report) -> ()]] }3132/// How long the worker sleeps between checks, in milliseconds.33const TICK: Int = 53435/// Batches of up to a thousand, every two seconds, the first event sent at once, a queue of a36/// hundred thousand, retries for up to ten minutes, and a disabled self-log.37export fn defaults() -> Options {38  Options{batchSizeLimit: 1000, bufferingTimeLimit: 2000, eagerlyEmitFirstEvent: true, queueLimit: Some(100000), retryTimeLimit: 600000, selfLog: SelfLog.create()}39}4041/// A sink queueing events for a worker thread that writes them to the target in batches. A failed42/// batch is retried with growing waits until the retry time runs out; events it loses are reported43/// to the self-log and to an attached listener. Flushing writes everything queued now.44export fn sink(target: Target, options: Options) -> Sink.Sink {45  let none: Array[Log.Event] = []46  let state = Sync.cell(State{queue: none, retrying: none, schedule: Schedule.create(options.bufferingTimeLimit, options.retryTimeLimit), eager: options.eagerlyEmitFirstEvent, closed: false})47  let lock = Sync.mutex()48  let writing = Sync.mutex()49  let noListener: Option[fn(&Sink.Report) -> ()] = None50  let listener = Sync.cell(noListener)51  let context = Context{target: target, options: options, state: state, lock: lock, listener: listener}52  let worker = Concurrent.start(fn() -> () { work(&context, &writing) })53  Sink.Sink {54    emit: fn(event: &Log.Event) -> Result[(), Str] { enqueue(&context, event) },55    flush: fn() -> () {56      let _written = Sync.withLock(&writing, fn() -> () { writeAll(&context) })57    },58    close: fn() -> () {59      change(&context, fn(held: State) -> State { State{..held, closed: true} })60      if let Ok(started) = worker { let _joined = Concurrent.join(&started) }61      let _written = Sync.withLock(&writing, fn() -> () { writeAll(&context) })62    },63    attach: fn(added: Sink.Listener) -> () { let _stored = Sync.set(&listener, Some(added)) }64  }65}6667/// The shared state.68fn stateOf(context: &Context) -> State {69  match Sync.get(&context.state) {70    case Ok(held) => held71    case Err(_) => State{queue: [], retrying: [], schedule: Schedule.create(0, 0), eager: false, closed: true}72  }73}7475/// Changes the shared state under its lock.76fn change(context: &Context, update: fn(State) -> State) -> () {77  let _changed = Sync.withLock(&context.lock, fn() -> () { let _stored = Sync.set(&context.state, update(stateOf(context))) })78}7980/// Adds an event to the queue, dropping it when the queue is at its limit or the sink is closed.81fn enqueue(context: &Context, event: &Log.Event) -> Result[(), Str] {82  let accepted = Sync.withLock(&context.lock, fn() -> Result[(), Str] {83      let held = stateOf(context)84      if held.closed { return Err("the batching sink is closed") }85      if let Some(limit) = context.options.queueLimit {86        if held.queue.length() >= limit { return Ok(()) }87      }88      let _stored = Sync.set(&context.state, State{..held, queue: held.queue.push(*event)})89      Ok(())90    })91  match accepted {92    case Ok(outcome) => outcome93    case Err(_) => Err("the batching sink's lock could not be taken")94  }95}9697/// The worker: waits for a full batch, the buffering time, or the first event when eager, then98/// writes one batch, until the sink closes.99fn work(context: &Context, writing: &Sync.Mutex) -> () {100  var waitedFrom = Time.elapsed()101  while !stateOf(context).closed {102    let _slept = Concurrent.sleep(TICK)103    let held = stateOf(context)104    let due = Time.elapsed() - waitedFrom >= Schedule.interval(&held.schedule)105    let full = held.queue.length() >= context.options.batchSizeLimit106    let eager = held.eager && !held.queue.isEmpty()107    if due || full || eager {108      let _written = Sync.withLock(writing, fn() -> Bool { writeOne(context) })109      waitedFrom = Time.elapsed()110    }111  }112}113114/// Writes batches until the queue is empty or a batch fails.115fn writeAll(context: &Context) -> () {116  while !stateOf(context).queue.isEmpty() || !stateOf(context).retrying.isEmpty() {117    if !writeOne(context) { return () }118  }119}120121/// Writes the batch being retried, or the next batch from the queue, and answers whether it was122/// written. An empty queue calls the target's `onEmptyBatch`.123fn writeOne(context: &Context) -> Bool {124  let held = stateOf(context)125  let batch = if held.retrying.isEmpty() { List.take(&held.queue, context.options.batchSizeLimit) } else { held.retrying }126  if batch.isEmpty() {127    (context.target.onEmptyBatch)()128    return true129  }130  let taken = if held.retrying.isEmpty() { batch.length() } else { 0 }131  change(context, fn(current: State) -> State { State{..current, queue: List.drop(&current.queue, taken), retrying: batch, eager: false} })132  match (context.target.emitBatch)(&batch) {133    case Ok(_) => {134      change(context, fn(current: State) -> State { State{..current, retrying: [], schedule: Schedule.succeeded(&current.schedule)} })135      true136    }137    case Err(problem) => {138      failed(context, problem, &batch)139      false140    }141  }142}143144/// Records a failed batch: it is kept for a retry until the schedule drops it, and the queue is145/// emptied after too many dropped batches.146fn failed(context: &Context, problem: Str, batch: &Array[Log.Event]) -> () {147  let verdict = Schedule.failed(&stateOf(context).schedule, Time.elapsed())148  SelfLog.write(&context.options.selfLog, "Failed to emit a batch of " + show(batch.length()) + " events: " + problem)149  let queued = stateOf(context).queue150  change(context, fn(current: State) -> State {151      let retrying = if verdict.dropBatch { [] } else { current.retrying }152      let queue = if verdict.dropQueue { [] } else { current.queue }153      State{..current, retrying: retrying, queue: queue, schedule: verdict.state}154    })155  let dropped = if verdict.dropBatch { *batch } else { [] }156  let lost = if verdict.dropQueue { dropped.concat(queued) } else { dropped }157  if !lost.isEmpty() {158    if let Ok(Some(listener)) = Sync.get(&context.listener) { listener(&Sink.Report{kind: Sink.Permanent, message: problem, events: lost}) }159  }160}161