
Exchange.pudu
Pudu236 lines10.2 KB
1/** @Transport.Exchange.Module — one request written and one response read over a connection */2module PuduLangHttpClient.Transport.Exchange34import Std.Bytes as ByteSeq5import Std.Http as Http6import Std.Http.Message as Message7import PuduLangHttpClient.Constants.Defaults as Defaults8import PuduLangHttpClient.Constants.Messages as Messages9import PuduLangHttpClient.Content as Content10import PuduLangHttpClient.Domain.Chunked as Chunked11import PuduLangHttpClient.Domain.Persistence as Persistence12import PuduLangHttpClient.Domain.Redirect as Redirect13import PuduLangHttpClient as HttpClient14import PuduLangHttpClient.Request as Request15import PuduLangHttpClient.Transport.Wire as Wire16import PuduLangHttpClient.Utils.Template as Template1718/** @Transport.Exchange.Reading — how much of a response is read before it is refused */19export type Reading = { headersLimit: Int, bodyLimit: Int }2021/** @Transport.Exchange.Head — a response's version, status, and headers */22export type Head = { version: Str, status: Http.Status, headers: Array[(Str, Str)] }2324/** @Transport.Exchange.Exchanged — a complete response and whether its connection may be reused */25export type Exchanged = { head: Head, body: Bytes, trailers: Array[(Str, Str)], reusable: Bool }2627/** @Transport.Exchange.Trouble — why an exchange failed, and whether any of the response had arrived */28export type Trouble = { fault: Wire.Fault, answered: Bool }2930/** @Transport.Exchange.Stream — bytes read but not yet consumed, and whether any arrived at all */31type Stream = { pending: Bytes, answered: Bool }323334const HEAD_END: Str = "\r\n\r\n"3536373839export fn exchange(40 link: &Wire.Link,41 written: &Bytes,42 method: &Http.Method,43 requestHeaders: &Array[(Str, Str)],44 reading: &Reading,45 deadline: Option[Int],46 receiver: Option[Request.Receiver],47 following: Bool,48 upload: Option[Content.Producer]49) -> Result[Exchanged, Trouble] {50 if let Err(fault) = Wire.write(link, written, deadline) { return Err(Trouble{fault: fault, answered: false}) }51 if let Some(producer) = upload {52 if let Err(fault) = streamBody(link, producer, deadline) { return Err(Trouble{fault: fault, answered: false}) }53 }54 var stream = Stream{pending: ByteSeq.empty(), answered: false}55 loop {56 let (head, headBytes, rest) = readHead(link, &stream, reading, deadline) ?57 stream = rest58 if head.status.code >= 200 || head.status.code == 101 {59 let framing = match Message.framingOf(&headBytes, method) {60 case Some(found) => found61 case None => { return Err(broken(HttpClient.Protocol(Template.fill(Messages.BAD_HEAD, &["no framing"])), true)) }62 }63 let taking = if following && Redirect.isRedirect(head.status.code) { None } else { receiver }64 let (body, trailers, complete) = readBody(link, &stream, &framing, reading, deadline, taking) ?65 let reusable = complete && Persistence.keepsAlive(head.version, requestHeaders, &head.headers, &framing) && head.status.code != 10166 return Ok(Exchanged{head: head, body: body, trailers: trailers, reusable: reusable})67 }68 }69}707172fn streamBody(link: &Wire.Link, producer: Content.Producer, deadline: Option[Int]) -> Result[(), Wire.Fault] {73 loop {74 match producer() {75 case Some(piece) => {76 if !piece.isEmpty() {77 Wire.write(link, &(hexOf(piece.length()) + "\r\n").toBytes().concat(piece).concat("\r\n".toBytes()), deadline) ?78 }79 }80 case None => { return Wire.write(link, &"0\r\n\r\n".toBytes(), deadline) }81 }82 }83}848586fn hexOf(count: Int) -> Str {87 if count == 0 { return "0" }88 var out = ""89 var left = count90 while left > 0 {91 let digit = left % 1692 out = "0123456789abcdef".slice(digit, digit + 1) + out93 left = left / 1694 }95 out96}979899fn broken(failure: HttpClient.Failure, answered: Bool) -> Trouble { Trouble{fault: Wire.Broken(failure), answered: answered} }100101102fn readHead(link: &Wire.Link, stream: &Stream, reading: &Reading, deadline: Option[Int]) -> Result[(Head, Bytes, Stream), Trouble] {103 var pending = stream.pending104 var answered = stream.answered105 loop {106 let end = pending.indexOf(HEAD_END.toBytes())107 if end + 4 > reading.headersLimit || (end < 0 && pending.length() > reading.headersLimit) { return Err(broken(HttpClient.HeadersTooLarge(reading.headersLimit), true)) }108 if end >= 0 {109 let headBytes = pending.take(end + 4)110 let head = parseHead(&headBytes) ?111 return Ok((head, headBytes, Stream{pending: pending.drop(end + 4), answered: true}))112 }113 match Wire.read(link, deadline) {114 case Ok(Some(piece)) => {115 pending = pending.concat(piece)116 answered = true117 }118 case Ok(None) => { return Err(Trouble{fault: Wire.Ended, answered: answered}) }119 case Err(fault) => { return Err(Trouble{fault: fault, answered: answered}) }120 }121 }122}123124125fn parseHead(headBytes: &Bytes) -> Result[Head, Trouble] {126 let text = match ByteSeq.toText(headBytes) {127 case Ok(found) => found128 case Err(_) => { return Err(broken(HttpClient.Protocol(Template.fill(Messages.BAD_HEAD, &["not text"])), true)) }129 }130 let firstLine = text.split("\r\n")[0]131 let version = firstLine.split(" ")[0]132 if !version.startsWith("HTTP/1.") { return Err(broken(HttpClient.Protocol(Template.fill(Messages.BAD_STATUS_LINE, &[firstLine])), true)) }133 match Message.parseResponse(text) {134 case Ok(parsed) => Ok(Head{version: version, status: parsed.status, headers: parsed.headers})135 case Err(problem) => Err(broken(HttpClient.Protocol(Template.fill(Messages.BAD_HEAD, &[Message.explain(&problem)])), true))136 }137}138139140141fn readBody(142 link: &Wire.Link,143 stream: &Stream,144 framing: &Message.Framing,145 reading: &Reading,146 deadline: Option[Int],147 receiver: Option[Request.Receiver]148) -> Result[(Bytes, Array[(Str, Str)], Bool), Trouble] {149 match framing {150 case Message.HeadOnly => Ok((ByteSeq.empty(), [], stream.pending.isEmpty()))151 case Message.Length(size) => {152 if receiver == None && size > reading.bodyLimit { return Err(broken(HttpClient.ContentTooLarge(reading.bodyLimit), true)) }153 readLength(link, stream, size, deadline, receiver)154 }155 case Message.Chunked => readChunked(link, stream, reading, deadline, receiver)156 case Message.UntilClose => readToClose(link, stream, reading, deadline, receiver)157 }158}159160161fn readLength(link: &Wire.Link, stream: &Stream, size: Int, deadline: Option[Int], receiver: Option[Request.Receiver]) -> Result[(Bytes, Array[(Str, Str)], Bool), Trouble] {162 var collected: Array[Bytes] = []163 var left = size164 var piece = stream.pending165 loop {166 let taken = piece.take(left)167 left = left - taken.length()168 if !taken.isEmpty() {169 match receiver {170 case Some(take) => if !take(taken) { return Ok((ByteSeq.empty(), [], false)) }171 case None => { collected = collected.push(taken) }172 }173 }174 if left == 0 { return Ok((ByteSeq.join(&collected), [], piece.length() == taken.length())) }175 match Wire.read(link, deadline) {176 case Ok(Some(next)) => { piece = next }177 case Ok(None) => { return Err(Trouble{fault: Wire.Broken(HttpClient.ResponseEnded), answered: true}) }178 case Err(fault) => { return Err(Trouble{fault: fault, answered: true}) }179 }180 }181}182183184fn readChunked(link: &Wire.Link, stream: &Stream, reading: &Reading, deadline: Option[Int], receiver: Option[Request.Receiver]) -> Result[(Bytes, Array[(Str, Str)], Bool), Trouble] {185 let limit = if receiver == None { reading.bodyLimit } else { Defaults.UNLIMITED }186 var decoder = Chunked.decoder()187 var collected: Array[Bytes] = []188 var piece = stream.pending189 loop {190 match Chunked.feed(&decoder, &piece, limit) {191 case Err(Chunked.Exceeded) => { return Err(broken(HttpClient.ContentTooLarge(reading.bodyLimit), true)) }192 case Err(Chunked.Malformed(reason)) => { return Err(broken(HttpClient.Protocol(Template.fill(Messages.BAD_CHUNK, &[reason])), true)) }193 case Ok((next, decoded)) => {194 decoder = next195 if !decoded.isEmpty() {196 match receiver {197 case Some(take) => if !take(decoded) { return Ok((ByteSeq.empty(), [], false)) }198 case None => { collected = collected.push(decoded) }199 }200 }201 }202 }203 if Chunked.isDone(&decoder) { return Ok((ByteSeq.join(&collected), decoder.trailers, Chunked.leftover(&decoder).isEmpty())) }204 match Wire.read(link, deadline) {205 case Ok(Some(next)) => { piece = next }206 case Ok(None) => { return Err(Trouble{fault: Wire.Broken(HttpClient.ResponseEnded), answered: true}) }207 case Err(fault) => { return Err(Trouble{fault: fault, answered: true}) }208 }209 }210}211212213fn readToClose(link: &Wire.Link, stream: &Stream, reading: &Reading, deadline: Option[Int], receiver: Option[Request.Receiver]) -> Result[(Bytes, Array[(Str, Str)], Bool), Trouble] {214 var collected: Array[Bytes] = []215 var size = 0216 var piece = stream.pending217 loop {218 if !piece.isEmpty() {219 match receiver {220 case Some(take) => if !take(piece) { return Ok((ByteSeq.empty(), [], false)) }221 case None => {222 size = size + piece.length()223 if size > reading.bodyLimit { return Err(broken(HttpClient.ContentTooLarge(reading.bodyLimit), true)) }224 collected = collected.push(piece)225 }226 }227 }228 match Wire.read(link, deadline) {229 case Ok(Some(next)) => { piece = next }230 case Ok(None) => { return Ok((ByteSeq.join(&collected), [], false)) }231 case Err(Wire.Ended) => { return Ok((ByteSeq.join(&collected), [], false)) }232 case Err(fault) => { return Err(Trouble{fault: fault, answered: true}) }233 }234 }235}236