Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-httpclient

Named HTTP clients for Pudu: a client factory, delegating handlers, pooled keep-alive connections, handler lifetimes, logging, and resilience

0.1.1Apache-2.01

InstallClose

Bus.pudu

Pudu90 lines4.2 KB

GitHub ↗
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 }2829/// A bus delivering each event to its listeners in the order they were added.30export 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}4445/// Publishes a pipeline change.46export 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}4950/// The observer announcing every request through a pipeline, when anyone listens to requests.51export 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}5657/// One registration per listener of a kind, named by position so each stays distinct.58fn 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}7071/// The handler publishing a request's start and end.72fn 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