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