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