pg_tviews 0.1.0-beta.24

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
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
use pgrx::prelude::*;
use std::collections::HashMap;
use std::sync::{LazyLock, Mutex, PoisonError};

use crate::cascade_path::CascadePath;

/// Cached information for a table managed by `pg_tviews`.
///
/// Beyond the entity name, this carries everything the issue #56 direct-patch
/// eligibility check needs, so the row trigger can decide the fast
/// path from a single cached lookup with **no SPI in the hot path** (populated once
/// per session on cache miss, invalidated on DDL).
#[derive(Clone, Debug)]
pub struct CachedEntityInfo {
    pub name: String,
    /// The TVIEW is keyed on a DISTINCT ON key: a row's own values may not be its
    /// group's, so the fast path declines.
    pub distinct_on: bool,
    /// Registered before the root table had a cascade path of its own (ADR
    /// 0169): how the trigger finds its key until it is re-registered.
    pub legacy_root: Option<LegacyRoot>,

    /// Direct-patch column→key map (issue #56): base column name → JSONB key it
    /// feeds in the entity's own `data`. Empty ⇒ the fast path never engages.
    pub direct_map: HashMap<String, String>,

    /// Integer FK columns of this entity's base table. A changed FK is a
    /// membership change ⇒ the fast path must decline (issue #56 eligibility).
    pub fk_columns: Vec<String>,

    /// UUID FK columns of this entity's base table (same membership rule).
    pub uuid_fk_columns: Vec<String>,

    /// Output columns the `tv_<entity>` table materialises (`pk_<entity>`, `id`,
    /// `data`, and any column projected outside `data`). A changed base column that
    /// is also a projected output column would go stale under a data-only patch ⇒
    /// the fast path declines (issue #56 eligibility).
    pub output_columns: Vec<String>,

    /// `true` when the backing view is a UNION / UNION ALL — different refresh
    /// machinery, so the fast path declines (issue #56 eligibility).
    pub is_union: bool,
}

/// How the row trigger keys a TVIEW registered before its root table had a
/// cascade path of its own (ADR 0169).
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum LegacyRoot {
    /// `pk_<entity>`, read by name off `tb_<entity>`.
    Pk,
    /// A DISTINCT ON TVIEW: refreshed in full.
    DistinctOn,
}

/// Global cache for `EntityDepGraph` to avoid repeated `pg_tview_meta` queries
static ENTITY_GRAPH_CACHE: LazyLock<Mutex<Option<super::graph::EntityDepGraph>>> =
    LazyLock::new(|| Mutex::new(None));

/// Global cache for table OID → entity info
/// Stores Option<CachedEntityInfo> to cache negative lookups (None results)
static TABLE_ENTITY_CACHE: LazyLock<Mutex<HashMap<pg_sys::Oid, Option<CachedEntityInfo>>>> =
    LazyLock::new(|| Mutex::new(HashMap::new()));

// Transaction-scoped cache for table OID → cascade paths.
// Cleared on transaction end to avoid stale data.
thread_local! {
    static CASCADE_PATH_CACHE: std::cell::RefCell<HashMap<pg_sys::Oid, Vec<CascadePath>>> =
        std::cell::RefCell::new(HashMap::new());
}

/// Cache operations for `EntityDepGraph`
pub mod graph_cache {
    #[allow(clippy::wildcard_imports)] // Reason: module-internal prelude import
    use super::*;

    /// Get cached `EntityDepGraph`, loading from database if not cached
    pub fn load_cached() -> crate::TViewResult<crate::queue::graph::EntityDepGraph> {
        // Check if caching is enabled
        if !crate::config::graph_cache_enabled() {
            return crate::queue::graph::EntityDepGraph::load();
        }

        let mut cache = ENTITY_GRAPH_CACHE
            .lock()
            .unwrap_or_else(PoisonError::into_inner);

        if let Some(graph) = cache.as_ref() {
            // Cache hit
            crate::metrics::metrics_api::record_graph_cache_hit();
            return Ok(graph.clone());
        }

        // Cache miss: load from database
        crate::metrics::metrics_api::record_graph_cache_miss();
        let graph = crate::queue::graph::EntityDepGraph::load()?;
        *cache = Some(graph.clone());
        drop(cache);
        Ok(graph)
    }

    /// Invalidate the `EntityDepGraph` cache
    /// Should be called when TVIEWs are created or dropped
    pub fn invalidate() {
        let mut cache = ENTITY_GRAPH_CACHE
            .lock()
            .unwrap_or_else(PoisonError::into_inner);
        *cache = None;
    }
}

