1mod executor;
16mod metadata;
17mod region_request;
18
19use std::collections::HashSet;
20use std::vec;
21
22use api::region::RegionResponse;
23use api::v1::alter_table_expr::Kind;
24use api::v1::{RenameTable, SetTableOptions, UnsetTableOptions};
25use async_trait::async_trait;
26use common_catalog::consts::{METRIC_ENGINE, MITO_ENGINE};
27use common_error::ext::BoxedError;
28use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu};
29use common_procedure::{
30 Context as ProcedureContext, ContextProvider, Error as ProcedureError, EventContext,
31 EventTrigger, LockKey, PoisonKey, PoisonKeys, Procedure, ProcedureId, Status, StringKey,
32};
33use common_telemetry::{error, info, warn};
34use serde::{Deserialize, Serialize};
35use snafu::{ResultExt, ensure};
36use store_api::metadata::ColumnMetadata;
37use store_api::metric_engine_consts::TABLE_COLUMN_METADATA_EXTENSION_KEY;
38use store_api::storage::RegionId;
39use strum::AsRefStr;
40use table::metadata::{TableId, TableInfo};
41use table::requests::{REPARTITION_COLUMN_HINT_KEY, SKIP_WAL_KEY};
42use table::table_reference::TableReference;
43
44use crate::ddl::DdlContext;
45use crate::ddl::alter_table::executor::AlterTableExecutor;
46use crate::ddl::event::table::{
47 TableDdlEvent, TableDdlEventType, TableDdlLocator, alter_table_kind_name,
48};
49use crate::ddl::utils::{
50 MultipleResults, extract_column_metadatas, handle_multiple_results, map_to_procedure_error,
51 sync_follower_regions,
52};
53use crate::error::{
54 AbortProcedureSnafu, NoLeaderSnafu, PutPoisonSnafu, Result, RetryLaterSnafu, UnsupportedSnafu,
55};
56use crate::key::table_info::TableInfoValue;
57use crate::key::{DeserializedValueWithBytes, RegionDistribution};
58use crate::lock_key::{CatalogLock, RegionLock, SchemaLock, TableLock, TableNameLock};
59use crate::metrics;
60use crate::poison_key::table_poison_key;
61use crate::rpc::ddl::AlterTableTask;
62use crate::rpc::router::{RegionRoute, find_leaders, region_distribution};
63
64pub struct AlterTableProcedure {
66 context: DdlContext,
68 data: AlterTableData,
70 new_table_info: Option<TableInfo>,
74 executor: AlterTableExecutor,
76}
77
78#[derive(Debug)]
79pub(crate) struct RegionRouteChanged;
80
81fn build_executor_from_alter_expr(alter_data: &AlterTableData) -> AlterTableExecutor {
86 let table_name = alter_data.table_ref().into();
87 let table_id = alter_data.table_id;
88 let alter_kind = alter_data.task.alter_table.kind.as_ref().unwrap();
89 let new_table_name = if let Kind::RenameTable(RenameTable { new_table_name }) = alter_kind {
90 Some(new_table_name.clone())
91 } else {
92 None
93 };
94 AlterTableExecutor::new(table_name, table_id, new_table_name)
95}
96
97impl AlterTableProcedure {
98 pub const TYPE_NAME: &'static str = "metasrv-procedure::AlterTable";
99
100 pub fn new(table_id: TableId, task: AlterTableTask, context: DdlContext) -> Result<Self> {
101 Self::new_with_region_locks(table_id, task, vec![], context)
102 }
103
104 pub(crate) fn new_with_region_locks(
105 table_id: TableId,
106 task: AlterTableTask,
107 region_locks: Vec<RegionId>,
108 context: DdlContext,
109 ) -> Result<Self> {
110 task.validate()?;
111 let data = AlterTableData::new(task, table_id, region_locks);
112 let executor = build_executor_from_alter_expr(&data);
113 Ok(Self {
114 context,
115 data,
116 new_table_info: None,
117 executor,
118 })
119 }
120
121 pub fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
122 let data: AlterTableData = serde_json::from_str(json).context(FromJsonSnafu)?;
123 let executor = build_executor_from_alter_expr(&data);
124
125 Ok(AlterTableProcedure {
126 context,
127 data,
128 new_table_info: None,
129 executor,
130 })
131 }
132
133 pub(crate) async fn on_prepare(&mut self) -> Result<Status> {
135 self.executor
136 .on_prepare(&self.context.table_metadata_manager)
137 .await?;
138 self.fill_table_info().await?;
139
140 let table_info_value = self.data.table_info_value.as_ref().unwrap();
142 let new_table_info = AlterTableExecutor::validate_alter_table_expr(
143 &table_info_value.table_info,
144 self.data.task.alter_table.clone(),
145 )?;
146 self.new_table_info = Some(new_table_info);
147
148 let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
150 if enables_skip_wal(alter_kind) {
151 ensure!(
152 only_enables_skip_wal(alter_kind),
153 UnsupportedSnafu {
154 operation: "combining skip_wal = 'true' with other table options".to_string()
155 }
156 );
157 let engine = table_info_value.table_info.meta.engine.as_str();
158 ensure!(
159 engine == MITO_ENGINE || engine == METRIC_ENGINE,
160 UnsupportedSnafu {
161 operation: format!("setting skip_wal on {engine} engine tables")
162 }
163 );
164 let (physical_table_id, physical_table_route) = self
166 .context
167 .table_metadata_manager
168 .table_route_manager()
169 .get_physical_table_route(self.data.table_id())
170 .await?;
171 ensure!(
172 physical_table_id == self.data.table_id(),
173 UnsupportedSnafu {
174 operation: "setting skip_wal on logical tables".to_string()
175 }
176 );
177 let current_region_ids = physical_table_route
178 .region_routes
179 .iter()
180 .map(|route| route.region.id)
181 .collect::<HashSet<_>>();
182 let locked_region_ids = self
183 .data
184 .region_locks
185 .iter()
186 .copied()
187 .collect::<HashSet<_>>();
188 if physical_table_route.region_routes.len() != self.data.region_locks.len()
189 || current_region_ids != locked_region_ids
190 {
191 return Ok(Status::done_with_output(RegionRouteChanged));
194 }
195 ensure!(
196 !find_leaders(&physical_table_route.region_routes).is_empty(),
197 NoLeaderSnafu {
198 table_id: physical_table_id
199 }
200 );
201 self.data.region_distribution =
202 Some(region_distribution(&physical_table_route.region_routes));
203 }
204 self.data.state = self.data.flow().after_prepare();
205 Ok(Status::executing(true))
206 }
207
208 fn table_poison_key(&self) -> PoisonKey {
209 table_poison_key(self.data.table_id())
210 }
211
212 async fn put_poison(
213 &self,
214 ctx_provider: &dyn ContextProvider,
215 procedure_id: ProcedureId,
216 ) -> Result<()> {
217 let poison_key = self.table_poison_key();
218 ctx_provider
219 .try_put_poison(&poison_key, procedure_id)
220 .await
221 .context(PutPoisonSnafu)
222 }
223
224 pub async fn submit_alter_region_requests(
225 &mut self,
226 procedure_id: ProcedureId,
227 ctx_provider: &dyn ContextProvider,
228 ) -> Result<Status> {
229 let table_id = self.data.table_id();
230 let (_, physical_table_route) = self
231 .context
232 .table_metadata_manager
233 .table_route_manager()
234 .get_physical_table_route(table_id)
235 .await?;
236
237 self.data.region_distribution =
238 Some(region_distribution(&physical_table_route.region_routes));
239 let leaders = find_leaders(&physical_table_route.region_routes);
240 let alter_kind = self.make_region_alter_kind()?;
241
242 info!(
243 "Submitting alter region requests for table {}, table_id: {}, alter_kind: {:?}",
244 self.data.table_ref(),
245 table_id,
246 alter_kind,
247 );
248
249 ensure!(!leaders.is_empty(), NoLeaderSnafu { table_id });
250 self.put_poison(ctx_provider, procedure_id).await?;
252 let flow = self.data.flow();
253 if flow == AlterTableFlow::MetadataFirst {
254 let results = self
255 .executor
256 .on_alter_skip_wal_regions(
257 &self.context.node_manager,
258 &physical_table_route.region_routes,
259 alter_kind,
260 )
261 .await;
262
263 return match handle_multiple_results(results) {
264 MultipleResults::PartialRetryable(error) => Err(error),
265 MultipleResults::PartialNonRetryable(error)
266 | MultipleResults::AllNonRetryable(error) => {
267 Err(BoxedError::new(error)).context(RetryLaterSnafu {
270 clean_poisons: true,
271 })
272 }
273 MultipleResults::AllRetryable(error) => {
274 Err(BoxedError::new(error)).context(RetryLaterSnafu {
275 clean_poisons: true,
276 })
277 }
278 MultipleResults::Ok(results) => {
279 self.handle_alter_region_response(results)?;
280 Ok(Status::executing_with_clean_poisons(true))
281 }
282 };
283 }
284
285 let results = self
286 .executor
287 .on_alter_regions(
288 &self.context.node_manager,
289 &physical_table_route.region_routes,
290 alter_kind,
291 )
292 .await;
293
294 match handle_multiple_results(results) {
295 MultipleResults::PartialRetryable(error) => {
296 Err(error)
298 }
299 MultipleResults::PartialNonRetryable(error) => {
300 error!(error; "Partial non-retryable errors occurred during alter table, table {}, table_id: {}", self.data.table_ref(), self.data.table_id());
301 Ok(Status::poisoned(
303 Some(self.table_poison_key()),
304 ProcedureError::external(error),
305 ))
306 }
307 MultipleResults::AllRetryable(error) => {
308 let err = BoxedError::new(error);
310 Err(err).context(RetryLaterSnafu {
311 clean_poisons: true,
312 })
313 }
314 MultipleResults::Ok(results) => {
315 self.submit_sync_region_requests(&results, &physical_table_route.region_routes)
316 .await;
317 self.handle_alter_region_response(results)?;
318 Ok(Status::executing_with_clean_poisons(true))
319 }
320 MultipleResults::AllNonRetryable(error) => {
321 error!(error; "All alter requests returned non-retryable errors for table {}, table_id: {}", self.data.table_ref(), self.data.table_id());
322 let err = BoxedError::new(error);
326 Err(err).context(AbortProcedureSnafu {
327 clean_poisons: true,
328 })
329 }
330 }
331 }
332
333 fn handle_alter_region_response(&mut self, mut results: Vec<RegionResponse>) -> Result<()> {
334 if let Some(column_metadatas) =
335 extract_column_metadatas(&mut results, TABLE_COLUMN_METADATA_EXTENSION_KEY)?
336 {
337 self.data.column_metadatas = column_metadatas;
338 } else {
339 warn!(
340 "altering table result doesn't contains extension key `{TABLE_COLUMN_METADATA_EXTENSION_KEY}`,leaving the table's column metadata unchanged"
341 );
342 }
343 self.data.state = self.data.flow().after_regions();
344 Ok(())
345 }
346
347 async fn submit_sync_region_requests(
348 &mut self,
349 results: &[RegionResponse],
350 region_routes: &[RegionRoute],
351 ) {
352 let table_info = self.data.table_info().unwrap();
354 if let Err(err) = sync_follower_regions(
355 &self.context,
356 self.data.table_id(),
357 results,
358 region_routes,
359 table_info.meta.engine.as_str(),
360 )
361 .await
362 {
363 error!(err; "Failed to sync regions for table {}, table_id: {}", self.data.table_ref(), self.data.table_id());
364 }
365 }
366
367 pub(crate) async fn on_update_metadata(&mut self) -> Result<Status> {
369 let table_id = self.data.table_id();
370 let table_ref = self.data.table_ref();
371 let table_info_value = self.data.table_info_value.as_ref().unwrap();
373 let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
375 let flow = self.data.flow();
376 let metadata_only_alter = flow == AlterTableFlow::MetadataOnly;
377
378 let new_info = match &self.new_table_info {
380 Some(cached) => cached.clone(),
381 None => AlterTableExecutor::validate_alter_table_expr(
382 &table_info_value.table_info,
383 self.data.task.alter_table.clone(),
384 )
385 .inspect_err(|e| {
386 error!(e; "Unable to build info for table {} in update metadata step, table_id: {}", table_ref, table_id);
388 })?,
389 };
390
391 self.executor
393 .on_alter_metadata(
394 &self.context.table_metadata_manager,
395 table_info_value,
396 self.data.region_distribution.as_ref(),
397 new_info,
398 &self.data.column_metadatas,
399 metadata_only_alter,
400 )
401 .await?;
402
403 info!(
404 "Updated table metadata for table {table_ref}, table_id: {table_id}, kind: {alter_kind:?}"
405 );
406 self.data.state = flow.after_metadata();
407 Ok(Status::executing(true))
408 }
409
410 async fn on_broadcast(&mut self) -> Result<Status> {
412 self.executor
413 .invalidate_table_cache(&self.context.cache_invalidator)
414 .await?;
415 Ok(Status::done())
416 }
417
418 fn lock_key_inner(&self) -> Vec<StringKey> {
419 let mut lock_key = vec![];
420 let table_ref = self.data.table_ref();
421 let table_id = self.data.table_id();
422 lock_key.push(CatalogLock::Read(table_ref.catalog).into());
423 lock_key.push(SchemaLock::read(table_ref.catalog, table_ref.schema).into());
424 lock_key.push(TableLock::Write(table_id).into());
425
426 for region_id in &self.data.region_locks {
427 lock_key.push(RegionLock::Write(*region_id).into());
428 }
429
430 let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
432 if let Kind::RenameTable(RenameTable { new_table_name }) = alter_kind {
433 lock_key.push(
434 TableNameLock::new(table_ref.catalog, table_ref.schema, new_table_name).into(),
435 )
436 }
437
438 lock_key
439 }
440
441 #[cfg(test)]
442 pub(crate) fn data(&self) -> &AlterTableData {
443 &self.data
444 }
445
446 #[cfg(test)]
447 pub(crate) fn mut_data(&mut self) -> &mut AlterTableData {
448 &mut self.data
449 }
450}
451
452fn enables_skip_wal(alter_kind: &Kind) -> bool {
453 let Kind::SetTableOptions(SetTableOptions { table_options }) = alter_kind else {
454 return false;
455 };
456
457 table_options
458 .iter()
459 .any(|option| option.key == SKIP_WAL_KEY && option.value == "true")
460}
461
462pub(crate) fn only_enables_skip_wal(alter_kind: &Kind) -> bool {
463 let Kind::SetTableOptions(SetTableOptions { table_options }) = alter_kind else {
464 return false;
465 };
466
467 table_options.len() == 1
468 && table_options[0].key == SKIP_WAL_KEY
469 && table_options[0].value == "true"
470}
471
472fn is_metadata_only_alter(alter_kind: &Kind) -> bool {
473 match alter_kind {
474 Kind::RenameTable { .. } => true,
475 Kind::SetTableOptions(SetTableOptions { table_options }) => {
476 table_options.len() == 1 && table_options[0].key.as_str() == REPARTITION_COLUMN_HINT_KEY
477 }
478 Kind::UnsetTableOptions(UnsetTableOptions { keys }) => {
479 keys.len() == 1 && keys[0].as_str() == REPARTITION_COLUMN_HINT_KEY
480 }
481 _ => false,
482 }
483}
484
485#[derive(Debug, Clone, Copy, PartialEq, Eq)]
486enum AlterTableFlow {
487 RegionFirst,
488 MetadataFirst,
489 MetadataOnly,
490}
491
492impl AlterTableFlow {
493 fn from_kind(kind: &Kind) -> Self {
494 if only_enables_skip_wal(kind) {
495 Self::MetadataFirst
496 } else if is_metadata_only_alter(kind) {
497 Self::MetadataOnly
498 } else {
499 Self::RegionFirst
500 }
501 }
502
503 fn after_prepare(self) -> AlterTableState {
504 match self {
505 Self::RegionFirst => AlterTableState::SubmitAlterRegionRequests,
506 Self::MetadataFirst | Self::MetadataOnly => AlterTableState::UpdateMetadata,
507 }
508 }
509
510 fn after_regions(self) -> AlterTableState {
511 match self {
512 Self::RegionFirst => AlterTableState::UpdateMetadata,
513 Self::MetadataFirst | Self::MetadataOnly => AlterTableState::InvalidateTableCache,
514 }
515 }
516
517 fn after_metadata(self) -> AlterTableState {
518 match self {
519 Self::MetadataFirst => AlterTableState::SubmitAlterRegionRequests,
520 Self::RegionFirst | Self::MetadataOnly => AlterTableState::InvalidateTableCache,
521 }
522 }
523}
524
525#[async_trait]
526impl Procedure for AlterTableProcedure {
527 fn type_name(&self) -> &str {
528 Self::TYPE_NAME
529 }
530
531 async fn execute(&mut self, ctx: &ProcedureContext) -> ProcedureResult<Status> {
532 let state = &self.data.state;
533
534 let step = state.as_ref();
535
536 let _timer = metrics::METRIC_META_PROCEDURE_ALTER_TABLE
537 .with_label_values(&[step])
538 .start_timer();
539
540 match state {
541 AlterTableState::Prepare => self.on_prepare().await,
542 AlterTableState::SubmitAlterRegionRequests => {
543 self.submit_alter_region_requests(ctx.procedure_id, ctx.provider.as_ref())
544 .await
545 }
546 AlterTableState::UpdateMetadata => self.on_update_metadata().await,
547 AlterTableState::InvalidateTableCache => self.on_broadcast().await,
548 }
549 .map_err(map_to_procedure_error)
550 }
551
552 fn dump(&self) -> ProcedureResult<String> {
553 serde_json::to_string(&self.data).context(ToJsonSnafu)
554 }
555
556 fn lock_key(&self) -> LockKey {
557 let key = self.lock_key_inner();
558
559 LockKey::new(key)
560 }
561
562 fn poison_keys(&self) -> PoisonKeys {
563 PoisonKeys::new(vec![self.table_poison_key()])
564 }
565
566 fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
567 if !ctx
568 .event_type_filter
569 .allows(TableDdlEventType::AlterTable.as_str())
570 {
571 return None;
572 }
573 let table_ref = self.data.table_ref();
574 let locator = TableDdlLocator::new(table_ref.catalog, table_ref.schema, table_ref.table)
575 .with_table_id(self.data.table_id());
576 let event = match &ctx.trigger {
577 EventTrigger::Submitted => {
578 let kind = self
579 .data
580 .task
581 .alter_table
582 .kind
583 .as_ref()
584 .and_then(alter_table_kind_name);
585 TableDdlEvent::alter_table_submitted(locator, kind)
586 }
587 _ => TableDdlEvent::lifecycle(TableDdlEventType::AlterTable, [locator]),
588 };
589
590 Some(Box::new(event))
591 }
592}
593
594#[derive(Debug, Serialize, Deserialize, AsRefStr)]
595enum AlterTableState {
596 Prepare,
598 SubmitAlterRegionRequests,
600 UpdateMetadata,
602 InvalidateTableCache,
604}
605
606#[derive(Debug, Serialize, Deserialize)]
608pub struct AlterTableData {
609 state: AlterTableState,
610 task: AlterTableTask,
611 table_id: TableId,
612 #[serde(default)]
613 column_metadatas: Vec<ColumnMetadata>,
614 table_info_value: Option<DeserializedValueWithBytes<TableInfoValue>>,
616 region_distribution: Option<RegionDistribution>,
618 #[serde(default)]
620 region_locks: Vec<RegionId>,
621}
622
623impl AlterTableData {
624 pub fn new(task: AlterTableTask, table_id: TableId, region_locks: Vec<RegionId>) -> Self {
625 Self {
626 state: AlterTableState::Prepare,
627 task,
628 table_id,
629 column_metadatas: vec![],
630 table_info_value: None,
631 region_distribution: None,
632 region_locks,
633 }
634 }
635
636 fn table_ref(&self) -> TableReference<'_> {
637 self.task.table_ref()
638 }
639
640 fn table_id(&self) -> TableId {
641 self.table_id
642 }
643
644 fn table_info(&self) -> Option<&TableInfo> {
645 self.table_info_value
646 .as_ref()
647 .map(|value| &value.table_info)
648 }
649
650 fn flow(&self) -> AlterTableFlow {
651 AlterTableFlow::from_kind(self.task.alter_table.kind.as_ref().unwrap())
653 }
654
655 #[cfg(test)]
656 pub(crate) fn column_metadatas(&self) -> &[ColumnMetadata] {
657 &self.column_metadatas
658 }
659
660 #[cfg(test)]
661 pub(crate) fn set_column_metadatas(&mut self, column_metadatas: Vec<ColumnMetadata>) {
662 self.column_metadatas = column_metadatas;
663 }
664}