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

Pool.pudu

Pudu137 lines6.4 KB

GitHub ↗
1/** @Transport.Pool.Module — idle connections kept per origin and handed out again */2module PuduLangHttpClient.Transport.Pool34import Std.Concurrent.Cancel as Cancel5import Std.Json as Json6import PuduLangHttpClient.Constants.Defaults as Defaults7import PuduLangHttpClient.Domain.Persistence as Persistence8import PuduLangHttpClient.Transport.Wire as Wire9import PuduLangHttpClient.Utils.Shared as Shared1011/** @Transport.Pool.Statistics — connections opened, reused, closed, in use, and idle so far */12export type Statistics = { opened: Int, reused: Int, closed: Int, active: Int, idle: Int } derives Json.Encode1314/** @Transport.Pool.Lease — a pooled connection to reuse, or permission to open a new one */15export type Lease = Reused(Persistence.Idle[Wire.Link]) | Fresh1617/** @Transport.Pool.Slot — one origin's idle connections and how many it holds open */18type Slot = { idle: Array[Persistence.Idle[Wire.Link]], open: Int }1920/** @Transport.Pool.State — every origin's slot, whether the pool is closed, and its counters */21type State = { slots: Map[Str, Slot], closed: Bool, counts: Statistics }2223/** @Transport.Pool.Pool — the shared state and the limits it keeps */24export type Pool = { state: Shared.Shared[State], limits: Persistence.Limits, maxPerServer: Int }2526/// An empty pool keeping connections within the limits, at most `maxPerServer` open per origin.27export fn create(limits: Persistence.Limits, maxPerServer: Int) -> Pool {28  let counts = Statistics{opened: 0, reused: 0, closed: 0, active: 0, idle: 0}29  Pool{state: Shared.shared(State{slots: mapOf([]), closed: false, counts: counts}), limits: limits, maxPerServer: maxPerServer}30}3132/// A lease for one request to the origin `key`, waiting while the origin is at its limit until a33/// connection is released, the token fires, or the deadline passes.34export fn acquire(pool: &Pool, key: Str, token: &Cancel.Token, deadline: Option[Int]) -> Result[Lease, Wire.Fault] {35  loop {36    let now = clock()37    let (granted, stale) = Shared.change(&pool.state, |state: State| grant(&state, key, &pool.limits, pool.maxPerServer, now))38    for one in stale { Wire.close(&one.link) }39    match granted {40      case Some(lease) => { return Ok(lease) }41      case None => {42        if Wire.budget(deadline) == None { return Err(Wire.Stalled) }43        if let Err(_) = Cancel.pause(token, Defaults.POOL_WAIT_SLICE) { return Err(Wire.Stalled) }44      }45    }46  }47}4849/// Returns a connection after its exchange: kept idle when it is reusable and the pool is open,50/// closed otherwise.51export fn release(pool: &Pool, key: Str, link: Wire.Link, created: Int, reusable: Bool) -> () {52  let now = clock()53  let kept = Shared.change(&pool.state, fn(state: State) -> (State, Bool) {54      let slot = slotOf(&state, key)55      let counts = Statistics{..state.counts, active: state.counts.active - 1}56      if reusable && !state.closed {57        let idle = Persistence.Idle{link: link, created: created, lastUsed: now}58        (State{..state, slots: state.slots.insert(key, Slot{..slot, idle: slot.idle.push(idle)}), counts: counts}, true)59      } else {60        (State{..state, slots: state.slots.insert(key, Slot{..slot, open: slot.open - 1}), counts: Statistics{..counts, closed: counts.closed + 1}}, false)61      }62    })63  if !kept { Wire.close(&link) }64}6566/// Gives back permission to open a connection that could not be opened.67export fn forfeit(pool: &Pool, key: Str) -> () {68  Shared.update(&pool.state, fn(state: State) -> State {69      let slot = slotOf(&state, key)70      State{..state, slots: state.slots.insert(key, Slot{..slot, open: slot.open - 1}), counts: Statistics{..state.counts, opened: state.counts.opened - 1, active: state.counts.active - 1}}71    })72}7374/// The counters as they stand, with the idle connections counted now.75export fn statistics(pool: &Pool) -> Statistics {76  let state = Shared.current(&pool.state)77  var idle = 078  for slot in state.slots.values() { idle = idle + slot.idle.length() }79  Statistics{..state.counts, idle: idle}80}8182/// Closes every idle connection and stops keeping connections: one released later is closed.83export fn close(pool: &Pool) -> () {84  let idle = Shared.change(&pool.state, fn(state: State) -> (State, Array[Persistence.Idle[Wire.Link]]) {85      var dropped: Array[Persistence.Idle[Wire.Link]] = []86      var slots: Map[Str, Slot] = mapOf([])87      for (key, slot) in state.slots.entries() {88        dropped = dropped.concat(slot.idle)89        slots = slots.insert(key, Slot{idle: [], open: slot.open - slot.idle.length()})90      }91      (State{slots: slots, closed: true, counts: Statistics{..state.counts, closed: state.counts.closed + dropped.length()}}, dropped)92    })93  for one in idle { Wire.close(&one.link) }94}9596/// Whether the pool has been closed.97export fn isClosed(pool: &Pool) -> Bool { Shared.current(&pool.state).closed }9899/// The state after sweeping every origin's idle connections and granting a lease for `key` when one100/// is available, with the lease and the connections swept out.101fn grant(state: &State, key: Str, limits: &Persistence.Limits, maxPerServer: Int, now: Int) -> (State, (Option[Lease], Array[Persistence.Idle[Wire.Link]])) {102  var stale: Array[Persistence.Idle[Wire.Link]] = []103  var slots: Map[Str, Slot] = mapOf([])104  for (name, slot) in state.slots.entries() {105    let (usable, expired) = Persistence.sweep(&slot.idle, limits, now)106    stale = stale.concat(expired)107    slots = slots.insert(name, Slot{idle: usable, open: slot.open - expired.length()})108  }109  let counts = Statistics{..state.counts, closed: state.counts.closed + stale.length()}110  let slot = match slots.get(key) {111    case Some(found) => found112    case None => Slot{idle: [], open: 0}113  }114  match slot.idle {115    case [..rest, last] => {116      let granted = Slot{..slot, idle: rest}117      (State{..*state, slots: slots.insert(key, granted), counts: Statistics{..counts, reused: counts.reused + 1, active: counts.active + 1}}, (Some(Reused(last)), stale))118    }119    case _ => {120      if slot.open < maxPerServer {121        let reserved = Slot{..slot, open: slot.open + 1}122        (State{..*state, slots: slots.insert(key, reserved), counts: Statistics{..counts, opened: counts.opened + 1, active: counts.active + 1}}, (Some(Fresh), stale))123      } else {124        (State{..*state, slots: slots, counts: counts}, (None, stale))125      }126    }127  }128}129130/// The slot of an origin, or an empty one.131fn slotOf(state: &State, key: Str) -> Slot {132  match state.slots.get(key) {133    case Some(found) => found134    case None => Slot{idle: [], open: 0}135  }136}137