/// Cache operations for table OID → entity mapping
pub mod table_cache {
    #[allow(clippy::wildcard_imports)] // Reason: module-internal prelude import
    use super::*;

    /// Get cached entity info for table OID
    /// Loads from database on first miss per session, caches negative results
    pub fn entity_info_cached(
        table_oid: pg_sys::Oid,
    ) -> crate::TViewResult<Option<CachedEntityInfo>> {
        // Check if caching is enabled
        if !crate::config::table_cache_enabled() {
            return load_entity_info_uncached(table_oid);
        }

        // Fast path: check cache (distinguishes cached None from not-in-cache)
        {
            let cache = TABLE_ENTITY_CACHE
                .lock()
                .unwrap_or_else(PoisonError::into_inner);
            if let Some(cached_value) = cache.get(&table_oid) {
                crate::metrics::metrics_api::record_table_cache_hit();
                return Ok(cached_value.clone());
            }
        }

        // Slow path: query and cache (including None)
        crate::metrics::metrics_api::record_table_cache_miss();
        let info = load_entity_info_uncached(table_oid)?;

        {
            let mut cache = TABLE_ENTITY_CACHE
                .lock()
                .unwrap_or_else(PoisonError::into_inner);
            crate::utils::bound_cache(&mut cache);
            cache.insert(table_oid, info.clone());
        }

        Ok(info)
    }

    /// Get cached entity name (backward compatibility)
    pub fn entity_for_table_cached(table_oid: pg_sys::Oid) -> crate::TViewResult<Option<String>> {
        entity_info_cached(table_oid).map(|info| info.map(|i| i.name))
    }

    /// Load entity info from the database on a cache miss.
    ///
    /// Loads the full `TviewMeta` for the entity plus the `tv_<entity>` output
    /// columns, so a single cached record answers every issue #56 eligibility
    /// question without further SPI in the trigger hot path.
    fn load_entity_info_uncached(
        table_oid: pg_sys::Oid,
    ) -> crate::TViewResult<Option<CachedEntityInfo>> {
        let Some(name) = crate::catalog::entity_for_table_uncached(table_oid)? else {
            return Ok(None);
        };

        // Name resolved but no meta row (shouldn't happen) → minimal safe info
        // with the fast path disabled (empty direct_map).
        let Some(meta) = crate::catalog::TviewMeta::load_by_entity(&name)? else {
            return Ok(Some(CachedEntityInfo {
                name,
                distinct_on: false,
                legacy_root: Some(LegacyRoot::Pk),
                direct_map: HashMap::new(),
                fk_columns: Vec::new(),
                uuid_fk_columns: Vec::new(),
                output_columns: Vec::new(),
                is_union: false,
            }));
        };

        let mut direct_map: HashMap<String, String> = meta
            .direct_map_columns
            .iter()
            .cloned()
            .zip(meta.direct_map_keys.iter().cloned())
            .collect();

        // Output columns the tv_<entity> table materialises. If they can't be
        // determined we can't verify the "projected column" eligibility rule, so
        // disable the fast path for this entity rather than risk a stale column.
        let output_columns = if let Ok(cols) = crate::utils::get_view_columns_by_oid(meta.tview_oid)
        {
            cols
        } else {
            direct_map.clear();
            Vec::new()
        };

        // A base column that also feeds a projected column outside `data` must
        // recompute: a data-only patch would leave that column stale (#98).
        // SAFETY: DatumWithOid wraps the entity name as a TEXT parameter.
        let args = [unsafe {
            pgrx::datum::DatumWithOid::new(
                name.as_str(),
                PgOid::BuiltIn(PgBuiltInOids::TEXTOID).value(),
            )
        }];
        let definition: Option<String> = Spi::get_one_with_args(
            &format!(
                "SELECT definition FROM {} WHERE entity = $1",
                crate::utils::meta_table()
            ),
            &args,
        )
        .unwrap_or(None);
        match definition
            .as_deref()
            .and_then(crate::schema::direct_map::columns_referenced_outside_data)
        {
            Some(projected) => direct_map.retain(|col, _| !projected.contains(&col.to_lowercase())),
            None => direct_map.clear(),
        }

        Ok(Some(CachedEntityInfo {
            name,
            distinct_on: meta.identity.kind == crate::lineage::IdentityKind::DistinctOn,
            legacy_root: if !meta.identity.legacy {
                None
            } else if meta.identity.legacy_distinct_on {
                Some(LegacyRoot::DistinctOn)
            } else {
                Some(LegacyRoot::Pk)
            },
            direct_map,
            fk_columns: meta.fk_columns,
            uuid_fk_columns: meta.uuid_fk_columns,
            output_columns,
            is_union: meta.is_union,
        }))
    }

