uqa-engine 0.5.0

Engine: schema-aware table store, catalog restore, transactions
//
// Unified Query Algebra
//
// Copyright (c) 2023-2026 Cognica, Inc.
//

use crate::Engine;
use std::sync::{atomic::Ordering, Arc};
use uqa_core::Value;
use uqa_sql::SQLError;

fn execute_at_message_clock(engine: &Engine, clock: i64, sql: &str) -> uqa_sql::SQLResult {
    engine
        .session
        .statement_started_at_micros
        .store(clock, Ordering::Relaxed);
    uqa_execution::statement::batch::execute(&engine.batch_execution_context(), sql, &[]).unwrap()
}

#[test]
fn cached_temporal_inputs_read_each_messages_transaction_clock() {
    let engine = Engine::new();
    let _caller_clock = uqa_sql::expr::TransactionClockScope::enter(5);
    let sql = "SELECT 'now'::timestamp AS input_clock, now()::timestamp AS transaction_clock, now() = 'now' AS implicit_input";
    for clock in [90_123_456_789, 190_123_456_789] {
        let result = execute_at_message_clock(&engine, clock, sql);
        let expected = Value::Temporal(uqa_core::TemporalValue::Timestamp { micros: clock });
        assert_eq!(result.rows[0]["input_clock"], expected);
        assert_eq!(result.rows[0]["transaction_clock"], expected);
        assert_eq!(result.rows[0]["implicit_input"], Value::Bool(true));
        assert!(engine.cached_sql_statement(sql).is_some());
        assert!(engine.cached_optimized_sql_plan(sql).is_none());
        assert_eq!(uqa_sql::expr::transaction_clock_micros(), Some(5));
    }
}

#[test]
fn cached_cursor_inputs_read_each_messages_transaction_clock() {
    let engine = Engine::new();
    let sql = "SELECT 'now'::timestamp AS value";
    for clock in [90_123_456_789, 190_123_456_789] {
        engine
            .session
            .statement_started_at_micros
            .store(clock, Ordering::Relaxed);
        let cursor =
            uqa_execution::statement::cursor::execute(&engine.batch_execution_context(), sql, &[])
                .unwrap();
        let rows = cursor
            .flat_map(|batch| batch.unwrap().into_rows())
            .collect::<Vec<_>>();
        assert_eq!(
            rows[0]["value"],
            Value::Temporal(uqa_core::TemporalValue::Timestamp { micros: clock })
        );
        assert!(engine.cached_sql_statement(sql).is_some());
        assert!(engine.cached_optimized_sql_plan(sql).is_none());
    }
}

#[test]
fn immutable_inputs_reuse_optimized_plans_while_parameter_types_are_reanalyzed() {
    let engine = Engine::new();
    let sql = "SELECT '42'::integer AS value";
    engine.sql(sql, &[]).unwrap();
    let first = engine.cached_optimized_sql_plan(sql).unwrap();
    engine.sql(sql, &[]).unwrap();
    assert!(Arc::ptr_eq(
        &first,
        &engine.cached_optimized_sql_plan(sql).unwrap()
    ));
    for value in [
        uqa_sql::SQLParam::Scalar(Value::Int(7)),
        uqa_sql::SQLParam::Scalar(Value::Str("text".into())),
    ] {
        engine.sql("SELECT $1 AS value", &[value]).unwrap();
        assert!(engine
            .cached_optimized_sql_plan("SELECT $1 AS value")
            .is_none());
    }
}

#[test]
fn simple_query_transaction_boundaries_refresh_temporal_input_clocks() {
    fn verify(engine: &Engine) {
        let _caller_clock = uqa_sql::expr::TransactionClockScope::enter(11);
        for (boundary, new_transaction) in [
            ("COMMIT", true),
            ("ROLLBACK", true),
            ("COMMIT AND CHAIN", true),
            ("ROLLBACK AND CHAIN", true),
            ("COMMIT; BEGIN", true),
            (
                "SAVEPOINT clock_savepoint; ROLLBACK TO clock_savepoint; RELEASE clock_savepoint",
                false,
            ),
            ("BEGIN", false),
        ] {
            engine.sql("BEGIN", &[]).unwrap();
            engine.session.transactions.lock()[0].started_at_micros = 90_123_456_789;
            let select = "SELECT 'now'::timestamp = now()::timestamp AS same_clock, now() = statement_timestamp() AS message_clock";
            let mut clocks = Vec::new();
            engine
                .sql_simple_query(&format!("{select}; {boundary}; {select}"), &[], |result| {
                    clocks.extend(
                        result
                            .rows
                            .iter()
                            .map(|row| (row["same_clock"].clone(), row["message_clock"].clone())),
                    );
                    Ok(())
                })
                .unwrap();
            assert_eq!(
                clocks,
                [
                    (Value::Bool(true), Value::Bool(false)),
                    (Value::Bool(true), Value::Bool(new_transaction)),
                ],
                "{boundary}"
            );
            assert_eq!(uqa_sql::expr::transaction_clock_micros(), Some(11));
            if engine.transaction_depth() != 0 {
                engine.sql("ROLLBACK", &[]).unwrap();
            }
        }
    }

    verify(&Engine::new());
    let directory = tempfile::tempdir().unwrap();
    verify(&Engine::open(&directory.path().join("clock-boundaries.db")).unwrap());
}

