pg_tviews 0.1.0-beta.26

Transactional materialized views with incremental refresh for PostgreSQL
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
//! Administrative SQL functions: refresh, migration, schema analysis, cascade path.

use crate::{TViewError, TViewResult, utils::quote_identifier};
use pgrx::JsonB;
use pgrx::datum::DatumWithOid;
use pgrx::prelude::*;

/// Analyze a SELECT statement and return inferred TVIEW schema as JSONB
///
/// Returns a JSON object with schema details on success, or `{"error": "..."}` on
/// failure. Never raises a `PostgreSQL` error so callers can use the result in
/// expressions (e.g., `IS NOT NULL`, `->>'error'`).
#[pg_extern]
fn pg_tviews_analyze_select(sql: &str) -> JsonB {
    match crate::schema::inference::infer_schema(sql) {
        Ok(schema) => match schema.to_jsonb() {
            Ok(jsonb) => jsonb,
            Err(e) => {
                JsonB(serde_json::json!({"error": format!("Failed to serialize schema: {e}")}))
            }
        },
        Err(e) => JsonB(serde_json::json!({"error": e.to_string()})),
    }
}

/// Infer column types from `PostgreSQL` catalog
#[pg_extern]
#[allow(clippy::needless_pass_by_value)] // Reason: pgrx #[pg_extern] requires Vec by value
fn pg_tviews_infer_types(table_name: &str, columns: Vec<String>) -> JsonB {
    match crate::schema::types::infer_column_types(table_name, &columns) {
        Ok(types) => match serde_json::to_value(&types) {
            Ok(json_value) => JsonB(json_value),
            Err(e) => {
                error!("Failed to serialize types to JSONB: {}", e);
            }
        },
        Err(e) => {
            error!("Type inference failed: {}", e);
        }
    }
}

/// Rebuild a TVIEW from its backing view, then every TVIEW whose view reads it,
/// directly or through others, in dependency order: a manual repair leaves
/// nothing stale.
///
/// Each rebuild is a `TRUNCATE` and an `INSERT … SELECT` with an explicit column
/// list (the view's own columns, so not the table-only `created_at`/`updated_at`),
/// and holds an ACCESS EXCLUSIVE lock on that TVIEW until the transaction ends.
/// `entity` is rebuilt with the caller's privileges; the TVIEWs that read it are
/// rebuilt as their owners, as a write's cascade is (issue #136).
///
/// # Errors
/// Returns error if the entity is not registered, the dependency graph cannot be
/// loaded, or a rebuild fails.
#[pg_extern]
fn pg_tviews_refresh(entity: &str) -> TViewResult<()> {
    crate::revision::check();
    rebuild_with_dependents(&[entity.to_string()], true)?;
    Ok(())
}

/// Bring the TVIEWs whose definitions read the current time up to date (#193):
/// `tview`, or every such TVIEW the caller owns (or may change as a member of its
/// owner's role). Each is refreshed in full as a write to a `full_refresh` table
/// would refresh it, then the TVIEWs reading it are, through the flush. For
/// `pg_cron` or the application to call at the boundary its rows depend on (the
/// day, for `CURRENT_DATE`). Returns the TVIEWs refreshed, dependencies first.
///
/// # Errors
/// Returns an error if `tview` is not a TVIEW, reads no time or is not the
/// caller's, or a refresh fails.
#[pg_extern]
fn pg_tviews_refresh_time_dependent(
    tview: default!(Option<&str>, "NULL"),
) -> Result<SetOfIterator<'static, String>, TViewError> {
    crate::revision::check();
    let rows: Vec<(String, pgrx::pg_sys::Oid, bool, bool)> = Spi::connect(|client| {
        let mut rows = Vec::new();
        for row in client.select(
            &format!(
                "SELECT m.entity::text, m.table_oid::oid, m.time_dependent, \
                        pg_catalog.pg_has_role(c.relowner, 'USAGE') \
                 FROM {} m JOIN pg_catalog.pg_class c ON c.oid = m.table_oid \
                 WHERE $1::text IS NULL \
                    OR m.table_oid::oid = $1::text::pg_catalog.regclass::pg_catalog.oid",
                crate::utils::meta_table()
            ),
            None,
            // SAFETY: the datum borrows `tview`, which outlives the select.
            &[unsafe { DatumWithOid::new(tview, PgOid::BuiltIn(PgBuiltInOids::TEXTOID).value()) }],
        )? {
            if let (Some(entity), Some(table)) =
                (row.get::<String>(1)?, row.get::<pgrx::pg_sys::Oid>(2)?)
            {
                rows.push((
                    entity,
                    table,
                    row.get::<bool>(3)?.unwrap_or(false),
                    row.get::<bool>(4)?.unwrap_or(false),
                ));
            }
        }
        Ok::<_, pgrx::spi::Error>(rows)
    })
    .map_err(|e| TViewError::CatalogError {
        operation: "Find the time-dependent TVIEWs".to_string(),
        pg_error: e.to_string(),
    })?;
    let chosen: Vec<(String, pgrx::pg_sys::Oid)> = match tview {
        Some(name) => {
            let Some((entity, table, dependent, _)) = rows.into_iter().next() else {
                return Err(TViewError::InvalidInput {
                    parameter: "tview".to_string(),
                    reason: format!("{name} is not a TVIEW"),
                });
            };
            if !dependent {
                return Err(TViewError::InvalidInput {
                    parameter: "tview".to_string(),
                    reason: format!("{name} does not read the time: nothing to refresh"),
                });
            }
            crate::owner::require_owner(table, name)?;
            vec![(entity, table)]
        }
        None => rows
            .into_iter()
            .filter(|(_, _, dependent, owned)| *dependent && *owned)
            .map(|(entity, table, _, _)| (entity, table))
            .collect(),
    };
    let order = crate::queue::graph::EntityDepGraph::load()?.topo_order;
    let mut chosen = chosen;
    chosen.sort_by_key(|(entity, _)| order.iter().position(|e| e == entity).unwrap_or(usize::MAX));
    let mut refreshed = Vec::new();
    for (entity, table) in &chosen {
        if crate::config::suspend_triggers() || crate::suspend::is_suspended() {
            crate::suspend::record_change(entity);
        } else {
            crate::queue::enqueue_refresh_all(entity);
        }
        refreshed.push(crate::utils::qualified_relname_from_oid(*table)?);
    }
    crate::queue::flush_refresh_queue()?;
    Ok(SetOfIterator::new(refreshed))
}

