///|
pub(all) struct StreamLimits {
max_buffer_bytes_value : Int
frame_limits_value : FrameLimits
} derive(Eq, Debug)
///|
pub struct StreamDecoder {
policy_value : ParsePolicy
limits_value : StreamLimits
buffer_value : Array[Byte]
mut record_index_value : Int
mut discarding_value : Bool
mut skip_lf_value : Bool
}
///|
fn stream_error(record_index : Int, message : String) -> NmeaError {
NmeaError::from_diagnostic(
Diagnostic::new(
StreamLimitExceeded,
Error,
message,
SourceRef::sentence(record_index),
),
)
}
///|
fn stream_sentence_error(record_index : Int, message : String) -> NmeaError {
NmeaError::from_diagnostic(
Diagnostic::new(
SentenceInvalid,
Error,
message,
SourceRef::sentence(record_index),
),
)
}
///|
pub fn StreamLimits::new(
max_buffer_bytes : Int,
frame_limits? : FrameLimits = FrameLimits::standard(),
) -> Result[StreamLimits, NmeaError] {
if max_buffer_bytes <= 0 {
return Err(stream_error(0, "stream buffer limit must be positive"))
}
Ok({
max_buffer_bytes_value: max_buffer_bytes,
frame_limits_value: frame_limits,
})
}
///|
pub fn StreamLimits::standard() -> StreamLimits {
{ max_buffer_bytes_value: 1024, frame_limits_value: FrameLimits::standard() }
}
///|
pub fn StreamLimits::max_buffer_bytes(self : StreamLimits) -> Int {
self.max_buffer_bytes_value
}
///|
pub fn StreamLimits::frame_limits(self : StreamLimits) -> FrameLimits {
self.frame_limits_value
}
///|
pub fn StreamDecoder::new(
policy? : ParsePolicy = Strict,
limits? : StreamLimits = StreamLimits::standard(),
) -> StreamDecoder {
{
policy_value: policy,
limits_value: limits,
buffer_value: [],
record_index_value: 0,
discarding_value: false,
skip_lf_value: false,
}
}
///|
pub fn StreamDecoder::buffered_bytes(self : StreamDecoder) -> Int {
self.buffer_value.length()
}
///|
fn StreamDecoder::complete_record(
self : StreamDecoder,
) -> SentenceResult[NmeaRecord] {
let record_index = self.record_index_value
self.record_index_value += 1
for byte in self.buffer_value {
if byte.to_int() > 127 {
return SentenceResult::err(
record_index,
NotChecked,
stream_sentence_error(record_index, "NMEA stream records must be ASCII"),
)
}
}
let text = @ascii.decode_lossy(Bytes::from_array(self.buffer_value).view())
match
parse_frame(
text,
policy=self.policy_value,
limits=self.limits_value.frame_limits_value,
) {
Err(error) => {
let checksum_status = if error.code() == "checksum.mismatch" {
PresentInvalid
} else {
NotChecked
}
SentenceResult::err(record_index, checksum_status, error)
}
Ok(frame) => {
let checksum_status = frame.checksum_status()
match parse_record(frame) {
Ok(record) =>
SentenceResult::ok(
record_index,
record.formatter(),
checksum_status,
value=Some(record),
)
Err(error) => SentenceResult::err(record_index, checksum_status, error)
}
}
}
}
///|
fn StreamDecoder::end_line(
self : StreamDecoder,
results : Array[SentenceResult[NmeaRecord]],
) -> Unit {
if self.discarding_value {
self.discarding_value = false
} else if self.buffer_value.length() > 0 {
results.push(self.complete_record())
}
self.buffer_value.clear()
}
///|
pub fn StreamDecoder::push(
self : StreamDecoder,
chunk : Bytes,
) -> Array[SentenceResult[NmeaRecord]] {
let results : Array[SentenceResult[NmeaRecord]] = []
for byte in chunk {
let value = byte.to_int()
if value == 13 {
self.end_line(results)
self.skip_lf_value = true
} else if value == 10 {
if self.skip_lf_value {
self.skip_lf_value = false
} else {
self.end_line(results)
}
} else {
self.skip_lf_value = false
if !self.discarding_value {
if self.buffer_value.length() >=
self.limits_value.max_buffer_bytes_value {
let record_index = self.record_index_value
self.record_index_value += 1
results.push(
SentenceResult::err(
record_index,
NotChecked,
stream_error(
record_index, "stream record exceeds configured buffer limit",
),
),
)
self.buffer_value.clear()
self.discarding_value = true
} else {
self.buffer_value.push(byte)
}
}
}
}
results
}
///|
pub fn StreamDecoder::finish(
self : StreamDecoder,
) -> Array[SentenceResult[NmeaRecord]] {
let results : Array[SentenceResult[NmeaRecord]] = []
if self.discarding_value {
self.discarding_value = false
self.buffer_value.clear()
} else if self.buffer_value.length() > 0 {
results.push(self.complete_record())
self.buffer_value.clear()
}
self.skip_lf_value = false
results
}