commit 3f96391734e948f576b99ef1a5d6c92cdb9844dc parent e2b796128acc6432e1449961051f01e244ce2b6f Author: triesap <tyson@radroots.org> Date: Wed, 15 Jul 2026 02:34:37 +0000 sql_core: replace native SQLite executor with SQLx - remove the rusqlite native executor and row utility surface from sql_core - add a SQLx-backed native executor with explicit dynamic-SQL and bind handling - update replica store and sync tests to use the SQLx executor directly - verify fmt, tests, clippy, and rusqlite inverse tree checks Diffstat:
20 files changed, 275 insertions(+), 233 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock @@ -4805,9 +4805,10 @@ name = "radroots_sql_core" version = "0.1.0-alpha.2" dependencies = [ "chrono", - "rusqlite", + "futures-executor", "serde", "serde_json", + "sqlx", "uuid", ] diff --git a/Cargo.toml b/Cargo.toml @@ -133,6 +133,7 @@ config = { version = "0.14" } directories = { version = "6" } ed25519-dalek = { version = "2.1.1", default-features = false } futures = { version = "0.3" } +futures-executor = { version = "0.3" } flate2 = { version = "1" } getrandom = { version = "0.2", default-features = false } hkdf = { version = "0.12", default-features = false } diff --git a/crates/replica_store/README b/crates/replica_store/README @@ -9,7 +9,7 @@ facades and shaped query helpers for the `radroots` core libraries. * migration, backup, restore, and export helpers for replica store state; * typed CRUD helpers over the shared replica schema models; * shaped query helpers for common trade, farm, and event freshness lookups; - * feature-gated web and native backend support through `radroots_sql_core`. + * feature-gated web and SQLx native backend support through `radroots_sql_core`. ## Copyright diff --git a/crates/replica_store/tests/error_paths.rs b/crates/replica_store/tests/error_paths.rs @@ -58,7 +58,7 @@ use radroots_replica_schema::trade_product::{ use radroots_replica_schema::trade_product_location::ITradeProductLocationRelation; use radroots_replica_schema::trade_product_media::ITradeProductMediaRelation; use radroots_replica_store::ReplicaSql; -use radroots_sql_core::{SqlError, SqlExecutor, SqliteExecutor}; +use radroots_sql_core::{SqlError, SqlExecutor, SqlxSqliteExecutor}; use serde::de::DeserializeOwned; use serde_json::json; @@ -70,14 +70,14 @@ fn hex64(ch: char) -> String { std::iter::repeat_n(ch, 64).collect() } -fn open_db() -> ReplicaSql<SqliteExecutor> { - let exec = SqliteExecutor::open_memory().expect("open sqlite memory"); +fn open_db() -> ReplicaSql<SqlxSqliteExecutor> { + let exec = SqlxSqliteExecutor::open_memory().expect("open sqlite memory"); let db = ReplicaSql::new(exec); db.migrate_up().expect("migrate up"); db } -fn drop_table(db: &ReplicaSql<SqliteExecutor>, table_name: &str) { +fn drop_table(db: &ReplicaSql<SqlxSqliteExecutor>, table_name: &str) { let sql = format!("DROP TABLE {table_name};"); db.executor().exec(&sql, "[]").expect("drop table"); } diff --git a/crates/replica_store/tests/full_mode.rs b/crates/replica_store/tests/full_mode.rs @@ -59,7 +59,7 @@ use radroots_replica_schema::trade_product::{ use radroots_replica_schema::trade_product_location::ITradeProductLocationRelation; use radroots_replica_schema::trade_product_media::ITradeProductMediaRelation; use radroots_replica_store::{ReplicaSql, export_manifest}; -use radroots_sql_core::{SqlError, SqliteExecutor}; +use radroots_sql_core::{SqlError, SqlxSqliteExecutor}; use serde::de::DeserializeOwned; use serde_json::json; @@ -87,8 +87,8 @@ fn assert_not_found<T>(result: Result<T, ReplicaSchemaError<SqlError>>) { assert!(matches!(err.error, SqlError::NotFound(_))); } -fn open_db() -> ReplicaSql<SqliteExecutor> { - let exec = SqliteExecutor::open_memory().expect("open sqlite memory"); +fn open_db() -> ReplicaSql<SqlxSqliteExecutor> { + let exec = SqlxSqliteExecutor::open_memory().expect("open sqlite memory"); let db = ReplicaSql::new(exec); db.migrate_up().expect("migrate up"); db diff --git a/crates/replica_store/tests/migration_repairs.rs b/crates/replica_store/tests/migration_repairs.rs @@ -1,13 +1,13 @@ use radroots_replica_store::migrations; -use radroots_sql_core::{SqlExecutor, SqliteExecutor}; +use radroots_sql_core::{SqlExecutor, SqlxSqliteExecutor}; use serde_json::Value; -fn query_rows(exec: &SqliteExecutor, sql: &str) -> Vec<Value> { +fn query_rows(exec: &SqlxSqliteExecutor, sql: &str) -> Vec<Value> { serde_json::from_str(&exec.query_raw(sql, "[]").expect("query should succeed")) .expect("query should decode") } -fn create_legacy_schema_without_secondary_indexes(exec: &SqliteExecutor) { +fn create_legacy_schema_without_secondary_indexes(exec: &SqlxSqliteExecutor) { let schema = [ "CREATE TABLE __migrations (id INTEGER PRIMARY KEY, name TEXT NOT NULL UNIQUE, applied_at TEXT NOT NULL DEFAULT (datetime('now')))", "CREATE TABLE farm (id CHAR(36) PRIMARY KEY NOT NULL UNIQUE CHECK(length(id) = 36), created_at DATETIME NOT NULL CHECK(length(created_at) = 24), updated_at DATETIME NOT NULL CHECK(length(updated_at) = 24), d_tag TEXT NOT NULL, pubkey TEXT NOT NULL, name TEXT NOT NULL, about TEXT, website TEXT, picture TEXT, banner TEXT, location_primary TEXT, location_city TEXT, location_region TEXT, location_country TEXT)", @@ -51,7 +51,7 @@ fn create_legacy_schema_without_secondary_indexes(exec: &SqliteExecutor) { #[test] fn run_all_up_repairs_missing_indexes_in_legacy_sqlite_dbs() { - let exec = SqliteExecutor::open_memory().expect("open sqlite memory"); + let exec = SqlxSqliteExecutor::open_memory().expect("open sqlite memory"); create_legacy_schema_without_secondary_indexes(&exec); let before = query_rows( diff --git a/crates/replica_sync/src/emit.rs b/crates/replica_sync/src/emit.rs @@ -265,7 +265,7 @@ pub fn radroots_replica_membership_claim_events( } let mut events = Vec::new(); - for (member_pubkey, _) in by_member.iter() { + for member_pubkey in by_member.keys() { let all_claims = load_member_claims_for_member(exec, member_pubkey)?; let mut farm_pubkeys = all_claims .into_iter() @@ -937,7 +937,7 @@ mod tests { farm, farm_gcs_location, farm_member, farm_member_claim, farm_tag, gcs_location, migrations, nostr_profile, plot, plot_gcs_location, plot_tag, }; - use radroots_sql_core::{ExecOutcome, SqlError, SqlExecutor, SqliteExecutor}; + use radroots_sql_core::{ExecOutcome, SqlError, SqlExecutor, SqlxSqliteExecutor}; struct ErrorExecutor; @@ -964,7 +964,7 @@ mod tests { } struct QueryFailExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, needle: &'static str, err: SqlError, } @@ -995,7 +995,7 @@ mod tests { } struct DuplicateFarmSelectorExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, duplicated_rows_json: String, } @@ -1022,7 +1022,7 @@ mod tests { } } - fn seed(exec: &SqliteExecutor) -> (Farm, Plot, Plot) { + fn seed(exec: &SqlxSqliteExecutor) -> (Farm, Plot, Plot) { migrations::run_all_up(exec).expect("migrations"); let farm = farm::create( exec, @@ -1286,7 +1286,12 @@ mod tests { (farm, plot_primary, plot_secondary) } - fn create_farm_record(exec: &SqliteExecutor, d_tag: &str, pubkey: &str, name: &str) -> Farm { + fn create_farm_record( + exec: &SqlxSqliteExecutor, + d_tag: &str, + pubkey: &str, + name: &str, + ) -> Farm { farm::create( exec, &IFarmFields { @@ -1307,7 +1312,7 @@ mod tests { .result } - fn create_plot_record(exec: &SqliteExecutor, farm_id: &str, d_tag: &str, name: &str) { + fn create_plot_record(exec: &SqlxSqliteExecutor, farm_id: &str, d_tag: &str, name: &str) { let _ = plot::create( exec, &IPlotFields { @@ -1324,7 +1329,12 @@ mod tests { .expect("plot"); } - fn add_member_record(exec: &SqliteExecutor, farm_id: &str, member_pubkey: &str, role: &str) { + fn add_member_record( + exec: &SqlxSqliteExecutor, + farm_id: &str, + member_pubkey: &str, + role: &str, + ) { let _ = farm_member::create( exec, &IFarmMemberFields { @@ -1338,7 +1348,7 @@ mod tests { #[test] fn emit_paths_cover_private_and_public_helpers() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_row, plot_primary, plot_secondary) = seed(&exec); let by_id = resolve_farm( @@ -1657,7 +1667,7 @@ mod tests { #[test] fn emit_option_toggles_and_empty_rows_cover_branches() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm = farm::create( @@ -1734,7 +1744,7 @@ mod tests { #[test] fn emit_profile_variants_and_missing_profiles_are_handled() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_row, _, _) = seed(&exec); let _ = farm_member::create( @@ -1812,7 +1822,7 @@ mod tests { #[test] fn emit_query_error_paths_are_reported() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm = Farm { @@ -1906,7 +1916,7 @@ mod tests { #[test] fn load_farm_location_omits_string_only_public_locations() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_row = farm::create( &exec, @@ -1933,7 +1943,7 @@ mod tests { #[test] fn emit_propagates_queryfail_and_builder_errors() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_row, _, _) = seed(&exec); let selector = RadrootsReplicaFarmSelector { id: Some(farm_row.id.clone()), @@ -2074,7 +2084,7 @@ mod tests { #[test] fn emit_additional_error_branches_are_reported() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_row, _, _) = seed(&exec); crate::canonical::failpoints::set_error(); @@ -2301,7 +2311,7 @@ mod tests { create_plot_record(&exec, &list_plot_error_farm.id, "", "plot-list-error"); assert!(radroots_replica_list_set_events(&exec, &list_plot_error_farm).is_err()); - let clean_exec = SqliteExecutor::open_memory().expect("db clean"); + let clean_exec = SqlxSqliteExecutor::open_memory().expect("db clean"); let (clean_farm, _, _) = seed(&clean_exec); super::failpoints::set_list_set_to_wire_error(); assert!(radroots_replica_list_set_events(&clean_exec, &clean_farm).is_err()); @@ -2342,7 +2352,7 @@ mod tests { #[test] fn emit_list_set_wire_error_paths_are_reported() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_row, _, _) = seed(&exec); super::failpoints::set_list_set_to_wire_error(); @@ -2361,7 +2371,7 @@ mod tests { #[test] fn emit_pass_through_executor_instantiation_paths_are_covered() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_row, _, plot_secondary) = seed(&exec); let pass = QueryFailExecutor { @@ -2459,7 +2469,7 @@ mod tests { #[test] fn emit_executor_trait_method_paths_are_covered() { - let sqlite = SqliteExecutor::open_memory().expect("db"); + let sqlite = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&sqlite).expect("migrations"); let err_exec = ErrorExecutor; diff --git a/crates/replica_sync/src/ingest.rs b/crates/replica_sync/src/ingest.rs @@ -1556,7 +1556,7 @@ mod tests { gcs_location, migrations, nostr_event_head, plot, plot_gcs_location, plot_tag, trade_product, }; - use radroots_sql_core::{ExecOutcome, SqlExecutor, SqliteExecutor}; + use radroots_sql_core::{ExecOutcome, SqlExecutor, SqlxSqliteExecutor}; fn test_event_envelope( id: u64, @@ -1612,7 +1612,7 @@ mod tests { } struct TxnExecutor<'a> { - inner: Option<&'a SqliteExecutor>, + inner: Option<&'a SqlxSqliteExecutor>, begin_err: Option<SqlError>, commit_err: Option<SqlError>, rollback_count: Arc<AtomicUsize>, @@ -1654,7 +1654,7 @@ mod tests { } struct DeleteErrorExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, table_name: &'static str, err: SqlError, } @@ -1686,7 +1686,7 @@ mod tests { } struct PassExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, } impl SqlExecutor for PassExecutor<'_> { @@ -1712,7 +1712,7 @@ mod tests { } struct QueryFailExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, needle: &'static str, err: SqlError, } @@ -2058,7 +2058,7 @@ mod tests { assert!(err.to_string().contains("last_event_id invalid")); } - fn seed_rows(exec: &SqliteExecutor) -> (String, String, String, String) { + fn seed_rows(exec: &SqlxSqliteExecutor) -> (String, String, String, String) { migrations::run_all_up(exec).expect("migrations"); let farm_row = farm::create( exec, @@ -2249,7 +2249,7 @@ mod tests { #[test] fn ingest_core_paths_cover_helpers_and_decisions() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let factory = RadrootsReplicaDefaultIdFactory; @@ -2582,7 +2582,7 @@ mod tests { #[test] fn ingest_listing_projects_trade_product_and_removes_archived_replacements() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let seller_pubkey = "c".repeat(64); @@ -2747,7 +2747,7 @@ mod tests { #[test] fn ingest_listing_preserves_fractional_exact_economics() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let seller_pubkey = "c".repeat(64); @@ -2811,7 +2811,7 @@ mod tests { #[test] fn upsert_location_none_paths_are_ok() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_row = farm::create( @@ -2853,7 +2853,7 @@ mod tests { #[test] fn ingest_delete_error_paths_are_covered() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_id, _farm_pubkey, farm_d_tag, _plot_d_tag) = seed_rows(&exec); let not_found_farm_tags = DeleteErrorExecutor { @@ -3053,7 +3053,7 @@ mod tests { #[test] fn ingest_pass_executor_and_parse_edge_paths_are_covered() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let pass = PassExecutor { inner: &exec }; @@ -3194,7 +3194,7 @@ mod tests { #[test] fn create_gcs_location_success_path_is_covered() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let id = create_gcs_location(&exec, sample_gcs(1.0, 2.0, "s0"), &FixedFactory) @@ -3204,7 +3204,7 @@ mod tests { #[test] fn ingest_default_factory_wrapper_paths_are_covered() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_pubkey = "f".repeat(64); @@ -3289,7 +3289,7 @@ mod tests { #[test] fn ingest_txn_executor_instantiation_error_paths_are_covered() { - let pass_db = SqliteExecutor::open_memory().expect("db"); + let pass_db = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&pass_db).expect("migrations"); let pass_txn = TxnExecutor { inner: Some(&pass_db), @@ -3455,7 +3455,7 @@ mod tests { #[test] fn ingest_sqlite_queryfail_and_parser_edges_are_covered() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let pass_through = QueryFailExecutor { inner: &exec, @@ -4029,7 +4029,7 @@ mod tests { #[test] fn ingest_insert_and_state_error_branches_are_covered() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_id, farm_pubkey, farm_d_tag, plot_d_tag) = seed_rows(&exec); let profile = profile_event( @@ -4218,7 +4218,7 @@ mod tests { #[test] fn upsert_member_helpers_ignore_empty_entry_values() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (farm_id, farm_pubkey, _, _) = seed_rows(&exec); let member_pubkey = "6".repeat(64); @@ -4309,7 +4309,7 @@ mod tests { #[test] fn ingest_error_paths_cover_missing_farm_and_bad_list_set_tags() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let (_, farm_pubkey, farm_d_tag, plot_d_tag) = seed_rows(&exec); let missing_farm_plot = plot_event( diff --git a/crates/replica_sync/src/sync_state.rs b/crates/replica_sync/src/sync_state.rs @@ -126,11 +126,11 @@ mod tests { use radroots_replica_schema::farm::IFarmFields; use radroots_replica_schema::nostr_event_head::INostrEventHeadFields; use radroots_replica_store::{farm, migrations, nostr_event_head}; - use radroots_sql_core::{SqlExecutor, SqliteExecutor}; + use radroots_sql_core::{SqlExecutor, SqlxSqliteExecutor}; #[test] fn sync_status_empty_db_is_zero() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let status = radroots_replica_sync_status(&exec).expect("status"); assert_eq!(status.expected_count, 0); @@ -139,7 +139,7 @@ mod tests { #[test] fn sync_status_tracks_expected_and_pending() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_row = farm::create( @@ -191,7 +191,7 @@ mod tests { #[test] fn pending_publish_batch_lists_only_missing_or_changed_expected_events() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_row = farm::create( @@ -247,14 +247,14 @@ mod tests { #[test] fn sync_status_reports_farm_query_errors() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); let err = radroots_replica_sync_status(&exec).expect_err("farm query error"); assert!(err.to_string().contains("invalid query")); } #[test] fn sync_status_reports_emit_errors() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let _ = farm::create( &exec, @@ -282,7 +282,7 @@ mod tests { #[test] fn sync_status_reports_content_hash_errors() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let _ = farm::create( &exec, @@ -308,7 +308,7 @@ mod tests { #[test] fn sync_status_reports_state_query_errors() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let _ = exec .exec("DROP TABLE nostr_event_head;", "[]") diff --git a/crates/replica_sync/src/tests.rs b/crates/replica_sync/src/tests.rs @@ -19,7 +19,7 @@ use radroots_replica_store::{ farm, farm_gcs_location, farm_member, farm_member_claim, farm_tag, gcs_location, migrations, nostr_profile, plot, plot_gcs_location, plot_tag, }; -use radroots_sql_core::SqliteExecutor; +use radroots_sql_core::SqlxSqliteExecutor; use radroots_sql_core::error::SqlError; use std::panic; @@ -32,7 +32,7 @@ fn unwrap_sql<T>(result: Result<T, ReplicaSchemaError<SqlError>>, label: &str) - #[test] fn sync_all_emits_expected_order() { - let exec = SqliteExecutor::open_memory().expect("exec"); + let exec = SqlxSqliteExecutor::open_memory().expect("exec"); migrations::run_all_up(&exec).expect("migrations"); let farm_pubkey = "f".repeat(64); diff --git a/crates/replica_sync/tests/ingest_roundtrip.rs b/crates/replica_sync/tests/ingest_roundtrip.rs @@ -41,7 +41,7 @@ use radroots_replica_sync::{ RadrootsReplicaSyncRequest, radroots_replica_ingest_event, radroots_replica_sync_all, radroots_replica_sync_status, }; -use radroots_sql_core::SqliteExecutor; +use radroots_sql_core::SqlxSqliteExecutor; use radroots_sql_core::error::SqlError; use radroots_sql_core::{ExecOutcome, SqlExecutor}; use std::panic; @@ -54,7 +54,7 @@ fn unwrap_sql<T>(result: Result<T, ReplicaSchemaError<SqlError>>, label: &str) - } struct BeginFailExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, } impl SqlExecutor for BeginFailExecutor<'_> { @@ -80,7 +80,7 @@ impl SqlExecutor for BeginFailExecutor<'_> { } struct CommitFailExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, } impl SqlExecutor for CommitFailExecutor<'_> { @@ -106,7 +106,7 @@ impl SqlExecutor for CommitFailExecutor<'_> { } struct DeleteFailExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, table_name: &'static str, err: SqlError, } @@ -137,7 +137,7 @@ impl SqlExecutor for DeleteFailExecutor<'_> { } struct QueryFailExecutor<'a> { - inner: &'a SqliteExecutor, + inner: &'a SqlxSqliteExecutor, needle: &'static str, err: SqlError, } @@ -191,7 +191,7 @@ fn draft_to_event(draft: &RadrootsReplicaEventDraft, index: u32) -> RadrootsEven } fn seed_source( - exec: &SqliteExecutor, + exec: &SqlxSqliteExecutor, ) -> ( RadrootsReplicaSyncRequest, String, @@ -460,11 +460,11 @@ fn seed_source( #[test] fn ingest_roundtrip_yields_zero_pending_sync() { - let source = SqliteExecutor::open_memory().expect("source db"); + let source = SqlxSqliteExecutor::open_memory().expect("source db"); let (_source_request, farm_d_tag, farm_pubkey, drafts) = seed_source(&source); assert_eq!(drafts.len(), 10); - let target = SqliteExecutor::open_memory().expect("target db"); + let target = SqlxSqliteExecutor::open_memory().expect("target db"); migrations::run_all_up(&target).expect("target migrations"); let mut skipped = 0usize; @@ -501,7 +501,7 @@ fn ingest_roundtrip_yields_zero_pending_sync() { #[test] fn sync_status_empty_db_is_zero() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let status = radroots_replica_sync_status(&exec).expect("status"); assert_eq!(status.expected_count, 0); @@ -510,7 +510,7 @@ fn sync_status_empty_db_is_zero() { #[test] fn sync_all_selector_and_options_paths_are_supported() { - let source = SqliteExecutor::open_memory().expect("source db"); + let source = SqlxSqliteExecutor::open_memory().expect("source db"); let (request, farm_d_tag, farm_pubkey, full_events) = seed_source(&source); let by_pair = radroots_replica_sync_all( @@ -544,7 +544,7 @@ fn sync_all_selector_and_options_paths_are_supported() { #[test] fn ingest_rejects_unsupported_kind() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let event = event_with_parts( 1, @@ -560,7 +560,7 @@ fn ingest_rejects_unsupported_kind() { #[test] fn ingest_reports_transaction_boundary_errors() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let author = "a".repeat(64); let profile = profile_event( @@ -580,7 +580,7 @@ fn ingest_reports_transaction_boundary_errors() { #[test] fn ingest_reports_delete_internal_errors() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_pubkey = "f".repeat(64); let farm_d_tag = "AAAAAAAAAAAAAAAAAAAAAA"; @@ -618,7 +618,7 @@ fn ingest_reports_delete_internal_errors() { #[test] fn ingest_reports_parse_and_state_error_paths_for_all_kinds() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let profile_pubkey = "a".repeat(64); @@ -764,7 +764,7 @@ fn ingest_reports_parse_and_state_error_paths_for_all_kinds() { #[test] fn ingest_reports_query_fail_paths_for_profile_farm_plot_and_list_sets() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let assert_query_fail = |needle: &'static str, event: &RadrootsEventEnvelope| { @@ -1081,7 +1081,7 @@ fn list_set_event( #[test] fn ingest_event_paths_cover_profile_farm_plot_and_list_set_variants() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let profile_pubkey = "9".repeat(64); @@ -1835,9 +1835,9 @@ fn ingest_event_paths_cover_profile_farm_plot_and_list_set_variants() { #[test] fn sync_status_reports_pending_when_not_all_events_are_ingested() { - let source = SqliteExecutor::open_memory().expect("source"); + let source = SqlxSqliteExecutor::open_memory().expect("source"); let (_request, _farm_d_tag, _farm_pubkey, drafts) = seed_source(&source); - let target = SqliteExecutor::open_memory().expect("target"); + let target = SqlxSqliteExecutor::open_memory().expect("target"); migrations::run_all_up(&target).expect("migrations"); for (index, draft) in drafts.iter().enumerate() { @@ -1858,7 +1858,7 @@ fn sync_status_reports_pending_when_not_all_events_are_ingested() { #[test] fn sync_all_rejects_invalid_selectors_and_resolves_unique_pair() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let missing_selector_err = radroots_replica_sync_all( @@ -1928,7 +1928,7 @@ fn sync_all_rejects_invalid_selectors_and_resolves_unique_pair() { #[test] fn sync_emit_handles_invalid_geojson_and_unknown_profile_type() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_pubkey = "0".repeat(64); @@ -2123,7 +2123,7 @@ fn sync_emit_handles_invalid_geojson_and_unknown_profile_type() { #[test] fn sync_emit_reports_encode_error_for_invalid_farm_record() { - let exec = SqliteExecutor::open_memory().expect("db"); + let exec = SqlxSqliteExecutor::open_memory().expect("db"); migrations::run_all_up(&exec).expect("migrations"); let farm_row = unwrap_sql( diff --git a/crates/sql_core/Cargo.toml b/crates/sql_core/Cargo.toml @@ -19,12 +19,13 @@ crate-type = ["rlib"] default = ["std"] std = ["dep:chrono", "dep:uuid", "serde/std", "serde_json/std"] web = ["std", "uuid/js"] -native = ["std", "dep:rusqlite"] +native = ["std", "dep:futures-executor", "dep:sqlx", "sqlx/sqlite-bundled"] embedded = [] [dependencies] serde_json = { workspace = true } -rusqlite = { workspace = true, features = ["bundled"], optional = true } +futures-executor = { workspace = true, optional = true } +sqlx = { workspace = true, optional = true, features = ["derive"] } chrono = { workspace = true, optional = true } serde = { workspace = true } uuid = { workspace = true, optional = true } diff --git a/crates/sql_core/README b/crates/sql_core/README @@ -7,7 +7,7 @@ migration primitives for the `radroots` core libraries. * the `SqlExecutor` trait and `ExecOutcome` type used by higher-level database crates; - * native SQLite execution and utility support behind the `native` feature; + * native SQLx SQLite execution and utility support behind the `native` feature; * WebAssembly execution, export locking, and bridge integration for browser-facing builds; * optional embedded-engine support for self-contained runtimes and migration diff --git a/crates/sql_core/src/error.rs b/crates/sql_core/src/error.rs @@ -51,9 +51,9 @@ impl From<serde_json::Error> for SqlError { } } -#[cfg(all(feature = "native", feature = "std"))] -impl From<rusqlite::Error> for SqlError { - fn from(e: rusqlite::Error) -> Self { +#[cfg(feature = "native")] +impl From<sqlx::Error> for SqlError { + fn from(e: sqlx::Error) -> Self { SqlError::InvalidQuery(e.to_string()) } } diff --git a/crates/sql_core/src/executor_sqlite.rs b/crates/sql_core/src/executor_sqlite.rs @@ -1,80 +0,0 @@ -use crate::sqlite_util; -use crate::{ExecOutcome, SqlExecutor, error::SqlError}; -use rusqlite::{Connection, params_from_iter}; -use serde_json::Value; -use std::path::Path; -use std::sync::{Arc, Mutex}; - -pub struct SqliteExecutor { - conn: Arc<Mutex<Connection>>, -} - -impl SqliteExecutor { - pub fn open<P: AsRef<Path>>(path: P) -> Result<Self, SqlError> { - let conn = Connection::open(path).map_err(SqlError::from)?; - Ok(Self { - conn: Arc::new(Mutex::new(conn)), - }) - } - - pub fn open_memory() -> Result<Self, SqlError> { - let conn = Connection::open_in_memory().map_err(SqlError::from)?; - Ok(Self { - conn: Arc::new(Mutex::new(conn)), - }) - } -} - -impl SqlExecutor for SqliteExecutor { - fn exec(&self, sql: &str, params_json: &str) -> Result<ExecOutcome, SqlError> { - let binds = sqlite_util::parse_params(params_json)?; - let conn = self.conn.lock().map_err(|_| SqlError::Internal)?; - if binds.is_empty() { - let total_changes_before = conn.total_changes(); - conn.execute_batch(sql).map_err(SqlError::from)?; - let total_changes_after = conn.total_changes(); - let last_insert_id = conn.last_insert_rowid(); - return Ok(ExecOutcome { - changes: (total_changes_after - total_changes_before) as i64, - last_insert_id, - }); - } - let n = conn - .execute(sql, params_from_iter(binds)) - .map_err(SqlError::from)?; - let last_insert_id = conn.last_insert_rowid(); - Ok(ExecOutcome { - changes: n as i64, - last_insert_id, - }) - } - - fn query_raw(&self, sql: &str, params_json: &str) -> Result<String, SqlError> { - let binds = sqlite_util::parse_params(params_json)?; - let rows = { - let conn = self.conn.lock().map_err(|_| SqlError::Internal)?; - let mut stmt = conn.prepare(sql).map_err(SqlError::from)?; - let mapped = stmt.query_map(params_from_iter(binds), sqlite_util::row_to_json)?; - mapped.collect::<Result<Vec<_>, _>>()? - }; - Ok(Value::from(rows).to_string()) - } - - fn begin(&self) -> Result<(), SqlError> { - let conn = self.conn.lock().map_err(|_| SqlError::Internal)?; - conn.execute("BEGIN", []).map_err(SqlError::from)?; - Ok(()) - } - - fn commit(&self) -> Result<(), SqlError> { - let conn = self.conn.lock().map_err(|_| SqlError::Internal)?; - conn.execute("COMMIT", []).map_err(SqlError::from)?; - Ok(()) - } - - fn rollback(&self) -> Result<(), SqlError> { - let conn = self.conn.lock().map_err(|_| SqlError::Internal)?; - conn.execute("ROLLBACK", []).map_err(SqlError::from)?; - Ok(()) - } -} diff --git a/crates/sql_core/src/executor_sqlx_sqlite.rs b/crates/sql_core/src/executor_sqlx_sqlite.rs @@ -0,0 +1,86 @@ +use std::path::Path; +use std::sync::{Arc, Mutex}; + +use sqlx::Connection; +use sqlx::sqlite::{SqliteConnectOptions, SqliteConnection}; + +use crate::sqlx_sqlite_util; +use crate::{ExecOutcome, SqlExecutor, error::SqlError}; + +pub struct SqlxSqliteExecutor { + conn: Arc<Mutex<SqliteConnection>>, +} + +impl SqlxSqliteExecutor { + pub fn open<P: AsRef<Path>>(path: P) -> Result<Self, SqlError> { + let options = SqliteConnectOptions::new() + .filename(path) + .create_if_missing(true); + Self::connect(options) + } + + pub fn open_memory() -> Result<Self, SqlError> { + Self::connect(SqliteConnectOptions::new().in_memory(true)) + } + + fn connect(options: SqliteConnectOptions) -> Result<Self, SqlError> { + let conn = futures_executor::block_on(SqliteConnection::connect_with(&options))?; + Ok(Self { + conn: Arc::new(Mutex::new(conn)), + }) + } +} + +impl SqlExecutor for SqlxSqliteExecutor { + fn exec(&self, sql: &str, params_json: &str) -> Result<ExecOutcome, SqlError> { + let binds = sqlx_sqlite_util::parse_params(params_json)?; + let mut conn = self.conn.lock().map_err(|_| SqlError::Internal)?; + if binds.is_empty() { + let result = futures_executor::block_on( + sqlx::raw_sql(sqlx::AssertSqlSafe(sql)).execute(&mut *conn), + )?; + return Ok(ExecOutcome { + changes: i64::try_from(result.rows_affected()).map_err(|_| SqlError::Internal)?, + last_insert_id: result.last_insert_rowid(), + }); + } + let query = sqlx_sqlite_util::bind_params(sqlx::query(sqlx::AssertSqlSafe(sql)), binds)?; + let result = futures_executor::block_on(query.execute(&mut *conn))?; + Ok(ExecOutcome { + changes: i64::try_from(result.rows_affected()).map_err(|_| SqlError::Internal)?, + last_insert_id: result.last_insert_rowid(), + }) + } + + fn query_raw(&self, sql: &str, params_json: &str) -> Result<String, SqlError> { + let binds = sqlx_sqlite_util::parse_params(params_json)?; + let query = sqlx_sqlite_util::bind_params(sqlx::query(sqlx::AssertSqlSafe(sql)), binds)?; + let rows = { + let mut conn = self.conn.lock().map_err(|_| SqlError::Internal)?; + futures_executor::block_on(query.fetch_all(&mut *conn))? + }; + let rows = rows + .iter() + .map(sqlx_sqlite_util::row_to_json) + .collect::<Result<Vec<_>, _>>()?; + Ok(serde_json::Value::from(rows).to_string()) + } + + fn begin(&self) -> Result<(), SqlError> { + let mut conn = self.conn.lock().map_err(|_| SqlError::Internal)?; + futures_executor::block_on(sqlx::query("BEGIN").execute(&mut *conn))?; + Ok(()) + } + + fn commit(&self) -> Result<(), SqlError> { + let mut conn = self.conn.lock().map_err(|_| SqlError::Internal)?; + futures_executor::block_on(sqlx::query("COMMIT").execute(&mut *conn))?; + Ok(()) + } + + fn rollback(&self) -> Result<(), SqlError> { + let mut conn = self.conn.lock().map_err(|_| SqlError::Internal)?; + futures_executor::block_on(sqlx::query("ROLLBACK").execute(&mut *conn))?; + Ok(()) + } +} diff --git a/crates/sql_core/src/lib.rs b/crates/sql_core/src/lib.rs @@ -6,11 +6,11 @@ pub mod error; pub mod migrations; #[cfg(all(feature = "native", feature = "std"))] -mod executor_sqlite; +mod executor_sqlx_sqlite; #[cfg(all(feature = "native", feature = "std"))] -pub use executor_sqlite::SqliteExecutor; +pub use executor_sqlx_sqlite::SqlxSqliteExecutor; #[cfg(all(feature = "native", feature = "std"))] -pub mod sqlite_util; +pub mod sqlx_sqlite_util; #[cfg(feature = "embedded")] mod executor_embedded; diff --git a/crates/sql_core/src/sqlite_util.rs b/crates/sql_core/src/sqlite_util.rs @@ -1,59 +0,0 @@ -#![forbid(unsafe_code)] - -use crate::error::SqlError; -use rusqlite::{Row, types::Value as SqlValue}; -use serde_json::{Map, Value}; - -pub fn parse_params(params_json: &str) -> Result<Vec<SqlValue>, SqlError> { - let vals: Vec<Value> = serde_json::from_str(params_json) - .map_err(|e| SqlError::SerializationError(e.to_string()))?; - vals.into_iter() - .map(|v| match v { - Value::Null => Ok(SqlValue::Null), - Value::Bool(b) => Ok(SqlValue::from(if b { 1 } else { 0 })), - Value::Number(n) => { - if let Some(i) = n.as_i64() { - Ok(SqlValue::from(i)) - } else if let Some(u) = n.as_u64() { - Ok(SqlValue::from(u as i64)) - } else if let Some(f) = n.as_f64() { - Ok(SqlValue::from(f)) - } else { - Err(SqlError::InvalidArgument("unsupported number".to_string())) - } - } - Value::String(s) => Ok(SqlValue::from(s)), - other => Err(SqlError::InvalidArgument(format!( - "unsupported bind value: {}", - other - ))), - }) - .collect() -} - -pub fn row_to_json(row: &Row) -> rusqlite::Result<Value> { - let stmt = row.as_ref(); - let mut obj = Map::new(); - for i in 0..stmt.column_count() { - let name = stmt.column_name(i).unwrap_or("").to_string(); - let v = row.get_ref(i)?; - let j = match v { - rusqlite::types::ValueRef::Null => Value::Null, - rusqlite::types::ValueRef::Integer(i) => Value::from(i), - rusqlite::types::ValueRef::Real(f) => Value::from(f), - rusqlite::types::ValueRef::Text(s) => { - let s = std::str::from_utf8(s).map_err(|e| { - rusqlite::Error::FromSqlConversionFailure( - i, - rusqlite::types::Type::Text, - Box::new(e), - ) - })?; - Value::from(s.to_string()) - } - rusqlite::types::ValueRef::Blob(_) => Value::Null, - }; - obj.insert(name, j); - } - Ok(Value::Object(obj)) -} diff --git a/crates/sql_core/src/sqlx_sqlite_util.rs b/crates/sql_core/src/sqlx_sqlite_util.rs @@ -0,0 +1,82 @@ +#![forbid(unsafe_code)] + +use crate::error::SqlError; +use serde_json::{Map, Value}; +use sqlx::sqlite::{SqliteArguments, SqliteRow}; +use sqlx::{Column, Row, TypeInfo, ValueRef}; + +#[derive(Debug)] +pub enum SqliteBindValue { + Null, + Integer(i64), + Real(f64), + Text(String), +} + +pub fn parse_params(params_json: &str) -> Result<Vec<SqliteBindValue>, SqlError> { + let vals: Vec<Value> = serde_json::from_str(params_json) + .map_err(|e| SqlError::SerializationError(e.to_string()))?; + vals.into_iter() + .map(|v| match v { + Value::Null => Ok(SqliteBindValue::Null), + Value::Bool(b) => Ok(SqliteBindValue::Integer(i64::from(b))), + Value::Number(n) => { + if let Some(i) = n.as_i64() { + Ok(SqliteBindValue::Integer(i)) + } else if let Some(u) = n.as_u64() { + let value = i64::try_from(u).map_err(|_| { + SqlError::InvalidArgument("integer bind exceeds i64".to_string()) + })?; + Ok(SqliteBindValue::Integer(value)) + } else if let Some(f) = n.as_f64() { + Ok(SqliteBindValue::Real(f)) + } else { + Err(SqlError::InvalidArgument("unsupported number".to_string())) + } + } + Value::String(s) => Ok(SqliteBindValue::Text(s)), + other => Err(SqlError::InvalidArgument(format!( + "unsupported bind value: {}", + other + ))), + }) + .collect() +} + +pub fn bind_params<'q>( + query: sqlx::query::Query<'q, sqlx::Sqlite, SqliteArguments>, + params: Vec<SqliteBindValue>, +) -> Result<sqlx::query::Query<'q, sqlx::Sqlite, SqliteArguments>, SqlError> { + let mut query = query; + for param in params { + query = match param { + SqliteBindValue::Null => query.bind(Option::<String>::None), + SqliteBindValue::Integer(value) => query.bind(value), + SqliteBindValue::Real(value) => query.bind(value), + SqliteBindValue::Text(value) => query.bind(value), + }; + } + Ok(query) +} + +pub fn row_to_json(row: &SqliteRow) -> Result<Value, SqlError> { + let mut obj = Map::new(); + for (index, column) in row.columns().iter().enumerate() { + let raw = row.try_get_raw(index)?; + let value = if raw.is_null() { + Value::Null + } else { + match raw.type_info().name() { + "INTEGER" | "BOOLEAN" => Value::from(row.try_get::<i64, _>(index)?), + "REAL" => Value::from(row.try_get::<f64, _>(index)?), + "TEXT" | "DATE" | "TIME" | "DATETIME" => { + Value::from(row.try_get::<String, _>(index)?) + } + "BLOB" => Value::Null, + _ => return Err(SqlError::InvalidQuery(raw.type_info().name().to_string())), + } + }; + obj.insert(column.name().to_string(), value); + } + Ok(Value::Object(obj)) +} diff --git a/crates/sql_core/tests/coverage.rs b/crates/sql_core/tests/coverage.rs @@ -1,5 +1,5 @@ #[cfg(feature = "native")] -use radroots_sql_core::SqliteExecutor; +use radroots_sql_core::SqlxSqliteExecutor; use radroots_sql_core::error::SqlError; use radroots_sql_core::migrations::{Migration, migrations_run_all_down, migrations_run_all_up}; use radroots_sql_core::utils::{ @@ -190,8 +190,8 @@ fn sql_executor_reference_impl_forwards_all_methods() { #[cfg(feature = "native")] #[test] -fn sqlite_executor_exec_runs_multi_statement_batches_without_params() { - let exec = SqliteExecutor::open_memory().expect("open sqlite memory"); +fn sqlx_sqlite_executor_exec_runs_multi_statement_batches_without_params() { + let exec = SqlxSqliteExecutor::open_memory().expect("open sqlite memory"); let outcome = exec .exec(