use std::io;
use crossbeam::sync::MsQueue;
const TOTAL_BUFFERS_MULTIPLICATIVE: usize = 2;
const TOTAL_BUFFERS_ADDITIVE: usize = 0;
pub struct PieceBuffers {
piece_queue: MsQueue<PieceBuffer>,
}
impl PieceBuffers {
pub fn new(piece_length: usize, num_workers: usize) -> PieceBuffers {
let piece_queue = MsQueue::new();
let total_buffers = calculate_total_buffers(num_workers);
for _ in 0..total_buffers {
piece_queue.push(PieceBuffer::new(piece_length));
}
PieceBuffers { piece_queue: piece_queue }
}
pub fn checkin(&self, mut buffer: PieceBuffer) {
buffer.bytes_read = 0;
self.piece_queue.push(buffer);
}
pub fn checkout(&self) -> PieceBuffer {
self.piece_queue.pop()
}
}
fn calculate_total_buffers(num_workers: usize) -> usize {
num_workers * TOTAL_BUFFERS_MULTIPLICATIVE + TOTAL_BUFFERS_ADDITIVE
}
#[derive(PartialEq, Eq)]
pub struct PieceBuffer {
buffer: Vec<u8>,
bytes_read: usize,
}
impl PieceBuffer {
fn new(piece_length: usize) -> PieceBuffer {
PieceBuffer {
buffer: vec![0u8; piece_length],
bytes_read: 0,
}
}
pub fn write_bytes<C>(&mut self, mut callback: C) -> io::Result<usize>
where C: FnMut(&mut [u8]) -> io::Result<usize>
{
let new_bytes_read = try!(callback(&mut self.buffer[self.bytes_read..]));
self.bytes_read += new_bytes_read;
Ok(new_bytes_read)
}
pub fn is_whole(&self) -> bool {
self.bytes_read == self.buffer.len()
}
pub fn is_empty(&self) -> bool {
self.bytes_read == 0
}
pub fn as_slice(&self) -> &[u8] {
&self.buffer[..self.bytes_read]
}
}