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 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}