Skip to main content

ursula_runtime/
runtime.rs

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    /// Per-group cap on raft-submitted-but-not-yet-applied payload bytes.
89    /// `None` disables the admission (default). Catches "raft replication slow"
90    /// before in-memory queues grow unbounded.
91    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    /// Build runtime configuration from a typed `ursula_config::RuntimeConfig`.
128    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                // Surfaces "this node can't write to S3" to the snapshot driver's
339                // health/yield logic, which drives leadership-yield off real cold
340                // flush failures rather than a stat probe a keep-alive connection
341                // can mask.
342                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    /// Drains the leader-side cold-GC queue for one group: physically reclaims
475    /// each queued target from cold storage, then replicates an ack that pops
476    /// the reclaimed entries. Deletions are idempotent, so a crash or leader
477    /// change between reclaim and ack simply re-runs them next tick.
478    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        // Entries are FIFO by seq; stop at the first failure so the ack never
493        // skips past an object that is still present in cold storage.
494        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    /// Routes a group command to its owning core and awaits the reply.
676    /// `admission` carries the incoming payload bytes for
677    /// raft-uncommitted-backpressure-guarded writes; the check itself runs on
678    /// the owning core.
679    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
755/// Expands the operation manifest into the uniform `ShardRuntime` client
756/// methods: locate the placement (by stream id or raft group id), open the
757/// reply channel, and submit the `GroupCommand` through [`ShardRuntime::
758/// group_rpc`]. Entries with `client { none }` keep hand-written methods; see
759/// the manifest grammar in [`crate::ops`].
760macro_rules! shard_runtime_operations {
761    // Stream-routed write with admission and an emptiness precheck.
762    (@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    // Stream-routed write with admission.
809    (@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    // Stream-routed operation without admission.
849    (@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    // Group-routed operation: an explicit `RaftGroupId` followed by the
887    // manifest fields in declaration order.
888    (@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    // Hand-written client method: nothing to generate.
930    (@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}