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