use crate::{
diagnostics::RequestSentStatus,
driver::routing::{CosmosEndpoint, LocationEffect, UnavailablePartition, UnavailableReason},
models::{CosmosOperation, CosmosResponseHeaders, CosmosStatus, SubStatusCode},
};
use std::sync::atomic::Ordering;
use super::components::{
OperationAction, OperationRetryState, TransportOutcome, TransportResult,
BACKEND_FAILOVER_RETRY_INTERVAL,
};
fn is_ppcb_managed(operation: &CosmosOperation, retry_state: &OperationRetryState) -> bool {
retry_state.ppcb_active
&& operation
.resource_type()
.is_partitioned(operation.operation_type())
&& (operation.is_read_only() || retry_state.can_use_multiple_write_locations)
}
fn make_partition_unavailable(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
is_read: bool,
) -> UnavailablePartition {
UnavailablePartition {
partition_key_range_id: retry_state.partition_key_range_id.clone(),
region: endpoint.region().cloned(),
is_read,
is_partitioned_resource: operation
.resource_type()
.is_partitioned(operation.operation_type()),
}
}
pub(crate) fn is_region_confirming_status(status: &CosmosStatus) -> bool {
let code = status.status_code();
if code.is_success() {
return true;
}
if code == azure_core::http::StatusCode::ServiceUnavailable
|| code == azure_core::http::StatusCode::RequestTimeout
|| code == azure_core::http::StatusCode::Gone
{
return false;
}
if status.is_throttled()
&& status.sub_status() == Some(SubStatusCode::SYSTEM_RESOURCE_UNAVAILABLE)
{
return false;
}
if status.is_write_forbidden() || status.is_database_account_not_found() {
return false;
}
if status.sub_status() == Some(SubStatusCode::CLIENT_OPERATION_TIMEOUT) {
return false;
}
true
}
pub(crate) fn partition_effects_for_deferral(
is_read_only: bool,
can_use_multiple_write_locations: bool,
ppaf_write_retry_allowed: bool,
effects: Vec<LocationEffect>,
) -> (Vec<LocationEffect>, Vec<LocationEffect>) {
if is_read_only || can_use_multiple_write_locations {
return (effects, Vec::new());
}
let mut immediate = Vec::with_capacity(effects.len());
let mut deferred = Vec::new();
for effect in effects {
match effect {
LocationEffect::MarkPartitionUnavailable(_) => deferred.push(effect),
LocationEffect::MarkEndpointUnavailable { .. } if ppaf_write_retry_allowed => {
deferred.push(effect);
}
other => immediate.push(other),
}
}
(immediate, deferred)
}
pub(crate) fn evaluate_transport_result(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
result: TransportResult,
retry_state: &OperationRetryState,
) -> (OperationAction, Vec<LocationEffect>) {
match result.outcome {
outcome @ TransportOutcome::Success { .. } => (
OperationAction::Complete(Box::new(TransportResult { outcome })),
Vec::new(),
),
TransportOutcome::HttpError {
status,
cosmos_headers,
body,
request_sent,
} => evaluate_http_outcome(
operation,
endpoint,
retry_state,
status,
cosmos_headers,
body,
request_sent,
),
TransportOutcome::TransportError {
status,
error,
request_sent,
} => evaluate_transport_layer_outcome(
operation,
endpoint,
retry_state,
status,
error,
request_sent,
),
TransportOutcome::DeadlineExceeded { request_sent } => {
evaluate_deadline_exceeded_outcome(request_sent)
}
}
}
#[derive(Debug, Default)]
pub(crate) struct HedgeLegEvaluation {
pub(crate) effects: Vec<LocationEffect>,
pub(crate) observed_session_unavailable: bool,
}
pub(crate) fn evaluate_hedge_leg_effects(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
result: &TransportResult,
) -> HedgeLegEvaluation {
let mut eval = HedgeLegEvaluation::default();
match &result.outcome {
TransportOutcome::Success { .. } => {}
TransportOutcome::HttpError {
status,
request_sent,
..
} => {
if status.is_read_session_not_available()
&& retry_state.can_retry_session()
&& retry_state.is_dataplane
&& !retry_state.can_use_multiple_write_locations
&& retry_state.session_token_retry_count == 0
&& !retry_state.hub_region_processing_only
{
eval.observed_session_unavailable = true;
}
if let Some((_action, effects)) =
try_handle_write_forbidden(operation, endpoint, retry_state, status)
{
eval.effects = effects;
} else if let Some((_action, effects)) =
try_handle_database_account_not_found(operation, endpoint, retry_state, status)
{
eval.effects = effects;
} else if let Some((_action, effects)) = try_handle_retry_trigger_group(
operation,
endpoint,
retry_state,
status,
*request_sent,
) {
eval.effects = effects;
} else if let Some((_action, effects)) =
try_handle_server_error(operation, endpoint, retry_state, status)
{
eval.effects = effects;
}
}
TransportOutcome::TransportError { request_sent, .. } => {
if !request_sent.definitely_not_sent() {
eval.effects.push(LocationEffect::MarkPartitionUnavailable(
make_partition_unavailable(
operation,
endpoint,
retry_state,
operation.is_read_only(),
),
));
if !is_ppcb_managed(operation, retry_state) {
eval.effects.push(LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint.clone(),
reason: UnavailableReason::TransportError,
});
}
}
}
TransportOutcome::DeadlineExceeded { .. } => {
}
}
eval
}
#[allow(clippy::too_many_arguments)]
fn evaluate_http_outcome(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
status: CosmosStatus,
cosmos_headers: CosmosResponseHeaders,
body: Vec<u8>,
request_sent: RequestSentStatus,
) -> (OperationAction, Vec<LocationEffect>) {
if let Some(result) = try_handle_write_forbidden(operation, endpoint, retry_state, &status) {
return result;
}
if let Some(result) =
try_handle_database_account_not_found(operation, endpoint, retry_state, &status)
{
return result;
}
if let Some(result) =
try_handle_read_session_not_available(retry_state, &status, &cosmos_headers, &body)
{
return result;
}
if let Some(result) =
try_handle_retry_trigger_group(operation, endpoint, retry_state, &status, request_sent)
{
return result;
}
if let Some(result) = try_handle_server_error(operation, endpoint, retry_state, &status) {
return result;
}
(
OperationAction::Abort {
error: build_service_error(&status, &cosmos_headers, &body),
},
Vec::new(),
)
}
fn try_handle_write_forbidden(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
status: &CosmosStatus,
) -> Option<(OperationAction, Vec<LocationEffect>)> {
if !status.is_write_forbidden() {
return None;
}
let (new_state, delay) = if retry_state.can_use_multiple_write_locations {
if !retry_state.can_retry_backend_failover() {
return None;
}
(
retry_state.clone().advance_backend_failover(),
Some(BACKEND_FAILOVER_RETRY_INTERVAL),
)
} else {
if !retry_state.can_retry_failover() {
return None;
}
(retry_state.clone().advance_failover(), None)
};
let mut effects = vec![
LocationEffect::RefreshAccountProperties,
LocationEffect::MarkPartitionUnavailable(make_partition_unavailable(
operation,
endpoint,
retry_state,
false,
)),
];
if !is_ppcb_managed(operation, retry_state) {
effects.push(LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint.clone(),
reason: UnavailableReason::WriteForbidden,
});
}
Some((OperationAction::FailoverRetry { new_state, delay }, effects))
}
fn try_handle_database_account_not_found(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
status: &CosmosStatus,
) -> Option<(OperationAction, Vec<LocationEffect>)> {
if !status.is_database_account_not_found() {
return None;
}
if !retry_state.can_retry_backend_failover() {
return None;
}
let new_state = retry_state.clone().advance_backend_failover();
let delay = Some(BACKEND_FAILOVER_RETRY_INTERVAL);
let mut effects = vec![
LocationEffect::RefreshAccountProperties,
LocationEffect::MarkPartitionUnavailable(make_partition_unavailable(
operation,
endpoint,
retry_state,
operation.is_read_only(),
)),
];
if !is_ppcb_managed(operation, retry_state) {
effects.push(LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint.clone(),
reason: UnavailableReason::DatabaseAccountNotFound,
});
}
Some((OperationAction::FailoverRetry { new_state, delay }, effects))
}
fn try_handle_read_session_not_available(
retry_state: &OperationRetryState,
status: &CosmosStatus,
cosmos_headers: &CosmosResponseHeaders,
body: &[u8],
) -> Option<(OperationAction, Vec<LocationEffect>)> {
if !(status.is_read_session_not_available() && retry_state.can_retry_session()) {
return None;
}
if !retry_state.can_use_multiple_write_locations && retry_state.session_token_retry_count >= 2 {
return Some((
OperationAction::Abort {
error: build_service_error(status, cosmos_headers, body),
},
Vec::new(),
));
}
Some((
OperationAction::SessionRetry {
new_state: build_session_retry_state(retry_state),
},
Vec::new(),
))
}
fn build_session_retry_state(retry_state: &OperationRetryState) -> OperationRetryState {
let mut new_state = retry_state.clone().advance_session_retry();
if retry_state.is_dataplane
&& !retry_state.can_use_multiple_write_locations
&& retry_state.session_token_retry_count == 0
&& !retry_state.hub_region_processing_only
{
new_state.hub_region_processing_only = true;
if let Some(shared) = new_state.shared_hub_region_latch.as_ref() {
shared.store(true, Ordering::Release);
}
}
new_state
}
fn try_handle_retry_trigger_group(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
status: &CosmosStatus,
request_sent: RequestSentStatus,
) -> Option<(OperationAction, Vec<LocationEffect>)> {
let is_system_resource_unavailable = status.is_throttled()
&& status.sub_status() == Some(SubStatusCode::SYSTEM_RESOURCE_UNAVAILABLE);
let is_service_unavailable =
status.status_code() == azure_core::http::StatusCode::ServiceUnavailable;
let is_gone = status.is_gone() && !status.is_partition_topology_change();
let is_request_timeout = status.status_code() == azure_core::http::StatusCode::RequestTimeout;
let in_trigger_group =
is_system_resource_unavailable || is_service_unavailable || is_gone || is_request_timeout;
if !(in_trigger_group && retry_state.can_retry_failover()) {
return None;
}
if request_sent.definitely_not_sent() {
return Some((
OperationAction::FailoverRetry {
new_state: retry_state.clone().advance_failover(),
delay: None,
},
Vec::new(),
));
}
let unavailable_reason = if is_request_timeout {
UnavailableReason::RequestTimeout
} else {
UnavailableReason::ServiceUnavailable
};
let mut effects = vec![LocationEffect::MarkPartitionUnavailable(
make_partition_unavailable(operation, endpoint, retry_state, operation.is_read_only()),
)];
if !is_ppcb_managed(operation, retry_state) {
effects.push(LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint.clone(),
reason: unavailable_reason,
});
}
Some((
OperationAction::FailoverRetry {
new_state: retry_state.clone().advance_failover(),
delay: None,
},
effects,
))
}
fn try_handle_server_error(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
status: &CosmosStatus,
) -> Option<(OperationAction, Vec<LocationEffect>)> {
let status_code = status.status_code();
let is_eligible_status = status_code.is_server_error()
|| status_code == azure_core::http::StatusCode::RequestTimeout;
if !(is_eligible_status && retry_state.can_retry_failover()) {
return None;
}
let mut effects = vec![LocationEffect::MarkPartitionUnavailable(
make_partition_unavailable(operation, endpoint, retry_state, operation.is_read_only()),
)];
if !is_ppcb_managed(operation, retry_state) {
effects.push(LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint.clone(),
reason: UnavailableReason::InternalServerError,
});
}
Some((
OperationAction::FailoverRetry {
new_state: retry_state.clone().advance_failover(),
delay: None,
},
effects,
))
}
fn evaluate_transport_layer_outcome(
operation: &CosmosOperation,
endpoint: &CosmosEndpoint,
retry_state: &OperationRetryState,
status: CosmosStatus,
error: crate::error::CosmosError,
request_sent: RequestSentStatus,
) -> (OperationAction, Vec<LocationEffect>) {
if request_sent.definitely_not_sent() {
let effects = vec![LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint.clone(),
reason: UnavailableReason::TransportError,
}];
if retry_state.can_retry_failover() {
return (
OperationAction::FailoverRetry {
new_state: retry_state.clone().advance_failover(),
delay: None,
},
effects,
);
}
return (
OperationAction::Abort {
error: build_transport_error(&status, error),
},
effects,
);
}
let effects = vec![LocationEffect::MarkPartitionUnavailable(
make_partition_unavailable(operation, endpoint, retry_state, operation.is_read_only()),
)];
if retry_state.can_retry_failover() {
return (
OperationAction::FailoverRetry {
new_state: retry_state.clone().advance_failover(),
delay: None,
},
effects,
);
}
(
OperationAction::Abort {
error: build_transport_error(&status, error),
},
effects,
)
}
fn evaluate_deadline_exceeded_outcome(
request_sent: RequestSentStatus,
) -> (OperationAction, Vec<LocationEffect>) {
let message: &'static str = if request_sent.definitely_not_sent() {
"end-to-end operation timeout exceeded before request was sent"
} else {
"end-to-end operation timeout exceeded"
};
let cosmos_err = crate::error::CosmosError::builder()
.with_status(CosmosStatus::from_parts(
azure_core::http::StatusCode::RequestTimeout,
Some(crate::models::SubStatusCode::CLIENT_OPERATION_TIMEOUT),
))
.with_message(message)
.build();
(OperationAction::Abort { error: cosmos_err }, Vec::new())
}
fn service_error_message(status: &CosmosStatus) -> String {
let sub_status_str = match status.sub_status() {
Some(s) => format!("/{}", s.value()),
None => String::new(),
};
format!(
"Cosmos DB returned HTTP {}{}: {}",
u16::from(status.status_code()),
sub_status_str,
status.name().unwrap_or("Unknown"),
)
}
pub(crate) fn build_service_error(
status: &CosmosStatus,
cosmos_headers: &CosmosResponseHeaders,
body: &[u8],
) -> crate::error::CosmosError {
let effective_status = synthesize_cross_partition_query_status(*status, body);
crate::error::CosmosError::builder()
.with_status(effective_status)
.with_message(service_error_message(&effective_status))
.with_response_parts(crate::models::CosmosResponsePayload::new(
body.to_vec(),
cosmos_headers.clone(),
))
.build()
}
fn synthesize_cross_partition_query_status(status: CosmosStatus, body: &[u8]) -> CosmosStatus {
use azure_core::http::StatusCode;
if status.status_code() != StatusCode::BadRequest || status.sub_status().is_some() {
return status;
}
let Ok(text) = std::str::from_utf8(body) else {
return status;
};
if text.contains("unsupported features") && text.contains("Upgrade your SDK") {
crate::error::CosmosStatus::CROSS_PARTITION_QUERY_NOT_SERVABLE
} else {
status
}
}
fn build_transport_error(
status: &CosmosStatus,
error: crate::error::CosmosError,
) -> crate::error::CosmosError {
let status_code = status.status_code();
let name = status.name().unwrap_or("Unknown");
let sub_status_str = match status.sub_status() {
Some(s) => format!("/{}", s.value()),
None => String::new(),
};
let detail_summary = crate::driver::error_chain_summary(&error);
let message = format!(
"Cosmos DB transport failure HTTP {}{}: {}. Details: {}",
u16::from(status_code),
sub_status_str,
name,
detail_summary,
);
let mut b = crate::error::CosmosError::builder()
.with_status(*status)
.with_message(message)
.with_arc_source(std::sync::Arc::new(error.clone()));
if let Some(diag) = error.diagnostics() {
b = b.with_diagnostics(diag);
}
b.build()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
diagnostics::RequestSentStatus,
models::{
AccountReference, ContainerProperties, ContainerReference, CosmosOperation,
CosmosResponseHeaders, CosmosStatus, DatabaseReference, ItemReference, PartitionKey,
PartitionKeyDefinition, SystemProperties,
},
};
use azure_core::http::StatusCode;
fn make_create_item_operation() -> CosmosOperation {
let account = AccountReference::with_master_key(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
"dGVzdA==",
);
let pk_def: PartitionKeyDefinition = serde_json::from_str(r#"{"paths":["/pk"]}"#).unwrap();
let props = ContainerProperties {
id: "testcontainer".into(),
partition_key: pk_def,
system_properties: SystemProperties::default(),
};
let container = ContainerReference::new(
account,
"testdb",
"testdb_rid",
"testcontainer",
"testcontainer_rid",
&props,
);
let item = ItemReference::from_name(&container, PartitionKey::from("pk1"), "doc1");
CosmosOperation::create_item(item).with_body(b"{}".to_vec())
}
fn make_read_operation() -> CosmosOperation {
let account = AccountReference::with_master_key(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
"dGVzdA==", );
let db_ref = DatabaseReference::from_name(account, "testdb".to_owned());
CosmosOperation::read_database(db_ref)
}
fn make_create_operation() -> CosmosOperation {
let account = AccountReference::with_master_key(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
"dGVzdA==",
);
CosmosOperation::create_database(account)
}
fn make_success_result() -> TransportResult {
TransportResult {
outcome: TransportOutcome::Success {
status: CosmosStatus::new(StatusCode::Ok),
cosmos_headers: CosmosResponseHeaders::default(),
body: b"{}".to_vec(),
},
}
}
fn make_transport_error(sent: RequestSentStatus) -> TransportResult {
TransportResult {
outcome: TransportOutcome::TransportError {
status: CosmosStatus::TRANSPORT_GENERATED_503,
error: crate::error::CosmosError::builder()
.with_status(CosmosStatus::TRANSPORT_GENERATED_503)
.with_message("connection refused")
.build(),
request_sent: sent,
},
}
}
fn make_http_error(status_code: StatusCode) -> TransportResult {
TransportResult {
outcome: TransportOutcome::HttpError {
status: CosmosStatus::new(status_code),
cosmos_headers: CosmosResponseHeaders::default(),
body: vec![],
request_sent: RequestSentStatus::Sent,
},
}
}
fn make_http_error_status(status: CosmosStatus) -> TransportResult {
TransportResult {
outcome: TransportOutcome::HttpError {
status,
cosmos_headers: CosmosResponseHeaders::default(),
body: vec![],
request_sent: RequestSentStatus::Sent,
},
}
}
#[test]
fn success_completes() {
let op = make_read_operation();
let result = make_success_result();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::Complete(_)));
assert!(effects.is_empty());
}
#[test]
fn transport_error_not_sent_retries() {
let op = make_create_operation();
let result = make_transport_error(RequestSentStatus::NotSent);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
}
#[test]
fn transport_error_sent_idempotent_retries() {
let op = make_read_operation();
let result = make_transport_error(RequestSentStatus::Sent);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn transport_error_sent_non_idempotent_retries() {
let op = make_create_operation();
let result = make_transport_error(RequestSentStatus::Sent);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn build_transport_error_forwards_inner_diagnostics() {
let diag: std::sync::Arc<crate::diagnostics::DiagnosticsContext> = std::sync::Arc::new(
crate::diagnostics::DiagnosticsContextBuilder::new(
crate::models::ActivityId::new_uuid(),
std::sync::Arc::new(crate::options::DiagnosticsOptions::default()),
)
.complete(),
);
let inner = crate::error::CosmosError::builder()
.with_status(CosmosStatus::TRANSPORT_GENERATED_503)
.with_message("inner transport failure")
.with_diagnostics(std::sync::Arc::clone(&diag))
.build();
let outer = build_transport_error(&CosmosStatus::TRANSPORT_GENERATED_503, inner);
let outer_diag = outer
.diagnostics()
.expect("outer error must inherit inner diagnostics");
assert!(
std::sync::Arc::ptr_eq(&outer_diag, &diag),
"outer diagnostics must be the same Arc as the inner's"
);
}
#[test]
fn transport_abort_error_includes_status_kind_and_details() {
let op = make_create_operation();
let result = TransportResult {
outcome: TransportOutcome::TransportError {
status: CosmosStatus::TRANSPORT_GENERATED_503,
error: crate::error::CosmosError::builder()
.with_status(CosmosStatus::TRANSPORT_GENERATED_503)
.with_message("failed to execute `reqwest` request")
.with_source(std::io::Error::new(
std::io::ErrorKind::BrokenPipe,
"socket reset",
))
.build(),
request_sent: RequestSentStatus::Unknown,
},
};
let state = OperationRetryState::initial(0, false, Vec::new(), 0, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(&op, &endpoint, result, &state);
match action {
OperationAction::Abort { error } => {
assert_eq!(error.status(), CosmosStatus::TRANSPORT_GENERATED_503);
let text = error.to_string();
assert!(text.contains("HTTP 503/20003"));
assert!(text.contains("TransportGenerated503"));
assert!(text.contains("failed to execute `reqwest` request"));
assert!(text.contains("socket reset"));
}
other => panic!("expected abort, got {other:?}"),
}
}
#[test]
fn transport_error_over_budget_aborts() {
let op = make_read_operation();
let result = make_transport_error(RequestSentStatus::NotSent);
let state = OperationRetryState {
location: crate::driver::routing::LocationIndex::initial(0),
failover_retry_count: 1,
session_token_retry_count: 0,
backend_failover_retry_count: 0,
max_failover_retries: 1,
max_backend_failover_retries: 120,
max_session_retries: 1,
can_use_multiple_write_locations: false,
is_dataplane: false,
hub_region_processing_only: false,
shared_hub_region_latch: None,
excluded_regions: Vec::new(),
session_retry_routing:
crate::driver::pipeline::components::SessionRetryRouting::PreferredEndpoints,
partition_key_range_id: None,
ppaf_write_retry_allowed: false,
ppcb_active: false,
pending_write_effects: Vec::new(),
hedge_already_fired: false,
};
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(&op, &endpoint, result, &state);
match action {
OperationAction::Abort { error } => {
assert_eq!(error.status(), CosmosStatus::TRANSPORT_GENERATED_503);
}
other => panic!("expected abort, got {other:?}"),
}
}
#[test]
fn http_error_aborts() {
let op = make_read_operation();
let result = make_http_error(StatusCode::BadRequest);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::Abort { .. }));
}
#[test]
fn partition_topology_gone_aborts_for_dataflow_handling() {
let op = make_read_operation();
let result = make_http_error_status(
CosmosStatus::new(StatusCode::Gone)
.with_sub_status(SubStatusCode::PARTITION_KEY_RANGE_GONE.value()),
);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
match action {
OperationAction::Abort { error, .. } => {
assert_eq!(
error.status(),
CosmosStatus::new(StatusCode::Gone)
.with_sub_status(SubStatusCode::PARTITION_KEY_RANGE_GONE.value())
);
}
other => panic!("expected abort, got {other:?}"),
}
assert!(effects.is_empty());
}
#[test]
fn non_topology_gone_still_retries() {
let op = make_read_operation();
let result = make_http_error_status(
CosmosStatus::new(StatusCode::Gone)
.with_sub_status(SubStatusCode::NAME_CACHE_STALE.value()),
);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn write_forbidden_triggers_failover_and_refresh_effect() {
let op = make_create_operation();
let result = TransportResult {
outcome: TransportOutcome::HttpError {
status: CosmosStatus::WRITE_FORBIDDEN,
cosmos_headers: CosmosResponseHeaders::default(),
body: vec![],
request_sent: RequestSentStatus::Sent,
},
};
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::RefreshAccountProperties)));
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })),
"non-PPCB 403/3 must mark the endpoint unavailable"
);
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))),
"403/3 must always mark the partition unavailable"
);
}
#[test]
fn write_forbidden_when_ppcb_managed_skips_endpoint_mark() {
let op = make_create_item_operation();
let mut state =
OperationRetryState::initial(0, true , Vec::new(), 3, 1);
state.ppcb_active = true;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::WRITE_FORBIDDEN),
&state,
);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::RefreshAccountProperties)),
"403/3 must still refresh account properties"
);
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))),
"PPCB-managed 403/3 must still mark the partition unavailable"
);
assert!(
!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })),
"PPCB-managed 403/3 must NOT mark the endpoint unavailable; \
per-partition counter drives failover"
);
}
fn http_error_status(status: CosmosStatus) -> TransportResult {
TransportResult {
outcome: TransportOutcome::HttpError {
status,
cosmos_headers: CosmosResponseHeaders::default(),
body: vec![],
request_sent: RequestSentStatus::Sent,
},
}
}
#[test]
fn database_account_not_found_on_write_emits_refresh_mark_endpoint_and_failover() {
let op = make_create_operation();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::DATABASE_ACCOUNT_NOT_FOUND),
&state,
);
assert!(
matches!(
action,
OperationAction::FailoverRetry { delay: Some(_), .. }
),
"expected FailoverRetry with backend-failover delay, got {:?}",
action
);
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::RefreshAccountProperties)),
"missing RefreshAccountProperties effect; effects={:?}",
effects
);
assert!(
effects.iter().any(|e| matches!(
e,
LocationEffect::MarkEndpointUnavailable {
reason: UnavailableReason::DatabaseAccountNotFound,
..
}
)),
"missing MarkEndpointUnavailable{{DatabaseAccountNotFound}}; effects={:?}",
effects
);
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))),
"missing MarkPartitionUnavailable; effects={:?}",
effects
);
}
#[test]
fn database_account_not_found_on_read_uses_op_read_only_for_partition_mark() {
let op = make_read_operation();
let state = OperationRetryState::initial(0, true, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::DATABASE_ACCOUNT_NOT_FOUND),
&state,
);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
let mark = effects.iter().find_map(|e| match e {
LocationEffect::MarkPartitionUnavailable(p) => Some(p),
_ => None,
});
let mark = mark.expect("expected MarkPartitionUnavailable for 1008 on read");
assert!(
mark.is_read,
"1008 on a read op must mark the partition with is_read=true so PPCB \
credits the read-direction failure counter, not the write counter"
);
}
#[test]
fn database_account_not_found_aborts_when_backend_failover_budget_exhausted() {
let op = make_create_operation();
let mut state = OperationRetryState::initial(0, true, Vec::new(), 3, 1);
state.backend_failover_retry_count = state.max_backend_failover_retries;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::DATABASE_ACCOUNT_NOT_FOUND),
&state,
);
match action {
OperationAction::Abort { error } => {
assert_eq!(
error.status(),
CosmosStatus::DATABASE_ACCOUNT_NOT_FOUND,
"1008 exhausted-budget bubble-up must surface the original status unchanged"
);
}
other => panic!(
"expected Abort once backend-failover budget is exhausted, got {:?}",
other
),
}
}
#[test]
fn write_forbidden_aborts_when_backend_failover_budget_exhausted() {
let op = make_create_operation();
let mut state = OperationRetryState::initial(0, true, Vec::new(), 3, 1);
state.backend_failover_retry_count = state.max_backend_failover_retries;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::WRITE_FORBIDDEN),
&state,
);
match action {
OperationAction::Abort { error } => {
assert_eq!(
error.status(),
CosmosStatus::WRITE_FORBIDDEN,
"403/3 exhausted-budget bubble-up must surface the original status unchanged"
);
}
other => panic!("expected Abort, got {:?}", other),
}
}
#[test]
fn database_account_not_found_does_not_consume_generic_failover_budget() {
let op = make_create_operation();
let state = OperationRetryState::initial(0, true, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::DATABASE_ACCOUNT_NOT_FOUND),
&state,
);
match action {
OperationAction::FailoverRetry { new_state, delay } => {
assert_eq!(new_state.failover_retry_count, 0);
assert_eq!(new_state.backend_failover_retry_count, 1);
assert_eq!(
delay,
Some(BACKEND_FAILOVER_RETRY_INTERVAL),
"multi-write 1008 must pace retries with BACKEND_FAILOVER_RETRY_INTERVAL"
);
}
other => panic!("expected FailoverRetry, got {:?}", other),
}
}
#[test]
fn write_forbidden_does_not_consume_generic_failover_budget() {
let op = make_create_operation();
let state = OperationRetryState::initial(0, true, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::WRITE_FORBIDDEN),
&state,
);
match action {
OperationAction::FailoverRetry { new_state, delay } => {
assert_eq!(new_state.failover_retry_count, 0);
assert_eq!(new_state.backend_failover_retry_count, 1);
assert_eq!(
delay,
Some(BACKEND_FAILOVER_RETRY_INTERVAL),
"multi-write 403/3 must pace retries with BACKEND_FAILOVER_RETRY_INTERVAL"
);
}
other => panic!("expected FailoverRetry, got {:?}", other),
}
}
#[test]
fn write_forbidden_on_single_write_uses_generic_failover_budget() {
let op = make_create_operation();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::WRITE_FORBIDDEN),
&state,
);
match action {
OperationAction::FailoverRetry { new_state, delay } => {
assert_eq!(new_state.failover_retry_count, 1);
assert_eq!(new_state.backend_failover_retry_count, 0);
assert_eq!(
delay, None,
"single-write 403/3 uses the generic budget and must not pace retries"
);
}
other => panic!("expected FailoverRetry, got {:?}", other),
}
}
#[test]
fn write_forbidden_on_single_write_aborts_when_generic_budget_exhausted() {
let op = make_create_operation();
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
state.failover_retry_count = state.max_failover_retries;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::WRITE_FORBIDDEN),
&state,
);
assert!(
matches!(action, OperationAction::Abort { .. }),
"single-write 403/3 must abort once generic failover budget is exhausted, got {:?}",
action
);
}
#[test]
fn database_account_not_found_on_single_write_uses_backend_failover_budget() {
let op = make_create_operation();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, _effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::DATABASE_ACCOUNT_NOT_FOUND),
&state,
);
match action {
OperationAction::FailoverRetry { new_state, delay } => {
assert_eq!(new_state.failover_retry_count, 0);
assert_eq!(new_state.backend_failover_retry_count, 1);
assert_eq!(
delay,
Some(BACKEND_FAILOVER_RETRY_INTERVAL),
"single-write 1008 must pace retries with BACKEND_FAILOVER_RETRY_INTERVAL"
);
}
other => panic!("expected FailoverRetry, got {:?}", other),
}
}
#[test]
fn database_account_not_found_when_ppcb_managed_skips_endpoint_mark() {
let op = make_create_item_operation();
let mut state = OperationRetryState::initial(0, true, Vec::new(), 3, 1);
state.ppcb_active = true;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(
&op,
&endpoint,
http_error_status(CosmosStatus::DATABASE_ACCOUNT_NOT_FOUND),
&state,
);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::RefreshAccountProperties)),
"1008 must still refresh account properties"
);
assert!(
effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))),
"PPCB-managed 1008 must still mark the partition unavailable"
);
assert!(
!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })),
"PPCB-managed 1008 must NOT mark the endpoint unavailable; \
per-partition counter drives failover"
);
}
#[test]
fn read_session_not_available_triggers_session_retry() {
let op = make_read_operation();
let result = TransportResult {
outcome: TransportOutcome::HttpError {
status: CosmosStatus::READ_SESSION_NOT_AVAILABLE,
cosmos_headers: CosmosResponseHeaders::default(),
body: vec![],
request_sent: RequestSentStatus::Sent,
},
};
let state = OperationRetryState::initial(0, true, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::SessionRetry { .. }));
assert!(effects.is_empty());
}
#[test]
fn service_unavailable_marks_endpoint_unavailable() {
let op = make_read_operation();
let result = make_http_error(StatusCode::ServiceUnavailable);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn throttle_substatus_gates_hedge_eligibility() {
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let throttle_result = |sub: SubStatusCode| TransportResult {
outcome: TransportOutcome::HttpError {
status: CosmosStatus::from_parts(StatusCode::TooManyRequests, Some(sub)),
cosmos_headers: CosmosResponseHeaders::default(),
body: vec![],
request_sent: RequestSentStatus::Sent,
},
};
let op = make_read_operation();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let (action, _) = evaluate_transport_result(
&op,
&endpoint,
throttle_result(SubStatusCode::SYSTEM_RESOURCE_UNAVAILABLE),
&state,
);
assert!(
matches!(action, OperationAction::FailoverRetry { .. }),
"429/3092 SystemResourceUnavailable must be failover-eligible; got {action:?}",
);
for sub in [
SubStatusCode::RU_BUDGET_EXCEEDED, SubStatusCode::RU_BUDGET_EXCEEDED_FOR_MASTER, SubStatusCode::HOT_PARTITION_KEY_THROTTLED, ] {
let op = make_read_operation();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let (action, _) =
evaluate_transport_result(&op, &endpoint, throttle_result(sub), &state);
assert!(
!matches!(
action,
OperationAction::FailoverRetry { .. } | OperationAction::SessionRetry { .. }
),
"429/{sub:?} must not become a region-changing retry; \
got {action:?}",
);
}
}
#[test]
fn service_unavailable_non_idempotent_write_retries() {
let op = make_create_operation();
let result = make_http_error(StatusCode::ServiceUnavailable);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn service_unavailable_non_idempotent_retries_with_ppaf() {
let op = make_create_operation();
let result = make_http_error(StatusCode::ServiceUnavailable);
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
state.ppaf_write_retry_allowed = true;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn transport_error_non_idempotent_retries_with_ppaf() {
let op = make_create_operation();
let result = make_transport_error(RequestSentStatus::Sent);
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
state.ppaf_write_retry_allowed = true;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn deadline_exceeded_aborts_with_timeout_status() {
let op = make_read_operation();
let result = TransportResult {
outcome: TransportOutcome::DeadlineExceeded {
request_sent: RequestSentStatus::Unknown,
},
};
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
match action {
OperationAction::Abort { error } => {
let status = error.status();
assert_eq!(status.status_code(), StatusCode::RequestTimeout);
assert_eq!(
status.sub_status(),
Some(SubStatusCode::CLIENT_OPERATION_TIMEOUT)
);
}
_ => panic!("expected timeout to abort"),
}
assert!(effects.is_empty());
}
#[test]
fn internal_server_error_on_read_fails_over() {
let op = make_read_operation();
let result = make_http_error(StatusCode::InternalServerError);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn internal_server_error_on_read_marks_partition_unavailable() {
let op = make_read_operation();
let result = make_http_error(StatusCode::InternalServerError);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
}
#[test]
fn transport_error_not_sent_marks_endpoint_only() {
let op = make_read_operation();
let result = make_transport_error(RequestSentStatus::NotSent);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
}
#[test]
fn transport_error_not_sent_with_ppcb_still_marks_endpoint() {
let op = make_read_operation();
let result = make_transport_error(RequestSentStatus::NotSent);
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
state.ppcb_active = true;
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
}
#[test]
fn transport_error_unknown_sent_status_marks_partition_only() {
let op = make_read_operation();
let result = make_transport_error(RequestSentStatus::Unknown);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn transport_error_not_sent_over_budget_aborts_with_marks() {
let op = make_read_operation();
let result = make_transport_error(RequestSentStatus::NotSent);
let state = OperationRetryState {
location: crate::driver::routing::LocationIndex::initial(0),
failover_retry_count: 1,
session_token_retry_count: 0,
backend_failover_retry_count: 0,
max_failover_retries: 1,
max_backend_failover_retries: 120,
max_session_retries: 1,
can_use_multiple_write_locations: false,
is_dataplane: false,
hub_region_processing_only: false,
shared_hub_region_latch: None,
excluded_regions: Vec::new(),
session_retry_routing:
crate::driver::pipeline::components::SessionRetryRouting::PreferredEndpoints,
partition_key_range_id: None,
ppaf_write_retry_allowed: false,
ppcb_active: false,
pending_write_effects: Vec::new(),
hedge_already_fired: false,
};
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::Abort { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
assert!(!effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
}
#[test]
fn request_timeout_from_server_marks_partition_and_endpoint_unavailable() {
let op = make_read_operation();
let result = make_http_error(StatusCode::RequestTimeout);
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 1);
let endpoint = CosmosEndpoint::global(
url::Url::parse("https://test.documents.azure.com:443/").unwrap(),
);
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, &state);
assert!(matches!(action, OperationAction::FailoverRetry { .. }));
assert!(effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(effects.iter().any(|e| matches!(
e,
LocationEffect::MarkEndpointUnavailable {
reason: UnavailableReason::RequestTimeout,
..
}
)));
}
fn status_with_substatus(code: StatusCode, sub: SubStatusCode) -> CosmosStatus {
CosmosStatus::from_parts(code, Some(sub))
}
#[test]
fn region_confirming_true_for_2xx() {
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::Ok
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::Created
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::Accepted
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::NoContent
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::from(207u16)
)));
}
#[test]
fn region_confirming_true_for_definitive_4xx() {
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::Conflict
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::PreconditionFailed
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::PayloadTooLarge
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::BadRequest
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::Unauthorized
)));
assert!(is_region_confirming_status(&CosmosStatus::new(
StatusCode::NotFound
)));
assert!(is_region_confirming_status(&status_with_substatus(
StatusCode::NotFound,
SubStatusCode::from(0u16)
)));
}
#[test]
fn region_confirming_false_for_retry_trigger_statuses() {
assert!(!is_region_confirming_status(&CosmosStatus::new(
StatusCode::ServiceUnavailable
)));
assert!(!is_region_confirming_status(&CosmosStatus::new(
StatusCode::RequestTimeout
)));
assert!(!is_region_confirming_status(&CosmosStatus::new(
StatusCode::Gone
)));
assert!(!is_region_confirming_status(&status_with_substatus(
StatusCode::TooManyRequests,
SubStatusCode::SYSTEM_RESOURCE_UNAVAILABLE
)));
assert!(!is_region_confirming_status(&status_with_substatus(
StatusCode::Forbidden,
SubStatusCode::WRITE_FORBIDDEN
)));
assert!(!is_region_confirming_status(&status_with_substatus(
StatusCode::Forbidden,
SubStatusCode::DATABASE_ACCOUNT_NOT_FOUND
)));
assert!(!is_region_confirming_status(&status_with_substatus(
StatusCode::Gone,
SubStatusCode::COMPLETING_PARTITION_MIGRATION
)));
}
#[test]
fn region_confirming_false_for_client_synthesized_timeout() {
assert!(!is_region_confirming_status(&status_with_substatus(
StatusCode::RequestTimeout,
SubStatusCode::CLIENT_OPERATION_TIMEOUT
)));
}
fn endpoint_for_test() -> CosmosEndpoint {
CosmosEndpoint::global(url::Url::parse("https://test.documents.azure.com:443/").unwrap())
}
#[test]
fn deferral_passes_all_effects_through_for_reads() {
let effects = vec![
LocationEffect::MarkPartitionUnavailable(UnavailablePartition {
partition_key_range_id: None,
region: None,
is_read: true,
is_partitioned_resource: true,
}),
LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint_for_test(),
reason: UnavailableReason::ServiceUnavailable,
},
LocationEffect::RefreshAccountProperties,
];
let (immediate, deferred) = partition_effects_for_deferral(true, false, false, effects);
assert_eq!(immediate.len(), 3);
assert!(deferred.is_empty());
}
#[test]
fn deferral_extracts_partition_marks_for_writes() {
let effects = vec![
LocationEffect::MarkPartitionUnavailable(UnavailablePartition {
partition_key_range_id: None,
region: None,
is_read: false,
is_partitioned_resource: true,
}),
LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint_for_test(),
reason: UnavailableReason::ServiceUnavailable,
},
LocationEffect::RefreshAccountProperties,
];
let (immediate, deferred) = partition_effects_for_deferral(false, false, false, effects);
assert_eq!(immediate.len(), 2);
assert!(immediate
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
assert!(immediate
.iter()
.any(|e| matches!(e, LocationEffect::RefreshAccountProperties)));
assert_eq!(deferred.len(), 1);
assert!(matches!(
deferred[0],
LocationEffect::MarkPartitionUnavailable(_)
));
}
#[test]
fn deferral_with_no_partition_marks_returns_empty_deferred() {
let effects = vec![
LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint_for_test(),
reason: UnavailableReason::ServiceUnavailable,
},
LocationEffect::RefreshAccountProperties,
];
let (immediate, deferred) = partition_effects_for_deferral(false, false, false, effects);
assert_eq!(immediate.len(), 2);
assert!(deferred.is_empty());
}
#[test]
fn deferral_defers_endpoint_mark_for_ppaf_single_master_writes() {
let effects = vec![
LocationEffect::MarkPartitionUnavailable(UnavailablePartition {
partition_key_range_id: None,
region: None,
is_read: false,
is_partitioned_resource: true,
}),
LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint_for_test(),
reason: UnavailableReason::TransportError,
},
LocationEffect::RefreshAccountProperties,
];
let (immediate, deferred) = partition_effects_for_deferral(false, false, true, effects);
assert_eq!(immediate.len(), 1);
assert!(matches!(
immediate[0],
LocationEffect::RefreshAccountProperties
));
assert_eq!(deferred.len(), 2);
assert!(deferred
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(deferred
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn deferral_passes_all_effects_through_for_multi_master_writes() {
let effects = vec![
LocationEffect::MarkPartitionUnavailable(UnavailablePartition {
partition_key_range_id: None,
region: None,
is_read: false,
is_partitioned_resource: true,
}),
LocationEffect::MarkEndpointUnavailable {
endpoint: endpoint_for_test(),
reason: UnavailableReason::ServiceUnavailable,
},
LocationEffect::RefreshAccountProperties,
];
let (immediate, deferred) = partition_effects_for_deferral(false, true, false, effects);
assert_eq!(immediate.len(), 3);
assert!(deferred.is_empty());
}
fn make_read_session_not_available_result() -> TransportResult {
TransportResult {
outcome: TransportOutcome::HttpError {
status: CosmosStatus::READ_SESSION_NOT_AVAILABLE,
cosmos_headers: CosmosResponseHeaders::default(),
body: vec![],
request_sent: RequestSentStatus::Sent,
},
}
}
fn test_endpoint() -> CosmosEndpoint {
CosmosEndpoint::global(url::Url::parse("https://test.documents.azure.com:443/").unwrap())
}
fn session_retry_state_for_1002(state: &OperationRetryState) -> OperationRetryState {
let op = make_read_operation();
let endpoint = test_endpoint();
let result = make_read_session_not_available_result();
let (action, effects) = evaluate_transport_result(&op, &endpoint, result, state);
assert!(
effects.is_empty(),
"1002 should not emit location effects, got {effects:?}",
);
match action {
OperationAction::SessionRetry { new_state } => new_state,
other => panic!("expected SessionRetry, got {other:?}"),
}
}
#[test]
fn hub_region_latch_sets_on_first_1002_single_master_dataplane() {
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
let new_state = session_retry_state_for_1002(&state);
assert!(
new_state.hub_region_processing_only,
"first 1002 on single-master data-plane should latch",
);
assert_eq!(new_state.session_token_retry_count, 1);
}
#[test]
fn hub_region_latch_does_not_set_on_multi_master_1002() {
let mut state = OperationRetryState::initial(0, true, Vec::new(), 3, 3);
state.is_dataplane = true;
let new_state = session_retry_state_for_1002(&state);
assert!(
!new_state.hub_region_processing_only,
"multi-master 1002 must not latch the hub-region header",
);
}
#[test]
fn hub_region_latch_stays_set_on_subsequent_1002() {
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
let after_first = session_retry_state_for_1002(&state);
assert!(after_first.hub_region_processing_only);
let after_second = session_retry_state_for_1002(&after_first);
assert!(
after_second.hub_region_processing_only,
"latch must persist across subsequent 1002 retries",
);
assert_eq!(after_second.session_token_retry_count, 2);
}
#[test]
fn hub_region_latch_does_not_set_on_non_1002_responses() {
let op = make_read_operation();
let endpoint = test_endpoint();
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
let (action, _) = evaluate_transport_result(&op, &endpoint, make_success_result(), &state);
assert!(matches!(action, OperationAction::Complete(_)));
let (action, _) =
evaluate_transport_result(&op, &endpoint, make_http_error(StatusCode::Gone), &state);
match action {
OperationAction::FailoverRetry { new_state, .. } => {
assert!(!new_state.hub_region_processing_only);
}
OperationAction::Abort { .. } => {
}
other => panic!("unexpected action for 410: {other:?}"),
}
let (action, _) = evaluate_transport_result(
&op,
&endpoint,
make_http_error(StatusCode::ServiceUnavailable),
&state,
);
match action {
OperationAction::FailoverRetry { new_state, .. } => {
assert!(!new_state.hub_region_processing_only);
}
OperationAction::Abort { .. } => {
}
other => panic!("unexpected action for 503: {other:?}"),
}
}
#[test]
fn hub_region_latch_state_at_budget_exhaustion() {
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
let after_first = session_retry_state_for_1002(&state);
assert!(after_first.hub_region_processing_only);
let after_second = session_retry_state_for_1002(&after_first);
assert!(after_second.hub_region_processing_only);
assert_eq!(after_second.session_token_retry_count, 2);
let op = make_read_operation();
let endpoint = test_endpoint();
let result = make_read_session_not_available_result();
let (action, _) = evaluate_transport_result(&op, &endpoint, result, &after_second);
assert!(
matches!(action, OperationAction::Abort { .. }),
"third 1002 must abort, got {action:?}",
);
}
#[test]
fn hub_region_latch_does_not_set_on_metadata_pipeline_1002() {
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
assert!(!state.is_dataplane);
let new_state = session_retry_state_for_1002(&state);
assert!(
!new_state.hub_region_processing_only,
"metadata-pipeline 1002 must not latch the hub-region header",
);
}
#[test]
fn hub_region_latch_independent_operations_do_not_share_state() {
let mut op_a = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
op_a.is_dataplane = true;
let mut op_b = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
op_b.is_dataplane = true;
let op_a_after = session_retry_state_for_1002(&op_a);
assert!(op_a_after.hub_region_processing_only);
assert!(!op_b.hub_region_processing_only);
let op_b_after = session_retry_state_for_1002(&op_b);
assert!(op_b_after.hub_region_processing_only);
assert!(op_a_after.hub_region_processing_only);
}
#[test]
fn hub_region_latch_survives_failover_after_latch() {
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
let after_1002 = session_retry_state_for_1002(&state);
assert!(after_1002.hub_region_processing_only);
let op = make_read_operation();
let endpoint = test_endpoint();
let (action, _) = evaluate_transport_result(
&op,
&endpoint,
make_http_error(StatusCode::ServiceUnavailable),
&after_1002,
);
match action {
OperationAction::FailoverRetry { new_state, .. } => {
assert!(
new_state.hub_region_processing_only,
"latch must propagate through `..self` in advance_failover",
);
}
OperationAction::Abort { .. } => {
}
other => panic!("unexpected action for 503 after latch: {other:?}"),
}
}
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc,
};
#[test]
fn shared_hub_region_latch_propagates_first_1002_across_hedges() {
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
let shared = Arc::new(AtomicBool::new(false));
state = state.with_shared_hub_region_latch(shared.clone());
let after = session_retry_state_for_1002(&state);
assert!(
after.hub_region_processing_only,
"per-state latch must still fire on the first 1002",
);
assert!(
shared.load(Ordering::Acquire),
"shared latch must be Release-stored when per-state latch fires",
);
assert!(
after.shared_hub_region_latch.as_ref().map(Arc::as_ptr) == Some(Arc::as_ptr(&shared)),
"advance_session_retry must propagate the shared latch via ..self",
);
}
#[test]
fn shared_hub_region_latch_does_not_set_on_multi_master_1002() {
let mut state = OperationRetryState::initial(0, true, Vec::new(), 3, 3); state.is_dataplane = true;
let shared = Arc::new(AtomicBool::new(false));
state = state.with_shared_hub_region_latch(shared.clone());
let after = session_retry_state_for_1002(&state);
assert!(
!after.hub_region_processing_only,
"multi-master never latches per-state",
);
assert!(
!shared.load(Ordering::Acquire),
"multi-master never flips the shared latch either",
);
}
#[test]
fn shared_hub_region_latch_does_not_set_on_metadata_pipeline_1002() {
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = false; let shared = Arc::new(AtomicBool::new(false));
state = state.with_shared_hub_region_latch(shared.clone());
let after = session_retry_state_for_1002(&state);
assert!(!after.hub_region_processing_only);
assert!(!shared.load(Ordering::Acquire));
}
#[test]
fn shared_hub_region_latch_is_monotonic_once_set() {
let mut state = OperationRetryState::initial(0, true, Vec::new(), 3, 3); state.is_dataplane = true;
let shared = Arc::new(AtomicBool::new(true)); state = state.with_shared_hub_region_latch(shared.clone());
let _ = session_retry_state_for_1002(&state);
assert!(
shared.load(Ordering::Acquire),
"shared latch must stay set across non-triggering retries",
);
}
#[test]
fn hedge_leg_effects_success_is_empty() {
let op = make_read_operation();
let endpoint = test_endpoint();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
let result = make_success_result();
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(eval.effects.is_empty());
assert!(!eval.observed_session_unavailable);
}
#[test]
fn hedge_leg_effects_503_emits_partition_and_endpoint_marks() {
let op = make_read_operation();
let endpoint = test_endpoint();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
let result = make_http_error(StatusCode::ServiceUnavailable);
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(eval
.effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(eval
.effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
assert!(!eval.observed_session_unavailable);
}
#[test]
fn hedge_leg_effects_transport_sent_emits_marks() {
let op = make_read_operation();
let endpoint = test_endpoint();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
let result = make_transport_error(RequestSentStatus::Sent);
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(eval
.effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkPartitionUnavailable(_))));
assert!(eval
.effects
.iter()
.any(|e| matches!(e, LocationEffect::MarkEndpointUnavailable { .. })));
}
#[test]
fn hedge_leg_effects_transport_not_sent_is_empty() {
let op = make_create_operation();
let endpoint = test_endpoint();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
let result = make_transport_error(RequestSentStatus::NotSent);
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(
eval.effects.is_empty(),
"not-sent transport error must not mark routing state",
);
}
#[test]
fn hedge_leg_effects_1002_signals_session_unavailable() {
let op = make_read_operation();
let endpoint = test_endpoint();
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
let result = make_read_session_not_available_result();
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(eval.effects.is_empty());
assert!(
eval.observed_session_unavailable,
"first 1002 on single-master dataplane should signal session-unavailable",
);
}
#[test]
fn hedge_leg_effects_1002_multi_master_no_signal() {
let op = make_read_operation();
let endpoint = test_endpoint();
let mut state = OperationRetryState::initial(0, true, Vec::new(), 3, 3); state.is_dataplane = true;
let result = make_read_session_not_available_result();
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(
!eval.observed_session_unavailable,
"multi-master 1002 must not flip the latch (AC-4)",
);
}
#[test]
fn hedge_leg_effects_1002_no_signal_when_already_latched() {
let op = make_read_operation();
let endpoint = test_endpoint();
let mut state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
state.is_dataplane = true;
state.hub_region_processing_only = true;
let result = make_read_session_not_available_result();
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(
!eval.observed_session_unavailable,
"already-latched state should not re-signal",
);
}
#[test]
fn hedge_leg_effects_deadline_exceeded_is_empty() {
let op = make_read_operation();
let endpoint = test_endpoint();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
let result = TransportResult {
outcome: TransportOutcome::DeadlineExceeded {
request_sent: RequestSentStatus::Sent,
},
};
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(eval.effects.is_empty());
assert!(!eval.observed_session_unavailable);
}
#[test]
fn hedge_leg_effects_409_conflict_is_empty() {
let op = make_create_operation();
let endpoint = test_endpoint();
let state = OperationRetryState::initial(0, false, Vec::new(), 3, 3);
let result = make_http_error(StatusCode::Conflict);
let eval = evaluate_hedge_leg_effects(&op, &endpoint, &state, &result);
assert!(
eval.effects.is_empty(),
"409 Conflict has no per-status handler; emits no effects",
);
}
}