1use 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 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 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#[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}