
Factory.pudu
Pudu196 lines9.8 KB
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 }333435export fn defaults() -> Options { Options{listeners: Events.listeners()} }36373839export fn build(builders: Array[Builder.Builder]) -> Result[Factory, Invalid] { buildWith(&defaults(), builders) }40414243export 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}636465export fn explain(invalid: &Invalid) -> Str { Template.fill(Messages.FACTORY_INVALID, &[invalid.problems.join(Messages.PROBLEM_SEPARATOR)]) }66676869export 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}757677export fn client(factory: &Factory) -> Client.Client { createClient(factory, Defaults.DEFAULT_NAME) }787980export fn typed[T](factory: &Factory, name: Str, make: fn(Client.Client) -> T) -> T { make(createClient(factory, name)) }818283export fn createHandler(factory: &Factory, name: Str) -> Handler.Send {84 let entry = liveEntry(factory, name)85 entry.send86}878889export fn names(factory: &Factory) -> Array[Str] { factory.builders.keys() }909192export 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}979899export fn pendingExpired(factory: &Factory) -> Int { Shared.current(&factory.state).expired.length() }100101102103export 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}107108109110export 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}114115116fn 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}122123124125fn 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}140141142143fn 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}157158159160fn 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}166167168fn closeEntry(factory: &Factory, entry: &Entry.Entry) -> () {169 Entry.close(entry)170 Bus.pipelineChanged(&factory.bus, entry.name, entry.generation, Events.Closed)171}172173174fn 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}180181182fn 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