
Bus.pudu
Pudu90 lines4.2 KB
1/** @Events.Bus.Module — events published to their listeners through a mediator the package owns */2module PuduLangHttpClient.Events.Bus34import Std.Concurrent.Cancel as Cancel5import Std.Http as Http6import PuduLangHttpClient.Events as Events7import PuduLangHttpClient.Factory.Builder as Builder8import PuduLangHttpClient.Handler as Handler9import PuduLangHttpClient as HttpClient10import PuduLangHttpClient.Request as Request11import PuduLangHttpClient.Response as Response12import PuduLangMediator.Context as Context13import PuduLangMediator.Mediator as Mediator14import PuduLangMediator as Messaging15import PuduLangMediator.Notification as Notification16import PuduLangMediator.Registration as Registration1718/** @Events.Bus.Kinds — the notification kind of each event, made once per bus */19type Kinds = {20 started: Notification.Kind[Events.RequestStarted, Str],21 completed: Notification.Kind[Events.RequestCompleted, Str],22 failed: Notification.Kind[Events.RequestFailed, Str],23 changed: Notification.Kind[Events.PipelineChanged, Str]24}2526/** @Events.Bus.Bus — the mediator carrying events, its kinds, and whether anyone hears requests */27export type Bus = { mediator: Mediator.Mediator, kinds: Kinds, hearsRequests: Bool }282930export fn create(given: &Events.Listeners) -> Bus {31 let kinds = Kinds {32 started: Notification.kind("httpclient.request.started"),33 completed: Notification.kind("httpclient.request.completed"),34 failed: Notification.kind("httpclient.request.failed"),35 changed: Notification.kind("httpclient.pipeline.changed")36 }37 let registrations = registered(&kinds.started, &given.started).concat(registered(&kinds.completed, &given.completed)).concat(registered(&kinds.failed, &given.failed)).concat(registered(&kinds.changed, &given.changed))38 let mediator = match Mediator.build(registrations) {39 case Ok(built) => built40 case Err(invalid) => panic(Mediator.explain(&invalid))41 }42 Bus{mediator: mediator, kinds: kinds, hearsRequests: Events.hearsRequests(given)}43}444546export fn pipelineChanged(bus: &Bus, client: Str, generation: Int, stage: Events.Stage) -> () {47 let _published = Mediator.publish(&bus.mediator, &bus.kinds.changed, Events.PipelineChanged{client: client, generation: generation, stage: stage})48}495051export fn observer(bus: &Bus) -> Option[Builder.Observer] {52 if !bus.hearsRequests { return None }53 let held = *bus54 Some(Builder.Observer{outer: |seen: Builder.Observation| announcing(held, seen.name), inner: |_seen: Builder.Observation| Handler.passThrough()})55}565758fn registered[N](kind: &Notification.Kind[N, Str], listeners: &Array[fn(N) -> ()]) -> Array[Registration.Registration] {59 var made: Array[Registration.Registration] = []60 var index = 061 for listen in *listeners {62 index = index + 163 made = made.push(Registration.notificationHandler(kind, "listener-" + show(index), fn(event: N, _context: Context.Context) -> Messaging.Outcome[(), Str] {64 listen(event)65 Ok(())66 }))67 }68 made69}707172fn announcing(bus: Bus, client: Str) -> Handler.Handler {73 fn(request: Request.Request, token: Cancel.Token, next: Handler.Send) -> HttpClient.Outcome[Response.Response] {74 let method = Http.methodName(&request.method)75 let _started = Mediator.publish(&bus.mediator, &bus.kinds.started, Events.RequestStarted{client: client, method: method, uri: request.uri})76 let began = clock()77 let outcome = next(request, token)78 let elapsed = clock() - began79 match outcome {80 case Ok(response) => {81 let _completed = Mediator.publish(&bus.mediator, &bus.kinds.completed, Events.RequestCompleted{client: client, method: method, uri: request.uri, status: response.status.code, elapsed: elapsed})82 }83 case Err(failure) => {84 let _failed = Mediator.publish(&bus.mediator, &bus.kinds.failed, Events.RequestFailed{client: client, method: method, uri: request.uri, failure: HttpClient.describe(&failure), elapsed: elapsed})85 }86 }87 outcome88 }89}90