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