use std::sync::atomic::{AtomicU64, Ordering};
use crate::bridge::envelope::Request;
static WRITES_BYPASSED_ADMISSION_GATE: AtomicU64 = AtomicU64::new(0);
pub fn writes_bypassed_admission_gate() -> u64 {
WRITES_BYPASSED_ADMISSION_GATE.load(Ordering::Relaxed)
}
#[inline]
pub(crate) fn assert_write_admitted(request: &Request) {
let bypassed = crate::control::server::shared::write_admission::plan_is_write(&request.plan)
&& request.admission.is_exempt_as_read();
if bypassed {
WRITES_BYPASSED_ADMISSION_GATE.fetch_add(1, Ordering::Relaxed);
}
debug_assert!(
!bypassed,
"a write-class plan reached the SPSC enqueue marked Exempt(Read) — it bypassed the write-admission gate"
);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bridge::envelope::{Admission, ExemptReason, PhysicalPlan, Priority};
use crate::control::server::shared::write_admission::plan_is_write;
use crate::types::{DatabaseId, ReadConsistency, RequestId, TenantId, TraceId, VShardId};
use nodedb_physical::physical_plan::{DocumentOp, KvOp};
use std::time::{Duration, Instant};
fn write_plan() -> PhysicalPlan {
PhysicalPlan::Kv(KvOp::Put {
collection: "c".into(),
key: b"k".to_vec(),
value: b"v".to_vec(),
ttl_ms: 0,
surrogate: nodedb_types::Surrogate::ZERO,
})
}
fn read_plan() -> PhysicalPlan {
PhysicalPlan::Document(DocumentOp::PointGet {
collection: "c".into(),
document_id: "d".into(),
surrogate: nodedb_types::Surrogate::ZERO,
pk_bytes: Vec::new(),
rls_filters: Vec::new(),
system_time: nodedb_types::SystemTimeScope::Current,
valid_at_ms: None,
})
}
fn request_with(plan: PhysicalPlan, admission: Admission) -> Request {
Request {
request_id: RequestId::new(1),
tenant_id: TenantId::new(1),
database_id: DatabaseId::DEFAULT,
vshard_id: VShardId::new(0),
plan,
deadline: Instant::now() + Duration::from_secs(5),
priority: Priority::Normal,
trace_id: TraceId::generate(),
consistency: ReadConsistency::Strong,
idempotency_key: None,
event_source: crate::event::EventSource::User,
user_roles: Vec::new(),
user_id: None,
statement_digest: None,
txn_id: None,
wal_lsn: None,
resolved_now_ms: None,
admission,
}
}
#[test]
#[should_panic(expected = "bypassed the write-admission gate")]
fn write_marked_exempt_as_read_trips_guard() {
let req = request_with(write_plan(), Admission::Exempt(ExemptReason::Read));
assert_write_admitted(&req);
}
#[test]
fn admitted_write_does_not_trip() {
let req = request_with(write_plan(), Admission::Admitted);
assert_write_admitted(&req);
}
#[test]
fn read_marked_exempt_as_read_does_not_trip() {
let req = request_with(read_plan(), Admission::Exempt(ExemptReason::Read));
assert_write_admitted(&req);
}
#[test]
fn write_marked_already_ordered_does_not_trip() {
let req = request_with(
write_plan(),
Admission::Exempt(ExemptReason::AlreadyOrdered),
);
assert_write_admitted(&req);
}
#[test]
fn plan_is_write_classification_and_bypass_predicate() {
let w = write_plan();
let r = read_plan();
assert!(plan_is_write(&w));
assert!(!plan_is_write(&r));
let exempt_read = Admission::Exempt(ExemptReason::Read);
assert!(plan_is_write(&w) && exempt_read.is_exempt_as_read());
assert!(!(plan_is_write(&r) && exempt_read.is_exempt_as_read()));
assert!(!(plan_is_write(&w) && Admission::Admitted.is_exempt_as_read()));
assert!(
!(plan_is_write(&w)
&& Admission::Exempt(ExemptReason::AlreadyOrdered).is_exempt_as_read())
);
}
}