fix(core): harden ABI cancellation contracts

This commit is contained in:
sechmachine
2026-08-12 15:43:20 +07:00
parent 67510b65b4
commit 72b3c54ed9
3 changed files with 561 additions and 128 deletions
+29
View File
@@ -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__)
+377 -68
View File
@@ -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<Mutex<HashMap<usize, Arc<CoreInner>>>> = OnceLock::new();
#[cfg(test)]
static CANCEL_ADMISSION_HOOK: OnceLock<Mutex<Option<CancelAdmissionHook>>> = OnceLock::new();
#[cfg(test)]
struct CancelAdmissionHook {
core: usize,
admitted: Sender<()>,
resume: Receiver<()>,
}
thread_local! {
static ACTIVE_CALLBACK: Cell<usize> = const { Cell::new(0) };
static ACTIVE_SIGNER: Cell<bool> = const { Cell::new(false) };
static CALLBACK_MODE: Cell<CallbackMode> = 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<Receiver<()>>,
worker: Mutex<Option<JoinHandle<()>>>,
in_flight: AtomicUsize,
cancelled: AtomicBool,
freeing: AtomicBool,
destroying: AtomicBool,
dropped_callbacks: AtomicU64,
@@ -216,8 +249,36 @@ fn handles() -> &'static Mutex<HashMap<usize, Arc<CoreInner>>> {
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<T>(mutex: &Mutex<T>, deadline: Instant) -> Option<MutexGuard<'_, T>> {
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<InputEvent, u32> {
@@ -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) {
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::<Config>()).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<super::CoreInner> {
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);
}
}
+154 -59
View File
@@ -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<Vec<u8>>,
tls_input: Mutex<Vec<u8>>,
states: Mutex<Vec<u32>>,
@@ -160,6 +165,8 @@ struct Context {
callback_active: AtomicUsize,
callback_max: AtomicUsize,
cancel_on_connecting: AtomicBool,
reentry_target: AtomicUsize,
reentry_results: Mutex<Vec<u32>>,
}
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();