// Copyright 2026 Leo Cheng
// SPDX-License-Identifier: Apache-2.0
// A QUIC connection's numbering and handshake-input core (RFC 9000 §12.3, RFC 9001
// §4.1.3). A connection keeps three independent packet-number spaces — Initial,
// Handshake, and Application (1-RTT) — and a CRYPTO byte stream per encryption level,
// each ordered on its own. This skeleton routes a received packet's number to its space
// (for acknowledgement) and a received CRYPTO frame to its level's reassembler (yielding
// the contiguous TLS handshake bytes to feed the handshake), and hands out send-side
// packet numbers per space. It holds no keys or socket — the key schedule and the async
// UDP transport wrap this pure core.
///|
/// A QUIC encryption level / packet-number space (RFC 9001 §4.1.1). 0-RTT shares the
/// Application space, so the three long-lived spaces are Initial, Handshake, and
/// Application.
pub(all) enum Level {
Initial
Handshake
Application
} derive(Eq, Debug)
///|
/// A QUIC connection's per-space numbering/ACK state and per-level CRYPTO streams.
pub struct Conn {
is_server : Bool
initial_space : @packet.Space
handshake_space : @packet.Space
application_space : @packet.Space
initial_crypto : @stream.Assembler
handshake_crypto : @stream.Assembler
application_crypto : @stream.Assembler
}
///|
/// A fresh connection (client or server) with empty spaces and CRYPTO streams.
pub fn Conn::new(is_server : Bool) -> Conn {
{
is_server,
initial_space: @packet.Space::new(),
handshake_space: @packet.Space::new(),
application_space: @packet.Space::new(),
initial_crypto: @stream.Assembler::new(),
handshake_crypto: @stream.Assembler::new(),
application_crypto: @stream.Assembler::new(),
}
}
///|
/// Whether this endpoint is the server.
pub fn Conn::is_server(self : Conn) -> Bool {
self.is_server
}
///|
/// The packet-number space for `level`.
pub fn Conn::space(self : Conn, level : Level) -> @packet.Space {
match level {
Initial => self.initial_space
Handshake => self.handshake_space
Application => self.application_space
}
}
///|
/// The CRYPTO reassembler for `level`.
pub fn Conn::crypto_stream(self : Conn, level : Level) -> @stream.Assembler {
match level {
Initial => self.initial_crypto
Handshake => self.handshake_crypto
Application => self.application_crypto
}
}
///|
/// Allocate the next packet number to send at `level`.
pub fn Conn::next_packet_number(self : Conn, level : Level) -> Int64 {
self.space(level).next_packet_number()
}
///|
/// Record that a packet numbered `pn` arrived at `level`; `ack_eliciting` marks one that
/// must be acknowledged.
pub fn Conn::on_packet_received(
self : Conn,
level : Level,
pn : Int64,
ack_eliciting : Bool,
) -> Unit {
self.space(level).on_packet_received(pn, ack_eliciting)
}
///|
/// Record the peer's acknowledgement of up to `largest` at `level`.
pub fn Conn::on_ack_received(
self : Conn,
level : Level,
largest : Int64,
) -> Unit {
self.space(level).on_ack_received(largest)
}
///|
/// Build the ACK frame owed at `level`, or `None` if nothing has been received there.
pub fn Conn::build_ack(
self : Conn,
level : Level,
ack_delay : UInt64,
) -> @frame.Frame? {
self.space(level).build_ack(ack_delay)
}
///|
/// Feed a CRYPTO frame's `data` at `offset` and `level` into that level's handshake
/// stream, returning the newly contiguous TLS handshake bytes now readable (empty while
/// an earlier gap is still outstanding). CRYPTO streams carry no fin.
pub fn Conn::on_crypto_frame(
self : Conn,
level : Level,
offset : UInt64,
data : Bytes,
) -> Bytes raise @stream.Refused {
let stream = self.crypto_stream(level)
stream.insert(offset, data, fin=false)
stream.read()
}
///|
pub extend Level with Debug::{to_repr}
///|
pub extend Level with Eq::{not_equal, equal}
// The QUIC send loop, assembled (RFC 9000 §13, RFC 9002): the deterministic core a connection
// drives to put stream data on the wire and recover from loss. It composes the stream scheduler
// (which stream sends how much, bounded by flow control), the loss-recovery state (RTT,
// congestion window, loss detection), and a record of which STREAM data each packet carried.
// `poll_send` yields the next packet — a retransmission first, else freshly scheduled data —
// charging the congestion window; `on_ack` clears acknowledged packets and re-queues the stream
// data of any the ACK reveals as lost. The async UDP transport is the thin glue over this: send
// what `poll_send` returns, feed received ACKs to `on_ack`.
///|
/// A connection's send state for one packet-number space.
pub struct Sender {
sched : @stream.Scheduler
recovery : @recovery.State
data : Map[UInt64, Bytes]
in_packet : Map[Int64, Array[@stream.Chunk]]
retransmit : Array[@stream.Chunk]
mut next_pn : Int64
}
///|
/// A fresh sender: `initial_max_data`/`initial_max_stream_data` are the peer's advertised flow
/// limits, `max_datagram_size` sizes the congestion window, `max_ack_delay`/`granularity` tune
/// the timers (microseconds), and `max_frame` caps a STREAM frame's payload.
pub fn Sender::new(
initial_max_data : UInt64,
initial_max_stream_data : UInt64,
max_datagram_size : Int64,
max_ack_delay : Int64,
granularity : Int64,
max_frame : UInt64,
) -> Sender {
{
sched: @stream.Scheduler::new(
@stream.Budget::new(initial_max_data, initial_max_stream_data),
frame=max_frame,
),
recovery: @recovery.State::new(
policy=@recovery.Policy::new(
datagram=max_datagram_size,
ack_delay=max_ack_delay,
granularity~,
),
),
data: Map([]),
in_packet: Map([]),
retransmit: [],
next_pn: 0L,
}
}
///|
/// Skip the send counter past `pn`, so a packet this sender did not build — a driver's
/// ACK-only packet in the same space — does not collide with a number it later assigns. A
/// `pn` the sender is already past leaves it alone.
pub fn Sender::advance_past(self : Sender, pn : Int64) -> Unit {
if pn >= self.next_pn {
self.next_pn = pn + 1L
}
}
///|
/// Queue `bytes` of application data to send on stream `id`.
pub fn Sender::queue_stream(self : Sender, id : UInt64, bytes : Bytes) -> Unit {
self.data[id] = bytes
self.sched.queue(id, bytes.length().to_uint64())
}
///|
/// Mark stream `id` finished.
pub fn Sender::queue_fin(self : Sender, id : UInt64) -> Unit {
self.sched.queue_fin(id)
}
///|
/// Raise the connection-wide send limit from a received MAX_DATA frame.
pub fn Sender::on_max_data(self : Sender, new_max : UInt64) -> Unit {
self.sched.on_max_data(new_max)
}
///|
fn Sender::materialize(self : Sender, send : @stream.Chunk) -> @frame.Frame {
let full = match self.data.get(send.stream) {
Some(b) => b
None => b""
}
let off = send.offset.to_int()
let len = send.length.to_int()
Stream(
id=send.stream,
offset=send.offset,
fin=send.fin,
data=full[off:off + len].to_owned(),
)
}
///|
/// The next packet to send at `now`, or `None` when the congestion window is closed or nothing
/// is queued (RFC 9002 §7). Retransmissions go first, then freshly scheduled stream data; the
/// packet is recorded in the recovery state and its stream sends remembered for retransmission.
pub fn Sender::poll_send(
self : Sender,
now : Int64,
) -> (Int64, Array[@frame.Frame])? raise {
if self.recovery.can_send(1L) {
let send = match self.retransmit.pop() {
Some(x) => Some(x)
None => self.sched.next()
}
match send {
None => None
Some(s) => {
let pn = self.next_pn
self.next_pn = self.next_pn + 1L
let frame = self.materialize(s)
self.recovery.on_sent(pn, at=now, size=s.length.reinterpret_as_int64())
self.in_packet[pn] = [s]
Some((pn, [frame]))
}
}
} else {
None
}
}
///|
/// Process a received ACK `frame` at `now`: clear acknowledged packets from the recovery state
/// and re-queue the stream data of any packet the ACK reveals as lost, to be retransmitted by a
/// later `poll_send` (RFC 9002 §6). Raises on a malformed ACK.
pub fn Sender::on_ack(
self : Sender,
frame : @frame.Frame,
now : Int64,
) -> Unit raise {
let lost = self.recovery.on_ack(frame, now, ack_delay=0L)
for pn in lost {
match self.in_packet.get(pn) {
Some(sends) =>
for s in sends {
self.retransmit.push(s)
}
None => ()
}
}
}
///|
/// Handle the probe timeout firing at `now` (RFC 9002 §6.2.4). If the PTO deadline has passed with
/// packets still outstanding, re-queue the oldest outstanding packet's stream data as a probe — a
/// later `poll_send` puts it back on the wire — and back off the timer for the next arming. Unlike
/// loss detection, this neither declares packets lost nor reduces the congestion window. Returns
/// whether a probe was armed.
pub fn Sender::on_pto_timeout(self : Sender, now : Int64) -> Bool {
guard self.recovery.pto_deadline() is Some(deadline) else { return false }
guard now >= deadline else { return false }
let out = self.recovery.outstanding()
guard out.length() > 0 else { return false }
let mut oldest = out[0]
for pn in out {
if pn < oldest {
oldest = pn
}
}
match self.in_packet.get(oldest) {
Some(sends) =>
for s in sends {
self.retransmit.push(s)
}
None => ()
}
self.recovery.on_pto()
true
}
///|
/// The packet numbers still outstanding (sent, not yet acknowledged or declared lost).
pub fn Sender::outstanding(self : Sender) -> Array[Int64] {
self.recovery.outstanding()
}
///|
/// The current congestion window (bytes).
pub fn Sender::window(self : Sender) -> Int64 {
self.recovery.window()
}
// A QUIC server connection's Initial-flight processing (RFC 9000 §7, RFC 9001 §5): it
// ties the Initial packet protection, the per-level CRYPTO reassembly, and the TLS 1.3
// handshake runner together. Receiving an Initial packet unprotects it with the keys
// derived from the client's connection id, records the packet number for acknowledgement,
// reassembles the CRYPTO stream, and drives the handshake with the contiguous handshake
// bytes. This is the pure connection core the async UDP transport wraps; feeding it a
// client's Initial packet advances the server handshake to RECVD_CH.
///|
/// A QUIC server connection mid-handshake: the Initial keys, the packet-number/CRYPTO
/// bookkeeping, and the TLS handshake it is driving.
pub struct Server {
initial_keys : @crypto.Keys
handshake : @hs.Handshake
conn : Conn
}
///|
/// A server connection for a client whose Destination Connection ID is `dcid` (the
/// Initial secret and keys derive from it, RFC 9001 §5.2).
pub fn Server::new(dcid : Bytes) -> Server {
{
initial_keys: @crypto.initial(dcid[:]).0,
handshake: @hs.Handshake::new(),
conn: Conn::new(true),
}
}
///|
/// The current TLS handshake state.
pub fn Server::handshake_state(self : Server) -> @hs.Server {
self.handshake.state()
}
///|
/// The transcript hash over the handshake messages processed so far.
pub fn Server::transcript_hash(self : Server) -> Bytes {
self.handshake.transcript()
}
///|
/// The raw ClientHello the connection received (empty before one arrives): the source a
/// server runs the ECDHE key_share from to derive its handshake secrets.
pub fn Server::client_hello(self : Server) -> Bytes {
self.handshake.hello()
}
///|
/// Negotiate the ClientHello this connection received (RFC 8446 §4.1.1): the group and key
/// share to run the ECDHE with, and the signature scheme to sign the CertificateVerify under,
/// read off what the client offered. Raises the §6 alert for a ClientHello this build cannot
/// serve, and returns `Retry` when the answer is a HelloRetryRequest.
pub fn Server::negotiate(self : Server) -> @hs.Choice raise @alert.Alert {
self.handshake.negotiate()
}
///|
/// Whether an ACK is owed in the Initial space.
pub fn Server::initial_ack_pending(self : Server) -> Bool {
self.conn.space(Level::Initial).ack_pending()
}
///|
/// Build the server's ServerHello response: encode it, fold it into the transcript, and
/// protect it into an Initial packet addressed to `out_dcid` (the client's source
/// connection id) from `out_scid`, at Initial packet number `pn` (RFC 9001 §5.3). The
/// state stays at RECVD_CH — the server flight is not complete until the Handshake-space
/// messages (EncryptedExtensions, Certificate, CertificateVerify, Finished) are sent by
/// `send_handshake_flight`, which advances it to WAIT_FINISHED.
pub fn Server::send_server_hello(
self : Server,
server_hello : @msg.Server,
out_dcid : Bytes,
out_scid : Bytes,
pn : Int64,
) -> Bytes {
let sh_msg = @msg.server_hello(server_hello)
self.handshake.sent(sh_msg[:])
send_long(
Initial,
1U,
out_dcid,
out_scid,
pn,
4,
[@frame.Crypto(offset=0, data=sh_msg)],
self.initial_keys,
)
}
///|
/// Send the server's Handshake-space flight (RFC 8446 §4, RFC 9001 §5): encode
/// EncryptedExtensions, Certificate, and CertificateVerify, fold them into the transcript
/// in order, compute the server Finished as `HMAC(finished_key, transcript hash through
/// CertificateVerify)` over `server_hs_secret`, fold the Finished too, and advance the
/// state past the whole flight (RECVD_CH → WAIT_FINISHED). The four messages form one
/// contiguous CRYPTO stream, protected into a Handshake packet with the handshake-space
/// keys `@crypto.Keys::of(server_hs_secret[:])`, addressed to `out_dcid` from `out_scid` at
/// Handshake packet number `pn`. `server_hs_secret` is the server handshake traffic secret
/// the key schedule derives once ECDHE completes.
pub fn Server::send_handshake_flight(
self : Server,
server_hs_secret : Bytes,
encrypted_extensions : Bytes,
certificate : Bytes,
verify : (Bytes) -> Bytes,
out_dcid : Bytes,
out_scid : Bytes,
pn : Int64,
) -> Bytes raise {
let ee = @msg.frame(EncryptedExtensions, encrypted_extensions[:])
let cert = @msg.frame(Certificate, certificate[:])
self.handshake.sent(ee[:])
self.handshake.sent(cert[:])
// The body is built here rather than passed in, because §4.4.3 signs the transcript
// as it stands at exactly this point: Certificate folded in, Finished not yet.
let cv = @msg.frame(CertificateVerify, verify(self.handshake.transcript())[:])
self.handshake.sent(cv[:])
let verify_data = @keys.finished(
server_hs_secret[:],
self.handshake.transcript()[:],
)
let fin = @msg.frame(Finished, verify_data[:])
self.handshake.sent(fin[:])
self.handshake.sent_flight(client_cert=false)
let crypto = join(join(join(ee, cert), cv), fin)
send_long(
Handshake,
1U,
out_dcid,
out_scid,
pn,
4,
[@frame.Crypto(offset=0, data=crypto)],
@crypto.Keys::of(server_hs_secret[:]),
)
}
///|
/// Drive an already-unprotected packet's `frames` at `level`, numbered `pn`: record the
/// packet for acknowledgement and feed its CRYPTO frames into that level's reassembler,
/// driving the handshake with whatever handshake bytes become contiguous. Returns the
/// handshake message types processed. A driver that unprotected the packet itself — because
/// it also wants the ACK and CONNECTION_CLOSE frames beside the CRYPTO — enters here rather
/// than through `receive_initial`.
pub fn Server::on_frames(
self : Server,
level : Level,
frames : Array[@frame.Frame],
pn : Int64,
) -> Array[@msg.Kind] raise {
self.conn.on_packet_received(level, pn, @packet.payload_elicits_ack(frames))
let processed : Array[@msg.Kind] = []
for frame in frames {
match frame {
Crypto(offset~, data~) => {
let contiguous = self.conn.on_crypto_frame(level, offset, data)
if contiguous.length() > 0 {
for msg_type in self.handshake.feed(contiguous) {
processed.push(msg_type)
}
}
}
_ => ()
}
}
processed
}
///|
/// Receive a client Initial packet: unprotect it, record its number for acknowledgement,
/// reassemble its CRYPTO frames, and drive the handshake with the contiguous handshake
/// bytes. Returns the handshake message types processed. Raises if the packet fails to
/// authenticate.
pub fn Server::receive_initial(
self : Server,
packet : Bytes,
) -> Array[@msg.Kind] raise {
let (frames, pn) = match recv_long(packet, self.initial_keys) {
Some(v) => v
None => raise @packet.Malformed("Initial packet failed to authenticate")
}
self.on_frames(Level::Initial, frames, pn)
}
///|
/// Receive the client's Handshake-space flight — in this no-client-authentication path,
/// the client Finished (RFC 8446 §4.4.4): unprotect the Handshake packet with the client
/// handshake keys, reassemble its CRYPTO stream, verify the Finished `verify_data` against
/// `client_hs_secret` over the transcript through the server Finished, then fold it in to
/// advance the handshake to CONNECTED. Returns whether the handshake is now connected with
/// a Finished that authenticated. Raises if the packet fails to authenticate.
pub fn Server::receive_handshake(
self : Server,
packet : Bytes,
client_hs_keys : @crypto.Keys,
client_hs_secret : Bytes,
) -> Bool raise {
let (frames, pn) = match recv_long(packet, client_hs_keys) {
Some(v) => v
None => raise @packet.Malformed("Handshake packet failed to authenticate")
}
self.on_handshake_frames(frames, pn, client_hs_secret)
}
///|
/// Drive the client's Handshake-space `frames`, numbered `pn`, already unprotected: record
/// the packet, reassemble its CRYPTO stream, verify any Finished in it against
/// `client_hs_secret` over the transcript as it stands, then fold the messages in. Returns
/// whether the handshake is now connected with a Finished that authenticated — the frame-level
/// half of `receive_handshake`, for a driver holding the decrypted packet.
pub fn Server::on_handshake_frames(
self : Server,
frames : Array[@frame.Frame],
pn : Int64,
client_hs_secret : Bytes,
) -> Bool raise {
self.conn.on_packet_received(
Level::Handshake,
pn,
@packet.payload_elicits_ack(frames),
)
let mut finished_ok = false
for frame in frames {
match frame {
Crypto(offset~, data~) => {
let contiguous = self.conn.on_crypto_frame(
Level::Handshake,
offset,
data,
)
if contiguous.length() > 0 {
// The Finished MACs the transcript as it stands now (through the server Finished),
// so verify against the current hash before folding the client Finished in.
let th = self.handshake.transcript()
let view = contiguous[:]
let mut off = 0
for ;; {
match @msg.unframe(view[off:]) {
Some((t, body)) => {
if t == @msg.Kind::Finished {
finished_ok = @keys.finished_ok(
client_hs_secret[:],
th[:],
body[:],
)
}
off = off + 4 + body.length()
}
None => break
}
}
let _ = self.handshake.feed(contiguous)
}
}
_ => ()
}
}
self.handshake.is_connected() && finished_ok
}
// A QUIC server's connection table and event loop (RFC 9000 §5.2, §8.1, §10, §14.1). Every
// piece below this file is a pure function — the packet protection, the frame codec, the
// handshake, the recovery loop — and none of them talk to each other. This is the driver that
// makes them a server: it demultiplexes a received datagram onto a connection by its
// Destination Connection ID, drives that connection's state machine at the encryption level
// its installed keys allow, fires the idle, probe and closing timers, and hands back the
// datagrams to put on the wire — padded to the Initial minimum, and withheld once an
// unvalidated peer's amplification budget is spent. It holds no socket and no clock: the peer
// address is a type parameter and every entry point takes `now`, so the tests drive it with a
// fake sink on an injected clock and `quic_serve.native.mbt` drives it with a UDP socket.
//
// Version negotiation, stateless reset, connection migration, path validation and coalesced
// packets (more than one QUIC packet in a datagram) are not implemented here.
///|
/// Where a connection is in its lifecycle (RFC 9000 §10.2): exchanging packets, closing —
/// this endpoint sent a CONNECTION_CLOSE and answers anything further with it — or draining,
/// where the peer closed and nothing more goes out.
pub(all) enum Phase {
Active
Closing
Draining
} derive(Eq, Debug)
///|
/// The smallest a datagram carrying an ack-eliciting Initial packet may be (RFC 9000 §14.1).
pub let min_datagram : Int = 1200
///|
/// The anti-amplification factor (RFC 9000 §8.1): until it has validated a peer's address, a
/// server may send it at most this many times the bytes it has received from it.
pub let amplification : Int64 = 3L
///|
/// What a server gives each new connection: the QUIC version its long headers carry, the idle
/// timeout it advertises, how long a closing or draining connection lingers, and the datagram
/// size, timer tuning and flow-control credit the send loop starts from. Every duration is in
/// microseconds, the unit the recovery code measures in.
pub(all) struct Config {
version : UInt
idle_timeout : Int64
close_period : Int64
max_datagram : Int64
max_ack_delay : Int64
granularity : Int64
initial_max_data : UInt64
initial_max_stream_data : UInt64
max_frame : UInt64
}
///|
/// QUIC version 1, a 30-second idle timeout, a 3-second closing period, 1200-byte datagrams,
/// a 25-millisecond ACK delay, and a mebibyte of flow-control credit on the connection and on
/// each stream.
pub fn Config::default() -> Config {
{
version: 1U,
idle_timeout: 30_000_000L,
close_period: 3_000_000L,
max_datagram: 1200L,
max_ack_delay: 25_000L,
granularity: 1_000L,
initial_max_data: 1_048_576UL,
initial_max_stream_data: 1_048_576UL,
max_frame: 1024UL,
}
}
///|
/// The idle timeout in force on a connection (RFC 9000 §10.1): the smaller of the two
/// endpoints' advertised `max_idle_timeout`s, where zero on either side means that end asks
/// for no limit at all.
pub fn idle_timeout(ours : Int64, theirs : Int64) -> Int64 {
if ours <= 0L {
theirs
} else if theirs <= 0L {
ours
} else if ours < theirs {
ours
} else {
theirs
}
}
///|
/// Build a datagram from `payload`, padding the payload with PADDING frames until the built
/// datagram reaches `min` bytes (RFC 9000 §14.1). The header's Length field is a varint, so
/// growing the payload can grow the header too; re-measure rather than guess the overhead.
fn pad(payload : Bytes, build : (Bytes) -> Bytes, min : Int) -> Bytes {
let mut body = payload
let mut datagram = build(body)
while datagram.length() < min {
body = @packet.pad(body, body.length() + min - datagram.length())
datagram = build(body)
}
datagram
}
///|
/// One server-side connection: the peer it answers, the connection ids the two ends address
/// each other by, the handshake and send state machines it drives, and the counters the
/// amplification limit, the idle timer and the closing period are read from. `A` is the
/// peer-address type the enclosing `Endpoint` was built with.
pub struct Session[A] {
peer : A
cid : Bytes
peer_cid : Bytes
core : Server
sender : Sender
cfg : Config
out : Array[Bytes]
received : Array[(Level, @frame.Frame)]
mut rx : Int64
mut tx : Int64
mut packets : Int
mut validated : Bool
mut last_rx : Int64
mut idle_timeout : Int64
mut phase : Phase
mut close_at : Int64
mut close_frame : @frame.Frame?
mut close_level : Level
mut hs_rx : @crypto.Keys?
mut hs_rx_secret : Bytes
mut hs_tx : @crypto.Keys?
mut app_rx : @crypto.Keys?
mut app_tx : @crypto.Keys?
}
///|
/// A connection for the peer at `peer` that addressed this server as `cid` and called itself
/// `peer_cid`. The server keeps the client's original Destination Connection ID as its own —
/// that is the id the Initial keys derive from (RFC 9001 §5.2), so reusing it keeps the
/// demultiplexing key and the key schedule talking about the same bytes.
fn[A] Session::new(
peer : A,
cid : Bytes,
peer_cid : Bytes,
cfg : Config,
now : Int64,
) -> Session[A] {
{
peer,
cid,
peer_cid,
core: Server::new(cid),
sender: Sender::new(
cfg.initial_max_data,
cfg.initial_max_stream_data,
cfg.max_datagram,
cfg.max_ack_delay,
cfg.granularity,
cfg.max_frame,
),
cfg,
out: [],
received: [],
rx: 0L,
tx: 0L,
packets: 0,
validated: false,
last_rx: now,
idle_timeout: cfg.idle_timeout,
phase: Active,
close_at: 0L,
close_frame: None,
close_level: Level::Initial,
hs_rx: None,
hs_rx_secret: b"",
hs_tx: None,
app_rx: None,
app_tx: None,
}
}
///|
/// The peer this connection answers.
pub fn[A] Session::peer(self : Session[A]) -> A {
self.peer
}
///|
/// The connection id the peer addresses this server by — the table's key.
pub fn[A] Session::cid(self : Session[A]) -> Bytes {
self.cid
}
///|
/// Where this connection is in its lifecycle.
pub fn[A] Session::phase(self : Session[A]) -> Phase {
self.phase
}
///|
/// The handshake state machine underneath, for a caller driving the TLS flights.
pub fn[A] Session::core(self : Session[A]) -> Server {
self.core
}
///|
/// Bytes received from the peer, the numerator of the amplification limit.
pub fn[A] Session::rx_bytes(self : Session[A]) -> Int64 {
self.rx
}
///|
/// Bytes sent to the peer, the quantity the amplification limit caps.
pub fn[A] Session::tx_bytes(self : Session[A]) -> Int64 {
self.tx
}
///|
/// Datagrams received on this connection.
pub fn[A] Session::packets_received(self : Session[A]) -> Int {
self.packets
}
///|
/// Whether the peer's address has been validated (RFC 9000 §8.1), which lifts the
/// amplification limit.
pub fn[A] Session::validated(self : Session[A]) -> Bool {
self.validated
}
///|
/// The idle timeout in force, in microseconds; zero means none.
pub fn[A] Session::idle_timeout(self : Session[A]) -> Int64 {
self.idle_timeout
}
///|
/// Datagrams built and waiting to go out — non-zero while the amplification limit or a closed
/// congestion window is holding them back.
pub fn[A] Session::pending(self : Session[A]) -> Int {
self.out.length()
}
///|
/// The largest packet number the peer has acknowledged at `level`, or `None` before its first
/// ACK there.
pub fn[A] Session::largest_acked(self : Session[A], level : Level) -> Int64? {
self.core.conn.space(level).largest_acked()
}
///|
/// Whether `n` more bytes may go to this peer: unrestricted once its address is validated,
/// otherwise at most three times what it has sent this server (RFC 9000 §8.1).
pub fn[A] Session::can_send(self : Session[A], n : Int) -> Bool {
self.validated || self.tx + n.to_int64() <= amplification * self.rx
}
///|
/// Install the Handshake-space keys, derived from the client's and this server's handshake
/// traffic secrets (RFC 9001 §5.1). Until they are installed, a Handshake packet cannot be
/// unprotected and is dropped; after, the connection reads and writes at that level. The
/// client's secret is kept so an arriving Finished can be verified against it.
pub fn[A] Session::install_handshake_keys(
self : Session[A],
client_secret : Bytes,
server_secret : Bytes,
) -> Unit {
self.hs_rx = Some(@crypto.Keys::of(client_secret[:]))
self.hs_rx_secret = client_secret
self.hs_tx = Some(@crypto.Keys::of(server_secret[:]))
}
///|
/// Install the Application-space (1-RTT) keys, derived from the two application traffic
/// secrets. Until they are installed, a short-header packet is dropped and the send loop has
/// nowhere to put stream data.
pub fn[A] Session::install_app_keys(
self : Session[A],
client_secret : Bytes,
server_secret : Bytes,
) -> Unit {
self.app_rx = Some(@crypto.Keys::of(client_secret[:]))
self.app_tx = Some(@crypto.Keys::of(server_secret[:]))
}
///|
/// Apply the peer's transport parameters (RFC 9000 §18): the idle timer settles at the
/// smaller of the two ends' values, and the peer's initial flow-control credit opens the
/// connection-wide send window, which §18.2 says is a MAX_DATA arriving at once.
pub fn[A] Session::apply_peer_params(
self : Session[A],
params : Array[Param],
) -> Unit {
let peer = Limits::read(params[:])
// The idle timer settles at the smaller of the two ends' values; the parameter is in
// milliseconds and the timers here are in microseconds.
if find(params[:], MaxIdleTimeout) is Some(_) {
self.idle_timeout = idle_timeout(
self.cfg.idle_timeout,
peer.idle_timeout.reinterpret_as_int64() * 1000L,
)
}
// §18.2: the initial flow-control parameters are a MAX_DATA and a MAX_STREAM_DATA on
// every stream, delivered the moment the handshake carries them.
self.sender.on_max_data(peer.data)
}
///|
/// The highest encryption level this connection can send at — Application once its 1-RTT keys
/// are in, else Handshake, else Initial, which is always available.
pub fn[A] Session::send_level(self : Session[A]) -> Level {
if self.app_tx is Some(_) {
Level::Application
} else if self.hs_tx is Some(_) {
Level::Handshake
} else {
Level::Initial
}
}
///|
/// Queue `frames` to go out at `level` on the next `poll_out`, numbered from that level's
/// packet-number space. An Initial packet carrying anything ack-eliciting takes its datagram
/// to the 1200-byte floor (RFC 9000 §14.1); a level whose keys are not installed drops the
/// frames, since there is nothing to protect them with.
pub fn[A] Session::send(
self : Session[A],
level : Level,
frames : Array[@frame.Frame],
) -> Unit {
let pn = self.core.conn.next_packet_number(level)
if level == Level::Application {
// The recovery loop numbers this space too; keep the two counters from colliding.
self.sender.advance_past(pn)
}
self.emit(level, frames, pn)
}
///|
fn[A] Session::emit(
self : Session[A],
level : Level,
frames : Array[@frame.Frame],
pn : Int64,
) -> Unit {
let keys = match level {
Initial => Some(self.core.initial_keys)
Handshake => self.hs_tx
Application => self.app_tx
}
guard keys is Some(k) else { return }
let payload = @packet.payload(frames)
let datagram = match level {
Initial => {
let build = p => {
seal_long(
Initial,
self.cfg.version,
self.peer_cid,
self.cid,
pn,
4,
p,
k,
)
}
if @packet.payload_elicits_ack(frames) {
pad(payload, build, min_datagram)
} else {
build(payload)
}
}
Handshake =>
send_long(
Handshake,
self.cfg.version,
self.peer_cid,
self.cid,
pn,
4,
frames,
k,
)
Application => send_short(self.peer_cid, pn, 4, frames, k)
}
self.out.push(datagram)
}
///|
/// Queue `bytes` to send on stream `id`; the send loop frames and paces them out through
/// `poll_out` once the Application keys are installed.
pub fn[A] Session::queue_stream(
self : Session[A],
id : UInt64,
bytes : Bytes,
) -> Unit {
self.sender.queue_stream(id, bytes)
}
///|
/// Take the frames received since the last call, each with the level it arrived at — where an
/// application layer reads its streams out of.
pub fn[A] Session::take_received(
self : Session[A],
) -> Array[(Level, @frame.Frame)] {
let out = self.received.copy()
self.received.clear()
out
}
///|
/// Close the connection (RFC 9000 §10.2): put a CONNECTION_CLOSE carrying `error_code` and
/// `reason` on the wire at the highest level with keys, drop whatever else was queued, and
/// enter the closing period — during which every packet the peer sends is answered with the
/// same CONNECTION_CLOSE, and at the end of which `tick` forgets the connection. Closing a
/// connection that is already closing or draining does nothing.
pub fn[A] Session::close(
self : Session[A],
error_code : UInt64,
reason : Bytes,
now : Int64,
) -> Unit {
guard self.phase == Active else { return }
// frame_type 0 names no triggering frame, which is what a close not provoked by one says.
let frame = @frame.ConnectionClose(error_code~, frame_type=Some(0UL), reason~)
self.phase = Closing
self.close_frame = Some(frame)
self.close_level = self.send_level()
self.close_at = now + self.cfg.close_period
self.out.clear()
self.send(self.close_level, [frame])
}
///|
/// The next datagram to put on the wire, or `None` when nothing is queued, the amplification
/// budget is spent, or the connection is draining. Pulls a fresh packet out of the send loop
/// when the queue has run dry.
fn[A] Session::next_datagram(self : Session[A], now : Int64) -> Bytes? raise {
guard self.phase != Draining else { return None }
if self.out.length() == 0 {
self.fill(now)
}
guard self.out.length() > 0 else { return None }
let datagram = self.out[0]
guard self.can_send(datagram.length()) else { return None }
let _ = self.out.remove(0)
self.tx = self.tx + datagram.length().to_int64()
Some(datagram)
}
///|
/// Pull the next packet the send loop has ready — a retransmission or freshly scheduled stream
/// data — into the outbound queue (RFC 9002 §7).
fn[A] Session::fill(self : Session[A], now : Int64) -> Unit raise {
guard self.phase == Active else { return }
guard self.app_tx is Some(_) else { return }
match self.sender.poll_send(now) {
Some((pn, frames)) => {
// The recovery loop assigned the number; keep the space's own counter past it.
self.core.conn.space(Level::Application).advance_past(pn)
self.emit(Level::Application, frames, pn)
}
None => ()
}
}
///|
/// Take a datagram from this connection's peer at `now`: count it against the amplification
/// budget and the idle timer, unprotect it at whichever level its header names, and drive the
/// state machine with the frames inside.
fn[A] Session::on_datagram(
self : Session[A],
datagram : Bytes,
now : Int64,
) -> Unit raise {
self.rx = self.rx + datagram.length().to_int64()
self.last_rx = now
self.packets = self.packets + 1
guard self.phase != Draining else { return }
match @packet.read_long(datagram[:]) {
Some((h, _)) =>
match h.packet_type {
@packet.Kind::Initial => self.on_initial(datagram, now)
@packet.Kind::Handshake => self.on_handshake(datagram, now)
// 0-RTT and Retry are not served.
_ => ()
}
None => self.on_app(datagram, now)
}
}
///|
fn[A] Session::on_initial(
self : Session[A],
datagram : Bytes,
now : Int64,
) -> Unit raise {
let opened = recv_long(datagram, self.core.initial_keys) catch { _ => None }
// A datagram that does not authenticate is discarded, not an error: an endpoint cannot
// tell a corrupted packet from an injected one (RFC 9000 §12.2).
guard opened is Some((frames, pn)) else { return }
let _ = self.core.on_frames(Level::Initial, frames, pn)
self.adopt_peer_params()
self.on_frames(Level::Initial, frames, now)
}
///|
fn[A] Session::on_handshake(
self : Session[A],
datagram : Bytes,
now : Int64,
) -> Unit raise {
guard self.hs_rx is Some(keys) else { return }
let opened = recv_long(datagram, keys) catch { _ => None }
guard opened is Some((frames, pn)) else { return }
// A Handshake packet the peer could only have sent from the address it claims proves that
// address, which lifts the amplification limit (RFC 9000 §8.1).
self.validated = true
let _ = self.core.on_handshake_frames(frames, pn, self.hs_rx_secret)
self.on_frames(Level::Handshake, frames, now)
}
///|
fn[A] Session::on_app(
self : Session[A],
datagram : Bytes,
now : Int64,
) -> Unit raise {
guard self.app_rx is Some(keys) else { return }
let opened = recv_short(datagram, self.cid.length(), keys) catch { _ => None }
guard opened is Some((frames, pn)) else { return }
self.core.conn.on_packet_received(
Level::Application,
pn,
@packet.payload_elicits_ack(frames),
)
self.on_frames(Level::Application, frames, now)
}
///|
/// Act on the frames of a packet that arrived at `level`: feed acknowledgements to the
/// recovery loop, take a CONNECTION_CLOSE into the draining period, raise the send limit from
/// a MAX_DATA, and record the rest for the application layer. Then answer — an ACK if one is
/// owed (RFC 9000 §13.2.1), or this endpoint's CONNECTION_CLOSE again if it is closing
/// (§10.2.1).
fn[A] Session::on_frames(
self : Session[A],
level : Level,
frames : Array[@frame.Frame],
now : Int64,
) -> Unit raise {
for frame in frames {
self.received.push((level, frame))
match frame {
Ack(largest~, ..) | AckEcn(largest~, ..) => {
self.core.conn.on_ack_received(level, largest.reinterpret_as_int64())
if level == Level::Application {
self.sender.on_ack(frame, now)
}
}
ConnectionClose(..) => {
self.phase = Draining
self.close_at = now + self.cfg.close_period
self.out.clear()
return
}
MaxData(v) => self.sender.on_max_data(v)
_ => ()
}
}
if self.phase == Closing {
match self.close_frame {
Some(frame) => self.send(self.close_level, [frame])
None => ()
}
return
}
if self.core.conn.space(level).ack_pending() {
match self.core.conn.build_ack(level, 0UL) {
Some(ack) => self.send(level, [ack])
None => ()
}
}
}
///|
/// Read the peer's transport parameters off the ClientHello the handshake has taken in, once
/// there is one to read (RFC 9001 §8.2). A hello that does not decode leaves the defaults
/// standing.
fn[A] Session::adopt_peer_params(self : Session[A]) -> Unit {
let hello = self.core.client_hello()
guard hello.length() > 0 else { return }
guard @msg.unframe(hello[:]) is Some((_, body)) else { return }
guard @msg.read_hello(body[:]) is Some(ch) else { return }
self.apply_peer_params(of_extensions(ch.extensions[:]))
}
///|
/// A QUIC server: the connections it is holding, keyed by the connection id peers address
/// each by, and the settings a new one starts from. `A` is the peer-address type — a socket
/// address under the native event loop, whatever a test finds convenient under a fake sink.
pub struct Endpoint[A] {
conns : Map[Bytes, Session[A]]
cfg : Config
}
///|
/// A server holding no connections, handing `cfg` to each one it opens.
pub fn[A] Endpoint::new(cfg : Config) -> Endpoint[A] {
{ conns: Map([]), cfg, }
}
///|
/// How many connections the server is holding.
pub fn[A] Endpoint::count(self : Endpoint[A]) -> Int {
self.conns.length()
}
///|
/// The connection ids the server is holding.
pub fn[A] Endpoint::conn_ids(self : Endpoint[A]) -> Array[Bytes] {
self.conns.keys().collect()
}
///|
/// The connection `cid` addresses, or `None` if the server is not holding one.
pub fn[A] Endpoint::conn(self : Endpoint[A], cid : Bytes) -> Session[A]? {
self.conns.get(cid)
}
///|
/// Take a datagram received at `now` from `from`: route it to the connection its Destination
/// Connection ID names, opening one when an Initial packet arrives for an id the server does
/// not hold, and drive that connection with it (RFC 9000 §5.2). A datagram naming no
/// connection is dropped — RFC 9000 §10.3's stateless reset is not implemented.
pub fn[A] Endpoint::recv(
self : Endpoint[A],
datagram : Bytes,
from : A,
now : Int64,
) -> Unit raise {
guard datagram.length() > 0 else { return }
match @packet.read_long(datagram[:]) {
Some((h, _)) => {
let conn = match self.conns.get(h.dcid) {
Some(c) => c
None => {
guard h.packet_type is @packet.Kind::Initial else { return }
let c = Session::new(from, h.dcid, h.scid, self.cfg, now)
self.conns[h.dcid] = c
c
}
}
conn.on_datagram(datagram, now)
}
None => {
guard self.lookup_short(datagram) is Some(conn) else { return }
conn.on_datagram(datagram, now)
}
}
}
///|
/// The connection a short-header datagram belongs to. A short header carries no connection-id
/// length (RFC 9000 §17.3), so the only way back to the connection is to test the datagram
/// against the ids this server has issued.
fn[A] Endpoint::lookup_short(
self : Endpoint[A],
datagram : Bytes,
) -> Session[A]? {
for cid in self.conns.keys().collect() {
let n = cid.length()
if datagram.length() >= 1 + n && datagram[1:1 + n].to_owned() == cid {
return self.conns.get(cid)
}
}
None
}
///|
/// The next datagram to put on the wire and the peer to send it to, or `None` when no
/// connection has one ready.
pub fn[A] Endpoint::poll_out(
self : Endpoint[A],
now : Int64,
) -> (Bytes, A)? raise {
for cid in self.conns.keys().collect() {
guard self.conns.get(cid) is Some(conn) else { continue }
match conn.next_datagram(now) {
Some(datagram) => return Some((datagram, conn.peer))
None => ()
}
}
None
}
///|
/// Close the connection `cid` addresses with `error_code` and `reason`, if the server holds
/// one (RFC 9000 §10.2).
pub fn[A] Endpoint::close(
self : Endpoint[A],
cid : Bytes,
error_code : UInt64,
reason : Bytes,
now : Int64,
) -> Unit {
match self.conns.get(cid) {
Some(conn) => conn.close(error_code, reason, now)
None => ()
}
}
///|
/// Run the timers at `now`. An idle connection is forgotten without a CONNECTION_CLOSE, which
/// is what RFC 9000 §10.1 asks for; a closing or draining one is forgotten when its period
/// ends (§10.2); an active one whose probe timeout has come re-queues its oldest outstanding
/// packet, to go out on the next `poll_out` (RFC 9002 §6.2.4).
pub fn[A] Endpoint::tick(self : Endpoint[A], now : Int64) -> Unit {
let expired : Array[Bytes] = []
for cid in self.conns.keys().collect() {
guard self.conns.get(cid) is Some(conn) else { continue }
match conn.phase {
Active =>
if conn.idle_timeout > 0L && now - conn.last_rx >= conn.idle_timeout {
expired.push(cid)
} else {
let _ = conn.sender.on_pto_timeout(now)
}
Closing | Draining => if now >= conn.close_at { expired.push(cid) }
}
}
for cid in expired {
self.conns.remove(cid)
}
}
///|
pub extend Phase with Debug::{to_repr}
///|
pub extend Phase with Eq::{not_equal, equal}
///|
/// A protected long-header packet carrying `frames` at this level.
///
/// The three lines it stands for — build the header, seal the payload, protect the
/// header — are worth writing out once to see, and worth not writing out every time
/// after that. `@packet.Long::header` and `@crypto.Keys::seal` are still there for a
/// sender that wants the pieces.
pub fn send_long(
kind : @packet.Kind,
version : UInt,
dcid : Bytes,
scid : Bytes,
number : Int64,
size : Int,
frames : Array[@frame.Frame],
keys : @crypto.Keys,
token? : Bytes = b"",
) -> Bytes {
seal_long(
kind,
version,
dcid,
scid,
number,
size,
@packet.payload(frames),
keys,
token~,
)
}
///|
/// The same from a payload already encoded, which is what padding an Initial to the
/// minimum datagram needs (RFC 9000 §14.1).
pub fn seal_long(
kind : @packet.Kind,
version : UInt,
dcid : Bytes,
scid : Bytes,
number : Int64,
size : Int,
payload : Bytes,
keys : @crypto.Keys,
token? : Bytes = b"",
) -> Bytes {
let pn = @packet.number(number, None, size~)
let h : @packet.Long = {
packet_type: kind,
type_specific: 0,
version,
dcid,
scid,
}
keys.seal(
h.header(pn[:], payload=payload.length() + 16, token~)[:],
payload[:],
number~,
)
}
///|
/// A protected 1-RTT packet carrying `frames`.
pub fn send_short(
dcid : Bytes,
number : Int64,
size : Int,
frames : Array[@frame.Frame],
keys : @crypto.Keys,
spin? : Bool = false,
key_phase? : Bool = false,
) -> Bytes {
let payload = @packet.payload(frames)
let pn = @packet.number(number, None, size~)
let h : @packet.Short = { spin, key_phase, pn_length: size, dcid, }
keys.seal(h.header(pn[:])[:], payload[:], number~)
}
///|
/// The frames a protected long-header packet carries, and its packet number. `None` when
/// the tag does not match, which is a packet to drop rather than an error to report.
pub fn recv_long(
packet : Bytes,
keys : @crypto.Keys,
) -> (Array[@frame.Frame], Int64)? raise @packet.Refused {
match keys.open(packet[:], largest=-1L) {
Some((payload, pn)) => Some((@packet.read_payload(payload), pn))
None => None
}
}
///|
/// The same for a 1-RTT packet, whose connection-id length only the receiver knows.
pub fn recv_short(
packet : Bytes,
dcid_len : Int,
keys : @crypto.Keys,
) -> (Array[@frame.Frame], Int64)? raise @packet.Refused {
match keys.open(packet[:], largest=-1L, at=1 + dcid_len) {
Some((payload, pn)) => Some((@packet.read_payload(payload), pn))
None => None
}
}
///|
/// A CertificateVerify body signed over the transcript as RFC 8446 §4.4.3 frames it.
///
/// This is the seam that keeps a signing key out of this library: a flight takes a
/// function from the transcript hash to a CertificateVerify body, and whoever holds the
/// key builds one with this.
pub fn signed(
signer : &@spec.Signer,
scheme? : Int = @ext.ecdsa_secp256r1_sha256,
) -> (Bytes) -> Bytes {
th => {
@msg.certificate_verify(
signer.sign(@msg.signed(@msg.server_context, th[:])[:]),
scheme,
)
}
}
///|
/// Two byte strings joined.
fn join(a : Bytes, b : Bytes) -> Bytes {
let buf = Buffer()
buf.write_bytes(a)
buf.write_bytes(b)
buf.to_bytes()
}