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}