Plan an `INSERT ... ON CONFLICT DO UPDATE SET` statement.
(
ins: &ast::Insert,
catalog: &dyn SqlCatalog,
on_conflict_updates: Vec<(String, SqlExpr)>,
)
| 234 | |
| 235 | /// Plan an `INSERT ... ON CONFLICT DO UPDATE SET` statement. |
| 236 | fn plan_upsert_with_on_conflict( |
| 237 | ins: &ast::Insert, |
| 238 | catalog: &dyn SqlCatalog, |
| 239 | on_conflict_updates: Vec<(String, SqlExpr)>, |
| 240 | ) -> Result<Vec<SqlPlan>> { |
| 241 | let table_name = match &ins.table { |
| 242 | ast::TableObject::TableName(name) => normalize_object_name_checked(name)?, |
| 243 | ast::TableObject::TableFunction(_) => { |
| 244 | return Err(SqlError::Unsupported { |
| 245 | detail: "INSERT ... ON CONFLICT on a table function is not supported".into(), |
| 246 | }); |
| 247 | } |
| 248 | }; |
| 249 | let info = catalog |
| 250 | .get_collection(DatabaseId::DEFAULT, &table_name)? |
| 251 | .ok_or_else(|| SqlError::UnknownTable { |
| 252 | name: table_name.clone(), |
| 253 | })?; |
| 254 | |
| 255 | let columns: Vec<String> = ins.columns.iter().map(normalize_ident).collect(); |
| 256 | |
| 257 | let source = ins.source.as_ref().ok_or_else(|| SqlError::Parse { |
| 258 | detail: "INSERT ... ON CONFLICT requires VALUES".into(), |
| 259 | })?; |
| 260 | let rows_ast = match &*source.body { |
| 261 | ast::SetExpr::Values(values) => &values.rows, |
| 262 | _ => { |
| 263 | return Err(SqlError::Unsupported { |
| 264 | detail: "INSERT ... ON CONFLICT source must be VALUES".into(), |
| 265 | }); |
| 266 | } |
| 267 | }; |
| 268 | |
| 269 | // KV: `INSERT ... ON CONFLICT (key) DO UPDATE SET ...` is an opt-in |
| 270 | // overwrite — same physical semantics as UPSERT, with the optional |
| 271 | // per-row assignments carried through for the Data Plane to apply |
| 272 | // against the existing row. |
| 273 | if info.engine == EngineType::KeyValue { |
| 274 | return build_kv_insert_plan( |
| 275 | table_name, |
| 276 | &columns, |
| 277 | rows_ast, |
| 278 | KvInsertIntent::Put, |
| 279 | on_conflict_updates, |
| 280 | info.primary_key.as_deref(), |
| 281 | ); |
| 282 | } |
| 283 | |
| 284 | let rows = convert_value_rows(&columns, rows_ast)?; |
| 285 | let column_defaults: Vec<(String, String)> = info |
| 286 | .columns |
| 287 | .iter() |
| 288 | .filter_map(|c| c.default.as_ref().map(|d| (c.name.clone(), d.clone()))) |
| 289 | .collect(); |
| 290 | let column_schema: Vec<(String, String)> = info |
| 291 | .columns |
| 292 | .iter() |
| 293 | .filter_map(|c| c.raw_type.as_ref().map(|t| (c.name.clone(), t.clone()))) |
no test coverage detected