Skip to main content

common_meta/ddl/
alter_database.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
15use async_trait::async_trait;
16use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu};
17use common_procedure::{
18    Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, Status,
19};
20use common_telemetry::tracing::info;
21use serde::{Deserialize, Serialize};
22use snafu::{ResultExt, ensure};
23use store_api::mito_engine_options::{TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_TRIGGER_FILE_NUM};
24use strum::AsRefStr;
25
26use crate::cache_invalidator::Context;
27use crate::ddl::DdlContext;
28use crate::ddl::event::database::{ALTER_DATABASE_EVENT_TYPE, DatabaseDdlEvent};
29use crate::ddl::utils::map_to_procedure_error;
30use crate::error::{Result, SchemaNotFoundSnafu};
31use crate::instruction::CacheIdent;
32use crate::key::DeserializedValueWithBytes;
33use crate::key::schema_name::{SchemaName, SchemaNameKey, SchemaNameValue};
34use crate::lock_key::{CatalogLock, SchemaLock};
35use crate::rpc::ddl::UnsetDatabaseOption::{self};
36use crate::rpc::ddl::{AlterDatabaseKind, AlterDatabaseTask, SetDatabaseOption};
37
38pub struct AlterDatabaseProcedure {
39    pub context: DdlContext,
40    pub data: AlterDatabaseData,
41}
42
43fn twcs_trigger_alias(key: &str) -> Option<&'static str> {
44    match key {
45        TWCS_TRIGGER_FILE_NUM => Some(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM),
46        TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM => Some(TWCS_TRIGGER_FILE_NUM),
47        _ => None,
48    }
49}
50
51fn build_new_schema_value(
52    mut value: SchemaNameValue,
53    alter_kind: &AlterDatabaseKind,
54) -> Result<SchemaNameValue> {
55    match alter_kind {
56        AlterDatabaseKind::SetDatabaseOptions(options) => {
57            for option in options.0.iter() {
58                match option {
59                    SetDatabaseOption::Ttl(ttl) => {
60                        value.ttl = Some(*ttl);
61                    }
62                    SetDatabaseOption::Other(key, val) => {
63                        // Keep the legacy key so older versions can read it after a downgrade.
64                        // Persisting both aliases would deserialize as a duplicate field.
65                        let persisted_key = if twcs_trigger_alias(key).is_some() {
66                            value.extra_options.remove(TWCS_TRIGGER_FILE_NUM);
67                            value
68                                .extra_options
69                                .remove(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM);
70                            TWCS_TRIGGER_FILE_NUM
71                        } else {
72                            key
73                        };
74                        value
75                            .extra_options
76                            .insert(persisted_key.to_string(), val.clone());
77                    }
78                }
79            }
80        }
81        AlterDatabaseKind::UnsetDatabaseOptions(keys) => {
82            for key in keys.0.iter() {
83                match key {
84                    UnsetDatabaseOption::Ttl => value.ttl = None,
85                    UnsetDatabaseOption::Other(key) => {
86                        value.extra_options.remove(key);
87                        if let Some(alias) = twcs_trigger_alias(key) {
88                            value.extra_options.remove(alias);
89                        }
90                    }
91                }
92            }
93        }
94    }
95    Ok(value)
96}
97
98impl AlterDatabaseProcedure {
99    pub const TYPE_NAME: &'static str = "metasrv-procedure::AlterDatabase";
100
101    pub fn new(task: AlterDatabaseTask, context: DdlContext) -> Result<Self> {
102        Ok(Self {
103            context,
104            data: AlterDatabaseData::new(task)?,
105        })
106    }
107
108    pub fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
109        let data = serde_json::from_str(json).context(FromJsonSnafu)?;
110
111        Ok(Self { context, data })
112    }
113
114    pub async fn on_prepare(&mut self) -> Result<Status> {
115        let value = self
116            .context
117            .table_metadata_manager
118            .schema_manager()
119            .get(SchemaNameKey::new(self.data.catalog(), self.data.schema()))
120            .await?;
121
122        ensure!(
123            value.is_some(),
124            SchemaNotFoundSnafu {
125                table_schema: self.data.schema(),
126            }
127        );
128
129        self.data.schema_value = value;
130        self.data.state = AlterDatabaseState::UpdateMetadata;
131
132        Ok(Status::executing(true))
133    }
134
135    pub async fn on_update_metadata(&mut self) -> Result<Status> {
136        let schema_name = SchemaNameKey::new(self.data.catalog(), self.data.schema());
137
138        // Safety: schema_value is not None.
139        let current_schema_value = self.data.schema_value.as_ref().unwrap();
140
141        let new_schema_value = build_new_schema_value(
142            current_schema_value.get_inner_ref().clone(),
143            &self.data.kind,
144        )?;
145
146        self.context
147            .table_metadata_manager
148            .schema_manager()
149            .update(schema_name, current_schema_value, &new_schema_value)
150            .await?;
151
152        info!("Updated database metadata for schema {schema_name}");
153        self.data.state = AlterDatabaseState::InvalidateSchemaCache;
154        Ok(Status::executing(true))
155    }
156
157    pub async fn on_invalidate_schema_cache(&mut self) -> Result<Status> {
158        let cache_invalidator = &self.context.cache_invalidator;
159        cache_invalidator
160            .invalidate(
161                &Context::default(),
162                &[CacheIdent::SchemaName(SchemaName {
163                    catalog_name: self.data.catalog().to_string(),
164                    schema_name: self.data.schema().to_string(),
165                })],
166            )
167            .await?;
168
169        Ok(Status::done())
170    }
171}
172
173#[async_trait]
174impl Procedure for AlterDatabaseProcedure {
175    fn type_name(&self) -> &str {
176        Self::TYPE_NAME
177    }
178
179    async fn execute(&mut self, _ctx: &ProcedureContext) -> ProcedureResult<Status> {
180        match self.data.state {
181            AlterDatabaseState::Prepare => self.on_prepare().await,
182            AlterDatabaseState::UpdateMetadata => self.on_update_metadata().await,
183            AlterDatabaseState::InvalidateSchemaCache => self.on_invalidate_schema_cache().await,
184        }
185        .map_err(map_to_procedure_error)
186    }
187
188    fn dump(&self) -> ProcedureResult<String> {
189        serde_json::to_string(&self.data).context(ToJsonSnafu)
190    }
191
192    fn lock_key(&self) -> LockKey {
193        let catalog = self.data.catalog();
194        let schema = self.data.schema();
195
196        let lock_key = vec![
197            CatalogLock::Read(catalog).into(),
198            SchemaLock::write(catalog, schema).into(),
199        ];
200
201        LockKey::new(lock_key)
202    }
203
204    fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
205        if !ctx.event_type_filter.allows(ALTER_DATABASE_EVENT_TYPE) {
206            return None;
207        }
208
209        let event = if matches!(&ctx.trigger, EventTrigger::Submitted) {
210            DatabaseDdlEvent::alter_submitted(
211                self.data.catalog(),
212                self.data.schema(),
213                &self.data.kind,
214            )
215        } else {
216            DatabaseDdlEvent::alter_lifecycle(self.data.catalog(), self.data.schema())
217        };
218        Some(Box::new(event))
219    }
220}
221
222#[derive(Debug, Serialize, Deserialize, AsRefStr)]
223enum AlterDatabaseState {
224    Prepare,
225    UpdateMetadata,
226    InvalidateSchemaCache,
227}
228
229/// The data of alter database procedure.
230#[derive(Debug, Serialize, Deserialize)]
231pub struct AlterDatabaseData {
232    state: AlterDatabaseState,
233    kind: AlterDatabaseKind,
234    catalog_name: String,
235    schema_name: String,
236    schema_value: Option<DeserializedValueWithBytes<SchemaNameValue>>,
237}
238
239impl AlterDatabaseData {
240    pub fn new(task: AlterDatabaseTask) -> Result<Self> {
241        Ok(Self {
242            state: AlterDatabaseState::Prepare,
243            kind: AlterDatabaseKind::try_from(task.alter_expr.kind.unwrap())?,
244            catalog_name: task.alter_expr.catalog_name,
245            schema_name: task.alter_expr.schema_name,
246            schema_value: None,
247        })
248    }
249
250    pub fn catalog(&self) -> &str {
251        &self.catalog_name
252    }
253
254    pub fn schema(&self) -> &str {
255        &self.schema_name
256    }
257}
258
259#[cfg(test)]
260mod tests {
261    use std::time::Duration;
262
263    use store_api::mito_engine_options::{
264        TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_TRIGGER_FILE_NUM,
265    };
266
267    use crate::ddl::alter_database::build_new_schema_value;
268    use crate::key::schema_name::SchemaNameValue;
269    use crate::rpc::ddl::{
270        AlterDatabaseKind, SetDatabaseOption, SetDatabaseOptions, UnsetDatabaseOption,
271        UnsetDatabaseOptions,
272    };
273
274    #[test]
275    fn test_build_new_schema_value() {
276        let set_ttl = AlterDatabaseKind::SetDatabaseOptions(SetDatabaseOptions(vec![
277            SetDatabaseOption::Ttl(Duration::from_secs(10).into()),
278        ]));
279        let current_schema_value = SchemaNameValue::default();
280        let new_schema_value =
281            build_new_schema_value(current_schema_value.clone(), &set_ttl).unwrap();
282        assert_eq!(new_schema_value.ttl, Some(Duration::from_secs(10).into()));
283
284        let unset_ttl_alter_kind =
285            AlterDatabaseKind::UnsetDatabaseOptions(UnsetDatabaseOptions(vec![
286                UnsetDatabaseOption::Ttl,
287            ]));
288        let new_schema_value =
289            build_new_schema_value(current_schema_value, &unset_ttl_alter_kind).unwrap();
290        assert_eq!(new_schema_value.ttl, None);
291    }
292
293    #[test]
294    fn test_build_new_schema_value_with_compaction_options() {
295        let set_compaction = AlterDatabaseKind::SetDatabaseOptions(SetDatabaseOptions(vec![
296            SetDatabaseOption::Other("compaction.type".to_string(), "twcs".to_string()),
297            SetDatabaseOption::Other("compaction.twcs.time_window".to_string(), "1d".to_string()),
298        ]));
299
300        let current_schema_value = SchemaNameValue::default();
301        let new_schema_value =
302            build_new_schema_value(current_schema_value.clone(), &set_compaction).unwrap();
303
304        assert_eq!(
305            new_schema_value.extra_options.get("compaction.type"),
306            Some(&"twcs".to_string())
307        );
308        assert_eq!(
309            new_schema_value
310                .extra_options
311                .get("compaction.twcs.time_window"),
312            Some(&"1d".to_string())
313        );
314
315        let unset_compaction = AlterDatabaseKind::UnsetDatabaseOptions(UnsetDatabaseOptions(vec![
316            UnsetDatabaseOption::Other("compaction.type".to_string()),
317        ]));
318
319        let new_schema_value = build_new_schema_value(new_schema_value, &unset_compaction).unwrap();
320
321        assert_eq!(new_schema_value.extra_options.get("compaction.type"), None);
322        assert_eq!(
323            new_schema_value
324                .extra_options
325                .get("compaction.twcs.time_window"),
326            Some(&"1d".to_string())
327        );
328    }
329
330    #[test]
331    fn test_set_twcs_trigger_persists_legacy_key() {
332        for key in [TWCS_TRIGGER_FILE_NUM, TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM] {
333            let mut current_schema_value = SchemaNameValue::default();
334            current_schema_value.extra_options.insert(
335                TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM.to_string(),
336                "8".to_string(),
337            );
338            let set = AlterDatabaseKind::SetDatabaseOptions(SetDatabaseOptions(vec![
339                SetDatabaseOption::Other(key.to_string(), "16".to_string()),
340            ]));
341
342            let new_schema_value = build_new_schema_value(current_schema_value, &set).unwrap();
343
344            assert_eq!(
345                new_schema_value
346                    .extra_options
347                    .get(TWCS_TRIGGER_FILE_NUM)
348                    .map(String::as_str),
349                Some("16")
350            );
351            assert!(
352                !new_schema_value
353                    .extra_options
354                    .contains_key(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM)
355            );
356        }
357    }
358
359    #[test]
360    fn test_unset_twcs_trigger_removes_both_aliases() {
361        for key in [TWCS_TRIGGER_FILE_NUM, TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM] {
362            let mut current_schema_value = SchemaNameValue::default();
363            current_schema_value
364                .extra_options
365                .insert(TWCS_TRIGGER_FILE_NUM.to_string(), "8".to_string());
366            current_schema_value.extra_options.insert(
367                TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM.to_string(),
368                "16".to_string(),
369            );
370            let unset = AlterDatabaseKind::UnsetDatabaseOptions(UnsetDatabaseOptions(vec![
371                UnsetDatabaseOption::Other(key.to_string()),
372            ]));
373
374            let new_schema_value = build_new_schema_value(current_schema_value, &unset).unwrap();
375
376            assert!(
377                !new_schema_value
378                    .extra_options
379                    .contains_key(TWCS_TRIGGER_FILE_NUM)
380            );
381            assert!(
382                !new_schema_value
383                    .extra_options
384                    .contains_key(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM)
385            );
386        }
387    }
388}