risingwave_frontend/handler/
alter_set_schema.rs1use pgwire::pg_response::StatementType;
16use risingwave_sqlparser::ast::{ObjectName, OperateFunctionArg};
17
18use super::alter_utils::validate_table_type_for_alter;
19use super::{HandlerArgs, RwPgResponse};
20use crate::catalog::root_catalog::SchemaPath;
21use crate::error::{ErrorCode, Result};
22use crate::{Binder, bind_data_type};
23
24pub async fn handle_alter_set_schema(
31 handler_args: HandlerArgs,
32 obj_name: ObjectName,
33 new_schema_name: ObjectName,
34 stmt_type: StatementType,
35 func_args: Option<Vec<OperateFunctionArg>>,
36) -> Result<RwPgResponse> {
37 let session = handler_args.session;
38 let db_name = &session.database();
39 let (schema_name, real_obj_name) = Binder::resolve_schema_qualified_name(db_name, &obj_name)?;
40 let search_path = session.config().search_path();
41 let user_name = &session.user_name();
42 let schema_path = SchemaPath::new(schema_name.as_deref(), &search_path, user_name);
43
44 let new_schema_name = Binder::resolve_schema_name(new_schema_name)?;
45 let object = {
46 let catalog_reader = session.env().catalog_reader().read_guard();
47
48 match stmt_type {
49 StatementType::ALTER_TABLE | StatementType::ALTER_MATERIALIZED_VIEW => {
50 let (table, old_schema_name) = catalog_reader.get_created_table_by_name(
51 db_name,
52 schema_path,
53 &real_obj_name,
54 )?;
55 validate_table_type_for_alter(table, stmt_type, "schema")?;
56 if old_schema_name == new_schema_name {
57 return Ok(RwPgResponse::empty_result(stmt_type));
58 }
59 session.check_privilege_for_drop_alter(old_schema_name, &**table)?;
60 catalog_reader.check_relation_name_duplicated(
61 db_name,
62 &new_schema_name,
63 table.name(),
64 )?;
65 table.id.into()
66 }
67 StatementType::ALTER_VIEW => {
68 let (view, old_schema_name) =
69 catalog_reader.get_view_by_name(db_name, schema_path, &real_obj_name)?;
70 if old_schema_name == new_schema_name {
71 return Ok(RwPgResponse::empty_result(stmt_type));
72 }
73 session.check_privilege_for_drop_alter(old_schema_name, &**view)?;
74 catalog_reader.check_relation_name_duplicated(
75 db_name,
76 &new_schema_name,
77 view.name(),
78 )?;
79 view.id.into()
80 }
81 StatementType::ALTER_SOURCE => {
82 let (source, old_schema_name) =
83 catalog_reader.get_source_by_name(db_name, schema_path, &real_obj_name)?;
84 if old_schema_name == new_schema_name {
85 return Ok(RwPgResponse::empty_result(stmt_type));
86 }
87 session.check_privilege_for_drop_alter(old_schema_name, &**source)?;
88 catalog_reader.check_relation_name_duplicated(
89 db_name,
90 &new_schema_name,
91 &source.name,
92 )?;
93 source.id.into()
94 }
95 StatementType::ALTER_SINK => {
96 let (sink, old_schema_name) = catalog_reader.get_created_sink_by_name(
97 db_name,
98 schema_path,
99 &real_obj_name,
100 )?;
101 if old_schema_name == new_schema_name {
102 return Ok(RwPgResponse::empty_result(stmt_type));
103 }
104 session.check_privilege_for_drop_alter(old_schema_name, &**sink)?;
105 catalog_reader.check_relation_name_duplicated(
106 db_name,
107 &new_schema_name,
108 &sink.name,
109 )?;
110 sink.id.into()
111 }
112 StatementType::ALTER_SUBSCRIPTION => {
113 let (subscription, old_schema_name) = catalog_reader.get_subscription_by_name(
114 db_name,
115 schema_path,
116 &real_obj_name,
117 )?;
118 if old_schema_name == new_schema_name {
119 return Ok(RwPgResponse::empty_result(stmt_type));
120 }
121 session.check_privilege_for_drop_alter(old_schema_name, &**subscription)?;
122 catalog_reader.check_relation_name_duplicated(
123 db_name,
124 &new_schema_name,
125 &subscription.name,
126 )?;
127 subscription.id.into()
128 }
129 StatementType::ALTER_CONNECTION => {
130 let (connection, old_schema_name) =
131 catalog_reader.get_connection_by_name(db_name, schema_path, &real_obj_name)?;
132 if old_schema_name == new_schema_name {
133 return Ok(RwPgResponse::empty_result(stmt_type));
134 }
135 session.check_privilege_for_drop_alter(old_schema_name, &**connection)?;
136 catalog_reader.check_connection_name_duplicated(
137 db_name,
138 &new_schema_name,
139 &connection.name,
140 )?;
141 connection.id.into()
142 }
143 StatementType::ALTER_FUNCTION => {
144 let (function, old_schema_name) = if let Some(args) = func_args {
145 let mut arg_types = Vec::with_capacity(args.len());
146 for arg in args {
147 arg_types.push(bind_data_type(&arg.data_type)?);
148 }
149 catalog_reader.get_function_by_name_args(
150 db_name,
151 schema_path,
152 &real_obj_name,
153 &arg_types,
154 )?
155 } else {
156 let (functions, old_schema_name) = catalog_reader.get_functions_by_name(
157 db_name,
158 schema_path,
159 &real_obj_name,
160 )?;
161 if functions.len() > 1 {
162 return Err(ErrorCode::CatalogError(format!("function name {real_obj_name:?} is not unique\nHINT: Specify the argument list to select the function unambiguously.").into()).into());
163 }
164 (
165 functions.into_iter().next().expect("no functions"),
166 old_schema_name,
167 )
168 };
169 if old_schema_name == new_schema_name {
170 return Ok(RwPgResponse::empty_result(stmt_type));
171 }
172 session.check_privilege_for_drop_alter(old_schema_name, &**function)?;
173 catalog_reader.check_function_name_duplicated(
174 db_name,
175 &new_schema_name,
176 &function.name,
177 &function.arg_types,
178 )?;
179 function.id.into()
180 }
181 _ => unreachable!(),
182 }
183 };
184
185 let (_, new_schema_id) =
186 session.get_database_and_schema_id_for_create(Some(new_schema_name))?;
187
188 let catalog_writer = session.catalog_writer()?;
189 catalog_writer
190 .alter_set_schema(object, new_schema_id)
191 .await?;
192
193 Ok(RwPgResponse::empty_result(stmt_type))
194}