    /// Invalidate the table entity cache
    /// Should be called when TVIEWs are created or dropped
    pub fn invalidate() {
        let mut cache = TABLE_ENTITY_CACHE
            .lock()
            .unwrap_or_else(PoisonError::into_inner);
        cache.clear();
    }
}

/// Cache operations for cascade paths (transaction-scoped)
pub mod cascade_cache {
    use super::{CASCADE_PATH_CACHE, CascadePath, pg_sys};

    /// Get cached cascade paths for a source table OID.
    /// Returns all `CascadePath` entries across all entities where `source_oid` matches.
    pub fn cascade_paths_for_table(table_oid: pg_sys::Oid) -> crate::TViewResult<Vec<CascadePath>> {
        CASCADE_PATH_CACHE.with(|cache| {
            let mut cache = cache.borrow_mut();

            if let Some(paths) = cache.get(&table_oid) {
                return Ok(paths.clone());
            }

            let paths = load_cascade_paths_for_table(table_oid)?;
            cache.insert(table_oid, paths.clone());
            Ok(paths)
        })
    }

    /// Load cascade paths matching a source table OID from `pg_tview_meta`
    fn load_cascade_paths_for_table(
        table_oid: pg_sys::Oid,
    ) -> crate::TViewResult<Vec<CascadePath>> {
        let meta_list = crate::catalog::TviewMeta::load_all()?;

        let mut relevant_paths = Vec::new();
        for meta in meta_list {
            for path in meta.cascade_paths {
                if path.source_oid == table_oid {
                    relevant_paths.push(path);
                }
            }
        }

        Ok(relevant_paths)
    }

    /// Clear the cascade path cache (called on transaction end)
    pub fn clear_cache() {
        CASCADE_PATH_CACHE.with(|cache| {
            cache.borrow_mut().clear();
        });
    }
}

// ── Cross-backend invalidation ───────────────────────────────────────────
//
// Every cache here is per backend. A relcache invalidation (DDL on a watched
// relation, or the `pg_tview_meta` statement trigger that invalidates the catalog
// itself) bumps the generation; the next `sync_generation()` then clears them all.
// The callback only touches a `Cell`: it can run in the middle of a catalog access,
// so it must neither take a lock nor call SPI.

thread_local! {
    static GENERATION: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
    static SEEN_GENERATION: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
    static WATCHED: std::cell::RefCell<std::collections::HashSet<pg_sys::Oid>> =
        std::cell::RefCell::new(std::collections::HashSet::new());
}

thread_local! {
    static CATALOG_WATCHED: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
}

/// Watch `pg_tview_meta` itself: its statement trigger invalidates it on every
/// write, which is how other backends learn about a created, changed or dropped
/// TVIEW. Resolved once per backend.
fn watch_catalog() {
    if CATALOG_WATCHED.with(std::cell::Cell::get) {
        return;
    }
    crate::metrics::metrics_api::record_catalog_lookup();
    if let Ok(Some(oid)) = Spi::connect(|client| {
        client
            .select(
                &format!("SELECT to_regclass('{}')::oid", crate::utils::meta_table()),
                None,
                &[],
            )?
            .first()
            .get_one::<pg_sys::Oid>()
    }) {
        watch(&[oid]);
        CATALOG_WATCHED.with(|w| w.set(true));
    }
}

/// Watch relations whose DDL must invalidate the caches (TVIEW tables and views,
/// `pg_tview_meta`).
pub fn watch(oids: &[pg_sys::Oid]) {
    WATCHED.with(|w| w.borrow_mut().extend(oids.iter().copied()));
}

/// Clear every cache if a relevant invalidation arrived since the last call.
pub fn sync_generation() {
    watch_catalog();
    let current = GENERATION.with(std::cell::Cell::get);
    if SEEN_GENERATION.with(std::cell::Cell::get) != current {
        SEEN_GENERATION.with(|s| s.set(current));
        invalidate_all_caches();
    }
}

#[pg_guard]
unsafe extern "C-unwind" fn relcache_callback(_arg: pg_sys::Datum, relid: pg_sys::Oid) {
    let relevant = relid == pg_sys::InvalidOid
        || WATCHED.with(|w| w.try_borrow().map_or(true, |w| w.contains(&relid)));
    if relevant {
        GENERATION.with(|g| g.set(g.get().wrapping_add(1)));
    }
}

