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

Exchange.pudu

Pudu236 lines10.2 KB

GitHub ↗
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 }3233/// The blank line that ends a head.34const HEAD_END: Str = "\r\n\r\n"3536/// Writes the request, streaming the producer's body in chunks when there is one, and reads its final37/// response, skipping interim `1xx` responses; the body is buffered, or handed to the receiver as it38/// arrives unless it belongs to a redirect being followed.39export 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}7071/// Writes every piece the producer gives as a chunk, then the terminating chunk.72fn 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}8485/// A count written in hexadecimal.86fn 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}9798/// A trouble carrying a failure.99fn broken(failure: HttpClient.Failure, answered: Bool) -> Trouble { Trouble{fault: Wire.Broken(failure), answered: answered} }100101/// Reads until the next head is complete, and answers it, its bytes, and what followed it.102fn 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}123124/// The version, status, and headers a head states.125fn 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}138139/// Reads the body the framing describes; answers it (empty when a receiver took it), its trailers,140/// and whether it was read to its end with nothing left over.141fn 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}159160/// Reads a body of a stated length.161fn 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}182183/// Reads a chunked body and its trailers.184fn 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}211212/// Reads a body that ends when the connection closes.213fn 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