Skip to main content

risingwave_frontend/handler/
alter_set_schema.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
24// Steps for validation:
25// 1. Check permission to alter original object.
26// 2. Check duplicate name in the new schema.
27// 3. Check permission to create in the new schema.
28
29/// Handle `ALTER [TABLE | [MATERIALIZED] VIEW | SOURCE | SINK | CONNECTION | FUNCTION] <name> SET SCHEMA <schema_name>` statements.
30pub 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}