rustyconover opened a new issue, #4817:
URL: https://github.com/apache/arrow-adbc/issues/4817

   ### What happened?
   
   The Rust driver manager's separate cancellation handles still acquire the 
same mutex held by the operation they are supposed to interrupt. As a result, 
the native cancellation callback cannot run until that operation returns:
   
   - `StatementCancelHandle::try_cancel()` locks `inner.statement`, while 
`execute()` holds that lock throughout `StatementExecuteQuery`.
   - `ConnectionCancelHandle::try_cancel()` locks `inner.connection`, while 
connection operations such as `commit()` hold it throughout the native call.
   
   I reproduced both cases against unmodified upstream `main` at 
`c942d481c6e083040c68676e3dd454dad89503e9`, using the standalone mock C-ABI 
driver below. The mock operation blocks until its cancellation callback 
releases it, with explicit timeouts/cleanup to keep the reproducer from 
hanging. Neither cancellation callback is reached while the operation is 
blocked; both are reached after the reproducer manually releases the operation.
   
   Expected: a cancellation handle obtained before execution should be able to 
invoke the native cancel callback concurrently with that execution, while 
ordinary operations remain serialized and the native handle remains alive. This 
does not require every driver to support cancellation or guarantee that 
cancellation wins a race with completion.
   
   Actual: cancellation itself waits for completion. For an operation that 
