
Transport.pudu
Pudu366 lines17.5 KB
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] }53545556export 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}737475export 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}89909192export 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}100101102export 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}106107108export fn bufferLimit() -> RequestOptions.Key[Int] { RequestOptions.intKey("pudu.httpclient.buffer-limit") }109110111export fn statistics(transport: &Transport) -> Pool.Statistics { Pool.statistics(&transport.pool) }112113114export fn cookies(transport: &Transport) -> Option[Cookies.Jar] { transport.jar }115116117118export fn close(transport: &Transport) -> () { Pool.close(&transport.pool) }119120121export 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 126 fn withPooledConnectionLifetime(self: &Self, millis: Int) -> Self127 128 fn withPooledConnectionIdleTimeout(self: &Self, millis: Int) -> Self129 130 fn withMaxConnectionsPerServer(self: &Self, count: Int) -> Self131 132 fn withConnectTimeout(self: &Self, millis: Int) -> Self133 134 fn withoutRedirects(self: &Self) -> Self135 136 fn withMaxRedirects(self: &Self, count: Int) -> Self137 138 fn withDecompression(self: &Self) -> Self139 140 fn usingCookies(self: &Self) -> Self141 142 fn withCookies(self: &Self, jar: Cookies.Jar) -> Self143 144 fn withProxy(self: &Self, proxy: Proxy.Proxy) -> Self145 146 fn withAddressPolicy(self: &Self, policy: AddressPolicy) -> Self147 148 fn withMaxResponseHeadersLength(self: &Self, bytes: Int) -> Self149 150 fn withMaxResponseContentLength(self: &Self, bytes: Int) -> Self151}152153impl Configuring for Options {154 155 fn withPooledConnectionLifetime(self: &Self, millis: Int) -> Self { Options{..*self, pooledConnectionLifetime: millis} }156157 158 fn withPooledConnectionIdleTimeout(self: &Self, millis: Int) -> Self { Options{..*self, pooledConnectionIdleTimeout: millis} }159160 161 fn withMaxConnectionsPerServer(self: &Self, count: Int) -> Self { Options{..*self, maxConnectionsPerServer: count} }162163 164 fn withConnectTimeout(self: &Self, millis: Int) -> Self { Options{..*self, connectTimeout: millis} }165166 167 fn withoutRedirects(self: &Self) -> Self { Options{..*self, allowAutoRedirect: false} }168169 170 fn withMaxRedirects(self: &Self, count: Int) -> Self { Options{..*self, allowAutoRedirect: true, maxAutomaticRedirections: count} }171172 173 fn withDecompression(self: &Self) -> Self { Options{..*self, automaticDecompression: true} }174175 176 fn usingCookies(self: &Self) -> Self { Options{..*self, useCookies: true} }177178 179 fn withCookies(self: &Self, jar: Cookies.Jar) -> Self { Options{..*self, useCookies: true, cookies: Some(jar)} }180181 182 fn withProxy(self: &Self, proxy: Proxy.Proxy) -> Self { Options{..*self, proxy: Some(proxy)} }183184 185 fn withAddressPolicy(self: &Self, policy: AddressPolicy) -> Self { Options{..*self, addressPolicy: policy} }186187 188 fn withMaxResponseHeadersLength(self: &Self, bytes: Int) -> Self { Options{..*self, maxResponseHeadersLength: bytes} }189190 191 fn withMaxResponseContentLength(self: &Self, bytes: Int) -> Self { Options{..*self, maxResponseContentLength: bytes} }192}193194195fn 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, ¤t)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, ¤t)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(¤t, Redirect.next(¤t.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(¤t, &next)254 }255 case None => { return decoded(transport, response, limit) }256 }257 }258}259260261262fn 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}272273274fn streamOf(request: &Request.Request) -> Option[Content.Producer] {275 let content = request.content ?276 content.stream277}278279280fn 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}284285286fn 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}297298299fn 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}310311312313fn 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}323324325fn 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}341342343fn 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}355356357fn 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}363364365fn isDuration(millis: Int) -> Bool { millis >= 0 || millis == Defaults.INFINITE }366