mod config;
pub mod decision;
mod table;
mod types;
use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
use std::sync::{OnceLock, RwLock};
use std::time::{Duration, Instant};
use super::{catalog, AttrValue};
use table::{Insert, Table};
pub use config::{
enabled, parse_probe_flag, set_probe_enabled, window_config_from_env, WindowConfig,
DEFAULT_MAX_PREFIX_BITS,
};
pub use decision::{Refusal, Verdict};
pub use types::{
AccessWeight, AccessWeightBuilder, BucketStats, RepeatBucket, SizeSource, WindowSummary,
};
#[derive(Clone, Copy, Debug)]
enum SlowPath {
None,
Allocate,
Downsample {
replay: bool,
},
}
struct WindowState {
recorded: AtomicU64,
dropped: AtomicU64,
table: Table,
prefix_bits: u32,
at_sampling_floor: bool,
started: Instant,
}
pub struct PartitionAccessRecorder {
state: RwLock<Option<WindowState>>,
close_sequence: AtomicU64,
gauge_order: std::sync::Mutex<u64>,
duration_nanos: AtomicU64,
max_accesses: AtomicU64,
max_prefix_bits: AtomicU32,
}
impl Default for PartitionAccessRecorder {
fn default() -> Self {
Self::new(WindowConfig::default())
}
}
impl PartitionAccessRecorder {
pub fn new(config: WindowConfig) -> Self {
let recorder = Self {
state: RwLock::new(None),
close_sequence: AtomicU64::new(0),
gauge_order: std::sync::Mutex::new(0),
duration_nanos: AtomicU64::new(0),
max_accesses: AtomicU64::new(0),
max_prefix_bits: AtomicU32::new(0),
};
recorder.set_window_config(config);
recorder
}
pub fn set_window_config(&self, config: WindowConfig) {
self.duration_nanos.store(
u64::try_from(config.duration.as_nanos()).unwrap_or(u64::MAX),
Ordering::Relaxed,
);
self.max_accesses
.store(config.max_accesses, Ordering::Relaxed);
self.max_prefix_bits
.store(config.max_prefix_bits, Ordering::Relaxed);
}
pub fn window_config(&self) -> WindowConfig {
WindowConfig {
duration: Duration::from_nanos(self.duration_nanos.load(Ordering::Relaxed)),
max_accesses: self.max_accesses.load(Ordering::Relaxed),
max_prefix_bits: self.max_prefix_bits.load(Ordering::Relaxed),
}
}
pub fn footprint_bytes(&self) -> usize {
match self.state.read() {
Ok(g) => g.as_ref().map_or(0, |s| s.table.footprint_bytes()),
Err(poisoned) => poisoned
.into_inner()
.as_ref()
.map_or(0, |s| s.table.footprint_bytes()),
}
}
pub const fn declared_footprint_bytes() -> usize {
table::TABLE_BYTES
}
pub fn record(
&self,
scope: TableScope<'_>,
key: &[u8],
weight: AccessWeight,
) -> Vec<WindowSummary> {
let hash = table::hash_partition(scope.keyspace, scope.table, key);
let bytes = weight.bytes();
let flags = match weight.source() {
SizeSource::SuccessorGap => table::FLAG_SIZE_FROM_GAP,
SizeSource::Index | SizeSource::Unavailable => 0,
};
let closed = self.close_if_expired();
let mut slow = SlowPath::None;
{
let guard = self.read_state();
if let Some(state) = guard.as_ref() {
match state
.table
.record_with_flags(hash, state.prefix_bits, bytes, flags)
{
Insert::Recorded => {
state.recorded.fetch_add(1, Ordering::Relaxed);
if state.table.occupancy() >= table::LOAD_FACTOR_LIMIT {
slow = SlowPath::Downsample { replay: false };
}
}
Insert::NotAdmitted => {
state.recorded.fetch_add(1, Ordering::Relaxed);
}
Insert::Full => slow = SlowPath::Downsample { replay: true },
}
} else {
slow = SlowPath::Allocate;
}
}
match slow {
SlowPath::None => {}
SlowPath::Allocate => self.grow_or_downsample(hash, bytes, flags, true),
SlowPath::Downsample { replay } => self.grow_or_downsample(hash, bytes, flags, replay),
}
let mut closed_windows = Vec::new();
closed_windows.extend(closed);
closed_windows.extend(self.close_if_over_count());
closed_windows
}
fn grow_or_downsample(&self, hash: u64, bytes: Option<u64>, flags: u8, replay: bool) {
let mut guard = self.write_state();
match guard.as_mut() {
None => {
let state = WindowState {
recorded: AtomicU64::new(0),
dropped: AtomicU64::new(0),
table: Table::new(),
prefix_bits: 0,
at_sampling_floor: false,
started: Instant::now(),
};
if replay {
state.table.record_with_flags(hash, 0, bytes, flags);
state.recorded.fetch_add(1, Ordering::Relaxed);
}
*guard = Some(state);
}
Some(state) => {
while state.table.occupancy() >= table::LOAD_FACTOR_LIMIT
&& state.prefix_bits < self.max_prefix_bits.load(Ordering::Relaxed)
{
state.prefix_bits += 1;
state.table.downsample(state.prefix_bits);
}
if state.prefix_bits >= self.max_prefix_bits.load(Ordering::Relaxed) {
state.at_sampling_floor = true;
}
if replay {
let mut seated =
state
.table
.record_with_flags(hash, state.prefix_bits, bytes, flags);
while seated == Insert::Full
&& state.prefix_bits < self.max_prefix_bits.load(Ordering::Relaxed)
{
state.prefix_bits += 1;
state.table.downsample(state.prefix_bits);
state.at_sampling_floor |=
state.prefix_bits >= self.max_prefix_bits.load(Ordering::Relaxed);
seated =
state
.table
.record_with_flags(hash, state.prefix_bits, bytes, flags);
}
state.recorded.fetch_add(1, Ordering::Relaxed);
if seated == Insert::Full {
state.dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
}
}
fn close_if_expired(&self) -> Option<WindowSummary> {
let duration_nanos = self.duration_nanos.load(Ordering::Relaxed);
self.close_window_if(|state| {
u64::try_from(state.started.elapsed().as_nanos()).unwrap_or(u64::MAX) >= duration_nanos
})
}
fn close_if_over_count(&self) -> Option<WindowSummary> {
let max_accesses = self.max_accesses.load(Ordering::Relaxed);
self.close_window_if(|state| state.recorded.load(Ordering::Relaxed) >= max_accesses)
}
fn close_window_if(
&self,
should_close: impl FnOnce(&WindowState) -> bool,
) -> Option<WindowSummary> {
let mut guard = self.write_state();
if !should_close(guard.as_ref()?) {
return None;
}
self.close_locked(&mut guard)
}
pub fn close_window(&self) -> Option<WindowSummary> {
let mut guard = self.write_state();
self.close_locked(&mut guard)
}
fn close_locked(&self, guard: &mut Option<WindowState>) -> Option<WindowSummary> {
let state = guard.as_mut()?;
let recorded = state.recorded.swap(0, Ordering::Relaxed);
let dropped = state.dropped.swap(0, Ordering::Relaxed);
let mut summary = WindowSummary {
close_sequence: self.close_sequence.fetch_add(1, Ordering::Relaxed),
sample_denominator: 1u64 << state.prefix_bits,
at_sampling_floor: state.at_sampling_floor,
recorded_accesses: recorded,
dropped_accesses: dropped,
..Default::default()
};
state.table.for_each_entry(|e| {
let idx = RepeatBucket::from_count(e.count).index();
let b = &mut summary.buckets[idx];
b.accesses = b.accesses.saturating_add(u64::from(e.count));
if e.size_unavailable {
b.distinct_unavailable += 1;
} else {
if e.size_from_gap {
b.distinct_successor_gap += 1;
} else {
b.distinct_index += 1;
}
b.bytes = b.bytes.saturating_add(e.bytes);
}
});
state.table.reset();
state.prefix_bits = 0;
state.at_sampling_floor = false;
state.started = Instant::now();
if summary.recorded_accesses == 0 {
return None;
}
Some(summary)
}
fn emit_window(&self, summary: &WindowSummary) {
emit_counters(summary);
let mut high_water = self
.gauge_order
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if summary.close_sequence < *high_water {
return;
}
*high_water = summary.close_sequence + 1;
super::record_gauge(
catalog::READ_PARTITION_ACCESS_SAMPLE_DENOMINATOR,
i64::try_from(summary.sample_denominator).unwrap_or(i64::MAX),
&[],
);
super::record_gauge(
catalog::READ_PARTITION_ACCESS_WINDOW_DROPPED,
i64::try_from(summary.dropped_accesses).unwrap_or(i64::MAX),
&[],
);
super::record_gauge(
catalog::READ_PARTITION_ACCESS_SAMPLING_FLOOR,
i64::from(summary.at_sampling_floor),
&[],
);
}
#[cfg(feature = "observability-testing")]
pub fn emit_for_test(&self, summary: &WindowSummary) {
self.emit_window(summary);
}
fn read_state(&self) -> std::sync::RwLockReadGuard<'_, Option<WindowState>> {
self.state.read().unwrap_or_else(|e| e.into_inner())
}
fn write_state(&self) -> std::sync::RwLockWriteGuard<'_, Option<WindowState>> {
self.state.write().unwrap_or_else(|e| e.into_inner())
}
}
pub fn global() -> &'static PartitionAccessRecorder {
static GLOBAL: OnceLock<PartitionAccessRecorder> = OnceLock::new();
GLOBAL.get_or_init(|| PartitionAccessRecorder::new(window_config_from_env()))
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct TableScope<'a> {
pub keyspace: &'a str,
pub table: &'a str,
}
impl<'a> TableScope<'a> {
pub fn new(keyspace: &'a str, table: &'a str) -> Self {
Self { keyspace, table }
}
pub fn from_qualified(name: &'a str) -> Self {
match name.split_once('.') {
Some((ks, table)) => Self {
keyspace: ks,
table,
},
None => Self {
keyspace: "",
table: name,
},
}
}
}
#[inline]
pub fn record_partition_access(scope: TableScope<'_>, key: &[u8], weight: AccessWeight) {
if !enabled() {
return;
}
for summary in global().record(scope, key, weight) {
global().emit_window(&summary);
}
}
pub fn close_window() -> Option<WindowSummary> {
let summary = global().close_window()?;
global().emit_window(&summary);
Some(summary)
}
fn emit_counters(summary: &WindowSummary) {
use super::add_counter;
for b in RepeatBucket::ALL {
let stats = summary.bucket(b);
let bucket: AttrValue = b.label().into();
for (count, source) in [
(stats.distinct_index, SizeSource::Index),
(stats.distinct_successor_gap, SizeSource::SuccessorGap),
(stats.distinct_unavailable, SizeSource::Unavailable),
] {
if count > 0 {
add_counter(
catalog::READ_PARTITION_ACCESS_DISTINCT_PARTITIONS,
count,
&[
(catalog::attr::REPEAT_BUCKET, bucket.clone()),
(catalog::attr::SIZE_SOURCE, source.label().into()),
],
);
}
}
if stats.accesses > 0 {
add_counter(
catalog::READ_PARTITION_ACCESS_ACCESSES,
stats.accesses,
&[(catalog::attr::REPEAT_BUCKET, bucket.clone())],
);
}
if stats.bytes > 0 {
add_counter(
catalog::READ_PARTITION_ACCESS_BYTES,
stats.bytes,
&[(catalog::attr::REPEAT_BUCKET, bucket)],
);
}
}
if summary.dropped_accesses > 0 {
add_counter(
catalog::READ_PARTITION_ACCESS_DROPPED,
summary.dropped_accesses,
&[],
);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn bucket_boundaries_are_exactly_the_six_declared_ranges() {
let expected: &[(u32, RepeatBucket)] = &[
(1, RepeatBucket::One),
(2, RepeatBucket::Two),
(3, RepeatBucket::ThreeToFour),
(4, RepeatBucket::ThreeToFour),
(5, RepeatBucket::FiveToEight),
(8, RepeatBucket::FiveToEight),
(9, RepeatBucket::NineToSixteen),
(16, RepeatBucket::NineToSixteen),
(17, RepeatBucket::SeventeenPlus),
(u32::MAX, RepeatBucket::SeventeenPlus),
];
for (count, want) in expected {
assert_eq!(RepeatBucket::from_count(*count), *want, "count {count}");
}
}
#[test]
fn bucket_labels_are_the_six_verbatim_values() {
let labels: Vec<&str> = RepeatBucket::ALL.iter().map(|b| b.label()).collect();
assert_eq!(labels, vec!["1", "2", "3-4", "5-8", "9-16", "17+"]);
}
#[test]
fn probe_flag_parsing_rejects_unrecognised_values() {
assert_eq!(parse_probe_flag("1"), Some(true));
assert_eq!(parse_probe_flag(" TRUE "), Some(true));
assert_eq!(parse_probe_flag("on"), Some(true));
assert_eq!(parse_probe_flag("0"), Some(false));
assert_eq!(parse_probe_flag("off"), Some(false));
assert_eq!(parse_probe_flag(""), Some(false));
assert_eq!(parse_probe_flag("yes please"), None);
assert_eq!(parse_probe_flag("maybe"), None);
}
#[test]
fn an_unpriced_builder_finishes_unavailable_rather_than_zero_bytes() {
assert_eq!(
AccessWeightBuilder::new().finish(),
AccessWeight::Unavailable,
"an access with nothing to price is not an access priced at zero"
);
let mut b = AccessWeightBuilder::new();
b.note_sized(100);
b.note_sized(200);
assert_eq!(b.finish(), AccessWeight::Index(300));
let mut b = AccessWeightBuilder::new();
b.note_measured(4_096);
assert_eq!(b.finish(), AccessWeight::SuccessorGap(4_096));
let mut b = AccessWeightBuilder::new();
b.note_sized(100);
b.note_sized(0);
assert_eq!(b.finish(), AccessWeight::Unavailable);
let mut b = AccessWeightBuilder::new();
b.note_sized(100);
b.note_unsized();
assert_eq!(b.finish(), AccessWeight::Unavailable);
}
#[test]
fn a_disabled_recorder_allocates_nothing() {
let r = PartitionAccessRecorder::default();
assert_eq!(r.footprint_bytes(), 0);
assert_eq!(r.close_window(), None);
assert_eq!(r.footprint_bytes(), 0);
}
#[test]
fn footprint_is_fixed_once_recording_starts() {
let r = PartitionAccessRecorder::default();
r.record(
TableScope::new("ks", "t"),
b"k",
AccessWeight::SuccessorGap(1),
);
assert_eq!(
r.footprint_bytes(),
PartitionAccessRecorder::declared_footprint_bytes()
);
assert_eq!(r.footprint_bytes(), 3 * 1024 * 1024);
}
}