
Hop.pudu
Pudu88 lines3.4 KB
1/** @Transport.Hop.Module — one request sent over a pooled connection */2module PuduLangHttpClient.Transport.Hop34import Std.Concurrent.Cancel as Cancel5import Std.Http as Http6import Std.Result as Result7import PuduLangHttpClient.Content as Content8import PuduLangHttpClient as HttpClient9import PuduLangHttpClient.Request as Request10import PuduLangHttpClient.Transport.Abort as Abort11import PuduLangHttpClient.Transport.Connect as Connect12import PuduLangHttpClient.Transport.Exchange as Exchange13import PuduLangHttpClient.Transport.Pool as Pool14import PuduLangHttpClient.Transport.Wire as Wire1516/** @Transport.Hop.Sending — everything one hop needs besides the pool and the route */17export type Sending = {18 written: Bytes,19 method: Http.Method,20 headers: Array[(Str, Str)],21 reading: Exchange.Reading,22 connectTimeout: Int,23 deadline: Option[Int],24 receiver: Option[Request.Receiver],25 following: Bool,26 upload: Option[Content.Producer]27}28293031323334export fn send(pool: &Pool.Pool, route: &Connect.Route, sending: &Sending, token: &Cancel.Token) -> Result[Exchange.Exchanged, Wire.Fault] {35 let first = attempt(pool, route, sending, token)36 let outcome = match first {37 case Err((_, true)) => attempt(pool, route, sending, token)38 case _ => first39 }40 Result.mapErr(outcome, |failed: (Wire.Fault, Bool)| failed[0])41}42434445fn attempt(pool: &Pool.Pool, route: &Connect.Route, sending: &Sending, token: &Cancel.Token) -> Result[Exchange.Exchanged, (Wire.Fault, Bool)] {46 let lease = match Pool.acquire(pool, route.key, token, sending.deadline) {47 case Ok(found) => found48 case Err(fault) => { return Err((fault, false)) }49 }50 let (link, created, reused) = match lease {51 case Pool.Reused(idle) => (idle.link, idle.created, true)52 case Pool.Fresh => {53 match Connect.open(route, sending.connectTimeout, sending.deadline) {54 case Ok(opened) => (opened, clock(), false)55 case Err(fault) => {56 Pool.forfeit(pool, route.key)57 return Err((fault, false))58 }59 }60 }61 }62 let guard = Abort.watch(&link, token)63 let outcome = Exchange.exchange(&link, &sending.written, &sending.method, &sending.headers, &sending.reading, sending.deadline, sending.receiver, sending.following, sending.upload)64 if Abort.finish(&guard) {65 Pool.release(pool, route.key, link, created, false)66 return Err((Wire.Stalled, false))67 }68 match outcome {69 case Ok(exchanged) => {70 Pool.release(pool, route.key, link, created, exchanged.reusable)71 Ok(exchanged)72 }73 case Err(trouble) => {74 Pool.release(pool, route.key, link, created, false)75 Err((trouble.fault, reused && !trouble.answered && sending.upload == None && stale(&trouble.fault)))76 }77 }78}798081fn stale(fault: &Wire.Fault) -> Bool {82 match fault {83 case Wire.Ended => true84 case Wire.Broken(HttpClient.Connection(_)) => true85 case _ => false86 }87}88