#[test]
fn consumer_failure_restores_statement_clock_and_depth_after_rollback() {
    let engine = Engine::new();
    engine.sql("CREATE TABLE pending (id BIGINT)", &[]).unwrap();
    let _caller_clock = uqa_sql::expr::TransactionClockScope::enter(17);
    engine
        .session
        .statement_started_at_micros
        .store(42, Ordering::Relaxed);
    let error = engine
        .sql_simple_query(
            "INSERT INTO pending VALUES (1); SELECT id FROM pending",
            &[],
            |_| {
                assert_eq!(
                    uqa_sql::expr::transaction_clock_micros(),
                    Some(engine.transaction_timestamp_micros())
                );
                Err(SQLError::Internal("consumer stopped".into()))
            },
        )
        .unwrap_err();
    assert!(matches!(error, SQLError::Internal(message) if message == "consumer stopped"));
    assert_eq!(
        engine.runtime.sql_execution_depth.load(Ordering::Relaxed),
        0
    );
    assert_eq!(
        engine
            .session
            .statement_started_at_micros
            .load(Ordering::Relaxed),
        42
    );
    assert_eq!(engine.transaction_depth(), 0);
    assert_eq!(uqa_sql::expr::transaction_clock_micros(), Some(17));
    assert!(engine
        .sql("SELECT id FROM pending", &[])
        .unwrap()
        .rows
        .is_empty());
}

#[test]
fn nested_callback_query_retains_the_outer_statement_scope() {
    let engine = Arc::new(Engine::new());
    let weak = Arc::downgrade(&engine);
    engine
        .register_scalar_function("nested_scope", move |_args: &[Value]| {
            let engine = weak.upgrade().unwrap();
            let clock = engine
                .session
                .statement_started_at_micros
                .load(Ordering::Relaxed);
            assert_eq!(
                engine.runtime.sql_execution_depth.load(Ordering::Relaxed),
                1
            );
            let transaction_clock = engine.transaction_timestamp_micros();
            assert_eq!(
                uqa_sql::expr::transaction_clock_micros(),
                Some(transaction_clock)
            );
            for _ in 0..2 {
                assert_eq!(
                    engine
                        .sql(
                            "SELECT 'now'::timestamp = now()::timestamp AS same_clock",
                            &[],
                        )?
                        .rows[0]["same_clock"],
                    Value::Bool(true)
                );
                assert_eq!(
                    uqa_sql::expr::transaction_clock_micros(),
                    Some(transaction_clock)
                );
            }
            assert_eq!(
                engine.runtime.sql_execution_depth.load(Ordering::Relaxed),
                1
            );
            assert_eq!(
                engine
                    .session
                    .statement_started_at_micros
                    .load(Ordering::Relaxed),
                clock
            );
            Ok(Value::Int(clock))
        })
        .unwrap();
    engine.sql("BEGIN", &[]).unwrap();
    engine.session.transactions.lock()[0].started_at_micros = 90_123_456_789;
    let _caller_clock = uqa_sql::expr::TransactionClockScope::enter(23);
    engine
        .session
        .statement_started_at_micros
        .store(73, Ordering::Relaxed);
    assert_eq!(
        engine
            .sql("COMMIT; SELECT nested_scope()", &[])
            .unwrap()
            .rows
            .len(),
        1
    );
    assert_eq!(uqa_sql::expr::transaction_clock_micros(), Some(23));
    assert_eq!(
        engine.runtime.sql_execution_depth.load(Ordering::Relaxed),
        0
    );
    assert_eq!(
        engine
            .session
            .statement_started_at_micros
            .load(Ordering::Relaxed),
        73
    );
}

#[test]
fn query_for_portals_retain_a_read_only_statement_scope_and_clean_up_after_failure() {
    fn verify(engine: &Engine) {
        engine.sql("CREATE TABLE loop_input(v integer); INSERT INTO loop_input VALUES (1), (2); CREATE FUNCTION loop_reader() RETURNS integer LANGUAGE plpgsql AS $$ DECLARE rec record; total integer := 0; BEGIN FOR rec IN SELECT v FROM loop_input ORDER BY v LOOP total := total + rec.v; END LOOP; RETURN total; END $$; CREATE FUNCTION loop_failure() RETURNS integer LANGUAGE plpgsql AS $$ DECLARE rec record; BEGIN FOR rec IN SELECT v FROM loop_input LOOP RAISE EXCEPTION 'loop failed'; END LOOP; RETURN 0; END $$", &[]).unwrap();
        for _ in 0..2 {
            assert_eq!(
                engine
                    .sql("SELECT loop_reader() AS total", &[])
                    .unwrap()
                    .rows[0]["total"],
                Value::Int(3)
            );
            assert_eq!(engine.transaction_depth(), 0);
            let cursor = engine
                .sql_cursor("SELECT loop_reader() AS total", &[])
                .unwrap();
            assert_eq!(cursor.row_count(), 1);
            drop(cursor);
            assert_eq!(engine.transaction_depth(), 0);
        }
        engine.sql("BEGIN READ ONLY", &[]).unwrap();
        assert_eq!(
            engine
                .sql("SELECT loop_reader() AS total", &[])
                .unwrap()
                .rows[0]["total"],
            Value::Int(3)
        );
        assert_eq!(engine.transaction_depth(), 1);
        assert!(engine.current_transaction_is_read_only());
        engine.sql("COMMIT", &[]).unwrap();
        assert_eq!(
            engine
                .sql("SELECT loop_failure()", &[])
                .unwrap_err()
                .sqlstate(),
            Some("P0001")
        );
        assert_eq!(engine.transaction_depth(), 0);
        assert!(engine.session.portals.lock().is_empty());
        assert_eq!(
            engine
                .sql("SELECT loop_reader() AS total", &[])
                .unwrap()
                .rows[0]["total"],
            Value::Int(3)
        );
    }
    verify(&Engine::new());
    for provider in 0..3 {
        let (_directory, engine, _) = crate::tests::relation_lock_support::sessions(provider);
        verify(&engine);
    }
}