use std::sync::Arc;
use rudb_common::{Cancel, Error, Result};
use rudb_metrics::Driver;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use rudb_pipeline::{Pipeline, Pool, RootReader, run_parallel};
use rudb_vector::Chunk;
use crate::schema::Schema;
#[derive(Debug)]
pub struct Query<'a> {
pipelines: Vec<Pipeline<'a>>,
drivers: Vec<Arc<Driver>>,
reader: RootReader,
schema: Schema,
worker_cpu_ns: AtomicU64,
widest: AtomicUsize,
}
impl<'a> Query<'a> {
pub(crate) fn new(
pipelines: Vec<Pipeline<'a>>,
drivers: Vec<Arc<Driver>>,
reader: RootReader,
schema: Schema,
) -> Result<Self> {
for (at, pipeline) in pipelines.iter().enumerate() {
for waited in pipeline.depends_on() {
let before = pipelines[..at].iter().any(|earlier| earlier.id() == *waited);
if !before {
return Err(Error::internal(format!(
"{} waits for {waited}, which the builder did not put before it",
pipeline.id()
)));
}
}
}
Ok(Self {
pipelines,
drivers,
reader,
schema,
worker_cpu_ns: AtomicU64::new(0),
widest: AtomicUsize::new(0),
})
}
#[must_use]
pub fn schema(&self) -> &Schema {
&self.schema
}
#[must_use]
pub fn pipelines(&self) -> usize {
self.pipelines.len()
}
pub fn run(&self, cancel: &Cancel, pool: &Pool) -> Result<()> {
for (pipeline, driver) in self.pipelines.iter().zip(&self.drivers) {
let lease = pool.lease(pipeline.degree(pool.threads()));
let degree = lease.degree();
let spent = {
let _running = driver.running();
run_parallel(pipeline, cancel, degree)?
};
driver.ran(degree, spent);
self.worker_cpu_ns.fetch_add(spent, Ordering::Relaxed);
self.widest.fetch_max(degree, Ordering::Relaxed);
}
Ok(())
}
#[must_use]
pub fn worker_cpu_ns(&self) -> u64 {
self.worker_cpu_ns.load(Ordering::Relaxed)
}
#[must_use]
pub fn widest(&self) -> usize {
self.widest.load(Ordering::Relaxed)
}
pub fn next_chunk(&self) -> Result<Option<Chunk>> {
self.reader.next_chunk()
}
pub fn collect(&self, cancel: &Cancel, pool: &Pool) -> Result<Vec<Chunk>> {
self.run(cancel, pool)?;
let mut chunks = Vec::new();
while let Some(chunk) = self.next_chunk()? {
chunks.push(chunk);
}
Ok(chunks)
}
}