1use std::collections::HashMap;
2use std::sync::Arc;
3use std::sync::atomic::AtomicU64;
4use std::sync::atomic::Ordering;
5
6use bytes::Bytes;
7#[cfg(not(madsim))]
8use tokio::task::JoinSet;
9use ursula_shard::BucketStreamId;
10use ursula_shard::CoreId;
11use ursula_shard::RaftGroupId;
12use ursula_shard::ShardId;
13use ursula_shard::ShardPlacement;
14use ursula_shard::StaticShardMap;
15use ursula_stream::ColdChunkRef;
16use ursula_stream::ColdFlushCandidate;
17use ursula_stream::ColdGcEntry;
18use ursula_stream::ColdGcTarget;
19use ursula_stream::StreamErrorCode;
20
21use crate::admission::RaftUncommittedAdmission;
22use crate::admission::RaftUncommittedBytesTracker;
23use crate::cold_index::cold_index_prefix;
24use crate::cold_store::ColdStoreHandle;
25use crate::cold_store::ColdStoreInfo;
26use crate::cold_store::cold_chunk_prefix;
27use crate::cold_store::new_cold_chunk_path;
28use crate::command::GroupSnapshot;
29use crate::core_worker::CoreCommand;
30use crate::core_worker::CoreMailbox;
31use crate::core_worker::CoreWorker;
32use crate::core_worker::WaitReadCancel;
33use crate::engine::GroupEngineError;
34use crate::engine::GroupEngineFactory;
35use crate::engine::in_memory::InMemoryGroupEngineFactory;
36use crate::error::RuntimeError;
37use crate::error::map_fork_source_ref_error;
38use crate::metrics::COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS;
39use crate::metrics::RuntimeMailboxSnapshot;
40use crate::metrics::RuntimeMetrics;
41use crate::metrics::RuntimeMetricsInner;
42use crate::metrics::elapsed_ns;
43use crate::metrics::is_stale_cold_flush_candidate_error;
44use crate::request::AckColdGcResponse;
45use crate::request::AppendBatchRequest;
46use crate::request::AppendBatchResponse;
47use crate::request::AppendExternalRequest;
48use crate::request::AppendRequest;
49use crate::request::AppendResponse;
50use crate::request::BootstrapStreamRequest;
51use crate::request::BootstrapStreamResponse;
52use crate::request::CloseStreamRequest;
53use crate::request::CloseStreamResponse;
54use crate::request::ColdWriteAdmission;
55use crate::request::CreateStreamExternalRequest;
56use crate::request::CreateStreamRequest;
57use crate::request::CreateStreamResponse;
58use crate::request::DeleteSnapshotRequest;
59use crate::request::DeleteStreamRequest;
60use crate::request::DeleteStreamResponse;
61use crate::request::FlushColdRequest;
62use crate::request::FlushColdResponse;
63use crate::request::ForkRefResponse;
64use crate::request::HeadStreamRequest;
65use crate::request::HeadStreamResponse;
66use crate::request::PlanColdFlushRequest;
67use crate::request::PlanGroupColdFlushRequest;
68use crate::request::PublishSnapshotRequest;
69use crate::request::PublishSnapshotResponse;
70use crate::request::ReadSnapshotRequest;
71use crate::request::ReadSnapshotResponse;
72use crate::request::ReadStreamRequest;
73use crate::request::ReadStreamResponse;
74use crate::rt::sync::Semaphore;
75use crate::rt::sync::mpsc;
76use crate::rt::sync::oneshot;
77use crate::rt::time::Instant;
78use crate::trace::Traced;
79
80#[derive(Debug, Clone)]
81pub struct RuntimeConfig {
82 pub core_count: usize,
83 pub raft_group_count: usize,
84 pub mailbox_capacity: usize,
85 pub threading: RuntimeThreading,
86 pub cold_max_hot_bytes_per_group: Option<u64>,
87 pub raft_max_uncommitted_bytes_per_group: Option<u64>,
91 pub live_read_max_waiters_per_core: Option<u64>,
92}
93
94impl RuntimeConfig {
95 pub fn new(core_count: usize, raft_group_count: usize) -> Self {
96 #[cfg(not(madsim))]
97 let threading = RuntimeThreading::ThreadPerCore;
98 #[cfg(madsim)]
99 let threading = RuntimeThreading::HostedTokio;
100 Self {
101 core_count,
102 raft_group_count,
103 mailbox_capacity: 1024,
104 threading,
105 cold_max_hot_bytes_per_group: None,
106 raft_max_uncommitted_bytes_per_group: None,
107 live_read_max_waiters_per_core: Some(65_536),
108 }
109 }
110
111 pub fn with_cold_max_hot_bytes_per_group(mut self, value: Option<u64>) -> Self {
112 self.cold_max_hot_bytes_per_group = value;
113 self
114 }
115
116 pub fn with_raft_max_uncommitted_bytes_per_group(mut self, value: Option<u64>) -> Self {
117 self.raft_max_uncommitted_bytes_per_group = value;
118 self
119 }
120
121 pub fn with_live_read_max_waiters_per_core(mut self, value: Option<u64>) -> Self {
122 self.live_read_max_waiters_per_core = value;
123 self
124 }
125}
126
127#[derive(Debug, Clone, Copy, PartialEq, Eq)]
128pub enum RuntimeThreading {
129 #[cfg(not(madsim))]
130 ThreadPerCore,
131 HostedTokio,
132}
133
134#[derive(Debug, Clone)]
135pub struct ShardRuntime {
136 shard_map: StaticShardMap,
137 mailboxes: Vec<CoreMailbox>,
138 metrics: Arc<RuntimeMetricsInner>,
139 next_waiter_id: Arc<AtomicU64>,
140 cold_store: Option<ColdStoreHandle>,
141}
142
143impl ShardRuntime {
144 pub fn spawn(config: RuntimeConfig) -> Result<Self, RuntimeError> {
145 Self::spawn_with_engine_factory(config, InMemoryGroupEngineFactory::default())
146 }
147
148 pub fn spawn_with_engine_factory(
149 config: RuntimeConfig,
150 engine_factory: impl GroupEngineFactory,
151 ) -> Result<Self, RuntimeError> {
152 Self::spawn_with_engine_factory_and_cold_store(config, engine_factory, None)
153 }
154
155 pub fn spawn_with_engine_factory_and_cold_store(
156 config: RuntimeConfig,
157 engine_factory: impl GroupEngineFactory,
158 cold_store: Option<ColdStoreHandle>,
159 ) -> Result<Self, RuntimeError> {
160 let shard_map = StaticShardMap::new(config.core_count, config.raft_group_count)?;
161 let metrics = Arc::new(RuntimeMetricsInner::new(
162 usize::from(shard_map.core_count()),
163 usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
164 ));
165 let cold_write_admission = ColdWriteAdmission {
166 max_hot_bytes_per_group: config.cold_max_hot_bytes_per_group,
167 };
168 let raft_uncommitted_admission = RaftUncommittedAdmission {
169 max_uncommitted_bytes_per_group: config.raft_max_uncommitted_bytes_per_group,
170 };
171 let raft_uncommitted_bytes = Arc::new(RaftUncommittedBytesTracker::new(
172 usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
173 ));
174 let engine_factory: Arc<dyn GroupEngineFactory> = Arc::new(engine_factory);
175 let read_materialization = Arc::new(Semaphore::new(config.mailbox_capacity.max(1)));
176 let mut mailboxes = Vec::with_capacity(usize::from(shard_map.core_count()));
177 for raw_core_id in 0..shard_map.core_count() {
178 let core_id = CoreId(raw_core_id);
179 let (tx, rx) = mpsc::channel(config.mailbox_capacity.max(1));
180 let worker = CoreWorker {
181 core_id,
182 rx,
183 engine_factory: engine_factory.clone(),
184 groups: HashMap::new(),
185 metrics: metrics.clone(),
186 group_mailbox_capacity: config.mailbox_capacity.max(1),
187 cold_write_admission,
188 raft_uncommitted_admission,
189 raft_uncommitted_bytes: raft_uncommitted_bytes.clone(),
190 live_read_max_waiters_per_core: config.live_read_max_waiters_per_core,
191 read_materialization: read_materialization.clone(),
192 };
193 spawn_core_worker(config.threading, worker)?;
194 mailboxes.push(CoreMailbox { core_id, tx });
195 }
196 Ok(Self {
197 shard_map,
198 mailboxes,
199 metrics,
200 next_waiter_id: Arc::new(AtomicU64::new(1)),
201 cold_store,
202 })
203 }
204
205 pub fn locate(&self, stream_id: &BucketStreamId) -> ShardPlacement {
206 self.shard_map.locate(stream_id)
207 }
208
209 pub fn has_cold_store(&self) -> bool {
210 self.cold_store.is_some()
211 }
212
213 pub fn cold_store(&self) -> Option<ColdStoreHandle> {
214 self.cold_store.clone()
215 }
216
217 pub fn cold_store_info(&self) -> Option<ColdStoreInfo> {
218 self.cold_store
219 .as_ref()
220 .map(|cold_store| cold_store.info().clone())
221 }
222
223 pub async fn create_stream(
224 &self,
225 request: CreateStreamRequest,
226 ) -> Result<CreateStreamResponse, RuntimeError> {
227 if request.forked_from.is_some() {
228 return self.create_fork_stream(request).await;
229 }
230 self.create_stream_on_owner(request).await
231 }
232
233 pub async fn create_stream_external(
234 &self,
235 request: CreateStreamExternalRequest,
236 ) -> Result<CreateStreamResponse, RuntimeError> {
237 let placement = self.shard_map.locate(&request.stream_id);
238 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
239 let (response_tx, response_rx) = oneshot::channel();
240 self.send_core_command(
241 mailbox,
242 CoreCommand::CreateExternal {
243 request,
244 placement,
245 response_tx,
246 },
247 response_rx,
248 )
249 .await
250 }
251
252 async fn create_stream_on_owner(
253 &self,
254 request: CreateStreamRequest,
255 ) -> Result<CreateStreamResponse, RuntimeError> {
256 let placement = self.shard_map.locate(&request.stream_id);
257 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
258 let (response_tx, response_rx) = oneshot::channel();
259 self.send_core_command(
260 mailbox,
261 CoreCommand::CreateStream {
262 request,
263 placement,
264 response_tx,
265 },
266 response_rx,
267 )
268 .await
269 }
270
271 async fn create_fork_stream(
272 &self,
273 mut request: CreateStreamRequest,
274 ) -> Result<CreateStreamResponse, RuntimeError> {
275 let source_id = request
276 .forked_from
277 .clone()
278 .expect("forked_from checked before create_fork_stream");
279 let now_ms = request.now_ms;
280 let source_placement = self.shard_map.locate(&source_id);
281 let source_head = self
282 .head_stream(HeadStreamRequest {
283 stream_id: source_id.clone(),
284 now_ms,
285 })
286 .await
287 .map_err(|err| map_fork_source_ref_error(err, source_placement))?;
288
289 if request.content_type_explicit {
290 if request.content_type != source_head.content_type {
291 return Err(RuntimeError::group_engine(
292 source_placement,
293 GroupEngineError::stream(
294 StreamErrorCode::ContentTypeMismatch,
295 format!(
296 "fork content type '{}' does not match source content type '{}'",
297 request.content_type, source_head.content_type
298 ),
299 ),
300 ));
301 }
302 } else {
303 request.content_type.clone_from(&source_head.content_type);
304 }
305
306 let fork_offset = request.fork_offset.unwrap_or(source_head.tail_offset);
307 if fork_offset > source_head.tail_offset {
308 return Err(RuntimeError::group_engine(
309 source_placement,
310 GroupEngineError::stream(
311 StreamErrorCode::InvalidFork,
312 format!(
313 "fork offset {fork_offset} is beyond source stream '{}' tail {}",
314 source_id, source_head.tail_offset
315 ),
316 ),
317 ));
318 }
319
320 let max_len = usize::try_from(fork_offset).map_err(|_| {
321 RuntimeError::group_engine(
322 source_placement,
323 GroupEngineError::stream(
324 StreamErrorCode::InvalidFork,
325 format!("fork offset {fork_offset} cannot fit in memory on this host"),
326 ),
327 )
328 })?;
329 request.initial_payload = if fork_offset == 0 {
330 Bytes::new()
331 } else {
332 self.read_stream(ReadStreamRequest {
333 stream_id: source_id.clone(),
334 offset: 0,
335 max_len,
336 now_ms,
337 })
338 .await?
339 .payload
340 .into()
341 };
342 self.add_fork_ref_on_owner(source_id.clone(), now_ms)
343 .await
344 .map_err(|err| map_fork_source_ref_error(err, source_placement))?;
345 request.close_after = false;
346 request.stream_seq = None;
347 request.producer = None;
348 if request.stream_ttl_seconds.is_none() && request.stream_expires_at_ms.is_none() {
349 request.stream_ttl_seconds = source_head.stream_ttl_seconds;
350 request.stream_expires_at_ms = source_head.stream_expires_at_ms;
351 }
352 request.fork_offset = Some(fork_offset);
353 match self.create_stream_on_owner(request).await {
354 Ok(response) if response.already_exists => {
355 self.release_fork_ref_cascade(source_id).await?;
356 Ok(response)
357 }
358 Ok(response) => Ok(response),
359 Err(err) => {
360 let _ = self.release_fork_ref_cascade(source_id).await;
361 Err(err)
362 }
363 }
364 }
365
366 pub async fn head_stream(
367 &self,
368 request: HeadStreamRequest,
369 ) -> Result<HeadStreamResponse, RuntimeError> {
370 let placement = self.shard_map.locate(&request.stream_id);
371 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
372 let (response_tx, response_rx) = oneshot::channel();
373 self.send_core_command(
374 mailbox,
375 CoreCommand::HeadStream {
376 request,
377 placement,
378 response_tx,
379 },
380 response_rx,
381 )
382 .await
383 }
384
385 pub async fn read_stream(
386 &self,
387 request: ReadStreamRequest,
388 ) -> Result<ReadStreamResponse, RuntimeError> {
389 let placement = self.shard_map.locate(&request.stream_id);
390 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
391 let (response_tx, response_rx) = oneshot::channel();
392 self.send_core_command(
393 mailbox,
394 CoreCommand::ReadStream {
395 request,
396 placement,
397 response_tx,
398 },
399 response_rx,
400 )
401 .await
402 }
403
404 pub async fn publish_snapshot(
405 &self,
406 request: PublishSnapshotRequest,
407 ) -> Result<PublishSnapshotResponse, RuntimeError> {
408 let placement = self.shard_map.locate(&request.stream_id);
409 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
410 let (response_tx, response_rx) = oneshot::channel();
411 self.send_core_command(
412 mailbox,
413 CoreCommand::PublishSnapshot {
414 request,
415 placement,
416 response_tx,
417 },
418 response_rx,
419 )
420 .await
421 }
422
423 pub async fn read_snapshot(
424 &self,
425 request: ReadSnapshotRequest,
426 ) -> Result<ReadSnapshotResponse, RuntimeError> {
427 let placement = self.shard_map.locate(&request.stream_id);
428 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
429 let (response_tx, response_rx) = oneshot::channel();
430 self.send_core_command(
431 mailbox,
432 CoreCommand::ReadSnapshot {
433 request,
434 placement,
435 response_tx,
436 },
437 response_rx,
438 )
439 .await
440 }
441
442 pub async fn delete_snapshot(
443 &self,
444 request: DeleteSnapshotRequest,
445 ) -> Result<(), RuntimeError> {
446 let placement = self.shard_map.locate(&request.stream_id);
447 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
448 let (response_tx, response_rx) = oneshot::channel();
449 self.send_core_command(
450 mailbox,
451 CoreCommand::DeleteSnapshot {
452 request,
453 placement,
454 response_tx,
455 },
456 response_rx,
457 )
458 .await
459 }
460
461 pub async fn bootstrap_stream(
462 &self,
463 request: BootstrapStreamRequest,
464 ) -> Result<BootstrapStreamResponse, RuntimeError> {
465 let placement = self.shard_map.locate(&request.stream_id);
466 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
467 let (response_tx, response_rx) = oneshot::channel();
468 self.send_core_command(
469 mailbox,
470 CoreCommand::BootstrapStream {
471 request,
472 placement,
473 response_tx,
474 },
475 response_rx,
476 )
477 .await
478 }
479
480 pub async fn wait_read_stream(
481 &self,
482 request: ReadStreamRequest,
483 ) -> Result<ReadStreamResponse, RuntimeError> {
484 let placement = self.shard_map.locate(&request.stream_id);
485 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
486 let waiter_id = self.next_waiter_id.fetch_add(1, Ordering::Relaxed);
487 let stream_id = request.stream_id.clone();
488 let (response_tx, response_rx) = oneshot::channel();
489 self.enqueue_core_command(mailbox, CoreCommand::WaitRead {
490 request,
491 placement,
492 waiter_id,
493 response_tx,
494 })
495 .await?;
496 let mut cancel = WaitReadCancel::new(mailbox.tx.clone(), stream_id, placement, waiter_id);
497 let response = response_rx
498 .await
499 .map_err(|_| RuntimeError::ResponseDropped {
500 core_id: mailbox.core_id,
501 })?;
502 cancel.disarm();
503 response
504 }
505
506 pub async fn require_local_live_read_owner(
507 &self,
508 stream_id: &BucketStreamId,
509 ) -> Result<(), RuntimeError> {
510 let placement = self.shard_map.locate(stream_id);
511 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
512 let (response_tx, response_rx) = oneshot::channel();
513 self.send_core_command(
514 mailbox,
515 CoreCommand::RequireLiveReadOwner {
516 placement,
517 response_tx,
518 },
519 response_rx,
520 )
521 .await
522 }
523
524 pub async fn close_stream(
525 &self,
526 request: CloseStreamRequest,
527 ) -> Result<CloseStreamResponse, RuntimeError> {
528 let placement = self.shard_map.locate(&request.stream_id);
529 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
530 let (response_tx, response_rx) = oneshot::channel();
531 self.send_core_command(
532 mailbox,
533 CoreCommand::CloseStream {
534 request,
535 placement,
536 response_tx,
537 },
538 response_rx,
539 )
540 .await
541 }
542
543 pub async fn delete_stream(
544 &self,
545 request: DeleteStreamRequest,
546 ) -> Result<DeleteStreamResponse, RuntimeError> {
547 let response = self.delete_stream_on_owner(request).await?;
548 if let Some(parent_to_release) = response.parent_to_release.clone() {
549 self.release_fork_ref_cascade(parent_to_release).await?;
550 }
551 Ok(response)
552 }
553
554 async fn delete_stream_on_owner(
555 &self,
556 request: DeleteStreamRequest,
557 ) -> Result<DeleteStreamResponse, RuntimeError> {
558 let placement = self.shard_map.locate(&request.stream_id);
559 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
560 let (response_tx, response_rx) = oneshot::channel();
561 self.send_core_command(
562 mailbox,
563 CoreCommand::DeleteStream {
564 request,
565 placement,
566 response_tx,
567 },
568 response_rx,
569 )
570 .await
571 }
572
573 async fn add_fork_ref_on_owner(
574 &self,
575 stream_id: BucketStreamId,
576 now_ms: u64,
577 ) -> Result<ForkRefResponse, RuntimeError> {
578 let placement = self.shard_map.locate(&stream_id);
579 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
580 let (response_tx, response_rx) = oneshot::channel();
581 self.send_core_command(
582 mailbox,
583 CoreCommand::AddForkRef {
584 stream_id,
585 now_ms,
586 placement,
587 response_tx,
588 },
589 response_rx,
590 )
591 .await
592 }
593
594 async fn release_fork_ref_on_owner(
595 &self,
596 stream_id: BucketStreamId,
597 ) -> Result<ForkRefResponse, RuntimeError> {
598 let placement = self.shard_map.locate(&stream_id);
599 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
600 let (response_tx, response_rx) = oneshot::channel();
601 self.send_core_command(
602 mailbox,
603 CoreCommand::ReleaseForkRef {
604 stream_id,
605 placement,
606 response_tx,
607 },
608 response_rx,
609 )
610 .await
611 }
612
613 async fn release_fork_ref_cascade(
614 &self,
615 stream_id: BucketStreamId,
616 ) -> Result<(), RuntimeError> {
617 let mut next = Some(stream_id);
618 while let Some(current) = next {
619 let response = self.release_fork_ref_on_owner(current).await?;
620 next = response.parent_to_release;
621 }
622 Ok(())
623 }
624
625 pub async fn flush_cold(
626 &self,
627 request: FlushColdRequest,
628 ) -> Result<FlushColdResponse, RuntimeError> {
629 let placement = self.shard_map.locate(&request.stream_id);
630 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
631 let (response_tx, response_rx) = oneshot::channel();
632 self.send_core_command(
633 mailbox,
634 CoreCommand::FlushCold {
635 request,
636 placement,
637 response_tx,
638 },
639 response_rx,
640 )
641 .await
642 }
643
644 pub async fn append_external(
645 &self,
646 request: AppendExternalRequest,
647 ) -> Result<AppendResponse, RuntimeError> {
648 let placement = self.shard_map.locate(&request.stream_id);
649 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
650 let (response_tx, response_rx) = oneshot::channel();
651 self.send_core_command(
652 mailbox,
653 CoreCommand::AppendExternal {
654 request,
655 placement,
656 response_tx,
657 },
658 response_rx,
659 )
660 .await
661 }
662
663 pub async fn plan_cold_flush(
664 &self,
665 request: PlanColdFlushRequest,
666 ) -> Result<Option<ColdFlushCandidate>, RuntimeError> {
667 let placement = self.shard_map.locate(&request.stream_id);
668 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
669 let (response_tx, response_rx) = oneshot::channel();
670 self.send_core_command(
671 mailbox,
672 CoreCommand::PlanColdFlush {
673 request,
674 placement,
675 response_tx,
676 },
677 response_rx,
678 )
679 .await
680 }
681
682 pub async fn flush_cold_once(
683 &self,
684 request: PlanColdFlushRequest,
685 ) -> Result<Option<FlushColdResponse>, RuntimeError> {
686 let Some(candidate) = self.plan_cold_flush(request).await? else {
687 return Ok(None);
688 };
689 self.flush_cold_candidate(candidate).await.map(Some)
690 }
691
692 pub async fn plan_next_cold_flush(
693 &self,
694 raft_group_id: RaftGroupId,
695 request: PlanGroupColdFlushRequest,
696 ) -> Result<Option<ColdFlushCandidate>, RuntimeError> {
697 let placement = self.placement_for_group(raft_group_id)?;
698 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
699 let (response_tx, response_rx) = oneshot::channel();
700 self.send_core_command(
701 mailbox,
702 CoreCommand::PlanNextColdFlush {
703 request,
704 placement,
705 response_tx,
706 },
707 response_rx,
708 )
709 .await
710 }
711
712 pub async fn plan_next_cold_flush_batch(
713 &self,
714 raft_group_id: RaftGroupId,
715 request: PlanGroupColdFlushRequest,
716 max_candidates: usize,
717 ) -> Result<Vec<ColdFlushCandidate>, RuntimeError> {
718 let placement = self.placement_for_group(raft_group_id)?;
719 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
720 let (response_tx, response_rx) = oneshot::channel();
721 self.send_core_command(
722 mailbox,
723 CoreCommand::PlanNextColdFlushBatch {
724 request,
725 placement,
726 max_candidates,
727 response_tx,
728 },
729 response_rx,
730 )
731 .await
732 }
733
734 pub async fn flush_cold_group_once(
735 &self,
736 raft_group_id: RaftGroupId,
737 request: PlanGroupColdFlushRequest,
738 ) -> Result<Option<FlushColdResponse>, RuntimeError> {
739 let Some(candidate) = self.plan_next_cold_flush(raft_group_id, request).await? else {
740 return Ok(None);
741 };
742 match self.flush_cold_candidate(candidate).await {
743 Ok(response) => Ok(Some(response)),
744 Err(err) if is_stale_cold_flush_candidate_error(&err) => Ok(None),
745 Err(err) => Err(err),
746 }
747 }
748
749 pub async fn flush_cold_group_batch_once(
750 &self,
751 raft_group_id: RaftGroupId,
752 request: PlanGroupColdFlushRequest,
753 max_candidates: usize,
754 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
755 let candidates = self
756 .plan_next_cold_flush_batch(raft_group_id, request, max_candidates)
757 .await?;
758 if candidates.is_empty() {
759 return Ok(Vec::new());
760 }
761 self.flush_cold_candidates_batch(candidates).await
762 }
763
764 async fn flush_cold_candidate(
765 &self,
766 candidate: ColdFlushCandidate,
767 ) -> Result<FlushColdResponse, RuntimeError> {
768 let Some(cold_store) = self.cold_store.as_ref() else {
769 return Err(RuntimeError::ColdStoreConfig {
770 message: "URSULA_COLD_BACKEND must be configured before flushing cold chunks"
771 .to_owned(),
772 });
773 };
774 let path = new_cold_chunk_path(
775 &candidate.stream_id,
776 candidate.start_offset,
777 candidate.end_offset,
778 );
779 let upload_started_at = Instant::now();
780 let object_size = match cold_store.write_chunk(&path, &candidate.payload).await {
781 Ok(object_size) => object_size,
782 Err(err) => {
783 self.metrics.record_cold_flush_write_error();
788 return Err(RuntimeError::ColdStoreIo {
789 message: err.to_string(),
790 });
791 }
792 };
793 self.metrics
794 .record_cold_upload(object_size, elapsed_ns(upload_started_at));
795 let chunk = ColdChunkRef {
796 start_offset: candidate.start_offset,
797 end_offset: candidate.end_offset,
798 s3_path: path.clone(),
799 object_size,
800 };
801 let publish_started_at = Instant::now();
802 let publish = self
803 .flush_cold(FlushColdRequest {
804 stream_id: candidate.stream_id,
805 chunk,
806 })
807 .await;
808 match publish {
809 Ok(response) => {
810 self.metrics
811 .record_cold_publish(object_size, elapsed_ns(publish_started_at));
812 Ok(response)
813 }
814 Err(err) => Err(err),
815 }
816 }
817
818 pub(crate) async fn flush_cold_candidates_batch(
819 &self,
820 candidates: Vec<ColdFlushCandidate>,
821 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
822 let mut responses = Vec::with_capacity(candidates.len());
823 for candidate in candidates {
824 match self.flush_cold_candidate(candidate).await {
825 Ok(response) => responses.push(response),
826 Err(err) if is_stale_cold_flush_candidate_error(&err) => {}
827 Err(err) => return Err(err),
828 }
829 }
830 Ok(responses)
831 }
832
833 #[cfg(madsim)]
834 pub async fn flush_cold_candidates_batch_for_simulation(
835 &self,
836 candidates: Vec<ColdFlushCandidate>,
837 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
838 self.flush_cold_candidates_batch(candidates).await
839 }
840
841 pub async fn flush_cold_all_groups_once(
842 &self,
843 request: PlanGroupColdFlushRequest,
844 ) -> Result<usize, RuntimeError> {
845 self.flush_cold_all_groups_once_bounded(request, 1).await
846 }
847
848 pub async fn flush_cold_all_groups_once_bounded(
849 &self,
850 request: PlanGroupColdFlushRequest,
851 max_concurrency: usize,
852 ) -> Result<usize, RuntimeError> {
853 let max_concurrency = max_concurrency.max(1);
854 if max_concurrency == 1 {
855 return self.flush_cold_all_groups_once_serial(request).await;
856 }
857 #[cfg(madsim)]
858 {
859 return self.flush_cold_all_groups_once_serial(request).await;
860 }
861 #[cfg(not(madsim))]
862 {
863 let mut flushed = 0;
864 let mut next_group_id = 0;
865 let group_count = self.shard_map.raft_group_count();
866 let mut tasks = JoinSet::new();
867
868 while next_group_id < group_count || !tasks.is_empty() {
869 while next_group_id < group_count && tasks.len() < max_concurrency {
870 let runtime = self.clone();
871 let request = request.clone();
872 let group_id = RaftGroupId(next_group_id);
873 next_group_id += 1;
874 tasks.spawn(async move {
875 runtime
876 .flush_cold_group_batch_once(
877 group_id,
878 request,
879 COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
880 )
881 .await
882 .map(|responses| responses.len())
883 });
884 }
885 if let Some(result) = tasks.join_next().await {
886 match result {
887 Ok(Ok(count)) => flushed += count,
888 Ok(Err(err)) => return Err(err),
889 Err(err) => {
890 return Err(RuntimeError::ColdStoreIo {
891 message: format!("cold flush task failed: {err}"),
892 });
893 }
894 }
895 }
896 }
897 Ok(flushed)
898 }
899 }
900
901 async fn flush_cold_all_groups_once_serial(
902 &self,
903 request: PlanGroupColdFlushRequest,
904 ) -> Result<usize, RuntimeError> {
905 let mut flushed = 0;
906 for group_id in 0..self.shard_map.raft_group_count() {
907 flushed += self
908 .flush_cold_group_batch_once(
909 RaftGroupId(group_id),
910 request.clone(),
911 COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
912 )
913 .await?
914 .len();
915 }
916 Ok(flushed)
917 }
918
919 async fn plan_cold_gc(
920 &self,
921 raft_group_id: RaftGroupId,
922 max: usize,
923 ) -> Result<Vec<ColdGcEntry>, RuntimeError> {
924 let placement = self.placement_for_group(raft_group_id)?;
925 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
926 let (response_tx, response_rx) = oneshot::channel();
927 self.send_core_command(
928 mailbox,
929 CoreCommand::PlanColdGc {
930 max,
931 placement,
932 response_tx,
933 },
934 response_rx,
935 )
936 .await
937 }
938
939 async fn ack_cold_gc(
940 &self,
941 raft_group_id: RaftGroupId,
942 up_to_seq: u64,
943 ) -> Result<AckColdGcResponse, RuntimeError> {
944 let placement = self.placement_for_group(raft_group_id)?;
945 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
946 let (response_tx, response_rx) = oneshot::channel();
947 self.send_core_command(
948 mailbox,
949 CoreCommand::AckColdGc {
950 up_to_seq,
951 placement,
952 response_tx,
953 },
954 response_rx,
955 )
956 .await
957 }
958
959 pub async fn run_cold_gc_group_once(
964 &self,
965 raft_group_id: RaftGroupId,
966 max_entries: usize,
967 ) -> Result<usize, RuntimeError> {
968 let Some(cold_store) = self.cold_store.as_ref() else {
969 return Ok(0);
970 };
971 let entries = self.plan_cold_gc(raft_group_id, max_entries).await?;
972 if entries.is_empty() {
973 return Ok(0);
974 }
975 let mut acked_seq = None;
976 let mut reclaimed = 0usize;
977 for entry in entries {
980 let result = match &entry.target {
981 ColdGcTarget::Stream(stream_id) => {
982 match cold_store.remove_all(&cold_chunk_prefix(stream_id)).await {
983 Ok(()) => cold_store.remove_all(&cold_index_prefix(stream_id)).await,
984 Err(err) => Err(err),
985 }
986 }
987 ColdGcTarget::Paths(paths) => {
988 let mut outcome = Ok(());
989 for path in paths {
990 if let Err(err) = cold_store.delete_chunk(path).await {
991 outcome = Err(err);
992 break;
993 }
994 }
995 outcome
996 }
997 };
998 match result {
999 Ok(()) => {
1000 acked_seq = Some(entry.seq);
1001 reclaimed += 1;
1002 }
1003 Err(err) => {
1004 self.metrics.record_cold_gc_error();
1005 if acked_seq.is_none() {
1006 return Err(RuntimeError::ColdStoreIo {
1007 message: err.to_string(),
1008 });
1009 }
1010 break;
1011 }
1012 }
1013 }
1014 if let Some(up_to_seq) = acked_seq {
1015 self.ack_cold_gc(raft_group_id, up_to_seq).await?;
1016 self.metrics
1017 .record_cold_gc_reclaimed(u64::try_from(reclaimed).expect("reclaimed fits u64"));
1018 }
1019 Ok(reclaimed)
1020 }
1021
1022 pub async fn run_cold_gc_all_groups_once(
1023 &self,
1024 max_entries_per_group: usize,
1025 ) -> Result<usize, RuntimeError> {
1026 if self.cold_store.is_none() {
1027 return Ok(0);
1028 }
1029 let mut reclaimed = 0;
1030 for group_id in 0..self.shard_map.raft_group_count() {
1031 reclaimed += self
1032 .run_cold_gc_group_once(RaftGroupId(group_id), max_entries_per_group)
1033 .await?;
1034 }
1035 Ok(reclaimed)
1036 }
1037
1038 pub async fn append(&self, request: AppendRequest) -> Result<AppendResponse, RuntimeError> {
1039 if request.payload.is_empty() {
1040 return Err(RuntimeError::EmptyAppend);
1041 }
1042 let placement = self.shard_map.locate(&request.stream_id);
1043 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1044 let (response_tx, response_rx) = oneshot::channel();
1045 self.send_core_command(
1046 mailbox,
1047 CoreCommand::Append {
1048 request,
1049 placement,
1050 response_tx,
1051 },
1052 response_rx,
1053 )
1054 .await
1055 }
1056
1057 pub async fn append_batch(
1058 &self,
1059 request: AppendBatchRequest,
1060 ) -> Result<AppendBatchResponse, RuntimeError> {
1061 if request.payloads.is_empty() {
1062 return Err(RuntimeError::EmptyAppend);
1063 }
1064 let placement = self.shard_map.locate(&request.stream_id);
1065 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1066 let (response_tx, response_rx) = oneshot::channel();
1067 self.send_core_command(
1068 mailbox,
1069 CoreCommand::AppendBatch {
1070 request,
1071 placement,
1072 response_tx,
1073 },
1074 response_rx,
1075 )
1076 .await
1077 }
1078
1079 pub async fn snapshot_group(
1080 &self,
1081 raft_group_id: RaftGroupId,
1082 ) -> Result<GroupSnapshot, RuntimeError> {
1083 let placement = self.placement_for_group(raft_group_id)?;
1084 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1085 let (response_tx, response_rx) = oneshot::channel();
1086 self.send_core_command(
1087 mailbox,
1088 CoreCommand::SnapshotGroup {
1089 placement,
1090 response_tx,
1091 },
1092 response_rx,
1093 )
1094 .await
1095 }
1096
1097 pub async fn install_group_snapshot(
1098 &self,
1099 snapshot: GroupSnapshot,
1100 ) -> Result<(), RuntimeError> {
1101 let expected = self.placement_for_group(snapshot.placement.raft_group_id)?;
1102 if snapshot.placement != expected {
1103 return Err(RuntimeError::SnapshotPlacementMismatch {
1104 expected,
1105 actual: snapshot.placement,
1106 });
1107 }
1108 let mailbox = &self.mailboxes[usize::from(expected.core_id.0)];
1109 let (response_tx, response_rx) = oneshot::channel();
1110 self.send_core_command(
1111 mailbox,
1112 CoreCommand::InstallGroupSnapshot {
1113 snapshot,
1114 response_tx,
1115 },
1116 response_rx,
1117 )
1118 .await
1119 }
1120
1121 #[cfg(madsim)]
1122 pub async fn shutdown_group_engine_for_simulation(
1123 &self,
1124 placement: ShardPlacement,
1125 ) -> Result<(), RuntimeError> {
1126 let expected = self.placement_for_group(placement.raft_group_id)?;
1127 if placement != expected {
1128 return Err(RuntimeError::SnapshotPlacementMismatch {
1129 expected,
1130 actual: placement,
1131 });
1132 }
1133 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1134 let (response_tx, response_rx) = oneshot::channel();
1135 self.send_core_command(
1136 mailbox,
1137 CoreCommand::ShutdownGroupEngine {
1138 placement,
1139 response_tx,
1140 },
1141 response_rx,
1142 )
1143 .await
1144 }
1145
1146 #[cfg(madsim)]
1147 pub async fn install_group_engine_for_simulation(
1148 &self,
1149 placement: ShardPlacement,
1150 engine: Box<dyn crate::engine::GroupEngine>,
1151 ) -> Result<(), RuntimeError> {
1152 let expected = self.placement_for_group(placement.raft_group_id)?;
1153 if placement != expected {
1154 return Err(RuntimeError::SnapshotPlacementMismatch {
1155 expected,
1156 actual: placement,
1157 });
1158 }
1159 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1160 let (response_tx, response_rx) = oneshot::channel();
1161 self.send_core_command(
1162 mailbox,
1163 CoreCommand::InstallGroupEngine {
1164 placement,
1165 engine,
1166 response_tx,
1167 },
1168 response_rx,
1169 )
1170 .await
1171 }
1172
1173 pub async fn warm_group(
1174 &self,
1175 raft_group_id: RaftGroupId,
1176 ) -> Result<ShardPlacement, RuntimeError> {
1177 let placement = self.placement_for_group(raft_group_id)?;
1178 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1179 let (response_tx, response_rx) = oneshot::channel();
1180 self.send_core_command(
1181 mailbox,
1182 CoreCommand::WarmGroup {
1183 placement,
1184 response_tx,
1185 },
1186 response_rx,
1187 )
1188 .await
1189 }
1190
1191 pub async fn warm_all_groups(&self) -> Result<(), RuntimeError> {
1192 for raw_group_id in 0..self.shard_map.raft_group_count() {
1193 self.warm_group(RaftGroupId(raw_group_id)).await?;
1194 }
1195 Ok(())
1196 }
1197
1198 fn placement_for_group(
1199 &self,
1200 raft_group_id: RaftGroupId,
1201 ) -> Result<ShardPlacement, RuntimeError> {
1202 if raft_group_id.0 >= self.shard_map.raft_group_count() {
1203 return Err(RuntimeError::InvalidRaftGroup {
1204 raft_group_id,
1205 raft_group_count: self.shard_map.raft_group_count(),
1206 });
1207 }
1208 Ok(ShardPlacement {
1209 core_id: CoreId(
1210 (raft_group_id.0 % u32::from(self.shard_map.core_count()))
1211 .try_into()
1212 .expect("core id fits u16"),
1213 ),
1214 shard_id: ShardId(raft_group_id.0),
1215 raft_group_id,
1216 })
1217 }
1218
1219 async fn send_core_command<T>(
1220 &self,
1221 mailbox: &CoreMailbox,
1222 command: CoreCommand,
1223 response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
1224 ) -> Result<T, RuntimeError> {
1225 self.enqueue_core_command(mailbox, command).await?;
1226 response_rx
1227 .await
1228 .map_err(|_| RuntimeError::ResponseDropped {
1229 core_id: mailbox.core_id,
1230 })?
1231 }
1232
1233 async fn enqueue_core_command(
1234 &self,
1235 mailbox: &CoreMailbox,
1236 command: CoreCommand,
1237 ) -> Result<(), RuntimeError> {
1238 if mailbox.tx.capacity() == 0 {
1239 self.metrics.record_mailbox_full(mailbox.core_id);
1240 }
1241 let started_at = Instant::now();
1242 mailbox
1243 .tx
1244 .send(Traced::capture(command))
1245 .await
1246 .map_err(|_| RuntimeError::MailboxClosed {
1247 core_id: mailbox.core_id,
1248 })?;
1249 self.metrics
1250 .record_routed_request(mailbox.core_id, elapsed_ns(started_at));
1251 Ok(())
1252 }
1253
1254 pub fn metrics(&self) -> RuntimeMetrics {
1255 RuntimeMetrics {
1256 inner: self.metrics.clone(),
1257 }
1258 }
1259
1260 pub fn mailbox_snapshot(&self) -> RuntimeMailboxSnapshot {
1261 let depths = self
1262 .mailboxes
1263 .iter()
1264 .map(CoreMailbox::depth)
1265 .collect::<Vec<_>>();
1266 let capacities = self
1267 .mailboxes
1268 .iter()
1269 .map(CoreMailbox::capacity)
1270 .collect::<Vec<_>>();
1271 RuntimeMailboxSnapshot { depths, capacities }
1272 }
1273}
1274
1275fn spawn_core_worker(threading: RuntimeThreading, worker: CoreWorker) -> Result<(), RuntimeError> {
1276 match threading {
1277 RuntimeThreading::HostedTokio => {
1278 crate::rt::spawn(worker.run());
1279 Ok(())
1280 }
1281 #[cfg(not(madsim))]
1282 RuntimeThreading::ThreadPerCore => {
1283 let core_id = worker.core_id;
1284 std::thread::Builder::new()
1285 .name(format!("ursula-core-{}", core_id.0))
1286 .spawn(move || {
1287 let runtime = tokio::runtime::Builder::new_current_thread()
1288 .enable_all()
1289 .build()
1290 .expect("build per-core tokio runtime");
1291 runtime.block_on(worker.run());
1292 })
1293 .map(|_| ())
1294 .map_err(|err| RuntimeError::SpawnCoreThread {
1295 core_id,
1296 message: err.to_string(),
1297 })
1298 }
1299 }
1300}