
Logging.pudu
Pudu87 lines3.8 KB
1/** @Mediator.Behaviors.Logging.Module — structured events for every request, stream, and notification */2module PuduLangMediator.Behaviors.Logging34import Std.Time as Time5import PuduLangLog.Logger as Logger6import PuduLangLog.Value as Value7import PuduLangMediator.Context as Context8import PuduLangMediator.Message as Message9import PuduLangMediator as Messaging10import PuduLangMediator.Open as Open11import PuduLangMediator.Request as Request12import PuduLangMediator.Stream as Stream13import PuduLangMediator.Utils.Shared as Shared1415/** @Mediator.Behaviors.Logging.Recorder — writes one event per message start and end */16export type Recorder = { logger: Logger.Logger }171819export const STARTED: Str = "Handling \{MessageName\} (\{MessageShape\})"202122export const SUCCEEDED: Str = "Handled \{MessageName\} in \{Elapsed\} ms"232425export const STREAMED: Str = "Streamed \{MessageName\}: \{Items\} items in \{Elapsed\} ms"262728export const FAILED: Str = "\{MessageName\} failed in \{Elapsed\} ms: \{Failure\}"293031export const HEARD: Str = "Published \{MessageName\}"323334export fn of(logger: Logger.Logger) -> Recorder { Recorder{logger: logger.forSource("PuduLangMediator")} }3536impl Open.Behavior for Recorder {37 38 fn handleRequest[Q, R, E](self: &Self, request: &Q, info: &Message.Info, context: Context.Context, next: Request.Next[R, E]) -> Messaging.Outcome[R, E] {39 started(self, info)40 let began = Time.elapsed()41 let outcome = next(context)42 let took = Time.elapsed() - began43 match outcome {44 case Ok(_) => self.logger.information(SUCCEEDED, [Value.text(info.name), Value.int(took)])45 case Err(failure) => failed(self, info, took, &failure)46 }47 outcome48 }49}5051impl Open.StreamBehavior for Recorder {52 53 fn handleStream[Q, T, E](self: &Self, request: &Q, info: &Message.Info, context: Context.Context, sink: Stream.Sink[T], next: Stream.Next[T, E]) -> Messaging.Outcome[(), E] {54 started(self, info)55 let began = Time.elapsed()56 let emitted = Shared.shared(0)57 let outcome = next(context, fn(item: T) -> Bool {58 Shared.update(&emitted, |count: Int| count + 1)59 sink(item)60 })61 let took = Time.elapsed() - began62 match outcome {63 case Ok(_) => self.logger.information(STREAMED, [Value.text(info.name), Value.int(Shared.current(&emitted)), Value.int(took)])64 case Err(failure) => failed(self, info, took, &failure)65 }66 outcome67 }68}6970impl Open.NotificationHandler for Recorder {71 72 fn handleNotification[N, E](self: &Self, notification: &N, info: &Message.Info, context: Context.Context) -> Messaging.Outcome[(), E] {73 self.logger.information(HEARD, [Value.text(info.name)])74 Ok(())75 }76}777879fn started(recorder: &Recorder, info: &Message.Info) -> () {80 recorder.logger.debug(STARTED, [Value.text(info.name), Value.text(Message.shapeName(&info.shape))])81}828384fn failed[E](recorder: &Recorder, info: &Message.Info, took: Int, failure: &Messaging.Failure[E]) -> () {85 recorder.logger.warning(FAILED, [Value.text(info.name), Value.int(took), Value.text(Messaging.describe(failure))])86}87