/// Rebuild `entities` and every TVIEW whose view reads one of them, transitively
/// (the readers of a TVIEW come from the complete dependency relation, not the
/// pruned one flush-time propagation follows), dependencies first. Readers are
/// rebuilt as their owners; `entities` too unless `requested_as_caller`. Returns
/// the rebuilt entities in order.
///
/// # Errors
/// Returns an error if the dependency graph cannot be loaded or a rebuild fails.
pub fn rebuild_with_dependents(
    entities: &[String],
    requested_as_caller: bool,
) -> TViewResult<Vec<String>> {
    let graph = crate::queue::graph::EntityDepGraph::load()?;
    let mut readers: std::collections::HashMap<&str, Vec<&str>> = std::collections::HashMap::new();
    for (reader, read) in &graph.children {
        for entity in read {
            readers.entry(entity).or_default().push(reader);
        }
    }
    let mut order: Vec<String> = Vec::new();
    let mut pending: std::collections::VecDeque<&str> =
        entities.iter().map(String::as_str).collect();
    while let Some(entity) = pending.pop_front() {
        if order.iter().any(|e| e == entity) {
            continue;
        }
        pending.extend(readers.get(entity).into_iter().flatten());
        order.push(entity.to_string());
    }
    order.sort_by_key(|e| {
        graph
            .topo_order
            .iter()
            .position(|t| t == e)
            .unwrap_or(usize::MAX)
    });
    for entity in &order {
        if requested_as_caller && entities.contains(entity) {
            rebuild_one(entity)?;
        } else {
            let _owner = crate::owner::AsOwner::of_entity(entity)?;
            rebuild_one(entity)?;
        }
    }
    flush_after_rebuilds()?;
    Ok(order)
}

/// Refresh what rebuilds queued before returning (#202): rewriting a `tv_*` table
/// that another TVIEW reads queues that reader's rows (#191). No statement-level
/// flush trigger follows a `SELECT` of a refresh function, so the work would
/// otherwise reach `COMMIT` still queued.
///
/// # Errors
/// Returns an error if the refresh fails.
pub fn flush_after_rebuilds() -> TViewResult<()> {
    crate::queue::flush_refresh_queue()
}

/// Rebuild one TVIEW from its backing view (`TRUNCATE` + `INSERT … SELECT`), with
/// the current role's privileges, and nothing that reads it.
///
/// # Errors
/// Returns error if the entity is not registered or the truncate/insert fails.
pub fn rebuild_one(entity: &str) -> TViewResult<()> {
    let _pin = crate::owner::RenderPin::new();
    let (qi_tv, insert) = rebuild_statements(entity)?;
    Spi::run(&format!("TRUNCATE {qi_tv}"))?;
    Spi::run(&insert)?;
    Ok(())
}

