
Pool.pudu
Pudu137 lines6.4 KB
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 }252627export 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}31323334export 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}48495051export 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}656667export 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}737475export 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}818283export 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}959697export fn isClosed(pool: &Pool) -> Bool { Shared.current(&pool.state).closed }9899100101fn 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}129130131fn 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