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

Transport.pudu

Pudu366 lines17.5 KB

GitHub ↗
1/** @HttpClient.Transport.Module — the primary handler: pooled keep-alive connections on the wire */2module PuduLangHttpClient.Transport34import Std.Compress.Gzip as Gzip5import Std.Concurrent.Cancel as Cancel6import Std.Http as Http7import Std.Http.Safe as Safe8import PuduLangHttpClient.Constants.Defaults as Defaults9import PuduLangHttpClient.Constants.Messages as Messages10import PuduLangHttpClient.Content as Content11import PuduLangHttpClient.Cookies as Cookies12import PuduLangHttpClient.Domain.Headers as Headers13import PuduLangHttpClient.Domain.Persistence as Persistence14import PuduLangHttpClient.Domain.Redirect as Redirect15import PuduLangHttpClient.Domain.Uri as Uri16import PuduLangHttpClient.Domain.Version as Version17import PuduLangHttpClient.Handler as Handler18import PuduLangHttpClient as HttpClient19import PuduLangHttpClient.Options as RequestOptions20import PuduLangHttpClient.Proxy as Proxy21import PuduLangHttpClient.Request as Request22import PuduLangHttpClient.Response as Response23import PuduLangHttpClient.Transport.Connect as Connect24import PuduLangHttpClient.Transport.Exchange as Exchange25import PuduLangHttpClient.Transport.Hop as Hop26import PuduLangHttpClient.Transport.Pool as Pool27import PuduLangHttpClient.Transport.Wire as Wire28import PuduLangHttpClient.Transport.Writer as Writer29import PuduLangHttpClient.Utils.Template as Template3031/** @HttpClient.Transport.AddressPolicy — which hosts requests may reach */32export type AddressPolicy = AnyAddress | PublicOnly(Array[Str])3334/** @HttpClient.Transport.Options — how connections are kept, and what a response may cost */35export type Options = {36  pooledConnectionLifetime: Int,37  pooledConnectionIdleTimeout: Int,38  maxConnectionsPerServer: Int,39  connectTimeout: Int,40  allowAutoRedirect: Bool,41  maxAutomaticRedirections: Int,42  automaticDecompression: Bool,43  useCookies: Bool,44  cookies: Option[Cookies.Jar],45  proxy: Option[Proxy.Proxy],46  addressPolicy: AddressPolicy,47  maxResponseHeadersLength: Int,48  maxResponseContentLength: Int49}5051/** @HttpClient.Transport.Transport — the options and the pool of one primary handler */52export type Transport = { options: Options, pool: Pool.Pool, jar: Option[Cookies.Jar] }5354/// Connections kept until idle for a minute, no per-server limit, redirects followed up to 50 times,55/// no decompression, no cookies, no proxy, any host, and a 64 MiB body.56export fn defaults() -> Options {57  Options {58    pooledConnectionLifetime: Defaults.POOLED_LIFETIME,59    pooledConnectionIdleTimeout: Defaults.POOLED_IDLE_TIMEOUT,60    maxConnectionsPerServer: Defaults.MAX_CONNECTIONS_PER_SERVER,61    connectTimeout: Defaults.CONNECT_TIMEOUT,62    allowAutoRedirect: true,63    maxAutomaticRedirections: Defaults.MAX_REDIRECTS,64    automaticDecompression: false,65    useCookies: false,66    cookies: None,67    proxy: None,68    addressPolicy: AnyAddress,69    maxResponseHeadersLength: Defaults.MAX_RESPONSE_HEADERS,70    maxResponseContentLength: Defaults.MAX_RESPONSE_BUFFER71  }72}7374/// Every reason the options cannot be followed.75export fn validate(options: &Options) -> Array[Str] {76  var problems: Array[Str] = []77  if !isDuration(options.pooledConnectionLifetime) { problems = problems.push("pooledConnectionLifetime must be at least 0 ms or INFINITE") }78  if !isDuration(options.pooledConnectionIdleTimeout) { problems = problems.push("pooledConnectionIdleTimeout must be at least 0 ms or INFINITE") }79  if !isDuration(options.connectTimeout) || options.connectTimeout == 0 { problems = problems.push("connectTimeout must be positive or INFINITE") }80  if options.maxConnectionsPerServer < 1 { problems = problems.push("maxConnectionsPerServer must be at least 1") }81  if options.maxAutomaticRedirections < 1 { problems = problems.push("maxAutomaticRedirections must be at least 1") }82  if options.maxResponseHeadersLength < 1 { problems = problems.push("maxResponseHeadersLength must be at least 1") }83  if options.maxResponseContentLength < 0 { problems = problems.push("maxResponseContentLength must be at least 0") }84  match options.proxy {85    case Some(proxy) => problems.concat(Proxy.validate(&proxy))86    case None => problems87  }88}8990/// A transport following the options, with its own empty pool and, when cookies are used without a91/// jar, its own jar.92export fn create(options: Options) -> Transport {93  let limits = Persistence.Limits{lifetime: options.pooledConnectionLifetime, idleTimeout: options.pooledConnectionIdleTimeout}94  let jar = match options.cookies {95    case Some(given) => Some(given)96    case None => if options.useCookies { Some(Cookies.jar()) } else { None }97  }98  Transport{options: options, pool: Pool.create(limits, options.maxConnectionsPerServer), jar: jar}99}100101/// The transport as the primary handler of a pipeline.102export fn send(transport: &Transport) -> Handler.Send {103  let held = *transport104  fn(request: Request.Request, token: Cancel.Token) -> HttpClient.Outcome[Response.Response] { sendOn(&held, request, &token) }105}106107/// The key a client sets to cap how many body bytes the transport buffers for one request.108export fn bufferLimit() -> RequestOptions.Key[Int] { RequestOptions.intKey("pudu.httpclient.buffer-limit") }109110/// The connections opened, reused, closed, in use, and idle so far.111export fn statistics(transport: &Transport) -> Pool.Statistics { Pool.statistics(&transport.pool) }112113/// The jar the transport stores cookies in, when it uses cookies.114export fn cookies(transport: &Transport) -> Option[Cookies.Jar] { transport.jar }115116/// Closes every idle connection and stops keeping new ones; requests still in flight finish and117/// requests sent later open connections they close afterwards.118export fn close(transport: &Transport) -> () { Pool.close(&transport.pool) }119120/// Whether the transport has been closed.121export fn isClosed(transport: &Transport) -> Bool { Pool.isClosed(&transport.pool) }122123/** @HttpClient.Transport.Configuring — transport options changed one decision at a time */124export trait Configuring {125  /// The options closing a pooled connection once it has lived this many milliseconds.126  fn withPooledConnectionLifetime(self: &Self, millis: Int) -> Self127  /// The options closing a pooled connection once it has sat unused this many milliseconds.128  fn withPooledConnectionIdleTimeout(self: &Self, millis: Int) -> Self129  /// The options holding at most this many connections open to one origin.130  fn withMaxConnectionsPerServer(self: &Self, count: Int) -> Self131  /// The options giving up on opening a connection after this many milliseconds.132  fn withConnectTimeout(self: &Self, millis: Int) -> Self133  /// The options answering redirects instead of following them.134  fn withoutRedirects(self: &Self) -> Self135  /// The options following at most this many redirects.136  fn withMaxRedirects(self: &Self, count: Int) -> Self137  /// The options asking for gzip and decompressing buffered bodies.138  fn withDecompression(self: &Self) -> Self139  /// The options storing and sending cookies in a jar of the transport's own.140  fn usingCookies(self: &Self) -> Self141  /// The options storing and sending cookies in a jar shared with others.142  fn withCookies(self: &Self, jar: Cookies.Jar) -> Self143  /// The options sending requests through a proxy.144  fn withProxy(self: &Self, proxy: Proxy.Proxy) -> Self145  /// The options limiting which hosts requests may reach.146  fn withAddressPolicy(self: &Self, policy: AddressPolicy) -> Self147  /// The options refusing a response head longer than this many bytes.148  fn withMaxResponseHeadersLength(self: &Self, bytes: Int) -> Self149  /// The options refusing to buffer more than this many body bytes when the request sets no limit.150  fn withMaxResponseContentLength(self: &Self, bytes: Int) -> Self151}152153impl Configuring for Options {154  /// The options closing a pooled connection once it has lived this many milliseconds.155  fn withPooledConnectionLifetime(self: &Self, millis: Int) -> Self { Options{..*self, pooledConnectionLifetime: millis} }156157  /// The options closing a pooled connection once it has sat unused this many milliseconds.158  fn withPooledConnectionIdleTimeout(self: &Self, millis: Int) -> Self { Options{..*self, pooledConnectionIdleTimeout: millis} }159160  /// The options holding at most this many connections open to one origin.161  fn withMaxConnectionsPerServer(self: &Self, count: Int) -> Self { Options{..*self, maxConnectionsPerServer: count} }162163  /// The options giving up on opening a connection after this many milliseconds.164  fn withConnectTimeout(self: &Self, millis: Int) -> Self { Options{..*self, connectTimeout: millis} }165166  /// The options answering redirects instead of following them.167  fn withoutRedirects(self: &Self) -> Self { Options{..*self, allowAutoRedirect: false} }168169  /// The options following at most this many redirects.170  fn withMaxRedirects(self: &Self, count: Int) -> Self { Options{..*self, allowAutoRedirect: true, maxAutomaticRedirections: count} }171172  /// The options asking for gzip and decompressing buffered bodies.173  fn withDecompression(self: &Self) -> Self { Options{..*self, automaticDecompression: true} }174175  /// The options storing and sending cookies in a jar of the transport's own.176  fn usingCookies(self: &Self) -> Self { Options{..*self, useCookies: true} }177178  /// The options storing and sending cookies in a jar shared with others.179  fn withCookies(self: &Self, jar: Cookies.Jar) -> Self { Options{..*self, useCookies: true, cookies: Some(jar)} }180181  /// The options sending requests through a proxy.182  fn withProxy(self: &Self, proxy: Proxy.Proxy) -> Self { Options{..*self, proxy: Some(proxy)} }183184  /// The options limiting which hosts requests may reach.185  fn withAddressPolicy(self: &Self, policy: AddressPolicy) -> Self { Options{..*self, addressPolicy: policy} }186187  /// The options refusing a response head longer than this many bytes.188  fn withMaxResponseHeadersLength(self: &Self, bytes: Int) -> Self { Options{..*self, maxResponseHeadersLength: bytes} }189190  /// The options refusing to buffer more than this many body bytes when the request sets no limit.191  fn withMaxResponseContentLength(self: &Self, bytes: Int) -> Self { Options{..*self, maxResponseContentLength: bytes} }192}193194/// The response to a request, following redirects.195fn sendOn(transport: &Transport, request: Request.Request, token: &Cancel.Token) -> HttpClient.Outcome[Response.Response] {196  let started = clock()197  let deadline = match Cancel.remaining(token) {198    case Some(left) => Some(started + left)199    case None => None200  }201  let version = match Version.negotiate(&request.version, &request.versionPolicy) {202    case Some(found) => found203    case None => { return Err(HttpClient.VersionUnsupported(Http.versionName(&request.version))) }204  }205  if let Some(name) = Headers.unwritable(&request.headers) { return HttpClient.invalid(Template.fill(Messages.BAD_HEADER, &[name])) }206  let limit = match Request.option(&request, &bufferLimit()) {207    case Some(given) => given208    case None => transport.options.maxResponseContentLength209  }210  let reading = Exchange.Reading{headersLimit: transport.options.maxResponseHeadersLength, bodyLimit: limit}211  var current = request212  var hops = 0213  loop {214    if let Err(why) = Cancel.check(token) { return Err(interrupted(&why, started)) }215    let endpoint = match Uri.endpointOf(current.uri) {216      case Some(found) => found217      case None => { return HttpClient.invalid(Template.fill(Messages.NOT_HTTP, &[current.uri])) }218    }219    if let PublicOnly(permitted) = transport.options.addressPolicy {220      if let Err(problem) = Safe.outboundHost(endpoint.host, &permitted) { return Err(HttpClient.NotPermitted(Safe.explain(&problem))) }221    }222    let prepared = prepare(transport, &current)223    let route = Connect.routeTo(&endpoint, &transport.options.proxy)224    let extra = if Connect.isProxied(&route) { proxyHeaders(&transport.options.proxy) } else { [] }225    let form = if Connect.isProxied(&route) { Writer.AbsoluteForm } else { Writer.OriginForm }226    let sending = Hop.Sending {227      written: Writer.render(&prepared, &version, &endpoint, &form, &extra),228      method: prepared.method,229      headers: prepared.headers,230      reading: reading,231      connectTimeout: transport.options.connectTimeout,232      deadline: deadline,233      receiver: prepared.receiver,234      following: transport.options.allowAutoRedirect,235      upload: streamOf(&prepared)236    }237    let exchanged = match Hop.send(&transport.pool, &route, &sending, token) {238      case Ok(found) => found239      case Err(fault) => { return Err(failureOf(&fault, token, started)) }240    }241    let response = responseOf(&exchanged, &current)242    if let Some(jar) = transport.jar { Cookies.store(&jar, current.uri, &Headers.values(&exchanged.head.headers, "set-cookie")) }243    let hop = if transport.options.allowAutoRedirect {244      match Response.header(&response, "location") {245        case Some(location) => replayable(&current, Redirect.next(&current.method, response.status.code, current.uri, location))246        case None => None247      }248    } else { None }249    match hop {250      case Some(next) => {251        hops = hops + 1252        if hops > transport.options.maxAutomaticRedirections { return Err(HttpClient.TooManyRedirects(transport.options.maxAutomaticRedirections)) }253        current = followed(&current, &next)254      }255      case None => { return decoded(transport, response, limit) }256    }257  }258}259260/// The request as it goes on the wire: with the jar's cookies, and asking for gzip when the261/// transport decompresses a buffered body.262fn prepare(transport: &Transport, request: &Request.Request) -> Request.Request {263  var headers = request.headers264  if let Some(jar) = transport.jar {265    if let Some(cookie) = Cookies.header(&jar, request.uri) { headers = Headers.set(&headers, "cookie", cookie) }266  }267  if transport.options.automaticDecompression && request.receiver == None && !Headers.has(&headers, "accept-encoding") {268    headers = Headers.add(&headers, "accept-encoding", "gzip")269  }270  Request.Request{..*request, headers: headers}271}272273/// The producer of a streamed body, when the request has one.274fn streamOf(request: &Request.Request) -> Option[Content.Producer] {275  let content = request.content ?276  content.stream277}278279/// The redirect hop, unless it would send a streamed body that was already sent.280fn replayable(request: &Request.Request, hop: Option[Redirect.Hop]) -> Option[Redirect.Hop] {281  let next = hop ?282  if next.keepsBody && streamOf(request) != None { None } else { Some(next) }283}284285/// The `proxy-authorization` header a plain proxy needs, when it has credentials.286fn proxyHeaders(proxy: &Option[Proxy.Proxy]) -> Array[(Str, Str)] {287  match proxy {288    case Some(found) => {289      match Proxy.authorization(&found) {290        case Some(value) => [("proxy-authorization", value)]291        case None => []292      }293    }294    case None => []295  }296}297298/// The response an exchange read, with its content headers kept apart.299fn responseOf(exchanged: &Exchange.Exchanged, request: &Request.Request) -> Response.Response {300  let (own, described) = Headers.splitContent(&exchanged.head.headers)301  Response.Response {302    status: exchanged.head.status,303    version: Http.versionFrom(exchanged.head.version),304    headers: own,305    content: Content.Content{payload: exchanged.body, headers: described, stream: None},306    trailers: exchanged.trailers,307    request: *request308  }309}310311/// The request a followed redirect sends: the new method and address, the body only when the method312/// kept it, and credentials only within one origin.313fn followed(request: &Request.Request, hop: &Redirect.Hop) -> Request.Request {314  let kept = if hop.keepsCredentials { request.headers } else { Headers.removeAll(&request.headers, &#{"authorization", "cookie", "proxy-authorization"}) }315  Request.Request {316    ..*request,317    method: hop.method,318    uri: hop.uri,319    headers: Headers.remove(&kept, "cookie"),320    content: if hop.keepsBody { request.content } else { None }321  }322}323324/// The response with a gzip body decompressed, when the transport decompresses.325fn decoded(transport: &Transport, response: Response.Response, limit: Int) -> HttpClient.Outcome[Response.Response] {326  let encoding = match Headers.get(&response.content.headers, "content-encoding") {327    case Some(found) => found328    case None => { return Ok(response) }329  }330  let applies = transport.options.automaticDecompression && response.request.receiver == None && Headers.lists(encoding, "gzip")331  if !applies || response.content.payload.isEmpty() { return Ok(response) }332  match Gzip.decompressWithin(&response.content.payload, limit) {333    case Ok(payload) => {334      let headers = Headers.removeAll(&response.content.headers, &#{"content-encoding", "content-length"})335      Ok(Response.Response{..response, content: Content.Content{payload: payload, headers: headers, stream: None}})336    }337    case Err(Gzip.OutputLimitExceeded(_)) => Err(HttpClient.ContentTooLarge(limit))338    case Err(problem) => Err(HttpClient.Protocol(Template.fill(Messages.BAD_ENCODING, &[show(problem)])))339  }340}341342/// The failure a wire fault stands for.343fn failureOf(fault: &Wire.Fault, token: &Cancel.Token, started: Int) -> HttpClient.Failure {344  match fault {345    case Wire.Stalled => {346      match Cancel.reason(token) {347        case Some(why) => interrupted(&why, started)348        case None => HttpClient.TimedOut(clock() - started)349      }350    }351    case Wire.Ended => HttpClient.ResponseEnded352    case Wire.Broken(failure) => failure353  }354}355356/// The failure a fired token stands for: a cancellation, or a timeout after the time spent.357fn interrupted(why: &Cancel.Reason, started: Int) -> HttpClient.Failure {358  match why {359    case Cancel.Requested(reason) => HttpClient.Cancelled(reason)360    case Cancel.DeadlineExceeded => HttpClient.TimedOut(clock() - started)361  }362}363364/// Whether milliseconds name a duration: at least zero, or `Defaults.INFINITE`.365fn isDuration(millis: Int) -> Bool { millis >= 0 || millis == Defaults.INFINITE }366