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 underlying: bool,
28	pub time: TimeSource,
29}
30
31impl RingBuffer {
32	pub fn name(&self) -> &str {
33		&self.name
34	}
35}
36
37#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
38pub struct RingBufferMetadata {
39	pub count: u64,
40	pub head: u64,
41	pub tail: u64,
42}
43
44#[derive(Debug, Clone, PartialEq)]
45pub struct PartitionedMetadata {
46	pub metadata: RingBufferMetadata,
47	pub partition_values: Vec<Value>,
48}
49
50impl RingBufferMetadata {
51	pub fn new() -> Self {
52		Self {
53			count: 0,
54			head: 1,
55			tail: 1,
56		}
57	}
58
59	pub fn is_full(&self, capacity: u64) -> bool {
60		self.count >= capacity
61	}
62
63	pub fn is_empty(&self) -> bool {
64		self.count == 0
65	}
66}
67
68impl Default for RingBufferMetadata {
69	fn default() -> Self {
70		Self::new()
71	}
72}
73
74const RINGBUFFER_METADATA_WIDTH: usize = 24;
75
76pub fn encode_ringbuffer_metadata(metadata: &RingBufferMetadata) -> EncodedPodRow {
77	let mut bytes = Vec::with_capacity(RINGBUFFER_METADATA_WIDTH);
78	bytes.extend_from_slice(&metadata.count.to_be_bytes());
79	bytes.extend_from_slice(&metadata.head.to_be_bytes());
80	bytes.extend_from_slice(&metadata.tail.to_be_bytes());
81	EncodedPodRow::new(&bytes)
82}
83
84pub fn decode_ringbuffer_metadata(row: &EncodedPodRow) -> Result<RingBufferMetadata> {
85	let bytes = row.body();
86	if bytes.len() != RINGBUFFER_METADATA_WIDTH {
87		return_internal_error!(
88			"Ring buffer metadata is {} bytes wide, expected {}. This indicates a corrupt metadata row.",
89			bytes.len(),
90			RINGBUFFER_METADATA_WIDTH
91		)
92	}
93	Ok(RingBufferMetadata {
94		count: u64::from_be_bytes(bytes[0..8].try_into().unwrap()),
95		head: u64::from_be_bytes(bytes[8..16].try_into().unwrap()),
96		tail: u64::from_be_bytes(bytes[16..24].try_into().unwrap()),
97	})
98}
99
100#[cfg(test)]
101mod tests {
102	use super::*;
103
104	#[test]
105	fn head_and_tail_survive_a_round_trip_because_they_are_the_scan_range() {
106		let metadata = RingBufferMetadata {
107			count: 7,
108			head: 3,
109			tail: 11,
110		};
111
112		let row = encode_ringbuffer_metadata(&metadata);
113
114		assert_eq!(row.len(), RINGBUFFER_METADATA_WIDTH);
115		assert_eq!(decode_ringbuffer_metadata(&row).unwrap(), metadata);
116	}
117
118	#[test]
119	fn a_wrapped_buffer_keeps_head_ahead_of_tail_rather_than_being_normalised() {
120		let metadata = RingBufferMetadata {
121			count: u64::MAX,
122			head: 99,
123			tail: 0,
124		};
125
126		assert_eq!(decode_ringbuffer_metadata(&encode_ringbuffer_metadata(&metadata)).unwrap(), metadata);
127	}
128
129	#[test]
130	fn a_row_of_the_wrong_width_is_rejected_rather_than_misread_as_eviction_bounds() {
131		assert!(decode_ringbuffer_metadata(&EncodedPodRow::new(&[0u8; 23])).is_err());
132		assert!(decode_ringbuffer_metadata(&EncodedPodRow::new(&[0u8; 25])).is_err());
133		assert!(decode_ringbuffer_metadata(&EncodedPodRow::new(&[0u8; 40])).is_err());
134	}
135}