1use super::{
10 migration_relation, params, Catalog, OptionalExtension, RelationIdentity, RelationKind, Result,
11 SQLiteError, SequenceOptions, SequenceReservationResult, SequenceRow, ViewRow,
12};
13use crate::catalog::{
14 sequence_value_reservation, SequenceOwner, SequenceOwnerDependency, SequenceValuePosition,
15};
16
17fn concrete_sequence_options(sequence: &SequenceRow) -> SequenceOptions {
18 let default_min = if sequence.increment > 0 { 1 } else { i64::MIN };
19 let default_max = if sequence.increment > 0 { i64::MAX } else { -1 };
20 SequenceOptions {
21 data_type: sequence.options.data_type.clone(),
22 min_value: Some(sequence.options.min_value.unwrap_or(default_min)),
23 max_value: Some(sequence.options.max_value.unwrap_or(default_max)),
24 cycle: sequence.options.cycle,
25 cache_size: sequence.options.cache_size,
26 }
27}
28
29fn decode_sequence_owner(
30 relation: &RelationIdentity,
31 table_object_id: Option<Vec<u8>>,
32 column_object_id: Option<Vec<u8>>,
33 dependency: Option<String>,
34) -> Result<Option<SequenceOwner>> {
35 let (table_object_id, column_object_id, dependency) =
36 match (table_object_id, column_object_id, dependency) {
37 (None, None, None) => return Ok(None),
38 (Some(table), Some(column), Some(dependency)) => (table, column, dependency),
39 _ => {
40 return Err(SQLiteError::StorageBackend(format!(
41 "corrupt sequence `{}` has an incomplete owner dependency",
42 relation.qualified_name()
43 )))
44 }
45 };
46 let table_object_id: [u8; 16] = table_object_id.try_into().map_err(|value: Vec<u8>| {
47 SQLiteError::StorageBackend(format!(
48 "corrupt sequence `{}` owner table identity has {} bytes",
49 relation.qualified_name(),
50 value.len()
51 ))
52 })?;
53 let column_object_id: [u8; 16] = column_object_id.try_into().map_err(|value: Vec<u8>| {
54 SQLiteError::StorageBackend(format!(
55 "corrupt sequence `{}` owner column identity has {} bytes",
56 relation.qualified_name(),
57 value.len()
58 ))
59 })?;
60 let dependency = match dependency.as_str() {
61 "a" => SequenceOwnerDependency::Automatic,
62 "i" => SequenceOwnerDependency::Internal,
63 other => {
64 return Err(SQLiteError::StorageBackend(format!(
65 "corrupt sequence `{}` owner dependency `{other}`",
66 relation.qualified_name()
67 )))
68 }
69 };
70 Ok(Some(SequenceOwner {
71 table_object_id,
72 column_object_id,
73 dependency,
74 }))
75}
76
77struct RawSequenceRow {
78 schema: String,
79 name: String,
80 object_id: Vec<u8>,
81 definition_generation: Vec<u8>,
82 start: i64,
83 increment: i64,
84 current: i64,
85 called: bool,
86 persistence: String,
87 data_type: String,
88 min_value: i64,
89 max_value: i64,
90 cycle: bool,
91 cache_size: i64,
92 owner_table_object_id: Option<Vec<u8>>,
93 owner_column_object_id: Option<Vec<u8>>,
94 owner_dependency: Option<String>,
95 role_owner: String,
96 acl_json: Option<String>,
97 log_count: i64,
98}
99
100fn read_raw_sequence_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<RawSequenceRow> {
101 Ok(RawSequenceRow {
102 schema: row.get(0)?,
103 name: row.get(1)?,
104 object_id: row.get(2)?,
105 definition_generation: row.get(3)?,
106 start: row.get(4)?,
107 increment: row.get(5)?,
108 current: row.get(6)?,
109 called: row.get(7)?,
110 persistence: row.get(8)?,
111 data_type: row.get(9)?,
112 min_value: row.get(10)?,
113 max_value: row.get(11)?,
114 cycle: row.get(12)?,
115 cache_size: row.get(13)?,
116 owner_table_object_id: row.get(14)?,
117 owner_column_object_id: row.get(15)?,
118 owner_dependency: row.get(16)?,
119 role_owner: row.get(17)?,
120 acl_json: row.get(18)?,
121 log_count: row.get(19)?,
122 })
123}
124
125fn decode_sequence_identity(
126 relation: &RelationIdentity,
127 label: &str,
128 value: Vec<u8>,
129) -> Result<[u8; 16]> {
130 value.try_into().map_err(|value: Vec<u8>| {
131 SQLiteError::StorageBackend(format!(
132 "corrupt sequence `{}` {label} has {} bytes",
133 relation.qualified_name(),
134 value.len()
135 ))
136 })
137}
138
139fn decode_raw_sequence_row(raw: RawSequenceRow) -> Result<SequenceRow> {
140 let relation = RelationIdentity::new(raw.schema, raw.name);
141 Ok(SequenceRow {
142 role_owner: raw.role_owner,
143 acl: raw
144 .acl_json
145 .map(|json| serde_json::from_str(&json))
146 .transpose()?,
147 owner: decode_sequence_owner(
148 &relation,
149 raw.owner_table_object_id,
150 raw.owner_column_object_id,
151 raw.owner_dependency,
152 )?,
153 object_id: decode_sequence_identity(&relation, "object identity", raw.object_id)?,
154 definition_generation: decode_sequence_identity(
155 &relation,
156 "definition generation",
157 raw.definition_generation,
158 )?,
159 relation,
160 start: raw.start,
161 increment: raw.increment,
162 current: raw.current,
163 called: raw.called,
164 log_count: raw.log_count,
165 persistence: raw.persistence,
166 options: SequenceOptions {
167 data_type: raw.data_type,
168 min_value: Some(raw.min_value),
169 max_value: Some(raw.max_value),
170 cycle: raw.cycle,
171 cache_size: raw.cache_size,
172 },
173 })
174}
175
176fn reserve_sequence_values_in_connection(
177 connection: &rusqlite::Connection,
178 relation: &RelationIdentity,
179 object_id: [u8; 16],
180 definition_generation: [u8; 16],
181) -> Result<SequenceReservationResult> {
182 let stored = connection
183 .query_row(
184 "SELECT object_id, definition_generation, current, called, increment, min_value, max_value, cycle, cache_size, log_count
185 FROM _sequences WHERE schema_name = ?1 AND relation_name = ?2",
186 params![relation.schema, relation.name],
187 |row| {
188 Ok((
189 row.get::<_, Vec<u8>>(0)?,
190 row.get::<_, Vec<u8>>(1)?,
191 row.get::<_, i64>(2)?,
192 row.get::<_, bool>(3)?,
193 row.get::<_, i64>(4)?,
194 row.get::<_, i64>(5)?,
195 row.get::<_, i64>(6)?,
196 row.get::<_, bool>(7)?,
197 row.get::<_, i64>(8)?,
198 row.get::<_, i64>(9)?,
199 ))
200 },
201 )
202 .optional()?;
203 let Some((
204 stored_object_id,
205 stored_generation,
206 current,
207 called,
208 increment,
209 min,
210 max,
211 cycle,
212 cache_size,
213 log_count,
214 )) = stored
215 else {
216 return Ok(SequenceReservationResult::Missing);
217 };
218 let stored_object_id: [u8; 16] = stored_object_id.try_into().map_err(|value: Vec<u8>| {
219 SQLiteError::StorageBackend(format!(
220 "corrupt sequence `{}` object identity has {} bytes",
221 relation.qualified_name(),
222 value.len()
223 ))
224 })?;
225 if stored_object_id != object_id {
226 return Ok(SequenceReservationResult::Missing);
227 }
228 let stored_generation: [u8; 16] = stored_generation.try_into().map_err(|value: Vec<u8>| {
229 SQLiteError::StorageBackend(format!(
230 "corrupt sequence `{}` definition generation has {} bytes",
231 relation.qualified_name(),
232 value.len()
233 ))
234 })?;
235 if stored_generation != definition_generation {
236 return Ok(SequenceReservationResult::DefinitionChanged);
237 }
238 if increment == 0 || cache_size <= 0 {
239 return Err(SQLiteError::StorageBackend(format!(
240 "corrupt sequence `{}` has increment {increment} and cache size {cache_size}",
241 relation.qualified_name()
242 )));
243 }
244 let Some(reservation) = sequence_value_reservation(
245 SequenceValuePosition {
246 current,
247 called,
248 log_count,
249 },
250 increment,
251 min,
252 max,
253 cycle,
254 cache_size,
255 ) else {
256 return Ok(SequenceReservationResult::Exhausted);
257 };
258 let updated = connection.execute(
259 "UPDATE _sequences SET current = ?5, called = 1, log_count = ?6
260 WHERE schema_name = ?1 AND relation_name = ?2 AND object_id = ?3 AND definition_generation = ?4",
261 params![
262 relation.schema,
263 relation.name,
264 object_id.as_slice(),
265 definition_generation.as_slice(),
266 reservation.last_value,
267 reservation.log_count,
268 ],
269 )?;
270 if updated != 1 {
271 return Err(SQLiteError::StorageBackend(format!(
272 "sequence `{}` changed while reserving cached values",
273 relation.qualified_name()
274 )));
275 }
276 Ok(SequenceReservationResult::Reserved(reservation))
277}
278
279impl Catalog {
280 pub fn create_sequence_row(&self, sequence: &SequenceRow) -> Result<bool> {
281 self.conn.with_mut(|connection| {
282 let tx = connection.savepoint()?;
283 let exists = tx
284 .query_row(
285 "SELECT 1 FROM _sequences
286 WHERE schema_name = ?1 AND relation_name = ?2",
287 params![sequence.relation.schema, sequence.relation.name],
288 |_| Ok(()),
289 )
290 .optional()?
291 .is_some();
292 if exists {
293 return Ok(false);
294 }
295 Self::claim_relation(&tx, &sequence.relation, RelationKind::Sequence)?;
296 let options = concrete_sequence_options(sequence);
297 let owner_table = sequence.owner.map(|owner| owner.table_object_id);
298 let owner_column = sequence.owner.map(|owner| owner.column_object_id);
299 let owner_dependency = sequence
300 .owner
301 .map(|owner| owner.dependency.catalog_code());
302 let acl_json = sequence
303 .acl
304 .as_ref()
305 .map(serde_json::to_string)
306 .transpose()?;
307 tx.execute(
308 "INSERT INTO _sequences
309 (schema_name, relation_name, kind, object_id, definition_generation, start, increment, current, called, persistence, data_type, min_value, max_value, cycle, cache_size, owner_table_object_id, owner_column_object_id, owner_dependency, role_owner, acl_json, log_count)
310 VALUES (?1, ?2, 'sequence', ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20)",
311 params![
312 sequence.relation.schema,
313 sequence.relation.name,
314 sequence.object_id.as_slice(),
315 sequence.definition_generation.as_slice(),
316 sequence.start,
317 sequence.increment,
318 sequence.current,
319 sequence.called,
320 sequence.persistence,
321 options.data_type,
322 options.min_value,
323 options.max_value,
324 options.cycle,
325 options.cache_size,
326 owner_table.as_ref().map(<[u8; 16]>::as_slice),
327 owner_column.as_ref().map(<[u8; 16]>::as_slice),
328 owner_dependency,
329 sequence.role_owner,
330 acl_json,
331 sequence.log_count,
332 ],
333 )?;
334 tx.commit()?;
335 Ok(true)
336 })
337 }
338
339 pub fn replace_sequence_row(&self, sequence: &SequenceRow) -> Result<bool> {
340 self.conn.with(|connection| {
341 let options = concrete_sequence_options(sequence);
342 let owner_table = sequence.owner.map(|owner| owner.table_object_id);
343 let owner_column = sequence.owner.map(|owner| owner.column_object_id);
344 let owner_dependency = sequence
345 .owner
346 .map(|owner| owner.dependency.catalog_code());
347 let acl_json = sequence
348 .acl
349 .as_ref()
350 .map(serde_json::to_string)
351 .transpose()?;
352 Ok(connection.execute(
353 "UPDATE _sequences
354 SET object_id = ?3, definition_generation = ?4, start = ?5, increment = ?6, current = ?7, called = ?8, persistence = ?9,
355 data_type = ?10, min_value = ?11, max_value = ?12, cycle = ?13, cache_size = ?14,
356 owner_table_object_id = ?15, owner_column_object_id = ?16, owner_dependency = ?17, role_owner = ?18, acl_json = ?19, log_count = ?20
357 WHERE schema_name = ?1 AND relation_name = ?2",
358 params![
359 sequence.relation.schema,
360 sequence.relation.name,
361 sequence.object_id.as_slice(),
362 sequence.definition_generation.as_slice(),
363 sequence.start,
364 sequence.increment,
365 sequence.current,
366 sequence.called,
367 sequence.persistence,
368 options.data_type,
369 options.min_value,
370 options.max_value,
371 options.cycle,
372 options.cache_size,
373 owner_table.as_ref().map(<[u8; 16]>::as_slice),
374 owner_column.as_ref().map(<[u8; 16]>::as_slice),
375 owner_dependency,
376 sequence.role_owner,
377 acl_json,
378 sequence.log_count,
379 ],
380 )? != 0)
381 })
382 }
383
384 pub fn rename_sequence_row(&self, from: &str, to: &str) -> Result<bool> {
385 let from_relation = migration_relation(from)?;
386 let to_relation = migration_relation(to)?;
387 self.conn.with_mut(|connection| {
388 let tx = connection.savepoint()?;
389 let source_exists = tx
390 .query_row(
391 "SELECT 1 FROM _sequences
392 WHERE schema_name = ?1 AND relation_name = ?2",
393 params![from_relation.schema, from_relation.name],
394 |_| Ok(()),
395 )
396 .optional()?
397 .is_some();
398 if !source_exists {
399 return Ok(false);
400 }
401 if from_relation == to_relation {
402 return Ok(true);
403 }
404 let target_kind = tx
405 .query_row(
406 "SELECT kind FROM _relations
407 WHERE schema_name = ?1 AND relation_name = ?2",
408 params![to_relation.schema, to_relation.name],
409 |row| row.get::<_, String>(0),
410 )
411 .optional()?;
412 if let Some(kind) = target_kind {
413 return Err(SQLiteError::StorageBackend(format!(
414 "relation `{}` already exists as {kind}",
415 to_relation.qualified_name()
416 )));
417 }
418 Self::claim_relation(&tx, &to_relation, RelationKind::Sequence)?;
419 let updated = tx.execute(
420 "UPDATE _sequences
421 SET schema_name = ?3, relation_name = ?4
422 WHERE schema_name = ?1 AND relation_name = ?2",
423 params![
424 from_relation.schema,
425 from_relation.name,
426 to_relation.schema,
427 to_relation.name
428 ],
429 )?;
430 if updated != 1 {
431 return Err(SQLiteError::StorageBackend(format!(
432 "sequence `{from}` changed while renaming"
433 )));
434 }
435 Self::release_relation(&tx, &from_relation, RelationKind::Sequence)?;
436 tx.commit()?;
437 Ok(true)
438 })
439 }
440
441 pub fn drop_sequence_row(&self, name: &str) -> Result<bool> {
442 let relation = migration_relation(name)?;
443 self.conn.with_mut(|connection| {
444 let tx = connection.savepoint()?;
445 let removed = tx.execute(
446 "DELETE FROM _sequences
447 WHERE schema_name = ?1 AND relation_name = ?2",
448 params![relation.schema, relation.name],
449 )? != 0;
450 if removed {
451 Self::release_relation(&tx, &relation, RelationKind::Sequence)?;
452 }
453 tx.commit()?;
454 Ok(removed)
455 })
456 }
457
458 pub fn load_sequence_rows(&self) -> Result<Vec<SequenceRow>> {
459 self.conn.with(|connection| {
460 let mut statement = connection.prepare(
461 "SELECT schema_name, relation_name, object_id, definition_generation, start, increment, current, called, persistence,
462 data_type, min_value, max_value, cycle, cache_size, owner_table_object_id, owner_column_object_id, owner_dependency, role_owner, acl_json, log_count
463 FROM _sequences ORDER BY schema_name, relation_name",
464 )?;
465 let sequences = statement
466 .query_map([], read_raw_sequence_row)?
467 .map(|row| decode_raw_sequence_row(row?))
468 .collect();
469 sequences
470 })
471 }
472
473 pub fn reserve_sequence_values(
474 &self,
475 name: &str,
476 object_id: [u8; 16],
477 definition_generation: [u8; 16],
478 ) -> Result<SequenceReservationResult> {
479 let relation = migration_relation(name)?;
480 self.conn.with_mut(|connection| {
481 if connection.is_autocommit() {
482 let tx = connection
483 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
484 let result = reserve_sequence_values_in_connection(
485 &tx,
486 &relation,
487 object_id,
488 definition_generation,
489 )?;
490 tx.commit()?;
491 return Ok(result);
492 }
493 let tx = connection.savepoint()?;
494 let result = reserve_sequence_values_in_connection(
495 &tx,
496 &relation,
497 object_id,
498 definition_generation,
499 )?;
500 tx.commit()?;
501 Ok(result)
502 })
503 }
504
505 pub fn set_sequence_value(
506 &self,
507 name: &str,
508 object_id: [u8; 16],
509 value: i64,
510 called: bool,
511 log_count: i64,
512 ) -> Result<Option<i64>> {
513 let relation = migration_relation(name)?;
514 self.conn.with(|connection| {
515 Ok(connection
516 .query_row(
517 "UPDATE _sequences SET current = ?4, called = ?5, log_count = ?6
518 WHERE schema_name = ?1 AND relation_name = ?2 AND object_id = ?3 RETURNING current",
519 params![
520 relation.schema,
521 relation.name,
522 object_id.as_slice(),
523 value,
524 called,
525 log_count,
526 ],
527 |row| row.get(0),
528 )
529 .optional()?)
530 })
531 }
532
533 pub fn save_view(&self, view: &ViewRow) -> Result<()> {
534 self.conn.with_mut(|connection| {
535 let tx = connection.savepoint()?;
536 Self::claim_relation(&tx, &view.relation, RelationKind::View)?;
537 let acl_json = view.acl.as_ref().map(serde_json::to_string).transpose()?;
538 let column_acls_json = serde_json::to_string(&view.column_acls)?;
539 tx.execute(
540 "INSERT OR REPLACE INTO _views
541 (schema_name, relation_name, kind, role_owner, acl_json, column_acls_json, definition_json)
542 VALUES (?1, ?2, 'view', ?3, ?4, ?5, ?6)",
543 params![
544 view.relation.schema,
545 view.relation.name,
546 view.role_owner,
547 acl_json,
548 column_acls_json,
549 view.definition_json
550 ],
551 )?;
552 tx.commit()?;
553 Ok(())
554 })
555 }
556
557 pub fn rename_view(&self, from: &RelationIdentity, to: &RelationIdentity) -> Result<bool> {
558 if from.schema != to.schema {
559 return Err(SQLiteError::StorageBackend(
560 "moving a view between schemas is not supported by the catalog".into(),
561 ));
562 }
563 self.conn.with_mut(|connection| {
564 let source_exists = connection.query_row(
565 "SELECT EXISTS(SELECT 1 FROM _views WHERE schema_name = ?1 AND relation_name = ?2)",
566 params![from.schema, from.name],
567 |row| row.get::<_, bool>(0),
568 )?;
569 if from == to || !source_exists {
570 return Ok(source_exists);
571 }
572 let target_exists = connection.query_row(
573 "SELECT EXISTS(SELECT 1 FROM _relations WHERE schema_name = ?1 AND relation_name = ?2)",
574 params![to.schema, to.name],
575 |row| row.get::<_, bool>(0),
576 )?;
577 if target_exists {
578 return Err(SQLiteError::StorageBackend(format!(
579 "relation `{}` already exists",
580 to.qualified_name()
581 )));
582 }
583 let tx = connection.savepoint()?;
584 Self::claim_relation(&tx, to, RelationKind::View)?;
585 let updated = tx.execute(
586 "UPDATE _views SET schema_name = ?3, relation_name = ?4 WHERE schema_name = ?1 AND relation_name = ?2",
587 params![from.schema, from.name, to.schema, to.name],
588 )?;
589 if updated != 1 {
590 return Err(SQLiteError::StorageBackend(format!(
591 "view `{}` disappeared during rename",
592 from.qualified_name()
593 )));
594 }
595 Self::release_relation(&tx, from, RelationKind::View)?;
596 tx.commit()?;
597 Ok(true)
598 })
599 }
600
601 pub fn drop_view(&self, relation: &RelationIdentity) -> Result<bool> {
602 self.conn.with_mut(|connection| {
603 let tx = connection.savepoint()?;
604 let removed = tx.execute(
605 "DELETE FROM _views WHERE schema_name = ?1 AND relation_name = ?2",
606 params![relation.schema, relation.name],
607 )? != 0;
608 if removed {
609 Self::release_relation(&tx, relation, RelationKind::View)?;
610 }
611 tx.commit()?;
612 Ok(removed)
613 })
614 }
615
616 pub fn load_views(&self) -> Result<Vec<ViewRow>> {
617 self.conn.with(|connection| {
618 let mut statement = connection.prepare(
619 "SELECT schema_name, relation_name, role_owner, acl_json, column_acls_json, definition_json
620 FROM _views ORDER BY schema_name, relation_name",
621 )?;
622 let rows = statement.query_map([], |row| {
623 Ok((
624 row.get::<_, String>(0)?,
625 row.get::<_, String>(1)?,
626 row.get::<_, String>(2)?,
627 row.get::<_, Option<String>>(3)?,
628 row.get::<_, Option<String>>(4)?,
629 row.get::<_, String>(5)?,
630 ))
631 })?;
632 let mut views = Vec::new();
633 for row in rows {
634 let (schema, name, role_owner, acl_json, column_acls_json, definition_json) = row?;
635 views.push(ViewRow {
636 relation: RelationIdentity::new(schema, name),
637 role_owner,
638 acl: acl_json
639 .as_deref()
640 .map(serde_json::from_str)
641 .transpose()?,
642 column_acls: column_acls_json
643 .as_deref()
644 .map(serde_json::from_str)
645 .transpose()?
646 .unwrap_or_default(),
647 definition_json,
648 });
649 }
650 Ok(views)
651 })
652 }
653}