Skip to main content

ursula_runtime/engine/
wal.rs

1use std::path::Path;
2use std::path::PathBuf;
3
4use serde::Deserialize;
5use serde::Serialize;
6use ursula_shard::BucketStreamId;
7use ursula_shard::ShardPlacement;
8use ursula_stream::StreamErrorCode;
9use ursula_stream::StreamSnapshot;
10
11use super::GroupAppendBatchFuture;
12use super::GroupAppendFuture;
13use super::GroupBootstrapStreamFuture;
14use super::GroupCloseStreamFuture;
15use super::GroupColdHotBacklogFuture;
16use super::GroupCreateStreamFuture;
17use super::GroupDeleteSnapshotFuture;
18use super::GroupDeleteStreamFuture;
19use super::GroupEngine;
20use super::GroupEngineCreateFuture;
21use super::GroupEngineError;
22use super::GroupEngineFactory;
23use super::GroupEngineMetrics;
24use super::GroupFlushColdFuture;
25use super::GroupForkRefFuture;
26use super::GroupGetStreamAttrsFuture;
27use super::GroupHeadStreamFuture;
28use super::GroupInstallSnapshotFuture;
29use super::GroupPlanColdFlushFuture;
30use super::GroupPlanNextColdFlushBatchFuture;
31use super::GroupPlanNextColdFlushFuture;
32use super::GroupPublishSnapshotFuture;
33use super::GroupReadSnapshotFuture;
34use super::GroupReadStreamFuture;
35use super::GroupSnapshotFuture;
36use super::GroupTouchStreamAccessFuture;
37use super::GroupUpdateStreamAttrsFuture;
38use super::GroupWriteResponse;
39use super::in_memory::InMemoryGroupEngine;
40use crate::cold_store::ColdStoreHandle;
41use crate::command::GroupSnapshot;
42use crate::command::GroupWriteCommand;
43use crate::journal;
44use crate::metrics::elapsed_ns;
45use crate::request::AppendBatchRequest;
46use crate::request::AppendRequest;
47use crate::request::BootstrapStreamRequest;
48use crate::request::CloseStreamRequest;
49use crate::request::ColdWriteAdmission;
50use crate::request::CreateStreamRequest;
51use crate::request::DeleteSnapshotRequest;
52use crate::request::DeleteStreamRequest;
53use crate::request::FlushColdRequest;
54use crate::request::GetStreamAttrsRequest;
55use crate::request::HeadStreamRequest;
56use crate::request::PlanColdFlushRequest;
57use crate::request::PlanGroupColdFlushRequest;
58use crate::request::PublishSnapshotRequest;
59use crate::request::ReadSnapshotRequest;
60use crate::request::ReadStreamRequest;
61use crate::request::StreamAppendCount;
62use crate::request::TouchStreamAccessResponse;
63use crate::request::UpdateStreamAttrsRequest;
64use crate::rt::time::Instant;
65
66#[derive(Debug, Clone)]
67pub struct WalGroupEngineFactory {
68    root: PathBuf,
69    cold_store: Option<ColdStoreHandle>,
70}
71
72impl WalGroupEngineFactory {
73    pub fn new(root: impl Into<PathBuf>) -> Self {
74        Self {
75            root: root.into(),
76            cold_store: None,
77        }
78    }
79
80    pub fn with_cold_store(root: impl Into<PathBuf>, cold_store: Option<ColdStoreHandle>) -> Self {
81        Self {
82            root: root.into(),
83            cold_store,
84        }
85    }
86}
87
88impl GroupEngineFactory for WalGroupEngineFactory {
89    fn create<'a>(
90        &'a self,
91        placement: ShardPlacement,
92        metrics: GroupEngineMetrics,
93    ) -> GroupEngineCreateFuture<'a> {
94        Box::pin(async move {
95            let engine: Box<dyn GroupEngine> = Box::new(WalGroupEngine::open(
96                &self.root,
97                placement,
98                metrics,
99                self.cold_store.clone(),
100            ));
101            Ok(engine)
102        })
103    }
104}
105
106pub struct WalGroupEngine {
107    inner: InMemoryGroupEngine,
108    log_path: PathBuf,
109    placement: ShardPlacement,
110    metrics: GroupEngineMetrics,
111    init_error: Option<String>,
112}
113
114#[derive(Debug, Clone, Serialize, Deserialize)]
115#[serde(tag = "wal_record", rename_all = "snake_case")]
116enum WalRecord {
117    Command {
118        command: Box<GroupWriteCommand>,
119    },
120    Snapshot {
121        group_commit_index: u64,
122        stream_snapshot: StreamSnapshot,
123        stream_append_counts: Vec<StreamAppendCount>,
124    },
125}
126
127impl WalGroupEngine {
128    fn open(
129        root: &Path,
130        placement: ShardPlacement,
131        metrics: GroupEngineMetrics,
132        cold_store: Option<ColdStoreHandle>,
133    ) -> Self {
134        let log_path = group_log_path(root, placement);
135        match replay_group_log(&log_path) {
136            Ok(mut inner) => {
137                inner.set_cold_store(cold_store);
138                Self {
139                    inner,
140                    log_path,
141                    placement,
142                    metrics,
143                    init_error: None,
144                }
145            }
146            Err(err) => Self {
147                inner: {
148                    let mut inner = InMemoryGroupEngine::default();
149                    inner.set_cold_store(cold_store);
150                    inner
151                },
152                log_path,
153                placement,
154                metrics,
155                init_error: Some(err.message().into_owned()),
156            },
157        }
158    }
159
160    fn ensure_ready(&self) -> Result<(), GroupEngineError> {
161        match &self.init_error {
162            Some(message) => Err(GroupEngineError::new(message.clone())),
163            None => Ok(()),
164        }
165    }
166
167    fn append_record(&self, command: &GroupWriteCommand) -> Result<(), GroupEngineError> {
168        self.append_records(std::slice::from_ref(command))
169    }
170
171    fn append_records(&self, commands: &[GroupWriteCommand]) -> Result<(), GroupEngineError> {
172        let records = commands
173            .iter()
174            .map(|command| WalRecord::Command {
175                command: Box::new(command.clone()),
176            })
177            .collect::<Vec<_>>();
178        self.append_wal_records(&records)
179    }
180
181    /// Frame each record into the WAL, `fsync` once, and meter the batch.
182    fn append_wal_records(&self, records: &[WalRecord]) -> Result<(), GroupEngineError> {
183        if records.is_empty() {
184            return Ok(());
185        }
186        let mut writer = journal::JournalWriter::new(!self.log_path.exists());
187        let write_started_at = Instant::now();
188        for record in records {
189            writer
190                .append::<journal::JsonCodec<WalRecord>>(&self.log_path, record)
191                .map_err(|err| {
192                    GroupEngineError::new(format!("write WAL '{}': {err}", self.log_path.display()))
193                })?;
194        }
195        let write_ns = elapsed_ns(write_started_at);
196        let sync_started_at = Instant::now();
197        writer.sync(&self.log_path).map_err(|err| {
198            GroupEngineError::new(format!("sync WAL '{}': {err}", self.log_path.display()))
199        })?;
200        self.metrics.record_wal_batch(
201            self.placement,
202            records.len(),
203            write_ns,
204            elapsed_ns(sync_started_at),
205        );
206        Ok(())
207    }
208
209    fn append_snapshot_record(&self, snapshot: &GroupSnapshot) -> Result<(), GroupEngineError> {
210        let record = WalRecord::Snapshot {
211            group_commit_index: snapshot.group_commit_index,
212            stream_snapshot: snapshot.stream_snapshot.clone(),
213            stream_append_counts: snapshot.stream_append_counts.clone(),
214        };
215        self.append_wal_records(std::slice::from_ref(&record))
216    }
217
218    fn commit_access_if_needed(
219        &mut self,
220        stream_id: &BucketStreamId,
221        now_ms: u64,
222        renew_ttl: bool,
223        placement: ShardPlacement,
224    ) -> Result<Option<TouchStreamAccessResponse>, GroupEngineError> {
225        if !self
226            .inner
227            .access_requires_write(stream_id, now_ms, renew_ttl)?
228        {
229            return Ok(None);
230        }
231        let command = GroupWriteCommand::TouchStreamAccess {
232            stream_id: stream_id.clone(),
233            now_ms,
234            renew_ttl,
235        };
236        let mut preview = self.inner.clone();
237        let response = match preview.apply_committed_write(command.clone(), placement)? {
238            GroupWriteResponse::TouchStreamAccess(response) => response,
239            other => {
240                return Err(GroupEngineError::new(format!(
241                    "unexpected touch stream access write response: {other:?}"
242                )));
243            }
244        };
245        if response.changed || response.expired {
246            self.append_record(&command)?;
247        }
248        self.inner = preview;
249        if response.expired {
250            return Err(GroupEngineError::stream(
251                StreamErrorCode::StreamNotFound,
252                format!("stream '{stream_id}' does not exist"),
253            ));
254        }
255        Ok(Some(response))
256    }
257}
258
259impl GroupEngine for WalGroupEngine {
260    fn create_stream<'a>(
261        &'a mut self,
262        request: CreateStreamRequest,
263        placement: ShardPlacement,
264    ) -> GroupCreateStreamFuture<'a> {
265        Box::pin(async move {
266            self.ensure_ready()?;
267            let command = GroupWriteCommand::from(request);
268            let mut preview = self.inner.clone();
269            let response = match preview.apply_committed_write(command.clone(), placement)? {
270                GroupWriteResponse::CreateStream(response) => response,
271                other => {
272                    return Err(GroupEngineError::new(format!(
273                        "unexpected create stream write response: {other:?}"
274                    )));
275                }
276            };
277            if !response.already_exists {
278                self.append_record(&command)?;
279            }
280            self.inner = preview;
281            Ok(response)
282        })
283    }
284
285    fn create_stream_with_cold_admission<'a>(
286        &'a mut self,
287        request: CreateStreamRequest,
288        placement: ShardPlacement,
289        admission: ColdWriteAdmission,
290    ) -> GroupCreateStreamFuture<'a> {
291        if !admission.is_enabled() {
292            return self.create_stream(request, placement);
293        }
294        Box::pin(async move {
295            self.ensure_ready()?;
296            let command = GroupWriteCommand::from(request.clone());
297            let mut preview = self.inner.clone();
298            let response =
299                preview.create_stream_with_admission_inner(request, placement, admission)?;
300            if !response.already_exists {
301                self.append_record(&command)?;
302            }
303            self.inner = preview;
304            Ok(response)
305        })
306    }
307
308    fn head_stream<'a>(
309        &'a mut self,
310        request: HeadStreamRequest,
311        placement: ShardPlacement,
312    ) -> GroupHeadStreamFuture<'a> {
313        Box::pin(async move {
314            self.ensure_ready()?;
315            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
316            self.inner.head_stream(request, placement).await
317        })
318    }
319
320    fn get_stream_attrs<'a>(
321        &'a mut self,
322        request: GetStreamAttrsRequest,
323        placement: ShardPlacement,
324    ) -> GroupGetStreamAttrsFuture<'a> {
325        Box::pin(async move {
326            self.ensure_ready()?;
327            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
328            self.inner.get_stream_attrs(request, placement).await
329        })
330    }
331
332    fn read_stream<'a>(
333        &'a mut self,
334        request: ReadStreamRequest,
335        placement: ShardPlacement,
336    ) -> GroupReadStreamFuture<'a> {
337        Box::pin(async move {
338            self.ensure_ready()?;
339            self.commit_access_if_needed(&request.stream_id, request.now_ms, true, placement)?;
340            self.inner.read_stream(request, placement).await
341        })
342    }
343
344    fn publish_snapshot<'a>(
345        &'a mut self,
346        request: PublishSnapshotRequest,
347        placement: ShardPlacement,
348    ) -> GroupPublishSnapshotFuture<'a> {
349        Box::pin(async move {
350            self.ensure_ready()?;
351            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
352            let command = GroupWriteCommand::from(request);
353            let mut preview = self.inner.clone();
354            let response = match preview.apply_committed_write(command.clone(), placement)? {
355                GroupWriteResponse::PublishSnapshot(response) => response,
356                other => {
357                    return Err(GroupEngineError::new(format!(
358                        "unexpected publish snapshot write response: {other:?}"
359                    )));
360                }
361            };
362            self.append_record(&command)?;
363            self.inner = preview;
364            Ok(response)
365        })
366    }
367
368    fn read_snapshot<'a>(
369        &'a mut self,
370        request: ReadSnapshotRequest,
371        placement: ShardPlacement,
372    ) -> GroupReadSnapshotFuture<'a> {
373        Box::pin(async move {
374            self.ensure_ready()?;
375            self.commit_access_if_needed(&request.stream_id, request.now_ms, true, placement)?;
376            self.inner.read_snapshot(request, placement).await
377        })
378    }
379
380    fn delete_snapshot<'a>(
381        &'a mut self,
382        request: DeleteSnapshotRequest,
383        placement: ShardPlacement,
384    ) -> GroupDeleteSnapshotFuture<'a> {
385        Box::pin(async move {
386            self.ensure_ready()?;
387            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
388            self.inner.delete_snapshot(request, placement).await
389        })
390    }
391
392    fn bootstrap_stream<'a>(
393        &'a mut self,
394        request: BootstrapStreamRequest,
395        placement: ShardPlacement,
396    ) -> GroupBootstrapStreamFuture<'a> {
397        Box::pin(async move {
398            self.ensure_ready()?;
399            self.commit_access_if_needed(&request.stream_id, request.now_ms, true, placement)?;
400            self.inner.bootstrap_stream(request, placement).await
401        })
402    }
403
404    fn touch_stream_access<'a>(
405        &'a mut self,
406        stream_id: BucketStreamId,
407        now_ms: u64,
408        renew_ttl: bool,
409        placement: ShardPlacement,
410    ) -> GroupTouchStreamAccessFuture<'a> {
411        Box::pin(async move {
412            self.ensure_ready()?;
413            let command = GroupWriteCommand::TouchStreamAccess {
414                stream_id,
415                now_ms,
416                renew_ttl,
417            };
418            let mut preview = self.inner.clone();
419            let response = match preview.apply_committed_write(command.clone(), placement)? {
420                GroupWriteResponse::TouchStreamAccess(response) => response,
421                other => {
422                    return Err(GroupEngineError::new(format!(
423                        "unexpected touch stream access write response: {other:?}"
424                    )));
425                }
426            };
427            if response.changed || response.expired {
428                self.append_record(&command)?;
429            }
430            self.inner = preview;
431            Ok(response)
432        })
433    }
434
435    fn update_stream_attrs<'a>(
436        &'a mut self,
437        request: UpdateStreamAttrsRequest,
438        placement: ShardPlacement,
439    ) -> GroupUpdateStreamAttrsFuture<'a> {
440        Box::pin(async move {
441            self.ensure_ready()?;
442            let command = GroupWriteCommand::from(request);
443            let mut preview = self.inner.clone();
444            let response = match preview.apply_committed_write(command.clone(), placement)? {
445                GroupWriteResponse::UpdateStreamAttrs(response) => response,
446                other => {
447                    return Err(GroupEngineError::new(format!(
448                        "unexpected update stream attrs write response: {other:?}"
449                    )));
450                }
451            };
452            if response.changed {
453                self.append_record(&command)?;
454            }
455            self.inner = preview;
456            Ok(response)
457        })
458    }
459
460    fn add_fork_ref<'a>(
461        &'a mut self,
462        stream_id: BucketStreamId,
463        now_ms: u64,
464        placement: ShardPlacement,
465    ) -> GroupForkRefFuture<'a> {
466        Box::pin(async move {
467            self.ensure_ready()?;
468            let command = GroupWriteCommand::AddForkRef { stream_id, now_ms };
469            let mut preview = self.inner.clone();
470            let response = match preview.apply_committed_write(command.clone(), placement)? {
471                GroupWriteResponse::AddForkRef(response) => response,
472                other => {
473                    return Err(GroupEngineError::new(format!(
474                        "unexpected add fork ref write response: {other:?}"
475                    )));
476                }
477            };
478            self.append_record(&command)?;
479            self.inner = preview;
480            Ok(response)
481        })
482    }
483
484    fn release_fork_ref<'a>(
485        &'a mut self,
486        stream_id: BucketStreamId,
487        placement: ShardPlacement,
488    ) -> GroupForkRefFuture<'a> {
489        Box::pin(async move {
490            self.ensure_ready()?;
491            let command = GroupWriteCommand::ReleaseForkRef { stream_id };
492            let mut preview = self.inner.clone();
493            let response = match preview.apply_committed_write(command.clone(), placement)? {
494                GroupWriteResponse::ReleaseForkRef(response) => response,
495                other => {
496                    return Err(GroupEngineError::new(format!(
497                        "unexpected release fork ref write response: {other:?}"
498                    )));
499                }
500            };
501            self.append_record(&command)?;
502            self.inner = preview;
503            Ok(response)
504        })
505    }
506
507    fn close_stream<'a>(
508        &'a mut self,
509        request: CloseStreamRequest,
510        placement: ShardPlacement,
511    ) -> GroupCloseStreamFuture<'a> {
512        Box::pin(async move {
513            self.ensure_ready()?;
514            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
515            let command = GroupWriteCommand::from(request);
516            let mut preview = self.inner.clone();
517            let response = match preview.apply_committed_write(command.clone(), placement)? {
518                GroupWriteResponse::CloseStream(response) => response,
519                other => {
520                    return Err(GroupEngineError::new(format!(
521                        "unexpected close stream write response: {other:?}"
522                    )));
523                }
524            };
525            self.append_record(&command)?;
526            self.inner = preview;
527            Ok(response)
528        })
529    }
530
531    fn delete_stream<'a>(
532        &'a mut self,
533        request: DeleteStreamRequest,
534        placement: ShardPlacement,
535    ) -> GroupDeleteStreamFuture<'a> {
536        Box::pin(async move {
537            self.ensure_ready()?;
538            let command = GroupWriteCommand::from(request);
539            let mut preview = self.inner.clone();
540            let response = match preview.apply_committed_write(command.clone(), placement)? {
541                GroupWriteResponse::DeleteStream(response) => response,
542                other => {
543                    return Err(GroupEngineError::new(format!(
544                        "unexpected delete stream write response: {other:?}"
545                    )));
546                }
547            };
548            self.append_record(&command)?;
549            self.inner = preview;
550            Ok(response)
551        })
552    }
553
554    fn append<'a>(
555        &'a mut self,
556        request: AppendRequest,
557        placement: ShardPlacement,
558    ) -> GroupAppendFuture<'a> {
559        Box::pin(async move {
560            self.ensure_ready()?;
561            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
562            let command = GroupWriteCommand::from(request);
563            let mut preview = self.inner.clone();
564            let response = match preview.apply_committed_write(command.clone(), placement)? {
565                GroupWriteResponse::Append(response) => response,
566                other => {
567                    return Err(GroupEngineError::new(format!(
568                        "unexpected append write response: {other:?}"
569                    )));
570                }
571            };
572            self.append_record(&command)?;
573            self.inner = preview;
574            Ok(response)
575        })
576    }
577
578    fn append_with_cold_admission<'a>(
579        &'a mut self,
580        request: AppendRequest,
581        placement: ShardPlacement,
582        admission: ColdWriteAdmission,
583    ) -> GroupAppendFuture<'a> {
584        if !admission.is_enabled() {
585            return self.append(request, placement);
586        }
587        Box::pin(async move {
588            self.ensure_ready()?;
589            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
590            let command = GroupWriteCommand::from(request.clone());
591            let mut preview = self.inner.clone();
592            let response = preview.append_with_admission_inner(request, placement, admission)?;
593            if !response.deduplicated {
594                self.append_record(&command)?;
595            }
596            self.inner = preview;
597            Ok(response)
598        })
599    }
600
601    fn append_batch<'a>(
602        &'a mut self,
603        request: AppendBatchRequest,
604        placement: ShardPlacement,
605    ) -> GroupAppendBatchFuture<'a> {
606        Box::pin(async move {
607            self.ensure_ready()?;
608            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
609            let command = GroupWriteCommand::from(request);
610            let mut preview = self.inner.clone();
611            let response = match preview.apply_committed_write(command.clone(), placement)? {
612                GroupWriteResponse::AppendBatch(response) => response,
613                other => {
614                    return Err(GroupEngineError::new(format!(
615                        "unexpected append batch write response: {other:?}"
616                    )));
617                }
618            };
619            if response
620                .items
621                .iter()
622                .any(|item| matches!(item, Ok(response) if !response.deduplicated))
623            {
624                self.append_record(&command)?;
625            }
626            self.inner = preview;
627            Ok(response)
628        })
629    }
630
631    fn append_batch_with_cold_admission<'a>(
632        &'a mut self,
633        request: AppendBatchRequest,
634        placement: ShardPlacement,
635        admission: ColdWriteAdmission,
636    ) -> GroupAppendBatchFuture<'a> {
637        if !admission.is_enabled() {
638            return self.append_batch(request, placement);
639        }
640        Box::pin(async move {
641            self.ensure_ready()?;
642            self.commit_access_if_needed(&request.stream_id, request.now_ms, false, placement)?;
643            let command = GroupWriteCommand::from(request.clone());
644            let mut preview = self.inner.clone();
645            let response =
646                preview.append_batch_with_admission_inner(request, placement, admission)?;
647            if response
648                .items
649                .iter()
650                .any(|item| matches!(item, Ok(response) if !response.deduplicated))
651            {
652                self.append_record(&command)?;
653            }
654            self.inner = preview;
655            Ok(response)
656        })
657    }
658
659    fn flush_cold<'a>(
660        &'a mut self,
661        request: FlushColdRequest,
662        placement: ShardPlacement,
663    ) -> GroupFlushColdFuture<'a> {
664        Box::pin(async move {
665            self.ensure_ready()?;
666            let command = GroupWriteCommand::from(request);
667            let mut preview = self.inner.clone();
668            let response = match preview.apply_committed_write(command.clone(), placement)? {
669                GroupWriteResponse::FlushCold(response) => response,
670                other => {
671                    return Err(GroupEngineError::new(format!(
672                        "unexpected flush cold write response: {other:?}"
673                    )));
674                }
675            };
676            self.append_record(&command)?;
677            self.inner = preview;
678            Ok(response)
679        })
680    }
681
682    fn plan_cold_flush<'a>(
683        &'a mut self,
684        request: PlanColdFlushRequest,
685        placement: ShardPlacement,
686    ) -> GroupPlanColdFlushFuture<'a> {
687        Box::pin(async move {
688            self.ensure_ready()?;
689            self.inner.plan_cold_flush(request, placement).await
690        })
691    }
692
693    fn plan_next_cold_flush<'a>(
694        &'a mut self,
695        request: PlanGroupColdFlushRequest,
696        placement: ShardPlacement,
697    ) -> GroupPlanNextColdFlushFuture<'a> {
698        Box::pin(async move {
699            self.ensure_ready()?;
700            self.inner.plan_next_cold_flush(request, placement).await
701        })
702    }
703
704    fn plan_next_cold_flush_batch<'a>(
705        &'a mut self,
706        request: PlanGroupColdFlushRequest,
707        placement: ShardPlacement,
708        max_candidates: usize,
709    ) -> GroupPlanNextColdFlushBatchFuture<'a> {
710        Box::pin(async move {
711            self.ensure_ready()?;
712            self.inner
713                .plan_next_cold_flush_batch(request, placement, max_candidates)
714                .await
715        })
716    }
717
718    fn cold_hot_backlog<'a>(
719        &'a mut self,
720        stream_id: BucketStreamId,
721        placement: ShardPlacement,
722    ) -> GroupColdHotBacklogFuture<'a> {
723        Box::pin(async move {
724            self.ensure_ready()?;
725            self.inner.cold_hot_backlog(stream_id, placement).await
726        })
727    }
728
729    fn snapshot<'a>(&'a mut self, placement: ShardPlacement) -> GroupSnapshotFuture<'a> {
730        Box::pin(async move {
731            self.ensure_ready()?;
732            self.inner.snapshot(placement).await
733        })
734    }
735
736    fn install_snapshot<'a>(
737        &'a mut self,
738        snapshot: GroupSnapshot,
739    ) -> GroupInstallSnapshotFuture<'a> {
740        Box::pin(async move {
741            self.ensure_ready()?;
742            let mut preview = self.inner.clone();
743            preview.install_snapshot(snapshot.clone()).await?;
744            self.append_snapshot_record(&snapshot)?;
745            self.inner = preview;
746            Ok(())
747        })
748    }
749}
750
751pub(crate) fn group_log_path(root: &Path, placement: ShardPlacement) -> PathBuf {
752    root.join(format!("core-{}", placement.core_id.0))
753        .join(format!("group-{}.jsonl", placement.raft_group_id.0))
754}
755
756fn replay_group_log(log_path: &Path) -> Result<InMemoryGroupEngine, GroupEngineError> {
757    let records = journal::replay::<journal::JsonCodec<WalRecord>>(log_path).map_err(|err| {
758        GroupEngineError::new(format!("read WAL '{}': {err}", log_path.display()))
759    })?;
760    let mut inner = InMemoryGroupEngine::default();
761    for (index, record) in records.into_iter().enumerate() {
762        match record {
763            WalRecord::Command { command } => {
764                inner
765                    .apply_replayed_write_command(*command)
766                    .map_err(|err| {
767                        GroupEngineError::new(format!(
768                            "replay WAL command '{}' record {}: {err}",
769                            log_path.display(),
770                            index + 1
771                        ))
772                    })?;
773            }
774            WalRecord::Snapshot {
775                group_commit_index,
776                stream_snapshot,
777                stream_append_counts,
778            } => {
779                inner
780                    .install_snapshot_parts(
781                        group_commit_index,
782                        stream_snapshot,
783                        stream_append_counts,
784                    )
785                    .map_err(|err| {
786                        GroupEngineError::new(format!(
787                            "replay WAL snapshot '{}' record {}: {err}",
788                            log_path.display(),
789                            index + 1
790                        ))
791                    })?;
792            }
793        }
794    }
795    Ok(inner)
796}