1use std::collections::VecDeque;
43use std::sync::atomic::{AtomicU64, Ordering};
44use std::sync::{Arc, Condvar, Mutex};
45use std::thread::JoinHandle;
46
47use rudb_common::Result;
48
49use crate::File;
50use crate::submit::{Completion, Filler, Request, Response};
51
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub struct Config {
55 pub threads: usize,
58 pub coalesce_gap: u64,
61 pub coalesce_span: u64,
65 pub coalesce: bool,
67}
68
69impl Config {
70 #[must_use]
95 pub fn local_disk() -> Self {
96 Self {
97 threads: cores().clamp(2, 8),
98 coalesce_gap: 16 << 10,
99 coalesce_span: 1 << 20,
100 coalesce: true,
101 }
102 }
103
104 #[must_use]
114 pub fn object_store() -> Self {
115 Self {
116 threads: (cores() * 8).clamp(32, 128),
117 coalesce_gap: 512 << 10,
118 coalesce_span: 8 << 20,
119 coalesce: true,
120 }
121 }
122
123 #[must_use]
125 pub fn with_threads(mut self, threads: usize) -> Self {
126 self.threads = threads.max(1);
127 self
128 }
129
130 #[must_use]
132 pub fn coalescing(mut self, gap: u64) -> Self {
133 self.coalesce = true;
134 self.coalesce_gap = gap;
135 self
136 }
137
138 #[must_use]
140 pub fn not_coalescing(mut self) -> Self {
141 self.coalesce = false;
142 self
143 }
144}
145
146impl Default for Config {
147 fn default() -> Self {
148 Self::local_disk()
149 }
150}
151
152fn cores() -> usize {
153 std::thread::available_parallelism().map_or(4, std::num::NonZero::get)
154}
155
156#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
164pub struct Stats {
165 pub requests: u64,
167 pub reads: u64,
169 pub wanted: u64,
171 pub read: u64,
173}
174
175#[derive(Debug, Default)]
176struct Counters {
177 requests: AtomicU64,
178 reads: AtomicU64,
179 wanted: AtomicU64,
180 read: AtomicU64,
181}
182
183impl Counters {
184 fn snapshot(&self) -> Stats {
185 Stats {
186 requests: self.requests.load(Ordering::Relaxed),
187 reads: self.reads.load(Ordering::Relaxed),
188 wanted: self.wanted.load(Ordering::Relaxed),
189 read: self.read.load(Ordering::Relaxed),
190 }
191 }
192}
193
194#[derive(Debug)]
196struct Part {
197 index: usize,
198 offset: u64,
199 buf: Vec<u8>,
200}
201
202#[derive(Debug)]
204struct Job {
205 file: Arc<dyn File>,
206 filler: Filler,
207 offset: u64,
208 span: usize,
209 parts: Vec<Part>,
210}
211
212impl Job {
213 fn run(self, counters: &Counters) {
215 let Self { file, filler, offset, span, mut parts } = self;
216 counters.reads.fetch_add(1, Ordering::Relaxed);
217
218 if parts.len() == 1 {
221 let part = parts.pop().unwrap_or_else(|| unreachable!("checked one part"));
222 let Part { index, offset, mut buf } = part;
223 let outcome = file.read_at(offset, &mut buf).map(|read| {
224 counters.read.fetch_add(read as u64, Ordering::Relaxed);
225 Response::new(index, offset, read, buf)
226 });
227 filler.finish(index, outcome);
228 return;
229 }
230
231 let mut scratch = vec![0u8; span];
232 match file.read_at(offset, &mut scratch) {
233 Ok(read) => {
234 counters.read.fetch_add(read as u64, Ordering::Relaxed);
235 for part in parts {
236 let Part { index, offset: at, mut buf } = part;
237 let start = (at - offset) as usize;
238 let got = read.saturating_sub(start).min(buf.len());
242 buf[..got].copy_from_slice(&scratch[start..start + got]);
243 filler.finish(index, Ok(Response::new(index, at, got, buf)));
244 }
245 }
246 Err(error) => {
247 for part in parts {
248 filler.finish(part.index, Err(error.clone()));
249 }
250 }
251 }
252 }
253}
254
255#[derive(Debug)]
256struct Queue {
257 jobs: VecDeque<Job>,
258 closed: bool,
259}
260
261#[derive(Debug)]
262struct Shared {
263 queue: Mutex<Queue>,
264 work: Condvar,
265 counters: Counters,
266}
267
268impl Shared {
269 fn lock(&self) -> std::sync::MutexGuard<'_, Queue> {
270 self.queue.lock().unwrap_or_else(std::sync::PoisonError::into_inner)
271 }
272}
273
274#[derive(Debug, Clone)]
279pub struct Pool {
280 shared: Arc<Shared>,
281 #[allow(dead_code, reason = "holding this is the point of it, reading it is not")]
284 threads: Arc<Threads>,
285 config: Config,
286}
287
288#[derive(Debug)]
291struct Threads {
292 shared: Arc<Shared>,
293 handles: Mutex<Vec<JoinHandle<()>>>,
294}
295
296impl Drop for Threads {
297 fn drop(&mut self) {
298 self.shared.lock().closed = true;
299 self.shared.work.notify_all();
300 let handles = std::mem::take(
301 &mut *self.handles.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
302 );
303 for handle in handles {
304 let _ = handle.join();
305 }
306 }
307}
308
309impl Pool {
310 #[must_use]
317 pub fn new(config: Config) -> Self {
318 let shared = Arc::new(Shared {
319 queue: Mutex::new(Queue { jobs: VecDeque::new(), closed: false }),
320 work: Condvar::new(),
321 counters: Counters::default(),
322 });
323 let mut handles = Vec::with_capacity(config.threads);
324 for n in 0..config.threads.max(1) {
325 let shared = Arc::clone(&shared);
326 handles.push(
327 std::thread::Builder::new()
328 .name(format!("rudb-io-{n}"))
329 .spawn(move || worker(&shared))
330 .expect("could not spawn an I/O thread"),
331 );
332 }
333 let threads =
334 Arc::new(Threads { shared: Arc::clone(&shared), handles: Mutex::new(handles) });
335 Self { shared, threads, config }
336 }
337
338 #[must_use]
340 pub fn config(&self) -> Config {
341 self.config
342 }
343
344 #[must_use]
346 pub fn stats(&self) -> Stats {
347 self.shared.counters.snapshot()
348 }
349
350 #[must_use]
356 pub fn submit(&self, file: &Arc<dyn File>, requests: Vec<Request>) -> Completion {
357 let count = requests.len();
358 self.shared.counters.requests.fetch_add(count as u64, Ordering::Relaxed);
359 let wanted: u64 = requests.iter().map(|r| r.len() as u64).sum();
360 self.shared.counters.wanted.fetch_add(wanted, Ordering::Relaxed);
361
362 let (completion, filler) = Completion::pending(count);
363 let mut wanted_nothing = Vec::new();
364 let mut real = Vec::with_capacity(count);
365 for (index, request) in requests.into_iter().enumerate() {
366 if request.is_empty() {
367 wanted_nothing.push((index, request));
368 } else {
369 real.push((index, request));
370 }
371 }
372 for (index, request) in wanted_nothing {
376 let offset = request.offset();
377 filler.finish(index, Ok(Response::new(index, offset, 0, request.into_buffer())));
378 }
379 let jobs = plan(real, self.config, file, &filler);
380 let mut queue = self.shared.lock();
381 if queue.closed {
382 drop(queue);
385 drop(filler);
386 return completion;
387 }
388 queue.jobs.extend(jobs);
389 drop(queue);
390 self.shared.work.notify_all();
391 completion
392 }
393
394 #[must_use]
396 pub fn queued(&self) -> usize {
397 self.shared.lock().jobs.len()
398 }
399}
400
401fn plan(
406 requests: Vec<(usize, Request)>,
407 config: Config,
408 file: &Arc<dyn File>,
409 filler: &Filler,
410) -> Vec<Job> {
411 let mut parts: Vec<Part> = requests
412 .into_iter()
413 .map(|(index, request)| {
414 let offset = request.offset();
415 Part { index, offset, buf: request.into_buffer() }
416 })
417 .collect();
418 parts.sort_by_key(|part| part.offset);
419
420 let mut jobs: Vec<Job> = Vec::with_capacity(parts.len());
421 for part in parts {
422 let end = part.offset + part.buf.len() as u64;
423 if config.coalesce {
424 if let Some(last) = jobs.last_mut() {
425 let last_end = last.offset + last.span as u64;
426 let gap = part.offset.saturating_sub(last_end);
427 let span = end.saturating_sub(last.offset);
428 if gap <= config.coalesce_gap && span <= config.coalesce_span {
430 last.span = span as usize;
431 last.parts.push(part);
432 continue;
433 }
434 }
435 }
436 jobs.push(Job {
437 file: Arc::clone(file),
438 filler: filler.clone(),
439 offset: part.offset,
440 span: part.buf.len(),
441 parts: vec![part],
442 });
443 }
444 jobs
445}
446
447fn worker(shared: &Arc<Shared>) {
448 loop {
449 let mut queue = shared.lock();
450 let job = loop {
451 if let Some(job) = queue.jobs.pop_front() {
452 break job;
453 }
454 if queue.closed {
455 return;
456 }
457 queue = shared.work.wait(queue).unwrap_or_else(std::sync::PoisonError::into_inner);
458 };
459 drop(queue);
460 job.run(&shared.counters);
461 }
462}
463
464#[derive(Debug)]
470pub struct Pooled {
471 file: Arc<dyn File>,
472 pool: Pool,
473}
474
475impl Pooled {
476 #[must_use]
478 pub fn new(file: Box<dyn File>, pool: Pool) -> Self {
479 Self { file: Arc::from(file), pool }
480 }
481
482 #[must_use]
484 pub fn pool(&self) -> &Pool {
485 &self.pool
486 }
487}
488
489impl File for Pooled {
490 fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<usize> {
491 self.file.read_at(offset, buf)
492 }
493
494 fn submit(&self, requests: Vec<Request>) -> Completion {
495 self.pool.submit(&self.file, requests)
496 }
497
498 fn write_at(&self, offset: u64, data: &[u8]) -> Result<()> {
499 self.file.write_at(offset, data)
500 }
501
502 fn sync(&self) -> Result<()> {
503 self.file.sync()
504 }
505
506 fn truncate(&self, len: u64) -> Result<()> {
507 self.file.truncate(len)
508 }
509
510 fn len(&self) -> Result<u64> {
511 self.file.len()
512 }
513}
514
515#[cfg(test)]
516mod tests {
517 use std::path::Path;
518
519 use super::{Config, Pool, Pooled};
520 use crate::submit::{Request, Response};
521 use crate::{File, Filesystem, OpenMode, SimFilesystem};
522
523 fn ramp(pool: Pool) -> Pooled {
526 let fs = SimFilesystem::new();
527 let file = fs.open(Path::new("/data"), OpenMode::Create).unwrap();
528 let bytes: Vec<u8> = (0..=255u8).collect();
529 file.write_at(0, &bytes).unwrap();
530 file.sync().unwrap();
531 Pooled::new(file, pool)
532 }
533
534 #[test]
535 fn a_batch_comes_back_answering_the_requests_it_was_given() {
536 let file = ramp(Pool::new(Config::local_disk()));
537 let responses = file
538 .submit(vec![Request::new(100, 4), Request::new(0, 4), Request::new(200, 4)])
539 .wait()
540 .unwrap();
541 assert_eq!(responses[0].bytes(), &[100, 101, 102, 103]);
542 assert_eq!(responses[1].bytes(), &[0, 1, 2, 3]);
543 assert_eq!(responses[2].bytes(), &[200, 201, 202, 203]);
544 }
545
546 #[test]
547 fn one_thread_answers_everything_just_as_well_as_eight() {
548 for threads in [1, 2, 8] {
551 let file = ramp(Pool::new(Config::local_disk().with_threads(threads)));
552 let requests = (0..32).map(|i| Request::new(i * 8, 8)).collect::<Vec<_>>();
553 let responses = file.submit(requests).wait().unwrap();
554 assert_eq!(responses.len(), 32, "at {threads} threads");
555 for (i, response) in responses.iter().enumerate() {
556 assert_eq!(response.bytes()[0], (i * 8) as u8, "at {threads} threads");
557 }
558 }
559 }
560
561 #[test]
562 fn adjacent_ranges_become_one_read_and_the_bytes_do_not_change() {
563 let pool = Pool::new(Config::local_disk().with_threads(1).coalescing(0));
564 let file = ramp(pool.clone());
565 let responses = file
566 .submit(vec![Request::new(0, 8), Request::new(8, 8), Request::new(16, 8)])
567 .wait()
568 .unwrap();
569 assert_eq!(pool.stats().reads, 1, "three adjacent ranges are one read");
570 assert_eq!(pool.stats().requests, 3);
571 assert_eq!(pool.stats().wanted, 24);
572 assert_eq!(pool.stats().read, 24, "and no byte was read that nobody asked for");
573 for (i, response) in responses.iter().enumerate() {
574 assert_eq!(response.bytes()[0], (i * 8) as u8);
575 }
576 }
577
578 #[test]
579 fn a_gap_is_read_and_thrown_away_and_the_byte_count_says_so() {
580 let pool = Pool::new(Config::local_disk().with_threads(1).coalescing(16));
584 let file = ramp(pool.clone());
585 let responses = file.submit(vec![Request::new(0, 4), Request::new(20, 4)]).wait().unwrap();
586 assert_eq!(pool.stats().reads, 1);
587 assert_eq!(pool.stats().wanted, 8);
588 assert_eq!(pool.stats().read, 24, "the sixteen byte gap was read too");
589 assert_eq!(responses[0].bytes(), &[0, 1, 2, 3]);
590 assert_eq!(responses[1].bytes(), &[20, 21, 22, 23]);
591 }
592
593 #[test]
594 fn a_gap_wider_than_the_policy_stays_two_reads() {
595 let pool = Pool::new(Config::local_disk().with_threads(1).coalescing(4));
596 let file = ramp(pool.clone());
597 file.submit(vec![Request::new(0, 4), Request::new(20, 4)]).wait().unwrap();
598 assert_eq!(pool.stats().reads, 2);
599 assert_eq!(pool.stats().read, 8);
600 }
601
602 #[test]
603 fn the_span_limit_stops_a_chain_of_small_gaps_becoming_one_huge_read() {
604 let mut config = Config::local_disk().with_threads(1).coalescing(12);
607 config.coalesce_span = 32;
608 let pool = Pool::new(config);
609 let file = ramp(pool.clone());
610 let requests = (0..8).map(|i| Request::new(i * 16, 4)).collect::<Vec<_>>();
611 file.submit(requests).wait().unwrap();
612 assert_eq!(pool.stats().reads, 4, "eight ranges over 128 bytes, capped at 32 a read");
613 }
614
615 #[test]
616 fn both_defaults_coalesce_and_the_local_one_does_it_far_more_narrowly() {
617 assert!(Config::local_disk().coalesce);
622 assert!(Config::object_store().coalesce);
623 assert!(Config::object_store().coalesce_gap >= Config::local_disk().coalesce_gap * 8);
624 assert!(Config::object_store().threads >= Config::local_disk().threads * 4);
626 }
627
628 #[test]
629 fn the_local_gap_is_too_small_to_swallow_the_space_between_two_columns() {
630 let pool = Pool::new(Config::local_disk().with_threads(1));
634 let fs = SimFilesystem::new();
635 let handle = fs.open(Path::new("/columns"), OpenMode::Create).unwrap();
636 handle.write_at(0, &vec![7u8; 2 << 20]).unwrap();
637 handle.sync().unwrap();
638 let file = Pooled::new(handle, pool.clone());
639 let requests = (0..4).map(|i| Request::new(i * (512 << 10), 64 << 10)).collect::<Vec<_>>();
640 file.submit(requests).wait().unwrap();
641 assert_eq!(pool.stats().reads, 4, "four columns 448KiB apart are four reads and not one");
642 assert_eq!(pool.stats().read, pool.stats().wanted, "and nothing else was read");
643 }
644
645 #[test]
646 fn a_short_read_stays_short_through_a_merge() {
647 let pool = Pool::new(Config::local_disk().with_threads(1).coalescing(0));
648 let file = ramp(pool.clone());
649 let responses =
651 file.submit(vec![Request::new(248, 8), Request::new(256, 8)]).wait().unwrap();
652 assert_eq!(pool.stats().reads, 1);
653 assert!(!responses[0].is_short());
654 assert!(responses[1].is_short());
655 assert_eq!(responses[1].read(), 0);
656 }
657
658 #[test]
659 fn a_failed_read_fails_every_request_it_was_merged_with_and_no_others() {
660 let fs = SimFilesystem::new();
661 let file = fs.open(Path::new("/data"), OpenMode::Create).unwrap();
662 file.write_at(0, &(0..=255u8).collect::<Vec<u8>>()).unwrap();
663 file.sync().unwrap();
664 let pool = Pool::new(Config::local_disk().with_threads(1).coalescing(0));
665 let pooled = Pooled::new(file, pool);
666 fs.fail_read_at(0);
669 let mut completion =
670 pooled.submit(vec![Request::new(0, 8), Request::new(8, 8), Request::new(128, 8)]);
671 let mut failed = 0;
672 let mut answered = 0;
673 while let Some(outcome) = completion.take() {
674 match outcome {
675 Ok(_) => answered += 1,
676 Err(_) => failed += 1,
677 }
678 }
679 assert_eq!((answered, failed), (1, 2));
680 }
681
682 #[test]
683 fn an_empty_request_is_answered_rather_than_queued() {
684 let pool = Pool::new(Config::local_disk().with_threads(1));
685 let file = ramp(pool.clone());
686 let responses = file.submit(vec![Request::new(0, 0), Request::new(4, 4)]).wait().unwrap();
687 assert_eq!(pool.stats().reads, 1, "nobody reads nothing");
688 assert_eq!(responses.len(), 2);
691 assert_eq!(responses[0].read(), 0);
692 assert_eq!(responses[1].bytes(), &[4, 5, 6, 7]);
693 }
694
695 #[test]
696 fn an_empty_batch_is_done_before_it_is_submitted() {
697 let pool = Pool::new(Config::local_disk());
698 let file = ramp(pool.clone());
699 let completion = file.submit(Vec::new());
700 assert!(completion.is_done());
701 assert!(completion.wait().unwrap().is_empty());
702 }
703
704 #[test]
705 fn a_submission_returns_before_the_reads_do() {
706 let pool = Pool::new(Config::local_disk().with_threads(1));
710 let file = ramp(pool.clone());
711 let requests = (0..64).map(|i| Request::new(i * 4, 4)).collect::<Vec<_>>();
712 let completion = file.submit(requests);
713 let responses = completion.wait().unwrap();
716 assert_eq!(responses.len(), 64);
717 assert_eq!(pool.stats().requests, 64);
718 }
719
720 #[test]
721 fn everything_submitted_is_answered_once_the_pool_is_stopping() {
722 let pool = Pool::new(Config::local_disk().with_threads(2));
723 let file = ramp(pool.clone());
724 let completion = file.submit(vec![Request::new(0, 4)]);
725 drop(pool);
728 let responses = completion.wait();
729 assert!(responses.is_ok() || responses.is_err());
730 }
731
732 #[test]
733 fn read_at_on_a_pooled_file_is_the_read_underneath_it() {
734 let file = ramp(Pool::new(Config::local_disk()));
735 let mut buf = [0u8; 4];
736 file.read_exact_at(64, &mut buf).unwrap();
737 assert_eq!(buf, [64, 65, 66, 67]);
738 assert_eq!(file.len().unwrap(), 256);
739 }
740
741 #[test]
742 fn responses_are_in_submission_order_whatever_order_the_threads_finished_in() {
743 let file = ramp(Pool::new(Config::local_disk().with_threads(8)));
744 let requests = (0..32).rev().map(|i| Request::new(i * 8, 8)).collect::<Vec<_>>();
747 let responses = file.submit(requests).wait().unwrap();
748 let indices: Vec<usize> = responses.iter().map(Response::index).collect();
749 assert_eq!(indices, (0..32).collect::<Vec<_>>());
750 for (i, response) in responses.iter().enumerate() {
751 assert_eq!(response.bytes()[0], ((31 - i) * 8) as u8);
752 }
753 }
754}