pub mod common;
pub mod file;
pub mod merge;
pub mod mmap;
use crate::common::IAsyncIO;
use crate::merge::merge_overlapping_writes;
use async_trait::async_trait;
use futures::stream::iter;
use futures::stream::StreamExt;
use off64::chrono::Off64AsyncReadChrono;
use off64::chrono::Off64AsyncWriteChrono;
use off64::int::Off64AsyncReadInt;
use off64::int::Off64AsyncWriteInt;
use off64::u64;
use off64::Off64AsyncRead;
use off64::Off64AsyncWrite;
use signal_future::SignalFuture;
use signal_future::SignalFutureController;
use std::io::SeekFrom;
use std::path::Path;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
use std::sync::Arc;
use std::time::Duration;
use tokio::fs::File;
use tokio::io;
use tokio::io::AsyncSeekExt;
use tokio::sync::Mutex;
use tokio::time::sleep;
use tokio::time::Instant;
pub async fn get_file_len_via_seek(path: &Path) -> io::Result<u64> {
let mut file = File::open(path).await?;
file.seek(SeekFrom::End(0)).await
}
fn dur_us(dur: Instant) -> u64 {
dur.elapsed().as_micros().try_into().unwrap()
}
pub struct WriteRequest<D: AsRef<[u8]> + Send + 'static> {
data: D,
offset: u64,
}
impl<D: AsRef<[u8]> + Send + 'static> WriteRequest<D> {
pub fn new(offset: u64, data: D) -> Self {
Self { data, offset }
}
}
struct PendingSyncState {
earliest_unsynced: Option<Instant>, latest_unsynced: Option<Instant>,
pending_sync_fut_states: Vec<SignalFutureController>,
}
#[derive(Default, Debug)]
pub struct SeekableAsyncFileMetrics {
sync_background_loops_counter: AtomicU64,
sync_counter: AtomicU64,
sync_delayed_counter: AtomicU64,
sync_longest_delay_us_counter: AtomicU64,
sync_shortest_delay_us_counter: AtomicU64,
sync_us_counter: AtomicU64,
write_bytes_counter: AtomicU64,
write_counter: AtomicU64,
write_us_counter: AtomicU64,
}
impl SeekableAsyncFileMetrics {
pub fn sync_background_loops_counter(&self) -> u64 {
self.sync_background_loops_counter.load(Ordering::Relaxed)
}
pub fn sync_counter(&self) -> u64 {
self.sync_counter.load(Ordering::Relaxed)
}
pub fn sync_delayed_counter(&self) -> u64 {
self.sync_delayed_counter.load(Ordering::Relaxed)
}
pub fn sync_longest_delay_us_counter(&self) -> u64 {
self.sync_longest_delay_us_counter.load(Ordering::Relaxed)
}
pub fn sync_shortest_delay_us_counter(&self) -> u64 {
self.sync_shortest_delay_us_counter.load(Ordering::Relaxed)
}
pub fn sync_us_counter(&self) -> u64 {
self.sync_us_counter.load(Ordering::Relaxed)
}
pub fn write_bytes_counter(&self) -> u64 {
self.write_bytes_counter.load(Ordering::Relaxed)
}
pub fn write_counter(&self) -> u64 {
self.write_counter.load(Ordering::Relaxed)
}
pub fn write_us_counter(&self) -> u64 {
self.write_us_counter.load(Ordering::Relaxed)
}
}
#[derive(Clone)]
pub struct SeekableAsyncFile {
io: Arc<dyn IAsyncIO>,
sync_delay_us: u64,
metrics: Arc<SeekableAsyncFileMetrics>,
pending_sync_state: Arc<Mutex<PendingSyncState>>,
}
impl SeekableAsyncFile {
pub async fn open(
io: Arc<dyn IAsyncIO>,
metrics: Arc<SeekableAsyncFileMetrics>,
sync_delay: Duration,
) -> Self {
SeekableAsyncFile {
io,
sync_delay_us: sync_delay.as_micros().try_into().unwrap(),
metrics,
pending_sync_state: Arc::new(Mutex::new(PendingSyncState {
earliest_unsynced: None,
latest_unsynced: None,
pending_sync_fut_states: Vec::new(),
})),
}
}
fn bump_write_metrics(&self, len: u64, call_us: u64) {
self
.metrics
.write_bytes_counter
.fetch_add(len, Ordering::Relaxed);
self.metrics.write_counter.fetch_add(1, Ordering::Relaxed);
self
.metrics
.write_us_counter
.fetch_add(call_us, Ordering::Relaxed);
}
pub async fn read_at(&self, offset: u64, len: u64) -> Vec<u8> {
self.io.read_at(offset, len).await
}
pub async fn write_at<D: AsRef<[u8]> + Send + 'static>(&self, offset: u64, data: D) {
let len = data.as_ref().len();
let started = Instant::now();
self.io.write_at(offset, data.as_ref()).await;
let call_us: u64 = started.elapsed().as_micros().try_into().unwrap();
self.bump_write_metrics(len.try_into().unwrap(), call_us);
}
pub async fn write_at_with_delayed_sync<D: AsRef<[u8]> + Send + 'static>(
&self,
writes: impl IntoIterator<Item = WriteRequest<D>>,
) {
let writes_vec: Vec<_> = writes.into_iter().collect();
let count = u64!(writes_vec.len());
let intervals = merge_overlapping_writes(writes_vec);
iter(intervals)
.for_each_concurrent(None, async |(offset, (_, data))| {
self.write_at(offset, data).await;
})
.await;
let (fut, fut_ctl) = SignalFuture::new();
{
let mut state = self.pending_sync_state.lock().await;
let now = Instant::now();
state.earliest_unsynced.get_or_insert(now);
state.latest_unsynced = Some(now);
state.pending_sync_fut_states.push(fut_ctl.clone());
};
self
.metrics
.sync_delayed_counter
.fetch_add(count, Ordering::Relaxed);
fut.await;
}
pub async fn start_delayed_data_sync_background_loop(&self) {
let mut futures_to_wake = Vec::new();
loop {
sleep(std::time::Duration::from_micros(self.sync_delay_us)).await;
struct SyncNow {
longest_delay_us: u64,
shortest_delay_us: u64,
}
let sync_now = {
let mut state = self.pending_sync_state.lock().await;
if !state.pending_sync_fut_states.is_empty() {
let longest_delay_us = dur_us(state.earliest_unsynced.unwrap());
let shortest_delay_us = dur_us(state.latest_unsynced.unwrap());
state.earliest_unsynced = None;
state.latest_unsynced = None;
futures_to_wake.extend(state.pending_sync_fut_states.drain(..));
Some(SyncNow {
longest_delay_us,
shortest_delay_us,
})
} else {
None
}
};
if let Some(SyncNow {
longest_delay_us,
shortest_delay_us,
}) = sync_now
{
self
.metrics
.sync_longest_delay_us_counter
.fetch_add(longest_delay_us, Ordering::Relaxed);
self
.metrics
.sync_shortest_delay_us_counter
.fetch_add(shortest_delay_us, Ordering::Relaxed);
self.sync_data().await;
for ft in futures_to_wake.drain(..) {
ft.signal();
}
};
self
.metrics
.sync_background_loops_counter
.fetch_add(1, Ordering::Relaxed);
}
}
pub async fn sync_data(&self) {
let started = Instant::now();
self.io.sync_data().await;
let sync_us: u64 = started.elapsed().as_micros().try_into().unwrap();
self.metrics.sync_counter.fetch_add(1, Ordering::Relaxed);
self
.metrics
.sync_us_counter
.fetch_add(sync_us, Ordering::Relaxed);
}
}
#[async_trait]
impl<'a> Off64AsyncRead<'a, Vec<u8>> for SeekableAsyncFile {
async fn read_at(&self, offset: u64, len: u64) -> Vec<u8> {
SeekableAsyncFile::read_at(self, offset, len).await
}
}
impl<'a> Off64AsyncReadChrono<'a, Vec<u8>> for SeekableAsyncFile {}
impl<'a> Off64AsyncReadInt<'a, Vec<u8>> for SeekableAsyncFile {}
#[async_trait]
impl Off64AsyncWrite for SeekableAsyncFile {
async fn write_at(&self, offset: u64, value: &[u8]) {
SeekableAsyncFile::write_at(self, offset, value.to_vec()).await
}
}
impl Off64AsyncWriteChrono for SeekableAsyncFile {}
impl Off64AsyncWriteInt for SeekableAsyncFile {}