/// Register the relcache callback. Called once from `_PG_init`.
pub fn register_relcache_callback() {
    // SAFETY: registers a static callback with no argument; valid in `_PG_init`.
    unsafe {
        pg_sys::CacheRegisterRelcacheCallback(Some(relcache_callback), pg_sys::Datum::from(0));
    }
}

/// Invalidate the `pg_tviews` caches of every backend (and this one) once the
/// current transaction commits, by invalidating `relid`'s relcache entry.
#[pg_extern]
fn pg_tviews_invalidate_caches(relid: pg_sys::Oid) {
    // SAFETY: the caller passes an existing relation (the trigger's TG_RELID).
    unsafe { pg_sys::CacheInvalidateRelcacheByRelid(relid) };
}

/// Combined cache invalidation for all caches
pub fn invalidate_all_caches() {
    graph_cache::invalidate();
    table_cache::invalidate();
    cascade_cache::clear_cache();
    crate::catalog::clear_meta_cache();
    super::ops::clear_crash_recovery_cache();
    crate::lifecycle::invalidate_jsonb_delta_cache();
    crate::utils::invalidate_oid_relname_cache();
    crate::utils::invalidate_view_columns_cache();
    crate::delta::clear_caches();
}

#[cfg(test)]
#[allow(clippy::wildcard_imports)] // Reason: test module prelude import
mod tests {
    use super::*;

    /// Serialises the tests that mutate the process-global `TABLE_ENTITY_CACHE`;
    /// the test harness runs tests in parallel threads.
    static TABLE_CACHE_TEST_LOCK: Mutex<()> = Mutex::new(());

    #[test]
    fn test_graph_cache_invalidation() {
        // Test that invalidate clears the cache
        graph_cache::invalidate();

        assert!(ENTITY_GRAPH_CACHE.lock().unwrap().is_none());
    }

    #[test]
    fn test_table_cache_invalidation() {
        let _guard = TABLE_CACHE_TEST_LOCK
            .lock()
            .unwrap_or_else(PoisonError::into_inner);

        // Add something to cache
        {
            let mut cache = TABLE_ENTITY_CACHE.lock().unwrap();
            cache.insert(
                pg_sys::Oid::from(123),
                Some(CachedEntityInfo {
                    name: "test".to_string(),
                    distinct_on: false,
                    legacy_root: None,
                    direct_map: HashMap::new(),
                    fk_columns: Vec::new(),
                    uuid_fk_columns: Vec::new(),
                    output_columns: Vec::new(),
                    is_union: false,
                }),
            );
        }

        // Verify it's there
        assert!(
            TABLE_ENTITY_CACHE
                .lock()
                .unwrap()
                .get(&pg_sys::Oid::from(123))
                .is_some()
        );

        // Invalidate
        table_cache::invalidate();

        // Verify it's gone
        assert!(TABLE_ENTITY_CACHE.lock().unwrap().is_empty());
    }

    #[test]
    fn test_cached_entity_info_carries_direct_patch_fields() {
        // The cached record exposes everything the issue #56 eligibility check
        // needs, so the trigger hot path never re-queries pg_tview_meta.
        let mut direct_map = HashMap::new();
        direct_map.insert("bio".to_string(), "bio".to_string());
        direct_map.insert("name".to_string(), "display_name".to_string());

        let info = CachedEntityInfo {
            name: "user".to_string(),
            distinct_on: false,
            legacy_root: None,
            direct_map,
            fk_columns: vec!["fk_org".to_string()],
            uuid_fk_columns: vec![],
            output_columns: vec!["pk_user".to_string(), "id".to_string(), "data".to_string()],
            is_union: false,
        };

        assert_eq!(info.direct_map.get("bio").map(String::as_str), Some("bio"));
        assert_eq!(
            info.direct_map.get("name").map(String::as_str),
            Some("display_name")
        );
        assert!(info.fk_columns.contains(&"fk_org".to_string()));
        assert!(info.output_columns.contains(&"data".to_string()));
        assert!(!info.is_union);
    }

    #[test]
    fn test_negative_cache_entry() {
        let _guard = TABLE_CACHE_TEST_LOCK
            .lock()
            .unwrap_or_else(PoisonError::into_inner);

        table_cache::invalidate();
        // Insert a None entry
        TABLE_ENTITY_CACHE
            .lock()
            .unwrap()
            .insert(pg_sys::Oid::from(999), None);
        // Verify it's cached as None (not a cache miss)
        let cache = TABLE_ENTITY_CACHE.lock().unwrap();
        assert!(cache.get(&pg_sys::Oid::from(999)).is_some()); // key exists
        assert!(cache.get(&pg_sys::Oid::from(999)).unwrap().is_none()); // value is None
    }
}