Skip to main content

kacrab_protocol/record/
entry.rs

1//! Individual records inside a [`crate::record::RecordBatch`].
2//!
3//! Wire layout (integers are zigzag varints, except `attributes` which is a
4//! fixed `i8`):
5//!
6//! ```text
7//! length (varint) | attributes (i8) | timestampDelta (varlong)
8//! offsetDelta (varint) | keyLen (varint) | key | valueLen (varint) | value
9//! headerCount (varint) | headers[…]
10//! ```
11
12use bytes::{Buf, Bytes, BytesMut};
13
14use super::{RecordError, RecordErrorKind, RecordHeader, Result};
15use crate::primitives::{
16    read_i8, read_signed_varint, read_signed_varlong, signed_varint_len, signed_varlong_len,
17    write_i8, write_signed_varint, write_signed_varlong,
18};
19
20/// A single record in a v2 [`crate::record::RecordBatch`].
21#[derive(Debug, Clone, PartialEq, Eq)]
22pub struct Record {
23    /// Reserved — currently always 0.
24    pub attributes: i8,
25    /// Signed varlong delta from the batch's `first_timestamp`.
26    pub timestamp_delta: i64,
27    /// Signed varint delta from the batch's `base_offset`.
28    pub offset_delta: i32,
29    /// Record key. `None` when the on-wire key length is `-1`.
30    pub key: Option<Bytes>,
31    /// Record value. `None` when the on-wire value length is `-1`.
32    pub value: Option<Bytes>,
33    /// Record headers.
34    pub headers: Vec<RecordHeader>,
35}
36
37impl Record {
38    /// Return the exact encoded length of this record, including its body
39    /// length prefix.
40    pub fn encoded_len(&self) -> Result<usize> {
41        let body_len = self.body_encoded_len()?;
42        let body_len_i32 = super::encode_len("record body", body_len)?;
43        super::add_encoded_len("record body", signed_varint_len(body_len_i32), body_len)
44    }
45
46    /// Encode this record into `buf` — body length + body bytes.
47    pub fn encode(&self, buf: &mut BytesMut) -> Result<()> {
48        let body_len = self.body_encoded_len()?;
49        write_signed_varint(buf, super::encode_len("record body", body_len)?);
50        self.encode_body(buf)
51    }
52
53    fn encode_body(&self, buf: &mut BytesMut) -> Result<()> {
54        write_i8(buf, self.attributes);
55        write_signed_varlong(buf, self.timestamp_delta);
56        write_signed_varint(buf, self.offset_delta);
57
58        super::write_nullable_bytes_field(buf, "record key", self.key.as_ref())?;
59        super::write_nullable_bytes_field(buf, "record value", self.value.as_ref())?;
60
61        let hdr_count = i32::try_from(self.headers.len()).map_err(|_| {
62            RecordError::unknown_offset(RecordErrorKind::LengthOverflow {
63                field: "header count",
64                got: self.headers.len(),
65                remaining: usize::try_from(i32::MAX).unwrap_or(usize::MAX),
66            })
67        })?;
68        write_signed_varint(buf, hdr_count);
69        for header in &self.headers {
70            header.encode(buf)?;
71        }
72        Ok(())
73    }
74
75    fn body_encoded_len(&self) -> Result<usize> {
76        let mut len = 1;
77        len = super::add_encoded_len(
78            "record timestamp delta",
79            len,
80            signed_varlong_len(self.timestamp_delta),
81        )?;
82        len = super::add_encoded_len(
83            "record offset delta",
84            len,
85            signed_varint_len(self.offset_delta),
86        )?;
87        len = super::add_encoded_len(
88            "record key",
89            len,
90            super::nullable_bytes_field_len("record key", self.key.as_ref())?,
91        )?;
92        len = super::add_encoded_len(
93            "record value",
94            len,
95            super::nullable_bytes_field_len("record value", self.value.as_ref())?,
96        )?;
97        let hdr_count = i32::try_from(self.headers.len()).map_err(|_| {
98            RecordError::unknown_offset(RecordErrorKind::LengthOverflow {
99                field: "header count",
100                got: self.headers.len(),
101                remaining: usize::try_from(i32::MAX).unwrap_or(usize::MAX),
102            })
103        })?;
104        len = super::add_encoded_len("header count", len, signed_varint_len(hdr_count))?;
105        for header in &self.headers {
106            len = super::add_encoded_len("record header", len, header.encoded_len()?)?;
107        }
108        Ok(len)
109    }
110
111    /// Decode one record from `buf`. Validates that the body length fits in
112    /// the remaining buffer before splitting.
113    pub fn decode(buf: &mut Bytes) -> Result<Self> {
114        let body_length = read_signed_varint(buf)?;
115        if body_length < 0 {
116            return Err(RecordError::unknown_offset(
117                RecordErrorKind::NegativeLength {
118                    field: "record body",
119                    length: body_length,
120                },
121            ));
122        }
123        let body_len = usize::try_from(body_length).map_err(|_| {
124            RecordError::unknown_offset(RecordErrorKind::LengthOverflow {
125                field: "record body",
126                got: usize::MAX,
127                remaining: buf.remaining(),
128            })
129        })?;
130        let remaining = buf.remaining();
131        if body_len > remaining {
132            return Err(RecordError::unknown_offset(
133                RecordErrorKind::LengthOverflow {
134                    field: "record body",
135                    got: body_len,
136                    remaining,
137                },
138            ));
139        }
140        let mut record_buf = buf.split_to(body_len);
141
142        let attributes = read_i8(&mut record_buf)?;
143        let timestamp_delta = read_signed_varlong(&mut record_buf)?;
144        let offset_delta = read_signed_varint(&mut record_buf)?;
145
146        let key = super::read_nullable_bytes_field(&mut record_buf, "record key")?;
147        let value = super::read_nullable_bytes_field(&mut record_buf, "record value")?;
148
149        let header_count = read_signed_varint(&mut record_buf)?;
150        if header_count < 0 {
151            return Err(RecordError::unknown_offset(
152                RecordErrorKind::NegativeLength {
153                    field: "header count",
154                    length: header_count,
155                },
156            ));
157        }
158        let header_count_usize = usize::try_from(header_count).map_err(|_| {
159            RecordError::unknown_offset(RecordErrorKind::LengthOverflow {
160                field: "header count",
161                got: usize::MAX,
162                remaining: record_buf.remaining(),
163            })
164        })?;
165        let mut headers = Vec::with_capacity(header_count_usize);
166        for _ in 0..header_count {
167            headers.push(RecordHeader::decode(&mut record_buf)?);
168        }
169
170        Ok(Self {
171            attributes,
172            timestamp_delta,
173            offset_delta,
174            key,
175            value,
176            headers,
177        })
178    }
179}
180
181#[cfg(test)]
182mod tests {
183    #![allow(
184        clippy::expect_used,
185        clippy::missing_assert_message,
186        reason = "Record encoding tests fail fastest with contextual expect calls."
187    )]
188
189    use bytes::{Bytes, BytesMut};
190
191    use super::{Record, RecordHeader};
192
193    #[test]
194    fn record_encoded_len_matches_encoded_bytes_with_headers_and_nulls() {
195        let record = Record {
196            attributes: 0,
197            timestamp_delta: 300,
198            offset_delta: 127,
199            key: Some(Bytes::from_static(b"key")),
200            value: None,
201            headers: vec![RecordHeader {
202                key: Bytes::from_static(b"h"),
203                value: Some(Bytes::from_static(b"value")),
204            }],
205        };
206        let encoded_len = record.encoded_len().expect("record encoded len");
207        let mut bytes = BytesMut::with_capacity(encoded_len);
208
209        record.encode(&mut bytes).expect("record encode");
210
211        assert_eq!(encoded_len, bytes.len());
212        assert_eq!(encoded_len, bytes.capacity());
213    }
214}