1use std::cell::RefCell;
16use std::collections::{HashMap, HashSet, VecDeque};
17use std::pin::pin;
18use std::sync::atomic::AtomicUsize;
19use std::sync::atomic::Ordering::Relaxed;
20use std::sync::{Arc, LazyLock};
21use std::time::Duration;
22
23use arc_swap::ArcSwap;
24use await_tree::{InstrumentAwait, SpanExt};
25use futures::FutureExt;
26use itertools::Itertools;
27use parking_lot::RwLock;
28use prometheus::{Histogram, IntGauge, IntGaugeVec};
29use risingwave_common::bitmap::Bitmap;
30use risingwave_common::catalog::TableId;
31use risingwave_common::config::Role;
32use risingwave_common::config::streaming::CacheRefillPolicy;
33use risingwave_common::metrics::UintGauge;
34use risingwave_hummock_sdk::compaction_group::hummock_version_ext::SstDeltaInfo;
35use risingwave_hummock_sdk::{HummockEpoch, SyncResult};
36use risingwave_pb::meta::table_cache_refill_policies::table_cache_refill_policy::PbCacheRefillPolicy;
37use risingwave_pb::meta::{PbServingTableVnodeMappings, PbTableRefillRuntimeConfig};
38use tokio::spawn;
39use tokio::sync::mpsc::error::SendError;
40use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel};
41use tokio::sync::oneshot;
42use tracing::{debug, error, info, trace, warn};
43
44use super::refiller::{CacheRefillConfig, CacheRefiller};
45use super::{LocalInstanceGuard, LocalInstanceId, ReadVersionMappingType};
46use crate::compaction_catalog_manager::CompactionCatalogManagerRef;
47use crate::hummock::compactor::{CompactorContext, await_tree_key, compact};
48use crate::hummock::event_handler::refiller::{CacheRefillerEvent, SpawnRefillTask};
49use crate::hummock::event_handler::uploader::{
50 HummockUploader, SpawnUploadTask, SyncedData, UploadTaskOutput,
51};
52use crate::hummock::event_handler::{
53 HummockEvent, HummockObserverEvent, HummockReadVersionRef, HummockVersionUpdate,
54 ReadOnlyReadVersionMapping, ReadOnlyRwLockRef,
55};
56use crate::hummock::local_version::pinned_version::PinnedVersion;
57use crate::hummock::local_version::recent_versions::RecentVersions;
58use crate::hummock::store::version::{HummockReadVersion, StagingSstableInfo, VersionUpdate};
59use crate::hummock::{HummockResult, MemoryLimiter, ObjectIdManager, SstableStoreRef};
60use crate::mem_table::ImmutableMemtable;
61use crate::monitor::HummockStateStoreMetrics;
62use crate::opts::StorageOpts;
63
64#[derive(Clone)]
65pub(crate) struct BufferTracker {
66 flush_threshold: usize,
67 min_batch_flush_size: usize,
68 global_uploading_memory_limiter: Arc<MemoryLimiter>,
69 uploader_imm_size: UintGauge,
70 uploader_uploading_task_size: UintGauge,
71}
72
73impl BufferTracker {
74 pub fn from_storage_opts(config: &StorageOpts, metrics: &HummockStateStoreMetrics) -> Self {
75 let capacity = config.shared_buffer_capacity_mb * (1 << 20);
76 let flush_threshold = (capacity as f32 * config.shared_buffer_flush_ratio) as usize;
77 let min_batch_flush_size = config.shared_buffer_min_batch_flush_size_mb * (1 << 20);
78 assert!(
79 flush_threshold < capacity,
80 "flush_threshold {} should be less or equal to capacity {}",
81 flush_threshold,
82 capacity
83 );
84 Self {
85 flush_threshold,
86 min_batch_flush_size,
87 global_uploading_memory_limiter: Arc::new(MemoryLimiter::new(capacity as u64)),
88 uploader_imm_size: metrics.uploader_imm_size.clone(),
89 uploader_uploading_task_size: metrics.uploader_uploading_task_size.clone(),
90 }
91 }
92
93 #[cfg(test)]
94 fn for_test_with_config(flush_threshold: usize, min_batch_flush_size: usize) -> Self {
95 Self {
96 flush_threshold,
97 min_batch_flush_size,
98 ..Self::for_test()
99 }
100 }
101
102 #[cfg(test)]
103 pub fn for_test() -> Self {
104 Self::from_storage_opts(&StorageOpts::default(), &HummockStateStoreMetrics::unused())
105 }
106
107 pub fn get_memory_limiter(&self) -> &Arc<MemoryLimiter> {
108 &self.global_uploading_memory_limiter
109 }
110
111 pub fn global_upload_task_size(&self) -> &UintGauge {
112 &self.uploader_uploading_task_size
113 }
114
115 pub fn need_flush(&self) -> bool {
118 self.uploader_imm_size.get()
119 > self.flush_threshold as u64 + self.uploader_uploading_task_size.get()
120 }
121
122 pub fn need_more_flush(&self, curr_batch_flush_size: usize) -> bool {
123 curr_batch_flush_size < self.min_batch_flush_size || self.need_flush()
124 }
125
126 #[cfg(test)]
127 pub(crate) fn flush_threshold(&self) -> usize {
128 self.flush_threshold
129 }
130}
131
132#[derive(Clone)]
133pub struct HummockEventSender {
134 inner: UnboundedSender<HummockEvent>,
135 event_count: IntGaugeVec,
136}
137
138pub fn event_channel(event_count: IntGaugeVec) -> (HummockEventSender, HummockEventReceiver) {
139 let (tx, rx) = unbounded_channel();
140 (
141 HummockEventSender {
142 inner: tx,
143 event_count: event_count.clone(),
144 },
145 HummockEventReceiver {
146 inner: rx,
147 event_count,
148 },
149 )
150}
151
152impl HummockEventSender {
153 pub fn send(&self, event: HummockEvent) -> Result<(), SendError<HummockEvent>> {
154 let event_type = event.event_name();
155 self.inner.send(event)?;
156 get_event_pending_gauge(&self.event_count, event_type).inc();
157 Ok(())
158 }
159}
160
161pub struct HummockEventReceiver {
162 inner: UnboundedReceiver<HummockEvent>,
163 event_count: IntGaugeVec,
164}
165
166impl HummockEventReceiver {
167 async fn recv(&mut self) -> Option<HummockEvent> {
168 let event = self.inner.recv().await?;
169 let event_type = event.event_name();
170 get_event_pending_gauge(&self.event_count, event_type).dec();
171 Some(event)
172 }
173}
174
175thread_local! {
176 static EVENT_PENDING_GAUGE_CACHE: RefCell<HashMap<&'static str, IntGauge>> = RefCell::new(HashMap::new());
177}
178
179fn get_event_pending_gauge(event_count: &IntGaugeVec, event_type: &'static str) -> IntGauge {
180 EVENT_PENDING_GAUGE_CACHE.with(|cache| {
181 cache
182 .borrow_mut()
183 .entry(event_type)
184 .or_insert_with(|| event_count.with_label_values(&[event_type]))
185 .clone()
186 })
187}
188
189struct HummockEventHandlerMetrics {
190 event_handler_on_upload_finish_latency: Histogram,
191 event_handler_on_apply_version_update: Histogram,
192 event_handler_on_recv_version_update: Histogram,
193 event_handler_on_spiller: Histogram,
194 event_handler_on_apply_table_refill_runtime_config: Histogram,
195}
196
197fn table_cache_refill_policy_from_protobuf(
198 pb_policy: risingwave_pb::meta::table_cache_refill_policies::PbTableCacheRefillPolicy,
199) -> Option<(TableId, CacheRefillPolicy)> {
200 let table_id = TableId::from(pb_policy.table_id);
201 let policy = match PbCacheRefillPolicy::try_from(pb_policy.policy) {
202 Ok(policy) => CacheRefillPolicy::from_protobuf(policy)?,
203 Err(_) => {
204 tracing::warn!(
205 table_id = pb_policy.table_id,
206 policy = pb_policy.policy,
207 "ignore unknown table cache refill policy"
208 );
209 return None;
210 }
211 };
212 Some((table_id, policy))
213}
214
215fn serving_table_vnode_mappings_to_map(
216 mappings: PbServingTableVnodeMappings,
217) -> HashMap<TableId, Bitmap> {
218 mappings
219 .mappings
220 .into_iter()
221 .filter_map(|mapping| {
222 let Some(bitmap) = mapping.bitmap else {
223 tracing::warn!(
224 table_id = mapping.table_id,
225 "ignore serving table vnode mapping without bitmap"
226 );
227 return None;
228 };
229 Some((TableId::from(mapping.table_id), Bitmap::from(bitmap)))
230 })
231 .collect()
232}
233
234#[derive(Default)]
235struct LocalReadVersions {
236 by_instance: HashMap<LocalInstanceId, (TableId, HummockReadVersionRef)>,
237 by_table: HashMap<TableId, HashMap<LocalInstanceId, HummockReadVersionRef>>,
238}
239
240impl LocalReadVersions {
241 fn insert(
242 &mut self,
243 instance_id: LocalInstanceId,
244 table_id: TableId,
245 read_version: HummockReadVersionRef,
246 ) {
247 assert!(
248 !self.by_instance.contains_key(&instance_id),
249 "duplicated local read version instance_id {}",
250 instance_id
251 );
252 assert!(
253 !self
254 .by_table
255 .get(&table_id)
256 .is_some_and(|read_versions| read_versions.contains_key(&instance_id)),
257 "duplicated local read version table_id {} instance_id {}",
258 table_id,
259 instance_id
260 );
261 self.by_instance
262 .insert(instance_id, (table_id, read_version.clone()));
263 self.by_table
264 .entry(table_id)
265 .or_default()
266 .insert(instance_id, read_version);
267 }
268
269 fn remove(&mut self, instance_id: LocalInstanceId) -> (TableId, HummockReadVersionRef) {
270 let (table_id, read_version) = self.by_instance.remove(&instance_id).unwrap_or_else(|| {
271 panic!(
272 "DestroyHummockInstance nonexistent instance instance_id {}",
273 instance_id
274 )
275 });
276 let table_read_versions = self.by_table.get_mut(&table_id).unwrap_or_else(|| {
277 panic!(
278 "DestroyHummockInstance table_id {} instance_id {} fail in local mapping",
279 table_id, instance_id
280 )
281 });
282 let removed = table_read_versions.remove(&instance_id).unwrap_or_else(|| {
283 panic!(
284 "DestroyHummockInstance nonexistent instance in local mapping table_id {} instance_id {}",
285 table_id, instance_id
286 )
287 });
288 debug_assert!(Arc::ptr_eq(&read_version, &removed));
289 if table_read_versions.is_empty() {
290 self.by_table.remove(&table_id);
291 }
292 (table_id, read_version)
293 }
294
295 fn get(&self, instance_id: &LocalInstanceId) -> Option<&(TableId, HummockReadVersionRef)> {
296 self.by_instance.get(instance_id)
297 }
298
299 fn contains_instance(&self, instance_id: &LocalInstanceId) -> bool {
300 self.by_instance.contains_key(instance_id)
301 }
302
303 fn instance_ids(&self) -> impl Iterator<Item = LocalInstanceId> + '_ {
304 self.by_instance.keys().copied()
305 }
306
307 fn is_empty(&self) -> bool {
308 self.by_instance.is_empty() && self.by_table.is_empty()
309 }
310
311 fn table_ids(&self) -> impl Iterator<Item = TableId> + '_ {
312 self.by_table.keys().copied()
313 }
314
315 fn table_vnodes(&self, table_id: TableId) -> Option<Bitmap> {
316 let read_versions = self.by_table.get(&table_id)?;
317 let mut vnodes = None;
318 for read_version in read_versions.values() {
319 let read_version_vnodes = read_version.read().vnodes();
320 let current_vnodes =
321 vnodes.get_or_insert_with(|| Bitmap::zeros(read_version_vnodes.as_ref().len()));
322 *current_vnodes |= read_version_vnodes.as_ref();
323 }
324 vnodes
325 }
326}
327
328pub(crate) struct HummockEventHandler {
329 hummock_event_tx: HummockEventSender,
330 hummock_event_rx: HummockEventReceiver,
331 observer_event_rx: UnboundedReceiver<HummockObserverEvent>,
332 read_version_mapping: Arc<RwLock<ReadVersionMappingType>>,
333 local_read_versions: LocalReadVersions,
335
336 version_update_notifier_tx: Arc<tokio::sync::watch::Sender<PinnedVersion>>,
337 recent_versions: Arc<ArcSwap<RecentVersions>>,
338
339 uploader: HummockUploader,
340 refiller: CacheRefiller,
341
342 last_instance_id: LocalInstanceId,
343
344 metrics: HummockEventHandlerMetrics,
345}
346
347async fn flush_imms(
348 payload: Vec<ImmutableMemtable>,
349 compactor_context: CompactorContext,
350 compaction_catalog_manager_ref: CompactionCatalogManagerRef,
351 object_id_manager: Arc<ObjectIdManager>,
352) -> HummockResult<UploadTaskOutput> {
353 compact(
354 compactor_context,
355 object_id_manager,
356 payload,
357 compaction_catalog_manager_ref,
358 )
359 .instrument_await("shared_buffer_compact".verbose())
360 .await
361}
362
363impl HummockEventHandler {
364 fn apply_table_refill_runtime_config(
367 refiller: &mut CacheRefiller,
368 config: PbTableRefillRuntimeConfig,
369 ) {
370 if let Some(policies) = config.table_cache_refill_policies {
371 let policies = policies
372 .table_policies
373 .into_iter()
374 .chain(policies.internal_table_policies)
375 .filter_map(table_cache_refill_policy_from_protobuf)
376 .collect();
377 refiller.replace_table_cache_refill_policies(policies);
378 }
379
380 if let Some(mappings) = config.serving_table_vnode_mappings {
381 refiller
382 .replace_serving_table_vnode_mapping(serving_table_vnode_mappings_to_map(mappings));
383 }
384 }
385
386 pub fn new(
387 role: Role,
388 observer_event_rx: UnboundedReceiver<HummockObserverEvent>,
389 pinned_version: PinnedVersion,
390 compactor_context: CompactorContext,
391 compaction_catalog_manager_ref: CompactionCatalogManagerRef,
392 object_id_manager: Arc<ObjectIdManager>,
393 state_store_metrics: Arc<HummockStateStoreMetrics>,
394 ) -> Self {
395 let upload_compactor_context = compactor_context.clone();
396 let upload_task_latency = state_store_metrics.uploader_upload_task_latency.clone();
397 let wait_poll_latency = state_store_metrics.uploader_wait_poll_latency.clone();
398 let recent_versions = RecentVersions::new(
399 pinned_version,
400 compactor_context
401 .storage_opts
402 .max_cached_recent_versions_number,
403 state_store_metrics.clone(),
404 );
405 let buffer_tracker =
406 BufferTracker::from_storage_opts(&compactor_context.storage_opts, &state_store_metrics);
407 Self::new_inner(
408 role,
409 observer_event_rx,
410 compactor_context.sstable_store.clone(),
411 state_store_metrics,
412 CacheRefillConfig::from_storage_opts(&compactor_context.storage_opts),
413 recent_versions,
414 buffer_tracker,
415 Arc::new(move |payload, task_info| {
416 static NEXT_UPLOAD_TASK_ID: LazyLock<AtomicUsize> =
417 LazyLock::new(|| AtomicUsize::new(0));
418 let tree_root = upload_compactor_context.await_tree_reg.as_ref().map(|reg| {
419 let upload_task_id = NEXT_UPLOAD_TASK_ID.fetch_add(1, Relaxed);
420 reg.register(
421 await_tree_key::SpawnUploadTask { id: upload_task_id },
422 format!("Spawn Upload Task: {}", task_info),
423 )
424 });
425 let upload_task_latency = upload_task_latency.clone();
426 let wait_poll_latency = wait_poll_latency.clone();
427 let upload_compactor_context = upload_compactor_context.clone();
428 let compaction_catalog_manager_ref = compaction_catalog_manager_ref.clone();
429 let object_id_manager = object_id_manager.clone();
430 spawn({
431 let future = async move {
432 let _timer = upload_task_latency.start_timer();
433 let mut output = flush_imms(
434 payload
435 .into_values()
436 .flat_map(|imms| imms.into_iter())
437 .collect(),
438 upload_compactor_context.clone(),
439 compaction_catalog_manager_ref.clone(),
440 object_id_manager.clone(),
441 )
442 .await?;
443 assert!(
444 output
445 .wait_poll_timer
446 .replace(wait_poll_latency.start_timer())
447 .is_none(),
448 "should not set timer before"
449 );
450 Ok(output)
451 };
452 if let Some(tree_root) = tree_root {
453 tree_root.instrument(future).left_future()
454 } else {
455 future.right_future()
456 }
457 })
458 }),
459 CacheRefiller::default_spawn_refill_task(),
460 )
461 }
462
463 fn new_inner(
464 role: Role,
465 observer_event_rx: UnboundedReceiver<HummockObserverEvent>,
466 sstable_store: SstableStoreRef,
467 state_store_metrics: Arc<HummockStateStoreMetrics>,
468 refill_config: CacheRefillConfig,
469 recent_versions: RecentVersions,
470 buffer_tracker: BufferTracker,
471 spawn_upload_task: SpawnUploadTask,
472 spawn_refill_task: SpawnRefillTask,
473 ) -> Self {
474 let (hummock_event_tx, hummock_event_rx) =
475 event_channel(state_store_metrics.event_handler_pending_event.clone());
476 let (version_update_notifier_tx, _) =
477 tokio::sync::watch::channel(recent_versions.latest_version().clone());
478 let version_update_notifier_tx = Arc::new(version_update_notifier_tx);
479 let read_version_mapping = Arc::new(RwLock::new(HashMap::default()));
480
481 let metrics = HummockEventHandlerMetrics {
482 event_handler_on_upload_finish_latency: state_store_metrics
483 .event_handler_latency
484 .with_label_values(&["on_upload_finish"]),
485 event_handler_on_apply_version_update: state_store_metrics
486 .event_handler_latency
487 .with_label_values(&["apply_version"]),
488 event_handler_on_recv_version_update: state_store_metrics
489 .event_handler_latency
490 .with_label_values(&["recv_version_update"]),
491 event_handler_on_spiller: state_store_metrics
492 .event_handler_latency
493 .with_label_values(&["spiller"]),
494 event_handler_on_apply_table_refill_runtime_config: state_store_metrics
495 .event_handler_latency
496 .with_label_values(&["apply_table_refill_runtime_config"]),
497 };
498
499 let uploader = HummockUploader::new(
500 state_store_metrics,
501 recent_versions.latest_version().clone(),
502 spawn_upload_task,
503 buffer_tracker,
504 );
505 let refiller = CacheRefiller::new(role, refill_config, sstable_store, spawn_refill_task);
506
507 Self {
508 hummock_event_tx,
509 hummock_event_rx,
510 observer_event_rx,
511 version_update_notifier_tx,
512 recent_versions: Arc::new(ArcSwap::from_pointee(recent_versions)),
513 read_version_mapping,
514 local_read_versions: Default::default(),
515 uploader,
516 refiller,
517 last_instance_id: 0,
518 metrics,
519 }
520 }
521
522 pub fn version_update_notifier_tx(&self) -> Arc<tokio::sync::watch::Sender<PinnedVersion>> {
523 self.version_update_notifier_tx.clone()
524 }
525
526 pub fn recent_versions(&self) -> Arc<ArcSwap<RecentVersions>> {
527 self.recent_versions.clone()
528 }
529
530 pub fn read_version_mapping(&self) -> ReadOnlyReadVersionMapping {
531 ReadOnlyRwLockRef::new(self.read_version_mapping.clone())
532 }
533
534 pub fn event_sender(&self) -> HummockEventSender {
535 self.hummock_event_tx.clone()
536 }
537
538 pub fn buffer_tracker(&self) -> &BufferTracker {
539 self.uploader.buffer_tracker()
540 }
541}
542
543impl HummockEventHandler {
545 fn for_each_read_version(
548 &self,
549 instances: impl IntoIterator<Item = LocalInstanceId>,
550 mut f: impl FnMut(LocalInstanceId, &mut HummockReadVersion),
551 ) {
552 let instances = {
553 #[cfg(debug_assertions)]
554 {
555 let mut id_set = std::collections::HashSet::new();
557 for instance in instances {
558 assert!(id_set.insert(instance));
559 }
560 id_set
561 }
562 #[cfg(not(debug_assertions))]
563 {
564 instances
565 }
566 };
567 let mut pending = VecDeque::new();
568 let mut total_count = 0;
569 for instance_id in instances {
570 let Some((_, read_version)) = self.local_read_versions.get(&instance_id) else {
571 continue;
572 };
573 total_count += 1;
574 match read_version.try_write() {
575 Some(mut write_guard) => {
576 f(instance_id, &mut write_guard);
577 }
578 _ => {
579 pending.push_back(instance_id);
580 }
581 }
582 }
583 if !pending.is_empty() {
584 if pending.len() * 10 > total_count {
585 warn!(
587 pending_count = pending.len(),
588 total_count, "cannot acquire lock for all read version"
589 );
590 } else {
591 debug!(
592 pending_count = pending.len(),
593 total_count, "cannot acquire lock for all read version"
594 );
595 }
596 }
597
598 const TRY_LOCK_TIMEOUT: Duration = Duration::from_millis(1);
599
600 while let Some(instance_id) = pending.pop_front() {
601 let (_, read_version) = self
602 .local_read_versions
603 .get(&instance_id)
604 .expect("have checked exist before");
605 match read_version.try_write_for(TRY_LOCK_TIMEOUT) {
606 Some(mut write_guard) => {
607 f(instance_id, &mut write_guard);
608 }
609 _ => {
610 warn!(instance_id, "failed to get lock again for instance");
611 pending.push_back(instance_id);
612 }
613 }
614 }
615 }
616
617 fn handle_uploaded_ssts_inner(&mut self, ssts: Vec<Arc<StagingSstableInfo>>) {
618 match ssts.as_slice() {
619 [] => {
620 if cfg!(debug_assertions) {
621 panic!("empty ssts")
622 }
623 }
624 [staging_sstable_info] => {
625 trace!("data_flushed. SST size {}", staging_sstable_info.imm_size());
626 self.for_each_read_version(
627 staging_sstable_info.imm_ids().keys().cloned(),
628 |_, read_version| {
629 read_version.update(VersionUpdate::Sst(staging_sstable_info.clone()))
630 },
631 )
632 }
633 ssts => {
634 warn!(
635 batch_size = ssts.len(),
636 "handle multiple uploaded ssts in batch"
637 );
638 let affected_instances: HashSet<_> = ssts
639 .iter()
640 .flat_map(|sst| {
641 trace!("data_flushed. SST size {}", sst.imm_size());
642 sst.imm_ids().keys()
643 })
644 .copied()
645 .collect();
646 self.for_each_read_version(affected_instances, |instance_id, read_version| {
647 for sst in ssts {
648 if sst.imm_ids().contains_key(&instance_id) {
649 read_version.update(VersionUpdate::Sst(sst.clone()));
650 }
651 }
652 })
653 }
654 }
655 }
656
657 fn handle_sync_epoch(
658 &mut self,
659 sync_table_epochs: Vec<(HummockEpoch, HashSet<TableId>)>,
660 sync_result_sender: oneshot::Sender<HummockResult<SyncedData>>,
661 ) {
662 debug!(?sync_table_epochs, "awaiting for epoch to be synced",);
663 self.uploader
664 .start_sync_epoch(sync_result_sender, sync_table_epochs);
665 }
666
667 fn handle_clear(&mut self, notifier: oneshot::Sender<()>, table_ids: Option<HashSet<TableId>>) {
668 info!(
669 current_version_id = ?self.uploader.hummock_version().id(),
670 ?table_ids,
671 "handle clear event"
672 );
673
674 self.uploader.clear(table_ids.clone());
675
676 if table_ids.is_none() {
677 assert!(
678 self.local_read_versions.is_empty(),
679 "read version mapping not empty when clear. remaining tables: {:?}",
680 self.local_read_versions.table_ids().collect_vec()
681 );
682 }
683
684 let _ = notifier.send(()).inspect_err(|e| {
686 error!("failed to notify completion of clear event: {:?}", e);
687 });
688
689 info!("clear finished");
690 }
691
692 fn handle_observer_event(&mut self, event: HummockObserverEvent) {
693 match event {
694 HummockObserverEvent::VersionUpdate(version_update) => {
695 self.handle_version_update(version_update);
696 }
697 HummockObserverEvent::TableRefillRuntimeConfig(_, config) => {
698 let _timer = self
699 .metrics
700 .event_handler_on_apply_table_refill_runtime_config
701 .start_timer();
702 Self::apply_table_refill_runtime_config(&mut self.refiller, config);
703 }
704 }
705 }
706
707 fn streaming_table_vnodes(&self, table_id: TableId) -> Option<Bitmap> {
708 self.local_read_versions.table_vnodes(table_id)
709 }
710
711 fn handle_version_update(&mut self, version_payload: HummockVersionUpdate) {
712 let _timer = self
713 .metrics
714 .event_handler_on_recv_version_update
715 .start_timer();
716 let pinned_version = self
717 .refiller
718 .last_new_pinned_version()
719 .cloned()
720 .unwrap_or_else(|| self.uploader.hummock_version().clone());
721
722 let mut sst_delta_infos = vec![];
723 if let Some(new_pinned_version) = Self::resolve_version_update_info(
724 &pinned_version,
725 version_payload,
726 Some(&mut sst_delta_infos),
727 ) {
728 self.refiller
729 .start_cache_refill(sst_delta_infos, pinned_version, new_pinned_version);
730 }
731 }
732
733 fn resolve_version_update_info(
734 pinned_version: &PinnedVersion,
735 version_payload: HummockVersionUpdate,
736 mut sst_delta_infos: Option<&mut Vec<SstDeltaInfo>>,
737 ) -> Option<PinnedVersion> {
738 match version_payload {
739 HummockVersionUpdate::VersionDeltas(version_deltas) => {
740 let mut version_to_apply = (**pinned_version).clone();
741 {
742 for version_delta in version_deltas {
743 assert_eq!(version_to_apply.id, version_delta.prev_id);
744 if let Some(sst_delta_infos) = &mut sst_delta_infos {
745 sst_delta_infos
746 .extend(version_to_apply.build_sst_delta_infos(&version_delta));
747 }
748
749 version_to_apply.apply_version_delta(&version_delta);
750 }
751 }
752
753 pinned_version.new_with_local_version(version_to_apply)
754 }
755 HummockVersionUpdate::PinnedVersion(version) => {
756 pinned_version.new_pin_version(*version)
757 }
758 }
759 }
760
761 fn apply_version_updates(&mut self, events: Vec<CacheRefillerEvent>) {
762 let Some(CacheRefillerEvent {
763 new_pinned_version: latest_pinned_version,
764 ..
765 }) = events.last()
766 else {
767 if cfg!(debug_assertions) {
768 panic!("empty events")
769 }
770 return;
771 };
772 if events.len() > 1 {
773 warn!(
774 count = events.len(),
775 "handle multiple version updates in batch"
776 );
777 }
778 let _timer = self
779 .metrics
780 .event_handler_on_apply_version_update
781 .start_timer();
782 self.recent_versions.rcu(|prev_recent_versions| {
783 let mut recent_versions = None;
784 for event in &events {
785 let CacheRefillerEvent {
786 new_pinned_version, ..
787 } = event;
788 recent_versions = Some(
789 recent_versions
790 .as_ref()
791 .unwrap_or(prev_recent_versions.as_ref())
792 .with_new_version(new_pinned_version.clone()),
793 );
794 }
795 recent_versions.expect("non-empty events")
796 });
797
798 {
799 self.for_each_read_version(
800 self.local_read_versions.instance_ids(),
801 |_, read_version| {
802 for CacheRefillerEvent {
803 new_pinned_version, ..
804 } in &events
805 {
806 read_version
807 .update(VersionUpdate::CommittedSnapshot(new_pinned_version.clone()))
808 }
809 },
810 );
811 }
812
813 self.version_update_notifier_tx.send_if_modified(|state| {
814 let mut modified = false;
815 for CacheRefillerEvent {
816 pinned_version,
817 new_pinned_version,
818 } in &events
819 {
820 assert_eq!(pinned_version.id(), state.id());
821 if state.id() == new_pinned_version.id() {
822 continue;
823 }
824 assert!(new_pinned_version.id() > state.id());
825 *state = new_pinned_version.clone();
826 modified = true;
827 }
828 modified
829 });
830
831 debug!("update to hummock version: {}", latest_pinned_version.id(),);
832
833 self.uploader
834 .update_pinned_version(latest_pinned_version.clone());
835 }
836}
837
838impl HummockEventHandler {
839 pub async fn start_hummock_event_handler_worker(mut self) {
840 loop {
841 tokio::select! {
842 ssts = self.uploader.next_uploaded_ssts() => {
843 self.handle_uploaded_ssts(ssts);
844 }
845 events = self.refiller.next_events() => {
846 self.apply_version_updates(events);
847 }
848 event = pin!(self.hummock_event_rx.recv()) => {
849 let Some(event) = event else { break };
850 match event {
851 HummockEvent::Shutdown => {
852 info!("event handler shutdown");
853 return;
854 },
855 event => {
856 self.handle_hummock_event(event);
857 }
858 }
859 }
860 observer_event = pin!(self.observer_event_rx.recv()) => {
861 let Some(observer_event) = observer_event else {
862 warn!("observer event stream ends. event handle shutdown");
863 return;
864 };
865 self.handle_observer_event(observer_event);
866 }
867 }
868 }
869 }
870
871 fn handle_uploaded_ssts(&mut self, ssts: Vec<Arc<StagingSstableInfo>>) {
872 let _timer = self
873 .metrics
874 .event_handler_on_upload_finish_latency
875 .start_timer();
876 self.handle_uploaded_ssts_inner(ssts);
877 }
878
879 fn handle_hummock_event(&mut self, event: HummockEvent) {
881 match event {
882 HummockEvent::BufferMayFlush => {
883 self.uploader
884 .may_flush(&self.metrics.event_handler_on_spiller);
885 }
886 HummockEvent::SyncEpoch {
887 sync_result_sender,
888 sync_table_epochs,
889 } => {
890 self.handle_sync_epoch(sync_table_epochs, sync_result_sender);
891 }
892 HummockEvent::Clear(notifier, table_ids) => {
893 self.handle_clear(notifier, table_ids);
894 }
895 HummockEvent::Shutdown => {
896 unreachable!("shutdown is handled specially")
897 }
898 HummockEvent::StartEpoch { epoch, table_ids } => {
899 self.uploader.start_epoch(epoch, table_ids);
900 }
901 HummockEvent::InitEpoch {
902 instance_id,
903 init_epoch,
904 } => {
905 let table_id = self
906 .local_read_versions
907 .get(&instance_id)
908 .expect("should exist")
909 .0;
910 self.uploader
911 .init_instance(instance_id, table_id, init_epoch);
912 }
913 HummockEvent::ImmToUploader { instance_id, imms } => {
914 assert!(
915 self.local_read_versions.contains_instance(&instance_id),
916 "add imm from non-existing read version instance: instance_id: {}, table_id {:?}",
917 instance_id,
918 imms.first().map(|(imm, _)| imm.table_id),
919 );
920 self.uploader.add_imms(instance_id, imms);
921 self.uploader
922 .may_flush(&self.metrics.event_handler_on_spiller);
923 }
924
925 HummockEvent::LocalSealEpoch {
926 next_epoch,
927 opts,
928 instance_id,
929 } => {
930 self.uploader
931 .local_seal_epoch(instance_id, next_epoch, opts);
932 }
933
934 #[cfg(any(test, feature = "test"))]
935 HummockEvent::FlushEvent(sender) => {
936 let _ = sender.send(()).inspect_err(|e| {
937 error!("unable to send flush result: {:?}", e);
938 });
939 }
940
941 HummockEvent::RegisterReadVersion {
942 table_id,
943 new_read_version_sender,
944 is_replicated,
945 vnodes,
946 } => {
947 let pinned_version = self.recent_versions.load().latest_version().clone();
948 let instance_id = self.generate_instance_id();
949 let basic_read_version = Arc::new(RwLock::new(
950 HummockReadVersion::new_with_replication_option(
951 table_id,
952 instance_id,
953 pinned_version,
954 is_replicated,
955 vnodes,
956 ),
957 ));
958
959 debug!(
960 "new read version registered: table_id: {}, instance_id: {}",
961 table_id, instance_id
962 );
963
964 {
965 self.local_read_versions.insert(
966 instance_id,
967 table_id,
968 basic_read_version.clone(),
969 );
970 let mut read_version_mapping_guard = self.read_version_mapping.write();
971
972 read_version_mapping_guard
973 .entry(table_id)
974 .or_default()
975 .insert(instance_id, basic_read_version.clone());
976 }
977
978 match new_read_version_sender.send((
979 basic_read_version,
980 LocalInstanceGuard {
981 table_id,
982 instance_id,
983 event_sender: Some(self.hummock_event_tx.clone()),
984 },
985 )) {
986 Ok(_) => {}
987 Err((_, mut guard)) => {
988 warn!(
989 "RegisterReadVersion send fail table_id {:?} instance_is {:?}",
990 table_id, instance_id
991 );
992 guard.event_sender.take().expect("sender is just set");
993 self.destroy_read_version(instance_id);
994 }
995 }
996
997 self.refiller
998 .update_streaming_table_vnodes(table_id, self.streaming_table_vnodes(table_id));
999 }
1000
1001 HummockEvent::DestroyReadVersion { instance_id } => {
1002 let table_id = self
1003 .local_read_versions
1004 .get(&instance_id)
1005 .unwrap_or_else(|| {
1006 panic!("query nonexistent instance instance_id {instance_id}")
1007 })
1008 .0;
1009 self.uploader.may_destroy_instance(instance_id);
1010 self.destroy_read_version(instance_id);
1011 self.refiller
1012 .update_streaming_table_vnodes(table_id, self.streaming_table_vnodes(table_id));
1013 }
1014 HummockEvent::GetMinUncommittedObjectId { result_tx } => {
1015 let _ = result_tx
1016 .send(self.uploader.min_uncommitted_object_id())
1017 .inspect_err(|e| {
1018 error!("unable to send get_min_uncommitted_sst_id result: {:?}", e);
1019 });
1020 }
1021 HummockEvent::GetTableCacheRefillMonitorSnapshot { result_tx } => {
1022 let _ = result_tx
1023 .send(self.refiller.table_cache_refill_monitor_snapshot())
1024 .inspect_err(|_| {
1025 error!("unable to send table cache refill monitor snapshot result");
1026 });
1027 }
1028 HummockEvent::RegisterVectorWriter {
1029 table_id,
1030 init_epoch,
1031 } => self.uploader.register_vector_writer(table_id, init_epoch),
1032 HummockEvent::VectorWriterSealEpoch {
1033 table_id,
1034 next_epoch,
1035 add,
1036 } => {
1037 self.uploader
1038 .vector_writer_seal_epoch(table_id, next_epoch, add);
1039 }
1040 HummockEvent::DropVectorWriter { table_id } => {
1041 self.uploader.drop_vector_writer(table_id);
1042 }
1043 }
1044 }
1045
1046 fn destroy_read_version(&mut self, instance_id: LocalInstanceId) {
1047 {
1048 {
1049 debug!("read version deregister: instance_id: {}", instance_id);
1050 let (table_id, _) = self.local_read_versions.remove(instance_id);
1051 let mut read_version_mapping_guard = self.read_version_mapping.write();
1052 let entry = read_version_mapping_guard
1053 .get_mut(&table_id)
1054 .unwrap_or_else(|| {
1055 panic!(
1056 "DestroyHummockInstance table_id {} instance_id {} fail",
1057 table_id, instance_id
1058 )
1059 });
1060 entry.remove(&instance_id).unwrap_or_else(|| {
1061 panic!(
1062 "DestroyHummockInstance nonexistent instance table_id {} instance_id {}",
1063 table_id, instance_id
1064 )
1065 });
1066 if entry.is_empty() {
1067 read_version_mapping_guard.remove(&table_id);
1068 }
1069 }
1070 }
1071 }
1072
1073 fn generate_instance_id(&mut self) -> LocalInstanceId {
1074 self.last_instance_id += 1;
1075 self.last_instance_id
1076 }
1077}
1078
1079pub(super) fn send_sync_result(
1080 sender: oneshot::Sender<HummockResult<SyncedData>>,
1081 result: HummockResult<SyncedData>,
1082) {
1083 let _ = sender.send(result).inspect_err(|e| {
1084 error!("unable to send sync result. Err: {:?}", e);
1085 });
1086}
1087
1088impl SyncedData {
1089 pub fn into_sync_result(self) -> SyncResult {
1090 {
1091 let SyncedData {
1092 uploaded_ssts,
1093 table_watermarks,
1094 vector_index_adds,
1095 } = self;
1096 let mut sync_size = 0;
1097 let mut uncommitted_ssts = Vec::new();
1098 let mut old_value_ssts = Vec::new();
1099 for sst in uploaded_ssts {
1102 sync_size += sst.imm_size();
1103 uncommitted_ssts.extend(sst.sstable_infos().iter().cloned());
1104 old_value_ssts.extend(sst.old_value_sstable_infos().iter().cloned());
1105 }
1106 SyncResult {
1107 sync_size,
1108 uncommitted_ssts,
1109 table_watermarks,
1110 old_value_ssts,
1111 vector_index_adds,
1112 }
1113 }
1114 }
1115}
1116
1117#[cfg(test)]
1118mod tests {
1119 use std::collections::{HashMap, HashSet};
1120 use std::future::poll_fn;
1121 use std::sync::Arc;
1122 use std::task::Poll;
1123
1124 use futures::FutureExt;
1125 use parking_lot::Mutex;
1126 use risingwave_common::bitmap::Bitmap;
1127 use risingwave_common::catalog::TableId;
1128 use risingwave_common::config::Role;
1129 use risingwave_common::config::streaming::CacheRefillPolicy;
1130 use risingwave_common::hash::VirtualNode;
1131 use risingwave_common::util::epoch::{EpochExt, test_epoch};
1132 use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
1133 use risingwave_hummock_sdk::version::HummockVersion;
1134 use risingwave_pb::hummock::{PbHummockVersion, StateTableInfo};
1135 use risingwave_pb::meta::serving_table_vnode_mappings::PbServingTableVnodeMapping;
1136 use risingwave_pb::meta::table_cache_refill_policies::PbTableCacheRefillPolicy;
1137 use risingwave_pb::meta::table_cache_refill_policies::table_cache_refill_policy::PbCacheRefillPolicy;
1138 use risingwave_pb::meta::{
1139 PbServingTableVnodeMappings, PbTableRefillRuntimeConfig, TableCacheRefillPolicies,
1140 };
1141 use tokio::spawn;
1142 use tokio::sync::mpsc::unbounded_channel;
1143 use tokio::sync::oneshot;
1144
1145 use crate::hummock::HummockError;
1146 use crate::hummock::event_handler::hummock_event_handler::BufferTracker;
1147 use crate::hummock::event_handler::refiller::{CacheRefillConfig, CacheRefiller};
1148 use crate::hummock::event_handler::uploader::UploadTaskOutput;
1149 use crate::hummock::event_handler::uploader::test_utils::{
1150 TEST_TABLE_ID, gen_imm_inner, gen_imm_with_unlimit,
1151 prepare_uploader_order_test_spawn_task_fn,
1152 };
1153 use crate::hummock::event_handler::{
1154 HummockEvent, HummockEventHandler, HummockReadVersionRef, LocalInstanceGuard,
1155 };
1156 use crate::hummock::iterator::test_utils::mock_sstable_store;
1157 use crate::hummock::local_version::pinned_version::PinnedVersion;
1158 use crate::hummock::local_version::recent_versions::RecentVersions;
1159 use crate::hummock::test_utils::default_opts_for_test;
1160 use crate::mem_table::ImmutableMemtable;
1161 use crate::monitor::HummockStateStoreMetrics;
1162 use crate::store::SealCurrentEpochOptions;
1163
1164 #[tokio::test]
1165 async fn test_runtime_config_fields_are_independent_full_replacements() {
1166 let old_table_id = TableId::new(233);
1167 let internal_table_id = TableId::new(234);
1168 let new_table_id = TableId::new(235);
1169 let old_vnodes = Bitmap::ones(VirtualNode::COUNT_FOR_TEST);
1170 let new_vnodes = Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 0..8);
1171 let storage_opt = default_opts_for_test();
1172 let mut refiller = CacheRefiller::new(
1173 Role::Both,
1174 CacheRefillConfig::from_storage_opts(&storage_opt),
1175 mock_sstable_store().await,
1176 CacheRefiller::default_spawn_refill_task(),
1177 );
1178
1179 HummockEventHandler::apply_table_refill_runtime_config(
1180 &mut refiller,
1181 PbTableRefillRuntimeConfig {
1182 table_cache_refill_policies: Some(TableCacheRefillPolicies {
1183 table_policies: vec![PbTableCacheRefillPolicy {
1184 table_id: old_table_id.as_raw_id(),
1185 policy: PbCacheRefillPolicy::Serving as i32,
1186 }],
1187 internal_table_policies: vec![PbTableCacheRefillPolicy {
1188 table_id: internal_table_id.as_raw_id(),
1189 policy: PbCacheRefillPolicy::Both as i32,
1190 }],
1191 }),
1192 serving_table_vnode_mappings: Some(PbServingTableVnodeMappings {
1193 mappings: vec![PbServingTableVnodeMapping {
1194 table_id: old_table_id.as_raw_id(),
1195 bitmap: Some(old_vnodes.to_protobuf()),
1196 }],
1197 }),
1198 ..Default::default()
1199 },
1200 );
1201 let snapshot = refiller.table_cache_refill_monitor_snapshot();
1202 assert_eq!(
1203 snapshot.policies,
1204 HashMap::from([
1205 (old_table_id, CacheRefillPolicy::Serving),
1206 (internal_table_id, CacheRefillPolicy::Both),
1207 ])
1208 );
1209 assert_eq!(
1210 snapshot.serving_table_vnode_mapping,
1211 HashMap::from([(old_table_id, old_vnodes)])
1212 );
1213
1214 HummockEventHandler::apply_table_refill_runtime_config(
1215 &mut refiller,
1216 PbTableRefillRuntimeConfig {
1217 serving_table_vnode_mappings: Some(PbServingTableVnodeMappings {
1218 mappings: vec![PbServingTableVnodeMapping {
1219 table_id: new_table_id.as_raw_id(),
1220 bitmap: Some(new_vnodes.to_protobuf()),
1221 }],
1222 }),
1223 ..Default::default()
1224 },
1225 );
1226 let snapshot = refiller.table_cache_refill_monitor_snapshot();
1227 assert_eq!(
1228 snapshot.policies,
1229 HashMap::from([
1230 (old_table_id, CacheRefillPolicy::Serving),
1231 (internal_table_id, CacheRefillPolicy::Both),
1232 ])
1233 );
1234 assert_eq!(
1235 snapshot.serving_table_vnode_mapping,
1236 HashMap::from([(new_table_id, new_vnodes.clone())])
1237 );
1238
1239 HummockEventHandler::apply_table_refill_runtime_config(
1240 &mut refiller,
1241 PbTableRefillRuntimeConfig {
1242 table_cache_refill_policies: Some(TableCacheRefillPolicies::default()),
1243 ..Default::default()
1244 },
1245 );
1246 let snapshot = refiller.table_cache_refill_monitor_snapshot();
1247 assert!(snapshot.policies.is_empty());
1248 assert_eq!(
1249 snapshot.serving_table_vnode_mapping,
1250 HashMap::from([(new_table_id, new_vnodes)])
1251 );
1252
1253 HummockEventHandler::apply_table_refill_runtime_config(
1254 &mut refiller,
1255 PbTableRefillRuntimeConfig {
1256 serving_table_vnode_mappings: Some(PbServingTableVnodeMappings::default()),
1257 ..Default::default()
1258 },
1259 );
1260 let snapshot = refiller.table_cache_refill_monitor_snapshot();
1261 assert!(snapshot.policies.is_empty());
1262 assert!(snapshot.serving_table_vnode_mapping.is_empty());
1263 }
1264
1265 #[test]
1266 fn test_runtime_config_ignores_invalid_entries() {
1267 let table_id = TableId::new(233);
1268 assert!(
1269 super::table_cache_refill_policy_from_protobuf(PbTableCacheRefillPolicy {
1270 table_id: table_id.as_raw_id(),
1271 policy: i32::MAX,
1272 })
1273 .is_none()
1274 );
1275
1276 assert!(
1277 super::serving_table_vnode_mappings_to_map(PbServingTableVnodeMappings {
1278 mappings: vec![PbServingTableVnodeMapping {
1279 table_id: table_id.as_raw_id(),
1280 bitmap: None,
1281 }],
1282 })
1283 .is_empty()
1284 );
1285 }
1286
1287 #[tokio::test]
1288 async fn test_read_version_events_update_refiller_streaming_mapping() {
1289 fn assert_vnodes(vnodes: &Bitmap, expected: std::ops::Range<usize>) {
1290 assert_eq!(vnodes.count_ones(), expected.len());
1291 assert!(expected.into_iter().all(|vnode| vnodes.is_set(vnode)));
1292 }
1293
1294 let table_id = TableId::new(233);
1295 let other_table_id = TableId::new(234);
1296 let (_observer_event_tx, observer_event_rx) = unbounded_channel();
1297 let metrics = Arc::new(HummockStateStoreMetrics::unused());
1298 let storage_opt = default_opts_for_test();
1299 let (spawn_upload_task, _) = prepare_uploader_order_test_spawn_task_fn(false);
1300 let pinned_version = PinnedVersion::new(
1301 HummockVersion::from_rpc_protobuf(&PbHummockVersion {
1302 id: 1.into(),
1303 ..Default::default()
1304 }),
1305 unbounded_channel().0,
1306 );
1307 let mut event_handler = HummockEventHandler::new_inner(
1308 Role::Streaming,
1309 observer_event_rx,
1310 mock_sstable_store().await,
1311 metrics.clone(),
1312 CacheRefillConfig::from_storage_opts(&storage_opt),
1313 RecentVersions::new(pinned_version, 10, metrics.clone()),
1314 BufferTracker::from_storage_opts(&storage_opt, &metrics),
1315 spawn_upload_task,
1316 CacheRefiller::default_spawn_refill_task(),
1317 );
1318
1319 let (tx, rx) = oneshot::channel();
1320 event_handler.handle_hummock_event(HummockEvent::RegisterReadVersion {
1321 table_id,
1322 new_read_version_sender: tx,
1323 is_replicated: false,
1324 vnodes: Arc::new(Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 0..8)),
1325 });
1326 let (_read_version_1, mut guard_1) = rx.await.unwrap();
1327
1328 let (tx, rx) = oneshot::channel();
1329 event_handler.handle_hummock_event(HummockEvent::RegisterReadVersion {
1330 table_id,
1331 new_read_version_sender: tx,
1332 is_replicated: false,
1333 vnodes: Arc::new(Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 8..16)),
1334 });
1335 let (_read_version_2, mut guard_2) = rx.await.unwrap();
1336
1337 let (tx, rx) = oneshot::channel();
1338 event_handler.handle_hummock_event(HummockEvent::RegisterReadVersion {
1339 table_id: other_table_id,
1340 new_read_version_sender: tx,
1341 is_replicated: false,
1342 vnodes: Arc::new(Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 32..40)),
1343 });
1344 let (_other_read_version, mut other_guard) = rx.await.unwrap();
1345
1346 let snapshot = event_handler.refiller.table_cache_refill_monitor_snapshot();
1347 let vnodes = snapshot
1348 .streaming_table_vnode_mapping
1349 .get(&table_id)
1350 .unwrap();
1351 assert_vnodes(vnodes, 0..16);
1352 assert_vnodes(
1353 snapshot
1354 .streaming_table_vnode_mapping
1355 .get(&other_table_id)
1356 .unwrap(),
1357 32..40,
1358 );
1359
1360 guard_1.event_sender.take();
1361 event_handler.handle_hummock_event(HummockEvent::DestroyReadVersion {
1362 instance_id: guard_1.instance_id,
1363 });
1364 let snapshot = event_handler.refiller.table_cache_refill_monitor_snapshot();
1365 assert_vnodes(
1366 snapshot
1367 .streaming_table_vnode_mapping
1368 .get(&table_id)
1369 .unwrap(),
1370 8..16,
1371 );
1372
1373 guard_2.event_sender.take();
1374 event_handler.handle_hummock_event(HummockEvent::DestroyReadVersion {
1375 instance_id: guard_2.instance_id,
1376 });
1377 let snapshot = event_handler.refiller.table_cache_refill_monitor_snapshot();
1378 assert!(
1379 !snapshot
1380 .streaming_table_vnode_mapping
1381 .contains_key(&table_id)
1382 );
1383 assert_vnodes(
1384 snapshot
1385 .streaming_table_vnode_mapping
1386 .get(&other_table_id)
1387 .unwrap(),
1388 32..40,
1389 );
1390
1391 other_guard.event_sender.take();
1392 event_handler.handle_hummock_event(HummockEvent::DestroyReadVersion {
1393 instance_id: other_guard.instance_id,
1394 });
1395 }
1396
1397 #[tokio::test]
1398 async fn test_old_epoch_sync_fail() {
1399 let epoch0 = test_epoch(233);
1400
1401 let initial_version = PinnedVersion::new(
1402 HummockVersion::from_rpc_protobuf(&PbHummockVersion {
1403 id: 1.into(),
1404 state_table_info: HashMap::from_iter([(
1405 TEST_TABLE_ID,
1406 StateTableInfo {
1407 committed_epoch: epoch0,
1408 compaction_group_id: StaticCompactionGroupId::StateDefault,
1409 },
1410 )]),
1411 ..Default::default()
1412 }),
1413 unbounded_channel().0,
1414 );
1415
1416 let (_observer_event_tx, observer_event_rx) = unbounded_channel();
1417
1418 let epoch1 = epoch0.next_epoch();
1419 let epoch2 = epoch1.next_epoch();
1420 let (tx, rx) = oneshot::channel();
1421 let rx = Arc::new(Mutex::new(Some(rx)));
1422
1423 let storage_opt = default_opts_for_test();
1424 let metrics = Arc::new(HummockStateStoreMetrics::unused());
1425
1426 let event_handler = HummockEventHandler::new_inner(
1427 Role::None,
1428 observer_event_rx,
1429 mock_sstable_store().await,
1430 metrics.clone(),
1431 CacheRefillConfig::from_storage_opts(&storage_opt),
1432 RecentVersions::new(initial_version.clone(), 10, metrics.clone()),
1433 BufferTracker::from_storage_opts(&storage_opt, &metrics),
1434 Arc::new(move |_, info| {
1435 assert_eq!(info.epochs.len(), 1);
1436 let epoch = info.epochs[0];
1437 match epoch {
1438 epoch if epoch == epoch1 => {
1439 let rx = rx.lock().take().unwrap();
1440 spawn(async move {
1441 rx.await.unwrap();
1442 Err(HummockError::other("fail"))
1443 })
1444 }
1445 epoch if epoch == epoch2 => spawn(async move {
1446 Ok(UploadTaskOutput {
1447 new_value_ssts: vec![],
1448 old_value_ssts: vec![],
1449 wait_poll_timer: None,
1450 })
1451 }),
1452 _ => unreachable!(),
1453 }
1454 }),
1455 CacheRefiller::default_spawn_refill_task(),
1456 );
1457
1458 let event_tx = event_handler.event_sender();
1459
1460 let send_event = |event| event_tx.send(event).unwrap();
1461
1462 let join_handle = spawn(event_handler.start_hummock_event_handler_worker());
1463
1464 let (read_version, guard) = {
1465 let (tx, rx) = oneshot::channel();
1466 send_event(HummockEvent::RegisterReadVersion {
1467 table_id: TEST_TABLE_ID,
1468 new_read_version_sender: tx,
1469 is_replicated: false,
1470 vnodes: Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1471 });
1472 rx.await.unwrap()
1473 };
1474
1475 send_event(HummockEvent::StartEpoch {
1476 epoch: epoch1,
1477 table_ids: HashSet::from_iter([TEST_TABLE_ID]),
1478 });
1479
1480 read_version.write().init();
1481
1482 send_event(HummockEvent::InitEpoch {
1483 instance_id: guard.instance_id,
1484 init_epoch: epoch1,
1485 });
1486
1487 let (imm1, tracker1) = gen_imm_with_unlimit(epoch1);
1488 read_version.write().add_pending_imm(imm1.clone(), tracker1);
1489
1490 send_event(HummockEvent::ImmToUploader {
1491 instance_id: guard.instance_id,
1492 imms: read_version.write().start_upload_pending_imms(),
1493 });
1494
1495 send_event(HummockEvent::StartEpoch {
1496 epoch: epoch2,
1497 table_ids: HashSet::from_iter([TEST_TABLE_ID]),
1498 });
1499
1500 send_event(HummockEvent::LocalSealEpoch {
1501 instance_id: guard.instance_id,
1502 next_epoch: epoch2,
1503 opts: SealCurrentEpochOptions::for_test(),
1504 });
1505
1506 {
1507 let (imm2, tracker2) = gen_imm_with_unlimit(epoch2);
1508 let mut read_version = read_version.write();
1509 read_version.add_pending_imm(imm2, tracker2);
1510
1511 send_event(HummockEvent::ImmToUploader {
1512 instance_id: guard.instance_id,
1513 imms: read_version.start_upload_pending_imms(),
1514 });
1515 }
1516
1517 let epoch3 = epoch2.next_epoch();
1518 send_event(HummockEvent::StartEpoch {
1519 epoch: epoch3,
1520 table_ids: HashSet::from_iter([TEST_TABLE_ID]),
1521 });
1522 send_event(HummockEvent::LocalSealEpoch {
1523 instance_id: guard.instance_id,
1524 next_epoch: epoch3,
1525 opts: SealCurrentEpochOptions::for_test(),
1526 });
1527
1528 let (tx1, mut rx1) = oneshot::channel();
1529 send_event(HummockEvent::SyncEpoch {
1530 sync_result_sender: tx1,
1531 sync_table_epochs: vec![(epoch1, HashSet::from_iter([TEST_TABLE_ID]))],
1532 });
1533 assert!(poll_fn(|cx| Poll::Ready(rx1.poll_unpin(cx).is_pending())).await);
1534 let (tx2, mut rx2) = oneshot::channel();
1535 send_event(HummockEvent::SyncEpoch {
1536 sync_result_sender: tx2,
1537 sync_table_epochs: vec![(epoch2, HashSet::from_iter([TEST_TABLE_ID]))],
1538 });
1539 assert!(poll_fn(|cx| Poll::Ready(rx2.poll_unpin(cx).is_pending())).await);
1540
1541 tx.send(()).unwrap();
1542 rx1.await.unwrap().unwrap_err();
1543 rx2.await.unwrap().unwrap_err();
1544
1545 send_event(HummockEvent::Shutdown);
1546 join_handle.await.unwrap();
1547 }
1548
1549 #[tokio::test]
1550 async fn test_clear_one_table_preserves_other_table_uploads() {
1551 let table_id1 = TableId::new(1);
1552 let table_id2 = TableId::new(2);
1553 let epoch0 = test_epoch(233);
1554
1555 let initial_version = PinnedVersion::new(
1556 HummockVersion::from_rpc_protobuf(&PbHummockVersion {
1557 id: 1.into(),
1558 state_table_info: HashMap::from_iter([
1559 (
1560 table_id1,
1561 StateTableInfo {
1562 committed_epoch: epoch0,
1563 compaction_group_id: StaticCompactionGroupId::StateDefault,
1564 },
1565 ),
1566 (
1567 table_id2,
1568 StateTableInfo {
1569 committed_epoch: epoch0,
1570 compaction_group_id: StaticCompactionGroupId::StateDefault,
1571 },
1572 ),
1573 ]),
1574 ..Default::default()
1575 }),
1576 unbounded_channel().0,
1577 );
1578
1579 let (_observer_event_tx, observer_event_rx) = unbounded_channel();
1580
1581 let epoch1 = epoch0.next_epoch();
1582 let epoch2 = epoch1.next_epoch();
1583 let epoch3 = epoch2.next_epoch();
1584
1585 let imm_size = gen_imm_inner(TEST_TABLE_ID, epoch1, 0).size();
1586
1587 let buffer_tracker = BufferTracker::for_test_with_config(imm_size * 2 - 1, 1);
1589 let memory_limiter = buffer_tracker.get_memory_limiter().clone();
1590
1591 let gen_imm = |table_id, epoch, spill_offset| {
1592 let imm = gen_imm_inner(table_id, epoch, spill_offset);
1593 assert_eq!(imm.size(), imm_size);
1594 imm
1595 };
1596 let imm1_1 = gen_imm(table_id1, epoch1, 0);
1597 let imm1_2_1 = gen_imm(table_id1, epoch2, 0);
1598
1599 let storage_opt = default_opts_for_test();
1600 let metrics = Arc::new(HummockStateStoreMetrics::unused());
1601
1602 let (spawn_task, new_task_notifier) = prepare_uploader_order_test_spawn_task_fn(false);
1603
1604 let event_handler = HummockEventHandler::new_inner(
1605 Role::None,
1606 observer_event_rx,
1607 mock_sstable_store().await,
1608 metrics.clone(),
1609 CacheRefillConfig::from_storage_opts(&storage_opt),
1610 RecentVersions::new(initial_version.clone(), 10, metrics.clone()),
1611 buffer_tracker,
1612 spawn_task,
1613 CacheRefiller::default_spawn_refill_task(),
1614 );
1615
1616 let event_tx = event_handler.event_sender();
1617
1618 let send_event = |event| event_tx.send(event).unwrap();
1619 let flush_event = || async {
1620 let (tx, rx) = oneshot::channel();
1621 send_event(HummockEvent::FlushEvent(tx));
1622 rx.await.unwrap();
1623 };
1624 let start_epoch = |table_id, epoch| {
1625 send_event(HummockEvent::StartEpoch {
1626 epoch,
1627 table_ids: HashSet::from_iter([table_id]),
1628 })
1629 };
1630 let init_epoch = |instance: &LocalInstanceGuard, init_epoch| {
1631 send_event(HummockEvent::InitEpoch {
1632 instance_id: instance.instance_id,
1633 init_epoch,
1634 })
1635 };
1636 let event_tx_clone = event_tx.clone();
1637 let write_imm = {
1638 let memory_limiter = memory_limiter.clone();
1639 move |read_version: &HummockReadVersionRef,
1640 instance: &LocalInstanceGuard,
1641 imm: &ImmutableMemtable| {
1642 let memory_limiter = memory_limiter.clone();
1643 let event_tx = event_tx_clone.clone();
1644 let read_version = read_version.clone();
1645 let imm = imm.clone();
1646 let instance_id = instance.instance_id;
1647 async move {
1648 let tracker = memory_limiter.require_memory(imm.size() as _).await;
1649 let mut read_version = read_version.write();
1650 read_version.add_pending_imm(imm.clone(), tracker);
1651
1652 event_tx
1653 .send(HummockEvent::ImmToUploader {
1654 instance_id,
1655 imms: read_version.start_upload_pending_imms(),
1656 })
1657 .unwrap();
1658 }
1659 }
1660 };
1661 let seal_epoch = |instance: &LocalInstanceGuard, next_epoch| {
1662 send_event(HummockEvent::LocalSealEpoch {
1663 instance_id: instance.instance_id,
1664 next_epoch,
1665 opts: SealCurrentEpochOptions::for_test(),
1666 })
1667 };
1668 let sync_epoch = |table_id, new_sync_epoch| {
1669 let (tx, rx) = oneshot::channel();
1670 send_event(HummockEvent::SyncEpoch {
1671 sync_result_sender: tx,
1672 sync_table_epochs: vec![(new_sync_epoch, HashSet::from_iter([table_id]))],
1673 });
1674 rx
1675 };
1676
1677 let join_handle = spawn(event_handler.start_hummock_event_handler_worker());
1678
1679 let (read_version1, guard1) = {
1680 let (tx, rx) = oneshot::channel();
1681 send_event(HummockEvent::RegisterReadVersion {
1682 table_id: table_id1,
1683 new_read_version_sender: tx,
1684 is_replicated: false,
1685 vnodes: Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1686 });
1687 rx.await.unwrap()
1688 };
1689
1690 let (read_version2, guard2) = {
1691 let (tx, rx) = oneshot::channel();
1692 send_event(HummockEvent::RegisterReadVersion {
1693 table_id: table_id2,
1694 new_read_version_sender: tx,
1695 is_replicated: false,
1696 vnodes: Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1697 });
1698 rx.await.unwrap()
1699 };
1700
1701 let (task1_1_finish_tx, task1_1_rx) = {
1703 start_epoch(table_id1, epoch1);
1704
1705 read_version1.write().init();
1706 init_epoch(&guard1, epoch1);
1707
1708 write_imm(&read_version1, &guard1, &imm1_1).await;
1709
1710 start_epoch(table_id1, epoch2);
1711
1712 seal_epoch(&guard1, epoch2);
1713
1714 let (wait_task_start, task_finish_tx) = new_task_notifier(HashMap::from_iter([(
1715 guard1.instance_id,
1716 vec![imm1_1.batch_id()],
1717 )]));
1718
1719 let mut rx = sync_epoch(table_id1, epoch1);
1720 wait_task_start.await;
1721 assert!(poll_fn(|cx| Poll::Ready(rx.poll_unpin(cx).is_pending())).await);
1722
1723 write_imm(&read_version1, &guard1, &imm1_2_1).await;
1724 flush_event().await;
1725
1726 (task_finish_tx, rx)
1727 };
1728 let (task1_2_finish_tx, _finish_txs) = {
1733 let mut finish_txs = vec![];
1734 let imm2_1_1 = gen_imm(table_id2, epoch1, 0);
1735 start_epoch(table_id2, epoch1);
1736 read_version2.write().init();
1737 init_epoch(&guard2, epoch1);
1738 let (wait_task_start, task1_2_finish_tx) = new_task_notifier(HashMap::from_iter([(
1739 guard1.instance_id,
1740 vec![imm1_2_1.batch_id()],
1741 )]));
1742 write_imm(&read_version2, &guard2, &imm2_1_1).await;
1743 wait_task_start.await;
1744
1745 let imm2_1_2 = gen_imm(table_id2, epoch1, 1);
1746 let (wait_task_start, finish_tx) = new_task_notifier(HashMap::from_iter([(
1747 guard2.instance_id,
1748 vec![imm2_1_2.batch_id(), imm2_1_1.batch_id()],
1749 )]));
1750 finish_txs.push(finish_tx);
1751 write_imm(&read_version2, &guard2, &imm2_1_2).await;
1752 wait_task_start.await;
1753
1754 let imm2_1_3 = gen_imm(table_id2, epoch1, 2);
1755 write_imm(&read_version2, &guard2, &imm2_1_3).await;
1756 start_epoch(table_id2, epoch2);
1757 seal_epoch(&guard2, epoch2);
1758 let (wait_task_start, finish_tx) = new_task_notifier(HashMap::from_iter([(
1759 guard2.instance_id,
1760 vec![imm2_1_3.batch_id()],
1761 )]));
1762 finish_txs.push(finish_tx);
1763 let _sync_rx = sync_epoch(table_id2, epoch1);
1764 wait_task_start.await;
1765
1766 let imm2_2_1 = gen_imm(table_id2, epoch2, 0);
1767 write_imm(&read_version2, &guard2, &imm2_2_1).await;
1768 flush_event().await;
1769 let imm2_2_2 = gen_imm(table_id2, epoch2, 1);
1770 write_imm(&read_version2, &guard2, &imm2_2_2).await;
1771 let (wait_task_start, finish_tx) = new_task_notifier(HashMap::from_iter([(
1772 guard2.instance_id,
1773 vec![imm2_2_2.batch_id(), imm2_2_1.batch_id()],
1774 )]));
1775 finish_txs.push(finish_tx);
1776 wait_task_start.await;
1777
1778 let imm2_2_3 = gen_imm(table_id2, epoch2, 2);
1779 write_imm(&read_version2, &guard2, &imm2_2_3).await;
1780
1781 drop(guard2);
1789 let (clear_tx, clear_rx) = oneshot::channel();
1790 send_event(HummockEvent::Clear(
1791 clear_tx,
1792 Some(HashSet::from_iter([table_id2])),
1793 ));
1794 clear_rx.await.unwrap();
1795 (task1_2_finish_tx, finish_txs)
1796 };
1797
1798 let imm1_2_2 = gen_imm(table_id1, epoch2, 1);
1799 write_imm(&read_version1, &guard1, &imm1_2_2).await;
1800 start_epoch(table_id1, epoch3);
1801 seal_epoch(&guard1, epoch3);
1802
1803 let (tx2, mut sync_rx2) = oneshot::channel();
1804 let (wait_task_start, task1_2_2_finish_tx) = new_task_notifier(HashMap::from_iter([(
1805 guard1.instance_id,
1806 vec![imm1_2_2.batch_id()],
1807 )]));
1808 send_event(HummockEvent::SyncEpoch {
1809 sync_result_sender: tx2,
1810 sync_table_epochs: vec![(epoch2, HashSet::from_iter([table_id1]))],
1811 });
1812 wait_task_start.await;
1813 assert!(poll_fn(|cx| Poll::Ready(sync_rx2.poll_unpin(cx).is_pending())).await);
1814
1815 task1_1_finish_tx.send(()).unwrap();
1816 let sync_data1 = task1_1_rx.await.unwrap().unwrap();
1817 assert!(!sync_data1.uploaded_ssts.is_empty());
1818 assert!(
1819 sync_data1
1820 .uploaded_ssts
1821 .iter()
1822 .all(|sst| sst.epochs().as_slice() == [epoch1])
1823 );
1824 task1_2_finish_tx.send(()).unwrap();
1825 assert!(poll_fn(|cx| Poll::Ready(sync_rx2.poll_unpin(cx).is_pending())).await);
1826 task1_2_2_finish_tx.send(()).unwrap();
1827 let sync_data2 = sync_rx2.await.unwrap().unwrap();
1828 assert!(!sync_data2.uploaded_ssts.is_empty());
1829 assert!(
1830 sync_data2
1831 .uploaded_ssts
1832 .iter()
1833 .all(|sst| sst.epochs().as_slice() == [epoch2])
1834 );
1835
1836 send_event(HummockEvent::Shutdown);
1837 join_handle.await.unwrap();
1838 }
1839}