1use std::collections::{HashMap, hash_map};
16use std::fmt::{Debug, Formatter};
17use std::sync::Arc;
18use std::time::Duration;
19
20use async_stream::stream;
21use common_runtime::{RepeatedTask, TaskFunction};
22use common_telemetry::{debug, error, info};
23use common_wal::config::raft_engine::RaftEngineConfig;
24use raft_engine::{Config, Engine, LogBatch, MessageExt, ReadableSize, RecoveryMode};
25use snafu::{OptionExt, ResultExt, ensure};
26use store_api::logstore::entry::{Entry, Id as EntryId, NaiveEntry};
27use store_api::logstore::provider::{Provider, RaftEngineProvider};
28use store_api::logstore::{AppendBatchResponse, LogStore, SendableEntryStream, WalIndex};
29use store_api::storage::RegionId;
30
31use crate::error::{
32 AddEntryLogBatchSnafu, DiscontinuousLogIndexSnafu, Error, FetchEntrySnafu,
33 IllegalNamespaceSnafu, IllegalStateSnafu, InvalidProviderSnafu, OverrideCompactedEntrySnafu,
34 RaftEngineSnafu, Result, StartWalTaskSnafu, StopWalTaskSnafu,
35};
36use crate::metrics;
37use crate::raft_engine::backend::SYSTEM_NAMESPACE;
38use crate::raft_engine::protos::logstore::{EntryImpl, NamespaceImpl};
39
40const NAMESPACE_PREFIX: &str = "$sys/";
41
42pub struct RaftEngineLogStore {
43 sync_write: bool,
44 sync_period: Option<Duration>,
45 read_batch_size: usize,
46 engine: Arc<Engine>,
47 gc_task: RepeatedTask<Error>,
48 sync_task: RepeatedTask<Error>,
49}
50
51pub struct PurgeExpiredFilesFunction {
52 pub engine: Arc<Engine>,
53}
54
55#[async_trait::async_trait]
56impl TaskFunction<Error> for PurgeExpiredFilesFunction {
57 fn name(&self) -> &str {
58 "RaftEngineLogStore-gc-task"
59 }
60
61 async fn call(&mut self) -> Result<()> {
62 match self.engine.purge_expired_files().context(RaftEngineSnafu) {
63 Ok(res) => {
64 let log_string = format!(
67 "Successfully purged logstore files, namespaces need compaction: {:?}",
68 res
69 );
70 if res.is_empty() {
71 debug!(log_string);
72 } else {
73 info!(log_string);
74 }
75 }
76 Err(e) => {
77 error!(e; "Failed to purge files in logstore");
78 }
79 }
80
81 Ok(())
82 }
83}
84
85pub struct SyncWalTaskFunction {
86 engine: Arc<Engine>,
87}
88
89#[async_trait::async_trait]
90impl TaskFunction<Error> for SyncWalTaskFunction {
91 async fn call(&mut self) -> std::result::Result<(), Error> {
92 let engine = self.engine.clone();
93 if let Err(e) = tokio::task::spawn_blocking(move || engine.sync()).await {
94 error!(e; "Failed to sync raft engine log files");
95 };
96 Ok(())
97 }
98
99 fn name(&self) -> &str {
100 "SyncWalTaskFunction"
101 }
102}
103
104impl SyncWalTaskFunction {
105 pub fn new(engine: Arc<Engine>) -> Self {
106 Self { engine }
107 }
108}
109
110impl RaftEngineLogStore {
111 pub async fn try_new(dir: String, config: &RaftEngineConfig) -> Result<Self> {
112 let raft_engine_config = Config {
113 dir,
114 purge_threshold: ReadableSize(config.purge_threshold.0),
115 recovery_mode: RecoveryMode::TolerateTailCorruption,
116 batch_compression_threshold: ReadableSize::kb(8),
117 target_file_size: ReadableSize(config.file_size.0),
118 enable_log_recycle: config.enable_log_recycle,
119 prefill_for_recycle: config.prefill_log_files,
120 recovery_threads: config.recovery_parallelism,
121 ..Default::default()
122 };
123 let engine = Arc::new(Engine::open(raft_engine_config).context(RaftEngineSnafu)?);
124 let gc_task = RepeatedTask::new(
125 config.purge_interval,
126 Box::new(PurgeExpiredFilesFunction {
127 engine: engine.clone(),
128 }),
129 );
130
131 let sync_task = RepeatedTask::new(
132 config.sync_period.unwrap_or(Duration::from_secs(5)),
133 Box::new(SyncWalTaskFunction::new(engine.clone())),
134 );
135
136 let log_store = Self {
137 sync_write: config.sync_write,
138 sync_period: config.sync_period,
139 read_batch_size: config.read_batch_size,
140 engine,
141 gc_task,
142 sync_task,
143 };
144 log_store.start()?;
145 Ok(log_store)
146 }
147
148 pub fn started(&self) -> bool {
149 self.gc_task.started()
150 }
151
152 fn start(&self) -> Result<()> {
153 self.gc_task
154 .start(common_runtime::global_runtime())
155 .context(StartWalTaskSnafu { name: "gc_task" })?;
156 self.sync_task
157 .start(common_runtime::global_runtime())
158 .context(StartWalTaskSnafu { name: "sync_task" })
159 }
160
161 pub fn span(&self, provider: &RaftEngineProvider) -> (Option<u64>, Option<u64>) {
162 (
163 self.engine.first_index(provider.id),
164 self.engine.last_index(provider.id),
165 )
166 }
167
168 fn entries_to_batch(
172 &self,
173 entries: Vec<Entry>,
174 ) -> Result<(LogBatch, HashMap<RegionId, EntryId>)> {
175 let mut entry_ids: HashMap<RegionId, EntryId> = HashMap::with_capacity(entries.len());
177 let mut batch = LogBatch::with_capacity(entries.len());
178
179 for e in entries {
180 let region_id = e.region_id();
181 let entry_id = e.entry_id();
182 match entry_ids.entry(region_id) {
183 hash_map::Entry::Occupied(mut o) => {
184 let prev = *o.get();
185 ensure!(
186 entry_id == prev + 1,
187 DiscontinuousLogIndexSnafu {
188 region_id,
189 last_index: prev,
190 attempt_index: entry_id
191 }
192 );
193 o.insert(entry_id);
194 }
195 hash_map::Entry::Vacant(v) => {
196 if let Some(first_index) = self.engine.first_index(region_id.as_u64()) {
198 ensure!(
200 entry_id > first_index,
201 OverrideCompactedEntrySnafu {
202 namespace: region_id,
203 first_index,
204 attempt_index: entry_id,
205 }
206 );
207 }
208 if let Some(last_index) = self.engine.last_index(region_id.as_u64()) {
210 ensure!(
211 entry_id == last_index + 1,
212 DiscontinuousLogIndexSnafu {
213 region_id,
214 last_index,
215 attempt_index: entry_id
216 }
217 );
218 }
219 v.insert(entry_id);
220 }
221 }
222 batch
223 .add_entries::<MessageType>(
224 region_id.as_u64(),
225 &[EntryImpl {
226 id: entry_id,
227 namespace_id: region_id.as_u64(),
228 data: e.into_bytes(),
229 ..Default::default()
230 }],
231 )
232 .context(AddEntryLogBatchSnafu)?;
233 }
234
235 Ok((batch, entry_ids))
236 }
237}
238
239impl Debug for RaftEngineLogStore {
240 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
241 f.debug_struct("RaftEngineLogsStore")
242 .field("sync_write", &self.sync_write)
243 .field("sync_period", &self.sync_period)
244 .field("read_batch_size", &self.read_batch_size)
245 .field("started", &self.gc_task.started())
246 .finish()
247 }
248}
249
250#[async_trait::async_trait]
251impl LogStore for RaftEngineLogStore {
252 type Error = Error;
253
254 async fn stop(&self) -> Result<()> {
255 self.gc_task
256 .stop()
257 .await
258 .context(StopWalTaskSnafu { name: "gc_task" })?;
259 self.sync_task
260 .stop()
261 .await
262 .context(StopWalTaskSnafu { name: "sync_task" })
263 }
264
265 async fn append_batch(&self, entries: Vec<Entry>) -> Result<AppendBatchResponse> {
268 metrics::METRIC_RAFT_ENGINE_APPEND_BATCH_BYTES_TOTAL.inc_by(
269 entries
270 .iter()
271 .map(|entry| entry.estimated_size())
272 .sum::<usize>() as u64,
273 );
274 let _timer = metrics::METRIC_RAFT_ENGINE_APPEND_BATCH_ELAPSED.start_timer();
275
276 ensure!(self.started(), IllegalStateSnafu);
277 if entries.is_empty() {
278 return Ok(AppendBatchResponse::default());
279 }
280
281 let (mut batch, last_entry_ids) = self.entries_to_batch(entries)?;
282 let _ = self
283 .engine
284 .write(&mut batch, self.sync_write)
285 .context(RaftEngineSnafu)?;
286
287 Ok(AppendBatchResponse { last_entry_ids })
288 }
289
290 async fn read(
293 &self,
294 provider: &Provider,
295 entry_id: EntryId,
296 _index: Option<WalIndex>,
297 ) -> Result<SendableEntryStream<'static, Entry, Self::Error>> {
298 let ns = provider
299 .as_raft_engine_provider()
300 .with_context(|| InvalidProviderSnafu {
301 expected: RaftEngineProvider::type_name(),
302 actual: provider.type_name(),
303 })?;
304 let namespace_id = ns.id;
305 let _timer = metrics::METRIC_RAFT_ENGINE_READ_ELAPSED.start_timer();
306
307 ensure!(self.started(), IllegalStateSnafu);
308 let engine = self.engine.clone();
309
310 let last_index = engine.last_index(namespace_id).unwrap_or(0);
311 let mut start_index =
312 entry_id.max(engine.first_index(namespace_id).unwrap_or(last_index + 1));
313
314 info!(
315 "Read logstore, namespace: {}, start: {}, span: {:?}",
316 namespace_id,
317 entry_id,
318 self.span(ns)
319 );
320 let max_batch_size = self.read_batch_size;
321 let (tx, mut rx) = tokio::sync::mpsc::channel(max_batch_size);
322 let _handle = common_runtime::spawn_global(async move {
323 while start_index <= last_index {
324 let mut vec = Vec::with_capacity(max_batch_size);
325 match engine
326 .fetch_entries_to::<MessageType>(
327 namespace_id,
328 start_index,
329 last_index + 1,
330 Some(max_batch_size),
331 &mut vec,
332 )
333 .context(FetchEntrySnafu {
334 ns: namespace_id,
335 start: start_index,
336 end: last_index,
337 max_size: max_batch_size,
338 }) {
339 Ok(_) => {
340 if let Some(last_entry) = vec.last() {
341 start_index = last_entry.id + 1;
342 }
343 if tx.send(Ok(vec)).await.is_err() {
345 break;
346 }
347 }
348 Err(e) => {
349 let _ = tx.send(Err(e)).await;
350 break;
351 }
352 }
353 }
354 });
355
356 let s = stream!({
357 while let Some(res) = rx.recv().await {
358 let res = res?;
359
360 yield Ok(res.into_iter().map(Entry::from).collect::<Vec<_>>());
361 }
362 });
363 Ok(Box::pin(s))
364 }
365
366 async fn create_namespace(&self, ns: &Provider) -> Result<()> {
367 let ns = ns
368 .as_raft_engine_provider()
369 .with_context(|| InvalidProviderSnafu {
370 expected: RaftEngineProvider::type_name(),
371 actual: ns.type_name(),
372 })?;
373 let namespace_id = ns.id;
374 ensure!(
375 namespace_id != SYSTEM_NAMESPACE,
376 IllegalNamespaceSnafu { ns: namespace_id }
377 );
378 ensure!(self.started(), IllegalStateSnafu);
379 let key = format!("{}{}", NAMESPACE_PREFIX, namespace_id)
380 .as_bytes()
381 .to_vec();
382 let mut batch = LogBatch::with_capacity(1);
383 batch
384 .put_message::<NamespaceImpl>(
385 SYSTEM_NAMESPACE,
386 key,
387 &NamespaceImpl {
388 id: namespace_id,
389 ..Default::default()
390 },
391 )
392 .context(RaftEngineSnafu)?;
393 let _ = self
394 .engine
395 .write(&mut batch, true)
396 .context(RaftEngineSnafu)?;
397 Ok(())
398 }
399
400 async fn delete_namespace(&self, ns: &Provider) -> Result<()> {
401 let ns = ns
402 .as_raft_engine_provider()
403 .with_context(|| InvalidProviderSnafu {
404 expected: RaftEngineProvider::type_name(),
405 actual: ns.type_name(),
406 })?;
407 let namespace_id = ns.id;
408 ensure!(
409 namespace_id != SYSTEM_NAMESPACE,
410 IllegalNamespaceSnafu { ns: namespace_id }
411 );
412 ensure!(self.started(), IllegalStateSnafu);
413 let key = format!("{}{}", NAMESPACE_PREFIX, namespace_id)
414 .as_bytes()
415 .to_vec();
416 let mut batch = LogBatch::with_capacity(1);
417 batch.delete(SYSTEM_NAMESPACE, key);
418 let _ = self
419 .engine
420 .write(&mut batch, true)
421 .context(RaftEngineSnafu)?;
422 Ok(())
423 }
424
425 async fn list_namespaces(&self) -> Result<Vec<Provider>> {
426 ensure!(self.started(), IllegalStateSnafu);
427 let mut namespaces: Vec<Provider> = vec![];
428 self.engine
429 .scan_messages::<NamespaceImpl, _>(
430 SYSTEM_NAMESPACE,
431 Some(NAMESPACE_PREFIX.as_bytes()),
432 None,
433 false,
434 |_, v| {
435 namespaces.push(Provider::RaftEngine(RaftEngineProvider { id: v.id }));
436 true
437 },
438 )
439 .context(RaftEngineSnafu)?;
440 Ok(namespaces)
441 }
442
443 fn entry(
444 &self,
445 data: Vec<u8>,
446 entry_id: EntryId,
447 region_id: RegionId,
448 provider: &Provider,
449 ) -> Result<Entry> {
450 debug_assert_eq!(
451 provider.as_raft_engine_provider().unwrap().id,
452 region_id.as_u64()
453 );
454 Ok(Entry::Naive(NaiveEntry {
455 provider: provider.clone(),
456 region_id,
457 entry_id,
458 data,
459 }))
460 }
461
462 async fn obsolete(
463 &self,
464 provider: &Provider,
465 _region_id: RegionId,
466 entry_id: EntryId,
467 ) -> Result<()> {
468 let ns = provider
469 .as_raft_engine_provider()
470 .with_context(|| InvalidProviderSnafu {
471 expected: RaftEngineProvider::type_name(),
472 actual: provider.type_name(),
473 })?;
474 let namespace_id = ns.id;
475 ensure!(self.started(), IllegalStateSnafu);
476 let obsoleted = self.engine.compact_to(namespace_id, entry_id + 1);
477 info!(
478 "Namespace {} obsoleted {} entries, compacted index: {}, span: {:?}",
479 namespace_id,
480 obsoleted,
481 entry_id,
482 self.span(ns)
483 );
484 Ok(())
485 }
486
487 async fn obsolete_all(&self, provider: &Provider, region_id: RegionId) -> Result<()> {
488 let latest_entry_id = self.latest_entry_id(provider)?;
489 self.obsolete(provider, region_id, latest_entry_id).await?;
490 self.delete_namespace(provider).await
491 }
492
493 fn latest_entry_id(&self, provider: &Provider) -> Result<EntryId> {
494 let ns = provider
495 .as_raft_engine_provider()
496 .with_context(|| InvalidProviderSnafu {
497 expected: RaftEngineProvider::type_name(),
498 actual: provider.type_name(),
499 })?;
500 let namespace_id = ns.id;
501 let last_index = self.engine.last_index(namespace_id).unwrap_or(0);
502 Ok(last_index)
503 }
504}
505
506#[derive(Debug, Clone)]
507struct MessageType;
508
509impl MessageExt for MessageType {
510 type Entry = EntryImpl;
511
512 fn index(e: &Self::Entry) -> u64 {
513 e.id
514 }
515}
516
517#[cfg(test)]
518impl RaftEngineLogStore {
519 async fn append(&self, entry: Entry) -> Result<store_api::logstore::AppendResponse> {
522 let response = self.append_batch(vec![entry]).await?;
523 if let Some((_, last_entry_id)) = response.last_entry_ids.into_iter().next() {
524 return Ok(store_api::logstore::AppendResponse { last_entry_id });
525 }
526 unreachable!()
527 }
528}
529
530#[cfg(test)]
531mod tests {
532 use std::collections::HashSet;
533 use std::time::Duration;
534
535 use common_base::readable_size::ReadableSize;
536 use common_telemetry::debug;
537 use common_test_util::temp_dir::{TempDir, create_temp_dir};
538 use futures_util::StreamExt;
539 use store_api::logstore::{LogStore, SendableEntryStream};
540
541 use super::*;
542 use crate::error::Error;
543 use crate::raft_engine::log_store::RaftEngineLogStore;
544 use crate::raft_engine::protos::logstore::EntryImpl;
545
546 #[tokio::test]
547 async fn test_open_logstore() {
548 let dir = create_temp_dir("raft-engine-logstore-test");
549 let logstore = RaftEngineLogStore::try_new(
550 dir.path().to_str().unwrap().to_string(),
551 &RaftEngineConfig::default(),
552 )
553 .await
554 .unwrap();
555 let namespaces = logstore.list_namespaces().await.unwrap();
556 assert_eq!(0, namespaces.len());
557 }
558
559 #[tokio::test]
560 async fn test_manage_namespace() {
561 let dir = create_temp_dir("raft-engine-logstore-test");
562 let logstore = RaftEngineLogStore::try_new(
563 dir.path().to_str().unwrap().to_string(),
564 &RaftEngineConfig::default(),
565 )
566 .await
567 .unwrap();
568 assert!(logstore.list_namespaces().await.unwrap().is_empty());
569
570 logstore
571 .create_namespace(&Provider::raft_engine_provider(42))
572 .await
573 .unwrap();
574 let namespaces = logstore.list_namespaces().await.unwrap();
575 assert_eq!(1, namespaces.len());
576 assert_eq!(Provider::raft_engine_provider(42), namespaces[0]);
577
578 logstore
579 .delete_namespace(&Provider::raft_engine_provider(42))
580 .await
581 .unwrap();
582 assert!(logstore.list_namespaces().await.unwrap().is_empty());
583 }
584
585 #[tokio::test]
586 async fn test_obsolete_all_removes_entries_and_namespace() {
587 let dir = create_temp_dir("raft-engine-logstore-test");
588 let logstore = RaftEngineLogStore::try_new(
589 dir.path().to_str().unwrap().to_string(),
590 &RaftEngineConfig::default(),
591 )
592 .await
593 .unwrap();
594 let region_id = RegionId::new(1, 1);
595 let provider = Provider::raft_engine_provider(region_id.as_u64());
596 logstore.create_namespace(&provider).await.unwrap();
597 for entry_id in 1..=3 {
598 logstore
599 .append(
600 EntryImpl::create(
601 entry_id,
602 region_id.as_u64(),
603 entry_id.to_string().into_bytes(),
604 )
605 .into(),
606 )
607 .await
608 .unwrap();
609 }
610
611 logstore.obsolete_all(&provider, region_id).await.unwrap();
612
613 assert_eq!(0, logstore.latest_entry_id(&provider).unwrap());
614 assert!(logstore.list_namespaces().await.unwrap().is_empty());
615 }
616
617 #[tokio::test]
618 async fn test_append_and_read() {
619 let dir = create_temp_dir("raft-engine-logstore-test");
620 let logstore = RaftEngineLogStore::try_new(
621 dir.path().to_str().unwrap().to_string(),
622 &RaftEngineConfig::default(),
623 )
624 .await
625 .unwrap();
626
627 let namespace_id = 1;
628 let cnt = 1024;
629 for i in 0..cnt {
630 let response = logstore
631 .append(
632 EntryImpl::create(i, namespace_id, i.to_string().as_bytes().to_vec()).into(),
633 )
634 .await
635 .unwrap();
636 assert_eq!(i, response.last_entry_id);
637 }
638 let mut entries = HashSet::with_capacity(1024);
639 let mut s = logstore
640 .read(&Provider::raft_engine_provider(1), 0, None)
641 .await
642 .unwrap();
643 while let Some(r) = s.next().await {
644 let vec = r.unwrap();
645 entries.extend(vec.into_iter().map(|e| e.entry_id()));
646 }
647 assert_eq!((0..cnt).collect::<HashSet<_>>(), entries);
648 }
649
650 async fn collect_entries(mut s: SendableEntryStream<'_, Entry, Error>) -> Vec<Entry> {
651 let mut res = vec![];
652 while let Some(r) = s.next().await {
653 res.extend(r.unwrap());
654 }
655 res
656 }
657
658 #[tokio::test]
659 async fn test_reopen() {
660 let dir = create_temp_dir("raft-engine-logstore-reopen-test");
661 {
662 let logstore = RaftEngineLogStore::try_new(
663 dir.path().to_str().unwrap().to_string(),
664 &RaftEngineConfig::default(),
665 )
666 .await
667 .unwrap();
668 assert!(
669 logstore
670 .append(EntryImpl::create(1, 1, "1".as_bytes().to_vec()).into())
671 .await
672 .is_ok()
673 );
674 let entries = logstore
675 .read(&Provider::raft_engine_provider(1), 1, None)
676 .await
677 .unwrap()
678 .collect::<Vec<_>>()
679 .await;
680 assert_eq!(1, entries.len());
681 logstore.stop().await.unwrap();
682 }
683
684 let logstore = RaftEngineLogStore::try_new(
685 dir.path().to_str().unwrap().to_string(),
686 &RaftEngineConfig::default(),
687 )
688 .await
689 .unwrap();
690
691 let entries = collect_entries(
692 logstore
693 .read(&Provider::raft_engine_provider(1), 1, None)
694 .await
695 .unwrap(),
696 )
697 .await;
698 assert_eq!(1, entries.len());
699 assert_eq!(1, entries[0].entry_id());
700 assert_eq!(1, entries[0].region_id().as_u64());
701 }
702
703 async fn wal_dir_usage(path: impl AsRef<str>) -> usize {
704 let mut size: usize = 0;
705 let mut read_dir = tokio::fs::read_dir(path.as_ref()).await.unwrap();
706 while let Ok(dir_entry) = read_dir.next_entry().await {
707 let Some(entry) = dir_entry else {
708 break;
709 };
710 if entry.file_type().await.unwrap().is_file() {
711 let file_name = entry.file_name();
712 let file_size = entry.metadata().await.unwrap().len() as usize;
713 debug!("File: {file_name:?}, size: {file_size}");
714 size += file_size;
715 }
716 }
717 size
718 }
719
720 async fn new_test_log_store(dir: &TempDir) -> RaftEngineLogStore {
721 let path = dir.path().to_str().unwrap().to_string();
722
723 let config = RaftEngineConfig {
724 file_size: ReadableSize::mb(2),
725 purge_threshold: ReadableSize::mb(4),
726 purge_interval: Duration::from_secs(5),
727 ..Default::default()
728 };
729
730 RaftEngineLogStore::try_new(path, &config).await.unwrap()
731 }
732
733 #[tokio::test]
734 async fn test_compaction() {
735 common_telemetry::init_default_ut_logging();
736 let dir = create_temp_dir("raft-engine-logstore-test");
737 let logstore = new_test_log_store(&dir).await;
738
739 let region_id = RegionId::new(1, 1);
740 let namespace_id = region_id.as_u64();
741 let namespace = Provider::raft_engine_provider(namespace_id);
742 for id in 0..4096 {
743 let entry = EntryImpl::create(id, namespace_id, [b'x'; 4096].to_vec()).into();
744 let _ = logstore.append(entry).await.unwrap();
745 }
746
747 let before_purge = wal_dir_usage(dir.path().to_str().unwrap()).await;
748 logstore
749 .obsolete(&namespace, region_id, 4000)
750 .await
751 .unwrap();
752
753 tokio::time::sleep(Duration::from_secs(6)).await;
754 let after_purge = wal_dir_usage(dir.path().to_str().unwrap()).await;
755 debug!(
756 "Before purge: {}, after purge: {}",
757 before_purge, after_purge
758 );
759 assert!(before_purge > after_purge);
760 }
761
762 #[tokio::test]
763 async fn test_obsolete() {
764 common_telemetry::init_default_ut_logging();
765 let dir = create_temp_dir("raft-engine-logstore-test");
766 let logstore = new_test_log_store(&dir).await;
767
768 let region_id = RegionId::new(1, 1);
769 let namespace_id = region_id.as_u64();
770 let namespace = Provider::raft_engine_provider(namespace_id);
771 for id in 0..1024 {
772 let entry = EntryImpl::create(id, namespace_id, [b'x'; 4096].to_vec()).into();
773 let _ = logstore.append(entry).await.unwrap();
774 }
775
776 logstore.obsolete(&namespace, region_id, 100).await.unwrap();
777 assert_eq!(101, logstore.engine.first_index(namespace_id).unwrap());
778
779 let res = logstore.read(&namespace, 100, None).await.unwrap();
780 let mut vec = collect_entries(res).await;
781 vec.sort_by(|a, b| a.entry_id().partial_cmp(&b.entry_id()).unwrap());
782 assert_eq!(101, vec.first().unwrap().entry_id());
783 }
784
785 #[tokio::test]
786 async fn test_append_batch() {
787 common_telemetry::init_default_ut_logging();
788 let dir = create_temp_dir("logstore-append-batch-test");
789 let logstore = new_test_log_store(&dir).await;
790
791 let entries = (0..8)
792 .flat_map(|ns_id| {
793 let data = [ns_id as u8].repeat(4096);
794 (0..16).map(move |idx| EntryImpl::create(idx, ns_id, data.clone()).into())
795 })
796 .collect();
797
798 logstore.append_batch(entries).await.unwrap();
799 for ns_id in 0..8 {
800 let namespace = &RaftEngineProvider::new(ns_id);
801 let (first, last) = logstore.span(namespace);
802 assert_eq!(0, first.unwrap());
803 assert_eq!(15, last.unwrap());
804 }
805 }
806
807 #[tokio::test]
808 async fn test_append_batch_interleaved() {
809 common_telemetry::init_default_ut_logging();
810 let dir = create_temp_dir("logstore-append-batch-test");
811 let logstore = new_test_log_store(&dir).await;
812 let entries = vec![
813 EntryImpl::create(0, 0, [b'0'; 4096].to_vec()).into(),
814 EntryImpl::create(1, 0, [b'0'; 4096].to_vec()).into(),
815 EntryImpl::create(0, 1, [b'1'; 4096].to_vec()).into(),
816 EntryImpl::create(2, 0, [b'0'; 4096].to_vec()).into(),
817 EntryImpl::create(1, 1, [b'1'; 4096].to_vec()).into(),
818 ];
819
820 logstore.append_batch(entries).await.unwrap();
821
822 assert_eq!(
823 (Some(0), Some(2)),
824 logstore.span(&RaftEngineProvider::new(0))
825 );
826 assert_eq!(
827 (Some(0), Some(1)),
828 logstore.span(&RaftEngineProvider::new(1))
829 );
830 }
831
832 #[tokio::test]
833 async fn test_append_batch_response() {
834 common_telemetry::init_default_ut_logging();
835 let dir = create_temp_dir("logstore-append-batch-test");
836 let logstore = new_test_log_store(&dir).await;
837
838 let entries = vec![
839 EntryImpl::create(0, 0, [b'0'; 4096].to_vec()).into(),
841 EntryImpl::create(0, 1, [b'1'; 4096].to_vec()).into(),
843 EntryImpl::create(1, 0, [b'1'; 4096].to_vec()).into(),
845 EntryImpl::create(1, 1, [b'0'; 4096].to_vec()).into(),
847 EntryImpl::create(2, 2, [b'2'; 4096].to_vec()).into(),
849 ];
850
851 let last_entry_ids = logstore.append_batch(entries).await.unwrap().last_entry_ids;
853 assert_eq!(last_entry_ids[&(0.into())], 1);
854 assert_eq!(last_entry_ids[&(1.into())], 1);
855 assert_eq!(last_entry_ids[&(2.into())], 2);
856 }
857}