requires cancellation to finish, this can deadlock without an independent 
timeout or release.
   
   Relevant source at the tested revision:
   
   - [Statement 
cancellation](https://github.com/apache/arrow-adbc/blob/c942d481c6e083040c68676e3dd454dad89503e9/rust/driver_manager/src/lib.rs#L1295)
 and [statement 
execution](https://github.com/apache/arrow-adbc/blob/c942d481c6e083040c68676e3dd454dad89503e9/rust/driver_manager/src/lib.rs#L1361).
   - [Connection 
cancellation](https://github.com/apache/arrow-adbc/blob/c942d481c6e083040c68676e3dd454dad89503e9/rust/driver_manager/src/lib.rs#L825)
 and 
[commit](https://github.com/apache/arrow-adbc/blob/c942d481c6e083040c68676e3dd454dad89503e9/rust/driver_manager/src/lib.rs#L985).
   - [ADBC cancellation 
contract](https://github.com/apache/arrow-adbc/blob/c942d481c6e083040c68676e3dd454dad89503e9/c/include/arrow-adbc/adbc.h):
 cancellation must be thread-safe and can be called during execution.
   
   This is a follow-up to #3454 / #3905: the separate cancellation handle makes 
the call expressible in Rust, but the manager's shared operation mutex still 
prevents it reaching the driver concurrently. It is distinct from #4474, which 
concerns reporting cancellation through an Arrow result stream.
   
   We originally encountered the statement case in Grainlift's real-driver 
tests: direct statement cancellation interrupted a DuckDB 1.5.5 query, whereas 
the path through the Rust driver manager timed out. The standalone reproduction 
below removes the proxy, network, DuckDB, and Rust driver exporter from the 
test.
   
   ### Stack Trace
   
   This is a blocked cancellation call rather than a crash. The bounded 
reproducer exits with this assertion after releasing both blocked operations:
   
   ```text
   statement execute: cancel callback reached while operation blocked = false
   connection commit: cancel callback reached while operation blocked = false
   
   thread 'main' (...) panicked at src/main.rs:170:5:
   assertion `left == right` failed: cancellation waited for the 
ordinary-operation mutex
     left: 2
    right: 0
   ```
   
   ### How can we reproduce the bug?
   
   Create an empty directory containing the following `Cargo.toml` and 
`src/main.rs`, then run:
   
   ```console
   cargo run
   ```
   
   The test uses a static mock C-ABI driver, so no database, external driver 
library, credentials, or service is needed. Each operation has a 10-second 
safety timeout; the caller waits up to one second for cancellation to reach the 
native callback before manually releasing the operation and joining both 
threads. The ordinary owners are retained until cancellation finishes.
   
   Expected output is `true` for both cases, with exit status 0. Actual output 
is `false` for both, with exit status 101.
   
   <details>
   <summary>Cargo.toml</summary>
   
   ```toml
   [package]
   name = "adbc-cancel-upstream-repro"
   version = "0.1.0"
   edition = "2024"
   publish = false
   
   [dependencies]
   adbc_core = { git = "https://github.com/apache/arrow-adbc.git";, rev = 
"c942d481c6e083040c68676e3dd454dad89503e9" }
   adbc_driver_manager = { git = "https://github.com/apache/arrow-adbc.git";, 
rev = "c942d481c6e083040c68676e3dd454dad89503e9" }
   adbc_ffi = { git = "https://github.com/apache/arrow-adbc.git";, rev = 
"c942d481c6e083040c68676e3dd454dad89503e9" }
   arrow-array = { version = "60", features = ["ffi"] }
   ```
   
   </details>
   
   <details>
   <summary>src/main.rs</summary>
   
   ```rust
   // Copyright (c) 2026 Query Farm LLC
   // SPDX-License-Identifier: Apache-2.0
   
   use std::ffi::{c_int, c_void};
   use std::sync::atomic::{AtomicBool, Ordering::SeqCst};
   use std::thread;
   use std::time::{Duration, Instant};
   
   use adbc_core::error::{AdbcStatusCode, Status};
   use adbc_core::options::AdbcVersion;
   use adbc_core::{Connection, Database, Driver, Statement};
   use adbc_driver_manager::ManagedDriver;
   use adbc_ffi::*;
   use arrow_array::ffi_stream::FFI_ArrowArrayStream;
   
   // One operation at a time; callbacks use no mutable handle state.
   static ENTERED: AtomicBool = AtomicBool::new(false);
   static RELEASE: AtomicBool = AtomicBool::new(false);
   static CANCEL_ENTERED: AtomicBool = AtomicBool::new(false);
   
   unsafe extern "C" fn database_noop(
       _: *mut FFI_AdbcDatabase,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       0
   }
   unsafe extern "C" fn connection_noop(
       _: *mut FFI_AdbcConnection,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       0
   }
   unsafe extern "C" fn statement_noop(
       _: *mut FFI_AdbcStatement,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       0
   }
   unsafe extern "C" fn connection_init(
       _: *mut FFI_AdbcConnection,
       _: *mut FFI_AdbcDatabase,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       0
   }
   unsafe extern "C" fn statement_new(
       _: *mut FFI_AdbcConnection,
       _: *mut FFI_AdbcStatement,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       0
   }
   
   fn wait_for(flag: &AtomicBool, duration: Duration) -> bool {
       let deadline = Instant::now() + duration;
       while !flag.load(SeqCst) && Instant::now() < deadline {
           thread::sleep(Duration::from_millis(1));
       }
       flag.load(SeqCst)
   }
   
   fn block_operation() -> AdbcStatusCode {
       ENTERED.store(true, SeqCst);
       // Safety timeout: the reproducer must terminate even with broken 
cancellation.
       wait_for(&RELEASE, Duration::from_secs(10));
       Status::Cancelled.into()
   }
   
   fn cancel_operation() -> AdbcStatusCode {
       CANCEL_ENTERED.store(true, SeqCst);
       RELEASE.store(true, SeqCst);
       0
   }
   
   unsafe extern "C" fn execute(
       _: *mut FFI_AdbcStatement,
       _: *mut FFI_ArrowArrayStream,
       _: *mut i64,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       block_operation()
   }
   unsafe extern "C" fn commit(_: *mut FFI_AdbcConnection, _: *mut 
FFI_AdbcError) -> AdbcStatusCode {
       block_operation()
   }
   unsafe extern "C" fn statement_cancel(
       _: *mut FFI_AdbcStatement,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       cancel_operation()
   }
   unsafe extern "C" fn connection_cancel(
       _: *mut FFI_AdbcConnection,
       _: *mut FFI_AdbcError,
   ) -> AdbcStatusCode {
       cancel_operation()
   }
   
   unsafe extern "C" fn init(_: c_int, driver: *mut c_void, _: *mut 
FFI_AdbcError) -> AdbcStatusCode {
       unsafe {
           *driver.cast::<FFI_AdbcDriver>() = FFI_AdbcDriver {
               DatabaseNew: Some(database_noop),
               DatabaseInit: Some(database_noop),
               DatabaseRelease: Some(database_noop),
               ConnectionNew: Some(connection_noop),
               ConnectionInit: Some(connection_init),
               ConnectionRelease: Some(connection_noop),
               ConnectionCommit: Some(commit),
               ConnectionCancel: Some(connection_cancel),
               StatementNew: Some(statement_new),
               StatementRelease: Some(statement_noop),
               StatementExecuteQuery: Some(execute),
               StatementCancel: Some(statement_cancel),
               ..FFI_AdbcDriver::default()
           };
       }
       0
   }
   
   fn main() {
       let mut failures = 0;
       for scope in ["statement execute", "connection commit"] {
           ENTERED.store(false, SeqCst);
           RELEASE.store(false, SeqCst);
           CANCEL_ENTERED.store(false, SeqCst);
           let mut driver =
               ManagedDriver::load_static(&(init as FFI_AdbcDriverInitFunc), 
AdbcVersion::V110)
                   .unwrap();
           let database = driver.new_database().unwrap();
           let mut connection = database.new_connection().unwrap();
           // Keep the ordinary owner alive until cancellation has joined.
           thread::scope(|threads| {
               let mut statement = connection.new_statement().unwrap();
               let cancel = if scope == "statement execute" {
                   statement.get_cancel_handle()
               } else {
                   connection.get_cancel_handle()
               };
               let execution = threads.spawn(move || {
                   let status = if scope == "statement execute" {
                       let result = statement.execute().err().unwrap().status;
                       // Return the owners to prevent release racing 
cancellation.
                       (result, statement, connection)
                   } else {
                       (
                           connection.commit().unwrap_err().status,
                           statement,
                           connection,
                       )
                   };
                   status
               });
               assert!(
                   wait_for(&ENTERED, Duration::from_secs(2)),
                   "operation did not start"
               );
               let cancelling = threads.spawn(move || cancel.try_cancel());
               let overlapped = wait_for(&CANCEL_ENTERED, 
Duration::from_secs(1));
               println!("{scope}: cancel callback reached while operation 
blocked = {overlapped}");
               RELEASE.store(true, SeqCst);
               let (status, _statement, _connection) = 
execution.join().unwrap();
               cancelling.join().unwrap().unwrap();
               assert_eq!(status, Status::Cancelled);
               assert!(CANCEL_ENTERED.load(SeqCst));
               if !overlapped {
                   failures += 1;
               }
           });
       }
       assert_eq!(
           failures, 0,
           "cancellation waited for the ordinary-operation mutex"
       );
   }
   ```
   
   </details>
   
   ### Environment/Setup
   
   - Amazon Linux 2023, Linux aarch64.
   - Rust `1.97.1 (8bab26f4f 2026-07-14)`, debug build.
   - `adbc_core`, `adbc_driver_manager`, and `adbc_ffi` 0.25.0 from upstream 
Git revision `c942d481c6e083040c68676e3dd454dad89503e9` (current `main` when 
checked).
   - Arrow Rust 60.0.0.
   - No local Cargo patches in the standalone reproduction.
   - The statement issue was also observed in our integration using upstream 
revision `616acfdfcea9b66956fdb3d11437b6cf24edbc39`.
   
   A fix should allow cancellation to overlap the native operation without 
permitting concurrent ordinary operations or releasing the native handle while 
cancellation is in flight. A regression test that waits for entry into a 
blocking native callback is useful here: cancellation tests against idle 
handles do not exercise the locking problem.
   
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to