diff --git a/core/include/versevdi_core.h b/core/include/versevdi_core.h new file mode 100644 index 0000000..9bcb2ae --- /dev/null +++ b/core/include/versevdi_core.h @@ -0,0 +1,175 @@ +#ifndef VERSEVDI_CORE_H +#define VERSEVDI_CORE_H + +#include +#include + +#ifdef __cplusplus +extern "C" { +#endif + +#define VERSE_CORE_ABI_VERSION_1 UINT32_C(1) + +typedef uint32_t verse_status_t; + +#define VERSE_STATUS_OK UINT32_C(0) +#define VERSE_STATUS_INVALID_ARGUMENT UINT32_C(1) +#define VERSE_STATUS_INVALID_STATE UINT32_C(2) +#define VERSE_STATUS_UNSUPPORTED_ABI UINT32_C(3) +#define VERSE_STATUS_AUTHORITY_REJECTED UINT32_C(4) +#define VERSE_STATUS_TLS UINT32_C(5) +#define VERSE_STATUS_TRANSPORT UINT32_C(6) +#define VERSE_STATUS_PROTOCOL UINT32_C(7) +#define VERSE_STATUS_EXPIRED UINT32_C(8) +#define VERSE_STATUS_QUEUE_FULL UINT32_C(9) +#define VERSE_STATUS_CANCELLED UINT32_C(10) +#define VERSE_STATUS_REENTRANT UINT32_C(11) +#define VERSE_STATUS_BUSY UINT32_C(12) +#define VERSE_STATUS_INTERNAL UINT32_C(13) + +#define VERSE_STATE_CONNECTING UINT32_C(1) +#define VERSE_STATE_CONNECTED UINT32_C(2) +#define VERSE_STATE_CANCELLED UINT32_C(3) + +#define VERSE_INPUT_KEYBOARD UINT32_C(1) +#define VERSE_INPUT_MOUSE_BUTTON UINT32_C(2) +#define VERSE_INPUT_RELATIVE_MOUSE UINT32_C(3) +#define VERSE_INPUT_TEXT UINT32_C(4) +#define VERSE_INPUT_CONTROLLER UINT32_C(5) +#define VERSE_INPUT_ABSOLUTE_MOUSE UINT32_C(6) +#define VERSE_INPUT_SCROLL UINT32_C(7) + +typedef struct verse_core verse_core_t; + +typedef struct verse_bytes_view { + const uint8_t *data; + size_t length; +} verse_bytes_view_t; + +typedef struct verse_state_event_v1 { + uint32_t struct_size; + uint32_t abi_version; + uint32_t state; + uint32_t reason; +} verse_state_event_v1_t; + +typedef struct verse_error_event_v1 { + uint32_t struct_size; + uint32_t abi_version; + uint32_t code; + uint32_t retryable; + uint32_t phase; + uint32_t reserved; +} verse_error_event_v1_t; + +typedef struct verse_stats_event_v1 { + uint32_t struct_size; + uint32_t abi_version; + uint64_t dropped_callbacks; + uint64_t dropped_media_units; + uint64_t dropped_input_events; +} verse_stats_event_v1_t; + +typedef struct verse_media_event_v1 { + uint32_t struct_size; + uint32_t abi_version; + uint32_t channel; + uint32_t sequence; + uint64_t timestamp_ms; + verse_bytes_view_t encoded_unit; +} verse_media_event_v1_t; + +typedef struct verse_control_event_v1 { + uint32_t struct_size; + uint32_t abi_version; + uint32_t kind; + uint32_t reserved; + verse_bytes_view_t payload; +} verse_control_event_v1_t; + +typedef verse_status_t (*verse_sign_admission_v1_fn)( + void *signer_context, + verse_bytes_view_t transcript, + uint8_t signature_out[64]); +typedef verse_status_t (*verse_sign_tls_ed25519_v1_fn)( + void *signer_context, + verse_bytes_view_t tls_message, + uint8_t signature_out[64]); +typedef void (*verse_state_event_v1_fn)( + void *context, + const verse_state_event_v1_t *event); +typedef void (*verse_error_event_v1_fn)( + void *context, + const verse_error_event_v1_t *event); +typedef void (*verse_stats_event_v1_fn)( + void *context, + const verse_stats_event_v1_t *event); +typedef void (*verse_media_event_v1_fn)( + void *context, + const verse_media_event_v1_t *event); +typedef void (*verse_control_event_v1_fn)( + void *context, + const verse_control_event_v1_t *event); + +typedef struct verse_core_config_v1 { + uint32_t struct_size; + uint32_t abi_version; + void *context; + verse_sign_admission_v1_fn sign_admission; + verse_sign_tls_ed25519_v1_fn sign_tls_ed25519; + verse_state_event_v1_fn on_state; + verse_error_event_v1_fn on_error; + verse_stats_event_v1_fn on_stats; + verse_media_event_v1_fn on_media; + verse_control_event_v1_fn on_control; +} verse_core_config_v1_t; + +typedef struct verse_connect_request_v1 { + uint32_t struct_size; + uint32_t abi_version; + verse_bytes_view_t manifest_json; + verse_bytes_view_t tunnel_credential_json; +} verse_connect_request_v1_t; + +/* values are kind-specific signed fields; every unused field and flags must be zero. */ +typedef struct verse_input_event_v1 { + uint32_t struct_size; + uint32_t abi_version; + uint32_t kind; + uint32_t flags; + int32_t values[12]; +} verse_input_event_v1_t; + +uint32_t verse_core_abi_version(void); +verse_status_t verse_core_create_v1( + const verse_core_config_v1_t *config, + verse_core_t **out_core); +verse_status_t verse_core_connect_v1( + verse_core_t *core, + const verse_connect_request_v1_t *request); +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); +verse_status_t verse_core_cancel_v1(verse_core_t *core); +verse_status_t verse_core_destroy_v1(verse_core_t *core, uint32_t timeout_ms); + +#if defined(__APPLE__) && defined(__aarch64__) +_Static_assert(sizeof(verse_bytes_view_t) == 16, "verse_bytes_view_t arm64 layout"); +_Static_assert(sizeof(verse_core_config_v1_t) == 72, "verse_core_config_v1_t arm64 layout"); +_Static_assert(offsetof(verse_core_config_v1_t, sign_admission) == 16, "config callback offset"); +_Static_assert(sizeof(verse_connect_request_v1_t) == 40, "verse_connect_request_v1_t arm64 layout"); +_Static_assert(offsetof(verse_connect_request_v1_t, manifest_json) == 8, "request view offset"); +_Static_assert(sizeof(verse_input_event_v1_t) == 64, "verse_input_event_v1_t arm64 layout"); +_Static_assert(sizeof(verse_state_event_v1_t) == 16, "verse_state_event_v1_t arm64 layout"); +_Static_assert(sizeof(verse_error_event_v1_t) == 24, "verse_error_event_v1_t arm64 layout"); +_Static_assert(sizeof(verse_stats_event_v1_t) == 32, "verse_stats_event_v1_t arm64 layout"); +_Static_assert(sizeof(verse_media_event_v1_t) == 40, "verse_media_event_v1_t arm64 layout"); +_Static_assert(sizeof(verse_control_event_v1_t) == 32, "verse_control_event_v1_t arm64 layout"); +#endif + +#ifdef __cplusplus +} +#endif + +#endif diff --git a/core/src/abi.rs b/core/src/abi.rs new file mode 100644 index 0000000..fce1013 --- /dev/null +++ b/core/src/abi.rs @@ -0,0 +1,786 @@ +use std::cell::Cell; +use std::collections::HashMap; +use std::ffi::c_void; +use std::mem::size_of; +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::thread::{self, JoinHandle}; +use std::time::{Duration, Instant}; + +use crate::error::CoreError; +use crate::input::{encode_input, ControllerState, InputEvent}; +use crate::session::SessionCore; +use crate::wire::{ConnectionManifest, NativeTunnelCredential}; + +const ABI_V1: u32 = 1; +const OK: u32 = 0; +const INVALID_ARGUMENT: u32 = 1; +const INVALID_STATE: u32 = 2; +const UNSUPPORTED_ABI: u32 = 3; +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; +const STATE_CANCELLED: u32 = 3; + +const INPUT_KEYBOARD: u32 = 1; +const INPUT_MOUSE_BUTTON: u32 = 2; +const INPUT_RELATIVE_MOUSE: u32 = 3; +const INPUT_TEXT: u32 = 4; +const INPUT_CONTROLLER: u32 = 5; +const INPUT_ABSOLUTE_MOUSE: u32 = 6; +const INPUT_SCROLL: u32 = 7; + +const MAX_CONNECT_BYTES: usize = 1024 * 1024; +const CALLBACK_QUEUE_CAPACITY: usize = 64; +const CONTROL_QUEUE_CAPACITY: usize = 64; + +static NEXT_ID: AtomicUsize = AtomicUsize::new(1); +static HANDLES: OnceLock>>> = OnceLock::new(); + +thread_local! { + static ACTIVE_CALLBACK: Cell = const { Cell::new(0) }; + static ACTIVE_SIGNER: Cell = const { Cell::new(false) }; +} + +#[repr(C)] +#[derive(Clone, Copy)] +struct BytesView { + data: *const u8, + length: usize, +} + +// Safety: this copied pair is only retained as an ABI value in records whose backing bytes are +// owned by the same queued event; it is never shared independently. +unsafe impl Send for BytesView {} + +#[repr(C)] +struct StateEvent { + struct_size: u32, + abi_version: u32, + state: u32, + reason: u32, +} + +#[repr(C)] +struct ErrorEvent { + struct_size: u32, + abi_version: u32, + code: u32, + retryable: u32, + phase: u32, + reserved: u32, +} + +#[repr(C)] +struct StatsEvent { + struct_size: u32, + abi_version: u32, + dropped_callbacks: u64, + dropped_media_units: u64, + dropped_input_events: u64, +} + +#[repr(C)] +struct MediaEvent { + struct_size: u32, + abi_version: u32, + channel: u32, + sequence: u32, + timestamp_ms: u64, + encoded_unit: BytesView, +} + +#[repr(C)] +struct ControlEvent { + struct_size: u32, + abi_version: u32, + kind: u32, + reserved: u32, + payload: BytesView, +} + +type SignFn = unsafe extern "C" fn(*mut c_void, BytesView, *mut u8) -> u32; +type StateFn = unsafe extern "C" fn(*mut c_void, *const StateEvent); +type ErrorFn = unsafe extern "C" fn(*mut c_void, *const ErrorEvent); +type StatsFn = unsafe extern "C" fn(*mut c_void, *const StatsEvent); +type MediaFn = unsafe extern "C" fn(*mut c_void, *const MediaEvent); +type ControlFn = unsafe extern "C" fn(*mut c_void, *const ControlEvent); + +#[repr(C)] +struct Config { + struct_size: u32, + abi_version: u32, + context: *mut c_void, + sign_admission: Option, + sign_tls_ed25519: Option, + on_state: Option, + on_error: Option, + on_stats: Option, + on_media: Option, + on_control: Option, +} + +#[repr(C)] +struct ConnectRequest { + struct_size: u32, + abi_version: u32, + manifest_json: BytesView, + tunnel_credential_json: BytesView, +} + +#[repr(C)] +struct AbiInputEvent { + struct_size: u32, + abi_version: u32, + kind: u32, + flags: u32, + values: [i32; 12], +} + +#[derive(Clone, Copy)] +struct Callbacks { + context: usize, + sign_admission: SignFn, + sign_tls_ed25519: SignFn, + on_state: Option, + on_error: Option, + on_stats: Option, + on_media: Option, + on_control: Option, +} + +enum CallbackEvent { + State(u32), +} + +struct SessionState { + session: SessionCore, + lifecycle: Lifecycle, + idr_queue: usize, +} + +#[derive(Clone, Copy, Eq, PartialEq)] +enum Lifecycle { + Created, + Connecting, + Connected, + Cancelled, + Destroying, +} + +struct CoreInner { + id: usize, + state: Mutex, + callbacks: Callbacks, + callback_tx: SyncSender, + callbacks_closed: Arc, + stop_tx: Sender<()>, + worker_done: Mutex>, + worker: Mutex>>, + in_flight: AtomicUsize, + freeing: AtomicBool, + destroying: AtomicBool, + dropped_callbacks: AtomicU64, +} + +#[repr(C)] +struct CoreHandle { + marker: u8, +} + +struct CallGuard { + inner: Arc, +} + +impl Drop for CallGuard { + fn drop(&mut self) { + self.inner.in_flight.fetch_sub(1, Ordering::AcqRel); + } +} + +fn lock(mutex: &Mutex) -> MutexGuard<'_, T> { + mutex + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + +fn handles() -> &'static Mutex>> { + HANDLES.get_or_init(|| Mutex::new(HashMap::new())) +} + +fn callback_id() -> usize { + ACTIVE_CALLBACK.with(Cell::get) +} + +fn status(error: CoreError) -> u32 { + match error { + CoreError::QueueFull => QUEUE_FULL, + CoreError::Cancelled => CANCELLED, + _ => INVALID_ARGUMENT, + } +} + +fn ffi_boundary(action: impl FnOnce() -> u32) -> u32 { + catch_unwind(AssertUnwindSafe(action)).unwrap_or(INTERNAL) +} + +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() { + cancel_inner(&inner); + } + INTERNAL + } +} + +unsafe fn header(pointer: *const u8) -> Result<(u32, u32), u32> { + if pointer.is_null() { + return Err(INVALID_ARGUMENT); + } + // Invariant: the caller must provide the ABI-mandated readable 8-byte prefix. + let struct_size = unsafe { ptr::read_unaligned(pointer.cast::()) }; + // Invariant: the caller must provide the ABI-mandated readable 8-byte prefix. + let abi_version = unsafe { ptr::read_unaligned(pointer.add(4).cast::()) }; + if struct_size < 8 { + return Err(INVALID_ARGUMENT); + } + if abi_version != ABI_V1 { + return Err(UNSUPPORTED_ABI); + } + Ok((struct_size, abi_version)) +} + +unsafe fn read_table(pointer: *const T) -> Result { + // Invariant: every sized table begins with the common readable 8-byte prefix. + let (struct_size, _) = unsafe { header(pointer.cast()) }?; + if usize::try_from(struct_size).map_err(|_| INVALID_ARGUMENT)? < size_of::() { + return Err(INVALID_ARGUMENT); + } + // Invariant: struct_size proves the complete v1 table is readable; trailing bytes are ignored. + Ok(unsafe { ptr::read_unaligned(pointer) }) +} + +unsafe fn copy_view(view: BytesView) -> Result, u32> { + if view.length > MAX_CONNECT_BYTES || (view.data.is_null() && view.length != 0) { + return Err(INVALID_ARGUMENT); + } + if view.length == 0 { + return Ok(Vec::new()); + } + // Invariant: non-null pointer plus caller-provided length is readable for this synchronous call. + Ok(unsafe { std::slice::from_raw_parts(view.data, view.length) }.to_vec()) +} + +unsafe fn begin_call(core: *mut CoreHandle) -> Result<(Arc, CallGuard), u32> { + let core = NonNull::new(core).ok_or(INVALID_ARGUMENT)?; + let inner = { + let live = lock(handles()); + let inner = live + .get(&(core.as_ptr() as usize)) + .cloned() + .ok_or(INVALID_ARGUMENT)?; + inner.in_flight.fetch_add(1, Ordering::AcqRel); + inner + }; + if inner.freeing.load(Ordering::Acquire) { + inner.in_flight.fetch_sub(1, Ordering::AcqRel); + return Err(INVALID_ARGUMENT); + } + Ok((Arc::clone(&inner), CallGuard { inner })) +} + +fn enqueue_callback_under_state_lock(inner: &CoreInner, event: CallbackEvent) { + if inner.callbacks_closed.load(Ordering::Acquire) { + return; + } + if let Err(TrySendError::Full(_)) = inner.callback_tx.try_send(event) { + inner.dropped_callbacks.fetch_add(1, Ordering::Relaxed); + } +} + +#[allow(clippy::needless_pass_by_value)] +fn callback_worker( + id: usize, + callbacks: Callbacks, + receiver: Receiver, + stop: Receiver<()>, + done: Sender<()>, + callbacks_closed: Arc, + callback_gate: Arc>, +) { + loop { + if stop.try_recv().is_ok() { + break; + } + match receiver.recv_timeout(Duration::from_millis(2)) { + Ok(CallbackEvent::State(state)) => { + if let Some(callback) = callbacks.on_state { + let _gate = lock(&callback_gate); + if callbacks_closed.load(Ordering::Acquire) { + continue; + } + let event = StateEvent { + struct_size: u32::try_from(size_of::()).expect("ABI size fits"), + abi_version: ABI_V1, + state, + reason: 0, + }; + ACTIVE_CALLBACK.with(|active| active.set(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) => {} + Err(RecvTimeoutError::Disconnected) => break, + } + } + let _ = done.send(()); +} + +fn cancel_inner(inner: &CoreInner) { + let mut state = lock(&inner.state); + let already_cancelled = state.session.is_cancelled(); + state.session.cancel(); + state.lifecycle = Lifecycle::Cancelled; + if !already_cancelled { + enqueue_callback_under_state_lock(inner, CallbackEvent::State(STATE_CANCELLED)); + } +} + +fn call_signer(id: usize, 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)); + // 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 +} + +fn convert_input(event: &AbiInputEvent) -> Result { + if event.flags != 0 { + return Err(INVALID_ARGUMENT); + } + let values = event.values; + let unused_are_zero = |used: usize| values[used..].iter().all(|value| *value == 0); + let value_u8 = |index: usize| u8::try_from(values[index]).map_err(|_| INVALID_ARGUMENT); + let unsigned_16 = |index: usize| u16::try_from(values[index]).map_err(|_| INVALID_ARGUMENT); + let signed_16 = |index: usize| i16::try_from(values[index]).map_err(|_| INVALID_ARGUMENT); + let pressed = || match values[0] { + 0 => Ok(false), + 1 => Ok(true), + _ => Err(INVALID_ARGUMENT), + }; + match event.kind { + INPUT_KEYBOARD if unused_are_zero(3) => Ok(InputEvent::Keyboard { + pressed: pressed()?, + modifiers: value_u8(1)?, + scancode: unsigned_16(2)?, + }), + INPUT_MOUSE_BUTTON if unused_are_zero(2) => Ok(InputEvent::MouseButton { + pressed: pressed()?, + button: value_u8(1)?, + }), + INPUT_RELATIVE_MOUSE if unused_are_zero(2) => Ok(InputEvent::RelativeMouse { + delta_x: signed_16(0)?, + delta_y: signed_16(1)?, + }), + INPUT_TEXT if unused_are_zero(1) => Ok(InputEvent::Text( + char::from_u32(u32::try_from(values[0]).map_err(|_| INVALID_ARGUMENT)?) + .ok_or(INVALID_ARGUMENT)?, + )), + INPUT_CONTROLLER if unused_are_zero(10) => Ok(InputEvent::Controller(ControllerState { + controller: value_u8(0)?, + active_mask: unsigned_16(1)?, + button_flags: unsigned_16(2)?, + left_trigger: value_u8(3)?, + right_trigger: value_u8(4)?, + left_x: signed_16(5)?, + left_y: signed_16(6)?, + right_x: signed_16(7)?, + right_y: signed_16(8)?, + extra_button_flags: unsigned_16(9)?, + })), + INPUT_ABSOLUTE_MOUSE if unused_are_zero(4) => Ok(InputEvent::AbsoluteMouse { + x: unsigned_16(0)?, + y: unsigned_16(1)?, + viewport_width: unsigned_16(2)?, + viewport_height: unsigned_16(3)?, + }), + INPUT_SCROLL if unused_are_zero(2) => Ok(InputEvent::Scroll { + vertical_delta: signed_16(0)?, + horizontal_delta: signed_16(1)?, + }), + _ => Err(INVALID_ARGUMENT), + } +} + +#[no_mangle] +pub extern "C" fn verse_core_abi_version() -> u32 { + ffi_boundary(|| ABI_V1) +} + +#[no_mangle] +unsafe extern "C" fn verse_core_create_v1( + config: *const Config, + out_core: *mut *mut CoreHandle, +) -> u32 { + ffi_boundary(|| { + if out_core.is_null() { + return INVALID_ARGUMENT; + } + // Invariant: caller supplies a writable out pointer for this synchronous call. + unsafe { out_core.write(ptr::null_mut()) }; + // Invariant: config points to a readable sized/versioned table. + let config = match unsafe { read_table(config) } { + Ok(config) => config, + Err(error) => return error, + }; + let (Some(sign_admission), Some(sign_tls_ed25519)) = + (config.sign_admission, config.sign_tls_ed25519) + else { + return INVALID_ARGUMENT; + }; + let callbacks = Callbacks { + context: config.context as usize, + sign_admission, + sign_tls_ed25519, + on_state: config.on_state, + on_error: config.on_error, + on_stats: config.on_stats, + on_media: config.on_media, + on_control: config.on_control, + }; + let (callback_tx, callback_rx) = mpsc::sync_channel(CALLBACK_QUEUE_CAPACITY); + let (stop_tx, stop_rx) = mpsc::channel(); + let (done_tx, done_rx) = mpsc::channel(); + let callbacks_closed = Arc::new(AtomicBool::new(false)); + let callback_gate = Arc::new(Mutex::new(())); + let id = NEXT_ID.fetch_add(1, Ordering::Relaxed); + let worker_callbacks_closed = Arc::clone(&callbacks_closed); + let worker_callback_gate = Arc::clone(&callback_gate); + let worker = thread::Builder::new() + .name("verse-core-callback".to_owned()) + .spawn(move || { + callback_worker( + id, + callbacks, + callback_rx, + stop_rx, + done_tx, + worker_callbacks_closed, + worker_callback_gate, + ); + }); + let Ok(worker) = worker else { + return INTERNAL; + }; + let inner = Arc::new(CoreInner { + id, + state: Mutex::new(SessionState { + session: SessionCore::new(), + lifecycle: Lifecycle::Created, + idr_queue: 0, + }), + callbacks, + callback_tx, + callbacks_closed, + stop_tx, + worker_done: Mutex::new(done_rx), + worker: Mutex::new(Some(worker)), + in_flight: AtomicUsize::new(0), + freeing: AtomicBool::new(false), + destroying: AtomicBool::new(false), + dropped_callbacks: AtomicU64::new(0), + }); + let handle = Box::new(CoreHandle { marker: 0 }); + let raw = Box::into_raw(handle); + lock(handles()).insert(raw as usize, inner); + // Invariant: out_core is writable and receives the newly owned opaque handle. + unsafe { out_core.write(raw) }; + OK + }) +} + +#[no_mangle] +unsafe extern "C" fn verse_core_connect_v1( + core: *mut CoreHandle, + request: *const ConnectRequest, +) -> u32 { + ffi_boundary_core(core, || { + // 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, + Err(error) => return error, + }; + // Invariant: views obey the pointer/length rules for this synchronous copy. + let manifest = match unsafe { copy_view(request.manifest_json) } { + Ok(bytes) => bytes, + Err(error) => return error, + }; + // Invariant: views obey the pointer/length rules for this synchronous copy. + let credential = match unsafe { copy_view(request.tunnel_credential_json) } { + Ok(bytes) => bytes, + Err(error) => return error, + }; + if manifest.is_empty() || credential.is_empty() { + return INVALID_ARGUMENT; + } + if ConnectionManifest::decode(&manifest).is_err() + || NativeTunnelCredential::decode(&credential).is_err() + { + return INVALID_ARGUMENT; + } + { + let mut state = lock(&handle.state); + if state.session.is_cancelled() { + return CANCELLED; + } + if state.lifecycle != Lifecycle::Created { + return INVALID_STATE; + } + state.lifecycle = Lifecycle::Connecting; + enqueue_callback_under_state_lock(&handle, CallbackEvent::State(STATE_CONNECTING)); + } + let admission = call_signer( + handle.id, + handle.callbacks.sign_admission, + handle.callbacks.context, + &manifest, + ); + if admission != OK { + cancel_inner(&handle); + return admission; + } + if lock(&handle.state).session.is_cancelled() { + return CANCELLED; + } + let tls = call_signer( + handle.id, + handle.callbacks.sign_tls_ed25519, + handle.callbacks.context, + &credential, + ); + if tls != OK { + cancel_inner(&handle); + return tls; + } + { + let mut state = lock(&handle.state); + if state.session.is_cancelled() || state.lifecycle == Lifecycle::Destroying { + state.lifecycle = Lifecycle::Cancelled; + return CANCELLED; + } + state.lifecycle = Lifecycle::Connected; + enqueue_callback_under_state_lock(&handle, CallbackEvent::State(STATE_CONNECTED)); + } + OK + }) +} + +#[no_mangle] +unsafe extern "C" fn verse_core_send_input_v1( + core: *mut CoreHandle, + event: *const AbiInputEvent, +) -> u32 { + ffi_boundary_core(core, || { + // 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() { + return CANCELLED; + } + if state.lifecycle != Lifecycle::Connected { + return INVALID_STATE; + } + let event = match convert_input(&event) { + Ok(event) => event, + Err(error) => return error, + }; + let features = ["input.absolute.v1", "input.scroll.v1"]; + let bytes = match encode_input(&event, &features) { + Ok(bytes) => bytes, + Err(error) => return status(error), + }; + match state.session.enqueue_input(bytes, &features) { + Ok(()) => OK, + Err(error) => status(error), + } + }) +} + +#[no_mangle] +unsafe extern "C" fn verse_core_request_idr_v1(core: *mut CoreHandle) -> u32 { + ffi_boundary_core(core, || { + // 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() { + return CANCELLED; + } + if state.lifecycle != Lifecycle::Connected { + return INVALID_STATE; + } + if state.idr_queue == CONTROL_QUEUE_CAPACITY { + return QUEUE_FULL; + } + state.idr_queue += 1; + OK + }) +} + +#[no_mangle] +unsafe extern "C" fn verse_core_cancel_v1(core: *mut CoreHandle) -> u32 { + ffi_boundary_core(core, || { + // 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; + } + cancel_inner(&handle); + OK + }) +} + +#[no_mangle] +unsafe extern "C" fn verse_core_destroy_v1(core: *mut CoreHandle, timeout_ms: u32) -> u32 { + ffi_boundary_core(core, || { + let Some(core) = NonNull::new(core) else { + return INVALID_ARGUMENT; + }; + let inner = { + let live = lock(handles()); + 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); + { + let mut state = lock(&inner.state); + 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; + } + } + { + let mut live = lock(handles()); + if inner.in_flight.load(Ordering::Acquire) != 0 { + inner.destroying.store(false, Ordering::Release); + return BUSY; + } + inner.freeing.store(true, Ordering::Release); + live.remove(&(core.as_ptr() as usize)); + } + // Invariant: all other calls and callback work have ended; this is the sole successful free. + unsafe { drop(Box::from_raw(core.as_ptr())) }; + OK + }) +} + +const _: () = { + assert!(size_of::() == 16); + assert!(size_of::() == 72); + assert!(size_of::() == 40); + assert!(size_of::() == 64); + assert!(size_of::() == 16); + assert!(size_of::() == 24); + assert!(size_of::() == 32); + assert!(size_of::() == 40); + assert!(size_of::() == 32); +}; + +// Keep provider-free callback slots part of the copied v1 table even before the fake session emits +// these record classes. +const _: fn(Callbacks) = |callbacks| { + let _ = ( + callbacks.on_error, + callbacks.on_stats, + callbacks.on_media, + callbacks.on_control, + ); +}; + +#[cfg(test)] +mod tests { + use super::{ffi_boundary, INTERNAL}; + + #[test] + fn panic_boundary_maps_to_internal_status() { + assert_eq!(ffi_boundary(|| panic!("contained test panic")), INTERNAL); + } +} diff --git a/core/src/lib.rs b/core/src/lib.rs index a478459..73f991a 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -1,4 +1,5 @@ -#![forbid(unsafe_code, unsafe_op_in_unsafe_fn)] +#![deny(unsafe_code)] +#![forbid(unsafe_op_in_unsafe_fn)] //! Provider-free wire codecs and bounded session primitives for `VerseVDI` clients. //! @@ -16,3 +17,6 @@ pub mod input; pub mod media; pub mod session; pub mod wire; + +#[allow(unsafe_code)] +mod abi; diff --git a/core/tests/abi_contract.rs b/core/tests/abi_contract.rs new file mode 100644 index 0000000..523fe8a --- /dev/null +++ b/core/tests/abi_contract.rs @@ -0,0 +1,692 @@ +#![allow( + unsafe_code, + clippy::borrow_as_ptr, + clippy::cast_possible_truncation, + clippy::items_after_statements +)] + +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::thread; +use std::time::{Duration, Instant}; + +const ABI_V1: u32 = 1; +const OK: u32 = 0; +const INVALID_ARGUMENT: u32 = 1; +const INVALID_STATE: u32 = 2; +const UNSUPPORTED_ABI: u32 = 3; +const QUEUE_FULL: u32 = 9; +const CANCELLED: u32 = 10; +const REENTRANT: u32 = 11; +const BUSY: u32 = 12; + +const STATE_CONNECTING: u32 = 1; +const STATE_CONNECTED: u32 = 2; +const STATE_CANCELLED: u32 = 3; +const INPUT_KEYBOARD: u32 = 1; + +#[repr(C)] +struct Core { + _private: [u8; 0], +} + +#[repr(C)] +#[derive(Clone, Copy)] +struct BytesView { + data: *const u8, + length: usize, +} + +#[repr(C)] +struct StateEvent { + struct_size: u32, + abi_version: u32, + state: u32, + reason: u32, +} + +#[repr(C)] +struct ErrorEvent { + struct_size: u32, + abi_version: u32, + code: u32, + retryable: u32, + phase: u32, + reserved: u32, +} + +#[repr(C)] +struct StatsEvent { + struct_size: u32, + abi_version: u32, + dropped_callbacks: u64, + dropped_media_units: u64, + dropped_input_events: u64, +} + +#[repr(C)] +struct MediaEvent { + struct_size: u32, + abi_version: u32, + channel: u32, + sequence: u32, + timestamp_ms: u64, + encoded_unit: BytesView, +} + +#[repr(C)] +struct ControlEvent { + struct_size: u32, + abi_version: u32, + kind: u32, + reserved: u32, + payload: BytesView, +} + +type SignFn = unsafe extern "C" fn(*mut c_void, BytesView, *mut u8) -> u32; +type StateFn = unsafe extern "C" fn(*mut c_void, *const StateEvent); +type ErrorFn = unsafe extern "C" fn(*mut c_void, *const ErrorEvent); +type StatsFn = unsafe extern "C" fn(*mut c_void, *const StatsEvent); +type MediaFn = unsafe extern "C" fn(*mut c_void, *const MediaEvent); +type ControlFn = unsafe extern "C" fn(*mut c_void, *const ControlEvent); + +#[repr(C)] +struct Config { + struct_size: u32, + abi_version: u32, + context: *mut c_void, + sign_admission: Option, + sign_tls_ed25519: Option, + on_state: Option, + on_error: Option, + on_stats: Option, + on_media: Option, + on_control: Option, +} + +#[repr(C)] +struct ConnectRequest { + struct_size: u32, + abi_version: u32, + manifest_json: BytesView, + tunnel_credential_json: BytesView, +} + +#[repr(C)] +struct InputEvent { + struct_size: u32, + abi_version: u32, + kind: u32, + flags: u32, + values: [i32; 12], +} + +unsafe extern "C" { + fn verse_core_abi_version() -> u32; + fn verse_core_create_v1(config: *const Config, out_core: *mut *mut Core) -> u32; + fn verse_core_connect_v1(core: *mut Core, request: *const ConnectRequest) -> u32; + fn verse_core_send_input_v1(core: *mut Core, event: *const InputEvent) -> u32; + fn verse_core_request_idr_v1(core: *mut Core) -> u32; + fn verse_core_cancel_v1(core: *mut Core) -> u32; + fn verse_core_destroy_v1(core: *mut Core, timeout_ms: u32) -> u32; +} + +const MANIFEST: &[u8] = br#"{ + "version":"1","purpose":"launch","session_id":"session","reconnect_sequence":0, + "gateway":{"id":"gateway","addresses":["gateway.test:443"],"public_identity":"gateway.test"}, + "tunnel":{"versions":["verse-gateway-v1/1"],"features":["control.v1","input.absolute.v1","input.scroll.v1"]}, + "profile":{"id":"standard","bounds":{"minimum_kbps":1000,"target_kbps":5000,"maximum_kbps":10000},"display_mode":{"resolution_width":1920,"resolution_height":1080,"fps":60}}, + "grant":{"opaque_value":"ggggggggggggggggggggggggggggggggggggggggggg","expires_at":"2099-01-01T00:00:00Z","audience":"audience"}, + "correlation_id":"correlation" +}"#; +const CREDENTIAL: &[u8] = br#"{"client_device_id":"device","device_key_id":"key","certificate_chain_pem":"-----BEGIN CERTIFICATE-----\nAQID\n-----END CERTIFICATE-----","trust_bundle_pem":"-----BEGIN CERTIFICATE-----\nAQID\n-----END CERTIFICATE-----","expires_at":"2099-01-01T00:00:00Z"}"#; + +struct Context { + core: AtomicUsize, + admission_calls: AtomicUsize, + tls_calls: AtomicUsize, + admission_input: Mutex>, + tls_input: Mutex>, + states: Mutex>, + wake: Condvar, + reentry_cancel: AtomicU32, + reentry_send: AtomicU32, + reentry_destroy: AtomicU32, + block_callbacks: AtomicBool, + release_callbacks: AtomicBool, + callback_active: AtomicUsize, + callback_max: AtomicUsize, + cancel_on_connecting: AtomicBool, +} + +impl Default for Context { + fn default() -> Self { + Self { + core: AtomicUsize::new(0), + admission_calls: AtomicUsize::new(0), + tls_calls: AtomicUsize::new(0), + admission_input: Mutex::new(Vec::new()), + tls_input: Mutex::new(Vec::new()), + states: Mutex::new(Vec::new()), + wake: Condvar::new(), + reentry_cancel: AtomicU32::new(u32::MAX), + reentry_send: AtomicU32::new(u32::MAX), + reentry_destroy: AtomicU32::new(u32::MAX), + block_callbacks: AtomicBool::new(false), + release_callbacks: AtomicBool::new(false), + callback_active: AtomicUsize::new(0), + callback_max: AtomicUsize::new(0), + cancel_on_connecting: AtomicBool::new(false), + } + } +} + +unsafe fn context<'a>(raw: *mut c_void) -> &'a Context { + // Test invariant: every callback receives the live Box supplied at create. + unsafe { &*raw.cast::() } +} + +unsafe extern "C" fn sign_admission(raw: *mut c_void, input: BytesView, output: *mut u8) -> u32 { + // Test invariant: ABI promises input is readable for input.length during this callback. + let bytes = unsafe { std::slice::from_raw_parts(input.data, input.length) }; + let ctx = unsafe { context(raw) }; + ctx.admission_calls.fetch_add(1, Ordering::SeqCst); + *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 +} + +unsafe extern "C" fn sign_tls(raw: *mut c_void, input: BytesView, output: *mut u8) -> u32 { + // Test invariant: ABI promises input is readable for input.length during this callback. + let bytes = unsafe { std::slice::from_raw_parts(input.data, input.length) }; + let ctx = unsafe { context(raw) }; + ctx.tls_calls.fetch_add(1, Ordering::SeqCst); + *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 +} + +unsafe extern "C" fn on_state(raw: *mut c_void, event: *const StateEvent) { + let ctx = unsafe { context(raw) }; + let active = ctx.callback_active.fetch_add(1, Ordering::SeqCst) + 1; + ctx.callback_max.fetch_max(active, Ordering::SeqCst); + // 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_CONNECTING && ctx.cancel_on_connecting.load(Ordering::SeqCst) { + let core = ctx.core.load(Ordering::SeqCst) as *mut Core; + ctx.reentry_cancel + .store(unsafe { verse_core_cancel_v1(core) }, Ordering::SeqCst); + } + + if ctx.block_callbacks.load(Ordering::SeqCst) && !ctx.release_callbacks.load(Ordering::SeqCst) { + let mut states = ctx.states.lock().expect("states lock"); + while !ctx.release_callbacks.load(Ordering::SeqCst) { + states = ctx.wake.wait(states).expect("callback wait"); + } + } + + if state == STATE_CONNECTED && ctx.reentry_cancel.load(Ordering::SeqCst) == u32::MAX { + let core = ctx.core.load(Ordering::SeqCst) as *mut Core; + // Test invariant: the stored handle is live until this callback and its destroy complete. + ctx.reentry_cancel + .store(unsafe { verse_core_cancel_v1(core) }, Ordering::SeqCst); + let event = keyboard_event(); + // Test invariant: event and handle remain valid for the synchronous call. + ctx.reentry_send.store( + unsafe { verse_core_send_input_v1(core, &event) }, + Ordering::SeqCst, + ); + // Test invariant: the callback intentionally probes the documented reentry rejection. + ctx.reentry_destroy + .store(unsafe { verse_core_destroy_v1(core, 1) }, Ordering::SeqCst); + } + ctx.callback_active.fetch_sub(1, Ordering::SeqCst); +} + +unsafe extern "C" fn sign_admission_reenters( + raw: *mut c_void, + input: BytesView, + output: *mut u8, +) -> u32 { + let ctx = unsafe { context(raw) }; + ctx.reentry_cancel.store( + unsafe { verse_core_cancel_v1(ctx.core.load(Ordering::SeqCst) as *mut Core) }, + Ordering::SeqCst, + ); + unsafe { sign_admission(raw, input, output) } +} + +unsafe extern "C" fn sign_admission_waits_for_cancel( + raw: *mut c_void, + input: BytesView, + output: *mut u8, +) -> u32 { + let ctx = unsafe { context(raw) }; + let deadline = Instant::now() + Duration::from_secs(2); + while ctx.reentry_cancel.load(Ordering::SeqCst) == u32::MAX { + assert!( + Instant::now() < deadline, + "connecting callback did not cancel" + ); + thread::yield_now(); + } + unsafe { sign_admission(raw, input, output) } +} + +fn config(ctx: &mut Context) -> Config { + Config { + struct_size: size_of::() as u32, + abi_version: ABI_V1, + context: ptr::from_mut(ctx).cast(), + sign_admission: Some(sign_admission), + sign_tls_ed25519: Some(sign_tls), + on_state: Some(on_state), + on_error: None, + on_stats: None, + on_media: None, + on_control: None, + } +} + +fn request(manifest: &[u8], credential: &[u8]) -> ConnectRequest { + ConnectRequest { + struct_size: size_of::() as u32, + abi_version: ABI_V1, + manifest_json: BytesView { + data: manifest.as_ptr(), + length: manifest.len(), + }, + tunnel_credential_json: BytesView { + data: credential.as_ptr(), + length: credential.len(), + }, + } +} + +fn keyboard_event() -> InputEvent { + let mut values = [0; 12]; + values[0] = 1; + values[2] = 30; + InputEvent { + struct_size: size_of::() as u32, + abi_version: ABI_V1, + kind: INPUT_KEYBOARD, + flags: 0, + values, + } +} + +fn create(ctx: &mut Context) -> *mut Core { + let mut core = ptr::null_mut(); + let config = config(ctx); + // Test invariant: config/out pointers remain valid for the synchronous create call. + assert_eq!(unsafe { verse_core_create_v1(&config, &mut core) }, OK); + assert!(!core.is_null()); + ctx.core.store(core as usize, Ordering::SeqCst); + core +} + +fn connect(core: *mut Core, manifest: &[u8], credential: &[u8]) -> u32 { + let request = request(manifest, credential); + // Test invariant: request and backing byte slices remain valid for the synchronous call. + unsafe { verse_core_connect_v1(core, &request) } +} + +fn wait_for(ctx: &Context, predicate: impl Fn(&[u32]) -> bool) { + let deadline = Instant::now() + Duration::from_secs(2); + let mut states = ctx.states.lock().expect("states lock"); + while !predicate(&states) { + let remaining = deadline.saturating_duration_since(Instant::now()); + assert!( + !remaining.is_zero(), + "callback deadline exceeded: {states:?}" + ); + (states, _) = ctx + .wake + .wait_timeout(states, remaining) + .expect("callback wait"); + } +} + +fn destroy(core: *mut Core) -> u32 { + // Test invariant: caller retains the handle until destroy reports success. + unsafe { verse_core_destroy_v1(core, 2_000) } +} + +#[test] +fn arm64_c_layout_is_exact_and_version_is_fixed() { + assert_eq!(versevdi_core::session::INPUT_QUEUE_CAPACITY, 64); + assert_eq!(size_of::(), 16); + assert_eq!(size_of::(), 72); + assert_eq!(offset_of!(Config, sign_admission), 16); + assert_eq!(size_of::(), 40); + assert_eq!(offset_of!(ConnectRequest, manifest_json), 8); + assert_eq!(size_of::(), 64); + assert_eq!(size_of::(), 16); + assert_eq!(size_of::(), 24); + assert_eq!(size_of::(), 32); + assert_eq!(size_of::(), 40); + assert_eq!(size_of::(), 32); + // Test invariant: no pointer arguments are involved. + assert_eq!(unsafe { verse_core_abi_version() }, ABI_V1); +} + +#[test] +fn create_validates_prefix_callbacks_output_and_trailing_bytes() { + let mut ctx = Context::default(); + let mut core = ptr::dangling_mut::(); + let mut cfg = config(&mut ctx); + + cfg.struct_size = 7; + // Test invariant: config/out are readable/writable for this call. + assert_eq!( + unsafe { verse_core_create_v1(&cfg, &mut core) }, + INVALID_ARGUMENT + ); + assert!(core.is_null()); + + cfg = config(&mut ctx); + cfg.abi_version = 2; + assert_eq!( + unsafe { verse_core_create_v1(&cfg, &mut core) }, + UNSUPPORTED_ABI + ); + assert!(core.is_null()); + + cfg = config(&mut ctx); + cfg.sign_tls_ed25519 = None; + assert_eq!( + unsafe { verse_core_create_v1(&cfg, &mut core) }, + INVALID_ARGUMENT + ); + assert!(core.is_null()); + + #[repr(C)] + struct Extended { + base: Config, + ignored: [u8; 32], + } + let extended = Extended { + base: config(&mut ctx), + ignored: [0xEE; 32], + }; + let mut extended = extended; + extended.base.struct_size = size_of::() as u32; + assert_eq!( + unsafe { verse_core_create_v1(&extended.base, &mut core) }, + OK + ); + assert_eq!(destroy(core), OK); +} + +#[test] +fn connect_copies_inputs_and_calls_purpose_specific_signers_once() { + let mut ctx = Context::default(); + let core = create(&mut ctx); + let mut manifest = MANIFEST.to_vec(); + let mut credential = CREDENTIAL.to_vec(); + assert_eq!(connect(core, &manifest, &credential), OK); + manifest.fill(b'x'); + credential.fill(b'y'); + + assert_eq!(ctx.admission_calls.load(Ordering::SeqCst), 1); + assert_eq!(ctx.tls_calls.load(Ordering::SeqCst), 1); + assert_eq!(&*ctx.admission_input.lock().expect("admission"), MANIFEST); + assert_eq!(&*ctx.tls_input.lock().expect("tls"), CREDENTIAL); + assert_eq!(destroy(core), OK); +} + +#[test] +fn tables_and_slices_reject_short_unsupported_null_and_oversized_inputs() { + let mut ctx = Context::default(); + ctx.reentry_cancel.store(OK, Ordering::SeqCst); + let core = create(&mut ctx); + let mut req = request(MANIFEST, CREDENTIAL); + + req.struct_size = 7; + assert_eq!( + unsafe { verse_core_connect_v1(core, &req) }, + INVALID_ARGUMENT + ); + req.struct_size = size_of::() as u32; + req.abi_version = 2; + assert_eq!( + unsafe { verse_core_connect_v1(core, &req) }, + UNSUPPORTED_ABI + ); + req.abi_version = ABI_V1; + req.manifest_json = BytesView { + data: ptr::null(), + length: 1, + }; + assert_eq!( + unsafe { verse_core_connect_v1(core, &req) }, + INVALID_ARGUMENT + ); + req.manifest_json = BytesView { + data: ptr::null(), + length: 0, + }; + assert_eq!( + unsafe { verse_core_connect_v1(core, &req) }, + INVALID_ARGUMENT + ); + req.manifest_json.length = 1_048_577; + assert_eq!( + unsafe { verse_core_connect_v1(core, &req) }, + INVALID_ARGUMENT + ); + + let mut input = keyboard_event(); + input.struct_size = 7; + assert_eq!( + unsafe { verse_core_send_input_v1(core, &input) }, + INVALID_ARGUMENT + ); + input.struct_size = size_of::() as u32; + input.abi_version = 2; + assert_eq!( + unsafe { verse_core_send_input_v1(core, &input) }, + UNSUPPORTED_ABI + ); + + #[repr(C)] + struct ExtendedRequest { + base: ConnectRequest, + ignored: [u8; 24], + } + let mut trailing_request = ExtendedRequest { + base: request(MANIFEST, CREDENTIAL), + ignored: [0xEE; 24], + }; + trailing_request.base.struct_size = size_of::() as u32; + assert_eq!( + unsafe { verse_core_connect_v1(core, &trailing_request.base) }, + OK + ); + + #[repr(C)] + struct ExtendedInput { + base: InputEvent, + ignored: [u8; 24], + } + let mut trailing_input = ExtendedInput { + base: keyboard_event(), + ignored: [0xEE; 24], + }; + trailing_input.base.struct_size = size_of::() as u32; + assert_eq!( + unsafe { verse_core_send_input_v1(core, &trailing_input.base) }, + OK + ); + assert_eq!(destroy(core), OK); +} + +#[test] +fn callback_order_is_serial_and_only_cancel_is_reentrant() { + let mut ctx = Context::default(); + let core = create(&mut ctx); + assert_eq!(connect(core, MANIFEST, CREDENTIAL), OK); + wait_for(&ctx, |states| states.contains(&STATE_CANCELLED)); + + assert_eq!( + ctx.states.lock().expect("states").as_slice(), + [STATE_CONNECTING, STATE_CONNECTED, STATE_CANCELLED] + ); + assert_eq!(ctx.reentry_cancel.load(Ordering::SeqCst), OK); + assert_eq!(ctx.reentry_send.load(Ordering::SeqCst), REENTRANT); + assert_eq!(ctx.reentry_destroy.load(Ordering::SeqCst), REENTRANT); + assert_eq!(ctx.callback_max.load(Ordering::SeqCst), 1); + assert_eq!(destroy(core), OK); +} + +#[test] +fn signer_callbacks_cannot_reenter_even_cancel() { + let mut ctx = Context::default(); + let mut cfg = config(&mut ctx); + cfg.sign_admission = Some(sign_admission_reenters); + 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_cancel.load(Ordering::SeqCst), REENTRANT); + assert_eq!(destroy(core), OK); +} + +#[test] +fn cancel_during_connect_preserves_state_order_and_stops_before_tls_signing() { + let mut ctx = Context::default(); + ctx.cancel_on_connecting.store(true, Ordering::SeqCst); + let mut cfg = config(&mut ctx); + cfg.sign_admission = Some(sign_admission_waits_for_cancel); + 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), CANCELLED); + wait_for(&ctx, |states| states.contains(&STATE_CANCELLED)); + assert_eq!( + ctx.states.lock().expect("states").as_slice(), + [STATE_CONNECTING, STATE_CANCELLED] + ); + assert_eq!(ctx.tls_calls.load(Ordering::SeqCst), 0); + assert_eq!(destroy(core), OK); +} + +#[test] +fn send_input_is_nonblocking_bounded_and_cancel_is_idempotent() { + 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 event = keyboard_event(); + for _ in 0..64 { + assert_eq!(unsafe { verse_core_send_input_v1(core, &event) }, OK); + } + assert_eq!( + unsafe { verse_core_send_input_v1(core, &event) }, + QUEUE_FULL + ); + assert_eq!(unsafe { verse_core_request_idr_v1(core) }, OK); + assert_eq!(unsafe { verse_core_cancel_v1(core) }, OK); + assert_eq!(unsafe { verse_core_cancel_v1(core) }, OK); + assert_eq!(unsafe { verse_core_send_input_v1(core, &event) }, CANCELLED); + assert_eq!(destroy(core), OK); +} + +#[test] +fn destroy_timeout_keeps_ownership_suppresses_late_callbacks_and_allows_retry() { + let mut ctx = Context::default(); + ctx.reentry_cancel.store(OK, Ordering::SeqCst); + ctx.block_callbacks.store(true, Ordering::SeqCst); + let core = create(&mut ctx); + assert_eq!(connect(core, MANIFEST, CREDENTIAL), OK); + wait_for(&ctx, |states| !states.is_empty()); + + assert_eq!(unsafe { verse_core_destroy_v1(core, 1) }, BUSY); + let count_at_timeout = ctx.states.lock().expect("states").len(); + 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(); + ctx.reentry_cancel.store(OK, Ordering::SeqCst); + let core = create(&mut ctx); + let event = keyboard_event(); + assert_eq!( + unsafe { verse_core_send_input_v1(core, &event) }, + INVALID_STATE + ); + assert_eq!(unsafe { verse_core_request_idr_v1(core) }, INVALID_STATE); + assert_eq!(unsafe { verse_core_cancel_v1(core) }, OK); + assert_eq!(connect(core, MANIFEST, CREDENTIAL), CANCELLED); + assert_eq!(destroy(core), OK); +} diff --git a/core/tests/ffi/abi_smoke.c b/core/tests/ffi/abi_smoke.c new file mode 100644 index 0000000..994a8d8 --- /dev/null +++ b/core/tests/ffi/abi_smoke.c @@ -0,0 +1,96 @@ +#include "versevdi_core.h" + +#include +#include +#include +#include + +static const char MANIFEST[] = + "{\"version\":\"1\",\"purpose\":\"launch\",\"session_id\":\"session\"," + "\"reconnect_sequence\":0,\"gateway\":{\"id\":\"gateway\",\"addresses\":[" + "\"gateway.test:443\"],\"public_identity\":\"gateway.test\"},\"tunnel\":{" + "\"versions\":[\"verse-gateway-v1/1\"],\"features\":[\"control.v1\"]}," + "\"profile\":{\"id\":\"standard\",\"bounds\":{\"minimum_kbps\":1000," + "\"target_kbps\":5000,\"maximum_kbps\":10000},\"display_mode\":null}," + "\"grant\":{\"opaque_value\":\"ggggggggggggggggggggggggggggggggggggggggggg\"," + "\"expires_at\":\"2099-01-01T00:00:00Z\",\"audience\":\"audience\"}," + "\"correlation_id\":\"correlation\"}"; +static const char CREDENTIAL[] = + "{\"client_device_id\":\"device\",\"device_key_id\":\"key\"," + "\"certificate_chain_pem\":\"-----BEGIN CERTIFICATE-----\\nAQID\\n-----END CERTIFICATE-----\"," + "\"trust_bundle_pem\":\"-----BEGIN CERTIFICATE-----\\nAQID\\n-----END CERTIFICATE-----\"," + "\"expires_at\":\"2099-01-01T00:00:00Z\"}"; + +typedef struct smoke_context { + atomic_uint admission_calls; + atomic_uint tls_calls; + atomic_uint state_calls; +} smoke_context_t; + +static verse_status_t sign_admission( + void *raw, + verse_bytes_view_t input, + uint8_t signature[64]) { + smoke_context_t *context = raw; + assert(input.length == sizeof(MANIFEST) - 1U); + memset(signature, 0xA5, 64U); + atomic_fetch_add(&context->admission_calls, 1U); + return VERSE_STATUS_OK; +} + +static verse_status_t sign_tls( + void *raw, + verse_bytes_view_t input, + uint8_t signature[64]) { + smoke_context_t *context = raw; + assert(input.length == sizeof(CREDENTIAL) - 1U); + memset(signature, 0x5A, 64U); + atomic_fetch_add(&context->tls_calls, 1U); + return VERSE_STATUS_OK; +} + +static void on_state(void *raw, const verse_state_event_v1_t *event) { + smoke_context_t *context = raw; + assert(event->struct_size == sizeof(*event)); + assert(event->abi_version == VERSE_CORE_ABI_VERSION_1); + atomic_fetch_add(&context->state_calls, 1U); +} + +int main(void) { + smoke_context_t context = {0}; + verse_core_config_v1_t config = { + .struct_size = sizeof(config), + .abi_version = VERSE_CORE_ABI_VERSION_1, + .context = &context, + .sign_admission = sign_admission, + .sign_tls_ed25519 = sign_tls, + .on_state = on_state, + }; + verse_core_t *core = NULL; + assert(verse_core_abi_version() == VERSE_CORE_ABI_VERSION_1); + assert(verse_core_create_v1(&config, &core) == VERSE_STATUS_OK); + assert(core != NULL); + + const verse_connect_request_v1_t request = { + .struct_size = sizeof(request), + .abi_version = VERSE_CORE_ABI_VERSION_1, + .manifest_json = {(const uint8_t *)MANIFEST, sizeof(MANIFEST) - 1U}, + .tunnel_credential_json = {(const uint8_t *)CREDENTIAL, sizeof(CREDENTIAL) - 1U}, + }; + assert(verse_core_connect_v1(core, &request) == VERSE_STATUS_OK); + + verse_input_event_v1_t input = { + .struct_size = sizeof(input), + .abi_version = VERSE_CORE_ABI_VERSION_1, + .kind = VERSE_INPUT_KEYBOARD, + .values = {1, 0, 30}, + }; + assert(verse_core_send_input_v1(core, &input) == VERSE_STATUS_OK); + assert(verse_core_request_idr_v1(core) == VERSE_STATUS_OK); + assert(verse_core_cancel_v1(core) == VERSE_STATUS_OK); + assert(verse_core_cancel_v1(core) == VERSE_STATUS_OK); + assert(verse_core_destroy_v1(core, 2000U) == VERSE_STATUS_OK); + assert(atomic_load(&context.admission_calls) == 1U); + assert(atomic_load(&context.tls_calls) == 1U); + return EXIT_SUCCESS; +}