1use std::collections::HashMap;
16use std::sync::Arc;
17
18use rskafka::record::Record as KafkaRecord;
19use serde::{Deserialize, Serialize};
20use snafu::{OptionExt, ResultExt, ensure};
21use store_api::logstore::entry::{Entry, MultiplePartEntry, MultiplePartHeader, NaiveEntry};
22use store_api::logstore::provider::{KafkaProvider, Provider};
23use store_api::storage::RegionId;
24
25use crate::error::{
26 DecodeJsonSnafu, EncodeJsonSnafu, IllegalSequenceSnafu, MetaLengthExceededLimitSnafu,
27 MissingKeySnafu, MissingValueSnafu, Result,
28};
29use crate::kafka::{EntryId, NamespaceImpl};
30
31pub(crate) const VERSION: u32 = 0;
33
34pub(crate) const ESTIMATED_META_SIZE: usize = 256;
37
38#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq)]
47pub enum RecordType {
48 Full,
50 First,
52 Middle(usize),
55 Last,
57}
58
59#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
61pub struct RecordMeta {
62 version: u32,
64 pub tp: RecordType,
66 pub entry_id: EntryId,
68 pub ns: NamespaceImpl,
70}
71
72#[derive(Debug, Clone, PartialEq)]
82pub(crate) struct Record {
83 pub(crate) meta: RecordMeta,
85 data: Vec<u8>,
87}
88
89impl TryFrom<Record> for KafkaRecord {
90 type Error = crate::error::Error;
91
92 fn try_from(record: Record) -> Result<Self> {
93 let key = serde_json::to_vec(&record.meta).context(EncodeJsonSnafu)?;
94 ensure!(
95 key.len() < ESTIMATED_META_SIZE,
96 MetaLengthExceededLimitSnafu {
97 limit: ESTIMATED_META_SIZE,
98 actual: key.len()
99 }
100 );
101 Ok(KafkaRecord {
102 key: Some(key),
103 value: Some(record.data),
104 timestamp: chrono::Utc::now(),
105 headers: Default::default(),
106 })
107 }
108}
109
110impl TryFrom<KafkaRecord> for Record {
112 type Error = crate::error::Error;
113
114 fn try_from(kafka_record: KafkaRecord) -> Result<Self> {
115 let key = kafka_record.key.context(MissingKeySnafu)?;
116 let meta = serde_json::from_slice(&key).context(DecodeJsonSnafu)?;
117 let data = kafka_record.value.context(MissingValueSnafu)?;
118 Ok(Self { meta, data })
119 }
120}
121
122pub(crate) fn convert_to_kafka_records(entry: Entry) -> Result<Vec<KafkaRecord>> {
123 match entry {
124 Entry::Naive(entry) => Ok(vec![KafkaRecord::try_from(Record {
125 meta: RecordMeta {
126 version: VERSION,
127 tp: RecordType::Full,
128 entry_id: 0,
130 ns: NamespaceImpl {
131 region_id: entry.region_id.as_u64(),
132 topic: String::new(),
134 },
135 },
136 data: entry.data,
137 })?]),
138 Entry::MultiplePart(entry) => {
139 let mut entries = Vec::with_capacity(entry.parts.len());
140
141 for (idx, part) in entry.parts.into_iter().enumerate() {
142 let tp = match entry.headers[idx] {
143 MultiplePartHeader::First => RecordType::First,
144 MultiplePartHeader::Middle(i) => RecordType::Middle(i),
145 MultiplePartHeader::Last => RecordType::Last,
146 };
147 entries.push(KafkaRecord::try_from(Record {
148 meta: RecordMeta {
149 version: VERSION,
150 tp,
151 entry_id: 0,
153 ns: NamespaceImpl {
154 region_id: entry.region_id.as_u64(),
155 topic: String::new(),
156 },
157 },
158 data: part,
159 })?)
160 }
161 Ok(entries)
162 }
163 }
164}
165
166fn convert_to_naive_entry(provider: Arc<KafkaProvider>, record: Record) -> Entry {
167 let region_id = RegionId::from_u64(record.meta.ns.region_id);
168
169 Entry::Naive(NaiveEntry {
170 provider: Provider::Kafka(provider),
171 region_id,
172 entry_id: record.meta.entry_id,
173 data: record.data,
174 })
175}
176
177fn convert_to_multiple_entry(
178 provider: Arc<KafkaProvider>,
179 region_id: RegionId,
180 records: Vec<Record>,
181) -> Entry {
182 let mut headers = Vec::with_capacity(records.len());
183 let mut parts = Vec::with_capacity(records.len());
184 let entry_id = records.last().map(|r| r.meta.entry_id).unwrap_or_default();
185
186 for record in records {
187 let header = match record.meta.tp {
188 RecordType::Full => unreachable!(),
189 RecordType::First => MultiplePartHeader::First,
190 RecordType::Middle(i) => MultiplePartHeader::Middle(i),
191 RecordType::Last => MultiplePartHeader::Last,
192 };
193 headers.push(header);
194 parts.push(record.data);
195 }
196
197 Entry::MultiplePart(MultiplePartEntry {
198 provider: Provider::Kafka(provider),
199 region_id,
200 entry_id,
201 headers,
202 parts,
203 })
204}
205
206pub fn remaining_entries(
208 provider: &Arc<KafkaProvider>,
209 buffered_records: &mut HashMap<RegionId, Vec<Record>>,
210) -> Option<Vec<Entry>> {
211 if buffered_records.is_empty() {
212 None
213 } else {
214 let mut entries = Vec::with_capacity(buffered_records.len());
215 for (region_id, records) in buffered_records.drain() {
216 entries.push(convert_to_multiple_entry(
217 provider.clone(),
218 region_id,
219 records,
220 ));
221 }
222 Some(entries)
223 }
224}
225
226pub(crate) fn maybe_emit_entry(
242 provider: &Arc<KafkaProvider>,
243 record: Record,
244 buffered_records: &mut HashMap<RegionId, Vec<Record>>,
245) -> Result<Option<Entry>> {
246 let mut entry = None;
247 match record.meta.tp {
248 RecordType::Full => entry = Some(convert_to_naive_entry(provider.clone(), record)),
249 RecordType::First => {
250 let region_id = record.meta.ns.region_id.into();
251 let duplicate = buffered_records.get(®ion_id).is_some_and(|records| {
252 records.len() == 1
253 && records[0].meta.tp == RecordType::First
254 && records[0].data == record.data
255 });
256 if let Some(records) = buffered_records.insert(region_id, vec![record]) {
257 if duplicate {
258 return Ok(None);
260 }
261 entry = Some(convert_to_multiple_entry(
263 provider.clone(),
264 region_id,
265 records,
266 ))
267 }
268 }
269 RecordType::Middle(seq) => {
270 let region_id = record.meta.ns.region_id.into();
271 let records = buffered_records.entry(region_id).or_default();
272
273 if !records.is_empty() {
275 let last_record = records.last().unwrap();
277 if matches!(last_record.meta.tp, RecordType::Middle(last_seq) if last_seq == seq)
278 && last_record.data == record.data
279 {
280 return Ok(None);
281 }
282 let legal = match last_record.meta.tp {
283 RecordType::First => seq == 1,
285 RecordType::Middle(last_seq) => last_seq + 1 == seq,
287 _ => false,
289 };
290 ensure!(
291 legal,
292 IllegalSequenceSnafu {
293 error: format!(
294 "Illegal sequence of a middle record, last record: {:?}, incoming record: {:?}",
295 last_record.meta.tp, record.meta.tp
296 )
297 }
298 );
299 }
300
301 records.push(record);
302 }
303 RecordType::Last => {
304 let region_id = record.meta.ns.region_id.into();
305 if let Some(mut records) = buffered_records.remove(®ion_id) {
306 records.push(record);
307 entry = Some(convert_to_multiple_entry(
308 provider.clone(),
309 region_id,
310 records,
311 ));
312 } else {
313 }
316 }
317 }
318 Ok(entry)
319}
320
321#[cfg(test)]
322mod tests {
323 use std::assert_matches;
324 use std::sync::Arc;
325
326 use super::*;
327 use crate::error;
328
329 fn new_test_record(tp: RecordType, entry_id: EntryId, region_id: u64, data: Vec<u8>) -> Record {
330 Record {
331 meta: RecordMeta {
332 version: VERSION,
333 tp,
334 ns: NamespaceImpl {
335 region_id,
336 topic: "greptimedb_wal_topic".to_string(),
337 },
338 entry_id,
339 },
340 data,
341 }
342 }
343
344 #[test]
345 fn test_maybe_emit_entry_emit_naive_entry() {
346 let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
347 let region_id = RegionId::new(1, 1);
348 let mut buffer = HashMap::new();
349 let record = new_test_record(RecordType::Full, 1, region_id.as_u64(), vec![1; 100]);
350 let entry = maybe_emit_entry(&provider, record, &mut buffer)
351 .unwrap()
352 .unwrap();
353 assert_eq!(
354 entry,
355 Entry::Naive(NaiveEntry {
356 provider: Provider::Kafka(provider),
357 region_id,
358 entry_id: 1,
359 data: vec![1; 100]
360 })
361 );
362 }
363
364 #[test]
365 fn test_maybe_emit_entry_emit_incomplete_entry() {
366 let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
367 let region_id = RegionId::new(1, 1);
368 let mut buffer = HashMap::new();
370 let record = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
371 assert!(
372 maybe_emit_entry(&provider, record, &mut buffer)
373 .unwrap()
374 .is_none()
375 );
376 let record = new_test_record(RecordType::First, 2, region_id.as_u64(), vec![1; 100]);
377 assert!(
379 maybe_emit_entry(&provider, record, &mut buffer)
380 .unwrap()
381 .is_none()
382 );
383 assert_eq!(buffer[®ion_id][0].meta.entry_id, 2);
384
385 let record = new_test_record(RecordType::First, 3, region_id.as_u64(), vec![2; 100]);
386 let incomplete_entry = maybe_emit_entry(&provider, record, &mut buffer)
387 .unwrap()
388 .unwrap();
389
390 assert_eq!(
391 incomplete_entry,
392 Entry::MultiplePart(MultiplePartEntry {
393 provider: Provider::Kafka(provider.clone()),
394 region_id,
395 entry_id: 2,
396 headers: vec![MultiplePartHeader::First],
397 parts: vec![vec![1; 100]],
398 })
399 );
400
401 let mut buffer = HashMap::new();
403 let record = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![1; 100]);
404 assert!(
405 maybe_emit_entry(&provider, record, &mut buffer)
406 .unwrap()
407 .is_none()
408 );
409 assert!(buffer.is_empty());
410
411 let mut buffer = HashMap::new();
413 let record = new_test_record(RecordType::Middle(0), 1, region_id.as_u64(), vec![1; 100]);
414 assert!(
415 maybe_emit_entry(&provider, record, &mut buffer)
416 .unwrap()
417 .is_none()
418 );
419 let record = new_test_record(RecordType::First, 2, region_id.as_u64(), vec![2; 100]);
420 let incomplete_entry = maybe_emit_entry(&provider, record, &mut buffer)
421 .unwrap()
422 .unwrap();
423
424 assert_eq!(
425 incomplete_entry,
426 Entry::MultiplePart(MultiplePartEntry {
427 provider: Provider::Kafka(provider),
428 region_id,
429 entry_id: 1,
430 headers: vec![MultiplePartHeader::Middle(0)],
431 parts: vec![vec![1; 100]],
432 })
433 );
434 }
435
436 #[test]
437 fn test_maybe_emit_entry_illegal_seq() {
438 let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
439 let region_id = RegionId::new(1, 1);
440 let mut buffer = HashMap::new();
441 let record = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
442 assert!(
443 maybe_emit_entry(&provider, record, &mut buffer)
444 .unwrap()
445 .is_none()
446 );
447 let record = new_test_record(RecordType::Middle(2), 1, region_id.as_u64(), vec![2; 100]);
448 let err = maybe_emit_entry(&provider, record, &mut buffer).unwrap_err();
449 assert_matches!(err, error::Error::IllegalSequence { .. });
450
451 let mut buffer = HashMap::new();
452 let record = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
453 assert!(
454 maybe_emit_entry(&provider, record, &mut buffer)
455 .unwrap()
456 .is_none()
457 );
458 let record = new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]);
459 assert!(
460 maybe_emit_entry(&provider, record, &mut buffer)
461 .unwrap()
462 .is_none()
463 );
464 let record = new_test_record(RecordType::Middle(3), 1, region_id.as_u64(), vec![2; 100]);
465 let err = maybe_emit_entry(&provider, record, &mut buffer).unwrap_err();
466 assert_matches!(err, error::Error::IllegalSequence { .. });
467 }
468
469 #[test]
470 fn test_maybe_emit_entry_deduplicates_multipart_records() {
471 let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
472 let region_id = RegionId::new(1, 1);
473 let mut buffer = HashMap::new();
474
475 for record in [
476 new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]),
477 new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]),
478 new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]),
479 ] {
480 assert!(
481 maybe_emit_entry(&provider, record, &mut buffer)
482 .unwrap()
483 .is_none()
484 );
485 }
486
487 let last = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![3; 100]);
488 let entry = maybe_emit_entry(&provider, last, &mut buffer)
489 .unwrap()
490 .unwrap();
491
492 assert_eq!(
493 entry,
494 Entry::MultiplePart(MultiplePartEntry {
495 provider: Provider::Kafka(provider.clone()),
496 region_id,
497 entry_id: 1,
498 headers: vec![
499 MultiplePartHeader::First,
500 MultiplePartHeader::Middle(1),
501 MultiplePartHeader::Last,
502 ],
503 parts: vec![vec![1; 100], vec![2; 100], vec![3; 100]],
504 })
505 );
506 let duplicate = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![3; 100]);
507 assert!(
508 maybe_emit_entry(&provider, duplicate, &mut buffer)
509 .unwrap()
510 .is_none()
511 );
512 }
513
514 #[test]
515 fn test_maybe_emit_entry_rejects_conflicting_duplicate_middle_record() {
516 let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
517 let region_id = RegionId::new(1, 1);
518 let mut buffer = HashMap::new();
519
520 for record in [
521 new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]),
522 new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]),
523 ] {
524 assert!(
525 maybe_emit_entry(&provider, record, &mut buffer)
526 .unwrap()
527 .is_none()
528 );
529 }
530
531 let duplicate = new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![3; 100]);
532 let err = maybe_emit_entry(&provider, duplicate, &mut buffer).unwrap_err();
533
534 assert_matches!(err, error::Error::IllegalSequence { .. });
535 }
536
537 #[test]
538 fn test_maybe_emit_entry_preserves_multipart_entry_order_before_full_record() {
539 let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
540 let region_id = RegionId::new(1, 1);
541 let mut buffer = HashMap::new();
542
543 let first = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
544 assert!(
545 maybe_emit_entry(&provider, first, &mut buffer)
546 .unwrap()
547 .is_none()
548 );
549 let last = new_test_record(RecordType::Last, 2, region_id.as_u64(), vec![2; 100]);
550 let multipart_entry = maybe_emit_entry(&provider, last, &mut buffer)
551 .unwrap()
552 .unwrap();
553 let full = new_test_record(RecordType::Full, 3, region_id.as_u64(), vec![3; 100]);
554 let full_entry = maybe_emit_entry(&provider, full, &mut buffer)
555 .unwrap()
556 .unwrap();
557
558 assert_eq!(multipart_entry.entry_id(), 2);
559 assert_eq!(full_entry.entry_id(), 3);
560 }
561
562 #[test]
563 fn test_meta_size() {
564 let meta = RecordMeta {
565 version: VERSION,
566 tp: RecordType::Middle(usize::MAX),
567 entry_id: u64::MAX,
568 ns: NamespaceImpl {
569 region_id: RegionId::new(u32::MAX, u32::MAX).as_u64(),
570 topic: format!("greptime_kafka_cluster/1024/2048/{}", uuid::Uuid::new_v4()),
571 },
572 };
573 let serialized = serde_json::to_vec(&meta).unwrap();
574 assert!(serialized.len() < ESTIMATED_META_SIZE);
576 }
577}