
Batching.pudu
Pudu161 lines7.1 KB
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) -> ()]] }313233const TICK: Int = 534353637export fn defaults() -> Options {38 Options{batchSizeLimit: 1000, bufferingTimeLimit: 2000, eagerlyEmitFirstEvent: true, queueLimit: Some(100000), retryTimeLimit: 600000, selfLog: SelfLog.create()}39}4041424344export 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}666768fn 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}747576fn 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}798081fn 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}96979899fn 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}113114115fn writeAll(context: &Context) -> () {116 while !stateOf(context).queue.isEmpty() || !stateOf(context).retrying.isEmpty() {117 if !writeOne(context) { return () }118 }119}120121122123fn 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(¤t.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(¤t.schedule)} })135 true136 }137 case Err(problem) => {138 failed(context, problem, &batch)139 false140 }141 }142}143144145146fn 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