/// Populate an **empty** `tv_<entity>` from its backing view without `TRUNCATE`.
///
/// Used when a TVIEW is found empty while its view is not (an UNLOGGED table reset
/// by a crash restart or promotion, or a TVIEW created empty). Unlike
/// [`rebuild_one`] it takes only a ROW EXCLUSIVE lock, so readers are never
/// blocked, even when the transaction stays prepared (2PC) for a while.
///
/// # Errors
/// Returns error if the entity is not registered or the insert fails.
pub fn fill_empty_tview(entity: &str) -> TViewResult<()> {
    let _pin = crate::owner::RenderPin::new();
    let (_, insert) = rebuild_statements(entity)?;
    Spi::run(&insert)?;
    Ok(())
}

/// The schema-qualified TVIEW table and the `INSERT … SELECT` that fills it from
/// its backing view. The explicit column list comes from the view's own
/// columns, which excludes the table-only `created_at`/`updated_at` columns.
fn rebuild_statements(entity: &str) -> TViewResult<(String, String)> {
    use crate::catalog::TviewMeta;

    let meta = TviewMeta::load_by_entity(entity)?.ok_or_else(|| TViewError::MetadataNotFound {
        entity: entity.to_string(),
    })?;
    let qi_tv = crate::utils::qualified_relname_from_oid(meta.tview_oid)?;
    let qi_view = crate::utils::qualified_relname_from_oid(meta.view_oid)?;
    let view_columns = crate::utils::get_view_columns_by_oid(meta.view_oid)?;
    if view_columns.is_empty() {
        return Err(TViewError::CatalogError {
            operation: format!("Get columns for view {qi_view}"),
            pg_error: "View has no selectable columns".to_string(),
        });
    }
    let col_list = view_columns
        .iter()
        .map(|c| quote_identifier(c))
        .collect::<Vec<_>>()
        .join(", ");
    let insert = format!("INSERT INTO {qi_tv} ({col_list}) SELECT {col_list} FROM {qi_view}");
    Ok((qi_tv, insert))
}

/// Create the propagation indexes `(fk_<x>, pk_<entity>)` that TVIEWs created
/// before they became part of TVIEW creation are missing.
///
/// Cascade propagation looks parent rows up by their integer `fk_*` columns;
/// without an index that lookup scans the whole TVIEW. An `fk_*` column counts
/// as covered when **any** index on the TVIEW leads with it, so user-created
/// indexes are respected.
///
/// Returns the DDL for each missing index: executed, or only reported when
/// `dry_run` is true. Idempotent: a second call returns no rows. On large
/// TVIEWs run the reported statements by hand with `CREATE INDEX CONCURRENTLY`
/// (which cannot run inside a function).
///
/// Usage:
///   `SELECT * FROM pg_tviews_ensure_propagation_indexes();`         -- all TVIEWs
///   `SELECT * FROM pg_tviews_ensure_propagation_indexes('post');`   -- one entity
///   `SELECT * FROM pg_tviews_ensure_propagation_indexes(NULL, true);` -- dry run
///
/// # Errors
/// Returns error if the catalog query or an index creation fails.
#[pg_extern]
fn pg_tviews_ensure_propagation_indexes(
    entity: default!(Option<&str>, "NULL"),
    dry_run: default!(bool, false),
) -> Result<SetOfIterator<'static, String>, TViewError> {
    crate::revision::check();
    let missing = Spi::connect(|client| {
        let args = vec![unsafe {
            DatumWithOid::new(entity, PgOid::BuiltIn(PgBuiltInOids::TEXTOID).value())
        }];
        let rows = client.select(
            &format!(
                "SELECT n.nspname::text, c.relname::text, a.attname::text, 'pk_' || m.entity \
             FROM {meta} m \
             JOIN pg_class c ON c.oid = m.table_oid \
             JOIN pg_namespace n ON n.oid = c.relnamespace \
             JOIN pg_attribute a ON a.attrelid = c.oid AND a.attnum > 0 AND NOT a.attisdropped \
             WHERE ($1::text IS NULL OR m.entity = $1) \
               AND a.attname LIKE 'fk\\_%' \
               AND a.atttypid IN ('int2'::regtype, 'int4'::regtype, 'int8'::regtype) \
               AND a.attname <> 'pk_' || m.entity \
               AND EXISTS (SELECT 1 FROM pg_attribute p \
                           WHERE p.attrelid = c.oid AND p.attname = 'pk_' || m.entity \
                             AND NOT p.attisdropped) \
               AND NOT EXISTS (SELECT 1 FROM pg_index i \
                               WHERE i.indrelid = c.oid AND i.indkey[0] = a.attnum) \
             ORDER BY 1, 2, 3",
                meta = crate::utils::meta_table()
            ),
            None,
            &args,
        )?;
        let mut ddl = Vec::new();
        for row in rows {
            if let (Some(schema), Some(table), Some(fk), Some(pk)) = (
                row[1].value::<String>()?,
                row[2].value::<String>()?,
                row[3].value::<String>()?,
                row[4].value::<String>()?,
            ) {
                ddl.push(crate::ddl::create::propagation_index_ddl(
                    &schema, &table, &fk, &pk,
                ));
            }
        }
        Ok::<_, spi::Error>(ddl)
    })?;

    if !dry_run {
        for ddl in &missing {
            Spi::run(ddl)?;
        }
    }
    Ok(SetOfIterator::new(missing))
}

