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