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

Factory.pudu

Pudu196 lines9.8 KB

GitHub ↗
1/** @HttpClient.Factory.Module — named clients over handler pipelines the factory owns and renews */2module PuduLangHttpClient.Factory34import Std.Concurrent.Cancel as Cancel5import PuduLangHttpClient.Client as Client6import PuduLangHttpClient.Constants.Defaults as Defaults7import PuduLangHttpClient.Constants.Messages as Messages8import PuduLangHttpClient.Domain.Expiry as Expiry9import PuduLangHttpClient.Events.Bus as Bus10import PuduLangHttpClient.Events as Events11import PuduLangHttpClient.Factory.Builder as Builder12import PuduLangHttpClient.Factory.Entry as Entry13import PuduLangHttpClient.Handler as Handler14import PuduLangHttpClient as HttpClient15import PuduLangHttpClient.Request as Request16import PuduLangHttpClient.Response as Response17import PuduLangHttpClient.Transport.Pool as Pool18import PuduLangHttpClient.Transport as Transport19import PuduLangHttpClient.Utils.Shared as Shared20import PuduLangHttpClient.Utils.Template as Template2122/** @HttpClient.Factory.Options — who hears about the factory's requests and pipelines */23export type Options = { listeners: Events.Listeners }2425/** @HttpClient.Factory.Invalid — every reason a factory could not be built */26export type Invalid = { problems: Array[Str] }2728/** @HttpClient.Factory.State — the live pipeline of each name, the expired ones, and a counter */29type State = { active: Map[Str, Entry.Entry], expired: Array[Entry.Entry], generation: Int }3031/** @HttpClient.Factory.Factory — the merged configuration of every name and the pipelines built */32export type Factory = { builders: Map[Str, Builder.Builder], defaults: Builder.Builder, state: Shared.Shared[State], bus: Bus.Bus }3334/// No listeners.35export fn defaults() -> Options { Options{listeners: Events.listeners()} }3637/// A factory configured by the builders; builders of the same name are combined in order, and the38/// defaults' builders apply to every name first. Every problem with any name is reported at once.39export fn build(builders: Array[Builder.Builder]) -> Result[Factory, Invalid] { buildWith(&defaults(), builders) }4041/// A factory configured by the builders, telling the options' listeners about every request through42/// its clients and every pipeline it creates, expires, and closes.43export fn buildWith(options: &Options, builders: Array[Builder.Builder]) -> Result[Factory, Invalid] {44  var shared = Builder.defaults()45  var named: Map[Str, Builder.Builder] = mapOf([])46  for builder in builders {47    if builder.everyClient {48      shared = Builder.combine(&shared, &builder)49    } else {50      named = named.insert(builder.name, match named.get(builder.name) {51          case Some(earlier) => Builder.combine(&earlier, &builder)52          case None => builder53        })54    }55  }56  var merged: Map[Str, Builder.Builder] = mapOf([])57  for (name, builder) in named.entries() { merged = merged.insert(name, Builder.combine(&shared, &builder)) }58  var problems: Array[Str] = []59  for builder in merged.values().push(shared) { problems = problems.concat(problemsOf(&builder)) }60  if !problems.isEmpty() { return Err(Invalid{problems: problems}) }61  Ok(Factory{builders: merged, defaults: shared, state: Shared.shared(State{active: mapOf([]), expired: [], generation: 0}), bus: Bus.create(&options.listeners)})62}6364/// The sentence describing every problem.65export fn explain(invalid: &Invalid) -> Str { Template.fill(Messages.FACTORY_INVALID, &[invalid.problems.join(Messages.PROBLEM_SEPARATOR)]) }6667/// A client of the given name over that name's current pipeline; a name never configured gets the68/// defaults. Clients are cheap: make one where it is needed rather than keeping it.69export fn createClient(factory: &Factory, name: Str) -> Client.Client {70  let builder = builderOf(factory, name)71  var made = Client.create(createHandler(factory, name)).named(name)72  for action in builder.clientActions { made = action(made) }73  made74}7576/// The client with the default name.77export fn client(factory: &Factory) -> Client.Client { createClient(factory, Defaults.DEFAULT_NAME) }7879/// A typed client: whatever `make` builds around a client of the given name.80export fn typed[T](factory: &Factory, name: Str, make: fn(Client.Client) -> T) -> T { make(createClient(factory, name)) }8182/// The current handler pipeline of a name, built when the name has none or its lifetime ran out.83export fn createHandler(factory: &Factory, name: Str) -> Handler.Send {84  let entry = liveEntry(factory, name)85  entry.send86}8788/// Every name configured on its own.89export fn names(factory: &Factory) -> Array[Str] { factory.builders.keys() }9091/// The connection counters of a name's current pipeline, when its primary handler is a transport.92export fn statistics(factory: &Factory, name: Str) -> Option[Pool.Statistics] {93  let entry = Shared.current(&factory.state).active.get(name) ?94  let transport = entry.transport ?95  Some(Transport.statistics(&transport))96}9798/// How many expired pipelines are still waiting to be closed.99export fn pendingExpired(factory: &Factory) -> Int { Shared.current(&factory.state).expired.length() }100101/// Closes every pipeline's transport, live or expired; clients made earlier keep working and open102/// connections they close afterwards.103export fn dispose(factory: &Factory) -> () {104  let entries = Shared.change(&factory.state, |state: State| (State{..state, active: mapOf([]), expired: []}, state.active.values().concat(state.expired)))105  for entry in entries { closeEntry(factory, &entry) }106}107108/// A sender that asks the factory for the name's current pipeline on every request, so code holding109/// it for a long time still moves to each new pipeline.110export fn rotatingHandler(factory: &Factory, name: Str) -> Handler.Send {111  let held = *factory112  fn(request: Request.Request, token: Cancel.Token) -> HttpClient.Outcome[Response.Response] { createHandler(&held, name)(request, token) }113}114115/// The merged configuration of a name.116fn builderOf(factory: &Factory, name: Str) -> Builder.Builder {117  match factory.builders.get(name) {118    case Some(found) => found119    case None => Builder.Builder{..factory.defaults, name: name, everyClient: false}120  }121}122123/// The live pipeline of a name; a new one replaces it once its lifetime ran out, and expired124/// pipelines whose grace has passed are closed.125fn liveEntry(factory: &Factory, name: Str) -> Entry.Entry {126  let now = clock()127  let live = Shared.current(&factory.state).active.get(name)128  if let Some(entry) = live {129    if !Expiry.elapsed(entry.created, entry.lifetime, now) { return entry }130  }131  let builder = builderOf(factory, name)132  let generation = Shared.change(&factory.state, |state: State| (State{..state, generation: state.generation + 1}, state.generation + 1))133  let fresh = Entry.build(&withBus(factory, &builder), generation)134  let (chosen, retired, closing) = Shared.change(&factory.state, |state: State| install(&state, &fresh, now))135  if chosen.generation != fresh.generation { Entry.close(&fresh) } else { Bus.pipelineChanged(&factory.bus, name, fresh.generation, Events.Created) }136  for entry in retired { Bus.pipelineChanged(&factory.bus, entry.name, entry.generation, Events.Expired) }137  for entry in closing { closeEntry(factory, &entry) }138  chosen139}140141/// The state with a freshly built pipeline installed unless another thread installed a live one142/// first; answers the pipeline to use, the one it replaced, and the expired ones to close now.143fn install(state: &State, fresh: &Entry.Entry, now: Int) -> (State, (Entry.Entry, Array[Entry.Entry], Array[Entry.Entry])) {144  let (keep, closing) = (state.expired.filter(|entry: Entry.Entry| !graceOver(&entry, now)), state.expired.filter(|entry: Entry.Entry| graceOver(&entry, now)))145  match state.active.get(fresh.name) {146    case Some(entry) => {147      if !Expiry.elapsed(entry.created, entry.lifetime, now) {148        (State{..*state, expired: keep}, (entry, [], closing))149      } else {150        let retired = Entry.Entry{..entry, expiredAt: Some(now)}151        (State{..*state, active: state.active.insert(fresh.name, *fresh), expired: keep.push(retired)}, (*fresh, [retired], closing))152      }153    }154    case None => (State{..*state, active: state.active.insert(fresh.name, *fresh), expired: keep}, (*fresh, [], closing))155  }156}157158/// Whether an expired pipeline has outlived one more lifetime, after which clients made from it159/// are assumed gone.160fn graceOver(entry: &Entry.Entry, now: Int) -> Bool {161  match entry.expiredAt {162    case Some(at) => Expiry.elapsed(at, entry.lifetime, now)163    case None => false164  }165}166167/// Closes a pipeline's transport and says so.168fn closeEntry(factory: &Factory, entry: &Entry.Entry) -> () {169  Entry.close(entry)170  Bus.pipelineChanged(&factory.bus, entry.name, entry.generation, Events.Closed)171}172173/// The builder with the bus's observer outermost, when anyone listens to requests.174fn withBus(factory: &Factory, builder: &Builder.Builder) -> Builder.Builder {175  match Bus.observer(&factory.bus) {176    case Some(observer) => Builder.Builder{..*builder, observers: [observer].concat(builder.observers)}177    case None => *builder178  }179}180181/// Every reason a name's merged configuration cannot be followed.182fn problemsOf(builder: &Builder.Builder) -> Array[Str] {183  var probe = Client.create(|request: Request.Request, _token: Cancel.Token| Ok(Response.answer(&request, 200)))184  for action in builder.clientActions { probe = action(probe) }185  var problems = Client.validate(&probe)186  let lifetime = Builder.lifetimeOf(builder)187  if lifetime < 0 && lifetime != Defaults.INFINITE { problems = problems.push("handlerLifetime must be at least 0 ms or INFINITE") }188  match builder.primary {189    case Some(Builder.Custom(_)) => {}190    case Some(Builder.Pooled(options)) => { problems = problems.concat(Transport.validate(&Entry.transportOptions(builder, options))) }191    case None => { problems = problems.concat(Transport.validate(&Entry.transportOptions(builder, Transport.defaults()))) }192  }193  let name = if builder.everyClient { "(defaults)" } else { builder.name }194  problems.map(|problem: Str| Template.fill(Messages.CLIENT_PROBLEM, &[name, problem]))195}196