Skip to main content

reifydb_core/interface/catalog/
ringbuffer.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_codec::row::pod::EncodedPodRow;
5use reifydb_value::{Result, value::Value};
6use serde::{Deserialize, Serialize};
7
8use crate::{
9	common::TimeSource,
10	interface::catalog::{
11		column::Column,
12		id::{NamespaceId, RingBufferId},
13		key::PrimaryKey,
14	},
15	return_internal_error,
16};
17
18#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
19pub struct RingBuffer {
20	pub id: RingBufferId,
21	pub namespace: NamespaceId,
22	pub name: String,
23	pub columns: Vec<Column>,
24	pub capacity: u64,
25	pub primary_key: Option<PrimaryKey>,
26	pub partition_by: Vec<String>,
27	pub time: TimeSource,
28}
29
30impl RingBuffer {
31	pub fn name(&self) -> &str {
32		&self.name
33	}
34}
35
36#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
37pub struct RingBufferMetadata {
38	pub count: u64,
39	pub head: u64,
40	pub tail: u64,
41}
42
43#[derive(Debug, Clone, PartialEq)]
44pub struct PartitionedMetadata {
45	pub metadata: RingBufferMetadata,
46	pub partition_values: Vec<Value>,
47}
48
49impl RingBufferMetadata {
50	pub fn new() -> Self {
51		Self {
52			count: 0,
53			head: 1,
54			tail: 1,
55		}
56	}
57
58	pub fn is_full(&self, capacity: u64) -> bool {
59		self.count >= capacity
60	}
61
62	pub fn is_empty(&self) -> bool {
63		self.count == 0
64	}
65}
66
67impl Default for RingBufferMetadata {
68	fn default() -> Self {
69		Self::new()
70	}
71}
72
73const RINGBUFFER_METADATA_WIDTH: usize = 24;
74
75pub fn encode_ringbuffer_metadata(metadata: &RingBufferMetadata) -> EncodedPodRow {
76	let mut bytes = Vec::with_capacity(RINGBUFFER_METADATA_WIDTH);
77	bytes.extend_from_slice(&metadata.count.to_be_bytes());
78	bytes.extend_from_slice(&metadata.head.to_be_bytes());
79	bytes.extend_from_slice(&metadata.tail.to_be_bytes());
80	EncodedPodRow::new(&bytes)
81}
82
83pub fn decode_ringbuffer_metadata(row: &EncodedPodRow) -> Result<RingBufferMetadata> {
84	let bytes = row.body();
85	if bytes.len() != RINGBUFFER_METADATA_WIDTH {
86		return_internal_error!(
87			"Ring buffer metadata is {} bytes wide, expected {}. This indicates a corrupt metadata row.",
88			bytes.len(),
89			RINGBUFFER_METADATA_WIDTH
90		)
91	}
92	Ok(RingBufferMetadata {
93		count: u64::from_be_bytes(bytes[0..8].try_into().unwrap()),
94		head: u64::from_be_bytes(bytes[8..16].try_into().unwrap()),
95		tail: u64::from_be_bytes(bytes[16..24].try_into().unwrap()),
96	})
97}
98
99#[cfg(test)]
100mod tests {
101	use super::*;
102
103	#[test]
104	fn head_and_tail_survive_a_round_trip_because_they_are_the_scan_range() {
105		let metadata = RingBufferMetadata {
106			count: 7,
107			head: 3,
108			tail: 11,
109		};
110
111		let row = encode_ringbuffer_metadata(&metadata);
112
113		assert_eq!(row.len(), RINGBUFFER_METADATA_WIDTH);
114		assert_eq!(decode_ringbuffer_metadata(&row).unwrap(), metadata);
115	}
116
117	#[test]
118	fn a_wrapped_buffer_keeps_head_ahead_of_tail_rather_than_being_normalised() {
119		let metadata = RingBufferMetadata {
120			count: u64::MAX,
121			head: 99,
122			tail: 0,
123		};
124
125		assert_eq!(decode_ringbuffer_metadata(&encode_ringbuffer_metadata(&metadata)).unwrap(), metadata);
126	}
127
128	#[test]
129	fn a_row_of_the_wrong_width_is_rejected_rather_than_misread_as_eviction_bounds() {
130		assert!(decode_ringbuffer_metadata(&EncodedPodRow::new(&[0u8; 23])).is_err());
131		assert!(decode_ringbuffer_metadata(&EncodedPodRow::new(&[0u8; 25])).is_err());
132		assert!(decode_ringbuffer_metadata(&EncodedPodRow::new(&[0u8; 40])).is_err());
133	}
134}