1use std::collections::HashMap;
2use std::sync::Arc;
3use std::sync::atomic::AtomicU64;
4use std::sync::atomic::Ordering;
5
6#[cfg(not(madsim))]
7use tokio::task::JoinSet;
8use ursula_shard::BucketStreamId;
9use ursula_shard::CoreId;
10use ursula_shard::RaftGroupId;
11use ursula_shard::ShardId;
12use ursula_shard::ShardPlacement;
13use ursula_shard::StaticShardMap;
14use ursula_stream::ColdChunkRef;
15use ursula_stream::ColdFlushCandidate;
16use ursula_stream::ColdGcEntry;
17use ursula_stream::ColdGcTarget;
18
19use crate::admission::RaftUncommittedAdmission;
20use crate::admission::RaftUncommittedBytesTracker;
21use crate::cold_index::cold_index_prefix;
22use crate::cold_store::ColdStoreHandle;
23use crate::cold_store::ColdStoreInfo;
24use crate::cold_store::cold_chunk_prefix;
25use crate::cold_store::new_cold_chunk_path;
26use crate::command::GroupSnapshot;
27use crate::core_worker::CoreCommand;
28use crate::core_worker::CoreMailbox;
29use crate::core_worker::CoreWorker;
30use crate::core_worker::WaitReadCancel;
31use crate::engine::GroupEngineFactory;
32use crate::engine::in_memory::InMemoryGroupEngineFactory;
33use crate::error::RuntimeError;
34use crate::group_actor::GroupCommand;
35use crate::metrics::COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS;
36use crate::metrics::RuntimeMailboxSnapshot;
37use crate::metrics::RuntimeMetrics;
38use crate::metrics::RuntimeMetricsInner;
39use crate::metrics::append_batch_payload_bytes;
40use crate::metrics::elapsed_ns;
41use crate::metrics::is_stale_cold_flush_candidate_error;
42use crate::request::AckColdGcResponse;
43use crate::request::AppendBatchRequest;
44use crate::request::AppendBatchResponse;
45use crate::request::AppendExternalRequest;
46use crate::request::AppendRequest;
47use crate::request::AppendResponse;
48use crate::request::BootstrapStreamRequest;
49use crate::request::BootstrapStreamResponse;
50use crate::request::CloseStreamRequest;
51use crate::request::CloseStreamResponse;
52use crate::request::ColdWriteAdmission;
53use crate::request::CreateStreamExternalRequest;
54use crate::request::CreateStreamRequest;
55use crate::request::CreateStreamResponse;
56use crate::request::DeleteSnapshotRequest;
57use crate::request::DeleteStreamRequest;
58use crate::request::DeleteStreamResponse;
59use crate::request::FlushColdRequest;
60use crate::request::FlushColdResponse;
61use crate::request::GetStreamAttrsRequest;
62use crate::request::GetStreamAttrsResponse;
63use crate::request::HeadStreamRequest;
64use crate::request::HeadStreamResponse;
65use crate::request::PlanColdFlushRequest;
66use crate::request::PlanGroupColdFlushRequest;
67use crate::request::PublishSnapshotRequest;
68use crate::request::PublishSnapshotResponse;
69use crate::request::ReadSnapshotRequest;
70use crate::request::ReadSnapshotResponse;
71use crate::request::ReadStreamRequest;
72use crate::request::ReadStreamResponse;
73use crate::request::UpdateStreamAttrsRequest;
74use crate::request::UpdateStreamAttrsResponse;
75use crate::rt::sync::Semaphore;
76use crate::rt::sync::mpsc;
77use crate::rt::sync::oneshot;
78use crate::rt::time::Instant;
79use crate::trace::Traced;
80
81#[derive(Debug, Clone)]
82pub struct RuntimeConfig {
83 pub core_count: usize,
84 pub raft_group_count: usize,
85 pub mailbox_capacity: usize,
86 pub threading: RuntimeThreading,
87 pub cold_max_hot_bytes_per_group: Option<u64>,
88 pub raft_max_uncommitted_bytes_per_group: Option<u64>,
92 pub live_read_max_waiters_per_core: Option<u64>,
93}
94
95impl RuntimeConfig {
96 pub fn new(core_count: usize, raft_group_count: usize) -> Self {
97 #[cfg(not(madsim))]
98 let threading = RuntimeThreading::ThreadPerCore;
99 #[cfg(madsim)]
100 let threading = RuntimeThreading::HostedTokio;
101 Self {
102 core_count,
103 raft_group_count,
104 mailbox_capacity: 1024,
105 threading,
106 cold_max_hot_bytes_per_group: None,
107 raft_max_uncommitted_bytes_per_group: None,
108 live_read_max_waiters_per_core: Some(65_536),
109 }
110 }
111
112 pub fn with_cold_max_hot_bytes_per_group(mut self, value: Option<u64>) -> Self {
113 self.cold_max_hot_bytes_per_group = value;
114 self
115 }
116
117 pub fn with_raft_max_uncommitted_bytes_per_group(mut self, value: Option<u64>) -> Self {
118 self.raft_max_uncommitted_bytes_per_group = value;
119 self
120 }
121
122 pub fn with_live_read_max_waiters_per_core(mut self, value: Option<u64>) -> Self {
123 self.live_read_max_waiters_per_core = value;
124 self
125 }
126
127 pub fn from_ursula_config(cfg: &ursula_config::RuntimeConfig, raft_group_count: usize) -> Self {
129 let mut config = Self::new(cfg.core_count, raft_group_count);
130 config.live_read_max_waiters_per_core = cfg
131 .live_read_max_waiters_per_core
132 .and_then(|n| if n == 0 { None } else { Some(n as u64) });
133 config
134 }
135}
136
137#[derive(Debug, Clone, Copy, PartialEq, Eq)]
138pub enum RuntimeThreading {
139 #[cfg(not(madsim))]
140 ThreadPerCore,
141 HostedTokio,
142}
143
144#[derive(Debug, Clone)]
145pub struct ShardRuntime {
146 shard_map: StaticShardMap,
147 mailboxes: Vec<CoreMailbox>,
148 metrics: Arc<RuntimeMetricsInner>,
149 next_waiter_id: Arc<AtomicU64>,
150 cold_store: Option<ColdStoreHandle>,
151}
152
153impl ShardRuntime {
154 pub fn spawn(config: RuntimeConfig) -> Result<Self, RuntimeError> {
155 Self::spawn_with_engine_factory(config, InMemoryGroupEngineFactory::default())
156 }
157
158 pub fn spawn_with_engine_factory(
159 config: RuntimeConfig,
160 engine_factory: impl GroupEngineFactory,
161 ) -> Result<Self, RuntimeError> {
162 Self::spawn_with_engine_factory_and_cold_store(config, engine_factory, None)
163 }
164
165 pub fn spawn_with_engine_factory_and_cold_store(
166 config: RuntimeConfig,
167 engine_factory: impl GroupEngineFactory,
168 cold_store: Option<ColdStoreHandle>,
169 ) -> Result<Self, RuntimeError> {
170 let shard_map = StaticShardMap::new(config.core_count, config.raft_group_count)?;
171 let metrics = Arc::new(RuntimeMetricsInner::new(
172 usize::from(shard_map.core_count()),
173 usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
174 ));
175 let cold_write_admission = ColdWriteAdmission {
176 max_hot_bytes_per_group: config.cold_max_hot_bytes_per_group,
177 };
178 let raft_uncommitted_admission = RaftUncommittedAdmission {
179 max_uncommitted_bytes_per_group: config.raft_max_uncommitted_bytes_per_group,
180 };
181 let raft_uncommitted_bytes = Arc::new(RaftUncommittedBytesTracker::new(
182 usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
183 ));
184 let engine_factory: Arc<dyn GroupEngineFactory> = Arc::new(engine_factory);
185 let read_materialization = Arc::new(Semaphore::new(config.mailbox_capacity.max(1)));
186 let mut mailboxes = Vec::with_capacity(usize::from(shard_map.core_count()));
187 for raw_core_id in 0..shard_map.core_count() {
188 let core_id = CoreId(raw_core_id);
189 let (tx, rx) = mpsc::channel(config.mailbox_capacity.max(1));
190 let worker = CoreWorker {
191 core_id,
192 rx,
193 engine_factory: engine_factory.clone(),
194 groups: HashMap::new(),
195 metrics: metrics.clone(),
196 group_mailbox_capacity: config.mailbox_capacity.max(1),
197 cold_write_admission,
198 raft_uncommitted_admission,
199 raft_uncommitted_bytes: raft_uncommitted_bytes.clone(),
200 live_read_max_waiters_per_core: config.live_read_max_waiters_per_core,
201 read_materialization: read_materialization.clone(),
202 };
203 spawn_core_worker(config.threading, worker)?;
204 mailboxes.push(CoreMailbox { core_id, tx });
205 }
206 Ok(Self {
207 shard_map,
208 mailboxes,
209 metrics,
210 next_waiter_id: Arc::new(AtomicU64::new(1)),
211 cold_store,
212 })
213 }
214
215 pub fn locate(&self, stream_id: &BucketStreamId) -> ShardPlacement {
216 self.shard_map.locate(stream_id)
217 }
218
219 pub fn has_cold_store(&self) -> bool {
220 self.cold_store.is_some()
221 }
222
223 pub fn cold_store(&self) -> Option<ColdStoreHandle> {
224 self.cold_store.clone()
225 }
226
227 pub fn cold_store_info(&self) -> Option<ColdStoreInfo> {
228 self.cold_store
229 .as_ref()
230 .map(|cold_store| cold_store.info().clone())
231 }
232
233 pub async fn wait_read_stream(
234 &self,
235 request: ReadStreamRequest,
236 ) -> Result<ReadStreamResponse, RuntimeError> {
237 let placement = self.shard_map.locate(&request.stream_id);
238 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
239 let waiter_id = self.next_waiter_id.fetch_add(1, Ordering::Relaxed);
240 let stream_id = request.stream_id.clone();
241 let (response_tx, response_rx) = oneshot::channel();
242 self.enqueue_core_command(mailbox, CoreCommand::Group {
243 placement,
244 admission: None,
245 command: GroupCommand::WaitRead {
246 request,
247 waiter_id,
248 response_tx,
249 },
250 })
251 .await?;
252 let mut cancel = WaitReadCancel::new(mailbox.tx.clone(), stream_id, placement, waiter_id);
253 let response = response_rx
254 .await
255 .map_err(|_| RuntimeError::ResponseDropped {
256 core_id: mailbox.core_id,
257 })?;
258 cancel.disarm();
259 response
260 }
261
262 pub async fn require_local_live_read_owner(
263 &self,
264 stream_id: &BucketStreamId,
265 ) -> Result<(), RuntimeError> {
266 let placement = self.shard_map.locate(stream_id);
267 let (response_tx, response_rx) = oneshot::channel();
268 self.group_rpc(
269 placement,
270 None,
271 GroupCommand::RequireLiveReadOwner { response_tx },
272 response_rx,
273 )
274 .await
275 }
276
277 pub async fn flush_cold_once(
278 &self,
279 request: PlanColdFlushRequest,
280 ) -> Result<Option<FlushColdResponse>, RuntimeError> {
281 let Some(candidate) = self.plan_cold_flush(request).await? else {
282 return Ok(None);
283 };
284 self.flush_cold_candidate(candidate).await.map(Some)
285 }
286
287 pub async fn flush_cold_group_once(
288 &self,
289 raft_group_id: RaftGroupId,
290 request: PlanGroupColdFlushRequest,
291 ) -> Result<Option<FlushColdResponse>, RuntimeError> {
292 let mut candidates = self
293 .plan_next_cold_flush_batch(raft_group_id, request, 1)
294 .await?;
295 let Some(candidate) = candidates.pop() else {
296 return Ok(None);
297 };
298 match self.flush_cold_candidate(candidate).await {
299 Ok(response) => Ok(Some(response)),
300 Err(err) if is_stale_cold_flush_candidate_error(&err) => Ok(None),
301 Err(err) => Err(err),
302 }
303 }
304
305 pub async fn flush_cold_group_batch_once(
306 &self,
307 raft_group_id: RaftGroupId,
308 request: PlanGroupColdFlushRequest,
309 max_candidates: usize,
310 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
311 let candidates = self
312 .plan_next_cold_flush_batch(raft_group_id, request, max_candidates)
313 .await?;
314 if candidates.is_empty() {
315 return Ok(Vec::new());
316 }
317 self.flush_cold_candidates_batch(candidates).await
318 }
319
320 async fn flush_cold_candidate(
321 &self,
322 candidate: ColdFlushCandidate,
323 ) -> Result<FlushColdResponse, RuntimeError> {
324 let Some(cold_store) = self.cold_store.as_ref() else {
325 return Err(RuntimeError::ColdStoreConfig {
326 message: "cold backend must be configured before flushing cold chunks".to_owned(),
327 });
328 };
329 let path = new_cold_chunk_path(
330 &candidate.stream_id,
331 candidate.start_offset,
332 candidate.end_offset,
333 );
334 let upload_started_at = Instant::now();
335 let object_size = match cold_store.write_chunk(&path, &candidate.payload).await {
336 Ok(object_size) => object_size,
337 Err(err) => {
338 self.metrics.record_cold_flush_write_error();
343 return Err(RuntimeError::ColdStoreIo {
344 message: err.to_string(),
345 });
346 }
347 };
348 self.metrics
349 .record_cold_upload(object_size, elapsed_ns(upload_started_at));
350 let chunk = ColdChunkRef {
351 start_offset: candidate.start_offset,
352 end_offset: candidate.end_offset,
353 s3_path: path.clone(),
354 object_size,
355 };
356 let publish_started_at = Instant::now();
357 let publish = self
358 .flush_cold(FlushColdRequest {
359 stream_id: candidate.stream_id,
360 chunk,
361 })
362 .await;
363 match publish {
364 Ok(response) => {
365 self.metrics
366 .record_cold_publish(object_size, elapsed_ns(publish_started_at));
367 Ok(response)
368 }
369 Err(err) => Err(err),
370 }
371 }
372
373 pub(crate) async fn flush_cold_candidates_batch(
374 &self,
375 candidates: Vec<ColdFlushCandidate>,
376 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
377 let mut responses = Vec::with_capacity(candidates.len());
378 for candidate in candidates {
379 match self.flush_cold_candidate(candidate).await {
380 Ok(response) => responses.push(response),
381 Err(err) if is_stale_cold_flush_candidate_error(&err) => {}
382 Err(err) => return Err(err),
383 }
384 }
385 Ok(responses)
386 }
387
388 #[cfg(madsim)]
389 pub async fn flush_cold_candidates_batch_for_simulation(
390 &self,
391 candidates: Vec<ColdFlushCandidate>,
392 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
393 self.flush_cold_candidates_batch(candidates).await
394 }
395
396 pub async fn flush_cold_all_groups_once(
397 &self,
398 request: PlanGroupColdFlushRequest,
399 ) -> Result<usize, RuntimeError> {
400 self.flush_cold_all_groups_once_bounded(request, 1).await
401 }
402
403 pub async fn flush_cold_all_groups_once_bounded(
404 &self,
405 request: PlanGroupColdFlushRequest,
406 max_concurrency: usize,
407 ) -> Result<usize, RuntimeError> {
408 let max_concurrency = max_concurrency.max(1);
409 if max_concurrency == 1 {
410 return self.flush_cold_all_groups_once_serial(request).await;
411 }
412 #[cfg(madsim)]
413 {
414 return self.flush_cold_all_groups_once_serial(request).await;
415 }
416 #[cfg(not(madsim))]
417 {
418 let mut flushed = 0;
419 let mut next_group_id = 0;
420 let group_count = self.shard_map.raft_group_count();
421 let mut tasks = JoinSet::new();
422
423 while next_group_id < group_count || !tasks.is_empty() {
424 while next_group_id < group_count && tasks.len() < max_concurrency {
425 let runtime = self.clone();
426 let request = request.clone();
427 let group_id = RaftGroupId(next_group_id);
428 next_group_id += 1;
429 tasks.spawn(async move {
430 runtime
431 .flush_cold_group_batch_once(
432 group_id,
433 request,
434 COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
435 )
436 .await
437 .map(|responses| responses.len())
438 });
439 }
440 if let Some(result) = tasks.join_next().await {
441 match result {
442 Ok(Ok(count)) => flushed += count,
443 Ok(Err(err)) => return Err(err),
444 Err(err) => {
445 return Err(RuntimeError::ColdStoreIo {
446 message: format!("cold flush task failed: {err}"),
447 });
448 }
449 }
450 }
451 }
452 Ok(flushed)
453 }
454 }
455
456 async fn flush_cold_all_groups_once_serial(
457 &self,
458 request: PlanGroupColdFlushRequest,
459 ) -> Result<usize, RuntimeError> {
460 let mut flushed = 0;
461 for group_id in 0..self.shard_map.raft_group_count() {
462 flushed += self
463 .flush_cold_group_batch_once(
464 RaftGroupId(group_id),
465 request.clone(),
466 COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
467 )
468 .await?
469 .len();
470 }
471 Ok(flushed)
472 }
473
474 pub async fn run_cold_gc_group_once(
479 &self,
480 raft_group_id: RaftGroupId,
481 max_entries: usize,
482 ) -> Result<usize, RuntimeError> {
483 let Some(cold_store) = self.cold_store.as_ref() else {
484 return Ok(0);
485 };
486 let entries = self.plan_cold_gc(raft_group_id, max_entries).await?;
487 if entries.is_empty() {
488 return Ok(0);
489 }
490 let mut acked_seq = None;
491 let mut reclaimed = 0usize;
492 for entry in entries {
495 let result = match &entry.target {
496 ColdGcTarget::Stream(stream_id) => {
497 match cold_store.remove_all(&cold_chunk_prefix(stream_id)).await {
498 Ok(()) => cold_store.remove_all(&cold_index_prefix(stream_id)).await,
499 Err(err) => Err(err),
500 }
501 }
502 ColdGcTarget::Paths(paths) => {
503 let mut outcome = Ok(());
504 for path in paths {
505 if let Err(err) = cold_store.delete_chunk(path).await {
506 outcome = Err(err);
507 break;
508 }
509 }
510 outcome
511 }
512 };
513 match result {
514 Ok(()) => {
515 acked_seq = Some(entry.seq);
516 reclaimed += 1;
517 }
518 Err(err) => {
519 self.metrics.record_cold_gc_error();
520 if acked_seq.is_none() {
521 return Err(RuntimeError::ColdStoreIo {
522 message: err.to_string(),
523 });
524 }
525 break;
526 }
527 }
528 }
529 if let Some(up_to_seq) = acked_seq {
530 self.ack_cold_gc(raft_group_id, up_to_seq).await?;
531 self.metrics
532 .record_cold_gc_reclaimed(u64::try_from(reclaimed).expect("reclaimed fits u64"));
533 }
534 Ok(reclaimed)
535 }
536
537 pub async fn run_cold_gc_all_groups_once(
538 &self,
539 max_entries_per_group: usize,
540 ) -> Result<usize, RuntimeError> {
541 if self.cold_store.is_none() {
542 return Ok(0);
543 }
544 let mut reclaimed = 0;
545 for group_id in 0..self.shard_map.raft_group_count() {
546 reclaimed += self
547 .run_cold_gc_group_once(RaftGroupId(group_id), max_entries_per_group)
548 .await?;
549 }
550 Ok(reclaimed)
551 }
552
553 pub async fn install_group_snapshot(
554 &self,
555 snapshot: GroupSnapshot,
556 ) -> Result<(), RuntimeError> {
557 let expected = self.placement_for_group(snapshot.placement.raft_group_id)?;
558 if snapshot.placement != expected {
559 return Err(RuntimeError::SnapshotPlacementMismatch {
560 expected,
561 actual: snapshot.placement,
562 });
563 }
564 let (response_tx, response_rx) = oneshot::channel();
565 self.group_rpc(
566 expected,
567 None,
568 GroupCommand::InstallGroupSnapshot {
569 snapshot,
570 response_tx,
571 },
572 response_rx,
573 )
574 .await
575 }
576
577 #[cfg(madsim)]
578 pub async fn shutdown_group_engine_for_simulation(
579 &self,
580 placement: ShardPlacement,
581 ) -> Result<(), RuntimeError> {
582 let expected = self.placement_for_group(placement.raft_group_id)?;
583 if placement != expected {
584 return Err(RuntimeError::SnapshotPlacementMismatch {
585 expected,
586 actual: placement,
587 });
588 }
589 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
590 let (response_tx, response_rx) = oneshot::channel();
591 self.send_core_command(
592 mailbox,
593 CoreCommand::ShutdownGroupEngine {
594 placement,
595 response_tx,
596 },
597 response_rx,
598 )
599 .await
600 }
601
602 #[cfg(madsim)]
603 pub async fn install_group_engine_for_simulation(
604 &self,
605 placement: ShardPlacement,
606 engine: Box<dyn crate::engine::GroupEngine>,
607 ) -> Result<(), RuntimeError> {
608 let expected = self.placement_for_group(placement.raft_group_id)?;
609 if placement != expected {
610 return Err(RuntimeError::SnapshotPlacementMismatch {
611 expected,
612 actual: placement,
613 });
614 }
615 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
616 let (response_tx, response_rx) = oneshot::channel();
617 self.send_core_command(
618 mailbox,
619 CoreCommand::InstallGroupEngine {
620 placement,
621 engine,
622 response_tx,
623 },
624 response_rx,
625 )
626 .await
627 }
628
629 pub async fn warm_group(
630 &self,
631 raft_group_id: RaftGroupId,
632 ) -> Result<ShardPlacement, RuntimeError> {
633 let placement = self.placement_for_group(raft_group_id)?;
634 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
635 let (response_tx, response_rx) = oneshot::channel();
636 self.send_core_command(
637 mailbox,
638 CoreCommand::WarmGroup {
639 placement,
640 response_tx,
641 },
642 response_rx,
643 )
644 .await
645 }
646
647 pub async fn warm_all_groups(&self) -> Result<(), RuntimeError> {
648 for raw_group_id in 0..self.shard_map.raft_group_count() {
649 self.warm_group(RaftGroupId(raw_group_id)).await?;
650 }
651 Ok(())
652 }
653
654 fn placement_for_group(
655 &self,
656 raft_group_id: RaftGroupId,
657 ) -> Result<ShardPlacement, RuntimeError> {
658 if raft_group_id.0 >= self.shard_map.raft_group_count() {
659 return Err(RuntimeError::InvalidRaftGroup {
660 raft_group_id,
661 raft_group_count: self.shard_map.raft_group_count(),
662 });
663 }
664 Ok(ShardPlacement {
665 core_id: CoreId(
666 (raft_group_id.0 % u32::from(self.shard_map.core_count()))
667 .try_into()
668 .expect("core id fits u16"),
669 ),
670 shard_id: ShardId(raft_group_id.0),
671 raft_group_id,
672 })
673 }
674
675 async fn group_rpc<T>(
680 &self,
681 placement: ShardPlacement,
682 admission: Option<u64>,
683 command: GroupCommand,
684 response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
685 ) -> Result<T, RuntimeError> {
686 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
687 self.send_core_command(
688 mailbox,
689 CoreCommand::Group {
690 placement,
691 admission,
692 command,
693 },
694 response_rx,
695 )
696 .await
697 }
698
699 async fn send_core_command<T>(
700 &self,
701 mailbox: &CoreMailbox,
702 command: CoreCommand,
703 response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
704 ) -> Result<T, RuntimeError> {
705 self.enqueue_core_command(mailbox, command).await?;
706 response_rx
707 .await
708 .map_err(|_| RuntimeError::ResponseDropped {
709 core_id: mailbox.core_id,
710 })?
711 }
712
713 async fn enqueue_core_command(
714 &self,
715 mailbox: &CoreMailbox,
716 command: CoreCommand,
717 ) -> Result<(), RuntimeError> {
718 if mailbox.tx.capacity() == 0 {
719 self.metrics.record_mailbox_full(mailbox.core_id);
720 }
721 let started_at = Instant::now();
722 mailbox
723 .tx
724 .send(Traced::capture(command))
725 .await
726 .map_err(|_| RuntimeError::MailboxClosed {
727 core_id: mailbox.core_id,
728 })?;
729 self.metrics
730 .record_routed_request(mailbox.core_id, elapsed_ns(started_at));
731 Ok(())
732 }
733
734 pub fn metrics(&self) -> RuntimeMetrics {
735 RuntimeMetrics {
736 inner: self.metrics.clone(),
737 }
738 }
739
740 pub fn mailbox_snapshot(&self) -> RuntimeMailboxSnapshot {
741 let depths = self
742 .mailboxes
743 .iter()
744 .map(CoreMailbox::depth)
745 .collect::<Vec<_>>();
746 let capacities = self
747 .mailboxes
748 .iter()
749 .map(CoreMailbox::capacity)
750 .collect::<Vec<_>>();
751 RuntimeMailboxSnapshot { depths, capacities }
752 }
753}
754
755macro_rules! shard_runtime_operations {
761 (@munch
763 methods { $($methods:tt)* }
764 rest {
765 $(#[$attr:meta])*
766 op $Variant:ident {
767 fields { $req:ident: $Req:ty $(,)? }
768 reply { $tx:ident: $Resp:ty }
769 guard { $g:ident }
770 handle { $($handle:tt)* }
771 client {
772 $vis:vis stream fn $method:ident,
773 non_empty: $ne:ident,
774 admit: $incoming:expr
775 }
776 }
777 $($rest:tt)*
778 }
779 ) => {
780 shard_runtime_operations! {
781 @munch
782 methods {
783 $($methods)*
784 $(#[$attr])*
785 $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
786 if $req.$ne.is_empty() {
787 return Err(RuntimeError::EmptyAppend);
788 }
789 let placement = self.shard_map.locate(&$req.stream_id);
790 let incoming_bytes = $incoming;
791 let (response_tx, response_rx) = oneshot::channel();
792 self.group_rpc(
793 placement,
794 Some(incoming_bytes),
795 GroupCommand::$Variant {
796 $req,
797 $tx: response_tx,
798 $g: None,
799 },
800 response_rx,
801 )
802 .await
803 }
804 }
805 rest { $($rest)* }
806 }
807 };
808 (@munch
810 methods { $($methods:tt)* }
811 rest {
812 $(#[$attr:meta])*
813 op $Variant:ident {
814 fields { $req:ident: $Req:ty $(,)? }
815 reply { $tx:ident: $Resp:ty }
816 guard { $g:ident }
817 handle { $($handle:tt)* }
818 client { $vis:vis stream fn $method:ident, admit: $incoming:expr }
819 }
820 $($rest:tt)*
821 }
822 ) => {
823 shard_runtime_operations! {
824 @munch
825 methods {
826 $($methods)*
827 $(#[$attr])*
828 $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
829 let placement = self.shard_map.locate(&$req.stream_id);
830 let incoming_bytes = $incoming;
831 let (response_tx, response_rx) = oneshot::channel();
832 self.group_rpc(
833 placement,
834 Some(incoming_bytes),
835 GroupCommand::$Variant {
836 $req,
837 $tx: response_tx,
838 $g: None,
839 },
840 response_rx,
841 )
842 .await
843 }
844 }
845 rest { $($rest)* }
846 }
847 };
848 (@munch
850 methods { $($methods:tt)* }
851 rest {
852 $(#[$attr:meta])*
853 op $Variant:ident {
854 fields { $req:ident: $Req:ty $(,)? }
855 reply { $tx:ident: $Resp:ty }
856 guard { none }
857 handle { $($handle:tt)* }
858 client { $vis:vis stream fn $method:ident }
859 }
860 $($rest:tt)*
861 }
862 ) => {
863 shard_runtime_operations! {
864 @munch
865 methods {
866 $($methods)*
867 $(#[$attr])*
868 $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
869 let placement = self.shard_map.locate(&$req.stream_id);
870 let (response_tx, response_rx) = oneshot::channel();
871 self.group_rpc(
872 placement,
873 None,
874 GroupCommand::$Variant {
875 $req,
876 $tx: response_tx,
877 },
878 response_rx,
879 )
880 .await
881 }
882 }
883 rest { $($rest)* }
884 }
885 };
886 (@munch
889 methods { $($methods:tt)* }
890 rest {
891 $(#[$attr:meta])*
892 op $Variant:ident {
893 fields { $($field:ident: $field_ty:ty),* $(,)? }
894 reply { $tx:ident: $Resp:ty }
895 guard { none }
896 handle { $($handle:tt)* }
897 client { $vis:vis group fn $method:ident }
898 }
899 $($rest:tt)*
900 }
901 ) => {
902 shard_runtime_operations! {
903 @munch
904 methods {
905 $($methods)*
906 $(#[$attr])*
907 $vis async fn $method(
908 &self,
909 raft_group_id: RaftGroupId,
910 $($field: $field_ty),*
911 ) -> Result<$Resp, RuntimeError> {
912 let placement = self.placement_for_group(raft_group_id)?;
913 let (response_tx, response_rx) = oneshot::channel();
914 self.group_rpc(
915 placement,
916 None,
917 GroupCommand::$Variant {
918 $($field,)*
919 $tx: response_tx,
920 },
921 response_rx,
922 )
923 .await
924 }
925 }
926 rest { $($rest)* }
927 }
928 };
929 (@munch
931 methods { $($methods:tt)* }
932 rest {
933 $(#[$attr:meta])*
934 op $Variant:ident {
935 fields { $($fields:tt)* }
936 reply { $($reply:tt)* }
937 guard { $($guard:tt)* }
938 handle { $($handle:tt)* }
939 client { none }
940 }
941 $($rest:tt)*
942 }
943 ) => {
944 shard_runtime_operations! {
945 @munch
946 methods { $($methods)* }
947 rest { $($rest)* }
948 }
949 };
950 (@munch
951 methods { $($methods:tt)* }
952 rest {}
953 ) => {
954 impl ShardRuntime {
955 $($methods)*
956 }
957 };
958 ($($manifest:tt)*) => {
959 shard_runtime_operations! {
960 @munch
961 methods {}
962 rest { $($manifest)* }
963 }
964 };
965}
966
967crate::ops::runtime_operations!(shard_runtime_operations);
968
969fn spawn_core_worker(threading: RuntimeThreading, worker: CoreWorker) -> Result<(), RuntimeError> {
970 match threading {
971 RuntimeThreading::HostedTokio => {
972 crate::rt::spawn(worker.run());
973 Ok(())
974 }
975 #[cfg(not(madsim))]
976 RuntimeThreading::ThreadPerCore => {
977 let core_id = worker.core_id;
978 std::thread::Builder::new()
979 .name(format!("ursula-core-{}", core_id.0))
980 .spawn(move || {
981 let runtime = tokio::runtime::Builder::new_current_thread()
982 .enable_all()
983 .build()
984 .expect("build per-core tokio runtime");
985 runtime.block_on(worker.run());
986 })
987 .map(|_| ())
988 .map_err(|err| RuntimeError::SpawnCoreThread {
989 core_id,
990 message: err.to_string(),
991 })
992 }
993 }
994}