uqa-engine 0.4.0

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

//! Persistent sessions and observed logical waits for relation-lock schedules.

use super::*;
use std::{
    sync::mpsc,
    thread,
    time::{Duration, Instant},
};
use uqa_execution::row_locks::{shared_objects::SharedCatalogLock, RowLockKey};

pub(super) fn sessions(provider: usize) -> (tempfile::TempDir, Engine, Engine) {
    let directory = tempfile::tempdir().unwrap();
    let path = directory.path().join("table-locks.db");
    let first = match provider {
        0 => Engine::open(&path).unwrap(),
        1 => Engine::from_persistent_provider(Arc::new(
            uqa_storage_sqlite::SQLiteKeyValueStorage::open(&path).unwrap(),
        ))
        .unwrap(),
        2 => Engine::from_persistent_provider(Arc::new(
            uqa_storage_redb::RedbStorage::open(&path).unwrap(),
        ))
        .unwrap(),
        _ => unreachable!(),
    };
    let second = first.new_session().unwrap();
    for engine in [&first, &second] {
        engine.release_automatic_statistics_client();
        engine
            .session
            .statistics_worker
            .store(true, Ordering::Release);
    }
    sql(
        &first,
        "CREATE TABLE t(v integer); INSERT INTO t VALUES (1)",
    );
    (directory, first, second)
}

pub(super) fn sql(engine: &Engine, statement: &str) -> SQLResult {
    engine
        .sql(statement, &[])
        .unwrap_or_else(|error| panic!("{statement}: {error}"))
}

pub(super) fn error(engine: &Engine, statement: &str, state: &str) {
    let error = engine.sql(statement, &[]).unwrap_err();
    assert_eq!(error.sqlstate(), Some(state), "{statement}: {error}");
}

pub(super) fn wait_for_relation(
    first: &Engine,
    session: u64,
    relation: &str,
    finished: impl Fn() -> bool,
) -> bool {
    let key = first.row_locks.table_key(relation);
    wait_for_relation_key(first, session, key, finished)
}

fn wait_for_relation_key(
    first: &Engine,
    session: u64,
    key: u64,
    finished: impl Fn() -> bool,
) -> bool {
    let deadline = Instant::now() + Duration::from_secs(30);
    while !first.row_locks.waiting_for_relation(session, key)
        && !finished()
        && Instant::now() < deadline
    {
        thread::sleep(Duration::from_millis(1));
    }
    first.row_locks.waiting_for_relation(session, key)
}

pub(super) fn after_wait(
    holder: &Engine,
    worker: Engine,
    statement: &str,
    relation: &str,
    release: &str,
) -> (Engine, Result<SQLResult, SQLError>) {
    let statement = statement.to_string();
    after_operation_wait(holder, worker, relation, release, move |worker| {
        worker.sql(&statement, &[])
    })
}

pub(super) fn after_operation_wait<T: Send + 'static>(
    holder: &Engine,
    worker: Engine,
    relation: &str,
    release: &str,
    operation: impl FnOnce(&Engine) -> Result<T, SQLError> + Send + 'static,
) -> (Engine, Result<T, SQLError>) {
    let key = holder.row_locks.table_key(relation);
    after_operation_key_wait(holder, worker, key, relation, release, operation)
}

pub(super) fn after_index_wait(
    holder: &Engine,
    worker: Engine,
    statement: &str,
    index: [u8; 16],
    release: &str,
) -> (Engine, Result<SQLResult, SQLError>) {
    let key = holder.row_locks.index_key(index);
    let statement = statement.to_owned();
    after_operation_key_wait(
        holder,
        worker,
        key,
        "index incarnation",
        release,
        move |worker| worker.sql(&statement, &[]),
    )
}

fn after_operation_key_wait<T: Send + 'static>(
    holder: &Engine,
    worker: Engine,
    key: u64,
    relation: &str,
    release: &str,
    operation: impl FnOnce(&Engine) -> Result<T, SQLError> + Send + 'static,
) -> (Engine, Result<T, SQLError>) {
    let session = worker.session_id;
    let cancel = worker.runtime.cancellation.clone();
    let (send, done) = mpsc::channel();
    let task = thread::spawn(move || {
        let result = operation(&worker);
        let _ = send.send(result);
        worker
    });
    let waited = wait_for_relation_key(holder, session, key, || task.is_finished());
    let released = holder.sql(release, &[]);
    if released.is_err() {
        cancel.cancel();
    }
    let result = done.recv_timeout(Duration::from_secs(30));
    if result.is_err() {
        cancel.cancel();
    }
    let worker = task.join().unwrap();
    released.unwrap();
    assert!(
        waited,
        "expected a logical wait on {relation}; operation error: {:?}",
        result.as_ref().map(|result| result.as_ref().err())
    );
    (worker, result.unwrap())
}

