
File.pudu
Pudu194 lines7.5 KB
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 }323334const GIBIBYTE: Int = 1073741824353637const CHECK_EVERY: Int = 5038394041export 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}5657585960export 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}868788fn 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}949596fn update(state: &Sync.Cell[State], change: fn(State) -> State) -> () {97 let _stored = Sync.set(state, change(stateOf(state)))98}99100101fn 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}120121122fn full(options: &Options, current: &State) -> Bool {123 match options.fileSizeLimitBytes {124 case Some(limit) => current.size >= limit125 case None => false126 }127}128129130131fn 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}157158159fn 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}169170171fn 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}178179180fn 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