use futures::future::BoxFuture;
use crate::models::{CosmosOperation, CosmosResponse, FeedRange};
use super::request::RequestTarget;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum PartitionRoutingRefresh {
UseCached,
ForceRefresh,
}
pub(crate) trait RequestExecutor: Send {
fn execute_request<'a>(
&'a mut self,
operation: &'a CosmosOperation,
target: RequestTarget,
partition_routing_refresh: PartitionRoutingRefresh,
continuation: Option<String>,
) -> BoxFuture<'a, crate::error::Result<CosmosResponse>>;
}
pub(crate) trait TopologyProvider: Send {
fn resolve_ranges<'a>(
&'a mut self,
range: &'a FeedRange,
refresh: PartitionRoutingRefresh,
) -> BoxFuture<'a, crate::error::Result<Vec<ResolvedRange>>>;
}
#[derive(Debug, Clone)]
pub(crate) struct ResolvedRange {
pub partition_key_range_id: String,
pub range: FeedRange,
}
pub(crate) struct PipelineContext<'a> {
request_executor: &'a mut dyn RequestExecutor,
topology_provider: Option<&'a mut dyn TopologyProvider>,
}
impl<'a> PipelineContext<'a> {
pub(crate) fn new(
request_executor: &'a mut dyn RequestExecutor,
topology_provider: Option<&'a mut dyn TopologyProvider>,
) -> Self {
Self {
request_executor,
topology_provider,
}
}
pub(crate) async fn execute_request(
&mut self,
operation: &CosmosOperation,
target: RequestTarget,
partition_routing_refresh: PartitionRoutingRefresh,
continuation: Option<String>,
) -> crate::error::Result<CosmosResponse> {
self.request_executor
.execute_request(operation, target, partition_routing_refresh, continuation)
.await
}
pub(crate) async fn resolve_ranges(
&mut self,
range: &FeedRange,
refresh: PartitionRoutingRefresh,
) -> crate::error::Result<Vec<ResolvedRange>> {
let provider = self.topology_provider.as_deref_mut().ok_or_else(|| {
crate::error::CosmosError::builder().with_status(crate::error::CosmosStatus::CLIENT_TOPOLOGY_PROVIDER_MISSING).with_message("topology resolution requested for a plan that was not given a topology provider").build()
})?;
provider.resolve_ranges(range, refresh).await
}
}