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
100fn is_legacy_cross_bucket_pack(stream_id: &BucketStreamId, chunk: &ColdChunkRef) -> bool {
101 chunk.shared_object
102 && !chunk
103 .s3_path
104 .starts_with(&format!("{}/_packs/", stream_id.bucket_id))
105}
106
107#[derive(Debug, Clone)]
108pub struct RuntimeConfig {
109 pub core_count: usize,
110 pub raft_group_count: usize,
111 pub mailbox_capacity: usize,
112 pub threading: RuntimeThreading,
113 pub cold_max_hot_bytes_per_group: Option<u64>,
114 pub raft_max_uncommitted_bytes_per_group: Option<u64>,
118 pub live_read_max_waiters_per_core: Option<u64>,
119}
120
121impl RuntimeConfig {
122 pub fn new(core_count: usize, raft_group_count: usize) -> Self {
123 #[cfg(not(madsim))]
124 let threading = RuntimeThreading::ThreadPerCore;
125 #[cfg(madsim)]
126 let threading = RuntimeThreading::HostedTokio;
127 Self {
128 core_count,
129 raft_group_count,
130 mailbox_capacity: 1024,
131 threading,
132 cold_max_hot_bytes_per_group: None,
133 raft_max_uncommitted_bytes_per_group: None,
134 live_read_max_waiters_per_core: Some(65_536),
135 }
136 }
137
138 pub fn with_cold_max_hot_bytes_per_group(mut self, value: Option<u64>) -> Self {
139 self.cold_max_hot_bytes_per_group = value;
140 self
141 }
142
143 pub fn with_raft_max_uncommitted_bytes_per_group(mut self, value: Option<u64>) -> Self {
144 self.raft_max_uncommitted_bytes_per_group = value;
145 self
146 }
147
148 pub fn with_live_read_max_waiters_per_core(mut self, value: Option<u64>) -> Self {
149 self.live_read_max_waiters_per_core = value;
150 self
151 }
152
153 pub fn from_ursula_config(cfg: &ursula_config::RuntimeConfig, raft_group_count: usize) -> Self {
155 let mut config = Self::new(cfg.core_count, raft_group_count);
156 config.live_read_max_waiters_per_core = cfg
157 .live_read_max_waiters_per_core
158 .and_then(|n| if n == 0 { None } else { Some(n as u64) });
159 config
160 }
161}
162
163#[derive(Debug, Clone, Copy, PartialEq, Eq)]
164pub enum RuntimeThreading {
165 #[cfg(not(madsim))]
166 ThreadPerCore,
167 HostedTokio,
168}
169
170#[derive(Debug, Clone)]
171pub struct ShardRuntime {
172 shard_map: StaticShardMap,
173 mailboxes: Vec<CoreMailbox>,
174 metrics: Arc<RuntimeMetricsInner>,
175 next_waiter_id: Arc<AtomicU64>,
176 cold_store: Option<ColdStoreHandle>,
177}
178
179#[derive(Debug, Clone, Default, PartialEq, Eq)]
181pub struct PurgeBucketReport {
182 pub removed_streams: u64,
183 pub groups_with_streams: Vec<usize>,
184 pub pending_cold_gc_entries: u64,
185}
186
187#[derive(Debug, Clone, Default, PartialEq, Eq)]
188pub struct LegacySharedMigrationReport {
189 pub observed_chunks: usize,
190 pub migrated_chunks: usize,
191 pub pending_chunks: usize,
192}
193
194impl ShardRuntime {
195 pub fn spawn(config: RuntimeConfig) -> Result<Self, RuntimeError> {
196 Self::spawn_with_engine_factory(config, InMemoryGroupEngineFactory::default())
197 }
198
199 pub fn spawn_with_engine_factory(
200 config: RuntimeConfig,
201 engine_factory: impl GroupEngineFactory,
202 ) -> Result<Self, RuntimeError> {
203 Self::spawn_with_engine_factory_and_cold_store(config, engine_factory, None)
204 }
205
206 pub fn spawn_with_engine_factory_and_cold_store(
207 config: RuntimeConfig,
208 engine_factory: impl GroupEngineFactory,
209 cold_store: Option<ColdStoreHandle>,
210 ) -> Result<Self, RuntimeError> {
211 let shard_map = StaticShardMap::new(config.core_count, config.raft_group_count)?;
212 let metrics = Arc::new(RuntimeMetricsInner::new(
213 usize::from(shard_map.core_count()),
214 usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
215 ));
216 let cold_write_admission = ColdWriteAdmission {
217 max_hot_bytes_per_group: config.cold_max_hot_bytes_per_group,
218 };
219 let raft_uncommitted_admission = RaftUncommittedAdmission {
220 max_uncommitted_bytes_per_group: config.raft_max_uncommitted_bytes_per_group,
221 };
222 let raft_uncommitted_bytes = Arc::new(RaftUncommittedBytesTracker::new(
223 usize::try_from(shard_map.raft_group_count()).expect("u32 fits usize"),
224 ));
225 let engine_factory: Arc<dyn GroupEngineFactory> = Arc::new(engine_factory);
226 let read_materialization = Arc::new(Semaphore::new(config.mailbox_capacity.max(1)));
227 let mut mailboxes = Vec::with_capacity(usize::from(shard_map.core_count()));
228 for raw_core_id in 0..shard_map.core_count() {
229 let core_id = CoreId(raw_core_id);
230 let (tx, rx) = mpsc::channel(config.mailbox_capacity.max(1));
231 let worker = CoreWorker {
232 core_id,
233 rx,
234 engine_factory: engine_factory.clone(),
235 groups: HashMap::new(),
236 metrics: metrics.clone(),
237 group_mailbox_capacity: config.mailbox_capacity.max(1),
238 cold_write_admission,
239 raft_uncommitted_admission,
240 raft_uncommitted_bytes: raft_uncommitted_bytes.clone(),
241 live_read_max_waiters_per_core: config.live_read_max_waiters_per_core,
242 read_materialization: read_materialization.clone(),
243 };
244 spawn_core_worker(config.threading, worker)?;
245 mailboxes.push(CoreMailbox { core_id, tx });
246 }
247 Ok(Self {
248 shard_map,
249 mailboxes,
250 metrics,
251 next_waiter_id: Arc::new(AtomicU64::new(1)),
252 cold_store,
253 })
254 }
255
256 pub fn locate(&self, stream_id: &BucketStreamId) -> ShardPlacement {
257 self.shard_map.locate(stream_id)
258 }
259
260 pub fn has_cold_store(&self) -> bool {
261 self.cold_store.is_some()
262 }
263
264 pub fn cold_store(&self) -> Option<ColdStoreHandle> {
265 self.cold_store.clone()
266 }
267
268 pub fn cold_store_info(&self) -> Option<ColdStoreInfo> {
269 self.cold_store
270 .as_ref()
271 .map(|cold_store| cold_store.info().clone())
272 }
273
274 pub async fn wait_read_stream(
275 &self,
276 request: ReadStreamRequest,
277 ) -> Result<ReadStreamResponse, RuntimeError> {
278 let placement = self.shard_map.locate(&request.stream_id);
279 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
280 let waiter_id = self.next_waiter_id.fetch_add(1, Ordering::Relaxed);
281 let stream_id = request.stream_id.clone();
282 let (response_tx, response_rx) = oneshot::channel();
283 self.enqueue_core_command(mailbox, CoreCommand::Group {
284 placement,
285 admission: None,
286 command: GroupCommand::WaitRead {
287 request,
288 waiter_id,
289 response_tx,
290 },
291 })
292 .await?;
293 let mut cancel = WaitReadCancel::new(mailbox.tx.clone(), stream_id, placement, waiter_id);
294 let response = response_rx
295 .await
296 .map_err(|_| RuntimeError::ResponseDropped {
297 core_id: mailbox.core_id,
298 })?;
299 cancel.disarm();
300 response
301 }
302
303 pub async fn require_local_live_read_owner(
304 &self,
305 stream_id: &BucketStreamId,
306 ) -> Result<(), RuntimeError> {
307 let placement = self.shard_map.locate(stream_id);
308 let (response_tx, response_rx) = oneshot::channel();
309 self.group_rpc(
310 placement,
311 None,
312 GroupCommand::RequireLiveReadOwner { response_tx },
313 response_rx,
314 )
315 .await
316 }
317
318 pub async fn flush_cold_once(
319 &self,
320 request: PlanColdFlushRequest,
321 ) -> Result<Option<FlushColdResponse>, RuntimeError> {
322 let Some(candidate) = self.plan_cold_flush(request).await? else {
323 return Ok(None);
324 };
325 self.flush_cold_candidate(candidate).await.map(Some)
326 }
327
328 pub async fn flush_cold_group_once(
329 &self,
330 raft_group_id: RaftGroupId,
331 request: PlanGroupColdFlushRequest,
332 ) -> Result<Option<FlushColdResponse>, RuntimeError> {
333 let mut candidates = self
334 .plan_next_cold_flush_batch(raft_group_id, request, 1)
335 .await?;
336 let Some(candidate) = candidates.pop() else {
337 return Ok(None);
338 };
339 match self.flush_cold_candidate(candidate).await {
340 Ok(response) => Ok(Some(response)),
341 Err(err) if is_stale_cold_flush_candidate_error(&err) => Ok(None),
342 Err(err) => Err(err),
343 }
344 }
345
346 pub async fn flush_cold_group_batch_once(
347 &self,
348 raft_group_id: RaftGroupId,
349 request: PlanGroupColdFlushRequest,
350 max_candidates: usize,
351 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
352 let candidates = self
353 .plan_next_cold_flush_batch(raft_group_id, request, max_candidates)
354 .await?;
355 if candidates.is_empty() {
356 return Ok(Vec::new());
357 }
358 self.flush_cold_candidates_batch(candidates).await
359 }
360
361 async fn flush_cold_candidate(
362 &self,
363 candidate: ColdFlushCandidate,
364 ) -> Result<FlushColdResponse, RuntimeError> {
365 let Some(cold_store) = self.cold_store.as_ref() else {
366 return Err(RuntimeError::ColdStoreConfig {
367 message: "cold backend must be configured before flushing cold chunks".to_owned(),
368 });
369 };
370 let path = new_cold_chunk_path(
371 &candidate.stream_id,
372 candidate.start_offset,
373 candidate.end_offset,
374 );
375 let upload_started_at = Instant::now();
376 let object_size = match cold_store.write_chunk(&path, &candidate.payload).await {
377 Ok(object_size) => object_size,
378 Err(err) => {
379 self.metrics.record_cold_flush_write_error();
384 return Err(RuntimeError::ColdStoreIo {
385 message: err.to_string(),
386 });
387 }
388 };
389 self.metrics
390 .record_cold_upload(object_size, elapsed_ns(upload_started_at));
391 let chunk = ColdChunkRef {
392 start_offset: candidate.start_offset,
393 end_offset: candidate.end_offset,
394 s3_path: path.clone(),
395 object_size,
396 object_offset: 0,
397 shared_object: false,
398 payload_digest: candidate.payload_digest,
399 };
400 let publish_started_at = Instant::now();
401 let publish = self
402 .flush_cold(FlushColdRequest {
403 stream_id: candidate.stream_id,
404 chunk,
405 })
406 .await;
407 match publish {
408 Ok(response) => {
409 self.metrics
410 .record_cold_publish(object_size, elapsed_ns(publish_started_at));
411 Ok(response)
412 }
413 Err(err) => Err(err),
414 }
415 }
416
417 pub(crate) async fn flush_cold_candidates_batch(
418 &self,
419 candidates: Vec<ColdFlushCandidate>,
420 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
421 let mut bucket_batches: Vec<Vec<ColdFlushCandidate>> = Vec::new();
425 for candidate in candidates {
426 if let Some(batch) = bucket_batches.iter_mut().find(|batch| {
427 batch
428 .first()
429 .is_some_and(|first| first.stream_id.bucket_id == candidate.stream_id.bucket_id)
430 }) {
431 batch.push(candidate);
432 } else {
433 bucket_batches.push(vec![candidate]);
434 }
435 }
436
437 let mut responses = Vec::new();
438 for mut batch in bucket_batches {
439 if batch.len() > 1 {
440 responses.extend(self.flush_cold_candidates_pack(batch).await?);
441 continue;
442 }
443 let candidate = batch.pop().expect("single-candidate batch");
444 match self.flush_cold_candidate(candidate).await {
445 Ok(response) => responses.push(response),
446 Err(err) if is_stale_cold_flush_candidate_error(&err) => {}
447 Err(err) => return Err(err),
448 }
449 }
450 Ok(responses)
451 }
452
453 async fn flush_cold_candidates_pack(
454 &self,
455 candidates: Vec<ColdFlushCandidate>,
456 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
457 let Some(cold_store) = self.cold_store.as_ref() else {
458 return Err(RuntimeError::ColdStoreConfig {
459 message: "cold backend must be configured before flushing cold chunks".to_owned(),
460 });
461 };
462 let first = candidates
463 .first()
464 .expect("packed cold flush requires at least two candidates");
465 let placement = self.shard_map.locate(&first.stream_id);
466 if candidates
467 .iter()
468 .any(|candidate| self.shard_map.locate(&candidate.stream_id) != placement)
469 {
470 return Err(RuntimeError::ColdStoreConfig {
471 message: "packed cold flush candidates must belong to one Raft group".to_owned(),
472 });
473 }
474 if candidates
475 .iter()
476 .any(|candidate| candidate.stream_id.bucket_id != first.stream_id.bucket_id)
477 {
478 return Err(RuntimeError::ColdStoreConfig {
479 message: "packed cold flush candidates must belong to one bucket erasure domain"
480 .to_owned(),
481 });
482 }
483 let payload_len = candidates.iter().try_fold(0usize, |total, candidate| {
484 total.checked_add(candidate.payload.len())
485 });
486 let Some(payload_len) = payload_len else {
487 return Err(RuntimeError::ColdStoreConfig {
488 message: "packed cold flush payload size overflow".to_owned(),
489 });
490 };
491 let mut payload = Vec::with_capacity(payload_len);
492 let mut object_offsets = Vec::with_capacity(candidates.len());
493 for candidate in &candidates {
494 object_offsets.push(u64::try_from(payload.len()).expect("pack offset fits u64"));
495 payload.extend_from_slice(&candidate.payload);
496 }
497 let path = new_cold_pack_path(&first.stream_id.bucket_id, placement.raft_group_id.0);
498 let upload_started_at = Instant::now();
499 let object_size = match cold_store.write_chunk(&path, &payload).await {
500 Ok(object_size) => object_size,
501 Err(err) => {
502 self.metrics.record_cold_flush_write_error();
503 return Err(RuntimeError::ColdStoreIo {
504 message: err.to_string(),
505 });
506 }
507 };
508 self.metrics
509 .record_cold_upload(object_size, elapsed_ns(upload_started_at));
510 self.metrics.record_cold_pack(
511 object_size,
512 u64::try_from(candidates.len()).expect("candidate count fits u64"),
513 );
514
515 let mut responses = Vec::with_capacity(candidates.len());
516 let mut published = 0usize;
517 for (candidate, object_offset) in candidates.into_iter().zip(object_offsets) {
518 let logical_size = candidate.end_offset.saturating_sub(candidate.start_offset);
519 let publish_started_at = Instant::now();
520 let publish = self
521 .flush_cold(FlushColdRequest {
522 stream_id: candidate.stream_id,
523 chunk: ColdChunkRef {
524 start_offset: candidate.start_offset,
525 end_offset: candidate.end_offset,
526 s3_path: path.clone(),
527 object_size,
528 object_offset,
529 shared_object: true,
530 payload_digest: candidate.payload_digest,
531 },
532 })
533 .await;
534 match publish {
535 Ok(response) => {
536 published = published.saturating_add(1);
537 self.metrics
538 .record_cold_publish(logical_size, elapsed_ns(publish_started_at));
539 responses.push(response);
540 }
541 Err(err) if is_stale_cold_flush_candidate_error(&err) => {}
542 Err(err) => return Err(err),
548 }
549 }
550 if published == 0 {
551 cold_store
552 .delete_chunk(&path)
553 .await
554 .map_err(|err| RuntimeError::ColdStoreIo {
555 message: err.to_string(),
556 })?;
557 }
558 Ok(responses)
559 }
560
561 #[cfg(madsim)]
562 pub async fn flush_cold_candidates_batch_for_simulation(
563 &self,
564 candidates: Vec<ColdFlushCandidate>,
565 ) -> Result<Vec<FlushColdResponse>, RuntimeError> {
566 self.flush_cold_candidates_batch(candidates).await
567 }
568
569 pub async fn purge_bucket_all_groups(
580 &self,
581 bucket_id: &str,
582 ) -> Result<PurgeBucketReport, RuntimeError> {
583 let mut report = PurgeBucketReport::default();
584 let group_count = self.shard_map.raft_group_count();
585 for group_id in 0..group_count {
586 let response = self
587 .purge_bucket(RaftGroupId(group_id), bucket_id.to_owned())
588 .await?;
589 if response.removed_streams > 0 {
590 report.groups_with_streams.push(group_id as usize);
591 }
592 report.removed_streams = report
593 .removed_streams
594 .saturating_add(response.removed_streams);
595 report.pending_cold_gc_entries = report
596 .pending_cold_gc_entries
597 .saturating_add(response.pending_cold_gc_entries);
598 }
599 Ok(report)
600 }
601
602 pub async fn erase_bucket_cold_prefix_and_prove(
608 &self,
609 bucket_id: &str,
610 ) -> Result<(), RuntimeError> {
611 let Some(cold_store) = self.cold_store.as_ref() else {
612 return Ok(());
613 };
614 let prefix = crate::cold_bucket_prefix(bucket_id);
615 cold_store
616 .remove_all(&prefix)
617 .await
618 .map_err(|err| RuntimeError::ColdStoreIo {
619 message: err.to_string(),
620 })?;
621 let absent =
622 cold_store
623 .prefix_is_empty(&prefix)
624 .await
625 .map_err(|err| RuntimeError::ColdStoreIo {
626 message: err.to_string(),
627 })?;
628 if !absent {
629 return Err(RuntimeError::ColdStoreIo {
630 message: format!("bucket prefix '{prefix}' still contains objects after erasure"),
631 });
632 }
633 Ok(())
634 }
635
636 pub async fn bucket_usage_all_groups(
637 &self,
638 ) -> Result<Vec<ursula_stream::BucketUsageSnapshot>, RuntimeError> {
639 let mut merged: std::collections::HashMap<String, ursula_stream::BucketUsage> =
640 std::collections::HashMap::new();
641 let group_count = self.shard_map.raft_group_count();
642 for group_id in 0..group_count {
643 let report = self.bucket_usage(RaftGroupId(group_id)).await?;
644 for entry in report {
645 let usage = merged.entry(entry.bucket_id).or_default();
646 usage.committed_append_bytes = usage
647 .committed_append_bytes
648 .saturating_add(entry.usage.committed_append_bytes);
649 usage.committed_records = usage
650 .committed_records
651 .saturating_add(entry.usage.committed_records);
652 usage.committed_write_units = usage
653 .committed_write_units
654 .saturating_add(entry.usage.committed_write_units);
655 usage.retained_bytes = usage
656 .retained_bytes
657 .saturating_add(entry.usage.retained_bytes);
658 usage.stream_count = usage.stream_count.saturating_add(entry.usage.stream_count);
659 }
660 }
661 let mut report = merged
662 .into_iter()
663 .map(|(bucket_id, usage)| ursula_stream::BucketUsageSnapshot { bucket_id, usage })
664 .collect::<Vec<_>>();
665 report.sort_by(|left, right| left.bucket_id.cmp(&right.bucket_id));
666 Ok(report)
667 }
668
669 pub async fn set_bucket_quota_all_groups(
673 &self,
674 bucket_id: &str,
675 max_streams: Option<u64>,
676 max_retained_bytes: Option<u64>,
677 ) -> Result<(), RuntimeError> {
678 let group_count = self.shard_map.raft_group_count();
679 for group_id in 0..group_count {
680 self.set_bucket_quota(RaftGroupId(group_id), SetBucketQuotaRequest {
681 bucket_id: bucket_id.to_owned(),
682 max_streams,
683 max_retained_bytes,
684 })
685 .await?;
686 }
687 Ok(())
688 }
689
690 pub async fn flush_cold_all_groups_once(
691 &self,
692 request: PlanGroupColdFlushRequest,
693 ) -> Result<usize, RuntimeError> {
694 self.flush_cold_all_groups_once_bounded(request, 1).await
695 }
696
697 pub async fn flush_cold_all_groups_once_bounded(
698 &self,
699 request: PlanGroupColdFlushRequest,
700 max_concurrency: usize,
701 ) -> Result<usize, RuntimeError> {
702 let max_concurrency = max_concurrency.max(1);
703 if max_concurrency == 1 {
704 return self.flush_cold_all_groups_once_serial(request).await;
705 }
706 #[cfg(madsim)]
707 {
708 return self.flush_cold_all_groups_once_serial(request).await;
709 }
710 #[cfg(not(madsim))]
711 {
712 let mut flushed = 0;
713 let mut next_group_id = 0;
714 let group_count = self.shard_map.raft_group_count();
715 let mut tasks = JoinSet::new();
716
717 while next_group_id < group_count || !tasks.is_empty() {
718 while next_group_id < group_count && tasks.len() < max_concurrency {
719 let runtime = self.clone();
720 let request = request.clone();
721 let group_id = RaftGroupId(next_group_id);
722 next_group_id += 1;
723 tasks.spawn(async move {
724 runtime
725 .flush_cold_group_batch_once(
726 group_id,
727 request,
728 COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
729 )
730 .await
731 .map(|responses| responses.len())
732 });
733 }
734 if let Some(result) = tasks.join_next().await {
735 match result {
736 Ok(Ok(count)) => flushed += count,
737 Ok(Err(err)) => return Err(err),
738 Err(err) => {
739 return Err(RuntimeError::ColdStoreIo {
740 message: format!("cold flush task failed: {err}"),
741 });
742 }
743 }
744 }
745 }
746 Ok(flushed)
747 }
748 }
749
750 async fn flush_cold_all_groups_once_serial(
751 &self,
752 request: PlanGroupColdFlushRequest,
753 ) -> Result<usize, RuntimeError> {
754 let mut flushed = 0;
755 for group_id in 0..self.shard_map.raft_group_count() {
756 flushed += self
757 .flush_cold_group_batch_once(
758 RaftGroupId(group_id),
759 request.clone(),
760 COLD_FLUSH_GROUP_BATCH_MAX_CHUNKS,
761 )
762 .await?
763 .len();
764 }
765 Ok(flushed)
766 }
767
768 pub async fn compact_cold_once(
771 &self,
772 target_bytes: u64,
773 max_bytes: u64,
774 max_streams: usize,
775 gc_grace_ms: u64,
776 ) -> Result<usize, RuntimeError> {
777 let Some(cold_store) = self.cold_store.as_ref() else {
778 return Ok(0);
779 };
780 let pages =
781 cold_store
782 .list_cold_index_pages()
783 .await
784 .map_err(|err| RuntimeError::ColdStoreIo {
785 message: err.to_string(),
786 })?;
787 let mut pages_by_stream: HashMap<BucketStreamId, Vec<_>> = HashMap::new();
788 for page in pages {
789 if page.generation == 0 {
790 pages_by_stream
791 .entry(page.stream_id.clone())
792 .or_default()
793 .push(page);
794 }
795 }
796 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
797 let mut compacted = 0;
798 for (stream_id, stream_pages) in pages_by_stream {
799 if compacted >= max_streams {
800 break;
801 }
802 if self
804 .require_local_live_read_owner(&stream_id)
805 .await
806 .is_err()
807 {
808 continue;
809 }
810 let chunks = load_cold_chunks_from_pages(&store, &stream_pages)
811 .await
812 .map_err(|err| RuntimeError::ColdStoreIo {
813 message: err.to_string(),
814 })?;
815 let Some(old_chunks) = select_cold_chunk_compaction(&chunks, target_bytes, max_bytes)
816 else {
817 continue;
818 };
819 let total_bytes = old_chunks
820 .iter()
821 .try_fold(0_u64, |total, chunk| {
822 total.checked_add(chunk.end_offset.saturating_sub(chunk.start_offset))
823 })
824 .ok_or_else(|| RuntimeError::ColdStoreIo {
825 message: "cold compaction byte count overflow".to_owned(),
826 })?;
827 let capacity = usize::try_from(total_bytes).map_err(|_| RuntimeError::ColdStoreIo {
828 message: "cold compaction object exceeds addressable memory".to_owned(),
829 })?;
830 let mut payload = Vec::with_capacity(capacity);
831 for chunk in &old_chunks {
832 let len = usize::try_from(chunk.end_offset.saturating_sub(chunk.start_offset))
833 .map_err(|_| RuntimeError::ColdStoreIo {
834 message: "cold chunk exceeds addressable memory".to_owned(),
835 })?;
836 let bytes = cold_store
837 .read_chunk_range(chunk, chunk.start_offset, len)
838 .await
839 .map_err(|err| RuntimeError::ColdStoreIo {
840 message: err.to_string(),
841 })?;
842 payload.extend_from_slice(&bytes);
843 }
844 let first = old_chunks
845 .first()
846 .expect("candidate contains at least two chunks");
847 let last = old_chunks
848 .last()
849 .expect("candidate contains at least two chunks");
850 let path = new_cold_chunk_path(&stream_id, first.start_offset, last.end_offset);
851 let object_size = cold_store
852 .write_chunk(&path, &payload)
853 .await
854 .map_err(|err| RuntimeError::ColdStoreIo {
855 message: err.to_string(),
856 })?;
857 let replacement = ColdChunkRef {
858 start_offset: first.start_offset,
859 end_offset: last.end_offset,
860 object_size,
861 s3_path: path,
862 object_offset: 0,
863 shared_object: false,
864 payload_digest: blake3::hash(&payload).to_hex().to_string(),
865 };
866 let replacement_path = replacement.s3_path.clone();
867 let gc_not_before_ms = unix_time_ms().saturating_add(gc_grace_ms);
868 let compact_result = self
869 .compact_cold(CompactColdRequest {
870 stream_id: stream_id.clone(),
871 old_chunks,
872 replacement,
873 gc_not_before_ms,
874 })
875 .await;
876 if let Err(err) = compact_result {
877 let rollback_safe =
878 err.leader_hint().is_some() || err.stream_error_code().is_some();
879 if !rollback_safe {
880 return Err(err);
881 }
882 if let Err(cleanup_err) = cold_store.delete_chunk(&replacement_path).await {
883 tracing::warn!(
884 stream = %stream_id,
885 path = %replacement_path,
886 error = %cleanup_err,
887 "failed to remove unpublished cold compaction replacement"
888 );
889 }
890 tracing::warn!(
891 stream = %stream_id,
892 error = %err,
893 "cold compaction publish failed; continuing with remaining streams"
894 );
895 continue;
896 }
897 compacted += 1;
898 }
899 Ok(compacted)
900 }
901
902 pub async fn migrate_legacy_shared_cold_once(
906 &self,
907 max_chunks: usize,
908 gc_grace_ms: u64,
909 ) -> Result<LegacySharedMigrationReport, RuntimeError> {
910 let Some(cold_store) = self.cold_store.as_ref() else {
911 return Ok(LegacySharedMigrationReport::default());
912 };
913 let mut candidates = Vec::new();
914 for group_id in 0..self.shard_map.raft_group_count() {
915 let snapshot = self.snapshot_group(RaftGroupId(group_id)).await?;
916 for stream in snapshot.stream_snapshot.streams {
917 let stream_id = stream.metadata.stream_id;
918 for chunk in stream.cold_chunks {
919 if is_legacy_cross_bucket_pack(&stream_id, &chunk) {
920 candidates.push((stream_id.clone(), chunk));
921 }
922 }
923 }
924 }
925 candidates.sort_by(|left, right| {
926 left.0
927 .bucket_id
928 .cmp(&right.0.bucket_id)
929 .then_with(|| left.0.affinity_key.cmp(&right.0.affinity_key))
930 .then_with(|| left.0.stream_id.cmp(&right.0.stream_id))
931 .then_with(|| left.1.start_offset.cmp(&right.1.start_offset))
932 });
933
934 let observed_chunks = candidates.len();
935 let mut migrated_chunks = 0usize;
936 for (stream_id, chunk) in candidates.into_iter().take(max_chunks) {
937 let logical_bytes = chunk.end_offset.saturating_sub(chunk.start_offset);
938 let len = usize::try_from(logical_bytes).map_err(|_| RuntimeError::ColdStoreIo {
939 message: "legacy shared chunk exceeds addressable memory".to_owned(),
940 })?;
941 let payload = cold_store
942 .read_chunk_range(&chunk, chunk.start_offset, len)
943 .await
944 .map_err(|err| RuntimeError::ColdStoreIo {
945 message: err.to_string(),
946 })?;
947 let path = new_cold_chunk_path(&stream_id, chunk.start_offset, chunk.end_offset);
948 let object_size = cold_store
949 .write_chunk(&path, &payload)
950 .await
951 .map_err(|err| RuntimeError::ColdStoreIo {
952 message: err.to_string(),
953 })?;
954 let replacement = ColdChunkRef {
955 start_offset: chunk.start_offset,
956 end_offset: chunk.end_offset,
957 object_size,
958 s3_path: path.clone(),
959 object_offset: 0,
960 shared_object: false,
961 payload_digest: blake3::hash(&payload).to_hex().to_string(),
962 };
963 if let Err(err) = self
964 .compact_cold(CompactColdRequest {
965 stream_id: stream_id.clone(),
966 old_chunks: vec![chunk],
967 replacement,
968 gc_not_before_ms: unix_time_ms().saturating_add(gc_grace_ms),
969 })
970 .await
971 {
972 if let Err(cleanup_err) = cold_store.delete_chunk(&path).await {
973 tracing::warn!(
974 stream = %stream_id,
975 path,
976 error = %cleanup_err,
977 "failed to remove unpublished legacy pack replacement"
978 );
979 }
980 return Err(err);
981 }
982 migrated_chunks = migrated_chunks.saturating_add(1);
983 }
984 Ok(LegacySharedMigrationReport {
985 observed_chunks,
986 migrated_chunks,
987 pending_chunks: observed_chunks.saturating_sub(migrated_chunks),
988 })
989 }
990
991 pub async fn run_cold_gc_group_once(
996 &self,
997 raft_group_id: RaftGroupId,
998 max_entries: usize,
999 ) -> Result<usize, RuntimeError> {
1000 let Some(cold_store) = self.cold_store.as_ref() else {
1001 return Ok(0);
1002 };
1003 let entries = self.plan_cold_gc(raft_group_id, max_entries).await?;
1004 if entries.is_empty() {
1005 return Ok(0);
1006 }
1007 let mut acked_seq = None;
1008 let mut reclaimed = 0usize;
1009 for entry in entries {
1012 if entry.not_before_ms > unix_time_ms() {
1013 break;
1014 }
1015 let result = match &entry.target {
1016 ColdGcTarget::Stream(stream_id) => {
1017 match cold_store.remove_all(&cold_chunk_prefix(stream_id)).await {
1018 Ok(()) => cold_store.remove_all(&cold_index_prefix(stream_id)).await,
1019 Err(err) => Err(err),
1020 }
1021 }
1022 ColdGcTarget::Paths(paths) => {
1023 let mut outcome = Ok(());
1024 for path in paths {
1025 if let Err(err) = cold_store.delete_chunk(path).await {
1026 outcome = Err(err);
1027 break;
1028 }
1029 }
1030 outcome
1031 }
1032 };
1033 match result {
1034 Ok(()) => {
1035 acked_seq = Some(entry.seq);
1036 reclaimed += 1;
1037 }
1038 Err(err) => {
1039 self.metrics.record_cold_gc_error();
1040 if acked_seq.is_none() {
1041 return Err(RuntimeError::ColdStoreIo {
1042 message: err.to_string(),
1043 });
1044 }
1045 break;
1046 }
1047 }
1048 }
1049 if let Some(up_to_seq) = acked_seq {
1050 self.ack_cold_gc(raft_group_id, up_to_seq).await?;
1051 self.metrics
1052 .record_cold_gc_reclaimed(u64::try_from(reclaimed).expect("reclaimed fits u64"));
1053 }
1054 Ok(reclaimed)
1055 }
1056
1057 pub async fn run_cold_gc_all_groups_once(
1058 &self,
1059 max_entries_per_group: usize,
1060 ) -> Result<usize, RuntimeError> {
1061 if self.cold_store.is_none() {
1062 return Ok(0);
1063 }
1064 let mut reclaimed = 0;
1065 for group_id in 0..self.shard_map.raft_group_count() {
1066 reclaimed += self
1067 .run_cold_gc_group_once(RaftGroupId(group_id), max_entries_per_group)
1068 .await?;
1069 }
1070 Ok(reclaimed)
1071 }
1072
1073 pub fn raft_group_count(&self) -> u32 {
1076 self.shard_map.raft_group_count()
1077 }
1078
1079 pub async fn install_group_snapshot(
1080 &self,
1081 snapshot: GroupSnapshot,
1082 ) -> Result<(), RuntimeError> {
1083 let expected = self.placement_for_group(snapshot.placement.raft_group_id)?;
1084 if snapshot.placement != expected {
1085 return Err(RuntimeError::SnapshotPlacementMismatch {
1086 expected,
1087 actual: snapshot.placement,
1088 });
1089 }
1090 let (response_tx, response_rx) = oneshot::channel();
1091 self.group_rpc(
1092 expected,
1093 None,
1094 GroupCommand::InstallGroupSnapshot {
1095 snapshot,
1096 response_tx,
1097 },
1098 response_rx,
1099 )
1100 .await
1101 }
1102
1103 pub async fn shutdown_group_engine(
1106 &self,
1107 placement: ShardPlacement,
1108 ) -> Result<(), RuntimeError> {
1109 let expected = self.placement_for_group(placement.raft_group_id)?;
1110 if placement != expected {
1111 return Err(RuntimeError::SnapshotPlacementMismatch {
1112 expected,
1113 actual: placement,
1114 });
1115 }
1116 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1117 let (response_tx, response_rx) = oneshot::channel();
1118 self.send_core_command(
1119 mailbox,
1120 CoreCommand::ShutdownGroupEngine {
1121 placement,
1122 response_tx,
1123 },
1124 response_rx,
1125 )
1126 .await
1127 }
1128
1129 #[cfg(madsim)]
1130 pub async fn shutdown_group_engine_for_simulation(
1131 &self,
1132 placement: ShardPlacement,
1133 ) -> Result<(), RuntimeError> {
1134 self.shutdown_group_engine(placement).await
1135 }
1136
1137 #[cfg(madsim)]
1138 pub async fn install_group_engine_for_simulation(
1139 &self,
1140 placement: ShardPlacement,
1141 engine: Box<dyn crate::engine::GroupEngine>,
1142 ) -> Result<(), RuntimeError> {
1143 let expected = self.placement_for_group(placement.raft_group_id)?;
1144 if placement != expected {
1145 return Err(RuntimeError::SnapshotPlacementMismatch {
1146 expected,
1147 actual: placement,
1148 });
1149 }
1150 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1151 let (response_tx, response_rx) = oneshot::channel();
1152 self.send_core_command(
1153 mailbox,
1154 CoreCommand::InstallGroupEngine {
1155 placement,
1156 engine,
1157 response_tx,
1158 },
1159 response_rx,
1160 )
1161 .await
1162 }
1163
1164 pub async fn warm_group(
1165 &self,
1166 raft_group_id: RaftGroupId,
1167 ) -> Result<ShardPlacement, RuntimeError> {
1168 let placement = self.placement_for_group(raft_group_id)?;
1169 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1170 let (response_tx, response_rx) = oneshot::channel();
1171 self.send_core_command(
1172 mailbox,
1173 CoreCommand::WarmGroup {
1174 placement,
1175 response_tx,
1176 },
1177 response_rx,
1178 )
1179 .await
1180 }
1181
1182 pub async fn warm_all_groups(&self) -> Result<(), RuntimeError> {
1183 let mut placements_by_core = vec![Vec::new(); self.mailboxes.len()];
1184 for raw_group_id in 0..self.shard_map.raft_group_count() {
1185 let placement = self.placement_for_group(RaftGroupId(raw_group_id))?;
1186 placements_by_core[usize::from(placement.core_id.0)].push(placement);
1187 }
1188 let mut responses = Vec::new();
1189 for (mailbox, placements) in self.mailboxes.iter().zip(placements_by_core) {
1190 let (response_tx, response_rx) = oneshot::channel();
1191 self.enqueue_core_command(mailbox, CoreCommand::WarmGroups {
1192 placements,
1193 response_tx,
1194 })
1195 .await?;
1196 responses.push((mailbox.core_id, response_rx));
1197 }
1198 for (core_id, response_rx) in responses {
1199 response_rx
1200 .await
1201 .map_err(|_| RuntimeError::ResponseDropped { core_id })??;
1202 }
1203 Ok(())
1204 }
1205
1206 fn placement_for_group(
1207 &self,
1208 raft_group_id: RaftGroupId,
1209 ) -> Result<ShardPlacement, RuntimeError> {
1210 if raft_group_id.0 >= self.shard_map.raft_group_count() {
1211 return Err(RuntimeError::InvalidRaftGroup {
1212 raft_group_id,
1213 raft_group_count: self.shard_map.raft_group_count(),
1214 });
1215 }
1216 Ok(ShardPlacement {
1217 core_id: CoreId(
1218 (raft_group_id.0 % u32::from(self.shard_map.core_count()))
1219 .try_into()
1220 .expect("core id fits u16"),
1221 ),
1222 shard_id: ShardId(raft_group_id.0),
1223 raft_group_id,
1224 })
1225 }
1226
1227 async fn group_rpc<T>(
1232 &self,
1233 placement: ShardPlacement,
1234 admission: Option<u64>,
1235 command: GroupCommand,
1236 response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
1237 ) -> Result<T, RuntimeError> {
1238 let mailbox = &self.mailboxes[usize::from(placement.core_id.0)];
1239 self.send_core_command(
1240 mailbox,
1241 CoreCommand::Group {
1242 placement,
1243 admission,
1244 command,
1245 },
1246 response_rx,
1247 )
1248 .await
1249 }
1250
1251 async fn send_core_command<T>(
1252 &self,
1253 mailbox: &CoreMailbox,
1254 command: CoreCommand,
1255 response_rx: oneshot::Receiver<Result<T, RuntimeError>>,
1256 ) -> Result<T, RuntimeError> {
1257 self.enqueue_core_command(mailbox, command).await?;
1258 response_rx
1259 .await
1260 .map_err(|_| RuntimeError::ResponseDropped {
1261 core_id: mailbox.core_id,
1262 })?
1263 }
1264
1265 async fn enqueue_core_command(
1266 &self,
1267 mailbox: &CoreMailbox,
1268 command: CoreCommand,
1269 ) -> Result<(), RuntimeError> {
1270 if mailbox.tx.capacity() == 0 {
1271 self.metrics.record_mailbox_full(mailbox.core_id);
1272 }
1273 let started_at = Instant::now();
1274 mailbox
1275 .tx
1276 .send(Traced::capture(command))
1277 .await
1278 .map_err(|_| RuntimeError::MailboxClosed {
1279 core_id: mailbox.core_id,
1280 })?;
1281 self.metrics
1282 .record_routed_request(mailbox.core_id, elapsed_ns(started_at));
1283 Ok(())
1284 }
1285
1286 pub async fn append_transaction(
1291 &self,
1292 request: AppendTransactionRequest,
1293 ) -> Result<AppendTransactionResponse, RuntimeError> {
1294 const MAX_OPERATIONS: usize = 64;
1295 let Some(first) = request.operations.first() else {
1296 return Err(RuntimeError::InvalidAppendTransaction {
1297 message: "at least one append operation is required".to_owned(),
1298 });
1299 };
1300 if request.operations.len() > MAX_OPERATIONS {
1301 return Err(RuntimeError::InvalidAppendTransaction {
1302 message: format!("at most {MAX_OPERATIONS} append operations are allowed"),
1303 });
1304 }
1305 if request
1306 .operations
1307 .iter()
1308 .any(|operation| operation.payload.is_empty())
1309 {
1310 return Err(RuntimeError::InvalidAppendTransaction {
1311 message: "every append operation must have a non-empty payload".to_owned(),
1312 });
1313 }
1314 let Some(affinity_key) = first.stream_id.affinity_key.as_deref() else {
1315 return Err(RuntimeError::InvalidAppendTransaction {
1316 message: "all streams must use an affinity path".to_owned(),
1317 });
1318 };
1319 let placement = self.shard_map.locate(&first.stream_id);
1320 for operation in &request.operations {
1321 if operation.stream_id.bucket_id != first.stream_id.bucket_id
1322 || operation.stream_id.affinity_key.as_deref() != Some(affinity_key)
1323 || self.shard_map.locate(&operation.stream_id) != placement
1324 {
1325 return Err(RuntimeError::InvalidAppendTransaction {
1326 message: "all streams must share one bucket and affinity key".to_owned(),
1327 });
1328 }
1329 }
1330 let incoming_bytes = request.payload_bytes();
1331 let (response_tx, response_rx) = oneshot::channel();
1332 self.group_rpc(
1333 placement,
1334 Some(incoming_bytes),
1335 GroupCommand::AppendTransaction {
1336 request,
1337 response_tx,
1338 raft_uncommitted: None,
1339 },
1340 response_rx,
1341 )
1342 .await
1343 }
1344
1345 pub fn metrics(&self) -> RuntimeMetrics {
1346 RuntimeMetrics {
1347 inner: self.metrics.clone(),
1348 }
1349 }
1350
1351 pub fn mailbox_snapshot(&self) -> RuntimeMailboxSnapshot {
1352 let depths = self
1353 .mailboxes
1354 .iter()
1355 .map(CoreMailbox::depth)
1356 .collect::<Vec<_>>();
1357 let capacities = self
1358 .mailboxes
1359 .iter()
1360 .map(CoreMailbox::capacity)
1361 .collect::<Vec<_>>();
1362 RuntimeMailboxSnapshot { depths, capacities }
1363 }
1364}
1365
1366#[cfg(not(madsim))]
1367fn unix_time_ms() -> u64 {
1368 SystemTime::now()
1369 .duration_since(UNIX_EPOCH)
1370 .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
1371 .unwrap_or(0)
1372}
1373
1374#[cfg(madsim)]
1375fn unix_time_ms() -> u64 {
1376 0
1377}
1378
1379macro_rules! shard_runtime_operations {
1385 (@munch
1387 methods { $($methods:tt)* }
1388 rest {
1389 $(#[$attr:meta])*
1390 op $Variant:ident {
1391 fields { $req:ident: $Req:ty $(,)? }
1392 reply { $tx:ident: $Resp:ty }
1393 guard { $g:ident }
1394 handle { $($handle:tt)* }
1395 client {
1396 $vis:vis stream fn $method:ident,
1397 non_empty: $ne:ident,
1398 admit: $incoming:expr
1399 }
1400 }
1401 $($rest:tt)*
1402 }
1403 ) => {
1404 shard_runtime_operations! {
1405 @munch
1406 methods {
1407 $($methods)*
1408 $(#[$attr])*
1409 $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
1410 if $req.$ne.is_empty() {
1411 return Err(RuntimeError::EmptyAppend);
1412 }
1413 let placement = self.shard_map.locate(&$req.stream_id);
1414 let incoming_bytes = $incoming;
1415 let (response_tx, response_rx) = oneshot::channel();
1416 self.group_rpc(
1417 placement,
1418 Some(incoming_bytes),
1419 GroupCommand::$Variant {
1420 $req,
1421 $tx: response_tx,
1422 $g: None,
1423 },
1424 response_rx,
1425 )
1426 .await
1427 }
1428 }
1429 rest { $($rest)* }
1430 }
1431 };
1432 (@munch
1434 methods { $($methods:tt)* }
1435 rest {
1436 $(#[$attr:meta])*
1437 op $Variant:ident {
1438 fields { $req:ident: $Req:ty $(,)? }
1439 reply { $tx:ident: $Resp:ty }
1440 guard { $g:ident }
1441 handle { $($handle:tt)* }
1442 client { $vis:vis stream fn $method:ident, admit: $incoming:expr }
1443 }
1444 $($rest:tt)*
1445 }
1446 ) => {
1447 shard_runtime_operations! {
1448 @munch
1449 methods {
1450 $($methods)*
1451 $(#[$attr])*
1452 $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
1453 let placement = self.shard_map.locate(&$req.stream_id);
1454 let incoming_bytes = $incoming;
1455 let (response_tx, response_rx) = oneshot::channel();
1456 self.group_rpc(
1457 placement,
1458 Some(incoming_bytes),
1459 GroupCommand::$Variant {
1460 $req,
1461 $tx: response_tx,
1462 $g: None,
1463 },
1464 response_rx,
1465 )
1466 .await
1467 }
1468 }
1469 rest { $($rest)* }
1470 }
1471 };
1472 (@munch
1474 methods { $($methods:tt)* }
1475 rest {
1476 $(#[$attr:meta])*
1477 op $Variant:ident {
1478 fields { $req:ident: $Req:ty $(,)? }
1479 reply { $tx:ident: $Resp:ty }
1480 guard { none }
1481 handle { $($handle:tt)* }
1482 client { $vis:vis stream fn $method:ident }
1483 }
1484 $($rest:tt)*
1485 }
1486 ) => {
1487 shard_runtime_operations! {
1488 @munch
1489 methods {
1490 $($methods)*
1491 $(#[$attr])*
1492 $vis async fn $method(&self, $req: $Req) -> Result<$Resp, RuntimeError> {
1493 let placement = self.shard_map.locate(&$req.stream_id);
1494 let (response_tx, response_rx) = oneshot::channel();
1495 self.group_rpc(
1496 placement,
1497 None,
1498 GroupCommand::$Variant {
1499 $req,
1500 $tx: response_tx,
1501 },
1502 response_rx,
1503 )
1504 .await
1505 }
1506 }
1507 rest { $($rest)* }
1508 }
1509 };
1510 (@munch
1513 methods { $($methods:tt)* }
1514 rest {
1515 $(#[$attr:meta])*
1516 op $Variant:ident {
1517 fields { $($field:ident: $field_ty:ty),* $(,)? }
1518 reply { $tx:ident: $Resp:ty }
1519 guard { none }
1520 handle { $($handle:tt)* }
1521 client { $vis:vis group fn $method:ident }
1522 }
1523 $($rest:tt)*
1524 }
1525 ) => {
1526 shard_runtime_operations! {
1527 @munch
1528 methods {
1529 $($methods)*
1530 $(#[$attr])*
1531 $vis async fn $method(
1532 &self,
1533 raft_group_id: RaftGroupId,
1534 $($field: $field_ty),*
1535 ) -> Result<$Resp, RuntimeError> {
1536 let placement = self.placement_for_group(raft_group_id)?;
1537 let (response_tx, response_rx) = oneshot::channel();
1538 self.group_rpc(
1539 placement,
1540 None,
1541 GroupCommand::$Variant {
1542 $($field,)*
1543 $tx: response_tx,
1544 },
1545 response_rx,
1546 )
1547 .await
1548 }
1549 }
1550 rest { $($rest)* }
1551 }
1552 };
1553 (@munch
1555 methods { $($methods:tt)* }
1556 rest {
1557 $(#[$attr:meta])*
1558 op $Variant:ident {
1559 fields { $($fields:tt)* }
1560 reply { $($reply:tt)* }
1561 guard { $($guard:tt)* }
1562 handle { $($handle:tt)* }
1563 client { none }
1564 }
1565 $($rest:tt)*
1566 }
1567 ) => {
1568 shard_runtime_operations! {
1569 @munch
1570 methods { $($methods)* }
1571 rest { $($rest)* }
1572 }
1573 };
1574 (@munch
1575 methods { $($methods:tt)* }
1576 rest {}
1577 ) => {
1578 impl ShardRuntime {
1579 $($methods)*
1580 }
1581 };
1582 ($($manifest:tt)*) => {
1583 shard_runtime_operations! {
1584 @munch
1585 methods {}
1586 rest { $($manifest)* }
1587 }
1588 };
1589}
1590
1591crate::ops::runtime_operations!(shard_runtime_operations);
1592
1593fn spawn_core_worker(threading: RuntimeThreading, worker: CoreWorker) -> Result<(), RuntimeError> {
1594 match threading {
1595 RuntimeThreading::HostedTokio => {
1596 crate::rt::spawn(worker.run());
1597 Ok(())
1598 }
1599 #[cfg(not(madsim))]
1600 RuntimeThreading::ThreadPerCore => {
1601 let core_id = worker.core_id;
1602 std::thread::Builder::new()
1603 .name(format!("ursula-core-{}", core_id.0))
1604 .spawn(move || {
1605 let runtime = tokio::runtime::Builder::new_current_thread()
1606 .enable_all()
1607 .build()
1608 .expect("build per-core tokio runtime");
1609 runtime.block_on(worker.run());
1610 })
1611 .map(|_| ())
1612 .map_err(|err| RuntimeError::SpawnCoreThread {
1613 core_id,
1614 message: err.to_string(),
1615 })
1616 }
1617 }
1618}