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