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}