Skip to main content

uqa_storage/sqlite/catalog/
sequences_views.rs

1//
2// Unified Query Algebra
3//
4// Copyright (c) 2023-2026 Cognica, Inc.
5//
6
7//! Sequence and SQL view catalog state.
8
9use 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}