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 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}