pub(super) fn after_tuple_wait(
    holder: &Engine,
    worker: Engine,
    statement: &str,
    catalog: &str,
    doc_id: u64,
    release: &str,
) -> (Engine, Result<SQLResult, SQLError>) {
    after_tuple_wait_with_release(holder, worker, statement, catalog, doc_id, || {
        holder.sql(release, &[])
    })
}

pub(super) fn after_tuple_wait_with_release(
    holder: &Engine,
    worker: Engine,
    statement: &str,
    catalog: &str,
    doc_id: u64,
    release: impl FnOnce() -> Result<SQLResult, SQLError>,
) -> (Engine, Result<SQLResult, SQLError>) {
    let session = worker.session_id;
    let key = RowLockKey {
        table: holder.row_locks.table_key(catalog),
        doc_id,
    };
    let cancel = worker.runtime.cancellation.clone();
    let statement = statement.to_string();
    let (send, done) = mpsc::channel();
    let task = thread::spawn(move || {
        let result = worker.sql(&statement, &[]);
        let _ = send.send(result);
        worker
    });
    let deadline = Instant::now() + Duration::from_secs(30);
    while !holder.row_locks.waiting_for_row(session, key)
        && !task.is_finished()
        && Instant::now() < deadline
    {
        thread::sleep(Duration::from_millis(1));
    }
    let waited = holder.row_locks.waiting_for_row(session, key);
    let released = release();
    if released.is_err() {
        cancel.cancel();
    }
    let result = done.recv_timeout(Duration::from_secs(30));
    if result.is_err() {
        cancel.cancel();
    }
    let worker = task.join().unwrap();
    assert!(
        waited,
        "expected catalog tuple wait: {catalog}/{doc_id}; received {result:?}"
    );
    released.unwrap();
    (worker, result.unwrap())
}

pub(super) fn after_shared_wait(
    holder: &Engine,
    worker: Engine,
    statement: &str,
    target: SharedCatalogLock<'_>,
    release: &str,
) -> (Engine, Result<SQLResult, SQLError>) {
    after_shared_wait_with_release(holder, worker, statement, target, || {
        holder.sql(release, &[])
    })
}

pub(super) fn after_shared_wait_with_release(
    holder: &Engine,
    worker: Engine,
    statement: &str,
    target: SharedCatalogLock<'_>,
    release: impl FnOnce() -> Result<SQLResult, SQLError>,
) -> (Engine, Result<SQLResult, SQLError>) {
    let key = holder.row_locks.shared_catalog_key(target);
    let session = worker.session_id;
    let cancel = worker.runtime.cancellation.clone();
    let statement = statement.to_string();
    let (send, done) = mpsc::channel();
    let task = thread::spawn(move || {
        eprintln!("shared-catalog worker executing: {statement}");
        let result = worker.sql(&statement, &[]);
        eprintln!("shared-catalog worker completed: {statement}");
        let _ = send.send(result);
        worker
    });
    let deadline = Instant::now() + Duration::from_secs(30);
    while !holder.row_locks.waiting_for_relation(session, key)
        && !task.is_finished()
        && Instant::now() < deadline
    {
        thread::sleep(Duration::from_millis(1));
    }
    let waited = holder.row_locks.waiting_for_relation(session, key);
    eprintln!("shared-catalog wait observed={waited}; releasing {target:?}");
    let released = release();
    eprintln!("shared-catalog holder release completed for {target:?}");
    if released.is_err() {
        cancel.cancel();
    }
    let result = done.recv_timeout(Duration::from_secs(30));
    if result.is_err() {
        cancel.cancel();
    }
    let worker = task.join().unwrap();
    released.unwrap();
    assert!(
        waited,
        "expected shared catalog wait on {target:?}, received {result:?}"
    );
    (worker, result.unwrap())
}

pub(super) fn reopen(provider: usize, path: &std::path::Path) -> Engine {
    match provider {
        0 => Engine::open(path).unwrap(),
        1 => Engine::from_persistent_provider(Arc::new(
            uqa_storage_sqlite::SQLiteKeyValueStorage::open(path).unwrap(),
        ))
        .unwrap(),
        2 => Engine::from_persistent_provider(Arc::new(
            uqa_storage_redb::RedbStorage::open(path).unwrap(),
        ))
        .unwrap(),
        _ => unreachable!(),
    }
}

pub(super) fn before_commit(holder: &Engine, worker: Engine, statement: &str) -> Engine {
    let label = statement.to_owned();
    let statement = statement.to_owned();
    let cancel = worker.runtime.cancellation.clone();
    let (send, receive) = std::sync::mpsc::channel();
    let task = std::thread::spawn(move || {
        let result = worker.sql(&statement, &[]);
        let _ = send.send(result);
        worker
    });
    let result = receive.recv_timeout(std::time::Duration::from_secs(30));
    if result.is_err() {
        cancel.cancel();
    }
    let released = holder.sql("COMMIT", &[]);
    let worker = task.join().unwrap();
    result
        .unwrap()
        .unwrap_or_else(|error| panic!("independent worker {label}: {error}"));
    released.unwrap_or_else(|error| panic!("holder commit after {label}: {error}"));
    worker
}