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#[cfg(not(madsim))]
6use std::time::SystemTime;
7#[cfg(not(madsim))]
8use std::time::UNIX_EPOCH;
9
10#[cfg(not(madsim))]
11use tokio::task::JoinSet;
12use ursula_shard::BucketStreamId;
13use ursula_shard::CoreId;
14use ursula_shard::RaftGroupId;
15use ursula_shard::ShardId;
16use ursula_shard::ShardPlacement;
17use ursula_shard::StaticShardMap;
18use ursula_stream::ColdChunkRef;
19use ursula_stream::ColdFlushCandidate;
20use ursula_stream::ColdGcEntry;
21use ursula_stream::ColdGcTarget;
22
23use crate::admission::RaftUncommittedAdmission;
24use crate::admission::RaftUncommittedBytesTracker;
25use crate::cold_index::ColdStoreColdIndexPageStore;
26use crate::cold_index::cold_index_prefix;
27use crate::cold_index::load_cold_chunks_from_pages;
28use crate::cold_index::select_cold_chunk_compaction;
29use crate::cold_store::ColdStoreHandle;
30use crate::cold_store::ColdStoreInfo;
31use crate::cold_store::cold_chunk_prefix;
32use crate::cold_store::new_cold_chunk_path;
33use crate::cold_store::new_cold_pack_path;
34use crate::command::GroupSnapshot;
35use crate::core_worker::CoreCommand;
36use crate::core_worker::CoreMailbox;
37use crate::core_worker::CoreWorker;
38use crate::core_worker::WaitReadCancel;
39use crate::engine::GroupEngineFactory;
40use crate::engine::in_memory::InMemoryGroupEngineFactory;
41use crate::error::RuntimeError;
42use crate::group_actor::GroupCommand;
43use crate::metrics::COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS;
44use crate::metrics::RuntimeMailboxSnapshot;
45use crate::metrics::RuntimeMetrics;
46use crate::metrics::RuntimeMetricsInner;
47use crate::metrics::append_batch_payload_bytes;
48use crate::metrics::elapsed_ns;
49use crate::metrics::is_stale_cold_flush_candidate_error;
50use crate::request::AckColdGcResponse;
51use crate::request::AdvanceRetentionRequest;
52use crate::request::AdvanceRetentionResponse;
53use crate::request::AppendBatchRequest;
54use crate::request::AppendBatchResponse;
55use crate::request::AppendExternalRequest;
56use crate::request::AppendRequest;
57use crate::request::AppendResponse;
58use crate::request::AppendTransactionRequest;
59use crate::request::AppendTransactionResponse;
60use crate::request::BootstrapStreamRequest;
61use crate::request::BootstrapStreamResponse;
62use crate::request::CloseStreamRequest;
63use crate::request::CloseStreamResponse;
64use crate::request::ColdWriteAdmission;
65use crate::request::CompactColdRequest;
66use crate::request::CompactColdResponse;
67use crate::request::CreateStreamExternalRequest;
68use crate::request::CreateStreamRequest;
69use crate::request::CreateStreamResponse;
70use crate::request::DeleteSnapshotRequest;
71use crate::request::DeleteStreamRequest;
72use crate::request::DeleteStreamResponse;
73use crate::request::FlushColdRequest;
74use crate::request::FlushColdResponse;
75use crate::request::GetStreamAttrsRequest;
76use crate::request::GetStreamAttrsResponse;
77use crate::request::HeadStreamRequest;
78use crate::request::HeadStreamResponse;
79use crate::request::ImportGroupStateRequest;
80use crate::request::ImportGroupStateResponse;
81use crate::request::PlanColdFlushRequest;
82use crate::request::PlanGroupColdFlushRequest;
83use crate::request::PublishSnapshotRequest;
84use crate::request::PublishSnapshotResponse;
85use crate::request::PurgeBucketResponse;
86use crate::request::ReadSnapshotRequest;
87use crate::request::ReadSnapshotResponse;
88use crate::request::ReadStreamRequest;
89use crate::request::ReadStreamResponse;
90use crate::request::SetBucketQuotaRequest;
91use crate::request::SetBucketQuotaResponse;
92use crate::request::UpdateStreamAttrsRequest;
93use crate::request::UpdateStreamAttrsResponse;
94use crate::rt::sync::Semaphore;
95use crate::rt::sync::mpsc;
96use crate::rt::sync::oneshot;
97use crate::rt::time::Instant;
98use crate::trace::Traced;
99
100#[derive(Debug, Clone)]
101pub struct RuntimeConfig {
102    pub core_count: usize,
103    pub raft_group_count: usize,
104    pub mailbox_capacity: usize,
105    pub threading: RuntimeThreading,
106    pub cold_max_hot_bytes_per_group: Option<u64>,
107    /// Per-group cap on raft-submitted-but-not-yet-applied payload bytes.
108    /// `None` disables the admission (default). Catches "raft replication slow"
109    /// before in-memory queues grow unbounded.
110    pub raft_max_uncommitted_bytes_per_group: Option<u64>,
111    pub live_read_max_waiters_per_core: Option<u64>,
112}
113
114impl RuntimeConfig {
115    pub fn new(core_count: usize, raft_group_count: usize) -> Self {
116        #[cfg(not(madsim))]
117        let threading = RuntimeThreading::ThreadPerCore;
118        #[cfg(madsim)]
119        let threading = RuntimeThreading::HostedTokio;
120        Self {
121            core_count,
122            raft_group_count,
123            mailbox_capacity: 1024,
124            threading,
125            cold_max_hot_bytes_per_group: None,
126            raft_max_uncommitted_bytes_per_group: None,
127            live_read_max_waiters_per_core: Some(65_536),
128        }
129    }
130
131    pub fn with_cold_max_hot_bytes_per_group(mut self, value: Option<u64>) -> Self {
132        self.cold_max_hot_bytes_per_group = value;
133        self
134    }
135
136    pub fn with_raft_max_uncommitted_bytes_per_group(mut self, value: Option<u64>) -> Self {
137        self.raft_max_uncommitted_bytes_per_group = value;
138        self
139    }
140
141    pub fn with_live_read_max_waiters_per_core(mut self, value: Option<u64>) -> Self {
142        self.live_read_max_waiters_per_core = value;
143        self
144    }
145
146    /// Build runtime configuration from a typed `ursula_config::RuntimeConfig`.
147    pub fn from_ursula_config(cfg: &ursula_config::RuntimeConfig, raft_group_count: usize) -> Self {
148        let mut config = Self::new(cfg.core_count, raft_group_count);
149        config.live_read_max_waiters_per_core = cfg
150            .live_read_max_waiters_per_core
151            .and_then(|n| if n == 0 { None } else { Some(n as u64) });
152        config
153    }
154}
155
156#[derive(Debug, Clone, Copy, PartialEq, Eq)]
157pub enum RuntimeThreading {
158    #[cfg(not(madsim))]
159    ThreadPerCore,
160    HostedTokio,
161}
162
163#[derive(Debug, Clone)]
164pub struct ShardRuntime {
165    shard_map: StaticShardMap,
166    mailboxes: Vec<CoreMailbox>,
167    metrics: Arc<RuntimeMetricsInner>,
168    next_waiter_id: Arc<AtomicU64>,
169    cold_store: Option<ColdStoreHandle>,
170}
171
172/// Cluster-local summary of a bucket purge across all Raft groups.
173#[derive(Debug, Clone, Default, PartialEq, Eq)]
174pub struct PurgeBucketReport {
175    pub removed_streams: u64,
176    pub groups_with_streams: Vec<usize>,
177}
178
179impl ShardRuntime {
180    pub fn spawn(config: RuntimeConfig) -> Result<Self, RuntimeError> {
181        Self::spawn_with_engine_factory(config, InMemoryGroupEngineFactory::default())
182    }
183
184    pub fn spawn_with_engine_factory(
185        config: RuntimeConfig,
186        engine_factory: impl GroupEngineFactory,
187    ) -> Result<Self, RuntimeError> {
188        Self::spawn_with_engine_factory_and_cold_store(config, engine_factory, None)
189    }
190
191    pub fn spawn_with_engine_factory_and_cold_store(
192        config: RuntimeConfig,
193        engine_factory: impl GroupEngineFactory,
194        cold_store: Option<ColdStoreHandle>,
195    ) -> Result<Self, RuntimeError> {
196        let shard_map = StaticShardMap::new(config.core_count, config.raft_group_count)?;
197        let metrics = Arc::new(RuntimeMetricsInner::new(
198            usize::from(shard_map.core_count()),
199            usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
200        ));
201        let cold_write_admission = ColdWriteAdmission {
202            max_hot_bytes_per_group: config.cold_max_hot_bytes_per_group,
203        };
204        let raft_uncommitted_admission = RaftUncommittedAdmission {
205            max_uncommitted_bytes_per_group: config.raft_max_uncommitted_bytes_per_group,
206        };
207        let raft_uncommitted_bytes = Arc::new(RaftUncommittedBytesTracker::new(
208            usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
209        ));
210        let engine_factory: Arc<dyn GroupEngineFactory> = Arc::new(engine_factory);
211        let read_materialization = Arc::new(Semaphore::new(config.mailbox_capacity.max(1)));
212        let mut mailboxes = Vec::with_capacity(usize::from(shard_map.core_count()));
213        for raw_core_id in 0..shard_map.core_count() {
214            let core_id = CoreId(raw_core_id);
215            let (tx, rx) = mpsc::channel(config.mailbox_capacity.max(1));
216            let worker = CoreWorker {
217                core_id,
218                rx,
219                engine_factory: engine_factory.clone(),
220                groups: HashMap::new(),
221                metrics: metrics.clone(),
222                group_mailbox_capacity: config.mailbox_capacity.max(1),
223                cold_write_admission,
224                raft_uncommitted_admission,
225                raft_uncommitted_bytes: raft_uncommitted_bytes.clone(),
226                live_read_max_waiters_per_core: config.live_read_max_waiters_per_core,
227                read_materialization: read_materialization.clone(),
228            };
229            spawn_core_worker(config.threading, worker)?;
230            mailboxes.push(CoreMailbox { core_id, tx });
231        }
232        Ok(Self {
233            shard_map,
234            mailboxes,
235            metrics,
236            next_waiter_id: Arc::new(AtomicU64::new(1)),
237            cold_store,
238        })
239    }
240
241    pub fn locate(&self, stream_id: &BucketStreamId) -> ShardPlacement {
242        self.shard_map.locate(stream_id)
243    }
244
245    pub fn has_cold_store(&self) -> bool {
246        self.cold_store.is_some()
247    }
248
249    pub fn cold_store(&self) -> Option<ColdStoreHandle> {
250        self.cold_store.clone()
251    }
252
253    pub fn cold_store_info(&self) -> Option<ColdStoreInfo> {
254        self.cold_store
255            .as_ref()
256            .map(|cold_store| cold_store.info().clone())
257    }
258
259    pub async fn wait_read_stream(
260        &self,
261        request: ReadStreamRequest,
262    ) -> Result<ReadStreamResponse, RuntimeError> {
263        let placement = self.shard_map.locate(&request.stream_id);
264        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
265        let waiter_id = self.next_waiter_id.fetch_add(1, Ordering::Relaxed);
266        let stream_id = request.stream_id.clone();
267        let (response_tx, response_rx) = oneshot::channel();
268        self.enqueue_core_command(mailbox, CoreCommand::Group {
269            placement,
270            admission: None,
271            command: GroupCommand::WaitRead {
272                request,
273                waiter_id,
274                response_tx,
275            },
276        })
277        .await?;
278        let mut cancel = WaitReadCancel::new(mailbox.tx.clone(), stream_id, placement, waiter_id);
279        let response = response_rx
280            .await
281            .map_err(|_| RuntimeError::ResponseDropped {
282                core_id: mailbox.core_id,
283            })?;
284        cancel.disarm();
285        response
286    }
287
288    pub async fn require_local_live_read_owner(
289        &self,
290        stream_id: &BucketStreamId,
291    ) -> Result<(), RuntimeError> {
292        let placement = self.shard_map.locate(stream_id);
293        let (response_tx, response_rx) = oneshot::channel();
294        self.group_rpc(
295            placement,
296            None,
297            GroupCommand::RequireLiveReadOwner { response_tx },
298            response_rx,
299        )
300        .await
301    }
302
303    pub async fn flush_cold_once(
304        &self,
305        request: PlanColdFlushRequest,
306    ) -> Result<Option<FlushColdResponse>, RuntimeError> {
307        let Some(candidate) = self.plan_cold_flush(request).await? else {
308            return Ok(None);
309        };
310        self.flush_cold_candidate(candidate).await.map(Some)
311    }
312
313    pub async fn flush_cold_group_once(
314        &self,
315        raft_group_id: RaftGroupId,
316        request: PlanGroupColdFlushRequest,
317    ) -> Result<Option<FlushColdResponse>, RuntimeError> {
318        let mut candidates = self
319            .plan_next_cold_flush_batch(raft_group_id, request, 1)
320            .await?;
321        let Some(candidate) = candidates.pop() else {
322            return Ok(None);
323        };
324        match self.flush_cold_candidate(candidate).await {
325            Ok(response) => Ok(Some(response)),
326            Err(err) if is_stale_cold_flush_candidate_error(&err) => Ok(None),
327            Err(err) => Err(err),
328        }
329    }
330
331    pub async fn flush_cold_group_batch_once(
332        &self,
333        raft_group_id: RaftGroupId,
334        request: PlanGroupColdFlushRequest,
335        max_candidates: usize,
336    ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
337        let candidates = self
338            .plan_next_cold_flush_batch(raft_group_id, request, max_candidates)
339            .await?;
340        if candidates.is_empty() {
341            return Ok(Vec::new());
342        }
343        self.flush_cold_candidates_batch(candidates).await
344    }
345
346    async fn flush_cold_candidate(
347        &self,
348        candidate: ColdFlushCandidate,
349    ) -> Result<FlushColdResponse, RuntimeError> {
350        let Some(cold_store) = self.cold_store.as_ref() else {
351            return Err(RuntimeError::ColdStoreConfig {
352                message: "cold backend must be configured before flushing cold chunks".to_owned(),
353            });
354        };
355        let path = new_cold_chunk_path(
356            &candidate.stream_id,
357            candidate.start_offset,
358            candidate.end_offset,
359        );
360        let upload_started_at = Instant::now();
361        let object_size = match cold_store.write_chunk(&path, &candidate.payload).await {
362            Ok(object_size) => object_size,
363            Err(err) => {
364                // Surfaces "this node can't write to S3" to the snapshot driver's
365                // health/yield logic, which drives leadership-yield off real cold
366                // flush failures rather than a stat probe a keep-alive connection
367                // can mask.
368                self.metrics.record_cold_flush_write_error();
369                return Err(RuntimeError::ColdStoreIo {
370                    message: err.to_string(),
371                });
372            }
373        };
374        self.metrics
375            .record_cold_upload(object_size, elapsed_ns(upload_started_at));
376        let chunk = ColdChunkRef {
377            start_offset: candidate.start_offset,
378            end_offset: candidate.end_offset,
379            s3_path: path.clone(),
380            object_size,
381            object_offset: 0,
382            shared_object: false,
383            payload_digest: candidate.payload_digest,
384        };
385        let publish_started_at = Instant::now();
386        let publish = self
387            .flush_cold(FlushColdRequest {
388                stream_id: candidate.stream_id,
389                chunk,
390            })
391            .await;
392        match publish {
393            Ok(response) => {
394                self.metrics
395                    .record_cold_publish(object_size, elapsed_ns(publish_started_at));
396                Ok(response)
397            }
398            Err(err) => Err(err),
399        }
400    }
401
402    pub(crate) async fn flush_cold_candidates_batch(
403        &self,
404        candidates: Vec<ColdFlushCandidate>,
405    ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
406        if candidates.len() > 1 {
407            return self.flush_cold_candidates_pack(candidates).await;
408        }
409        let mut responses = Vec::with_capacity(candidates.len());
410        for candidate in candidates {
411            match self.flush_cold_candidate(candidate).await {
412                Ok(response) => responses.push(response),
413                Err(err) if is_stale_cold_flush_candidate_error(&err) => {}
414                Err(err) => return Err(err),
415            }
416        }
417        Ok(responses)
418    }
419
420    async fn flush_cold_candidates_pack(
421        &self,
422        candidates: Vec<ColdFlushCandidate>,
423    ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
424        let Some(cold_store) = self.cold_store.as_ref() else {
425            return Err(RuntimeError::ColdStoreConfig {
426                message: "cold backend must be configured before flushing cold chunks".to_owned(),
427            });
428        };
429        let first = candidates
430            .first()
431            .expect("packed cold flush requires at least two candidates");
432        let placement = self.shard_map.locate(&first.stream_id);
433        if candidates
434            .iter()
435            .any(|candidate| self.shard_map.locate(&candidate.stream_id) != placement)
436        {
437            return Err(RuntimeError::ColdStoreConfig {
438                message: "packed cold flush candidates must belong to one Raft group".to_owned(),
439            });
440        }
441        let payload_len = candidates.iter().try_fold(0usize, |total, candidate| {
442            total.checked_add(candidate.payload.len())
443        });
444        let Some(payload_len) = payload_len else {
445            return Err(RuntimeError::ColdStoreConfig {
446                message: "packed cold flush payload size overflow".to_owned(),
447            });
448        };
449        let mut payload = Vec::with_capacity(payload_len);
450        let mut object_offsets = Vec::with_capacity(candidates.len());
451        for candidate in &candidates {
452            object_offsets.push(u64::try_from(payload.len()).expect("pack offset fits u64"));
453            payload.extend_from_slice(&candidate.payload);
454        }
455        let path = new_cold_pack_path(placement.raft_group_id.0);
456        let upload_started_at = Instant::now();
457        let object_size = match cold_store.write_chunk(&path, &payload).await {
458            Ok(object_size) => object_size,
459            Err(err) => {
460                self.metrics.record_cold_flush_write_error();
461                return Err(RuntimeError::ColdStoreIo {
462                    message: err.to_string(),
463                });
464            }
465        };
466        self.metrics
467            .record_cold_upload(object_size, elapsed_ns(upload_started_at));
468        self.metrics.record_cold_pack(
469            object_size,
470            u64::try_from(candidates.len()).expect("candidate count fits u64"),
471        );
472
473        let mut responses = Vec::with_capacity(candidates.len());
474        let mut published = 0usize;
475        for (candidate, object_offset) in candidates.into_iter().zip(object_offsets) {
476            let logical_size = candidate.end_offset.saturating_sub(candidate.start_offset);
477            let publish_started_at = Instant::now();
478            let publish = self
479                .flush_cold(FlushColdRequest {
480                    stream_id: candidate.stream_id,
481                    chunk: ColdChunkRef {
482                        start_offset: candidate.start_offset,
483                        end_offset: candidate.end_offset,
484                        s3_path: path.clone(),
485                        object_size,
486                        object_offset,
487                        shared_object: true,
488                        payload_digest: candidate.payload_digest,
489                    },
490                })
491                .await;
492            match publish {
493                Ok(response) => {
494                    published = published.saturating_add(1);
495                    self.metrics
496                        .record_cold_publish(logical_size, elapsed_ns(publish_started_at));
497                    responses.push(response);
498                }
499                Err(err) if is_stale_cold_flush_candidate_error(&err) => {}
500                // A transport or leadership error may happen after the Raft
501                // commit became durable. Never delete the pack on an
502                // ambiguous response; a later orphan sweep may prove it
503                // unreferenced, while eager deletion could corrupt a
504                // successfully published slice.
505                Err(err) => return Err(err),
506            }
507        }
508        if published == 0 {
509            cold_store
510                .delete_chunk(&path)
511                .await
512                .map_err(|err| RuntimeError::ColdStoreIo {
513                    message: err.to_string(),
514                })?;
515        }
516        Ok(responses)
517    }
518
519    #[cfg(madsim)]
520    pub async fn flush_cold_candidates_batch_for_simulation(
521        &self,
522        candidates: Vec<ColdFlushCandidate>,
523    ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
524        self.flush_cold_candidates_batch(candidates).await
525    }
526
527    /// Sums per-bucket committed usage across every Raft group on this node.
528    ///
529    /// Groups are read serially from their local applied state (leader or
530    /// follower); usage export consumers tolerate replication lag, so this
531    /// never requires leadership and never blocks on quorum.
532    /// Purges one tenant bucket from every Raft group: streams, bucket, and
533    /// usage entries. Idempotent — a re-run over an already purged bucket
534    /// reports zero removals. Cold objects are reclaimed by the enqueued GC
535    /// entries; callers wanting synchronous reclamation run the cold GC pass
536    /// afterwards.
537    pub async fn purge_bucket_all_groups(
538        &self,
539        bucket_id: &str,
540    ) -> Result<PurgeBucketReport, RuntimeError> {
541        let mut report = PurgeBucketReport::default();
542        let group_count = self.shard_map.raft_group_count();
543        for group_id in 0..group_count {
544            let response = self
545                .purge_bucket(RaftGroupId(group_id), bucket_id.to_owned())
546                .await?;
547            if response.removed_streams > 0 {
548                report.groups_with_streams.push(group_id as usize);
549            }
550            report.removed_streams = report
551                .removed_streams
552                .saturating_add(response.removed_streams);
553        }
554        Ok(report)
555    }
556
557    pub async fn bucket_usage_all_groups(
558        &self,
559    ) -> Result<Vec<ursula_stream::BucketUsageSnapshot>, RuntimeError> {
560        let mut merged: std::collections::HashMap<String, ursula_stream::BucketUsage> =
561            std::collections::HashMap::new();
562        let group_count = self.shard_map.raft_group_count();
563        for group_id in 0..group_count {
564            let report = self.bucket_usage(RaftGroupId(group_id)).await?;
565            for entry in report {
566                let usage = merged.entry(entry.bucket_id).or_default();
567                usage.committed_append_bytes = usage
568                    .committed_append_bytes
569                    .saturating_add(entry.usage.committed_append_bytes);
570                usage.committed_records = usage
571                    .committed_records
572                    .saturating_add(entry.usage.committed_records);
573                usage.committed_write_units = usage
574                    .committed_write_units
575                    .saturating_add(entry.usage.committed_write_units);
576                usage.retained_bytes = usage
577                    .retained_bytes
578                    .saturating_add(entry.usage.retained_bytes);
579                usage.stream_count = usage.stream_count.saturating_add(entry.usage.stream_count);
580            }
581        }
582        let mut report = merged
583            .into_iter()
584            .map(|(bucket_id, usage)| ursula_stream::BucketUsageSnapshot { bucket_id, usage })
585            .collect::<Vec<_>>();
586        report.sort_by(|left, right| left.bucket_id.cmp(&right.bucket_id));
587        Ok(report)
588    }
589
590    /// Replicates one bucket's quota record to every Raft group so each
591    /// enforces the same local backstop. Serial like the other all-group
592    /// admin sweeps; quota changes are rare control-plane writes.
593    pub async fn set_bucket_quota_all_groups(
594        &self,
595        bucket_id: &str,
596        max_streams: Option<u64>,
597        max_retained_bytes: Option<u64>,
598    ) -> Result<(), RuntimeError> {
599        let group_count = self.shard_map.raft_group_count();
600        for group_id in 0..group_count {
601            self.set_bucket_quota(RaftGroupId(group_id), SetBucketQuotaRequest {
602                bucket_id: bucket_id.to_owned(),
603                max_streams,
604                max_retained_bytes,
605            })
606            .await?;
607        }
608        Ok(())
609    }
610
611    pub async fn flush_cold_all_groups_once(
612        &self,
613        request: PlanGroupColdFlushRequest,
614    ) -> Result<usize, RuntimeError> {
615        self.flush_cold_all_groups_once_bounded(request, 1).await
616    }
617
618    pub async fn flush_cold_all_groups_once_bounded(
619        &self,
620        request: PlanGroupColdFlushRequest,
621        max_concurrency: usize,
622    ) -> Result<usize, RuntimeError> {
623        let max_concurrency = max_concurrency.max(1);
624        if max_concurrency == 1 {
625            return self.flush_cold_all_groups_once_serial(request).await;
626        }
627        #[cfg(madsim)]
628        {
629            return self.flush_cold_all_groups_once_serial(request).await;
630        }
631        #[cfg(not(madsim))]
632        {
633            let mut flushed = 0;
634            let mut next_group_id = 0;
635            let group_count = self.shard_map.raft_group_count();
636            let mut tasks = JoinSet::new();
637
638            while next_group_id < group_count || !tasks.is_empty() {
639                while next_group_id < group_count && tasks.len() < max_concurrency {
640                    let runtime = self.clone();
641                    let request = request.clone();
642                    let group_id = RaftGroupId(next_group_id);
643                    next_group_id += 1;
644                    tasks.spawn(async move {
645                        runtime
646                            .flush_cold_group_batch_once(
647                                group_id,
648                                request,
649                                COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
650                            )
651                            .await
652                            .map(|responses| responses.len())
653                    });
654                }
655                if let Some(result) = tasks.join_next().await {
656                    match result {
657                        Ok(Ok(count)) => flushed += count,
658                        Ok(Err(err)) => return Err(err),
659                        Err(err) => {
660                            return Err(RuntimeError::ColdStoreIo {
661                                message: format!("cold flush task failed: {err}"),
662                            });
663                        }
664                    }
665                }
666            }
667            Ok(flushed)
668        }
669    }
670
671    async fn flush_cold_all_groups_once_serial(
672        &self,
673        request: PlanGroupColdFlushRequest,
674    ) -> Result<usize, RuntimeError> {
675        let mut flushed = 0;
676        for group_id in 0..self.shard_map.raft_group_count() {
677            flushed += self
678                .flush_cold_group_batch_once(
679                    RaftGroupId(group_id),
680                    request.clone(),
681                    COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
682                )
683                .await?
684                .len();
685        }
686        Ok(flushed)
687    }
688
689    /// Rewrites undersized, contiguous objects from the same stream into
690    /// target-sized immutable chunks. Discovery reads only cold-index pages.
691    pub async fn compact_cold_once(
692        &self,
693        target_bytes: u64,
694        max_bytes: u64,
695        max_streams: usize,
696        gc_grace_ms: u64,
697    ) -> Result<usize, RuntimeError> {
698        let Some(cold_store) = self.cold_store.as_ref() else {
699            return Ok(0);
700        };
701        let pages =
702            cold_store
703                .list_cold_index_pages()
704                .await
705                .map_err(|err| RuntimeError::ColdStoreIo {
706                    message: err.to_string(),
707                })?;
708        let mut pages_by_stream: HashMap<BucketStreamId, Vec<_>> = HashMap::new();
709        for page in pages {
710            if page.generation == 0 {
711                pages_by_stream
712                    .entry(page.stream_id.clone())
713                    .or_default()
714                    .push(page);
715            }
716        }
717        let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
718        let mut compacted = 0;
719        for (stream_id, stream_pages) in pages_by_stream {
720            if compacted >= max_streams {
721                break;
722            }
723            // Only the local Raft leader may publish a replacement.
724            if self
725                .require_local_live_read_owner(&stream_id)
726                .await
727                .is_err()
728            {
729                continue;
730            }
731            let chunks = load_cold_chunks_from_pages(&store, &stream_pages)
732                .await
733                .map_err(|err| RuntimeError::ColdStoreIo {
734                    message: err.to_string(),
735                })?;
736            let Some(old_chunks) = select_cold_chunk_compaction(&chunks, target_bytes, max_bytes)
737            else {
738                continue;
739            };
740            let total_bytes = old_chunks
741                .iter()
742                .try_fold(0_u64, |total, chunk| total.checked_add(chunk.object_size))
743                .ok_or_else(|| RuntimeError::ColdStoreIo {
744                    message: "cold compaction byte count overflow".to_owned(),
745                })?;
746            let capacity = usize::try_from(total_bytes).map_err(|_| RuntimeError::ColdStoreIo {
747                message: "cold compaction object exceeds addressable memory".to_owned(),
748            })?;
749            let mut payload = Vec::with_capacity(capacity);
750            for chunk in &old_chunks {
751                let len =
752                    usize::try_from(chunk.object_size).map_err(|_| RuntimeError::ColdStoreIo {
753                        message: "cold chunk exceeds addressable memory".to_owned(),
754                    })?;
755                let bytes = cold_store
756                    .read_chunk_range(chunk, chunk.start_offset, len)
757                    .await
758                    .map_err(|err| RuntimeError::ColdStoreIo {
759                        message: err.to_string(),
760                    })?;
761                payload.extend_from_slice(&bytes);
762            }
763            let first = old_chunks
764                .first()
765                .expect("candidate contains at least two chunks");
766            let last = old_chunks
767                .last()
768                .expect("candidate contains at least two chunks");
769            let path = new_cold_chunk_path(&stream_id, first.start_offset, last.end_offset);
770            let object_size = cold_store
771                .write_chunk(&path, &payload)
772                .await
773                .map_err(|err| RuntimeError::ColdStoreIo {
774                    message: err.to_string(),
775                })?;
776            let replacement = ColdChunkRef {
777                start_offset: first.start_offset,
778                end_offset: last.end_offset,
779                object_size,
780                s3_path: path,
781                object_offset: 0,
782                shared_object: false,
783                payload_digest: blake3::hash(&payload).to_hex().to_string(),
784            };
785            let replacement_path = replacement.s3_path.clone();
786            let gc_not_before_ms = unix_time_ms().saturating_add(gc_grace_ms);
787            let compact_result = self
788                .compact_cold(CompactColdRequest {
789                    stream_id: stream_id.clone(),
790                    old_chunks,
791                    replacement,
792                    gc_not_before_ms,
793                })
794                .await;
795            if let Err(err) = compact_result {
796                let rollback_safe =
797                    err.leader_hint().is_some() || err.stream_error_code().is_some();
798                if !rollback_safe {
799                    return Err(err);
800                }
801                if let Err(cleanup_err) = cold_store.delete_chunk(&replacement_path).await {
802                    tracing::warn!(
803                        stream = %stream_id,
804                        path = %replacement_path,
805                        error = %cleanup_err,
806                        "failed to remove unpublished cold compaction replacement"
807                    );
808                }
809                tracing::warn!(
810                    stream = %stream_id,
811                    error = %err,
812                    "cold compaction publish failed; continuing with remaining streams"
813                );
814                continue;
815            }
816            compacted += 1;
817        }
818        Ok(compacted)
819    }
820
821    /// Drains the leader-side cold-GC queue for one group: physically reclaims
822    /// each queued target from cold storage, then replicates an ack that pops
823    /// the reclaimed entries. Deletions are idempotent, so a crash or leader
824    /// change between reclaim and ack simply re-runs them next tick.
825    pub async fn run_cold_gc_group_once(
826        &self,
827        raft_group_id: RaftGroupId,
828        max_entries: usize,
829    ) -> Result<usize, RuntimeError> {
830        let Some(cold_store) = self.cold_store.as_ref() else {
831            return Ok(0);
832        };
833        let entries = self.plan_cold_gc(raft_group_id, max_entries).await?;
834        if entries.is_empty() {
835            return Ok(0);
836        }
837        let mut acked_seq = None;
838        let mut reclaimed = 0usize;
839        // Entries are FIFO by seq; stop at the first failure so the ack never
840        // skips past an object that is still present in cold storage.
841        for entry in entries {
842            if entry.not_before_ms > unix_time_ms() {
843                break;
844            }
845            let result = match &entry.target {
846                ColdGcTarget::Stream(stream_id) => {
847                    match cold_store.remove_all(&cold_chunk_prefix(stream_id)).await {
848                        Ok(()) => cold_store.remove_all(&cold_index_prefix(stream_id)).await,
849                        Err(err) => Err(err),
850                    }
851                }
852                ColdGcTarget::Paths(paths) => {
853                    let mut outcome = Ok(());
854                    for path in paths {
855                        if let Err(err) = cold_store.delete_chunk(path).await {
856                            outcome = Err(err);
857                            break;
858                        }
859                    }
860                    outcome
861                }
862            };
863            match result {
864                Ok(()) => {
865                    acked_seq = Some(entry.seq);
866                    reclaimed += 1;
867                }
868                Err(err) => {
869                    self.metrics.record_cold_gc_error();
870                    if acked_seq.is_none() {
871                        return Err(RuntimeError::ColdStoreIo {
872                            message: err.to_string(),
873                        });
874                    }
875                    break;
876                }
877            }
878        }
879        if let Some(up_to_seq) = acked_seq {
880            self.ack_cold_gc(raft_group_id, up_to_seq).await?;
881            self.metrics
882                .record_cold_gc_reclaimed(u64::try_from(reclaimed).expect("reclaimed fits u64"));
883        }
884        Ok(reclaimed)
885    }
886
887    pub async fn run_cold_gc_all_groups_once(
888        &self,
889        max_entries_per_group: usize,
890    ) -> Result<usize, RuntimeError> {
891        if self.cold_store.is_none() {
892            return Ok(0);
893        }
894        let mut reclaimed = 0;
895        for group_id in 0..self.shard_map.raft_group_count() {
896            reclaimed += self
897                .run_cold_gc_group_once(RaftGroupId(group_id), max_entries_per_group)
898                .await?;
899        }
900        Ok(reclaimed)
901    }
902
903    /// Number of raft groups this runtime is sharded into. Backup tooling
904    /// iterates `0..raft_group_count()` to cover the whole keyspace.
905    pub fn raft_group_count(&self) -> u32 {
906        self.shard_map.raft_group_count()
907    }
908
909    pub async fn install_group_snapshot(
910        &self,
911        snapshot: GroupSnapshot,
912    ) -> Result<(), RuntimeError> {
913        let expected = self.placement_for_group(snapshot.placement.raft_group_id)?;
914        if snapshot.placement != expected {
915            return Err(RuntimeError::SnapshotPlacementMismatch {
916                expected,
917                actual: snapshot.placement,
918            });
919        }
920        let (response_tx, response_rx) = oneshot::channel();
921        self.group_rpc(
922            expected,
923            None,
924            GroupCommand::InstallGroupSnapshot {
925                snapshot,
926                response_tx,
927            },
928            response_rx,
929        )
930        .await
931    }
932
933    #[cfg(madsim)]
934    pub async fn shutdown_group_engine_for_simulation(
935        &self,
936        placement: ShardPlacement,
937    ) -> Result<(), RuntimeError> {
938        let expected = self.placement_for_group(placement.raft_group_id)?;
939        if placement != expected {
940            return Err(RuntimeError::SnapshotPlacementMismatch {
941                expected,
942                actual: placement,
943            });
944        }
945        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
946        let (response_tx, response_rx) = oneshot::channel();
947        self.send_core_command(
948            mailbox,
949            CoreCommand::ShutdownGroupEngine {
950                placement,
951                response_tx,
952            },
953            response_rx,
954        )
955        .await
956    }
957
958    #[cfg(madsim)]
959    pub async fn install_group_engine_for_simulation(
960        &self,
961        placement: ShardPlacement,
962        engine: Box<dyn crate::engine::GroupEngine>,
963    ) -> Result<(), RuntimeError> {
964        let expected = self.placement_for_group(placement.raft_group_id)?;
965        if placement != expected {
966            return Err(RuntimeError::SnapshotPlacementMismatch {
967                expected,
968                actual: placement,
969            });
970        }
971        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
972        let (response_tx, response_rx) = oneshot::channel();
973        self.send_core_command(
974            mailbox,
975            CoreCommand::InstallGroupEngine {
976                placement,
977                engine,
978                response_tx,
979            },
980            response_rx,
981        )
982        .await
983    }
984
985    pub async fn warm_group(
986        &self,
987        raft_group_id: RaftGroupId,
988    ) -> Result<ShardPlacement, RuntimeError> {
989        let placement = self.placement_for_group(raft_group_id)?;
990        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
991        let (response_tx, response_rx) = oneshot::channel();
992        self.send_core_command(
993            mailbox,
994            CoreCommand::WarmGroup {
995                placement,
996                response_tx,
997            },
998            response_rx,
999        )
1000        .await
1001    }
1002
1003    pub async fn warm_all_groups(&self) -> Result<(), RuntimeError> {
1004        let mut placements_by_core = vec![Vec::new(); self.mailboxes.len()];
1005        for raw_group_id in 0..self.shard_map.raft_group_count() {
1006            let placement = self.placement_for_group(RaftGroupId(raw_group_id))?;
1007            placements_by_core[usize::from(placement.core_id.0)].push(placement);
1008        }
1009        let mut responses = Vec::new();
1010        for (mailbox, placements) in self.mailboxes.iter().zip(placements_by_core) {
1011            let (response_tx, response_rx) = oneshot::channel();
1012            self.enqueue_core_command(mailbox, CoreCommand::WarmGroups {
1013                placements,
1014                response_tx,
1015            })
1016            .await?;
1017            responses.push((mailbox.core_id, response_rx));
1018        }
1019        for (core_id, response_rx) in responses {
1020            response_rx
1021                .await
1022                .map_err(|_| RuntimeError::ResponseDropped { core_id })??;
1023        }
1024        Ok(())
1025    }
1026
1027    fn placement_for_group(
1028        &self,
1029        raft_group_id: RaftGroupId,
1030    ) -> Result<ShardPlacement, RuntimeError> {
1031        if raft_group_id.0 >= self.shard_map.raft_group_count() {
1032            return Err(RuntimeError::InvalidRaftGroup {
1033                raft_group_id,
1034                raft_group_count: self.shard_map.raft_group_count(),
1035            });
1036        }
1037        Ok(ShardPlacement {
1038            core_id: CoreId(
1039                (raft_group_id.0 % u32::from(self.shard_map.core_count()))
1040                    .try_into()
1041                    .expect("core id fits u16"),
1042            ),
1043            shard_id: ShardId(raft_group_id.0),
1044            raft_group_id,
1045        })
1046    }
1047
1048    /// Routes a group command to its owning core and awaits the reply.
1049    /// `admission` carries the incoming payload bytes for
1050    /// raft-uncommitted-backpressure-guarded writes; the check itself runs on
1051    /// the owning core.
1052    async fn group_rpc<T>(
1053        &self,
1054        placement: ShardPlacement,
1055        admission: Option<u64>,
1056        command: GroupCommand,
1057        response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
1058    ) -> Result<T, RuntimeError> {
1059        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1060        self.send_core_command(
1061            mailbox,
1062            CoreCommand::Group {
1063                placement,
1064                admission,
1065                command,
1066            },
1067            response_rx,
1068        )
1069        .await
1070    }
1071
1072    async fn send_core_command<T>(
1073        &self,
1074        mailbox: &CoreMailbox,
1075        command: CoreCommand,
1076        response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
1077    ) -> Result<T, RuntimeError> {
1078        self.enqueue_core_command(mailbox, command).await?;
1079        response_rx
1080            .await
1081            .map_err(|_| RuntimeError::ResponseDropped {
1082                core_id: mailbox.core_id,
1083            })?
1084    }
1085
1086    async fn enqueue_core_command(
1087        &self,
1088        mailbox: &CoreMailbox,
1089        command: CoreCommand,
1090    ) -> Result<(), RuntimeError> {
1091        if mailbox.tx.capacity() == 0 {
1092            self.metrics.record_mailbox_full(mailbox.core_id);
1093        }
1094        let started_at = Instant::now();
1095        mailbox
1096            .tx
1097            .send(Traced::capture(command))
1098            .await
1099            .map_err(|_| RuntimeError::MailboxClosed {
1100                core_id: mailbox.core_id,
1101            })?;
1102        self.metrics
1103            .record_routed_request(mailbox.core_id, elapsed_ns(started_at));
1104        Ok(())
1105    }
1106
1107    /// Atomically appends to affinity-grouped streams in one Raft group.
1108    ///
1109    /// This operation deliberately rejects ungrouped or differently grouped
1110    /// streams: the runtime has no cross-group transaction coordinator.
1111    pub async fn append_transaction(
1112        &self,
1113        request: AppendTransactionRequest,
1114    ) -> Result<AppendTransactionResponse, RuntimeError> {
1115        const MAX_OPERATIONS: usize = 64;
1116        let Some(first) = request.operations.first() else {
1117            return Err(RuntimeError::InvalidAppendTransaction {
1118                message: "at least one append operation is required".to_owned(),
1119            });
1120        };
1121        if request.operations.len() > MAX_OPERATIONS {
1122            return Err(RuntimeError::InvalidAppendTransaction {
1123                message: format!("at most {MAX_OPERATIONS} append operations are allowed"),
1124            });
1125        }
1126        if request
1127            .operations
1128            .iter()
1129            .any(|operation| operation.payload.is_empty())
1130        {
1131            return Err(RuntimeError::InvalidAppendTransaction {
1132                message: "every append operation must have a non-empty payload".to_owned(),
1133            });
1134        }
1135        let Some(affinity_key) = first.stream_id.affinity_key.as_deref() else {
1136            return Err(RuntimeError::InvalidAppendTransaction {
1137                message: "all streams must use an affinity path".to_owned(),
1138            });
1139        };
1140        let placement = self.shard_map.locate(&first.stream_id);
1141        for operation in &request.operations {
1142            if operation.stream_id.bucket_id != first.stream_id.bucket_id
1143                || operation.stream_id.affinity_key.as_deref() != Some(affinity_key)
1144                || self.shard_map.locate(&operation.stream_id) != placement
1145            {
1146                return Err(RuntimeError::InvalidAppendTransaction {
1147                    message: "all streams must share one bucket and affinity key".to_owned(),
1148                });
1149            }
1150        }
1151        let incoming_bytes = request.payload_bytes();
1152        let (response_tx, response_rx) = oneshot::channel();
1153        self.group_rpc(
1154            placement,
1155            Some(incoming_bytes),
1156            GroupCommand::AppendTransaction {
1157                request,
1158                response_tx,
1159                raft_uncommitted: None,
1160            },
1161            response_rx,
1162        )
1163        .await
1164    }
1165
1166    pub fn metrics(&self) -> RuntimeMetrics {
1167        RuntimeMetrics {
1168            inner: self.metrics.clone(),
1169        }
1170    }
1171
1172    pub fn mailbox_snapshot(&self) -> RuntimeMailboxSnapshot {
1173        let depths = self
1174            .mailboxes
1175            .iter()
1176            .map(CoreMailbox::depth)
1177            .collect::<Vec<_>>();
1178        let capacities = self
1179            .mailboxes
1180            .iter()
1181            .map(CoreMailbox::capacity)
1182            .collect::<Vec<_>>();
1183        RuntimeMailboxSnapshot { depths, capacities }
1184    }
1185}
1186
1187#[cfg(not(madsim))]
1188fn unix_time_ms() -> u64 {
1189    SystemTime::now()
1190        .duration_since(UNIX_EPOCH)
1191        .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
1192        .unwrap_or(0)
1193}
1194
1195#[cfg(madsim)]
1196fn unix_time_ms() -> u64 {
1197    0
1198}
1199
1200/// Expands the operation manifest into the uniform `ShardRuntime` client
1201/// methods: locate the placement (by stream id or raft group id), open the
1202/// reply channel, and submit the `GroupCommand` through [`ShardRuntime::
1203/// group_rpc`]. Entries with `client { none }` keep hand-written methods; see
1204/// the manifest grammar in [`crate::ops`].
1205macro_rules! shard_runtime_operations {
1206    // Stream-routed write with admission and an emptiness precheck.
1207    (@munch
1208        methods { $($methods:tt)* }
1209        rest {
1210            $(#[$attr:meta])*
1211            op $Variant:ident {
1212                fields { $req:ident: $Req:ty $(,)? }
1213                reply { $tx:ident: $Resp:ty }
1214                guard { $g:ident }
1215                handle { $($handle:tt)* }
1216                client {
1217                    $vis:vis stream fn $method:ident,
1218                    non_empty: $ne:ident,
1219                    admit: $incoming:expr
1220                }
1221            }
1222            $($rest:tt)*
1223        }
1224    ) => {
1225        shard_runtime_operations! {
1226            @munch
1227            methods {
1228                $($methods)*
1229                $(#[$attr])*
1230                $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
1231                    if $req.$ne.is_empty() {
1232                        return Err(RuntimeError::EmptyAppend);
1233                    }
1234                    let placement = self.shard_map.locate(&$req.stream_id);
1235                    let incoming_bytes = $incoming;
1236                    let (response_tx, response_rx) = oneshot::channel();
1237                    self.group_rpc(
1238                        placement,
1239                        Some(incoming_bytes),
1240                        GroupCommand::$Variant {
1241                            $req,
1242                            $tx: response_tx,
1243                            $g: None,
1244                        },
1245                        response_rx,
1246                    )
1247                    .await
1248                }
1249            }
1250            rest { $($rest)* }
1251        }
1252    };
1253    // Stream-routed write with admission.
1254    (@munch
1255        methods { $($methods:tt)* }
1256        rest {
1257            $(#[$attr:meta])*
1258            op $Variant:ident {
1259                fields { $req:ident: $Req:ty $(,)? }
1260                reply { $tx:ident: $Resp:ty }
1261                guard { $g:ident }
1262                handle { $($handle:tt)* }
1263                client { $vis:vis stream fn $method:ident, admit: $incoming:expr }
1264            }
1265            $($rest:tt)*
1266        }
1267    ) => {
1268        shard_runtime_operations! {
1269            @munch
1270            methods {
1271                $($methods)*
1272                $(#[$attr])*
1273                $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
1274                    let placement = self.shard_map.locate(&$req.stream_id);
1275                    let incoming_bytes = $incoming;
1276                    let (response_tx, response_rx) = oneshot::channel();
1277                    self.group_rpc(
1278                        placement,
1279                        Some(incoming_bytes),
1280                        GroupCommand::$Variant {
1281                            $req,
1282                            $tx: response_tx,
1283                            $g: None,
1284                        },
1285                        response_rx,
1286                    )
1287                    .await
1288                }
1289            }
1290            rest { $($rest)* }
1291        }
1292    };
1293    // Stream-routed operation without admission.
1294    (@munch
1295        methods { $($methods:tt)* }
1296        rest {
1297            $(#[$attr:meta])*
1298            op $Variant:ident {
1299                fields { $req:ident: $Req:ty $(,)? }
1300                reply { $tx:ident: $Resp:ty }
1301                guard { none }
1302                handle { $($handle:tt)* }
1303                client { $vis:vis stream fn $method:ident }
1304            }
1305            $($rest:tt)*
1306        }
1307    ) => {
1308        shard_runtime_operations! {
1309            @munch
1310            methods {
1311                $($methods)*
1312                $(#[$attr])*
1313                $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
1314                    let placement = self.shard_map.locate(&$req.stream_id);
1315                    let (response_tx, response_rx) = oneshot::channel();
1316                    self.group_rpc(
1317                        placement,
1318                        None,
1319                        GroupCommand::$Variant {
1320                            $req,
1321                            $tx: response_tx,
1322                        },
1323                        response_rx,
1324                    )
1325                    .await
1326                }
1327            }
1328            rest { $($rest)* }
1329        }
1330    };
1331    // Group-routed operation: an explicit `RaftGroupId` followed by the
1332    // manifest fields in declaration order.
1333    (@munch
1334        methods { $($methods:tt)* }
1335        rest {
1336            $(#[$attr:meta])*
1337            op $Variant:ident {
1338                fields { $($field:ident: $field_ty:ty),* $(,)? }
1339                reply { $tx:ident: $Resp:ty }
1340                guard { none }
1341                handle { $($handle:tt)* }
1342                client { $vis:vis group fn $method:ident }
1343            }
1344            $($rest:tt)*
1345        }
1346    ) => {
1347        shard_runtime_operations! {
1348            @munch
1349            methods {
1350                $($methods)*
1351                $(#[$attr])*
1352                $vis async fn $method(
1353                    &self,
1354                    raft_group_id: RaftGroupId,
1355                    $($field: $field_ty),*
1356                ) -> Result<$Resp, RuntimeError> {
1357                    let placement = self.placement_for_group(raft_group_id)?;
1358                    let (response_tx, response_rx) = oneshot::channel();
1359                    self.group_rpc(
1360                        placement,
1361                        None,
1362                        GroupCommand::$Variant {
1363                            $($field,)*
1364                            $tx: response_tx,
1365                        },
1366                        response_rx,
1367                    )
1368                    .await
1369                }
1370            }
1371            rest { $($rest)* }
1372        }
1373    };
1374    // Hand-written client method: nothing to generate.
1375    (@munch
1376        methods { $($methods:tt)* }
1377        rest {
1378            $(#[$attr:meta])*
1379            op $Variant:ident {
1380                fields { $($fields:tt)* }
1381                reply { $($reply:tt)* }
1382                guard { $($guard:tt)* }
1383                handle { $($handle:tt)* }
1384                client { none }
1385            }
1386            $($rest:tt)*
1387        }
1388    ) => {
1389        shard_runtime_operations! {
1390            @munch
1391            methods { $($methods)* }
1392            rest { $($rest)* }
1393        }
1394    };
1395    (@munch
1396        methods { $($methods:tt)* }
1397        rest {}
1398    ) => {
1399        impl ShardRuntime {
1400            $($methods)*
1401        }
1402    };
1403    ($($manifest:tt)*) => {
1404        shard_runtime_operations! {
1405            @munch
1406            methods {}
1407            rest { $($manifest)* }
1408        }
1409    };
1410}
1411
1412crate::ops::runtime_operations!(shard_runtime_operations);
1413
1414fn spawn_core_worker(threading: RuntimeThreading, worker: CoreWorker) -> Result<(), RuntimeError> {
1415    match threading {
1416        RuntimeThreading::HostedTokio => {
1417            crate::rt::spawn(worker.run());
1418            Ok(())
1419        }
1420        #[cfg(not(madsim))]
1421        RuntimeThreading::ThreadPerCore => {
1422            let core_id = worker.core_id;
1423            std::thread::Builder::new()
1424                .name(format!("ursula-core-{}", core_id.0))
1425                .spawn(move || {
1426                    let runtime = tokio::runtime::Builder::new_current_thread()
1427                        .enable_all()
1428                        .build()
1429                        .expect("build per-core tokio runtime");
1430                    runtime.block_on(worker.run());
1431                })
1432                .map(|_| ())
1433                .map_err(|err| RuntimeError::SpawnCoreThread {
1434                    core_id,
1435                    message: err.to_string(),
1436                })
1437        }
1438    }
1439}