///|
pub(all) enum DecoderState {
AwaitingPrefix
ReadingV1Line
ReadingV2Fixed
ReadingV2Body(Int)
Done
Failed
} derive(Eq, Debug)
///|
pub(all) enum DecodeProgress {
NeedMore
Done(DecodedFrame)
Failed(ProxyError)
} derive(Debug)
///|
pub struct Decoder {
policy : DecodePolicy
buffer : ChunkBuffer
mut state_value : DecoderState
mut failure : ProxyError?
} derive(Debug)
///|
pub fn Decoder::new(policy : DecodePolicy) -> Decoder {
{
policy,
buffer: ChunkBuffer::new(),
state_value: AwaitingPrefix,
failure: None,
}
}
///|
pub fn Decoder::state(self : Decoder) -> DecoderState {
self.state_value
}
///|
pub fn Decoder::buffered_len(self : Decoder) -> Int {
self.buffer.available()
}
///|
pub fn Decoder::reset(self : Decoder) -> Unit {
self.buffer.clear()
self.state_value = AwaitingPrefix
self.failure = None
}
///|
fn buffer_prefix_matches(input : ChunkBuffer, prefix : Bytes) -> Bool {
let mut different = 0
let limit = if input.available() < prefix.length() {
input.available()
} else {
prefix.length()
}
for i = 0; i < limit; i = i + 1 {
match input.peek_byte(i) {
Some(byte) => different = different | (byte.to_int() ^ prefix[i].to_int())
None => return false
}
}
different == 0
}
///|
fn v1_prefix() -> Bytes {
Bytes::from_array([b'P', b'R', b'O', b'X', b'Y'])
}
///|
fn Decoder::decoder_failure(
self : Decoder,
error : ProxyError,
) -> DecodeProgress {
self.state_value = Failed
self.failure = Some(error)
Failed(error)
}
///|
pub fn Decoder::feed(self : Decoder, chunk : Bytes) -> DecodeProgress {
if self.state_value == Done {
return Failed(
proxy_error(
DecoderAlreadyFinished,
0,
"reset before feeding a completed decoder",
),
)
}
if self.state_value == Failed {
match self.failure {
Some(err) =>
return Failed(proxy_error(DecoderFailed, err.offset, err.context))
None => return Failed(proxy_error(DecoderFailed, 0, "decoder is failed"))
}
}
self.buffer.append(chunk)
let available = self.buffer.available()
if available == 0 {
return NeedMore
}
let is_v2_prefix = buffer_prefix_matches(self.buffer, v2_signature())
let is_v1_prefix = buffer_prefix_matches(self.buffer, v1_prefix())
if !is_v2_prefix && !is_v1_prefix {
return self.decoder_failure(
proxy_error(InvalidPrefix, 0, "input is neither a v1 nor v2 prefix"),
)
}
if is_v2_prefix && available < 12 {
self.state_value = ReadingV2Fixed
return NeedMore
}
if is_v1_prefix && available < 5 {
self.state_value = ReadingV1Line
return NeedMore
}
let parsed = if is_v2_prefix {
if available < 16 {
self.state_value = ReadingV2Fixed
return NeedMore
}
let high = self.buffer.peek_byte(14).unwrap().to_int()
let low = self.buffer.peek_byte(15).unwrap().to_int()
let body_length = (high << 8) | low
if body_length > self.policy.max_v2_payload_bytes {
return self.decoder_failure(
proxy_error(
HeaderTooLarge,
14,
"declared v2 payload exceeds configured bound",
),
)
}
if available < 16 + body_length {
self.state_value = ReadingV2Body(body_length)
return NeedMore
}
self.state_value = ReadingV2Body(body_length)
decode_v2(self.buffer.copy_exact(available).unwrap(), self.policy)
} else {
self.state_value = ReadingV1Line
let mut end = -1
let limit = if available < self.policy.max_v1_line_bytes {
available
} else {
self.policy.max_v1_line_bytes
}
for i = 0; i < limit; i = i + 1 {
let byte = self.buffer.peek_byte(i).unwrap()
if byte == b'\r' {
if i + 1 >= available {
if i + 1 >= self.policy.max_v1_line_bytes {
return self.decoder_failure(
proxy_error(
V1LineTooLong,
i,
"CRLF would exceed the configured v1 line bound",
),
)
}
return NeedMore
}
if self.buffer.peek_byte(i + 1).unwrap() != b'\n' {
return self.decoder_failure(
proxy_error(MissingCrlf, i, "v1 requires CRLF"),
)
}
if i + 2 > self.policy.max_v1_line_bytes {
return self.decoder_failure(
proxy_error(
V1LineTooLong,
i + 2,
"v1 line exceeded configured bound",
),
)
}
end = i
break
}
if byte == b'\n' {
return self.decoder_failure(
proxy_error(MissingCrlf, i, "v1 requires CRLF"),
)
}
}
if end < 0 {
if available >= self.policy.max_v1_line_bytes {
return self.decoder_failure(
proxy_error(
V1LineTooLong,
self.policy.max_v1_line_bytes,
"v1 line exceeded configured bound",
),
)
}
return NeedMore
}
decode_v1(self.buffer.copy_exact(available).unwrap(), self.policy)
}
match parsed {
Ok(frame) => {
self.state_value = Done
Done(frame)
}
Err(err) =>
if err.kind == NeedMoreData {
NeedMore
} else {
self.decoder_failure(err)
}
}
}