Skip to main content

mito2/worker/
handle_open.rs

1// Copyright 2023 Greptime Team
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
15//! Handling open request.
16
17use std::sync::Arc;
18use std::time::Instant;
19
20use common_telemetry::{info, warn};
21use object_store::util::{join_path, normalize_dir};
22use snafu::{OptionExt, ResultExt};
23use store_api::logstore::LogStore;
24use store_api::region_request::{AffectedRows, RegionCleanUpRequest, RegionOpenRequest};
25use store_api::storage::RegionId;
26use table::requests::STORAGE_KEY;
27
28use crate::access_layer::AccessLayer;
29use crate::engine::region_hook::{RegionGcInfo, RegionHookRef};
30use crate::error::{
31    ObjectStoreNotFoundSnafu, OpenDalSnafu, OpenRegionSnafu, RegionBusySnafu, RegionNotFoundSnafu,
32    Result,
33};
34use crate::region::opener::{
35    RegionOpener, get_object_store, provider_from_wal_options, sanitize_open_request_options,
36};
37use crate::region::options::RegionOptions;
38use crate::request::OptionOutputTx;
39use crate::sst::location::region_dir_from_table_dir;
40use crate::wal::entry_distributor::WalEntryReceiver;
41use crate::worker::handle_drop::{
42    cleanup_region_file_artifacts, remove_region_dir_for_full_drop, remove_region_dir_once,
43};
44use crate::worker::{DROPPING_MARKER_FILE, RegionWorkerLoop};
45
46impl<S: LogStore> RegionWorkerLoop<S> {
47    pub(crate) async fn handle_offline_cleanup_request(
48        &mut self,
49        region_id: RegionId,
50        mut request: RegionCleanUpRequest,
51    ) -> Result<AffectedRows> {
52        info!(
53            "Try to clean region {} offline, worker: {}",
54            region_id, self.id
55        );
56
57        if self.regions.is_region_exists(region_id) {
58            return RegionBusySnafu { region_id }.fail();
59        }
60
61        sanitize_open_request_options(&mut request.options);
62
63        let options = RegionOptions::try_from_options(region_id, &request.options)?;
64        let object_store = get_object_store(&options.storage, &self.object_store_manager)?;
65        let provider = provider_from_wal_options::<S>(region_id, &options.wal_options)?;
66        self.wal.obsolete_all(region_id, &provider).await?;
67
68        let table_dir = normalize_dir(&request.table_dir);
69        let region_dir = region_dir_from_table_dir(&table_dir, region_id, request.path_type);
70        remove_region_dir_for_full_drop(&region_dir, &object_store).await?;
71
72        // The region directory is gone. Offline cleanup (soft-drop PURGE) is a
73        // terminal physical removal, identical in effect to the drop GC
74        // worker's directory deletion — but unlike the drop worker it never
75        // fires `on_region_files_removed`, because the region is offline and
76        // carries no `RegionMetadataRef`. Fire `on_region_gc` instead: it
77        // accepts an absent region and lets extensions with sidecar files
78        // outside the region dir (e.g. an Iceberg export) reclaim their
79        // per-region state.
80        //
81        // Propagate the hook's error so the caller retries the cleanup, mirroring
82        // how the global GC worker keeps a region for the next pass when
83        // `on_region_gc` fails. The whole handler is idempotent (the existing
84        // `test_engine_offline_cleanup_closed_region` already drives it twice in
85        // a row), so a retry re-runs the directory removal, WAL obsolete and the
86        // hook safely. Runs inline, consistent with the blocking directory
87        // removal above; offline cleanup is already a slow path.
88        if let Some(hook) = self.plugins.get::<RegionHookRef>() {
89            let access_layer = Arc::new(AccessLayer::new(
90                table_dir.clone(),
91                request.path_type,
92                object_store.clone(),
93                self.puffin_manager_factory.clone(),
94                self.intermediate_manager.clone(),
95            ));
96            let gc_info = RegionGcInfo {
97                removed_files: &[],
98                is_region_dropped: true,
99                full_file_listing: true,
100            };
101            if let Err(err) = hook
102                .on_region_gc(region_id, None, &access_layer, &gc_info)
103                .await
104            {
105                warn!(
106                    err;
107                    "Region hook on_region_gc failed during offline cleanup for region {}; \
108                     returning the error so the caller retries the cleanup",
109                    region_id,
110                );
111                return Err(err);
112            }
113        }
114
115        self.cleanup_dropped_region_runtime_state(region_id).await;
116        self.dropping_regions.remove_region(region_id);
117        cleanup_region_file_artifacts(
118            region_id,
119            &table_dir,
120            &self.intermediate_manager,
121            &self.cache_manager,
122        )
123        .await;
124
125        Ok(0)
126    }
127
128    async fn check_and_cleanup_region(
129        &self,
130        region_id: RegionId,
131        request: &RegionOpenRequest,
132    ) -> Result<()> {
133        let object_store = if let Some(storage_name) = request.options.get(STORAGE_KEY) {
134            self.object_store_manager
135                .find(storage_name)
136                .with_context(|| ObjectStoreNotFoundSnafu {
137                    object_store: storage_name.clone(),
138                })?
139        } else {
140            self.object_store_manager.default_object_store()
141        };
142        // Check if this region is pending drop. And clean the entire dir if so.
143        let region_dir =
144            region_dir_from_table_dir(&request.table_dir, region_id, request.path_type);
145        if !self.dropping_regions.is_region_exists(region_id)
146            && object_store
147                .exists(&join_path(&region_dir, DROPPING_MARKER_FILE))
148                .await
149                .context(OpenDalSnafu)?
150        {
151            let result = remove_region_dir_once(&region_dir, object_store, true).await;
152            info!(
153                "Region {} is dropped, worker: {}, result: {:?}",
154                region_id, self.id, result
155            );
156            return RegionNotFoundSnafu { region_id }.fail();
157        }
158
159        Ok(())
160    }
161
162    pub(crate) async fn handle_open_request(
163        &mut self,
164        region_id: RegionId,
165        mut request: RegionOpenRequest,
166        wal_entry_receiver: Option<WalEntryReceiver>,
167        sender: OptionOutputTx,
168    ) {
169        if self.regions.is_region_exists(region_id) {
170            sender.send(Ok(0));
171            return;
172        }
173        let Some(sender) = self
174            .opening_regions
175            .wait_for_opening_region(region_id, sender)
176        else {
177            return;
178        };
179        info!("Try to open region {}, worker: {}", region_id, self.id);
180        sanitize_open_request_options(&mut request.options);
181
182        // Open region from specific region dir.
183        let requirements = request.requirements;
184        let opener = match RegionOpener::new(
185            region_id,
186            &request.table_dir,
187            request.path_type,
188            self.memtable_builder_provider.clone(),
189            self.object_store_manager.clone(),
190            self.purge_scheduler.clone(),
191            self.puffin_manager_factory.clone(),
192            self.intermediate_manager.clone(),
193            self.time_provider.clone(),
194            self.file_ref_manager.clone(),
195            self.partition_expr_fetcher.clone(),
196        )
197        .skip_wal_replay(request.skip_wal_replay)
198        .cache(Some(self.cache_manager.clone()))
199        .hook(self.plugins.get())
200        .wal_entry_reader(wal_entry_receiver.map(|receiver| Box::new(receiver) as _))
201        .replay_checkpoint(request.checkpoint.map(|checkpoint| checkpoint.entry_id))
202        .parse_options(request.options.clone())
203        {
204            Ok(opener) => opener,
205            Err(err) => {
206                sender.send(Err(err));
207                return;
208            }
209        };
210
211        if let Err(err) = opener.ensure_region_requirements(requirements) {
212            sender.send(Err(err));
213            return;
214        }
215
216        if let Err(err) = self.check_and_cleanup_region(region_id, &request).await {
217            sender.send(Err(err));
218            return;
219        }
220
221        let now = Instant::now();
222        let regions = self.regions.clone();
223        let wal = self.wal.clone();
224        let config = self.config.clone();
225        let opening_regions = self.opening_regions.clone();
226        let region_count = self.region_count.clone();
227        let worker_id = self.id;
228        opening_regions.insert_sender(region_id, sender);
229        common_runtime::spawn_global(async move {
230            match opener.open(&config, &wal).await {
231                Ok(region) => {
232                    info!(
233                        "Region {} is opened with requirements {:?}, worker: {}, elapsed: {:?}",
234                        region_id,
235                        requirements,
236                        worker_id,
237                        now.elapsed()
238                    );
239                    region_count.inc();
240
241                    // Notify the region hook that the region has been opened.
242                    // Fires before registration; allocates nothing when no hook
243                    // is registered.
244                    if let Some(hook) = region.manifest_ctx.hook() {
245                        hook.on_region_opened(region_id, &region.metadata()).await;
246                    }
247
248                    // Insert the Region into the RegionMap.
249                    regions.insert_region(region);
250
251                    let senders = opening_regions.remove_sender(region_id);
252                    for sender in senders {
253                        sender.send(Ok(0));
254                    }
255                }
256                Err(err) => {
257                    let senders = opening_regions.remove_sender(region_id);
258                    let err = Arc::new(err);
259                    for sender in senders {
260                        sender.send(Err(err.clone()).context(OpenRegionSnafu));
261                    }
262                }
263            }
264        });
265    }
266}