source.nct
see ./index.nct
use std/io.BlockingReader
use std/vec.Vec
const SOURCE_BUFFER_BYTES: usize = 8192
enum SourceLimitKind {
compressed_input
decoded_output
}
struct LimitedSource<R> {
source: R
buffer: Vec<u8>
next: usize
len: usize
accepted: u64
maximum: u64
limit_kind: SourceLimitKind
finished: bool
failed: bool
}
func source_buffer(length: usize): Vec<u8> {
var bytes: Vec<u8> = Vec.with_capacity(length)
while bytes.len() < length { bytes.push(0) }
return move bytes
}
construct LimitedSource<R> {
func compressed(source: R, maximum: usize): Self where R impl BlockingReader {
return LimitedSource.with_limit(
move source,
maximum as u64,
SourceLimitKind.compressed_input,
)
}
func decoded(source: R, maximum: u64): Self where R impl BlockingReader {
return LimitedSource.with_limit(move source, maximum, SourceLimitKind.decoded_output)
}
func with_limit(source: R, maximum: u64, limit_kind: SourceLimitKind): Self where R impl BlockingReader {
let buffer_bytes = SOURCE_BUFFER_BYTES as u64
let capacity = if maximum < buffer_bytes {
(maximum as usize) + 1
} else {
buffer_bytes as usize
}
return LimitedSource<R> {
source: move source,
buffer: source_buffer(capacity),
next: 0,
len: 0,
accepted: 0,
maximum: maximum,
limit_kind: limit_kind,
finished: false,
failed: false,
}
}
}
noalloc func compressed_input_limit_error(): error {
return error.new("archive-inspect.input_limit", "compressed input exceeds its configured limit")
}
noalloc func decoded_output_limit_error(): error {
return error.new("archive-inspect.output_limit", "decoded output exceeds its configured limit")
}
noalloc func source_limit_error(kind: SourceLimitKind): error {
match kind {
SourceLimitKind.compressed_input { return compressed_input_limit_error() }
SourceLimitKind.decoded_output { return decoded_output_limit_error() }
}
}
noalloc func invalid_source_progress(): error {
return error.new("archive-inspect.invalid_progress", "source reported bytes outside its buffer")
}
blocking func refill_limited_source<R>(reader: &+LimitedSource<R>): void! where R impl BlockingReader {
let remaining = reader.maximum - reader.accepted
let buffer_bytes = SOURCE_BUFFER_BYTES as u64
let requested = if remaining < buffer_bytes {
(remaining as usize) + 1
} else {
SOURCE_BUFFER_BYTES
}
reader.buffer.truncate(requested)
while reader.buffer.len() < requested { reader.buffer.push(0) }
let received = reader.source.read_blocking(&+reader.buffer) catch failure {
reader.failed = true
return move failure
}
if received > requested {
reader.failed = true
return invalid_source_progress()
}
if (received as u64) > remaining {
reader.failed = true
return source_limit_error(reader.limit_kind)
}
reader.next = 0
reader.len = received
reader.accepted += received as u64
if received == 0 { reader.finished = true }
return
}
instance LimitedSource<R> where R impl BlockingReader {
blocking method &+self.read_blocking(output: &+[u8]): usize! {
if self.failed { return invalid_source_progress() }
if output.len() == 0 { return 0 }
if self.next == self.len {
if self.finished { return 0 }
refill_limited_source(self)?
if self.finished { return 0 }
}
let available = self.len - self.next
let count = if available < output.len() { available } else { output.len() }
var offset: usize = 0
while offset < count {
output[offset] = self.buffer[self.next + offset]
offset += 1
}
self.next += count
return count
}
}