From 72b3c54ed917221db7d068a45581aa400fd48482 Mon Sep 17 00:00:00 2001 From: sechmachine <97589681+sechmachine727@users.noreply.github.com> Date: Wed, 12 Aug 2026 15:43:20 +0700 Subject: [PATCH] fix(core): harden ABI cancellation contracts --- core/include/versevdi_core.h | 29 +++ core/src/abi.rs | 447 +++++++++++++++++++++++++++++------ core/tests/abi_contract.rs | 213 ++++++++++++----- 3 files changed, 561 insertions(+), 128 deletions(-) diff --git a/core/include/versevdi_core.h b/core/include/versevdi_core.h index 9bcb2ae..9affa83 100644 --- a/core/include/versevdi_core.h +++ b/core/include/versevdi_core.h @@ -41,6 +41,7 @@ typedef uint32_t verse_status_t; typedef struct verse_core verse_core_t; +/* data may be NULL only when length is zero. The view never transfers ownership. */ typedef struct verse_bytes_view { const uint8_t *data; size_t length; @@ -87,6 +88,12 @@ typedef struct verse_control_event_v1 { verse_bytes_view_t payload; } verse_control_event_v1_t; +/* + * Signers run synchronously. transcript/tls_message is borrowed only for the call; + * signature_out is exactly 64 writable bytes. Admission may return OK, + * AUTHORITY_REJECTED, CANCELLED, or INTERNAL. TLS may return OK, TLS, CANCELLED, + * or INTERNAL. Any other value is normalized to INTERNAL. + */ typedef verse_status_t (*verse_sign_admission_v1_fn)( void *signer_context, verse_bytes_view_t transcript, @@ -95,6 +102,10 @@ typedef verse_status_t (*verse_sign_tls_ed25519_v1_fn)( void *signer_context, verse_bytes_view_t tls_message, uint8_t signature_out[64]); +/* + * Event callbacks are serialized with one another. Each event and nested byte view + * is borrowed only for its callback and must not be retained. + */ typedef void (*verse_state_event_v1_fn)( void *context, const verse_state_event_v1_t *event); @@ -111,6 +122,10 @@ typedef void (*verse_control_event_v1_fn)( void *context, const verse_control_event_v1_t *event); +/* + * create copies this table. context and every non-NULL callback must remain valid + * until destroy succeeds; BUSY does not end that lifetime. + */ typedef struct verse_core_config_v1 { uint32_t struct_size; uint32_t abi_version; @@ -141,6 +156,13 @@ typedef struct verse_input_event_v1 { } verse_input_event_v1_t; uint32_t verse_core_abi_version(void); +/* + * During a signer callback, every API below except verse_core_abi_version returns + * REENTRANT for every handle. During an event callback, only cancel on that event's + * originating handle is allowed; all other calls and cross-handle cancel return + * REENTRANT. Connect copies both byte inputs before returning, so callers may + * mutate or release their buffers afterward. + */ verse_status_t verse_core_create_v1( const verse_core_config_v1_t *config, verse_core_t **out_core); @@ -151,7 +173,14 @@ verse_status_t verse_core_send_input_v1( verse_core_t *core, const verse_input_event_v1_t *event); verse_status_t verse_core_request_idr_v1(verse_core_t *core); +/* Idempotent and nonblocking; it does not wait for an internal state lock. */ verse_status_t verse_core_cancel_v1(verse_core_t *core); +/* + * OK suppresses all later callbacks, releases session resources, and invalidates + * core. Calls admitted before destruction are accounted safely, but no API call + * may begin after OK. BUSY retains core, context, and callback ownership and + * requires a later retry. + */ verse_status_t verse_core_destroy_v1(verse_core_t *core, uint32_t timeout_ms); #if defined(__APPLE__) && defined(__aarch64__) diff --git a/core/src/abi.rs b/core/src/abi.rs index fce1013..a9f71dc 100644 --- a/core/src/abi.rs +++ b/core/src/abi.rs @@ -6,7 +6,7 @@ use std::panic::{catch_unwind, AssertUnwindSafe}; use std::ptr::{self, NonNull}; use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}; use std::sync::mpsc::{self, Receiver, RecvTimeoutError, Sender, SyncSender, TrySendError}; -use std::sync::{Arc, Mutex, MutexGuard, OnceLock}; +use std::sync::{Arc, Mutex, MutexGuard, OnceLock, TryLockError}; use std::thread::{self, JoinHandle}; use std::time::{Duration, Instant}; @@ -20,6 +20,8 @@ const OK: u32 = 0; const INVALID_ARGUMENT: u32 = 1; const INVALID_STATE: u32 = 2; const UNSUPPORTED_ABI: u32 = 3; +const AUTHORITY_REJECTED: u32 = 4; +const TLS: u32 = 5; const QUEUE_FULL: u32 = 9; const CANCELLED: u32 = 10; const REENTRANT: u32 = 11; @@ -45,9 +47,39 @@ const CONTROL_QUEUE_CAPACITY: usize = 64; static NEXT_ID: AtomicUsize = AtomicUsize::new(1); static HANDLES: OnceLock>>> = OnceLock::new(); +#[cfg(test)] +static CANCEL_ADMISSION_HOOK: OnceLock>> = OnceLock::new(); + +#[cfg(test)] +struct CancelAdmissionHook { + core: usize, + admitted: Sender<()>, + resume: Receiver<()>, +} + thread_local! { - static ACTIVE_CALLBACK: Cell = const { Cell::new(0) }; - static ACTIVE_SIGNER: Cell = const { Cell::new(false) }; + static CALLBACK_MODE: Cell = const { Cell::new(CallbackMode::None) }; +} + +#[derive(Clone, Copy, Eq, PartialEq)] +enum CallbackMode { + None, + Event(usize), + Signer, +} + +struct CallbackModeGuard(CallbackMode); + +impl CallbackModeGuard { + fn enter(mode: CallbackMode) -> Self { + Self(CALLBACK_MODE.with(|current| current.replace(mode))) + } +} + +impl Drop for CallbackModeGuard { + fn drop(&mut self) { + CALLBACK_MODE.with(|current| current.set(self.0)); + } } #[repr(C)] @@ -186,6 +218,7 @@ struct CoreInner { worker_done: Mutex>, worker: Mutex>>, in_flight: AtomicUsize, + cancelled: AtomicBool, freeing: AtomicBool, destroying: AtomicBool, dropped_callbacks: AtomicU64, @@ -216,8 +249,36 @@ fn handles() -> &'static Mutex>> { HANDLES.get_or_init(|| Mutex::new(HashMap::new())) } -fn callback_id() -> usize { - ACTIVE_CALLBACK.with(Cell::get) +fn callback_mode() -> CallbackMode { + CALLBACK_MODE.with(Cell::get) +} + +fn try_lock_until(mutex: &Mutex, deadline: Instant) -> Option> { + loop { + match mutex.try_lock() { + Ok(guard) => return Some(guard), + Err(TryLockError::Poisoned(error)) => return Some(error.into_inner()), + Err(TryLockError::WouldBlock) if Instant::now() >= deadline => return None, + Err(TryLockError::WouldBlock) => thread::yield_now(), + } + } +} + +#[cfg(test)] +fn pause_after_cancel_admission(core: *mut CoreHandle) { + let hook = CANCEL_ADMISSION_HOOK.get_or_init(|| Mutex::new(None)); + let selected = { + let mut hook = lock(hook); + if hook.as_ref().is_some_and(|hook| hook.core == core as usize) { + hook.take() + } else { + None + } + }; + if let Some(hook) = selected { + hook.admitted.send(()).expect("announce cancel admission"); + hook.resume.recv().expect("resume admitted cancel"); + } } fn status(error: CoreError) -> u32 { @@ -236,7 +297,12 @@ fn ffi_boundary_core(core: *mut CoreHandle, action: impl FnOnce() -> u32) -> u32 if let Ok(result) = catch_unwind(AssertUnwindSafe(action)) { result } else { - if let Some(inner) = lock(handles()).get(&(core as usize)).cloned() { + let inner = match handles().try_lock() { + Ok(live) => live.get(&(core as usize)).cloned(), + Err(TryLockError::Poisoned(error)) => error.into_inner().get(&(core as usize)).cloned(), + Err(TryLockError::WouldBlock) => None, + }; + if let Some(inner) = inner { cancel_inner(&inner); } INTERNAL @@ -335,11 +401,10 @@ fn callback_worker( state, reason: 0, }; - ACTIVE_CALLBACK.with(|active| active.set(id)); + let _mode = CallbackModeGuard::enter(CallbackMode::Event(id)); // Invariant: copied callback/context remain caller-owned and valid until destroy succeeds; // event is readable for this synchronous callback only. unsafe { callback(callbacks.context as *mut c_void, &raw const event) }; - ACTIVE_CALLBACK.with(|active| active.set(0)); } } Err(RecvTimeoutError::Timeout) => {} @@ -349,29 +414,47 @@ fn callback_worker( let _ = done.send(()); } -fn cancel_inner(inner: &CoreInner) { - let mut state = lock(&inner.state); - let already_cancelled = state.session.is_cancelled(); +fn apply_cancellation(state: &mut SessionState) { state.session.cancel(); state.lifecycle = Lifecycle::Cancelled; - if !already_cancelled { +} + +fn cancel_inner(inner: &CoreInner) { + let first = !inner.cancelled.swap(true, Ordering::AcqRel); + match inner.state.try_lock() { + Ok(mut state) => apply_cancellation(&mut state), + Err(TryLockError::Poisoned(error)) => apply_cancellation(&mut error.into_inner()), + Err(TryLockError::WouldBlock) => {} + } + if first { enqueue_callback_under_state_lock(inner, CallbackEvent::State(STATE_CANCELLED)); } } -fn call_signer(id: usize, callback: SignFn, context: usize, bytes: &[u8]) -> u32 { +#[derive(Clone, Copy)] +enum SignerPurpose { + Admission, + Tls, +} + +fn normalize_signer_status(purpose: SignerPurpose, result: u32) -> u32 { + match (purpose, result) { + (SignerPurpose::Admission, OK | AUTHORITY_REJECTED | CANCELLED | INTERNAL) + | (SignerPurpose::Tls, OK | TLS | CANCELLED | INTERNAL) => result, + _ => INTERNAL, + } +} + +fn call_signer(purpose: SignerPurpose, callback: SignFn, context: usize, bytes: &[u8]) -> u32 { let mut signature = [0_u8; 64]; let view = BytesView { data: bytes.as_ptr(), length: bytes.len(), }; - ACTIVE_CALLBACK.with(|active| active.set(id)); - ACTIVE_SIGNER.with(|active| active.set(true)); + let _mode = CallbackModeGuard::enter(CallbackMode::Signer); // Invariant: bytes and the writable 64-byte signature buffer live for the synchronous callback. let result = unsafe { callback(context as *mut c_void, view, signature.as_mut_ptr()) }; - ACTIVE_SIGNER.with(|active| active.set(false)); - ACTIVE_CALLBACK.with(|active| active.set(0)); - result + normalize_signer_status(purpose, result) } fn convert_input(event: &AbiInputEvent) -> Result { @@ -443,6 +526,9 @@ unsafe extern "C" fn verse_core_create_v1( out_core: *mut *mut CoreHandle, ) -> u32 { ffi_boundary(|| { + if callback_mode() != CallbackMode::None { + return REENTRANT; + } if out_core.is_null() { return INVALID_ARGUMENT; } @@ -506,6 +592,7 @@ unsafe extern "C" fn verse_core_create_v1( worker_done: Mutex::new(done_rx), worker: Mutex::new(Some(worker)), in_flight: AtomicUsize::new(0), + cancelled: AtomicBool::new(false), freeing: AtomicBool::new(false), destroying: AtomicBool::new(false), dropped_callbacks: AtomicU64::new(0), @@ -525,14 +612,14 @@ unsafe extern "C" fn verse_core_connect_v1( request: *const ConnectRequest, ) -> u32 { ffi_boundary_core(core, || { + if callback_mode() != CallbackMode::None { + return REENTRANT; + } // Invariant: caller owns a live opaque handle for the complete call. let (handle, _guard) = match unsafe { begin_call(core) } { Ok(value) => value, Err(error) => return error, }; - if callback_id() == handle.id { - return REENTRANT; - } // Invariant: request points to a readable sized/versioned table. let request = match unsafe { read_table(request) } { Ok(request) => request, @@ -558,7 +645,8 @@ unsafe extern "C" fn verse_core_connect_v1( } { let mut state = lock(&handle.state); - if state.session.is_cancelled() { + if handle.cancelled.load(Ordering::Acquire) { + apply_cancellation(&mut state); return CANCELLED; } if state.lifecycle != Lifecycle::Created { @@ -568,7 +656,7 @@ unsafe extern "C" fn verse_core_connect_v1( enqueue_callback_under_state_lock(&handle, CallbackEvent::State(STATE_CONNECTING)); } let admission = call_signer( - handle.id, + SignerPurpose::Admission, handle.callbacks.sign_admission, handle.callbacks.context, &manifest, @@ -577,11 +665,11 @@ unsafe extern "C" fn verse_core_connect_v1( cancel_inner(&handle); return admission; } - if lock(&handle.state).session.is_cancelled() { + if handle.cancelled.load(Ordering::Acquire) { return CANCELLED; } let tls = call_signer( - handle.id, + SignerPurpose::Tls, handle.callbacks.sign_tls_ed25519, handle.callbacks.context, &credential, @@ -592,8 +680,9 @@ unsafe extern "C" fn verse_core_connect_v1( } { let mut state = lock(&handle.state); - if state.session.is_cancelled() || state.lifecycle == Lifecycle::Destroying { - state.lifecycle = Lifecycle::Cancelled; + if handle.cancelled.load(Ordering::Acquire) || state.lifecycle == Lifecycle::Destroying + { + apply_cancellation(&mut state); return CANCELLED; } state.lifecycle = Lifecycle::Connected; @@ -609,21 +698,22 @@ unsafe extern "C" fn verse_core_send_input_v1( event: *const AbiInputEvent, ) -> u32 { ffi_boundary_core(core, || { + if callback_mode() != CallbackMode::None { + return REENTRANT; + } // Invariant: caller owns a live opaque handle for the complete call. let (handle, _guard) = match unsafe { begin_call(core) } { Ok(value) => value, Err(error) => return error, }; - if callback_id() == handle.id { - return REENTRANT; - } // Invariant: event points to a readable sized/versioned table. let event = match unsafe { read_table(event) } { Ok(event) => event, Err(error) => return error, }; let mut state = lock(&handle.state); - if state.session.is_cancelled() { + if handle.cancelled.load(Ordering::Acquire) { + apply_cancellation(&mut state); return CANCELLED; } if state.lifecycle != Lifecycle::Connected { @@ -648,16 +738,17 @@ unsafe extern "C" fn verse_core_send_input_v1( #[no_mangle] unsafe extern "C" fn verse_core_request_idr_v1(core: *mut CoreHandle) -> u32 { ffi_boundary_core(core, || { + if callback_mode() != CallbackMode::None { + return REENTRANT; + } // Invariant: caller owns a live opaque handle for the complete call. let (handle, _guard) = match unsafe { begin_call(core) } { Ok(value) => value, Err(error) => return error, }; - if callback_id() == handle.id { - return REENTRANT; - } let mut state = lock(&handle.state); - if state.session.is_cancelled() { + if handle.cancelled.load(Ordering::Acquire) { + apply_cancellation(&mut state); return CANCELLED; } if state.lifecycle != Lifecycle::Connected { @@ -674,74 +765,138 @@ unsafe extern "C" fn verse_core_request_idr_v1(core: *mut CoreHandle) -> u32 { #[no_mangle] unsafe extern "C" fn verse_core_cancel_v1(core: *mut CoreHandle) -> u32 { ffi_boundary_core(core, || { + if callback_mode() == CallbackMode::Signer { + return REENTRANT; + } // Invariant: caller owns a live opaque handle for the complete call. let (handle, _guard) = match unsafe { begin_call(core) } { Ok(value) => value, Err(error) => return error, }; - if callback_id() == handle.id && ACTIVE_SIGNER.with(Cell::get) { - return REENTRANT; + if let CallbackMode::Event(origin) = callback_mode() { + if origin != handle.id { + return REENTRANT; + } } + #[cfg(test)] + pause_after_cancel_admission(core); cancel_inner(&handle); OK }) } +fn destroy_busy(inner: &CoreInner) -> u32 { + inner.destroying.store(false, Ordering::Release); + BUSY +} + +fn wait_for_destroy_quiescence(inner: &CoreInner, deadline: Instant) -> Result<(), u32> { + while inner.in_flight.load(Ordering::Acquire) != 0 { + if Instant::now() >= deadline { + return Err(destroy_busy(inner)); + } + thread::yield_now(); + } + + let worker_present = { + if Instant::now() > deadline { + return Err(destroy_busy(inner)); + } + let Some(worker) = try_lock_until(&inner.worker, deadline) else { + return Err(destroy_busy(inner)); + }; + worker.is_some() + }; + if worker_present { + if Instant::now() > deadline { + return Err(destroy_busy(inner)); + } + let Some(done) = try_lock_until(&inner.worker_done, deadline) else { + return Err(destroy_busy(inner)); + }; + if matches!( + done.recv_timeout(deadline.saturating_duration_since(Instant::now())), + Err(RecvTimeoutError::Timeout) + ) { + return Err(destroy_busy(inner)); + } + } + loop { + if Instant::now() > deadline { + return Err(destroy_busy(inner)); + } + let Some(worker) = try_lock_until(&inner.worker, deadline) else { + return Err(destroy_busy(inner)); + }; + if worker.as_ref().is_none_or(JoinHandle::is_finished) { + break; + } + drop(worker); + thread::yield_now(); + } + let worker = { + if Instant::now() > deadline { + return Err(destroy_busy(inner)); + } + let Some(mut worker) = try_lock_until(&inner.worker, deadline) else { + return Err(destroy_busy(inner)); + }; + worker.take() + }; + if worker.is_some_and(|worker| worker.join().is_err()) { + inner.destroying.store(false, Ordering::Release); + return Err(INTERNAL); + } + Ok(()) +} + #[no_mangle] unsafe extern "C" fn verse_core_destroy_v1(core: *mut CoreHandle, timeout_ms: u32) -> u32 { ffi_boundary_core(core, || { + if callback_mode() != CallbackMode::None { + return REENTRANT; + } + let deadline = Instant::now() + Duration::from_millis(u64::from(timeout_ms)); let Some(core) = NonNull::new(core) else { return INVALID_ARGUMENT; }; let inner = { - let live = lock(handles()); + let Some(live) = try_lock_until(handles(), deadline) else { + return BUSY; + }; let Some(inner) = live.get(&(core.as_ptr() as usize)).cloned() else { return INVALID_ARGUMENT; }; inner }; - if callback_id() == inner.id { - return REENTRANT; - } if inner.destroying.swap(true, Ordering::AcqRel) { return BUSY; } - let deadline = Instant::now() + Duration::from_millis(u64::from(timeout_ms)); inner.callbacks_closed.store(true, Ordering::Release); + inner.cancelled.store(true, Ordering::Release); { - let mut state = lock(&inner.state); + if Instant::now() > deadline { + return destroy_busy(&inner); + } + let Some(mut state) = try_lock_until(&inner.state, deadline) else { + return destroy_busy(&inner); + }; state.lifecycle = Lifecycle::Destroying; state.session.cancel(); } let _ = inner.stop_tx.send(()); - - while inner.in_flight.load(Ordering::Acquire) != 0 { - if Instant::now() >= deadline { - inner.destroying.store(false, Ordering::Release); - return BUSY; - } - thread::yield_now(); - } - - let worker_finished = lock(&inner.worker).is_none() - || lock(&inner.worker_done) - .recv_timeout(deadline.saturating_duration_since(Instant::now())) - .is_ok(); - if !worker_finished { - inner.destroying.store(false, Ordering::Release); - return BUSY; - } - if let Some(worker) = lock(&inner.worker).take() { - if worker.join().is_err() { - inner.destroying.store(false, Ordering::Release); - return INTERNAL; - } + if let Err(error) = wait_for_destroy_quiescence(&inner, deadline) { + return error; } { - let mut live = lock(handles()); + if Instant::now() > deadline { + return destroy_busy(&inner); + } + let Some(mut live) = try_lock_until(handles(), deadline) else { + return destroy_busy(&inner); + }; if inner.in_flight.load(Ordering::Acquire) != 0 { - inner.destroying.store(false, Ordering::Release); - return BUSY; + return destroy_busy(&inner); } inner.freeing.store(true, Ordering::Release); live.remove(&(core.as_ptr() as usize)); @@ -777,10 +932,164 @@ const _: fn(Callbacks) = |callbacks| { #[cfg(test)] mod tests { - use super::{ffi_boundary, INTERNAL}; + use std::ffi::c_void; + use std::mem::size_of; + use std::ptr; + use std::sync::mpsc; + use std::sync::Arc; + use std::thread; + use std::time::Duration; + + use super::{ + ffi_boundary, handles, lock, BytesView, CancelAdmissionHook, Config, CoreHandle, SignFn, + BUSY, CANCELLED, CANCEL_ADMISSION_HOOK, INTERNAL, OK, + }; + + unsafe extern "C" fn sign(_context: *mut c_void, _input: BytesView, signature: *mut u8) -> u32 { + // Test invariant: create supplies this callback only to the ABI, which provides 64 bytes. + unsafe { ptr::write_bytes(signature, 0, 64) }; + OK + } + + fn create_core() -> *mut CoreHandle { + let config = Config { + struct_size: u32::try_from(size_of::()).expect("config size"), + abi_version: 1, + context: ptr::null_mut(), + sign_admission: Some(sign as SignFn), + sign_tls_ed25519: Some(sign as SignFn), + on_state: None, + on_error: None, + on_stats: None, + on_media: None, + on_control: None, + }; + let mut core = ptr::null_mut(); + // Test invariant: config and out pointer are valid for the synchronous call. + assert_eq!( + unsafe { super::verse_core_create_v1(&raw const config, &raw mut core) }, + OK + ); + core + } + + fn inner(core: *mut CoreHandle) -> Arc { + lock(handles()) + .get(&(core as usize)) + .expect("live core") + .clone() + } + + fn hold_state(core: *mut CoreHandle) -> (mpsc::Sender<()>, thread::JoinHandle<()>) { + let inner = inner(core); + let (locked_tx, locked_rx) = mpsc::channel(); + let (release_tx, release_rx) = mpsc::channel(); + let holder = thread::spawn(move || { + let _state = lock(&inner.state); + locked_tx.send(()).expect("announce state lock"); + release_rx.recv().expect("release state lock"); + }); + locked_rx.recv().expect("state lock acquired"); + (release_tx, holder) + } #[test] fn panic_boundary_maps_to_internal_status() { assert_eq!(ffi_boundary(|| panic!("contained test panic")), INTERNAL); } + + #[test] + fn cancel_does_not_wait_for_contended_state_lock() { + let core = create_core(); + let (release, holder) = hold_state(core); + let (result_tx, result_rx) = mpsc::channel(); + let address = core as usize; + let caller = thread::spawn(move || { + result_tx + .send(unsafe { super::verse_core_cancel_v1(address as *mut CoreHandle) }) + .expect("send cancel result"); + }); + + let timely = result_rx.recv_timeout(Duration::from_millis(100)); + let returned_before_release = timely.is_ok(); + release.send(()).expect("release state"); + holder.join().expect("state holder"); + let result = timely.unwrap_or_else(|_| { + result_rx + .recv_timeout(Duration::from_secs(1)) + .expect("cancel eventually returns") + }); + caller.join().expect("cancel caller"); + assert!(returned_before_release, "cancel waited for the state mutex"); + assert_eq!(result, OK, "cancel waited for the state mutex"); + assert_eq!( + unsafe { super::verse_core_request_idr_v1(core) }, + CANCELLED, + "the admitted cancellation must be observed by state operations" + ); + assert_eq!(unsafe { super::verse_core_destroy_v1(core, 2_000) }, OK); + } + + #[test] + fn zero_and_short_timeout_destroy_do_not_wait_for_contended_state_lock() { + for timeout_ms in [0, 5] { + let core = create_core(); + let inner = inner(core); + let (release, holder) = hold_state(core); + let (result_tx, result_rx) = mpsc::channel(); + let address = core as usize; + let caller = thread::spawn(move || { + result_tx + .send(unsafe { + super::verse_core_destroy_v1(address as *mut CoreHandle, timeout_ms) + }) + .expect("send destroy result"); + }); + + let timely = result_rx.recv_timeout(Duration::from_millis(100)); + let returned_before_release = timely.is_ok(); + release.send(()).expect("release state"); + holder.join().expect("state holder"); + let result = timely.unwrap_or_else(|_| { + result_rx + .recv_timeout(Duration::from_secs(1)) + .expect("destroy eventually returns") + }); + caller.join().expect("destroy caller"); + assert!( + returned_before_release, + "destroy({timeout_ms}) waited for the state mutex" + ); + assert_eq!(result, BUSY); + assert!(inner + .callbacks_closed + .load(std::sync::atomic::Ordering::Acquire)); + assert!(lock(handles()).contains_key(&(core as usize))); + assert_eq!(unsafe { super::verse_core_destroy_v1(core, 2_000) }, OK); + } + } + + #[test] + fn destroy_accounts_for_cancel_admitted_before_destruction() { + let core = create_core(); + let (admitted_tx, admitted_rx) = mpsc::channel(); + let (resume_tx, resume_rx) = mpsc::channel(); + *lock(CANCEL_ADMISSION_HOOK.get_or_init(|| std::sync::Mutex::new(None))) = + Some(CancelAdmissionHook { + core: core as usize, + admitted: admitted_tx, + resume: resume_rx, + }); + let address = core as usize; + let cancel = thread::spawn(move || unsafe { + super::verse_core_cancel_v1(address as *mut CoreHandle) + }); + admitted_rx.recv().expect("cancel admitted"); + + assert_eq!(unsafe { super::verse_core_destroy_v1(core, 0) }, BUSY); + assert!(lock(handles()).contains_key(&(core as usize))); + resume_tx.send(()).expect("resume cancel"); + assert_eq!(cancel.join().expect("cancel caller"), OK); + assert_eq!(unsafe { super::verse_core_destroy_v1(core, 2_000) }, OK); + } } diff --git a/core/tests/abi_contract.rs b/core/tests/abi_contract.rs index 523fe8a..cfd13b0 100644 --- a/core/tests/abi_contract.rs +++ b/core/tests/abi_contract.rs @@ -9,7 +9,7 @@ use std::ffi::c_void; use std::mem::{offset_of, size_of}; use std::ptr; use std::sync::atomic::{AtomicBool, AtomicU32, AtomicUsize, Ordering}; -use std::sync::{Arc, Barrier, Condvar, Mutex}; +use std::sync::{Condvar, Mutex}; use std::thread; use std::time::{Duration, Instant}; @@ -18,10 +18,13 @@ const OK: u32 = 0; const INVALID_ARGUMENT: u32 = 1; const INVALID_STATE: u32 = 2; const UNSUPPORTED_ABI: u32 = 3; +const AUTHORITY_REJECTED: u32 = 4; +const TLS: u32 = 5; const QUEUE_FULL: u32 = 9; const CANCELLED: u32 = 10; const REENTRANT: u32 = 11; const BUSY: u32 = 12; +const INTERNAL: u32 = 13; const STATE_CONNECTING: u32 = 1; const STATE_CONNECTED: u32 = 2; @@ -148,6 +151,8 @@ struct Context { core: AtomicUsize, admission_calls: AtomicUsize, tls_calls: AtomicUsize, + admission_status: AtomicU32, + tls_status: AtomicU32, admission_input: Mutex>, tls_input: Mutex>, states: Mutex>, @@ -160,6 +165,8 @@ struct Context { callback_active: AtomicUsize, callback_max: AtomicUsize, cancel_on_connecting: AtomicBool, + reentry_target: AtomicUsize, + reentry_results: Mutex>, } impl Default for Context { @@ -168,6 +175,8 @@ impl Default for Context { core: AtomicUsize::new(0), admission_calls: AtomicUsize::new(0), tls_calls: AtomicUsize::new(0), + admission_status: AtomicU32::new(OK), + tls_status: AtomicU32::new(OK), admission_input: Mutex::new(Vec::new()), tls_input: Mutex::new(Vec::new()), states: Mutex::new(Vec::new()), @@ -180,6 +189,8 @@ impl Default for Context { callback_active: AtomicUsize::new(0), callback_max: AtomicUsize::new(0), cancel_on_connecting: AtomicBool::new(false), + reentry_target: AtomicUsize::new(0), + reentry_results: Mutex::new(Vec::new()), } } } @@ -197,7 +208,7 @@ unsafe extern "C" fn sign_admission(raw: *mut c_void, input: BytesView, output: *ctx.admission_input.lock().expect("admission lock") = bytes.to_vec(); // Test invariant: ABI promises a writable 64-byte Rust-owned signature buffer. unsafe { ptr::write_bytes(output, 0xA5, 64) }; - OK + ctx.admission_status.load(Ordering::SeqCst) } unsafe extern "C" fn sign_tls(raw: *mut c_void, input: BytesView, output: *mut u8) -> u32 { @@ -208,7 +219,57 @@ unsafe extern "C" fn sign_tls(raw: *mut c_void, input: BytesView, output: *mut u *ctx.tls_input.lock().expect("tls lock") = bytes.to_vec(); // Test invariant: ABI promises a writable 64-byte Rust-owned signature buffer. unsafe { ptr::write_bytes(output, 0x5A, 64) }; - OK + ctx.tls_status.load(Ordering::SeqCst) +} + +unsafe extern "C" fn sign_admission_probes_global_reentry( + raw: *mut c_void, + input: BytesView, + output: *mut u8, +) -> u32 { + let ctx = unsafe { context(raw) }; + let target = ctx.reentry_target.load(Ordering::SeqCst) as *mut Core; + let results = [ + unsafe { verse_core_abi_version() }, + unsafe { verse_core_create_v1(ptr::null(), ptr::null_mut()) }, + unsafe { verse_core_connect_v1(target, ptr::null()) }, + unsafe { verse_core_send_input_v1(target, ptr::null()) }, + unsafe { verse_core_request_idr_v1(target) }, + unsafe { verse_core_cancel_v1(target) }, + unsafe { verse_core_destroy_v1(target, 0) }, + ]; + ctx.reentry_results + .lock() + .expect("signer reentry results") + .extend(results); + unsafe { sign_admission(raw, input, output) } +} + +unsafe extern "C" fn on_state_probes_global_reentry(raw: *mut c_void, event: *const StateEvent) { + let ctx = unsafe { context(raw) }; + // Test invariant: ABI promises a readable state record for the callback duration. + let state = unsafe { (*event).state }; + ctx.states.lock().expect("states lock").push(state); + ctx.wake.notify_all(); + if state != STATE_CONNECTED { + return; + } + let origin = ctx.core.load(Ordering::SeqCst) as *mut Core; + let other = ctx.reentry_target.load(Ordering::SeqCst) as *mut Core; + let results = [ + unsafe { verse_core_abi_version() }, + unsafe { verse_core_create_v1(ptr::null(), ptr::null_mut()) }, + unsafe { verse_core_connect_v1(origin, ptr::null()) }, + unsafe { verse_core_send_input_v1(origin, ptr::null()) }, + unsafe { verse_core_request_idr_v1(origin) }, + unsafe { verse_core_cancel_v1(other) }, + unsafe { verse_core_destroy_v1(origin, 0) }, + unsafe { verse_core_cancel_v1(origin) }, + ]; + ctx.reentry_results + .lock() + .expect("event reentry results") + .extend(results); } unsafe extern "C" fn on_state(raw: *mut c_void, event: *const StateEvent) { @@ -561,6 +622,96 @@ fn signer_callbacks_cannot_reenter_even_cancel() { assert_eq!(destroy(core), OK); } +#[test] +fn signer_callback_rejects_every_stateful_api_across_handles() { + let mut other_ctx = Context::default(); + other_ctx.reentry_cancel.store(OK, Ordering::SeqCst); + let other = create(&mut other_ctx); + + let mut ctx = Context::default(); + ctx.reentry_target.store(other as usize, Ordering::SeqCst); + let mut cfg = config(&mut ctx); + cfg.sign_admission = Some(sign_admission_probes_global_reentry); + let mut core = ptr::null_mut(); + assert_eq!(unsafe { verse_core_create_v1(&cfg, &mut core) }, OK); + ctx.core.store(core as usize, Ordering::SeqCst); + assert_eq!(connect(core, MANIFEST, CREDENTIAL), OK); + assert_eq!( + ctx.reentry_results + .lock() + .expect("signer results") + .as_slice(), + [ABI_V1, REENTRANT, REENTRANT, REENTRANT, REENTRANT, REENTRANT, REENTRANT] + ); + + assert_eq!(connect(other, MANIFEST, CREDENTIAL), OK); + assert_eq!(destroy(core), OK); + assert_eq!(destroy(other), OK); +} + +#[test] +fn event_callback_allows_only_originating_handle_cancel() { + let mut other_ctx = Context::default(); + other_ctx.reentry_cancel.store(OK, Ordering::SeqCst); + let other = create(&mut other_ctx); + + let mut ctx = Context::default(); + ctx.reentry_target.store(other as usize, Ordering::SeqCst); + let mut cfg = config(&mut ctx); + cfg.on_state = Some(on_state_probes_global_reentry); + let mut core = ptr::null_mut(); + assert_eq!(unsafe { verse_core_create_v1(&cfg, &mut core) }, OK); + ctx.core.store(core as usize, Ordering::SeqCst); + assert_eq!(connect(core, MANIFEST, CREDENTIAL), OK); + wait_for(&ctx, |states| states.contains(&STATE_CANCELLED)); + assert_eq!( + ctx.reentry_results + .lock() + .expect("event results") + .as_slice(), + [ABI_V1, REENTRANT, REENTRANT, REENTRANT, REENTRANT, REENTRANT, REENTRANT, OK] + ); + + assert_eq!(connect(other, MANIFEST, CREDENTIAL), OK); + assert_eq!(destroy(core), OK); + assert_eq!(destroy(other), OK); +} + +#[test] +fn signer_statuses_are_purpose_specific_and_unknown_values_are_internal() { + for (admission, expected) in [ + (OK, OK), + (AUTHORITY_REJECTED, AUTHORITY_REJECTED), + (CANCELLED, CANCELLED), + (INTERNAL, INTERNAL), + (TLS, INTERNAL), + (u32::MAX, INTERNAL), + ] { + let mut ctx = Context::default(); + ctx.reentry_cancel.store(OK, Ordering::SeqCst); + ctx.admission_status.store(admission, Ordering::SeqCst); + let core = create(&mut ctx); + assert_eq!(connect(core, MANIFEST, CREDENTIAL), expected); + assert_eq!(destroy(core), OK); + } + + for (tls, expected) in [ + (OK, OK), + (TLS, TLS), + (CANCELLED, CANCELLED), + (INTERNAL, INTERNAL), + (AUTHORITY_REJECTED, INTERNAL), + (u32::MAX, INTERNAL), + ] { + let mut ctx = Context::default(); + ctx.reentry_cancel.store(OK, Ordering::SeqCst); + ctx.tls_status.store(tls, Ordering::SeqCst); + let core = create(&mut ctx); + assert_eq!(connect(core, MANIFEST, CREDENTIAL), expected); + assert_eq!(destroy(core), OK); + } +} + #[test] fn cancel_during_connect_preserves_state_order_and_stops_before_tls_signing() { let mut ctx = Context::default(); @@ -616,65 +767,9 @@ fn destroy_timeout_keeps_ownership_suppresses_late_callbacks_and_allows_retry() ctx.release_callbacks.store(true, Ordering::SeqCst); ctx.wake.notify_all(); assert_eq!(destroy(core), OK); - thread::sleep(Duration::from_millis(20)); assert_eq!(ctx.states.lock().expect("states").len(), count_at_timeout); } -#[test] -fn concurrent_cancel_and_destroy_do_not_race_lifetime() { - let mut ctx = Context::default(); - ctx.reentry_cancel.store(OK, Ordering::SeqCst); - let core = create(&mut ctx); - assert_eq!(connect(core, MANIFEST, CREDENTIAL), OK); - let barrier = Arc::new(Barrier::new(3)); - let core_address = core as usize; - let cancel_barrier = Arc::clone(&barrier); - let cancel = thread::spawn(move || { - cancel_barrier.wait(); - // Test invariant: the main thread retains ownership across this overlapping call. - unsafe { verse_core_cancel_v1(core_address as *mut Core) } - }); - let destroy_barrier = Arc::clone(&barrier); - let destroyer = thread::spawn(move || { - destroy_barrier.wait(); - // Test invariant: this is the sole successful destroy attempt in the race. - unsafe { verse_core_destroy_v1(core_address as *mut Core, 2_000) } - }); - barrier.wait(); - let cancel_status = cancel.join().expect("cancel thread"); - let destroy_status = destroyer.join().expect("destroy thread"); - assert!(matches!(cancel_status, OK | INVALID_ARGUMENT)); - assert_eq!(destroy_status, OK); -} - -#[test] -fn repeated_concurrent_cancel_destroy_cycles_are_safe() { - for _ in 0..100 { - let mut ctx = Context::default(); - ctx.reentry_cancel.store(OK, Ordering::SeqCst); - let core = create(&mut ctx); - assert_eq!(connect(core, MANIFEST, CREDENTIAL), OK); - let barrier = Arc::new(Barrier::new(3)); - let address = core as usize; - let cancel_barrier = Arc::clone(&barrier); - let cancel = thread::spawn(move || { - cancel_barrier.wait(); - unsafe { verse_core_cancel_v1(address as *mut Core) } - }); - let destroy_barrier = Arc::clone(&barrier); - let destroyer = thread::spawn(move || { - destroy_barrier.wait(); - unsafe { verse_core_destroy_v1(address as *mut Core, 2_000) } - }); - barrier.wait(); - assert!(matches!( - cancel.join().expect("cancel"), - OK | INVALID_ARGUMENT - )); - assert_eq!(destroyer.join().expect("destroy"), OK); - } -} - #[test] fn preconnect_and_postcancel_state_checks_are_stable() { let mut ctx = Context::default();