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