From 5b5da4adb4301472c9b4843adecab8cab7629fcc Mon Sep 17 00:00:00 2001 From: Ralph Kuepper Date: Mon, 24 Aug 2026 19:38:36 +0200 Subject: [PATCH] fix(mysql2): isolate prepared operations and pool transactions --- Cargo.lock | 1 + crates/perry-ext-mysql2/Cargo.toml | 3 + crates/perry-ext-mysql2/src/lib.rs | 509 +++++++++++++----- .../perry-ext-mysql2/src/test_async_shims.rs | 103 ++++ ...ue_8745_8746_mysql2_operation_isolation.ts | 69 +++ 5 files changed, 536 insertions(+), 149 deletions(-) create mode 100644 crates/perry-ext-mysql2/src/test_async_shims.rs create mode 100644 test-files/test_issue_8745_8746_mysql2_operation_isolation.ts diff --git a/Cargo.lock b/Cargo.lock index 4645cbc6ac..2a9e494b60 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6077,6 +6077,7 @@ version = "0.5.1519" dependencies = [ "chrono", "perry-ffi", + "perry-runtime", "sqlx", "tokio", ] diff --git a/crates/perry-ext-mysql2/Cargo.toml b/crates/perry-ext-mysql2/Cargo.toml index 5dc4a3be90..dcd6219727 100644 --- a/crates/perry-ext-mysql2/Cargo.toml +++ b/crates/perry-ext-mysql2/Cargo.toml @@ -24,3 +24,6 @@ chrono.workspace = true [dev-dependencies] perry-ffi = { workspace = true, features = ["runtime-link"] } +# Standalone extension tests need the runtime half of the test-only async FFI +# shims; production code still depends on perry-ffi only. +perry-runtime = { workspace = true, features = ["default", "stdlib"] } diff --git a/crates/perry-ext-mysql2/src/lib.rs b/crates/perry-ext-mysql2/src/lib.rs index 960ad26621..7ee3fe054b 100644 --- a/crates/perry-ext-mysql2/src/lib.rs +++ b/crates/perry-ext-mysql2/src/lib.rs @@ -15,15 +15,20 @@ //! adapter; followup once a wrapper actually demands it). use perry_ffi::{ - alloc_string, build_object_shape, get_handle_mut, js_array_alloc, js_array_get, js_array_push, + alloc_string, build_object_shape, js_array_alloc, js_array_get, js_array_push, js_object_alloc_with_shape, js_object_get_field, js_object_set_field, register_handle, - spawn_blocking, take_handle, ArrayHeader, Handle, JsPromise, JsValue, ObjectHeader, Promise, - StringHeader, + spawn_blocking, take_handle, with_handle, ArrayHeader, Handle, JsPromise, JsValue, + ObjectHeader, Promise, StringHeader, }; use sqlx::mysql::{MySqlConnection, MySqlPool, MySqlPoolOptions, MySqlRow}; use sqlx::pool::PoolConnection; use sqlx::{Column, Connection, MySql, Row, TypeInfo}; +use std::sync::Arc; use std::time::Duration; +use tokio::sync::Mutex; + +#[cfg(test)] +mod test_async_shims; const DEFAULT_CONNECT_TIMEOUT_SECS: u64 = 10; const DEFAULT_QUERY_TIMEOUT_SECS: u64 = 30; @@ -470,7 +475,7 @@ fn is_row_returning_query(sql: &str) -> bool { || upper.starts_with("WITH") } -#[derive(Clone, Debug)] +#[derive(Clone, Debug, PartialEq)] enum ParamValue { Null, String(String), @@ -479,6 +484,44 @@ enum ParamValue { Bool(bool), } +/// Everything needed to execute one mysql2 call, copied off the Perry heap +/// before the asynchronous work is scheduled. Keeping the SQL and its bind +/// values in one owned object makes it impossible for a later call to replace +/// either half while this request is waiting for a pool connection. +#[derive(Clone, Debug, PartialEq)] +struct QueryRequest { + sql: String, + params: Vec, + rows_as_array: bool, + /// `mysql2.query()` uses the text protocol when it has no values, whereas + /// `execute()` always represents a prepared statement. + force_prepared: bool, +} + +impl QueryRequest { + fn new( + sql: String, + params: Vec, + rows_as_array: bool, + force_prepared: bool, + ) -> Self { + Self { + sql, + params, + rows_as_array, + force_prepared, + } + } + + fn is_row_returning(&self) -> bool { + is_row_returning_query(&self.sql) + } + + fn uses_prepared_statement(&self) -> bool { + self.force_prepared || !self.params.is_empty() + } +} + unsafe fn extract_params_from_jsvalue(params: JsValue) -> Vec { let arr_ptr = params.as_pointer::(); if arr_ptr.is_null() { @@ -527,13 +570,135 @@ unsafe fn read_sql(sql_ptr: *const u8) -> String { // ── Connection ──────────────────────────────────────────────────── pub struct MysqlConnectionHandle { - pub connection: Option, + pub connection: Arc>>, } impl MysqlConnectionHandle { pub fn new(conn: MySqlConnection) -> Self { Self { - connection: Some(conn), + connection: Arc::new(Mutex::new(Some(conn))), + } + } +} + +#[derive(Clone)] +enum MysqlConnectionTarget { + Direct(Arc>>), + Pool(Arc>>>), +} + +/// Resolve either mysql2 connection handle family without returning a +/// registry-backed `'static` reference. The old `get_handle_mut` calls dropped +/// DashMap's guard before async work began, so overlapping workers could hold +/// aliased mutable references to the same connection wrapper. +fn connection_target(handle: Handle) -> Option { + with_handle::(handle, |wrapper| { + MysqlConnectionTarget::Direct(Arc::clone(&wrapper.connection)) + }) + .or_else(|| { + with_handle::(handle, |wrapper| { + MysqlConnectionTarget::Pool(Arc::clone(&wrapper.connection)) + }) + }) +} + +async fn execute_query_on_connection( + conn: &mut MySqlConnection, + request: &QueryRequest, +) -> Result { + let is_select = request.is_row_returning(); + + if !request.uses_prepared_statement() { + // SQLx's `query()` prepares even when there are no bind values. That + // needlessly put mysql2 `query("DROP ...")` / `query("CREATE ...")` + // calls into the per-connection statement cache beside parameterized + // `execute()` calls. Use MySQL's text protocol for the no-param query + // shape, matching mysql2 and keeping those statements out of the cache. + let raw = sqlx::raw_sql(sqlx::AssertSqlSafe(request.sql.clone())); + if is_select { + let rows = tokio::time::timeout( + Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), + raw.fetch_all(conn), + ) + .await + .map_err(|_| "Query timed out".to_string())? + .map_err(|e| format!("Query failed: {}", e))?; + return Ok(QueryOutcome::Rows(raws_from_mysql_rows(rows))); + } + + let res = tokio::time::timeout( + Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), + raw.execute(conn), + ) + .await + .map_err(|_| "Query timed out".to_string())? + .map_err(|e| format!("Query failed: {}", e))?; + return Ok(QueryOutcome::Executed { + affected_rows: res.rows_affected(), + last_insert_id: res.last_insert_id(), + }); + } + + // Build the SQLx query and all of its arguments from the same owned + // request immediately before execution. Nothing is shared with another + // mysql2 call, even while this future is waiting on I/O. + // Keep the prepared statement scoped to this request. SQLx's connection + // cache is where #8745 observed metadata from a neighboring statement + // being paired with this request's arguments; an ephemeral statement + // preserves mysql2 execute semantics without reusing that association. + let mut query = sqlx::query(sqlx::AssertSqlSafe(request.sql.clone())).persistent(false); + for param in &request.params { + query = match param { + ParamValue::Null => query.bind(Option::::None), + ParamValue::String(s) => query.bind(s.clone()), + ParamValue::Number(n) => query.bind(*n), + ParamValue::Int(i) => query.bind(*i), + ParamValue::Bool(b) => query.bind(*b), + }; + } + + if is_select { + let rows = tokio::time::timeout( + Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), + query.fetch_all(conn), + ) + .await + .map_err(|_| "Query timed out".to_string())? + .map_err(|e| format!("Query failed: {}", e))?; + Ok(QueryOutcome::Rows(raws_from_mysql_rows(rows))) + } else { + let res = tokio::time::timeout( + Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), + query.execute(conn), + ) + .await + .map_err(|_| "Query timed out".to_string())? + .map_err(|e| format!("Query failed: {}", e))?; + Ok(QueryOutcome::Executed { + affected_rows: res.rows_affected(), + last_insert_id: res.last_insert_id(), + }) + } +} + +async fn execute_query_on_target( + target: MysqlConnectionTarget, + request: &QueryRequest, +) -> Result { + match target { + MysqlConnectionTarget::Direct(connection) => { + let mut slot = connection.lock().await; + let conn = slot + .as_mut() + .ok_or_else(|| "Connection already closed".to_string())?; + execute_query_on_connection(conn, request).await + } + MysqlConnectionTarget::Pool(connection) => { + let mut slot = connection.lock().await; + let conn = slot + .as_mut() + .ok_or_else(|| "Pool connection released".to_string())?; + execute_query_on_connection(conn, request).await } } } @@ -575,8 +740,11 @@ pub extern "C" fn js_mysql2_connection_end(conn_handle: Handle) -> *mut Promise let promise = JsPromise::new(); let raw = promise.as_raw(); spawn_blocking(move || { - if let Some(mut wrapper) = take_handle::(conn_handle) { - if let Some(conn) = wrapper.connection.take() { + if let Some(wrapper) = take_handle::(conn_handle) { + let connection = Arc::clone(&wrapper.connection); + let conn = tokio::runtime::Handle::current() + .block_on(async move { connection.lock().await.take() }); + if let Some(conn) = conn { let result = tokio::runtime::Handle::current().block_on(conn.close()); match result { Ok(()) => promise.resolve_undefined(), @@ -597,56 +765,23 @@ unsafe fn run_connection_query( sql_ptr: *const u8, params_f: f64, rows_as_array: bool, + force_prepared: bool, ) -> *mut Promise { let sql = read_sql(sql_ptr); let params = JsValue::from_bits(params_f.to_bits()); let param_values = extract_params_from_jsvalue(params); - let is_select = is_row_returning_query(&sql); + let request = QueryRequest::new(sql, param_values, rows_as_array, force_prepared); + let target = connection_target(conn_handle); let promise = JsPromise::new(); let raw = promise.as_raw(); spawn_blocking(move || { + let rows_as_array = request.rows_as_array; let outcome: Result = tokio::runtime::Handle::current().block_on(async move { - let wrapper = get_handle_mut::(conn_handle) - .ok_or_else(|| "Invalid connection handle".to_string())?; - let conn = wrapper - .connection - .as_mut() - .ok_or_else(|| "Connection already closed".to_string())?; - let mut q = sqlx::query(sqlx::AssertSqlSafe(sql.clone())); - for p in ¶m_values { - q = match p { - ParamValue::Null => q.bind(Option::::None), - ParamValue::String(s) => q.bind(s.clone()), - ParamValue::Number(n) => q.bind(*n), - ParamValue::Int(i) => q.bind(*i), - ParamValue::Bool(b) => q.bind(*b), - }; - } - if is_select { - let rows = tokio::time::timeout( - Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), - q.fetch_all(conn), - ) - .await - .map_err(|_| "Query timed out".to_string())? - .map_err(|e| format!("Query failed: {}", e))?; - Ok(QueryOutcome::Rows(raws_from_mysql_rows(rows))) - } else { - let res = tokio::time::timeout( - Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), - q.execute(conn), - ) - .await - .map_err(|_| "Query timed out".to_string())? - .map_err(|e| format!("Query failed: {}", e))?; - Ok(QueryOutcome::Executed { - affected_rows: res.rows_affected(), - last_insert_id: res.last_insert_id(), - }) - } + let target = target.ok_or_else(|| "Invalid connection handle".to_string())?; + execute_query_on_target(target, &request).await }); match outcome { // #1824: build the JS result on the MAIN thread. outcome_to_jsvalue @@ -670,7 +805,7 @@ pub unsafe extern "C" fn js_mysql2_connection_query( sql_ptr: *const u8, params_f: f64, ) -> *mut Promise { - run_connection_query(conn_handle, sql_ptr, params_f, false) + run_connection_query(conn_handle, sql_ptr, params_f, false, false) } /// `connection.execute(sql, params) -> Promise<[rows, fields]>`. @@ -684,25 +819,40 @@ pub unsafe extern "C" fn js_mysql2_connection_execute( sql_ptr: *const u8, params_f: f64, ) -> *mut Promise { - run_connection_query(conn_handle, sql_ptr, params_f, false) + run_connection_query(conn_handle, sql_ptr, params_f, false, true) } fn run_simple_command(conn_handle: Handle, sql: &'static str) -> *mut Promise { + let target = connection_target(conn_handle); let promise = JsPromise::new(); let raw = promise.as_raw(); spawn_blocking(move || { let result = tokio::runtime::Handle::current().block_on(async move { - let wrapper = get_handle_mut::(conn_handle) - .ok_or_else(|| "Invalid connection handle".to_string())?; - let conn = wrapper - .connection - .as_mut() - .ok_or_else(|| "Connection already closed".to_string())?; - sqlx::query(sql) - .execute(conn) - .await - .map(|_| ()) - .map_err(|e| format!("{}: {}", sql, e)) + let target = target.ok_or_else(|| "Invalid connection handle".to_string())?; + match target { + MysqlConnectionTarget::Direct(connection) => { + let mut slot = connection.lock().await; + let conn = slot + .as_mut() + .ok_or_else(|| "Connection already closed".to_string())?; + sqlx::raw_sql(sql) + .execute(conn) + .await + .map(|_| ()) + .map_err(|e| format!("{}: {}", sql, e)) + } + MysqlConnectionTarget::Pool(connection) => { + let mut slot = connection.lock().await; + let conn = slot + .as_mut() + .ok_or_else(|| "Pool connection released".to_string())?; + sqlx::raw_sql(sql) + .execute(&mut **conn) + .await + .map(|_| ()) + .map_err(|e| format!("{}: {}", sql, e)) + } + } }); match result { Ok(()) => promise.resolve_undefined(), @@ -712,6 +862,15 @@ fn run_simple_command(conn_handle: Handle, sql: &'static str) -> *mut Promise { raw } +fn transaction_sql_for_method(method: &str) -> Option<&'static str> { + match method { + "beginTransaction" => Some("START TRANSACTION"), + "commit" => Some("COMMIT"), + "rollback" => Some("ROLLBACK"), + _ => None, + } +} + #[no_mangle] pub extern "C" fn js_mysql2_connection_begin_transaction(conn_handle: Handle) -> *mut Promise { run_simple_command(conn_handle, "START TRANSACTION") @@ -740,13 +899,13 @@ impl MysqlPoolHandle { } pub struct MysqlPoolConnectionHandle { - pub connection: Option>, + pub connection: Arc>>>, } impl MysqlPoolConnectionHandle { pub fn new(conn: PoolConnection) -> Self { Self { - connection: Some(conn), + connection: Arc::new(Mutex::new(Some(conn))), } } } @@ -898,9 +1057,9 @@ unsafe extern "C" fn js_mysql2_handle_method_dispatch( }; // Only claim methods for handles we actually own. - let is_pool = perry_ffi::get_handle::(handle).is_some(); - let is_pool_conn = perry_ffi::get_handle::(handle).is_some(); - let is_conn = perry_ffi::get_handle::(handle).is_some(); + let is_pool = with_handle::(handle, |_| ()).is_some(); + let is_pool_conn = with_handle::(handle, |_| ()).is_some(); + let is_conn = with_handle::(handle, |_| ()).is_some(); if !is_pool && !is_pool_conn && !is_conn { return 0; } @@ -914,12 +1073,13 @@ unsafe extern "C" fn js_mysql2_handle_method_dispatch( .get(1) .copied() .unwrap_or(f64::from_bits(DISPATCH_TAG_UNDEFINED)); + let force_prepared = method == "execute"; let promise = if is_pool { - run_pool_query(handle, sql_ptr, params_f, rows_as_array) + run_pool_query(handle, sql_ptr, params_f, rows_as_array, force_prepared) } else if is_pool_conn { - run_pool_conn_query(handle, sql_ptr, params_f, rows_as_array) + run_pool_conn_query(handle, sql_ptr, params_f, rows_as_array, force_prepared) } else { - run_connection_query(handle, sql_ptr, params_f, rows_as_array) + run_connection_query(handle, sql_ptr, params_f, rows_as_array, force_prepared) }; dispatch_nanbox_ptr(promise) } @@ -929,6 +1089,12 @@ unsafe extern "C" fn js_mysql2_handle_method_dispatch( js_mysql2_pool_connection_release(handle); f64::from_bits(DISPATCH_TAG_UNDEFINED) } + method if (is_pool_conn || is_conn) && transaction_sql_for_method(method).is_some() => { + let Some(sql) = transaction_sql_for_method(method) else { + return 0; + }; + dispatch_nanbox_ptr(run_simple_command(handle, sql)) + } "end" if is_conn => dispatch_nanbox_ptr(js_mysql2_connection_end(handle)), // `mysql2/promise` pools are already promise-based: `pool.promise()` // returns the pool itself. Drizzle's `isCallbackClient` only reaches this @@ -963,52 +1129,32 @@ unsafe fn run_pool_query( sql_ptr: *const u8, params_f: f64, rows_as_array: bool, + force_prepared: bool, ) -> *mut Promise { let sql = read_sql(sql_ptr); let params = JsValue::from_bits(params_f.to_bits()); let param_values = extract_params_from_jsvalue(params); - let is_select = is_row_returning_query(&sql); + let request = QueryRequest::new(sql, param_values, rows_as_array, force_prepared); + let pool = with_handle::(pool_handle, |wrapper| wrapper.pool.clone()); let promise = JsPromise::new(); let raw = promise.as_raw(); spawn_blocking(move || { + let rows_as_array = request.rows_as_array; let outcome: Result = tokio::runtime::Handle::current().block_on(async move { - let wrapper = get_handle_mut::(pool_handle) - .ok_or_else(|| "Invalid pool handle".to_string())?; - let pool = &wrapper.pool; - let mut q = sqlx::query(sqlx::AssertSqlSafe(sql.clone())); - for p in ¶m_values { - q = match p { - ParamValue::Null => q.bind(Option::::None), - ParamValue::String(s) => q.bind(s.clone()), - ParamValue::Number(n) => q.bind(*n), - ParamValue::Int(i) => q.bind(*i), - ParamValue::Bool(b) => q.bind(*b), - }; - } - if is_select { - let rows = tokio::time::timeout( - Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), - q.fetch_all(pool), - ) - .await - .map_err(|_| "Query timed out".to_string())? - .map_err(|e| format!("Query failed: {}", e))?; - Ok(QueryOutcome::Rows(raws_from_mysql_rows(rows))) - } else { - let res = tokio::time::timeout( - Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), - q.execute(pool), - ) - .await - .map_err(|_| "Query timed out".to_string())? - .map_err(|e| format!("Query failed: {}", e))?; - Ok(QueryOutcome::Executed { - affected_rows: res.rows_affected(), - last_insert_id: res.last_insert_id(), - }) - } + let pool = pool.ok_or_else(|| "Invalid pool handle".to_string())?; + // Explicitly check out one connection for the whole request so + // statement preparation, bind encoding, execution, and result + // draining cannot be split across independent pool operations. + let mut conn = tokio::time::timeout( + Duration::from_secs(DEFAULT_ACQUIRE_TIMEOUT_SECS), + pool.acquire(), + ) + .await + .map_err(|_| "Pool acquire timed out".to_string())? + .map_err(|e| format!("Pool acquire failed: {}", e))?; + execute_query_on_connection(&mut conn, &request).await }); match outcome { // #1824: build the JS result on the MAIN thread. outcome_to_jsvalue @@ -1030,7 +1176,7 @@ pub unsafe extern "C" fn js_mysql2_pool_query( sql_ptr: *const u8, params_f: f64, ) -> *mut Promise { - run_pool_query(pool_handle, sql_ptr, params_f, false) + run_pool_query(pool_handle, sql_ptr, params_f, false, false) } /// # Safety @@ -1041,20 +1187,20 @@ pub unsafe extern "C" fn js_mysql2_pool_execute( sql_ptr: *const u8, params_f: f64, ) -> *mut Promise { - run_pool_query(pool_handle, sql_ptr, params_f, false) + run_pool_query(pool_handle, sql_ptr, params_f, false, true) } #[no_mangle] pub extern "C" fn js_mysql2_pool_get_connection(pool_handle: Handle) -> *mut Promise { + let pool = with_handle::(pool_handle, |wrapper| wrapper.pool.clone()); let promise = JsPromise::new(); let raw = promise.as_raw(); spawn_blocking(move || { let result = tokio::runtime::Handle::current().block_on(async move { - let wrapper = get_handle_mut::(pool_handle) - .ok_or_else(|| "Invalid pool handle".to_string())?; + let pool = pool.ok_or_else(|| "Invalid pool handle".to_string())?; tokio::time::timeout( Duration::from_secs(DEFAULT_ACQUIRE_TIMEOUT_SECS), - wrapper.pool.acquire(), + pool.acquire(), ) .await .map_err(|_| "Pool acquire timed out".to_string())? @@ -1075,7 +1221,15 @@ pub extern "C" fn js_mysql2_pool_get_connection(pool_handle: Handle) -> *mut Pro /// underlying `PoolConnection` returns to the pool via Drop. #[no_mangle] pub extern "C" fn js_mysql2_pool_connection_release(conn_handle: Handle) { - take_handle::(conn_handle); + if let Some(wrapper) = take_handle::(conn_handle) { + // A query already in flight owns another Arc and holds this mutex. Wait + // for it to finish before dropping the checkout back into the pool. + spawn_blocking(move || { + tokio::runtime::Handle::current().block_on(async move { + wrapper.connection.lock().await.take(); + }); + }); + } } unsafe fn run_pool_conn_query( @@ -1083,55 +1237,29 @@ unsafe fn run_pool_conn_query( sql_ptr: *const u8, params_f: f64, rows_as_array: bool, + force_prepared: bool, ) -> *mut Promise { let sql = read_sql(sql_ptr); let params = JsValue::from_bits(params_f.to_bits()); let param_values = extract_params_from_jsvalue(params); - let is_select = is_row_returning_query(&sql); + let request = QueryRequest::new(sql, param_values, rows_as_array, force_prepared); + let connection = with_handle::(conn_handle, |wrapper| { + Arc::clone(&wrapper.connection) + }); let promise = JsPromise::new(); let raw = promise.as_raw(); spawn_blocking(move || { + let rows_as_array = request.rows_as_array; let outcome: Result = tokio::runtime::Handle::current().block_on(async move { - let wrapper = get_handle_mut::(conn_handle) - .ok_or_else(|| "Invalid pool-connection handle".to_string())?; - let conn = wrapper - .connection + let connection = + connection.ok_or_else(|| "Invalid pool-connection handle".to_string())?; + let mut slot = connection.lock().await; + let conn = slot .as_mut() .ok_or_else(|| "Pool connection released".to_string())?; - let mut q = sqlx::query(sqlx::AssertSqlSafe(sql.clone())); - for p in ¶m_values { - q = match p { - ParamValue::Null => q.bind(Option::::None), - ParamValue::String(s) => q.bind(s.clone()), - ParamValue::Number(n) => q.bind(*n), - ParamValue::Int(i) => q.bind(*i), - ParamValue::Bool(b) => q.bind(*b), - }; - } - if is_select { - let rows = tokio::time::timeout( - Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), - q.fetch_all(&mut **conn), - ) - .await - .map_err(|_| "Query timed out".to_string())? - .map_err(|e| format!("Query failed: {}", e))?; - Ok(QueryOutcome::Rows(raws_from_mysql_rows(rows))) - } else { - let res = tokio::time::timeout( - Duration::from_secs(DEFAULT_QUERY_TIMEOUT_SECS), - q.execute(&mut **conn), - ) - .await - .map_err(|_| "Query timed out".to_string())? - .map_err(|e| format!("Query failed: {}", e))?; - Ok(QueryOutcome::Executed { - affected_rows: res.rows_affected(), - last_insert_id: res.last_insert_id(), - }) - } + execute_query_on_connection(conn, &request).await }); match outcome { // #1824: build the JS result on the MAIN thread. outcome_to_jsvalue @@ -1153,7 +1281,7 @@ pub unsafe extern "C" fn js_mysql2_pool_connection_query( sql_ptr: *const u8, params_f: f64, ) -> *mut Promise { - run_pool_conn_query(conn_handle, sql_ptr, params_f, false) + run_pool_conn_query(conn_handle, sql_ptr, params_f, false, false) } /// # Safety @@ -1164,7 +1292,7 @@ pub unsafe extern "C" fn js_mysql2_pool_connection_execute( sql_ptr: *const u8, params_f: f64, ) -> *mut Promise { - run_pool_conn_query(conn_handle, sql_ptr, params_f, false) + run_pool_conn_query(conn_handle, sql_ptr, params_f, false, true) } #[cfg(test)] @@ -1235,4 +1363,87 @@ mod tests { assert!(!is_row_returning_query("UPDATE t SET x = 1")); assert!(!is_row_returning_query("DELETE FROM t")); } + + #[test] + fn query_request_keeps_each_statement_with_its_own_params() { + let ddl = QueryRequest::new("DROP TABLE IF EXISTS t".into(), Vec::new(), false, false); + let insert = QueryRequest::new( + "INSERT INTO t (name, cents) VALUES (?, ?)".into(), + vec![ParamValue::String("x".into()), ParamValue::Int(100)], + false, + true, + ); + let select = QueryRequest::new( + "SELECT * FROM t WHERE id = ?".into(), + vec![ParamValue::Int(1)], + false, + true, + ); + + assert_eq!(ddl.params, Vec::::new()); + assert_eq!(insert.params.len(), 2); + assert_eq!(select.params, vec![ParamValue::Int(1)]); + assert!(!ddl.uses_prepared_statement()); + assert!(insert.uses_prepared_statement()); + assert!(select.uses_prepared_statement()); + } + + #[test] + fn execute_stays_prepared_even_without_params() { + let execute = QueryRequest::new("SELECT 1".into(), Vec::new(), false, true); + assert!(execute.uses_prepared_statement()); + } + + #[test] + fn both_connection_handle_families_resolve_to_serialized_targets() { + let direct_connection = Arc::new(Mutex::new(None)); + let direct_handle = register_handle(MysqlConnectionHandle { + connection: Arc::clone(&direct_connection), + }); + let pool_connection = Arc::new(Mutex::new(None)); + let pool_handle = register_handle(MysqlPoolConnectionHandle { + connection: Arc::clone(&pool_connection), + }); + + match connection_target(direct_handle) { + Some(MysqlConnectionTarget::Direct(resolved)) => { + assert!(Arc::ptr_eq(&resolved, &direct_connection)); + let _guard = resolved + .try_lock() + .expect("first operation locks connection"); + assert!( + direct_connection.try_lock().is_err(), + "a second operation on the same connection must serialize" + ); + } + _ => panic!("direct connection handle was not resolved"), + } + match connection_target(pool_handle) { + Some(MysqlConnectionTarget::Pool(resolved)) => { + assert!(Arc::ptr_eq(&resolved, &pool_connection)); + let _guard = resolved + .try_lock() + .expect("first operation locks pool connection"); + assert!( + pool_connection.try_lock().is_err(), + "a second operation on the same checkout must serialize" + ); + } + _ => panic!("pool connection handle was not resolved"), + } + + take_handle::(direct_handle); + take_handle::(pool_handle); + } + + #[test] + fn pool_connections_expose_the_full_transaction_command_set() { + assert_eq!( + transaction_sql_for_method("beginTransaction"), + Some("START TRANSACTION") + ); + assert_eq!(transaction_sql_for_method("commit"), Some("COMMIT")); + assert_eq!(transaction_sql_for_method("rollback"), Some("ROLLBACK")); + assert_eq!(transaction_sql_for_method("release"), None); + } } diff --git a/crates/perry-ext-mysql2/src/test_async_shims.rs b/crates/perry-ext-mysql2/src/test_async_shims.rs new file mode 100644 index 0000000000..47945057b6 --- /dev/null +++ b/crates/perry-ext-mysql2/src/test_async_shims.rs @@ -0,0 +1,103 @@ +//! Test-only host shims for the standalone extension test binary. +//! +//! Production binaries receive these symbols from perry-stdlib's async bridge. + +use perry_ffi::{NativeAsyncCompletion, Promise}; +use std::ffi::c_void; + +#[no_mangle] +pub extern "C" fn perry_ffi_promise_new() -> *mut Promise { + perry_runtime::promise::js_promise_new() as *mut Promise +} + +#[no_mangle] +pub extern "C" fn perry_ffi_promise_resolve_bits(promise: *mut Promise, bits: u64) { + perry_runtime::promise::js_promise_resolve( + promise as *mut perry_runtime::Promise, + f64::from_bits(bits), + ); +} + +#[no_mangle] +pub extern "C" fn perry_ffi_promise_reject_bits(promise: *mut Promise, bits: u64) { + perry_runtime::promise::js_promise_reject( + promise as *mut perry_runtime::Promise, + f64::from_bits(bits), + ); +} + +#[no_mangle] +pub extern "C" fn perry_ffi_promise_resolve_deferred( + promise: *mut Promise, + ctx: *mut c_void, + invoke: extern "C" fn(*mut c_void) -> u64, +) { + perry_ffi_promise_resolve_bits(promise, invoke(ctx)); +} + +#[no_mangle] +pub extern "C" fn perry_ffi_spawn_blocking(ctx: *mut c_void, invoke: extern "C" fn(*mut c_void)) { + invoke(ctx); +} + +#[no_mangle] +pub extern "C" fn perry_ffi_spawn_blocking_with_reactor( + ctx: *mut c_void, + invoke: extern "C" fn(*mut c_void), +) { + invoke(ctx); +} + +#[no_mangle] +pub extern "C" fn perry_ffi_native_async_new(_flags: u32) -> *mut NativeAsyncCompletion { + std::ptr::null_mut() +} + +#[no_mangle] +pub extern "C" fn perry_ffi_native_async_promise( + _token: *mut NativeAsyncCompletion, +) -> *mut Promise { + std::ptr::null_mut() +} + +#[no_mangle] +pub extern "C" fn perry_ffi_native_async_resolve_bits( + _token: *mut NativeAsyncCompletion, + _bits: u64, +) -> i32 { + 0 +} + +#[no_mangle] +pub extern "C" fn perry_ffi_native_async_reject_bits( + _token: *mut NativeAsyncCompletion, + _bits: u64, +) -> i32 { + 0 +} + +#[no_mangle] +pub extern "C" fn perry_ffi_native_async_reject_string( + _token: *mut NativeAsyncCompletion, + _data: *const u8, + _len: usize, +) -> i32 { + 0 +} + +#[no_mangle] +pub extern "C" fn perry_ffi_native_async_cancel(_token: *mut NativeAsyncCompletion) -> i32 { + 0 +} + +#[no_mangle] +pub extern "C" fn perry_ffi_native_async_attach_handle( + _token: *mut NativeAsyncCompletion, + _handle_bits: u64, + _cleanup_flags: u32, +) -> i32 { + 0 +} + +#[no_mangle] +pub extern "C" fn perry_ffi_run_pending(_budget_ms: u64) {} diff --git a/test-files/test_issue_8745_8746_mysql2_operation_isolation.ts b/test-files/test_issue_8745_8746_mysql2_operation_isolation.ts new file mode 100644 index 0000000000..eb53011232 --- /dev/null +++ b/test-files/test_issue_8745_8746_mysql2_operation_isolation.ts @@ -0,0 +1,69 @@ +// parity-skip: requires a live MySQL fixture; native wrapper unit-tested +// Regression coverage for issues #8745 and #8746. +// +// Run against a local MySQL database after removing the skip marker. The loop +// mixes text-protocol DDL with prepared statements of different arities; every +// statement must retain its own parameter vector. The checked-out connection +// then exercises the canonical row-lock transaction lifecycle. +// +// platforms: skip + +import mysql from 'mysql2/promise'; + +const pool = mysql.createPool({ + host: 'localhost', port: 3306, user: 'perry', password: 'perry', + database: 'perry_hub', +}); + +async function main(): Promise { + for (let round = 0; round < 25; round++) { + await pool.query('DROP TABLE IF EXISTS perry_issue_8745'); + await pool.query( + 'CREATE TABLE perry_issue_8745 (' + + 'id INT AUTO_INCREMENT PRIMARY KEY, name VARCHAR(50), cents INT)', + ); + await pool.execute( + 'INSERT INTO perry_issue_8745 (name, cents) VALUES (?, ?)', + ['round-' + round, 100 + round], + ); + const selected: any = await pool.execute( + 'SELECT * FROM perry_issue_8745 WHERE id = ?', + [1], + ); + if (selected[0][0].cents !== 100 + round) { + throw new Error('wrong prepared-statement parameters in round ' + round); + } + } + + const connection = await pool.getConnection(); + try { + await connection.beginTransaction(); + const locked: any = await connection.execute( + 'SELECT cents FROM perry_issue_8745 WHERE id = ? FOR UPDATE', + [1], + ); + await connection.execute( + 'UPDATE perry_issue_8745 SET cents = ? WHERE id = ?', + [locked[0][0].cents + 1, 1], + ); + await connection.commit(); + + await connection.beginTransaction(); + await connection.execute( + 'UPDATE perry_issue_8745 SET cents = ? WHERE id = ?', + [9999, 1], + ); + await connection.rollback(); + } finally { + connection.release(); + } + + const finalRows: any = await pool.execute( + 'SELECT cents FROM perry_issue_8745 WHERE id = ?', + [1], + ); + console.log(finalRows[0][0].cents); + await pool.end(); +} + +main();