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