reifydb_core/interface/catalog/
ringbuffer.rs1use 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}