Skip to main content

risingwave_meta/hummock/manager/
mod.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::{BTreeMap, HashMap};
16use std::ops::{Deref, DerefMut};
17use std::sync::Arc;
18use std::sync::atomic::AtomicBool;
19
20use anyhow::anyhow;
21use bytes::Bytes;
22use futures::FutureExt;
23use itertools::Itertools;
24use parking_lot::lock_api::RwLock;
25use risingwave_common::catalog::{TableId, TableOption};
26use risingwave_common::monitor::MonitoredRwLock;
27use risingwave_common::system_param::reader::SystemParamsRead;
28use risingwave_hummock_sdk::change_log::TableChangeLog;
29use risingwave_hummock_sdk::version::{HummockVersion, HummockVersionDelta};
30use risingwave_hummock_sdk::{
31    CompactionGroupId, HummockCompactionTaskId, HummockContextId, HummockVersionId,
32    version_archive_dir, version_checkpoint_path,
33};
34use risingwave_meta_model::{
35    compaction_status, compaction_task, hummock_pinned_version, hummock_version_delta,
36    hummock_version_stats,
37};
38use risingwave_pb::hummock::compact_task::TaskStatus;
39use risingwave_pb::hummock::{
40    HummockVersionStats, PbCompactionGroupInfo, SubscribeCompactionEventRequest,
41};
42use table_write_throughput_statistic::TableWriteThroughputStatisticManager;
43use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel};
44use tokio::sync::{Mutex, Semaphore, oneshot};
45use tonic::Streaming;
46
47use crate::MetaResult;
48use crate::hummock::CompactorManagerRef;
49use crate::hummock::compaction::CompactStatus;
50use crate::hummock::error::Result;
51use crate::hummock::manager::checkpoint::HummockVersionCheckpoint;
52use crate::hummock::manager::context::ContextInfo;
53use crate::hummock::manager::gc::{FullGcState, GcManager};
54use crate::hummock::manager::sequence::PrefetchedSequence;
55use crate::hummock::model::ext::{compaction_task_model_to_assignment, to_table_change_log};
56use crate::manager::{MetaSrvEnv, MetadataManager};
57use crate::model::{ClusterId, MetadataModelError};
58use crate::rpc::metrics::MetaMetrics;
59
60mod context;
61mod gc;
62mod tests;
63mod versioning;
64pub use context::HummockVersionSafePoint;
65use versioning::*;
66pub(crate) mod checkpoint;
67mod commit_epoch;
68mod compaction;
69pub mod sequence;
70pub mod table_write_throughput_statistic;
71pub mod time_travel;
72mod timer_task;
73mod transaction;
74mod utils;
75mod worker;
76
77pub use commit_epoch::{CommitEpochInfo, NewTableFragmentInfo};
78pub use compaction::compaction_event_loop::*;
79use compaction::*;
80pub use compaction::{GroupState, GroupStateValidator, ManualCompactionTriggerResult};
81pub(crate) use utils::*;
82
83struct TableCommittedEpochNotifiers {
84    txs: HashMap<TableId, Vec<UnboundedSender<u64>>>,
85}
86
87impl TableCommittedEpochNotifiers {
88    fn notify_deltas(&mut self, deltas: &[HummockVersionDelta]) {
89        self.txs.retain(|table_id, txs| {
90            let mut is_dropped = false;
91            let mut committed_epoch = None;
92            for delta in deltas {
93                if delta.removed_table_ids.contains(table_id) {
94                    is_dropped = true;
95                    break;
96                }
97                if let Some(info) = delta.state_table_info_delta.get(table_id) {
98                    committed_epoch = Some(info.committed_epoch);
99                }
100            }
101            if is_dropped {
102                false
103            } else if let Some(committed_epoch) = committed_epoch {
104                txs.retain(|tx| tx.send(committed_epoch).is_ok());
105                !txs.is_empty()
106            } else {
107                true
108            }
109        })
110    }
111}
112
113#[derive(Clone, Debug)]
114struct CompactionTaskReportResult {
115    task_id: HummockCompactionTaskId,
116
117    task_status: TaskStatus,
118    reported: bool,
119}
120
121struct CompactionTaskReportNotifiers {
122    txs: HashMap<HummockCompactionTaskId, Vec<oneshot::Sender<CompactionTaskReportResult>>>,
123}
124
125impl CompactionTaskReportNotifiers {
126    fn register(
127        &mut self,
128        task_id: HummockCompactionTaskId,
129        tx: oneshot::Sender<CompactionTaskReportResult>,
130    ) {
131        self.txs.entry(task_id).or_default().push(tx);
132    }
133
134    fn remove(&mut self, task_id: HummockCompactionTaskId) {
135        self.txs.remove(&task_id);
136    }
137
138    fn notify(&mut self, result: CompactionTaskReportResult) {
139        if let Some(txs) = self.txs.remove(&result.task_id) {
140            for tx in txs {
141                let _ = tx.send(result.clone());
142            }
143        }
144    }
145}
146// Update to states are performed as follow:
147// - Initialize ValTransaction for the meta state to update
148// - Make changes on the ValTransaction.
149// - Call `commit_multi_var` to commit the changes via meta store transaction. If transaction
150//   succeeds, the in-mem state will be updated by the way.
151pub struct HummockManager {
152    pub env: MetaSrvEnv,
153
154    metadata_manager: MetadataManager,
155    /// Lock order: `compaction`, `versioning`, `compaction_group_manager`, `context_info`
156    /// - Lock `compaction` first, then `versioning`, then `compaction_group_manager` and finally `context_info`.
157    /// - This order should be strictly followed to prevent deadlock.
158    compaction: MonitoredRwLock<Compaction>,
159    versioning: MonitoredRwLock<Versioning>,
160    /// `CompactionGroupManager` manages compaction configs for compaction groups.
161    compaction_group_manager: MonitoredRwLock<CompactionGroupManager>,
162    context_info: MonitoredRwLock<ContextInfo>,
163
164    pub metrics: Arc<MetaMetrics>,
165
166    pub compactor_manager: CompactorManagerRef,
167    pub iceberg_compactor_manager: Arc<IcebergCompactorManager>,
168    event_sender: HummockManagerEventSender,
169    object_store: ObjectStoreRef,
170    version_checkpoint_path: String,
171    version_archive_dir: String,
172    pause_version_checkpoint: AtomicBool,
173    table_write_throughput_statistic_manager:
174        parking_lot::RwLock<TableWriteThroughputStatisticManager>,
175    table_committed_epoch_notifiers: parking_lot::Mutex<TableCommittedEpochNotifiers>,
176    compaction_task_report_notifiers: parking_lot::Mutex<CompactionTaskReportNotifiers>,
177    version_stat_tx: UnboundedSender<Arc<HummockVersion>>,
178
179    // for compactor
180    // `compactor_streams_change_tx` is used to pass the mapping from `context_id` to event_stream
181    // and is maintained in memory. All event_streams are consumed through a separate event loop
182    compactor_streams_change_tx:
183        UnboundedSender<(HummockContextId, Streaming<SubscribeCompactionEventRequest>)>,
184
185    // `compaction_state` will record the types of compact tasks that can be triggered in `hummock`
186    // and suggest types with a certain priority.
187    pub compaction_state: CompactionState,
188    full_gc_state: Arc<FullGcState>,
189    now: Mutex<u64>,
190    inflight_time_travel_query: Semaphore,
191    gc_manager: GcManager,
192
193    /// In-memory cache of prefetched compaction task ids to reduce per-task DB round-trips.
194    prefetched_compaction_task_ids: PrefetchedSequence,
195    table_id_to_table_option: parking_lot::RwLock<HashMap<TableId, TableOption>>,
196}
197
198pub type HummockManagerRef = Arc<HummockManager>;
199
200use risingwave_object_store::object::{ObjectError, ObjectStoreRef, build_remote_object_store};
201use risingwave_pb::catalog::Table;
202
203use super::IcebergCompactorManager;
204use crate::controller::SqlMetaStore;
205use crate::hummock::manager::compaction_group_manager::CompactionGroupManager;
206use crate::hummock::manager::worker::HummockManagerEventSender;
207
208impl HummockManager {
209    pub async fn new(
210        env: MetaSrvEnv,
211        metadata_manager: MetadataManager,
212        metrics: Arc<MetaMetrics>,
213        compactor_manager: CompactorManagerRef,
214        compactor_streams_change_tx: UnboundedSender<(
215            HummockContextId,
216            Streaming<SubscribeCompactionEventRequest>,
217        )>,
218    ) -> Result<HummockManagerRef> {
219        let compaction_group_manager = CompactionGroupManager::new(&env).await?;
220        Self::new_impl(
221            env,
222            metadata_manager,
223            metrics,
224            compactor_manager,
225            compaction_group_manager,
226            compactor_streams_change_tx,
227        )
228        .await
229    }
230
231    #[cfg(any(test, feature = "test"))]
232    pub(super) async fn with_config(
233        env: MetaSrvEnv,
234        cluster_controller: crate::controller::cluster::ClusterControllerRef,
235        catalog_controller: crate::controller::catalog::CatalogControllerRef,
236        metrics: Arc<MetaMetrics>,
237        compactor_manager: CompactorManagerRef,
238        config: risingwave_pb::hummock::CompactionConfig,
239        compactor_streams_change_tx: UnboundedSender<(
240            HummockContextId,
241            Streaming<SubscribeCompactionEventRequest>,
242        )>,
243    ) -> HummockManagerRef {
244        let compaction_group_manager = CompactionGroupManager::new_with_config(&env, config)
245            .await
246            .unwrap();
247        let metadata_manager = MetadataManager::new(cluster_controller, catalog_controller);
248        Self::new_impl(
249            env,
250            metadata_manager,
251            metrics,
252            compactor_manager,
253            compaction_group_manager,
254            compactor_streams_change_tx,
255        )
256        .await
257        .unwrap()
258    }
259
260    async fn new_impl(
261        env: MetaSrvEnv,
262        metadata_manager: MetadataManager,
263        metrics: Arc<MetaMetrics>,
264        compactor_manager: CompactorManagerRef,
265        compaction_group_manager: CompactionGroupManager,
266        compactor_streams_change_tx: UnboundedSender<(
267            HummockContextId,
268            Streaming<SubscribeCompactionEventRequest>,
269        )>,
270    ) -> Result<HummockManagerRef> {
271        let sys_params = env.system_params_reader().await;
272        let state_store_url = sys_params.state_store();
273        let state_store_url = state_store_url.expose();
274
275        let state_store_dir: &str = sys_params.data_directory();
276        let use_new_object_prefix_strategy: bool = sys_params.use_new_object_prefix_strategy();
277        let deterministic_mode = env.opts.compaction_deterministic_test;
278        let mut object_store_config = env.opts.object_store_config.clone();
279        // For fs and hdfs object store, operations are not always atomic.
280        // We should manually enable atomicity guarantee by setting the atomic_write_dir config when building services.
281        object_store_config.set_atomic_write_dir();
282        let object_store = Arc::new(
283            build_remote_object_store(
284                state_store_url.strip_prefix("hummock+").unwrap_or("memory"),
285                metrics.object_store_metric.clone(),
286                "Version Checkpoint",
287                Arc::new(object_store_config),
288            )
289            .await,
290        );
291        // Make sure data dir is not used by another cluster.
292        // Skip this check in e2e compaction test, which needs to start a secondary cluster with
293        // same bucket
294        if !deterministic_mode {
295            write_exclusive_cluster_id(
296                state_store_dir,
297                env.cluster_id().clone(),
298                object_store.clone(),
299            )
300            .await?;
301
302            // config bucket lifecycle for new cluster.
303            if let risingwave_object_store::object::ObjectStoreImpl::S3(s3) = object_store.as_ref()
304                && !env.opts.do_not_config_object_storage_lifecycle
305            {
306                let is_bucket_expiration_configured =
307                    s3.inner().configure_bucket_lifecycle(state_store_dir).await;
308                if is_bucket_expiration_configured {
309                    return Err(ObjectError::internal("Cluster cannot start with object expiration configured for bucket because RisingWave data will be lost when object expiration kicks in.
310                    Please disable object expiration and restart the cluster.")
311                    .into());
312                }
313            }
314        }
315        let version_checkpoint_path = version_checkpoint_path(state_store_dir);
316        let version_archive_dir = version_archive_dir(state_store_dir);
317        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
318        let (version_stat_tx, mut version_stat_rx) = tokio::sync::mpsc::unbounded_channel();
319        let inflight_time_travel_query = env.opts.max_inflight_time_travel_query;
320        let gc_manager = GcManager::new(
321            object_store.clone(),
322            state_store_dir,
323            use_new_object_prefix_strategy,
324        );
325
326        let max_table_statistic_expired_time = std::cmp::max(
327            env.opts.table_stat_throuput_window_seconds_for_split,
328            env.opts.table_stat_throuput_window_seconds_for_merge,
329        ) as i64;
330
331        let iceberg_compactor_manager = Arc::new(IcebergCompactorManager::new());
332
333        let instance = HummockManager {
334            env,
335            versioning: MonitoredRwLock::new(
336                metrics.hummock_manager_lock_time.clone(),
337                metrics.hummock_manager_real_process_time.clone(),
338                Default::default(),
339                "hummock_manager::versioning",
340            ),
341            compaction: MonitoredRwLock::new(
342                metrics.hummock_manager_lock_time.clone(),
343                metrics.hummock_manager_real_process_time.clone(),
344                Default::default(),
345                "hummock_manager::compaction",
346            ),
347            compaction_group_manager: MonitoredRwLock::new(
348                metrics.hummock_manager_lock_time.clone(),
349                metrics.hummock_manager_real_process_time.clone(),
350                compaction_group_manager,
351                "hummock_manager::compaction_group_manager",
352            ),
353            context_info: MonitoredRwLock::new(
354                metrics.hummock_manager_lock_time.clone(),
355                metrics.hummock_manager_real_process_time.clone(),
356                Default::default(),
357                "hummock_manager::context_info",
358            ),
359            metrics,
360            metadata_manager,
361            compactor_manager,
362            iceberg_compactor_manager,
363            event_sender: tx,
364            object_store,
365            version_checkpoint_path,
366            version_archive_dir,
367            pause_version_checkpoint: AtomicBool::new(false),
368            table_write_throughput_statistic_manager: parking_lot::RwLock::new(
369                TableWriteThroughputStatisticManager::new(max_table_statistic_expired_time),
370            ),
371            table_committed_epoch_notifiers: parking_lot::Mutex::new(
372                TableCommittedEpochNotifiers {
373                    txs: Default::default(),
374                },
375            ),
376            compaction_task_report_notifiers: parking_lot::Mutex::new(
377                CompactionTaskReportNotifiers {
378                    txs: Default::default(),
379                },
380            ),
381            version_stat_tx,
382            compactor_streams_change_tx,
383            compaction_state: CompactionState::new(),
384            full_gc_state: FullGcState::new().into(),
385            now: Mutex::new(0),
386            inflight_time_travel_query: Semaphore::new(inflight_time_travel_query as usize),
387            gc_manager,
388            prefetched_compaction_task_ids: PrefetchedSequence::new(),
389            table_id_to_table_option: RwLock::new(HashMap::new()),
390        };
391        let instance = Arc::new(instance);
392        let version_stat_metrics = instance.metrics.clone();
393        tokio::spawn(async move {
394            while let Some(mut version) = version_stat_rx.recv().await {
395                while let Some(Some(next_version)) = version_stat_rx.recv().now_or_never() {
396                    version = next_version;
397                }
398                transaction::trigger_version_stat(&version_stat_metrics, version.as_ref());
399            }
400        });
401        instance.init_time_travel_state().await?;
402        instance.load_meta_store_state().await?;
403        instance.start_worker(rx);
404        instance.release_invalid_contexts().await?;
405        // Release snapshots pinned by meta on restarting.
406        instance.release_meta_context().await?;
407        Ok(instance)
408    }
409
410    fn meta_store_ref(&self) -> &SqlMetaStore {
411        self.env.meta_store_ref()
412    }
413
414    /// Load state from meta store.
415    async fn load_meta_store_state(&self) -> Result<()> {
416        let now = self.load_now().await?;
417        *self.now.lock().await = now.unwrap_or(0);
418
419        let mut compaction_guard = self.compaction.write().await;
420        let mut versioning_guard = self.versioning.write().await;
421        let mut context_info_guard = self.context_info.write().await;
422        self.load_meta_store_state_impl(
423            &mut compaction_guard,
424            &mut versioning_guard,
425            &mut context_info_guard,
426        )
427        .await
428    }
429
430    /// Load state from meta store.
431    async fn load_meta_store_state_impl(
432        &self,
433        compaction_guard: &mut Compaction,
434        versioning_guard: &mut Versioning,
435        context_info: &mut ContextInfo,
436    ) -> Result<()> {
437        use sea_orm::EntityTrait;
438        let meta_store = self.meta_store_ref();
439        let compaction_statuses: BTreeMap<CompactionGroupId, CompactStatus> =
440            compaction_status::Entity::find()
441                .all(&meta_store.conn)
442                .await
443                .map_err(MetadataModelError::from)?
444                .into_iter()
445                .map(|m| (m.compaction_group_id as CompactionGroupId, m.into()))
446                .collect();
447        if !compaction_statuses.is_empty() {
448            compaction_guard.compaction_statuses = compaction_statuses;
449        }
450
451        compaction_guard.compact_task_assignment = compaction_task::Entity::find()
452            .all(&meta_store.conn)
453            .await
454            .map_err(MetadataModelError::from)?
455            .into_iter()
456            .map(|m| {
457                (
458                    m.id as HummockCompactionTaskId,
459                    compaction_task_model_to_assignment(m),
460                )
461            })
462            .collect();
463
464        let hummock_version_deltas: BTreeMap<HummockVersionId, HummockVersionDelta> =
465            hummock_version_delta::Entity::find()
466                .all(&meta_store.conn)
467                .await
468                .map_err(MetadataModelError::from)?
469                .into_iter()
470                .map(|m| {
471                    (
472                        m.id,
473                        HummockVersionDelta::from_persisted_protobuf_owned(m.into()),
474                    )
475                })
476                .collect();
477
478        let checkpoint = self.try_read_checkpoint().await?;
479        let mut redo_state = if let Some(c) = checkpoint {
480            versioning_guard.checkpoint = c;
481            versioning_guard.checkpoint.version.as_ref().clone()
482        } else {
483            let default_compaction_config = self
484                .compaction_group_manager
485                .read()
486                .await
487                .default_compaction_config();
488            let checkpoint_version = HummockVersion::create_init_version(default_compaction_config);
489            tracing::info!("init hummock version checkpoint");
490            versioning_guard.checkpoint = HummockVersionCheckpoint {
491                version: Arc::new(checkpoint_version.clone()),
492                stale_objects: Default::default(),
493            };
494            self.write_checkpoint(&versioning_guard.checkpoint).await?;
495            checkpoint_version
496        };
497        let mut applied_delta_count = 0;
498        let total_to_apply = hummock_version_deltas.range(redo_state.id + 1..).count();
499        tracing::info!(
500            total_delta = hummock_version_deltas.len(),
501            total_to_apply,
502            "Start redo Hummock version."
503        );
504        for version_delta in hummock_version_deltas
505            .range(redo_state.id + 1..)
506            .map(|(_, v)| v)
507        {
508            assert_eq!(
509                version_delta.prev_id, redo_state.id,
510                "delta prev_id {}, redo state id {}",
511                version_delta.prev_id, redo_state.id
512            );
513            redo_state.apply_version_delta(version_delta);
514            applied_delta_count += 1;
515            if applied_delta_count % 1000 == 0 {
516                tracing::info!("Redo progress {applied_delta_count}/{total_to_apply}.");
517            }
518        }
519        tracing::info!("Finish redo Hummock version.");
520        let pruned_stale_table_id_count = redo_state.prune_stale_table_ids_from_ssts();
521        if pruned_stale_table_id_count > 0 {
522            tracing::warn!(
523                pruned_stale_table_id_count,
524                version_id = ?redo_state.id,
525                "Pruned stale table ids from recovered Hummock SST metadata."
526            );
527        }
528        versioning_guard.version_stats = hummock_version_stats::Entity::find()
529            .one(&meta_store.conn)
530            .await
531            .map_err(MetadataModelError::from)?
532            .map(HummockVersionStats::from)
533            .unwrap_or_else(|| HummockVersionStats {
534                // version_stats.hummock_version_id is always 0 in meta store.
535                hummock_version_id: 0.into(),
536                ..Default::default()
537            });
538
539        versioning_guard.current_version = Arc::new(redo_state);
540        versioning_guard.hummock_version_deltas = hummock_version_deltas;
541        versioning_guard.table_change_log =
542            risingwave_meta_model::hummock_table_change_log::Entity::find()
543                .all(&self.env.meta_store_ref().conn)
544                .await
545                .map_err(MetadataModelError::from)?
546                .into_iter()
547                .map(|m| (m.table_id, to_table_change_log(m)))
548                .into_group_map()
549                .into_iter()
550                .map(|(table_id, unordered_change_logs)| {
551                    (
552                        table_id,
553                        TableChangeLog::new(
554                            unordered_change_logs
555                                .into_iter()
556                                .sorted_by_key(|l| l.checkpoint_epoch),
557                        ),
558                    )
559                })
560                .collect();
561
562        context_info.pinned_versions = hummock_pinned_version::Entity::find()
563            .all(&meta_store.conn)
564            .await
565            .map_err(MetadataModelError::from)?
566            .into_iter()
567            .map(|m| (m.context_id as HummockContextId, m.into()))
568            .collect();
569
570        self.initial_compaction_group_config_after_load(
571            versioning_guard,
572            self.compaction_group_manager.write().await.deref_mut(),
573        )
574        .await?;
575
576        Ok(())
577    }
578
579    pub fn init_metadata_for_version_replay(
580        &self,
581        _table_catalogs: Vec<Table>,
582        _compaction_groups: Vec<PbCompactionGroupInfo>,
583    ) -> Result<()> {
584        unimplemented!("kv meta store is deprecated");
585    }
586
587    /// Replay a version delta to current hummock version.
588    /// Returns the `version_id`, `max_committed_epoch` of the new version and the modified
589    /// compaction groups
590    #[cfg(any(test, feature = "test"))]
591    pub async fn replay_version_delta(
592        &self,
593        mut version_delta: HummockVersionDelta,
594    ) -> Result<(HummockVersion, Vec<CompactionGroupId>)> {
595        let mut versioning_guard = self.versioning.write().await;
596        // ensure the version id is ascending after replay
597        version_delta.id = versioning_guard.current_version.next_version_id();
598        version_delta.prev_id = versioning_guard.current_version.id;
599        let mut version_new = versioning_guard.current_version.as_ref().clone();
600        version_new.apply_version_delta(&version_delta);
601
602        let compaction_group_ids = version_delta.group_deltas.keys().cloned().collect();
603        versioning_guard.current_version = Arc::new(version_new.clone());
604        Ok((version_new, compaction_group_ids))
605    }
606
607    pub async fn disable_commit_epoch(&self) -> Arc<HummockVersion> {
608        let mut versioning_guard = self.versioning.write().await;
609        versioning_guard.disable_commit_epochs = true;
610        versioning_guard.current_version.clone()
611    }
612
613    pub fn metadata_manager(&self) -> &MetadataManager {
614        &self.metadata_manager
615    }
616
617    pub fn object_store_media_type(&self) -> &'static str {
618        self.object_store.media_type()
619    }
620
621    pub fn update_table_id_to_table_option(
622        &self,
623        new_table_id_to_table_option: HashMap<TableId, TableOption>,
624    ) {
625        *self.table_id_to_table_option.write() = new_table_id_to_table_option;
626    }
627
628    pub fn metadata_manager_ref(&self) -> &MetadataManager {
629        &self.metadata_manager
630    }
631
632    pub async fn subscribe_table_committed_epoch(
633        &self,
634        table_id: TableId,
635    ) -> MetaResult<(u64, UnboundedReceiver<u64>)> {
636        let version = self.versioning.read().await;
637        if let Some(epoch) = version.current_version.table_committed_epoch(table_id) {
638            let (tx, rx) = unbounded_channel();
639            self.table_committed_epoch_notifiers
640                .lock()
641                .txs
642                .entry(table_id)
643                .or_default()
644                .push(tx);
645            Ok((epoch, rx))
646        } else {
647            Err(anyhow!("table {} not exist", table_id).into())
648        }
649    }
650}
651
652async fn write_exclusive_cluster_id(
653    state_store_dir: &str,
654    cluster_id: ClusterId,
655    object_store: ObjectStoreRef,
656) -> Result<()> {
657    const CLUSTER_ID_DIR: &str = "cluster_id";
658    const CLUSTER_ID_NAME: &str = "0";
659    let cluster_id_dir = format!("{}/{}/", state_store_dir, CLUSTER_ID_DIR);
660    let cluster_id_full_path = format!("{}{}", cluster_id_dir, CLUSTER_ID_NAME);
661    tracing::info!("try reading cluster_id");
662    match object_store.read(&cluster_id_full_path, ..).await {
663        Ok(stored_cluster_id) => {
664            let stored_cluster_id = String::from_utf8(stored_cluster_id.to_vec()).unwrap();
665            if cluster_id.deref() == stored_cluster_id {
666                return Ok(());
667            }
668
669            Err(ObjectError::internal(format!(
670                "Data directory is already used by another cluster with id {:?}, path {}.",
671                stored_cluster_id, cluster_id_full_path,
672            ))
673            .into())
674        }
675        Err(e) => {
676            if e.is_object_not_found_error() {
677                tracing::info!("cluster_id not found, writing cluster_id");
678                object_store
679                    .upload(&cluster_id_full_path, Bytes::from(String::from(cluster_id)))
680                    .await?;
681                return Ok(());
682            }
683            Err(e.into())
684        }
685    }
686}