use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ConsistencyMode {
Strong,
ReadYourWrites,
BoundedStaleness,
ReplicaBounded,
Eventual,
ProjectionOk,
CacheOk,
}
impl Default for ConsistencyMode {
fn default() -> Self {
Self::Strong
}
}
impl ConsistencyMode {
pub fn as_str(self) -> &'static str {
match self {
Self::Strong => "strong",
Self::ReadYourWrites => "read_your_writes",
Self::BoundedStaleness => "bounded_staleness",
Self::ReplicaBounded => "replica_bounded",
Self::Eventual => "eventual",
Self::ProjectionOk => "projection_ok",
Self::CacheOk => "cache_ok",
}
}
pub fn parse(token: &str) -> Option<Self> {
match token.trim().to_ascii_lowercase().as_str() {
"strong" | "linearizable" | "primary" => Some(Self::Strong),
"read_your_writes" | "ryw" | "read-your-writes" => Some(Self::ReadYourWrites),
"bounded_staleness" | "bounded-staleness" => Some(Self::BoundedStaleness),
"replica_bounded" | "replica-bounded" => Some(Self::ReplicaBounded),
"eventual" | "eventual_consistency" => Some(Self::Eventual),
"projection_ok" | "projection-ok" => Some(Self::ProjectionOk),
"cache_ok" | "cache-ok" => Some(Self::CacheOk),
_ => None,
}
}
pub fn parse_or_default(token: &str) -> Self {
Self::parse(token).unwrap_or_default()
}
pub fn allows_replica(self) -> bool {
!matches!(self, Self::Strong | Self::ReadYourWrites)
}
pub fn allows_projection(self) -> bool {
matches!(
self,
Self::Eventual | Self::ProjectionOk | Self::CacheOk | Self::BoundedStaleness
)
}
pub fn allows_cache(self) -> bool {
matches!(self, Self::CacheOk | Self::Eventual)
}
pub fn honours_fence(self) -> bool {
!matches!(self, Self::CacheOk)
}
pub(crate) fn from_proto_i32(mode: i32) -> Option<Self> {
match mode {
1 => Some(Self::Strong),
2 => Some(Self::ReadYourWrites),
3 => Some(Self::BoundedStaleness),
4 => Some(Self::ReplicaBounded),
5 => Some(Self::Eventual),
6 => Some(Self::ProjectionOk),
7 => Some(Self::CacheOk),
_ => None,
}
}
#[cfg(test)]
pub(crate) fn to_proto_i32(self) -> i32 {
match self {
Self::Strong => 1,
Self::ReadYourWrites => 2,
Self::BoundedStaleness => 3,
Self::ReplicaBounded => 4,
Self::Eventual => 5,
Self::ProjectionOk => 6,
Self::CacheOk => 7,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WriteReceipt {
pub source_lsn: String,
pub outbox_seq: u64,
pub projection_task_ids: Vec<String>,
pub manifest_checksum: String,
pub written_at_unix_ms: i64,
}
impl WriteReceipt {
pub fn empty() -> Self {
Self {
source_lsn: String::new(),
outbox_seq: 0,
projection_task_ids: Vec::new(),
manifest_checksum: String::new(),
written_at_unix_ms: 0,
}
}
pub fn is_empty(&self) -> bool {
self.source_lsn.is_empty()
&& self.outbox_seq == 0
&& self.projection_task_ids.is_empty()
&& self.manifest_checksum.is_empty()
}
pub(crate) fn to_proto(&self) -> crate::proto::WriteReceipt {
crate::proto::WriteReceipt {
source_lsn: self.source_lsn.clone(),
outbox_seq: self.outbox_seq,
projection_task_ids: self.projection_task_ids.clone(),
manifest_checksum: self.manifest_checksum.clone(),
written_at_unix_ms: self.written_at_unix_ms,
}
}
pub(crate) fn from_proto(proto: &crate::proto::WriteReceipt) -> Self {
Self {
source_lsn: proto.source_lsn.clone(),
outbox_seq: proto.outbox_seq,
projection_task_ids: proto.projection_task_ids.clone(),
manifest_checksum: proto.manifest_checksum.clone(),
written_at_unix_ms: proto.written_at_unix_ms,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
pub struct ReadFence {
#[serde(default, skip_serializing_if = "String::is_empty")]
pub min_outbox_lsn: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub projection_task_ids: Vec<String>,
#[serde(default)]
pub max_wait_ms: u64,
}
impl ReadFence {
pub fn from_receipt(receipt: &WriteReceipt, max_wait_ms: u64) -> Self {
Self {
min_outbox_lsn: receipt.source_lsn.clone(),
projection_task_ids: receipt.projection_task_ids.clone(),
max_wait_ms,
}
}
pub fn is_empty(&self) -> bool {
self.min_outbox_lsn.is_empty() && self.projection_task_ids.is_empty()
}
#[cfg(test)]
pub(crate) fn to_proto(&self) -> crate::proto::ReadFence {
crate::proto::ReadFence {
min_outbox_lsn: self.min_outbox_lsn.clone(),
projection_task_ids: self.projection_task_ids.clone(),
max_wait_ms: self.max_wait_ms,
}
}
pub(crate) fn from_proto(proto: &crate::proto::ReadFence) -> Self {
Self {
min_outbox_lsn: proto.min_outbox_lsn.clone(),
projection_task_ids: proto.projection_task_ids.clone(),
max_wait_ms: proto.max_wait_ms,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum StaleReadWarning {
FenceTimedOut {
backend: String,
instance: String,
lag_ms: u64,
},
ReplicaLagExceeded {
instance: String,
lag_ms: u64,
budget_ms: u64,
},
ProjectionMissing { backend: String, resource: String },
CacheStale {
backend: String,
cache_key_prefix: String,
},
}
impl StaleReadWarning {
pub fn kind_token(&self) -> &'static str {
match self {
Self::FenceTimedOut { .. } => "fence_timed_out",
Self::ReplicaLagExceeded { .. } => "replica_lag_exceeded",
Self::ProjectionMissing { .. } => "projection_missing",
Self::CacheStale { .. } => "cache_stale",
}
}
}
#[derive(Debug, Clone, Default)]
pub struct ConsistencyPolicy {
pub mode: ConsistencyMode,
pub fence: ReadFence,
pub max_replica_lag_ms: u64,
}
impl ConsistencyPolicy {
pub fn from_request_context(
consistency_token: &str,
max_replica_lag_ms: u64,
primary_read: bool,
eventual_consistency_allowed: bool,
) -> Self {
let mode = if primary_read {
ConsistencyMode::Strong
} else if consistency_token.trim().is_empty() && eventual_consistency_allowed {
ConsistencyMode::Eventual
} else {
ConsistencyMode::parse_or_default(consistency_token)
};
Self {
mode,
fence: ReadFence::default(),
max_replica_lag_ms,
}
}
pub fn with_fence(mut self, fence: ReadFence) -> Self {
self.fence = fence;
self
}
pub fn force_primary(&self) -> bool {
matches!(self.mode, ConsistencyMode::Strong)
|| (matches!(self.mode, ConsistencyMode::ReadYourWrites) && self.fence.is_empty())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DurabilityTokenClass {
RealPosition,
WallClock,
}
pub fn durability_token_class(backend_label: &str) -> DurabilityTokenClass {
match backend_label.trim().to_ascii_lowercase().as_str() {
"postgres" | "postgresql" | "pg" | "mysql" | "mariadb" | "mssql" | "sqlserver"
| "sqlite" | "mongodb" | "mongo" | "redis" | "qdrant" | "clickhouse" | "neo4j"
| "cassandra" | "scylla" | "weaviate" | "pinecone" | "elasticsearch" => {
DurabilityTokenClass::RealPosition
}
_ => DurabilityTokenClass::WallClock,
}
}
pub fn supports_bounded_reads(backend_label: &str) -> bool {
matches!(
durability_token_class(backend_label),
DurabilityTokenClass::RealPosition
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BoundedReadRefused {
pub backend: String,
pub reason: &'static str,
}
impl std::fmt::Display for BoundedReadRefused {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"bounded-staleness read refused for backend '{}': {}",
self.backend, self.reason
)
}
}
impl std::error::Error for BoundedReadRefused {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReadRouting {
Primary,
ReplicaBounded {
max_staleness_ms: u64,
min_lsn: Option<String>,
},
ReplicaUnfenced,
RefusedBounded(BoundedReadRefused),
}
impl ConsistencyPolicy {
pub fn route_read(&self, is_write: bool, backend_label: &str) -> ReadRouting {
if is_write {
return ReadRouting::Primary;
}
if self.force_primary() {
return ReadRouting::Primary;
}
match self.mode {
ConsistencyMode::ReadYourWrites => self.bounded_or_refuse(backend_label, false),
ConsistencyMode::BoundedStaleness => self.bounded_or_refuse(backend_label, true),
ConsistencyMode::ReplicaBounded => self.bounded_or_refuse(backend_label, false),
ConsistencyMode::Eventual
| ConsistencyMode::ProjectionOk
| ConsistencyMode::CacheOk => ReadRouting::ReplicaUnfenced,
ConsistencyMode::Strong => ReadRouting::Primary,
}
}
fn bounded_or_refuse(&self, backend_label: &str, refuse_on_wallclock: bool) -> ReadRouting {
let min_lsn = {
let lsn = self.fence.min_outbox_lsn.trim();
(!lsn.is_empty()).then(|| lsn.to_string())
};
let max_staleness_ms = match self.mode {
ConsistencyMode::BoundedStaleness | ConsistencyMode::ReplicaBounded => {
self.max_replica_lag_ms
}
_ => self.fence.max_wait_ms,
};
if supports_bounded_reads(backend_label) {
ReadRouting::ReplicaBounded {
max_staleness_ms,
min_lsn,
}
} else if refuse_on_wallclock {
ReadRouting::RefusedBounded(BoundedReadRefused {
backend: backend_label.to_string(),
reason: "backend mints no real replication-position token; a bounded-staleness \
fence on it would be a vacuous wall-clock fence",
})
} else {
ReadRouting::Primary
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn consistency_golden_value() -> serde_json::Value {
let receipt = WriteReceipt {
source_lsn: "0/1A2B3C4D".to_string(),
outbox_seq: 42,
projection_task_ids: vec![
"projection-task-a".to_string(),
"projection-task-b".to_string(),
],
manifest_checksum: "sha256:0123456789abcdef".to_string(),
written_at_unix_ms: 1_735_689_600_000,
};
let fence = ReadFence::from_receipt(&receipt, 2_500);
serde_json::json!({
"write_receipt": receipt,
"read_fence": fence,
})
}
#[test]
fn consistency_golden_json_matches_serde_contract() {
let root = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"));
let committed: serde_json::Value = serde_json::from_str(
&std::fs::read_to_string(root.join("docs/generated/consistency-golden.json"))
.expect("consistency golden JSON should be readable"),
)
.expect("consistency golden JSON should parse");
assert_eq!(
committed,
consistency_golden_value(),
"docs/generated/consistency-golden.json drifted from WriteReceipt/ReadFence serde"
);
}
#[test]
fn mode_tokens_are_pinned() {
assert_eq!(ConsistencyMode::Strong.as_str(), "strong");
assert_eq!(ConsistencyMode::ReadYourWrites.as_str(), "read_your_writes");
assert_eq!(
ConsistencyMode::BoundedStaleness.as_str(),
"bounded_staleness"
);
assert_eq!(ConsistencyMode::ReplicaBounded.as_str(), "replica_bounded");
assert_eq!(ConsistencyMode::Eventual.as_str(), "eventual");
assert_eq!(ConsistencyMode::ProjectionOk.as_str(), "projection_ok");
assert_eq!(ConsistencyMode::CacheOk.as_str(), "cache_ok");
}
#[test]
fn legacy_aliases_parse() {
assert_eq!(
ConsistencyMode::parse("linearizable"),
Some(ConsistencyMode::Strong)
);
assert_eq!(
ConsistencyMode::parse("primary"),
Some(ConsistencyMode::Strong)
);
assert_eq!(
ConsistencyMode::parse("eventual_consistency"),
Some(ConsistencyMode::Eventual)
);
assert_eq!(
ConsistencyMode::parse("ryw"),
Some(ConsistencyMode::ReadYourWrites)
);
assert_eq!(ConsistencyMode::parse("unknown_mode"), None);
assert_eq!(ConsistencyMode::parse(""), None);
assert_eq!(
ConsistencyMode::parse_or_default("unknown"),
ConsistencyMode::Strong
);
}
#[test]
fn mode_capability_matrix_is_correct() {
assert!(!ConsistencyMode::Strong.allows_replica());
assert!(!ConsistencyMode::Strong.allows_projection());
assert!(!ConsistencyMode::Strong.allows_cache());
assert!(!ConsistencyMode::ReadYourWrites.allows_replica());
assert!(!ConsistencyMode::ReadYourWrites.allows_projection());
assert!(!ConsistencyMode::ReadYourWrites.allows_cache());
assert!(ConsistencyMode::BoundedStaleness.allows_replica());
assert!(ConsistencyMode::BoundedStaleness.allows_projection());
assert!(!ConsistencyMode::BoundedStaleness.allows_cache());
assert!(ConsistencyMode::Eventual.allows_replica());
assert!(ConsistencyMode::Eventual.allows_projection());
assert!(ConsistencyMode::Eventual.allows_cache());
assert!(ConsistencyMode::ProjectionOk.allows_replica());
assert!(ConsistencyMode::ProjectionOk.allows_projection());
assert!(!ConsistencyMode::ProjectionOk.allows_cache());
assert!(ConsistencyMode::CacheOk.allows_replica());
assert!(ConsistencyMode::CacheOk.allows_projection());
assert!(ConsistencyMode::CacheOk.allows_cache());
}
#[test]
fn cache_ok_skips_fence_others_honour_it() {
assert!(!ConsistencyMode::CacheOk.honours_fence());
assert!(ConsistencyMode::Strong.honours_fence());
assert!(ConsistencyMode::ReadYourWrites.honours_fence());
assert!(ConsistencyMode::Eventual.honours_fence());
}
#[test]
fn write_receipt_converts_to_read_fence() {
let receipt = WriteReceipt {
source_lsn: "0/1A2B3C".into(),
outbox_seq: 42,
projection_task_ids: vec!["task-a".into(), "task-b".into()],
manifest_checksum: "abc123".into(),
written_at_unix_ms: 1_700_000_000_000,
};
let fence = ReadFence::from_receipt(&receipt, 5_000);
assert_eq!(fence.min_outbox_lsn, "0/1A2B3C");
assert_eq!(
fence.projection_task_ids,
vec!["task-a".to_string(), "task-b".to_string()]
);
assert_eq!(fence.max_wait_ms, 5_000);
assert!(!fence.is_empty());
}
#[test]
fn write_receipt_proto_round_trip_preserves_serde_shape() {
let receipt = WriteReceipt {
source_lsn: "0/1A2B3C".into(),
outbox_seq: 42,
projection_task_ids: vec!["task-a".into(), "task-b".into()],
manifest_checksum: "abc123".into(),
written_at_unix_ms: 1_700_000_000_000,
};
let proto = receipt.to_proto();
let decoded = WriteReceipt::from_proto(&proto);
assert_eq!(decoded, receipt);
assert_eq!(
serde_json::to_value(decoded).unwrap(),
serde_json::json!({
"source_lsn": "0/1A2B3C",
"outbox_seq": 42,
"projection_task_ids": ["task-a", "task-b"],
"manifest_checksum": "abc123",
"written_at_unix_ms": 1_700_000_000_000i64,
})
);
}
#[test]
fn read_fence_proto_round_trip_preserves_serde_shape() {
let fence = ReadFence {
min_outbox_lsn: "0/1A2B3C".into(),
projection_task_ids: vec!["task-a".into()],
max_wait_ms: 2_500,
};
let proto = fence.to_proto();
let decoded = ReadFence::from_proto(&proto);
assert_eq!(decoded, fence);
assert_eq!(
serde_json::to_value(decoded).unwrap(),
serde_json::json!({
"min_outbox_lsn": "0/1A2B3C",
"projection_task_ids": ["task-a"],
"max_wait_ms": 2_500,
})
);
}
#[test]
fn consistency_mode_proto_numbers_match_pinned_wire_tokens() {
for mode in [
ConsistencyMode::Strong,
ConsistencyMode::ReadYourWrites,
ConsistencyMode::BoundedStaleness,
ConsistencyMode::ReplicaBounded,
ConsistencyMode::Eventual,
ConsistencyMode::ProjectionOk,
ConsistencyMode::CacheOk,
] {
assert_eq!(
ConsistencyMode::from_proto_i32(mode.to_proto_i32()),
Some(mode),
"{} must round-trip through the proto enum number",
mode.as_str()
);
}
assert_eq!(ConsistencyMode::from_proto_i32(0), None);
assert_eq!(ConsistencyMode::from_proto_i32(999), None);
}
#[test]
fn empty_receipts_and_fences_are_detectable() {
assert!(WriteReceipt::empty().is_empty());
assert!(ReadFence::default().is_empty());
}
#[test]
fn from_request_context_honours_legacy_flags() {
let p = ConsistencyPolicy::from_request_context("eventual", 0, true, false);
assert_eq!(p.mode, ConsistencyMode::Strong);
let p = ConsistencyPolicy::from_request_context("", 0, false, true);
assert_eq!(p.mode, ConsistencyMode::Eventual);
let p = ConsistencyPolicy::from_request_context("read_your_writes", 0, false, true);
assert_eq!(p.mode, ConsistencyMode::ReadYourWrites);
}
#[test]
fn force_primary_respects_ryw_fence_state() {
let strong = ConsistencyPolicy {
mode: ConsistencyMode::Strong,
..Default::default()
};
assert!(strong.force_primary());
let ryw_no_fence = ConsistencyPolicy {
mode: ConsistencyMode::ReadYourWrites,
..Default::default()
};
assert!(
ryw_no_fence.force_primary(),
"RYW with no fence behaves like Strong"
);
let ryw_with_fence = ConsistencyPolicy {
mode: ConsistencyMode::ReadYourWrites,
fence: ReadFence {
min_outbox_lsn: "0/100".into(),
..Default::default()
},
..Default::default()
};
assert!(
!ryw_with_fence.force_primary(),
"RYW with fence allows replica routing once the fence is satisfied"
);
let eventual = ConsistencyPolicy {
mode: ConsistencyMode::Eventual,
..Default::default()
};
assert!(!eventual.force_primary());
}
#[test]
fn stale_read_warning_tokens_are_pinned() {
let cases = [
(
StaleReadWarning::FenceTimedOut {
backend: "mongodb".into(),
instance: "default".into(),
lag_ms: 1234,
},
"fence_timed_out",
),
(
StaleReadWarning::ReplicaLagExceeded {
instance: "replica-1".into(),
lag_ms: 5_000,
budget_ms: 1_000,
},
"replica_lag_exceeded",
),
(
StaleReadWarning::ProjectionMissing {
backend: "qdrant".into(),
resource: "customers_vec".into(),
},
"projection_missing",
),
(
StaleReadWarning::CacheStale {
backend: "redis".into(),
cache_key_prefix: "udb:billing".into(),
},
"cache_stale",
),
];
for (warning, token) in cases {
assert_eq!(warning.kind_token(), token);
}
}
#[test]
fn serde_round_trip_matches_wire_tokens() {
for mode in [
ConsistencyMode::Strong,
ConsistencyMode::ReadYourWrites,
ConsistencyMode::BoundedStaleness,
ConsistencyMode::ReplicaBounded,
ConsistencyMode::Eventual,
ConsistencyMode::ProjectionOk,
ConsistencyMode::CacheOk,
] {
let json = serde_json::to_string(&mode).unwrap();
let token = json.trim_matches('"');
assert_eq!(token, mode.as_str());
let back: ConsistencyMode = serde_json::from_str(&json).unwrap();
assert_eq!(back, mode);
}
}
#[test]
fn durability_token_class_is_pinned() {
for real in [
"postgres",
"postgresql",
"mysql",
"mariadb",
"mssql",
"sqlserver",
"sqlite",
"mongodb",
"redis",
"qdrant",
"clickhouse",
"neo4j",
"cassandra",
"weaviate",
"pinecone",
"elasticsearch",
] {
assert_eq!(
durability_token_class(real),
DurabilityTokenClass::RealPosition,
"{real} mints a real position token"
);
assert!(
supports_bounded_reads(real),
"{real} supports bounded reads"
);
}
for wall in [
"s3",
"minio",
"azureblob",
"gcs",
"memcached",
"weird-future-backend",
] {
assert_eq!(
durability_token_class(wall),
DurabilityTokenClass::WallClock,
"{wall} has no real position token"
);
assert!(
!supports_bounded_reads(wall),
"{wall} must be refused a bounded read"
);
}
assert!(supports_bounded_reads(" Postgres "));
}
#[test]
fn writes_always_route_to_primary() {
for mode in [
ConsistencyMode::Strong,
ConsistencyMode::ReadYourWrites,
ConsistencyMode::BoundedStaleness,
ConsistencyMode::ReplicaBounded,
ConsistencyMode::Eventual,
ConsistencyMode::ProjectionOk,
ConsistencyMode::CacheOk,
] {
let policy = ConsistencyPolicy {
mode,
max_replica_lag_ms: 1_000,
..Default::default()
};
assert_eq!(
policy.route_read(true, "postgres"),
ReadRouting::Primary,
"write under {mode:?} must route to primary"
);
}
}
#[test]
fn strong_and_unfenced_ryw_route_to_primary() {
let strong = ConsistencyPolicy {
mode: ConsistencyMode::Strong,
..Default::default()
};
assert_eq!(strong.route_read(false, "postgres"), ReadRouting::Primary);
let ryw = ConsistencyPolicy {
mode: ConsistencyMode::ReadYourWrites,
..Default::default()
};
assert_eq!(ryw.route_read(false, "postgres"), ReadRouting::Primary);
}
#[test]
fn bounded_staleness_routes_to_replica_with_real_token() {
let policy = ConsistencyPolicy {
mode: ConsistencyMode::BoundedStaleness,
max_replica_lag_ms: 750,
fence: ReadFence {
min_outbox_lsn: "0/1A2B3C".into(),
max_wait_ms: 9_999,
..Default::default()
},
};
assert_eq!(
policy.route_read(false, "postgres"),
ReadRouting::ReplicaBounded {
max_staleness_ms: 750,
min_lsn: Some("0/1A2B3C".to_string()),
}
);
let policy = ConsistencyPolicy {
mode: ConsistencyMode::BoundedStaleness,
max_replica_lag_ms: 200,
fence: ReadFence::default(),
};
assert_eq!(
policy.route_read(false, "mongodb"),
ReadRouting::ReplicaBounded {
max_staleness_ms: 200,
min_lsn: None,
}
);
}
#[test]
fn bounded_staleness_refused_on_wallclock_backend() {
let policy = ConsistencyPolicy {
mode: ConsistencyMode::BoundedStaleness,
max_replica_lag_ms: 500,
fence: ReadFence::default(),
};
for backend in ["s3", "azureblob", "gcs", "memcached", "unknown"] {
match policy.route_read(false, backend) {
ReadRouting::RefusedBounded(refused) => {
assert_eq!(refused.backend, backend);
assert!(!refused.reason.is_empty());
let err: &dyn std::error::Error = &refused;
assert!(err.to_string().contains(backend));
}
other => panic!("expected RefusedBounded for {backend}, got {other:?}"),
}
}
}
#[test]
fn fenced_ryw_routes_to_replica_or_primary_never_refused() {
let policy = ConsistencyPolicy {
mode: ConsistencyMode::ReadYourWrites,
max_replica_lag_ms: 0,
fence: ReadFence {
min_outbox_lsn: "0/200".into(),
max_wait_ms: 1_500,
..Default::default()
},
};
assert_eq!(
policy.route_read(false, "postgres"),
ReadRouting::ReplicaBounded {
max_staleness_ms: 1_500,
min_lsn: Some("0/200".to_string()),
}
);
assert_eq!(policy.route_read(false, "s3"), ReadRouting::Primary);
}
#[test]
fn eventual_family_routes_to_unfenced_replica() {
for mode in [
ConsistencyMode::Eventual,
ConsistencyMode::ProjectionOk,
ConsistencyMode::CacheOk,
] {
let policy = ConsistencyPolicy {
mode,
..Default::default()
};
assert_eq!(
policy.route_read(false, "s3"),
ReadRouting::ReplicaUnfenced,
"{mode:?} routes to an unfenced replica"
);
}
}
#[test]
fn replica_bounded_token_and_aliases_parse() {
assert_eq!(ConsistencyMode::ReplicaBounded.as_str(), "replica_bounded");
assert_eq!(
ConsistencyMode::parse("replica_bounded"),
Some(ConsistencyMode::ReplicaBounded)
);
assert_eq!(
ConsistencyMode::parse("replica-bounded"),
Some(ConsistencyMode::ReplicaBounded)
);
assert_ne!(
ConsistencyMode::ReplicaBounded,
ConsistencyMode::BoundedStaleness
);
}
#[test]
fn replica_bounded_is_replica_only() {
assert!(ConsistencyMode::ReplicaBounded.allows_replica());
assert!(ConsistencyMode::ReplicaBounded.honours_fence());
assert!(
!ConsistencyMode::ReplicaBounded.allows_projection(),
"REPLICA_BOUNDED targets SQL replicas only, never projections"
);
assert!(!ConsistencyMode::ReplicaBounded.allows_cache());
}
#[test]
fn replica_bounded_routes_to_replica_then_fails_over_to_primary() {
let policy = ConsistencyPolicy {
mode: ConsistencyMode::ReplicaBounded,
max_replica_lag_ms: 500,
fence: ReadFence {
min_outbox_lsn: "0/1A2B3C".into(),
max_wait_ms: 9_999,
..Default::default()
},
};
assert_eq!(
policy.route_read(false, "postgres"),
ReadRouting::ReplicaBounded {
max_staleness_ms: 500,
min_lsn: Some("0/1A2B3C".to_string()),
}
);
assert_eq!(policy.route_read(false, "s3"), ReadRouting::Primary);
assert_eq!(policy.route_read(true, "postgres"), ReadRouting::Primary);
}
}