/// Refresh all TVIEWs in the database, dependencies first.
/// This is a convenience function for bulk operations like schema migrations
/// or data seeding workflows.
///
/// # Errors
/// Returns error if any TVIEW cannot be refreshed
#[pg_extern]
fn pg_tviews_refresh_all_entities() -> TViewResult<()> {
    crate::revision::check();
    let order = refresh_all_in_dependency_order()?;
    if order.is_empty() {
        info!("No TVIEWs found to refresh");
    } else {
        info!("Successfully refreshed {} TVIEWs", order.len());
    }
    Ok(())
}

/// Rebuild every TVIEW from its backing view in dependency order: a TVIEW whose
/// view reads another `tv_*` table is rebuilt after it. Returns the entities in
/// the order they were rebuilt.
///
/// # Errors
/// Returns error if the dependency graph cannot be loaded or a rebuild fails.
pub fn refresh_all_in_dependency_order() -> TViewResult<Vec<String>> {
    let graph = crate::queue::graph::EntityDepGraph::load()?;
    for entity in &graph.topo_order {
        rebuild_one(entity)?;
    }
    flush_after_rebuilds()?;
    Ok(graph.topo_order)
}

/// Migrate all existing TVIEW triggers from the old PL/pgSQL handler to the
/// Rust `pg_tview_trigger_handler()`.
///
/// Call this once after upgrading `pg_tviews` to convert triggers installed by
/// prior versions. The operation is idempotent and safe to re-run.
///
/// Raises a `PostgreSQL` ERROR if any trigger cannot be migrated.
#[pg_extern]
fn pg_tviews_migrate_triggers() {
    crate::revision::check();
    if let Err(e) = crate::dependency::triggers::migrate_all_triggers_to_rust_handler() {
        error!("Failed to migrate triggers: {:?}", e);
    }
}

/// Show cascade dependency path for a given entity
///
/// Returns the dependency chain showing which TVIEWs depend on this entity
#[pg_extern]
fn pg_tviews_show_cascade_path(
    entity: &str,
) -> TableIterator<
    'static,
    (
        name!(depth, i32),
        name!(entity_name, String),
        name!(depends_on, String),
    ),
> {
    crate::revision::check();
    let results = Spi::connect(|client| {
        let args = vec![unsafe {
            DatumWithOid::new(entity, PgOid::BuiltIn(PgBuiltInOids::TEXTOID).value())
        }];
        match client.select(
            &format!(
                "WITH RECURSIVE dep_tree AS (
                SELECT
                    pg_tview_meta.entity,
                    0 as depth,
                    ARRAY[pg_tview_meta.entity] as path,
                    pg_tview_meta.entity as depends_on
                FROM {meta} pg_tview_meta
                WHERE pg_tview_meta.entity = $1

                UNION ALL

                SELECT
                    m.entity,
                    dt.depth + 1,
                    dt.path || m.entity,
                    dt.entity as depends_on
                FROM dep_tree dt
                JOIN {meta} m ON ('fk_' || dt.entity) = ANY(m.fk_columns)
                WHERE NOT (m.entity = ANY(dt.path))
                  AND dt.depth < 10
            )
            SELECT depth, entity AS entity_name, depends_on
            FROM dep_tree
            ORDER BY depth, entity_name",
                meta = crate::utils::meta_table()
            ),
            None,
            &args,
        ) {
            Ok(rows) => {
                let mut paths = Vec::new();
                for row in rows {
                    let depth = row["depth"].value::<i32>()?.unwrap_or(0);
                    let entity_name = row["entity_name"].value::<String>()?.unwrap_or_default();
                    let depends_on = row["depends_on"].value::<String>()?.unwrap_or_default();
                    paths.push((depth, entity_name, depends_on));
                }
                Ok::<_, spi::Error>(paths)
            }
            Err(e) => {
                warning!("Failed to query cascade path: {}", e);
                Ok(Vec::new())
            }
        }
    })
    .unwrap_or_default();

    TableIterator::new(results)
}