Skip to main content

ursula_runtime/engine/
wal.rs

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