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