A pure, bring-your-own-I/O implementation of HTTP/1.1 (port of python-hyper/h11)
Dependencies
moon add bobzhang/h11import {
"bobzhang/h11",
}///|
test "quick start: one request/response cycle" {
let client = @h11.Connection::new(Client)
let server = @h11.Connection::new(Server)
// The client turns events into bytes...
let request = @h11.Request::new(method_=b"GET", target=b"/hello", headers=[
(b"Host", b"example.com"),
])
let wire = client.send(Request(request)).unwrap() +
client.send(EndOfMessage(@h11.EndOfMessage::new())).unwrap()
inspect(
@utf8.decode(wire),
content="GET /hello HTTP/1.1\r\nHost: example.com\r\n\r\n",
)
// ...and the server turns bytes back into events.
server.receive_data(wire)
guard server.next_event() is Event(Request(req)) else { fail("no request") }
assert_eq(req.target, b"/hello")
guard server.next_event() is Event(EndOfMessage(_)) else { fail("no EOM") }
assert_eq(server.next_event(), NeedData)
// The server replies. No Content-Length was given, so h11 picks chunked
// transfer encoding automatically because the client speaks HTTP/1.1.
let response = @h11.Response::new(status_code=200, headers=[
(b"Content-Type", b"text/plain"),
])
let wire = server.send(Response(response)).unwrap() +
server.send(Data(@h11.Data::new(b"hi!"))).unwrap() +
server.send(EndOfMessage(@h11.EndOfMessage::new())).unwrap()
inspect(
@utf8.decode(wire),
content="HTTP/1.1 200 \r\nContent-Type: text/plain\r\nTransfer-Encoding: chunked\r\n\r\n3\r\nhi!\r\n0\r\n\r\n",
)
// The client parses the response.
client.receive_data(wire)
guard client.next_event() is Event(Response(resp)) else { fail("no resp") }
assert_eq(resp.status_code, 200)
guard client.next_event() is Event(Data(body)) else { fail("no data") }
assert_eq(body.data, b"hi!")
guard client.next_event() is Event(EndOfMessage(_)) else { fail("no EOM") }
// Both sides are DONE, so the connection can be reused.
assert_eq(client.states(), { client: Done, server: Done, })
client.start_next_cycle()
server.start_next_cycle()
}| Operation | What it does |
|---|---|
| conn.receive_data(bytes) | Append bytes you read from the network to the internal buffer. Pass b"" to signal end-of-file. |
| conn.next_event() | Parse the next event from the buffer. Returns Event(event), NeedData (read more from the socket), or Paused (the peer is done for now; see Keep-alive). |
| conn.send(event) | Validate event against the state machine and return the bytes to write (None for ConnectionClosed). |
| Event | Meaning |
|---|---|
| Request(Request) | Start of a request: method_, target, headers, http_version |
| InformationalResponse(InformationalResponse) | A 1xx response |
| Response(Response) | Start of a final response: status_code, headers, http_version, reason |
| Data(Data) | A piece of a message body |
| EndOfMessage(EndOfMessage) | End of a message body, with optional trailers |
| ConnectionClosed | The peer closed their side of the connection |
///|
/// Handle one connection whose incoming bytes arrive as `chunks`, returning
/// everything the server wrote.
fn serve(chunks : Array[Bytes]) -> Bytes raise {
let conn = @h11.Connection::new(Server)
let out = @buffer.Buffer()
let mut next_chunk = 0
for ;; {
match conn.next_event() {
NeedData =>
// read from the socket; b"" means EOF
if next_chunk < chunks.length() {
conn.receive_data(chunks[next_chunk])
next_chunk 1
} else {
conn.receive_data(b"")
}
Event(Request(req)) => {
let body = b"you asked for " + req.target
let length = @utf8.encode(body.length().to_string())
let resp = @h11.Response::new(status_code=200, headers=[
(b"Content-Length", length),
])
out.write_bytes(conn.send(Response(resp)).unwrap())
out.write_bytes(conn.send(Data(@h11.Data::new(body))).unwrap())
out.write_bytes(
conn.send(EndOfMessage(@h11.EndOfMessage::new())).unwrap(),
)
}
Event(ConnectionClosed) => break
Event(_) => () // request body chunks, end of request, ...
Paused =>
// Both sides finished one request/response cycle.
if conn.our_state() is Done && conn.their_state() is Done {
conn.start_next_cycle()
} else {
break // e.g. MUST_CLOSE: we should close the socket
}
}
}
out.to_bytes()
}
///|
test "server loop with a pipelined, fragmented request stream" {
let replies = serve([
b"GET /a HTTP/1.1\r\nHost: x\r\n\r\nGET /b HT", b"TP/1.1\r\nHost: x\r\n\r\n",
])
inspect(
@utf8.decode(replies),
content="HTTP/1.1 200 \r\nContent-Length: 16\r\n\r\nyou asked for /aHTTP/1.1 200 \r\nContent-Length: 16\r\n\r\nyou asked for /b",
)
}///|
test "HTTP/1.0 peers get close-delimited bodies" {
let server = @h11.Connection::new(Server)
server.receive_data(b"GET / HTTP/1.0\r\n\r\n")
guard server.next_event() is Event(Request(_)) else { fail("no request") }
guard server.next_event() is Event(EndOfMessage(_)) else { fail("no EOM") }
let wire = server.send(
Response(@h11.Response::new(status_code=200, headers=[])),
)
inspect(
@utf8.decode(wire.unwrap()),
content="HTTP/1.1 200 \r\nConnection: close\r\n\r\n",
)
assert_eq(server.our_state(), SendBody)
}///|
test "headers keep their casing on the wire" {
let headers = @h11.Headers::new([
(b"Content-Type", b"text/html"),
(b"X-Custom", b"1"),
])
assert_eq(headers[0], (b"content-type", b"text/html"))
assert_eq(headers.raw_items()[1], (b"X-Custom", b"1"))
// Comma-separated headers can be read case-insensitively:
let h = @h11.Headers::new([(b"Connection", b"Keep-Alive, Upgrade")])
assert_eq(@h11.get_comma_header(h, b"connection"), [b"keep-alive", b"upgrade"])
}///|
test "100-continue" {
let server = @h11.Connection::new(Server)
server.receive_data(
b"POST /upload HTTP/1.1\r\nHost: x\r\nContent-Length: 4\r\nExpect: 100-continue\r\n\r\n",
)
guard server.next_event() is Event(Request(_)) else { fail("no request") }
assert_true(server.they_are_waiting_for_100_continue())
let wire = server.send(
InformationalResponse(
@h11.InformationalResponse::new(status_code=100, headers=[]),
),
)
inspect(@utf8.decode(wire.unwrap()), content="HTTP/1.1 100 \r\n\r\n")
assert_true(!server.they_are_waiting_for_100_continue())
}///|
test "upgrading to another protocol" {
let server = @h11.Connection::new(Server)
server.receive_data(
b"GET /chat HTTP/1.1\r\nHost: x\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\r\n\x81\x05hello",
)
guard server.next_event() is Event(Request(_)) else { fail("no request") }
guard server.next_event() is Event(EndOfMessage(_)) else { fail("no EOM") }
assert_eq(server.next_event(), Paused)
assert_eq(server.their_state(), MightSwitchProtocol)
let accept = @h11.InformationalResponse::new(status_code=101, headers=[
(b"Upgrade", b"websocket"),
(b"Connection", b"Upgrade"),
])
ignore(server.send(InformationalResponse(accept)))
assert_eq(server.states(), {
client: SwitchedProtocol,
server: SwitchedProtocol,
})
// The bytes that followed the handshake belong to the new protocol.
assert_eq(server.trailing_data(), (b"\x81\x05hello", false))
}///|
test "errors carry a suggested status code" {
let server = @h11.Connection::new(Server)
server.receive_data(
b"GET / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: gzip\r\n\r\n",
)
try server.next_event() catch {
RemoteProtocolError(msg, error_status_hint~) => {
inspect(msg, content="Only Transfer-Encoding: chunked is supported")
assert_eq(error_status_hint, 501)
}
e => fail("unexpected error \{e}")
} noraise {
_ => fail("expected an error")
}
assert_eq(server.their_state(), Error)
// We can still tell the client what went wrong:
let wire = server.send(
Response(@h11.Response::new(status_code=501, headers=[])),
)
inspect(
@utf8.decode(wire.unwrap()),
content="HTTP/1.1 501 \r\nConnection: close\r\n\r\n",
)
}| Python | MoonBit |
|---|---|
| h11.Connection(our_role=h11.CLIENT) | @h11.Connection::new(Client) |
| conn.next_event() → event, h11.NEED_DATA, h11.PAUSED | conn.next_event() → Event(e), NeedData, Paused |
| conn.send(event) → bytes or None | conn.send(event) → Some(bytes) or None |
| conn.states, conn.our_state, conn.their_state | conn.states(), conn.our_state(), conn.their_state() |
| conn.their_http_version, conn.trailing_data | conn.their_http_version(), conn.trailing_data() |
| h11.Request(method=..., target=..., headers=[...]) | @h11.Request::new(method_=..., target=..., headers=[...]) |
| h11.Data(data=b"..."), h11.EndOfMessage() | @h11.Data::new(b"..."), @h11.EndOfMessage::new() |
| h11.ConnectionClosed() | ConnectionClosed |
| h11.CLIENT, h11.SEND_BODY, h11.MUST_CLOSE, ... | Client, SendBody, MustClose, ... |
| except h11.RemoteProtocolError as e: e.error_status_hint | catch { RemoteProtocolError(msg, error_status_hint~) => ... } |
moon testpython3 fuzz/difftest.py --cases 100000 --seed 1pub(all) suberror ProtocolError {
LocalProtocolError(String, error_status_hint~ : Int)
RemoteProtocolError(String, error_status_hint~ : Int)
}impl Show for ProtocolErrorpub(all) suberror RuntimeError {
RuntimeError(String)
}impl Show for RuntimeErrorpub struct Connection {
// private fields
}fn Connection::send_with_data_passthrough(self : Connection, event : Event) -> Array[Bytes]? raise ProtocolErrorimpl Debug for EndOfMessagepub(all) enum Event {
Request(Request)
InformationalResponse(InformationalResponse)
Response(Response)
Data(Data)
EndOfMessage(EndOfMessage)
ConnectionClosed
} derive(Eq, Debug)pub struct Headers {
// private fields
}impl Debug for InformationalResponsefn InformationalResponse::new(status_code~ : Int, headers~ : ArrayView[(Bytes, Bytes)], http_version? : Bytes, reason? : Bytes) -> InformationalResponse raise ProtocolErrorfn Request::new(method_~ : Bytes, target~ : Bytes, headers~ : ArrayView[(Bytes, Bytes)], http_version? : Bytes) -> Request raise ProtocolErrortest {
let req = @h11.Request::new(method_=b"GET", target=b"/", headers=[
(b"Host", b"example.com"),
])
inspect(req.headers.length(), content="1")
}fn Response::new(status_code~ : Int, headers~ : ArrayView[(Bytes, Bytes)], http_version? : Bytes, reason? : Bytes) -> Response raise ProtocolErrortest {
let resp = @h11.Response::new(status_code=200, headers=[], reason=b"OK")
inspect(resp.status_code, content="200")
}let DEFAULT_MAX_INCOMPLETE_EVENT_SIZE : Intfn set_comma_header(headers : Headers, name : Bytes, new_values : ArrayView[Bytes]) -> Headers raise ProtocolErrorInstall
Download zipA pure, bring-your-own-I/O implementation of HTTP/1.1 (port of python-hyper/h11)
Dependencies