use crate::reorder_buffer::ReorderBuffer;
use crate::write::file_writer::FileWriter;
use crate::write::metrics::WRITER_METRICS;
use std::io::{self, Write};
use std::sync::atomic::Ordering::Relaxed;
use std::sync::mpsc::{Receiver, SyncSender, sync_channel};
use std::thread::JoinHandle;
pub(crate) const WRITE_AHEAD: usize = 32;
pub(crate) const PIPELINE_DISPATCH_PERMITS: usize = 64;
pub(crate) enum OutputChunk {
Framed(FramedBlobParts),
Raw(Vec<u8>),
RawChunks(Vec<Vec<u8>>),
#[cfg(feature = "linux-direct-io")]
CopyRange {
in_fd: std::os::unix::io::RawFd,
offset: u64,
len: u64,
},
}
pub(crate) struct FramedBlobParts {
pub(crate) prefix: [u8; 4],
pub(crate) header: Vec<u8>,
pub(crate) blob: Vec<u8>,
}
impl FramedBlobParts {
pub(crate) fn total_len(&self) -> u64 {
(self.prefix.len() + self.header.len() + self.blob.len()) as u64
}
pub(crate) fn into_vec(self) -> Vec<u8> {
let total_len = self.prefix.len() + self.header.len() + self.blob.len();
let mut out = Vec::with_capacity(total_len);
out.extend_from_slice(&self.prefix);
out.extend_from_slice(&self.header);
out.extend_from_slice(&self.blob);
out
}
}
pub(crate) struct PipelineItem {
pub(crate) seq: usize,
pub(crate) data: io::Result<OutputChunk>,
}
pub(crate) trait OutputSink {
fn write_chunk(&mut self, chunk: OutputChunk) -> io::Result<()>;
fn flush(&mut self) -> io::Result<()>;
}
pub(crate) struct FileOutputSink {
writer: FileWriter,
}
impl FileOutputSink {
pub(crate) fn new(writer: FileWriter) -> Self {
Self { writer }
}
}
impl OutputSink for FileOutputSink {
fn write_chunk(&mut self, chunk: OutputChunk) -> io::Result<()> {
match chunk {
OutputChunk::Framed(parts) => {
let len = parts.total_len();
let bytes = parts.into_vec();
let t_write = std::time::Instant::now();
self.writer.write_all(&bytes)?;
WRITER_METRICS
.write_ns
.fetch_add(elapsed_ns_u64(t_write), Relaxed);
WRITER_METRICS.bytes_written.fetch_add(len, Relaxed);
Ok(())
}
OutputChunk::Raw(bytes) => {
let len = bytes.len() as u64;
let t_write = std::time::Instant::now();
self.writer.write_all(&bytes)?;
WRITER_METRICS
.write_ns
.fetch_add(elapsed_ns_u64(t_write), Relaxed);
WRITER_METRICS.bytes_written.fetch_add(len, Relaxed);
Ok(())
}
OutputChunk::RawChunks(chunks) => {
let total_bytes: u64 = chunks.iter().map(|chunk| chunk.len() as u64).sum();
let t_write = std::time::Instant::now();
for chunk in &chunks {
self.writer.write_all(chunk)?;
}
WRITER_METRICS
.write_ns
.fetch_add(elapsed_ns_u64(t_write), Relaxed);
WRITER_METRICS.bytes_written.fetch_add(total_bytes, Relaxed);
Ok(())
}
#[cfg(feature = "linux-direct-io")]
OutputChunk::CopyRange { in_fd, offset, len } => {
let t_write = std::time::Instant::now();
let out_fd = self.writer.flush_and_raw_fd()?.ok_or_else(|| {
io::Error::new(
io::ErrorKind::Unsupported,
"copy_file_range incompatible with O_DIRECT output",
)
})?;
super::copy_range::copy_range(in_fd, out_fd, offset, len)?;
WRITER_METRICS
.write_ns
.fetch_add(elapsed_ns_u64(t_write), Relaxed);
WRITER_METRICS.bytes_written.fetch_add(len, Relaxed);
Ok(())
}
}
}
fn flush(&mut self) -> io::Result<()> {
self.writer.flush()
}
}
pub(super) struct WritePipeline {
pub(super) tx: SyncSender<PipelineItem>,
pub(super) seq: usize,
pub(super) join_handle: Option<JoinHandle<io::Result<()>>>,
pub(super) permit_tx: SyncSender<()>,
pub(super) permit_rx: Receiver<()>,
}
pub(super) fn new_permit_pool() -> (SyncSender<()>, Receiver<()>) {
let (permit_tx, permit_rx) = sync_channel::<()>(PIPELINE_DISPATCH_PERMITS);
for _ in 0..PIPELINE_DISPATCH_PERMITS {
permit_tx
.send(())
.expect("seeding permit pool on fresh channel");
}
(permit_tx, permit_rx)
}
pub(super) fn record_send_wait(send_start: std::time::Instant) {
WRITER_METRICS
.pipeline_send_wait_ns
.fetch_add(elapsed_ns_u64(send_start), Relaxed);
}
pub(crate) fn elapsed_ns_u64(start: std::time::Instant) -> u64 {
u64::try_from(start.elapsed().as_nanos()).unwrap_or(u64::MAX)
}
#[allow(clippy::needless_pass_by_value)] pub(super) fn writer_thread<S: OutputSink>(
rx: std::sync::mpsc::Receiver<PipelineItem>,
mut sink: S,
) -> io::Result<()> {
let mut pending: ReorderBuffer<io::Result<OutputChunk>> =
ReorderBuffer::with_capacity(WRITE_AHEAD);
loop {
let t_recv = std::time::Instant::now();
let item = match rx.recv() {
Ok(item) => item,
Err(_) => break,
};
WRITER_METRICS
.recv_wait_ns
.fetch_add(elapsed_ns_u64(t_recv), Relaxed);
pending.push(item.seq, item.data);
WRITER_METRICS.record_reorder_high_water(pending.pending_len());
while let Some(result) = pending.pop_ready() {
let chunk = result?;
sink.write_chunk(chunk)?;
}
}
sink.flush()?;
Ok(())
}