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
6use bytes::Bytes;
7#[cfg(not(madsim))]
8use tokio::task::JoinSet;
9use ursula_shard::BucketStreamId;
10use ursula_shard::CoreId;
11use ursula_shard::RaftGroupId;
12use ursula_shard::ShardId;
13use ursula_shard::ShardPlacement;
14use ursula_shard::StaticShardMap;
15use ursula_stream::ColdChunkRef;
16use ursula_stream::ColdFlushCandidate;
17use ursula_stream::ColdGcEntry;
18use ursula_stream::ColdGcTarget;
19use ursula_stream::StreamErrorCode;
20
21use crate::admission::RaftUncommittedAdmission;
22use crate::admission::RaftUncommittedBytesTracker;
23use crate::cold_index::cold_index_prefix;
24use crate::cold_store::ColdStoreHandle;
25use crate::cold_store::ColdStoreInfo;
26use crate::cold_store::cold_chunk_prefix;
27use crate::cold_store::new_cold_chunk_path;
28use crate::command::GroupSnapshot;
29use crate::core_worker::CoreCommand;
30use crate::core_worker::CoreMailbox;
31use crate::core_worker::CoreWorker;
32use crate::core_worker::WaitReadCancel;
33use crate::engine::GroupEngineError;
34use crate::engine::GroupEngineFactory;
35use crate::engine::in_memory::InMemoryGroupEngineFactory;
36use crate::error::RuntimeError;
37use crate::error::map_fork_source_ref_error;
38use crate::metrics::COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS;
39use crate::metrics::RuntimeMailboxSnapshot;
40use crate::metrics::RuntimeMetrics;
41use crate::metrics::RuntimeMetricsInner;
42use crate::metrics::elapsed_ns;
43use crate::metrics::is_stale_cold_flush_candidate_error;
44use crate::request::AckColdGcResponse;
45use crate::request::AppendBatchRequest;
46use crate::request::AppendBatchResponse;
47use crate::request::AppendExternalRequest;
48use crate::request::AppendRequest;
49use crate::request::AppendResponse;
50use crate::request::BootstrapStreamRequest;
51use crate::request::BootstrapStreamResponse;
52use crate::request::CloseStreamRequest;
53use crate::request::CloseStreamResponse;
54use crate::request::ColdWriteAdmission;
55use crate::request::CreateStreamExternalRequest;
56use crate::request::CreateStreamRequest;
57use crate::request::CreateStreamResponse;
58use crate::request::DeleteSnapshotRequest;
59use crate::request::DeleteStreamRequest;
60use crate::request::DeleteStreamResponse;
61use crate::request::FlushColdRequest;
62use crate::request::FlushColdResponse;
63use crate::request::ForkRefResponse;
64use crate::request::GetStreamAttrsRequest;
65use crate::request::GetStreamAttrsResponse;
66use crate::request::HeadStreamRequest;
67use crate::request::HeadStreamResponse;
68use crate::request::PlanColdFlushRequest;
69use crate::request::PlanGroupColdFlushRequest;
70use crate::request::PublishSnapshotRequest;
71use crate::request::PublishSnapshotResponse;
72use crate::request::ReadSnapshotRequest;
73use crate::request::ReadSnapshotResponse;
74use crate::request::ReadStreamRequest;
75use crate::request::ReadStreamResponse;
76use crate::request::UpdateStreamAttrsRequest;
77use crate::request::UpdateStreamAttrsResponse;
78use crate::rt::sync::Semaphore;
79use crate::rt::sync::mpsc;
80use crate::rt::sync::oneshot;
81use crate::rt::time::Instant;
82use crate::trace::Traced;
83
84#[derive(Debug, Clone)]
85pub struct RuntimeConfig {
86    pub core_count: usize,
87    pub raft_group_count: usize,
88    pub mailbox_capacity: usize,
89    pub threading: RuntimeThreading,
90    pub cold_max_hot_bytes_per_group: Option<u64>,
91    /// Per-group cap on raft-submitted-but-not-yet-applied payload bytes.
92    /// `None` disables the admission (default). Catches "raft replication slow"
93    /// before in-memory queues grow unbounded.
94    pub raft_max_uncommitted_bytes_per_group: Option<u64>,
95    pub live_read_max_waiters_per_core: Option<u64>,
96}
97
98impl RuntimeConfig {
99    pub fn new(core_count: usize, raft_group_count: usize) -> Self {
100        #[cfg(not(madsim))]
101        let threading = RuntimeThreading::ThreadPerCore;
102        #[cfg(madsim)]
103        let threading = RuntimeThreading::HostedTokio;
104        Self {
105            core_count,
106            raft_group_count,
107            mailbox_capacity: 1024,
108            threading,
109            cold_max_hot_bytes_per_group: None,
110            raft_max_uncommitted_bytes_per_group: None,
111            live_read_max_waiters_per_core: Some(65_536),
112        }
113    }
114
115    pub fn with_cold_max_hot_bytes_per_group(mut self, value: Option<u64>) -> Self {
116        self.cold_max_hot_bytes_per_group = value;
117        self
118    }
119
120    pub fn with_raft_max_uncommitted_bytes_per_group(mut self, value: Option<u64>) -> Self {
121        self.raft_max_uncommitted_bytes_per_group = value;
122        self
123    }
124
125    pub fn with_live_read_max_waiters_per_core(mut self, value: Option<u64>) -> Self {
126        self.live_read_max_waiters_per_core = value;
127        self
128    }
129
130    /// Build runtime configuration from a typed `ursula_config::RuntimeConfig`.
131    pub fn from_ursula_config(cfg: &ursula_config::RuntimeConfig, raft_group_count: usize) -> Self {
132        let mut config = Self::new(cfg.core_count, raft_group_count);
133        config.live_read_max_waiters_per_core = cfg
134            .live_read_max_waiters_per_core
135            .and_then(|n| if n == 0 { None } else { Some(n as u64) });
136        config
137    }
138}
139
140#[derive(Debug, Clone, Copy, PartialEq, Eq)]
141pub enum RuntimeThreading {
142    #[cfg(not(madsim))]
143    ThreadPerCore,
144    HostedTokio,
145}
146
147#[derive(Debug, Clone)]
148pub struct ShardRuntime {
149    shard_map: StaticShardMap,
150    mailboxes: Vec<CoreMailbox>,
151    metrics: Arc<RuntimeMetricsInner>,
152    next_waiter_id: Arc<AtomicU64>,
153    cold_store: Option<ColdStoreHandle>,
154}
155
156impl ShardRuntime {
157    pub fn spawn(config: RuntimeConfig) -> Result<Self, RuntimeError> {
158        Self::spawn_with_engine_factory(config, InMemoryGroupEngineFactory::default())
159    }
160
161    pub fn spawn_with_engine_factory(
162        config: RuntimeConfig,
163        engine_factory: impl GroupEngineFactory,
164    ) -> Result<Self, RuntimeError> {
165        Self::spawn_with_engine_factory_and_cold_store(config, engine_factory, None)
166    }
167
168    pub fn spawn_with_engine_factory_and_cold_store(
169        config: RuntimeConfig,
170        engine_factory: impl GroupEngineFactory,
171        cold_store: Option<ColdStoreHandle>,
172    ) -> Result<Self, RuntimeError> {
173        let shard_map = StaticShardMap::new(config.core_count, config.raft_group_count)?;
174        let metrics = Arc::new(RuntimeMetricsInner::new(
175            usize::from(shard_map.core_count()),
176            usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
177        ));
178        let cold_write_admission = ColdWriteAdmission {
179            max_hot_bytes_per_group: config.cold_max_hot_bytes_per_group,
180        };
181        let raft_uncommitted_admission = RaftUncommittedAdmission {
182            max_uncommitted_bytes_per_group: config.raft_max_uncommitted_bytes_per_group,
183        };
184        let raft_uncommitted_bytes = Arc::new(RaftUncommittedBytesTracker::new(
185            usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
186        ));
187        let engine_factory: Arc<dyn GroupEngineFactory> = Arc::new(engine_factory);
188        let read_materialization = Arc::new(Semaphore::new(config.mailbox_capacity.max(1)));
189        let mut mailboxes = Vec::with_capacity(usize::from(shard_map.core_count()));
190        for raw_core_id in 0..shard_map.core_count() {
191            let core_id = CoreId(raw_core_id);
192            let (tx, rx) = mpsc::channel(config.mailbox_capacity.max(1));
193            let worker = CoreWorker {
194                core_id,
195                rx,
196                engine_factory: engine_factory.clone(),
197                groups: HashMap::new(),
198                metrics: metrics.clone(),
199                group_mailbox_capacity: config.mailbox_capacity.max(1),
200                cold_write_admission,
201                raft_uncommitted_admission,
202                raft_uncommitted_bytes: raft_uncommitted_bytes.clone(),
203                live_read_max_waiters_per_core: config.live_read_max_waiters_per_core,
204                read_materialization: read_materialization.clone(),
205            };
206            spawn_core_worker(config.threading, worker)?;
207            mailboxes.push(CoreMailbox { core_id, tx });
208        }
209        Ok(Self {
210            shard_map,
211            mailboxes,
212            metrics,
213            next_waiter_id: Arc::new(AtomicU64::new(1)),
214            cold_store,
215        })
216    }
217
218    pub fn locate(&self, stream_id: &BucketStreamId) -> ShardPlacement {
219        self.shard_map.locate(stream_id)
220    }
221
222    pub fn has_cold_store(&self) -> bool {
223        self.cold_store.is_some()
224    }
225
226    pub fn cold_store(&self) -> Option<ColdStoreHandle> {
227        self.cold_store.clone()
228    }
229
230    pub fn cold_store_info(&self) -> Option<ColdStoreInfo> {
231        self.cold_store
232            .as_ref()
233            .map(|cold_store| cold_store.info().clone())
234    }
235
236    pub async fn create_stream(
237        &self,
238        request: CreateStreamRequest,
239    ) -> Result<CreateStreamResponse, RuntimeError> {
240        if request.forked_from.is_some() {
241            return self.create_fork_stream(request).await;
242        }
243        self.create_stream_on_owner(request).await
244    }
245
246    pub async fn create_stream_external(
247        &self,
248        request: CreateStreamExternalRequest,
249    ) -> Result<CreateStreamResponse, RuntimeError> {
250        let placement = self.shard_map.locate(&request.stream_id);
251        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
252        let (response_tx, response_rx) = oneshot::channel();
253        self.send_core_command(
254            mailbox,
255            CoreCommand::CreateExternal {
256                request,
257                placement,
258                response_tx,
259            },
260            response_rx,
261        )
262        .await
263    }
264
265    async fn create_stream_on_owner(
266        &self,
267        request: CreateStreamRequest,
268    ) -> Result<CreateStreamResponse, RuntimeError> {
269        let placement = self.shard_map.locate(&request.stream_id);
270        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
271        let (response_tx, response_rx) = oneshot::channel();
272        self.send_core_command(
273            mailbox,
274            CoreCommand::CreateStream {
275                request,
276                placement,
277                response_tx,
278            },
279            response_rx,
280        )
281        .await
282    }
283
284    async fn create_fork_stream(
285        &self,
286        mut request: CreateStreamRequest,
287    ) -> Result<CreateStreamResponse, RuntimeError> {
288        let source_id = request
289            .forked_from
290            .clone()
291            .expect("forked_from checked before create_fork_stream");
292        let now_ms = request.now_ms;
293        let source_placement = self.shard_map.locate(&source_id);
294        let source_head = self
295            .head_stream(HeadStreamRequest {
296                stream_id: source_id.clone(),
297                now_ms,
298            })
299            .await
300            .map_err(|err| map_fork_source_ref_error(err, source_placement))?;
301
302        if request.content_type_explicit {
303            if request.content_type != source_head.content_type {
304                return Err(RuntimeError::group_engine(
305                    source_placement,
306                    GroupEngineError::stream(
307                        StreamErrorCode::ContentTypeMismatch,
308                        format!(
309                            "fork content type '{}' does not match source content type '{}'",
310                            request.content_type, source_head.content_type
311                        ),
312                    ),
313                ));
314            }
315        } else {
316            request.content_type.clone_from(&source_head.content_type);
317        }
318
319        let fork_offset = request.fork_offset.unwrap_or(source_head.tail_offset);
320        if fork_offset > source_head.tail_offset {
321            return Err(RuntimeError::group_engine(
322                source_placement,
323                GroupEngineError::stream(
324                    StreamErrorCode::InvalidFork,
325                    format!(
326                        "fork offset {fork_offset} is beyond source stream '{}' tail {}",
327                        source_id, source_head.tail_offset
328                    ),
329                ),
330            ));
331        }
332
333        let max_len = usize::try_from(fork_offset).map_err(|_| {
334            RuntimeError::group_engine(
335                source_placement,
336                GroupEngineError::stream(
337                    StreamErrorCode::InvalidFork,
338                    format!("fork offset {fork_offset} cannot fit in memory on this host"),
339                ),
340            )
341        })?;
342        request.initial_payload = if fork_offset == 0 {
343            Bytes::new()
344        } else {
345            self.read_stream(ReadStreamRequest {
346                stream_id: source_id.clone(),
347                offset: 0,
348                max_len,
349                now_ms,
350            })
351            .await?
352            .payload
353            .into()
354        };
355        self.add_fork_ref_on_owner(source_id.clone(), now_ms)
356            .await
357            .map_err(|err| map_fork_source_ref_error(err, source_placement))?;
358        request.close_after = false;
359        request.stream_seq = None;
360        request.producer = None;
361        if request.stream_ttl_seconds.is_none() && request.stream_expires_at_ms.is_none() {
362            request.stream_ttl_seconds = source_head.stream_ttl_seconds;
363            request.stream_expires_at_ms = source_head.stream_expires_at_ms;
364        }
365        request.fork_offset = Some(fork_offset);
366        match self.create_stream_on_owner(request).await {
367            Ok(response) if response.already_exists => {
368                self.release_fork_ref_cascade(source_id).await?;
369                Ok(response)
370            }
371            Ok(response) => Ok(response),
372            Err(err) => {
373                let _ = self.release_fork_ref_cascade(source_id).await;
374                Err(err)
375            }
376        }
377    }
378
379    pub async fn head_stream(
380        &self,
381        request: HeadStreamRequest,
382    ) -> Result<HeadStreamResponse, RuntimeError> {
383        let placement = self.shard_map.locate(&request.stream_id);
384        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
385        let (response_tx, response_rx) = oneshot::channel();
386        self.send_core_command(
387            mailbox,
388            CoreCommand::HeadStream {
389                request,
390                placement,
391                response_tx,
392            },
393            response_rx,
394        )
395        .await
396    }
397
398    pub async fn get_stream_attrs(
399        &self,
400        request: GetStreamAttrsRequest,
401    ) -> Result<GetStreamAttrsResponse, RuntimeError> {
402        let placement = self.shard_map.locate(&request.stream_id);
403        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
404        let (response_tx, response_rx) = oneshot::channel();
405        self.send_core_command(
406            mailbox,
407            CoreCommand::GetStreamAttrs {
408                request,
409                placement,
410                response_tx,
411            },
412            response_rx,
413        )
414        .await
415    }
416
417    pub async fn update_stream_attrs(
418        &self,
419        request: UpdateStreamAttrsRequest,
420    ) -> Result<UpdateStreamAttrsResponse, RuntimeError> {
421        let placement = self.shard_map.locate(&request.stream_id);
422        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
423        let (response_tx, response_rx) = oneshot::channel();
424        self.send_core_command(
425            mailbox,
426            CoreCommand::UpdateStreamAttrs {
427                request,
428                placement,
429                response_tx,
430            },
431            response_rx,
432        )
433        .await
434    }
435
436    pub async fn read_stream(
437        &self,
438        request: ReadStreamRequest,
439    ) -> Result<ReadStreamResponse, RuntimeError> {
440        let placement = self.shard_map.locate(&request.stream_id);
441        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
442        let (response_tx, response_rx) = oneshot::channel();
443        self.send_core_command(
444            mailbox,
445            CoreCommand::ReadStream {
446                request,
447                placement,
448                response_tx,
449            },
450            response_rx,
451        )
452        .await
453    }
454
455    pub async fn publish_snapshot(
456        &self,
457        request: PublishSnapshotRequest,
458    ) -> Result<PublishSnapshotResponse, RuntimeError> {
459        let placement = self.shard_map.locate(&request.stream_id);
460        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
461        let (response_tx, response_rx) = oneshot::channel();
462        self.send_core_command(
463            mailbox,
464            CoreCommand::PublishSnapshot {
465                request,
466                placement,
467                response_tx,
468            },
469            response_rx,
470        )
471        .await
472    }
473
474    pub async fn read_snapshot(
475        &self,
476        request: ReadSnapshotRequest,
477    ) -> Result<ReadSnapshotResponse, RuntimeError> {
478        let placement = self.shard_map.locate(&request.stream_id);
479        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
480        let (response_tx, response_rx) = oneshot::channel();
481        self.send_core_command(
482            mailbox,
483            CoreCommand::ReadSnapshot {
484                request,
485                placement,
486                response_tx,
487            },
488            response_rx,
489        )
490        .await
491    }
492
493    pub async fn delete_snapshot(
494        &self,
495        request: DeleteSnapshotRequest,
496    ) -> Result<(), RuntimeError> {
497        let placement = self.shard_map.locate(&request.stream_id);
498        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
499        let (response_tx, response_rx) = oneshot::channel();
500        self.send_core_command(
501            mailbox,
502            CoreCommand::DeleteSnapshot {
503                request,
504                placement,
505                response_tx,
506            },
507            response_rx,
508        )
509        .await
510    }
511
512    pub async fn bootstrap_stream(
513        &self,
514        request: BootstrapStreamRequest,
515    ) -> Result<BootstrapStreamResponse, RuntimeError> {
516        let placement = self.shard_map.locate(&request.stream_id);
517        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
518        let (response_tx, response_rx) = oneshot::channel();
519        self.send_core_command(
520            mailbox,
521            CoreCommand::BootstrapStream {
522                request,
523                placement,
524                response_tx,
525            },
526            response_rx,
527        )
528        .await
529    }
530
531    pub async fn wait_read_stream(
532        &self,
533        request: ReadStreamRequest,
534    ) -> Result<ReadStreamResponse, RuntimeError> {
535        let placement = self.shard_map.locate(&request.stream_id);
536        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
537        let waiter_id = self.next_waiter_id.fetch_add(1, Ordering::Relaxed);
538        let stream_id = request.stream_id.clone();
539        let (response_tx, response_rx) = oneshot::channel();
540        self.enqueue_core_command(mailbox, CoreCommand::WaitRead {
541            request,
542            placement,
543            waiter_id,
544            response_tx,
545        })
546        .await?;
547        let mut cancel = WaitReadCancel::new(mailbox.tx.clone(), stream_id, placement, waiter_id);
548        let response = response_rx
549            .await
550            .map_err(|_| RuntimeError::ResponseDropped {
551                core_id: mailbox.core_id,
552            })?;
553        cancel.disarm();
554        response
555    }
556
557    pub async fn require_local_live_read_owner(
558        &self,
559        stream_id: &BucketStreamId,
560    ) -> Result<(), RuntimeError> {
561        let placement = self.shard_map.locate(stream_id);
562        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
563        let (response_tx, response_rx) = oneshot::channel();
564        self.send_core_command(
565            mailbox,
566            CoreCommand::RequireLiveReadOwner {
567                placement,
568                response_tx,
569            },
570            response_rx,
571        )
572        .await
573    }
574
575    pub async fn close_stream(
576        &self,
577        request: CloseStreamRequest,
578    ) -> Result<CloseStreamResponse, RuntimeError> {
579        let placement = self.shard_map.locate(&request.stream_id);
580        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
581        let (response_tx, response_rx) = oneshot::channel();
582        self.send_core_command(
583            mailbox,
584            CoreCommand::CloseStream {
585                request,
586                placement,
587                response_tx,
588            },
589            response_rx,
590        )
591        .await
592    }
593
594    pub async fn delete_stream(
595        &self,
596        request: DeleteStreamRequest,
597    ) -> Result<DeleteStreamResponse, RuntimeError> {
598        let response = self.delete_stream_on_owner(request).await?;
599        if let Some(parent_to_release) = response.parent_to_release.clone() {
600            self.release_fork_ref_cascade(parent_to_release).await?;
601        }
602        Ok(response)
603    }
604
605    async fn delete_stream_on_owner(
606        &self,
607        request: DeleteStreamRequest,
608    ) -> Result<DeleteStreamResponse, RuntimeError> {
609        let placement = self.shard_map.locate(&request.stream_id);
610        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
611        let (response_tx, response_rx) = oneshot::channel();
612        self.send_core_command(
613            mailbox,
614            CoreCommand::DeleteStream {
615                request,
616                placement,
617                response_tx,
618            },
619            response_rx,
620        )
621        .await
622    }
623
624    async fn add_fork_ref_on_owner(
625        &self,
626        stream_id: BucketStreamId,
627        now_ms: u64,
628    ) -> Result<ForkRefResponse, RuntimeError> {
629        let placement = self.shard_map.locate(&stream_id);
630        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
631        let (response_tx, response_rx) = oneshot::channel();
632        self.send_core_command(
633            mailbox,
634            CoreCommand::AddForkRef {
635                stream_id,
636                now_ms,
637                placement,
638                response_tx,
639            },
640            response_rx,
641        )
642        .await
643    }
644
645    async fn release_fork_ref_on_owner(
646        &self,
647        stream_id: BucketStreamId,
648    ) -> Result<ForkRefResponse, RuntimeError> {
649        let placement = self.shard_map.locate(&stream_id);
650        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
651        let (response_tx, response_rx) = oneshot::channel();
652        self.send_core_command(
653            mailbox,
654            CoreCommand::ReleaseForkRef {
655                stream_id,
656                placement,
657                response_tx,
658            },
659            response_rx,
660        )
661        .await
662    }
663
664    async fn release_fork_ref_cascade(
665        &self,
666        stream_id: BucketStreamId,
667    ) -> Result<(), RuntimeError> {
668        let mut next = Some(stream_id);
669        while let Some(current) = next {
670            let response = self.release_fork_ref_on_owner(current).await?;
671            next = response.parent_to_release;
672        }
673        Ok(())
674    }
675
676    pub async fn flush_cold(
677        &self,
678        request: FlushColdRequest,
679    ) -> Result<FlushColdResponse, RuntimeError> {
680        let placement = self.shard_map.locate(&request.stream_id);
681        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
682        let (response_tx, response_rx) = oneshot::channel();
683        self.send_core_command(
684            mailbox,
685            CoreCommand::FlushCold {
686                request,
687                placement,
688                response_tx,
689            },
690            response_rx,
691        )
692        .await
693    }
694
695    pub async fn append_external(
696        &self,
697        request: AppendExternalRequest,
698    ) -> Result<AppendResponse, RuntimeError> {
699        let placement = self.shard_map.locate(&request.stream_id);
700        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
701        let (response_tx, response_rx) = oneshot::channel();
702        self.send_core_command(
703            mailbox,
704            CoreCommand::AppendExternal {
705                request,
706                placement,
707                response_tx,
708            },
709            response_rx,
710        )
711        .await
712    }
713
714    pub async fn plan_cold_flush(
715        &self,
716        request: PlanColdFlushRequest,
717    ) -> Result<Option<ColdFlushCandidate>, RuntimeError> {
718        let placement = self.shard_map.locate(&request.stream_id);
719        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
720        let (response_tx, response_rx) = oneshot::channel();
721        self.send_core_command(
722            mailbox,
723            CoreCommand::PlanColdFlush {
724                request,
725                placement,
726                response_tx,
727            },
728            response_rx,
729        )
730        .await
731    }
732
733    pub async fn flush_cold_once(
734        &self,
735        request: PlanColdFlushRequest,
736    ) -> Result<Option<FlushColdResponse>, RuntimeError> {
737        let Some(candidate) = self.plan_cold_flush(request).await? else {
738            return Ok(None);
739        };
740        self.flush_cold_candidate(candidate).await.map(Some)
741    }
742
743    pub async fn plan_next_cold_flush(
744        &self,
745        raft_group_id: RaftGroupId,
746        request: PlanGroupColdFlushRequest,
747    ) -> Result<Option<ColdFlushCandidate>, RuntimeError> {
748        let placement = self.placement_for_group(raft_group_id)?;
749        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
750        let (response_tx, response_rx) = oneshot::channel();
751        self.send_core_command(
752            mailbox,
753            CoreCommand::PlanNextColdFlush {
754                request,
755                placement,
756                response_tx,
757            },
758            response_rx,
759        )
760        .await
761    }
762
763    pub async fn plan_next_cold_flush_batch(
764        &self,
765        raft_group_id: RaftGroupId,
766        request: PlanGroupColdFlushRequest,
767        max_candidates: usize,
768    ) -> Result<Vec<ColdFlushCandidate>, RuntimeError> {
769        let placement = self.placement_for_group(raft_group_id)?;
770        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
771        let (response_tx, response_rx) = oneshot::channel();
772        self.send_core_command(
773            mailbox,
774            CoreCommand::PlanNextColdFlushBatch {
775                request,
776                placement,
777                max_candidates,
778                response_tx,
779            },
780            response_rx,
781        )
782        .await
783    }
784
785    pub async fn flush_cold_group_once(
786        &self,
787        raft_group_id: RaftGroupId,
788        request: PlanGroupColdFlushRequest,
789    ) -> Result<Option<FlushColdResponse>, RuntimeError> {
790        let Some(candidate) = self.plan_next_cold_flush(raft_group_id, request).await? else {
791            return Ok(None);
792        };
793        match self.flush_cold_candidate(candidate).await {
794            Ok(response) => Ok(Some(response)),
795            Err(err) if is_stale_cold_flush_candidate_error(&err) => Ok(None),
796            Err(err) => Err(err),
797        }
798    }
799
800    pub async fn flush_cold_group_batch_once(
801        &self,
802        raft_group_id: RaftGroupId,
803        request: PlanGroupColdFlushRequest,
804        max_candidates: usize,
805    ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
806        let candidates = self
807            .plan_next_cold_flush_batch(raft_group_id, request, max_candidates)
808            .await?;
809        if candidates.is_empty() {
810            return Ok(Vec::new());
811        }
812        self.flush_cold_candidates_batch(candidates).await
813    }
814
815    async fn flush_cold_candidate(
816        &self,
817        candidate: ColdFlushCandidate,
818    ) -> Result<FlushColdResponse, RuntimeError> {
819        let Some(cold_store) = self.cold_store.as_ref() else {
820            return Err(RuntimeError::ColdStoreConfig {
821                message: "cold backend must be configured before flushing cold chunks".to_owned(),
822            });
823        };
824        let path = new_cold_chunk_path(
825            &candidate.stream_id,
826            candidate.start_offset,
827            candidate.end_offset,
828        );
829        let upload_started_at = Instant::now();
830        let object_size = match cold_store.write_chunk(&path, &candidate.payload).await {
831            Ok(object_size) => object_size,
832            Err(err) => {
833                // Surfaces "this node can't write to S3" to the snapshot driver's
834                // health/yield logic, which drives leadership-yield off real cold
835                // flush failures rather than a stat probe a keep-alive connection
836                // can mask.
837                self.metrics.record_cold_flush_write_error();
838                return Err(RuntimeError::ColdStoreIo {
839                    message: err.to_string(),
840                });
841            }
842        };
843        self.metrics
844            .record_cold_upload(object_size, elapsed_ns(upload_started_at));
845        let chunk = ColdChunkRef {
846            start_offset: candidate.start_offset,
847            end_offset: candidate.end_offset,
848            s3_path: path.clone(),
849            object_size,
850        };
851        let publish_started_at = Instant::now();
852        let publish = self
853            .flush_cold(FlushColdRequest {
854                stream_id: candidate.stream_id,
855                chunk,
856            })
857            .await;
858        match publish {
859            Ok(response) => {
860                self.metrics
861                    .record_cold_publish(object_size, elapsed_ns(publish_started_at));
862                Ok(response)
863            }
864            Err(err) => Err(err),
865        }
866    }
867
868    pub(crate) async fn flush_cold_candidates_batch(
869        &self,
870        candidates: Vec<ColdFlushCandidate>,
871    ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
872        let mut responses = Vec::with_capacity(candidates.len());
873        for candidate in candidates {
874            match self.flush_cold_candidate(candidate).await {
875                Ok(response) => responses.push(response),
876                Err(err) if is_stale_cold_flush_candidate_error(&err) => {}
877                Err(err) => return Err(err),
878            }
879        }
880        Ok(responses)
881    }
882
883    #[cfg(madsim)]
884    pub async fn flush_cold_candidates_batch_for_simulation(
885        &self,
886        candidates: Vec<ColdFlushCandidate>,
887    ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
888        self.flush_cold_candidates_batch(candidates).await
889    }
890
891    pub async fn flush_cold_all_groups_once(
892        &self,
893        request: PlanGroupColdFlushRequest,
894    ) -> Result<usize, RuntimeError> {
895        self.flush_cold_all_groups_once_bounded(request, 1).await
896    }
897
898    pub async fn flush_cold_all_groups_once_bounded(
899        &self,
900        request: PlanGroupColdFlushRequest,
901        max_concurrency: usize,
902    ) -> Result<usize, RuntimeError> {
903        let max_concurrency = max_concurrency.max(1);
904        if max_concurrency == 1 {
905            return self.flush_cold_all_groups_once_serial(request).await;
906        }
907        #[cfg(madsim)]
908        {
909            return self.flush_cold_all_groups_once_serial(request).await;
910        }
911        #[cfg(not(madsim))]
912        {
913            let mut flushed = 0;
914            let mut next_group_id = 0;
915            let group_count = self.shard_map.raft_group_count();
916            let mut tasks = JoinSet::new();
917
918            while next_group_id < group_count || !tasks.is_empty() {
919                while next_group_id < group_count && tasks.len() < max_concurrency {
920                    let runtime = self.clone();
921                    let request = request.clone();
922                    let group_id = RaftGroupId(next_group_id);
923                    next_group_id += 1;
924                    tasks.spawn(async move {
925                        runtime
926                            .flush_cold_group_batch_once(
927                                group_id,
928                                request,
929                                COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
930                            )
931                            .await
932                            .map(|responses| responses.len())
933                    });
934                }
935                if let Some(result) = tasks.join_next().await {
936                    match result {
937                        Ok(Ok(count)) => flushed += count,
938                        Ok(Err(err)) => return Err(err),
939                        Err(err) => {
940                            return Err(RuntimeError::ColdStoreIo {
941                                message: format!("cold flush task failed: {err}"),
942                            });
943                        }
944                    }
945                }
946            }
947            Ok(flushed)
948        }
949    }
950
951    async fn flush_cold_all_groups_once_serial(
952        &self,
953        request: PlanGroupColdFlushRequest,
954    ) -> Result<usize, RuntimeError> {
955        let mut flushed = 0;
956        for group_id in 0..self.shard_map.raft_group_count() {
957            flushed += self
958                .flush_cold_group_batch_once(
959                    RaftGroupId(group_id),
960                    request.clone(),
961                    COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
962                )
963                .await?
964                .len();
965        }
966        Ok(flushed)
967    }
968
969    async fn plan_cold_gc(
970        &self,
971        raft_group_id: RaftGroupId,
972        max: usize,
973    ) -> Result<Vec<ColdGcEntry>, RuntimeError> {
974        let placement = self.placement_for_group(raft_group_id)?;
975        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
976        let (response_tx, response_rx) = oneshot::channel();
977        self.send_core_command(
978            mailbox,
979            CoreCommand::PlanColdGc {
980                max,
981                placement,
982                response_tx,
983            },
984            response_rx,
985        )
986        .await
987    }
988
989    async fn ack_cold_gc(
990        &self,
991        raft_group_id: RaftGroupId,
992        up_to_seq: u64,
993    ) -> Result<AckColdGcResponse, RuntimeError> {
994        let placement = self.placement_for_group(raft_group_id)?;
995        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
996        let (response_tx, response_rx) = oneshot::channel();
997        self.send_core_command(
998            mailbox,
999            CoreCommand::AckColdGc {
1000                up_to_seq,
1001                placement,
1002                response_tx,
1003            },
1004            response_rx,
1005        )
1006        .await
1007    }
1008
1009    /// Drains the leader-side cold-GC queue for one group: physically reclaims
1010    /// each queued target from cold storage, then replicates an ack that pops
1011    /// the reclaimed entries. Deletions are idempotent, so a crash or leader
1012    /// change between reclaim and ack simply re-runs them next tick.
1013    pub async fn run_cold_gc_group_once(
1014        &self,
1015        raft_group_id: RaftGroupId,
1016        max_entries: usize,
1017    ) -> Result<usize, RuntimeError> {
1018        let Some(cold_store) = self.cold_store.as_ref() else {
1019            return Ok(0);
1020        };
1021        let entries = self.plan_cold_gc(raft_group_id, max_entries).await?;
1022        if entries.is_empty() {
1023            return Ok(0);
1024        }
1025        let mut acked_seq = None;
1026        let mut reclaimed = 0usize;
1027        // Entries are FIFO by seq; stop at the first failure so the ack never
1028        // skips past an object that is still present in cold storage.
1029        for entry in entries {
1030            let result = match &entry.target {
1031                ColdGcTarget::Stream(stream_id) => {
1032                    match cold_store.remove_all(&cold_chunk_prefix(stream_id)).await {
1033                        Ok(()) => cold_store.remove_all(&cold_index_prefix(stream_id)).await,
1034                        Err(err) => Err(err),
1035                    }
1036                }
1037                ColdGcTarget::Paths(paths) => {
1038                    let mut outcome = Ok(());
1039                    for path in paths {
1040                        if let Err(err) = cold_store.delete_chunk(path).await {
1041                            outcome = Err(err);
1042                            break;
1043                        }
1044                    }
1045                    outcome
1046                }
1047            };
1048            match result {
1049                Ok(()) => {
1050                    acked_seq = Some(entry.seq);
1051                    reclaimed += 1;
1052                }
1053                Err(err) => {
1054                    self.metrics.record_cold_gc_error();
1055                    if acked_seq.is_none() {
1056                        return Err(RuntimeError::ColdStoreIo {
1057                            message: err.to_string(),
1058                        });
1059                    }
1060                    break;
1061                }
1062            }
1063        }
1064        if let Some(up_to_seq) = acked_seq {
1065            self.ack_cold_gc(raft_group_id, up_to_seq).await?;
1066            self.metrics
1067                .record_cold_gc_reclaimed(u64::try_from(reclaimed).expect("reclaimed fits u64"));
1068        }
1069        Ok(reclaimed)
1070    }
1071
1072    pub async fn run_cold_gc_all_groups_once(
1073        &self,
1074        max_entries_per_group: usize,
1075    ) -> Result<usize, RuntimeError> {
1076        if self.cold_store.is_none() {
1077            return Ok(0);
1078        }
1079        let mut reclaimed = 0;
1080        for group_id in 0..self.shard_map.raft_group_count() {
1081            reclaimed += self
1082                .run_cold_gc_group_once(RaftGroupId(group_id), max_entries_per_group)
1083                .await?;
1084        }
1085        Ok(reclaimed)
1086    }
1087
1088    pub async fn append(&self, request: AppendRequest) -> Result<AppendResponse, RuntimeError> {
1089        if request.payload.is_empty() {
1090            return Err(RuntimeError::EmptyAppend);
1091        }
1092        let placement = self.shard_map.locate(&request.stream_id);
1093        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1094        let (response_tx, response_rx) = oneshot::channel();
1095        self.send_core_command(
1096            mailbox,
1097            CoreCommand::Append {
1098                request,
1099                placement,
1100                response_tx,
1101            },
1102            response_rx,
1103        )
1104        .await
1105    }
1106
1107    pub async fn append_batch(
1108        &self,
1109        request: AppendBatchRequest,
1110    ) -> Result<AppendBatchResponse, RuntimeError> {
1111        if request.payloads.is_empty() {
1112            return Err(RuntimeError::EmptyAppend);
1113        }
1114        let placement = self.shard_map.locate(&request.stream_id);
1115        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1116        let (response_tx, response_rx) = oneshot::channel();
1117        self.send_core_command(
1118            mailbox,
1119            CoreCommand::AppendBatch {
1120                request,
1121                placement,
1122                response_tx,
1123            },
1124            response_rx,
1125        )
1126        .await
1127    }
1128
1129    pub async fn snapshot_group(
1130        &self,
1131        raft_group_id: RaftGroupId,
1132    ) -> Result<GroupSnapshot, RuntimeError> {
1133        let placement = self.placement_for_group(raft_group_id)?;
1134        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1135        let (response_tx, response_rx) = oneshot::channel();
1136        self.send_core_command(
1137            mailbox,
1138            CoreCommand::SnapshotGroup {
1139                placement,
1140                response_tx,
1141            },
1142            response_rx,
1143        )
1144        .await
1145    }
1146
1147    pub async fn install_group_snapshot(
1148        &self,
1149        snapshot: GroupSnapshot,
1150    ) -> Result<(), RuntimeError> {
1151        let expected = self.placement_for_group(snapshot.placement.raft_group_id)?;
1152        if snapshot.placement != expected {
1153            return Err(RuntimeError::SnapshotPlacementMismatch {
1154                expected,
1155                actual: snapshot.placement,
1156            });
1157        }
1158        let mailbox = &self.mailboxes[usize::from(expected.core_id.0)];
1159        let (response_tx, response_rx) = oneshot::channel();
1160        self.send_core_command(
1161            mailbox,
1162            CoreCommand::InstallGroupSnapshot {
1163                snapshot,
1164                response_tx,
1165            },
1166            response_rx,
1167        )
1168        .await
1169    }
1170
1171    #[cfg(madsim)]
1172    pub async fn shutdown_group_engine_for_simulation(
1173        &self,
1174        placement: ShardPlacement,
1175    ) -> Result<(), RuntimeError> {
1176        let expected = self.placement_for_group(placement.raft_group_id)?;
1177        if placement != expected {
1178            return Err(RuntimeError::SnapshotPlacementMismatch {
1179                expected,
1180                actual: placement,
1181            });
1182        }
1183        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1184        let (response_tx, response_rx) = oneshot::channel();
1185        self.send_core_command(
1186            mailbox,
1187            CoreCommand::ShutdownGroupEngine {
1188                placement,
1189                response_tx,
1190            },
1191            response_rx,
1192        )
1193        .await
1194    }
1195
1196    #[cfg(madsim)]
1197    pub async fn install_group_engine_for_simulation(
1198        &self,
1199        placement: ShardPlacement,
1200        engine: Box<dyn crate::engine::GroupEngine>,
1201    ) -> Result<(), RuntimeError> {
1202        let expected = self.placement_for_group(placement.raft_group_id)?;
1203        if placement != expected {
1204            return Err(RuntimeError::SnapshotPlacementMismatch {
1205                expected,
1206                actual: placement,
1207            });
1208        }
1209        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1210        let (response_tx, response_rx) = oneshot::channel();
1211        self.send_core_command(
1212            mailbox,
1213            CoreCommand::InstallGroupEngine {
1214                placement,
1215                engine,
1216                response_tx,
1217            },
1218            response_rx,
1219        )
1220        .await
1221    }
1222
1223    pub async fn warm_group(
1224        &self,
1225        raft_group_id: RaftGroupId,
1226    ) -> Result<ShardPlacement, RuntimeError> {
1227        let placement = self.placement_for_group(raft_group_id)?;
1228        let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1229        let (response_tx, response_rx) = oneshot::channel();
1230        self.send_core_command(
1231            mailbox,
1232            CoreCommand::WarmGroup {
1233                placement,
1234                response_tx,
1235            },
1236            response_rx,
1237        )
1238        .await
1239    }
1240
1241    pub async fn warm_all_groups(&self) -> Result<(), RuntimeError> {
1242        for raw_group_id in 0..self.shard_map.raft_group_count() {
1243            self.warm_group(RaftGroupId(raw_group_id)).await?;
1244        }
1245        Ok(())
1246    }
1247
1248    fn placement_for_group(
1249        &self,
1250        raft_group_id: RaftGroupId,
1251    ) -> Result<ShardPlacement, RuntimeError> {
1252        if raft_group_id.0 >= self.shard_map.raft_group_count() {
1253            return Err(RuntimeError::InvalidRaftGroup {
1254                raft_group_id,
1255                raft_group_count: self.shard_map.raft_group_count(),
1256            });
1257        }
1258        Ok(ShardPlacement {
1259            core_id: CoreId(
1260                (raft_group_id.0 % u32::from(self.shard_map.core_count()))
1261                    .try_into()
1262                    .expect("core id fits u16"),
1263            ),
1264            shard_id: ShardId(raft_group_id.0),
1265            raft_group_id,
1266        })
1267    }
1268
1269    async fn send_core_command<T>(
1270        &self,
1271        mailbox: &CoreMailbox,
1272        command: CoreCommand,
1273        response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
1274    ) -> Result<T, RuntimeError> {
1275        self.enqueue_core_command(mailbox, command).await?;
1276        response_rx
1277            .await
1278            .map_err(|_| RuntimeError::ResponseDropped {
1279                core_id: mailbox.core_id,
1280            })?
1281    }
1282
1283    async fn enqueue_core_command(
1284        &self,
1285        mailbox: &CoreMailbox,
1286        command: CoreCommand,
1287    ) -> Result<(), RuntimeError> {
1288        if mailbox.tx.capacity() == 0 {
1289            self.metrics.record_mailbox_full(mailbox.core_id);
1290        }
1291        let started_at = Instant::now();
1292        mailbox
1293            .tx
1294            .send(Traced::capture(command))
1295            .await
1296            .map_err(|_| RuntimeError::MailboxClosed {
1297                core_id: mailbox.core_id,
1298            })?;
1299        self.metrics
1300            .record_routed_request(mailbox.core_id, elapsed_ns(started_at));
1301        Ok(())
1302    }
1303
1304    pub fn metrics(&self) -> RuntimeMetrics {
1305        RuntimeMetrics {
1306            inner: self.metrics.clone(),
1307        }
1308    }
1309
1310    pub fn mailbox_snapshot(&self) -> RuntimeMailboxSnapshot {
1311        let depths = self
1312            .mailboxes
1313            .iter()
1314            .map(CoreMailbox::depth)
1315            .collect::<Vec<_>>();
1316        let capacities = self
1317            .mailboxes
1318            .iter()
1319            .map(CoreMailbox::capacity)
1320            .collect::<Vec<_>>();
1321        RuntimeMailboxSnapshot { depths, capacities }
1322    }
1323}
1324
1325fn spawn_core_worker(threading: RuntimeThreading, worker: CoreWorker) -> Result<(), RuntimeError> {
1326    match threading {
1327        RuntimeThreading::HostedTokio => {
1328            crate::rt::spawn(worker.run());
1329            Ok(())
1330        }
1331        #[cfg(not(madsim))]
1332        RuntimeThreading::ThreadPerCore => {
1333            let core_id = worker.core_id;
1334            std::thread::Builder::new()
1335                .name(format!("ursula-core-{}", core_id.0))
1336                .spawn(move || {
1337                    let runtime = tokio::runtime::Builder::new_current_thread()
1338                        .enable_all()
1339                        .build()
1340                        .expect("build per-core tokio runtime");
1341                    runtime.block_on(worker.run());
1342                })
1343                .map(|_| ())
1344                .map_err(|err| RuntimeError::SpawnCoreThread {
1345                    core_id,
1346                    message: err.to_string(),
1347                })
1348        }
1349    }
1350}