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

File.pudu

Pudu194 lines7.5 KB

GitHub ↗
1/** @Log.Sinks.File.Sink — events appended to rolling files */2module PuduLangLog.Sinks.File34import Std.Concurrent as Concurrent5import Std.Io as Io6import Std.Option as Option7import Std.Result as Result8import Std.Sync as Sync9import PuduLangLog.Constants.Names as Names10import PuduLangLog.Domain.Rolling as Rolling11import PuduLangLog.Formatting.Text as Text12import PuduLangLog as Log13import PuduLangLog.Sink as Sink1415/** @Log.Sinks.File.Options — path, layout, rolling, retention, and buffering */16export type Options = {17  path: Str,18  formatter: fn(&Log.Event) -> Str,19  fileSizeLimitBytes: Option[Int],20  rollingInterval: Log.RollingInterval,21  rollOnFileSizeLimit: Bool,22  retainedFileCountLimit: Option[Int],23  retainedFileTimeLimit: Option[Int],24  buffered: Bool,25  flushInterval: Option[Int],26  onOpened: Option[fn(Str) -> Option[Str]],27  onDeleting: Option[fn(Str) -> ()]28}2930/** @Log.Sinks.File.State — the open file and what is waiting for it */31type State = { file: Str, period: Option[Int], sequence: Option[Int], size: Int, pending: Array[Str], closed: Bool }3233/// One gibibyte, the default size limit of a file.34const GIBIBYTE: Int = 10737418243536/// How often a flushing thread checks whether the sink closed, in milliseconds.37const CHECK_EVERY: Int = 503839/// The file template, one file without rolling, a one-gibibyte limit, thirty-one retained files,40/// and unbuffered writes.41export fn defaults(path: Str) -> Options {42  Options {43    path: path,44    formatter: Text.template(Names.FILE_TEMPLATE),45    fileSizeLimitBytes: Some(GIBIBYTE),46    rollingInterval: Log.Infinite,47    rollOnFileSizeLimit: false,48    retainedFileCountLimit: Some(31),49    retainedFileTimeLimit: None,50    buffered: false,51    flushInterval: None,52    onOpened: None,53    onDeleting: None54  }55}5657/// A sink appending each event to the file of its period. A file that reached its size limit takes58/// no more events, or gives way to the next sequence when rolling on size; opening a file removes59/// the files retention no longer keeps. Buffered text is written on flush and close.60export fn sink(options: Options) -> Sink.Sink {61  let none: Array[Str] = []62  let state = Sync.cell(State{file: "", period: None, sequence: None, size: 0, pending: none, closed: false})63  let lock = Sync.mutex()64  let flush = fn() -> () {65    let _flushed = Sync.withLock(&lock, fn() -> () { let _written = writePending(&state) })66  }67  if let Some(every) = options.flushInterval { startFlushing(every, &state, flush) }68  Sink {69    ..Sink.none(),70    emit: fn(event: &Log.Event) -> Result[(), Str] {71      let text = (options.formatter)(event)72      match Sync.withLock(&lock, fn() -> Result[(), Str] { write(&options, &state, event, text) }) {73        case Ok(outcome) => outcome74        case Err(_) => Err("the file lock could not be taken")75      }76    },77    flush: flush,78    close: fn() -> () {79      let _closed = Sync.withLock(&lock, fn() -> () {80          let _written = writePending(&state)81          update(&state, fn(current: State) -> State { State{..current, closed: true} })82        })83    }84  }85}8687/// The state held by the cell.88fn stateOf(state: &Sync.Cell[State]) -> State {89  match Sync.get(state) {90    case Ok(current) => current91    case Err(_) => State{file: "", period: None, sequence: None, size: 0, pending: [], closed: true}92  }93}9495/// Stores a changed state.96fn update(state: &Sync.Cell[State], change: fn(State) -> State) -> () {97  let _stored = Sync.set(state, change(stateOf(state)))98}99100/// Writes one event's text, opening the file of its period and rolling on size as configured.101fn write(options: &Options, state: &Sync.Cell[State], event: &Log.Event, text: Str) -> Result[(), Str] {102  if stateOf(state).closed { return Err("the file sink is closed") }103  let period = Rolling.checkpoint(options.rollingInterval, &event.timestamp)104  let current = stateOf(state)105  if current.file.isEmpty() || current.period != period { open(options, state, event, period, None) ? }106  if full(options, &stateOf(state)) {107    if !options.rollOnFileSizeLimit { return Ok(()) }108    writePending(state) ?109    open(options, state, event, period, Some(Option.unwrapOr(stateOf(state).sequence, 0) + 1)) ?110  }111  let size = text.toBytes().length()112  if options.buffered {113    update(state, fn(held: State) -> State { State{..held, pending: held.pending.push(text), size: held.size + size} })114    return Ok(())115  }116  Io.append(stateOf(state).file, text) ?117  update(state, fn(held: State) -> State { State{..held, size: held.size + size} })118  Ok(())119}120121/// Whether the open file has reached the size limit.122fn full(options: &Options, current: &State) -> Bool {123  match options.fileSizeLimitBytes {124    case Some(limit) => current.size >= limit125    case None => false126  }127}128129/// Opens the file of a period: the given sequence, or the latest one already on disk. Writes the130/// header of a new file and applies retention.131fn open(options: &Options, state: &Sync.Cell[State], event: &Log.Event, period: Option[Int], wanted: Option[Int]) -> Result[(), Str] {132  writePending(state) ?133  let directory = Io.directoryOf(options.path)134  let listed = if directory.isEmpty() { Io.list(".") } else { Io.list(directory) }135  let names = Result.unwrapOr(listed, [])136  var rolled: Array[Rolling.Rolled] = []137  for name in names {138    if let Some(found) = Rolling.matches(options.path, options.rollingInterval, name) { rolled = rolled.push(found) }139  }140  let sequence = match wanted {141    case Some(number) => Some(number)142    case None => Rolling.latestSequence(&rolled, period)143  }144  let file = Rolling.fileName(options.path, options.rollingInterval, &event.timestamp, sequence)145  if !directory.isEmpty() { Io.makeDirectory(directory) ? }146  let existing = if Io.exists(file) { Result.unwrapOr(Io.countBytes(file), 0) } else { 0 }147  if existing == 0 {148    if let Some(hook) = options.onOpened {149      if let Some(header) = hook(file) { Io.append(file, header) ? }150    }151  }152  let size = if Io.exists(file) { Result.unwrapOr(Io.countBytes(file), 0) } else { 0 }153  update(state, fn(held: State) -> State { State{..held, file: file, period: period, sequence: sequence, size: size} })154  retire(options, &rolled, Io.nameOf(file), event)155  Ok(())156}157158/// Deletes the files retention no longer keeps, calling the deleting hook first.159fn retire(options: &Options, rolled: &Array[Rolling.Rolled], current: Str, event: &Log.Event) -> () {160  let local = event.timestamp.millis + event.timestamp.offset * 60000161  let oldest = Option.map(options.retainedFileTimeLimit, |limit: Int| local - limit)162  let directory = Io.directoryOf(options.path)163  for name in Rolling.retired(rolled, current, options.retainedFileCountLimit, oldest) {164    let path = if directory.isEmpty() { name } else { Io.join(directory, name) }165    if let Some(hook) = options.onDeleting { hook(path) }166    let _removed = Io.removeIfPresent(path)167  }168}169170/// Writes buffered text to the open file.171fn writePending(state: &Sync.Cell[State]) -> Result[(), Str] {172  let current = stateOf(state)173  if current.pending.isEmpty() { return Ok(()) }174  Io.append(current.file, current.pending.join("")) ?175  update(state, fn(held: State) -> State { State{..held, pending: []} })176  Ok(())177}178179/// Starts a thread that flushes every `every` milliseconds until the sink closes.180fn startFlushing(every: Int, state: &Sync.Cell[State], flush: fn() -> ()) -> () {181  let shared = *state182  let _started = Concurrent.start(fn() -> () {183      var waited = 0184      while !stateOf(&shared).closed {185        let _slept = Concurrent.sleep(CHECK_EVERY)186        waited = waited + CHECK_EVERY187        if waited >= every {188          flush()189          waited = 0190        }191      }192    })193}194