mod context;
mod drain;
mod drained;
#[cfg(test)]
mod integration_tests;
#[cfg(test)]
pub(crate) mod mocks;
mod node;
mod pipeline;
pub(crate) mod planner;
pub(crate) mod query_plan;
mod request;
mod snapshot;
mod topology;
pub(crate) use context::{
PartitionRoutingRefresh, PipelineContext, RequestExecutor, ResolvedRange, TopologyProvider,
};
pub(crate) use drain::SequentialDrain;
pub(crate) use drained::DrainedLeaf;
pub(crate) use node::{PageResult, PipelineNode};
pub use pipeline::OperationPlan;
pub(crate) use pipeline::Pipeline;
pub(crate) use request::{intersect_feed_ranges, Request, RequestTarget};
pub(crate) use snapshot::{PipelineNodeState, RangedToken};
pub(crate) use topology::CachedTopologyProvider;
#[cfg(test)]
mod tests {
use super::mocks::*;
use super::*;
#[tokio::test]
async fn pipeline_forwards_pages_from_root() {
let mut pipeline =
Pipeline::new(Box::new(MockLeaf::with_pages(vec![Ok(PageResult::Page {
response: response(b"page"),
is_terminal: false,
})])));
let mut executor = NoopRequestExecutor;
let mut topology = NoopTopologyProvider;
let mut context = PipelineContext::new(&mut executor, Some(&mut topology));
let page = pipeline.next_page(&mut context).await.unwrap().unwrap();
assert_eq!(page.body_bytes(), b"page");
}
}