Skip to main content

node_app_sdk_rust/
lib.rs

1//! # Node-App SDK for Rust
2//!
3//! Build native [Node-App] plugins as shared libraries (`cdylib`) with a
4//! safe, ergonomic Rust API. The SDK targets the **Node Host API v1**
5//! (the canonical C header is `core/host-abi-v1/include/node-host-api-v1.h`
6//! in the host repository).
7//!
8//! [Node-App]: https://github.com/econ-v1/node-app-distribution
9//!
10//! ## What this crate provides
11//!
12//! - The [`NodeApp`] trait — implement it on your plugin type to get HTTP,
13//!   event, and capability handling.
14//! - The [`declare_node_app!`] macro — generates all FFI boilerplate
15//!   (vtable, panic guards, JSON serialization) from your trait impl.
16//! - Helpers for talking to the host: [`log`], [`invoke_capability`],
17//!   [`publish_event`], [`get_config`], [`get_storage`], [`set_storage`].
18//! - Distributed-trace propagation across capability invocations
19//!   (thread-local span context — see [`CurrentTrace`]).
20//! - Hard limits matching the host: [`MAX_CAPABILITY_RESPONSE_SIZE`]
21//!   (16 MiB), [`MAX_EVENT_NAME_LEN`] (256 B), [`MAX_EVENT_DATA_LEN`]
22//!   (64 KiB).
23//!
24//! ## Stability
25//!
26//! This crate is published as **`1.0.1-experimental`**. The trait surface
27//! and the underlying C ABI are FROZEN for the v1 series, but breaking
28//! changes between `1.0.x-experimental` releases are possible until at
29//! least two first-party packages have proven the surface in production.
30//! Pin to an exact version in `Cargo.toml`.
31//!
32//! ## Quick start
33//!
34//! ```toml
35//! # Cargo.toml
36//! [lib]
37//! crate-type = ["cdylib"]
38//!
39//! [dependencies]
40//! node-app-sdk-rust = "1.0.1-experimental"
41//! serde_json = "1.0"
42//! ```
43//!
44//! ```rust,ignore
45//! use node_app_sdk_rust::*;
46//!
47//! #[derive(Default)]
48//! pub struct MyApp;
49//!
50//! impl NodeApp for MyApp {
51//!     fn metadata() -> NodeAppInfo {
52//!         NodeAppInfo {
53//!             name: "my-app".into(),
54//!             version: "0.1.0".into(),
55//!             author: "Me".into(),
56//!             description: "Hello-world Node-App".into(),
57//!             capabilities: vec!["http_handler".into()],
58//!         }
59//!     }
60//!
61//!     fn handle_request(&self, _req: AppRequest) -> Result<AppResponse, NodeAppError> {
62//!         Ok(AppResponse {
63//!             status: 200,
64//!             headers: Default::default(),
65//!             body: serde_json::json!({ "hello": "world" }),
66//!         })
67//!     }
68//! }
69//!
70//! declare_node_app!(MyApp);
71//! ```
72//!
73//! Build with `cargo build --release`; the resulting shared library plus a
74//! `manifest.json` is installed into the host's app directory.
75//!
76//! ## Calling host capabilities
77//!
78//! Use [`invoke_capability`] to call any capability registered with the
79//! host's capability router (e.g. `core.storage.get`,
80//! `core.lightning.payment.send`):
81//!
82//! ```rust,ignore
83//! use node_app_sdk_rust::{invoke_capability, CapabilityRequest};
84//!
85//! let response = invoke_capability(&CapabilityRequest {
86//!     id: "req-1".into(),
87//!     capability: "core.storage.get".into(),
88//!     payload: serde_json::json!({ "key": "user_pref" }),
89//!     caller_node_id: None,
90//!     trace_id: None,
91//!     span_id: None,
92//!     parent_span_id: None,
93//!     trace_depth: None,
94//! })?;
95//! # Ok::<(), node_app_sdk_rust::NodeAppError>(())
96//! ```
97//!
98//! Trace context (`trace_id`, `span_id`, `trace_depth`) is propagated
99//! automatically when invoking from inside `handle_capability` — you do
100//! not need to thread it manually.
101//!
102//! ## Publishing events
103//!
104//! ```rust,ignore
105//! use node_app_sdk_rust::publish_event;
106//!
107//! publish_event(
108//!     "my-app.something_happened",
109//!     &serde_json::json!({ "user_id": 42 }),
110//! )?;
111//! # Ok::<(), node_app_sdk_rust::NodeAppError>(())
112//! ```
113//!
114//! Event names **must** be namespaced with the app name (`my-app.*`); the
115//! host rejects un-namespaced events.
116
117#![warn(missing_docs)]
118#![warn(rustdoc::broken_intra_doc_links)]
119
120pub use node_app_api::types::{
121    AppEvent, AppRequest, AppResponse, Capabilities, CapabilityExample, CapabilityRequest,
122    CapabilityResponse, NodeAppInfo, ProvidedCapability,
123};
124pub use node_app_api::context::NodeAppContext;
125pub use node_app_api::ffi::{FfiResult, NodeAppMetadata, NodeAppVTable};
126pub use node_app_api::API_VERSION;
127
128use std::cell::RefCell;
129use std::ffi::CString;
130use std::sync::atomic::{AtomicPtr, Ordering};
131use std::sync::{Condvar, Mutex, MutexGuard, PoisonError};
132use std::time::Duration;
133
134/// Maximum response size for capability handlers (16 MiB).
135///
136/// The host enforces this limit on every capability response. Apps that
137/// produce a larger response will see [`FfiResult::error`] returned with
138/// error code `-6`. Use streaming or pagination for larger payloads.
139pub const MAX_CAPABILITY_RESPONSE_SIZE: usize = 16 * 1024 * 1024;
140
141/// Trace fields for the currently-executing capability span (per-thread).
142///
143/// Set at the start of `handle_capability` by the [`declare_node_app!`]
144/// macro and cleared on return. Allows [`invoke_capability`] to propagate
145/// distributed trace context automatically without the caller having to
146/// thread `trace_id` / `span_id` through every call site.
147#[doc(hidden)]
148#[derive(Clone)]
149pub struct CurrentTrace {
150    /// Trace ID inherited from the inbound capability request.
151    pub trace_id: String,
152    /// Span ID of the current execution — becomes `parent_span_id` for
153    /// sub-invocations made from this thread.
154    pub span_id: String,
155    /// Depth of the current span in the trace tree (root = 0).
156    pub depth: u8,
157}
158
159thread_local! {
160    /// Current trace context for the executing capability (set by
161    /// [`declare_node_app!`]).
162    #[doc(hidden)]
163    pub static CURRENT_TRACE: RefCell<Option<CurrentTrace>> = const { RefCell::new(None) };
164    #[doc(hidden)]
165    pub static CURRENT_INVOCATION_CONTEXT: RefCell<Option<String>> = const { RefCell::new(None) };
166}
167
168#[doc(hidden)]
169pub struct CurrentInvocationContextGuard(Option<String>);
170
171impl CurrentInvocationContextGuard {
172    #[doc(hidden)]
173    pub fn enter(invocation_context_id: Option<String>) -> Self {
174        let previous = CURRENT_INVOCATION_CONTEXT.with(|current| {
175            std::mem::replace(&mut *current.borrow_mut(), invocation_context_id)
176        });
177        Self(previous)
178    }
179}
180
181impl Drop for CurrentInvocationContextGuard {
182    fn drop(&mut self) {
183        let previous = self.0.take();
184        CURRENT_INVOCATION_CONTEXT.with(|current| {
185            *current.borrow_mut() = previous;
186        });
187    }
188}
189
190/// Global storage for the host context pointer.
191///
192/// Uses `AtomicPtr` instead of `OnceLock` so that the pointer can be
193/// **updated on every init call**. On macOS, `dlclose` does not actually
194/// unload user libraries (`man dlclose`: "Mac OS X does not support
195/// dynamic unloading"), so the same library instance is reused across
196/// hot-reloads. With a `OnceLock` the first-load context pointer would
197/// be retained permanently; after the first-load `HostData` is dropped
198/// that pointer becomes dangling, causing a SIGSEGV on the next reload
199/// when any SDK function (e.g. `log`) tries to read it.
200///
201/// `AtomicPtr` allows `__store_context` to atomically replace the pointer
202/// on each init, ensuring it always points to the current live `HostData`.
203///
204/// Safety invariant: the pointer is set to a valid `NodeAppContext` during
205/// `__node_app_init` and is only read while the app is alive. The host
206/// (NativeLoader) keeps `HostData` + `NodeAppContext` alive for the entire
207/// lifetime of the loaded app instance.
208static APP_CONTEXT: AtomicPtr<NodeAppContext> = AtomicPtr::new(std::ptr::null_mut());
209
210/// Error code a call gets when it arrives after the app's shutdown started.
211///
212/// From the moment the host calls the app's shutdown entry, the
213/// [`declare_node_app!`] wrappers refuse new `handle_request`,
214/// `handle_event` and `handle_capability` calls with this code instead of
215/// running them. The next `init` accepts calls again. See ADR-083.
216pub const ERROR_SHUTTING_DOWN: i32 = -11;
217
218/// How long the shutdown entry waits for in-flight calls before it logs a
219/// warning that names how many are still running.
220///
221/// The entry keeps waiting after this. It must not return while a call is
222/// still running inside the library: the host releases the library owner
223/// when the shutdown entry returns, and it calls the app without holding
224/// that owner. The host's own shutdown budget bounds the whole wait. See
225/// ADR-083.
226pub const SHUTDOWN_DRAIN_WARN_AFTER: Duration = Duration::from_secs(2);
227
228/// Admission gate for calls into the app (generated by [`declare_node_app!`]).
229///
230/// Counts the calls that run inside the app and refuses new ones after
231/// [`CallGate::close`], so the shutdown entry can tell the app to stop at
232/// once and take the write lock only after the in-flight calls return.
233#[doc(hidden)]
234pub struct CallGate {
235    state: Mutex<CallGateState>,
236    drained: Condvar,
237}
238
239struct CallGateState {
240    closed: bool,
241    in_flight: usize,
242}
243
244impl CallGate {
245    #[doc(hidden)]
246    pub const fn new() -> Self {
247        Self {
248            state: Mutex::new(CallGateState {
249                closed: false,
250                in_flight: 0,
251            }),
252            drained: Condvar::new(),
253        }
254    }
255
256    fn lock(&self) -> MutexGuard<'_, CallGateState> {
257        // No code panics while it holds this mutex, so poison carries no
258        // information; recover the state rather than wedge every call.
259        self.state.lock().unwrap_or_else(PoisonError::into_inner)
260    }
261
262    /// Admit one call, or `None` once the gate is closed. The call counts
263    /// as in flight until the returned permit drops.
264    #[doc(hidden)]
265    pub fn enter(&self) -> Option<CallPermit<'_>> {
266        let mut state = self.lock();
267        if state.closed {
268            return None;
269        }
270        state.in_flight += 1;
271        Some(CallPermit(self))
272    }
273
274    /// Accept calls again (the app is being initialized).
275    #[doc(hidden)]
276    pub fn open(&self) {
277        self.lock().closed = false;
278    }
279
280    /// Refuse every call from now on.
281    #[doc(hidden)]
282    pub fn close(&self) {
283        self.lock().closed = true;
284    }
285
286    /// Wait until no call is in flight, or until `timeout` passes when it is
287    /// `Some`. Returns the number of calls still in flight.
288    #[doc(hidden)]
289    pub fn wait_drained(&self, timeout: Option<Duration>) -> usize {
290        let state = self.lock();
291        let state = match timeout {
292            Some(timeout) => {
293                self.drained
294                    .wait_timeout_while(state, timeout, |state| state.in_flight > 0)
295                    .unwrap_or_else(PoisonError::into_inner)
296                    .0
297            }
298            None => self
299                .drained
300                .wait_while(state, |state| state.in_flight > 0)
301                .unwrap_or_else(PoisonError::into_inner),
302        };
303        state.in_flight
304    }
305}
306
307impl Default for CallGate {
308    fn default() -> Self {
309        Self::new()
310    }
311}
312
313/// One admitted call; see [`CallGate::enter`].
314#[doc(hidden)]
315pub struct CallPermit<'a>(&'a CallGate);
316
317impl Drop for CallPermit<'_> {
318    fn drop(&mut self) {
319        let mut state = self.0.lock();
320        state.in_flight -= 1;
321        if state.in_flight == 0 {
322            self.0.drained.notify_all();
323        }
324    }
325}
326
327// ── Public api-store LLM SDK (feature 472) ──────────────────────────────────
328//
329// Type-safe surface for the host's `/api/v2/public/api_store/llm/{call,stream}`
330// endpoints. Lives in its own submodule so plugins that don't talk to the
331// public LLM endpoint never see these symbols in their import surface.
332pub mod llm;
333
334/// Log levels accepted by the [`log`] function.
335///
336/// Use these constants instead of magic numbers when calling
337/// [`log`] directly. The convenience macros ([`log_info!`], etc.)
338/// take care of this for you.
339pub mod log_level {
340    /// Most verbose level — use for fine-grained tracing.
341    pub const TRACE: u32 = 0;
342    /// Debug-level diagnostics, typically not shown in production.
343    pub const DEBUG: u32 = 1;
344    /// Informational messages indicating normal operation.
345    pub const INFO: u32 = 2;
346    /// Warnings — recoverable issues or unusual conditions.
347    pub const WARN: u32 = 3;
348    /// Errors — operations that failed and require attention.
349    pub const ERROR: u32 = 4;
350}
351
352/// Log a message to the host using the stored context.
353///
354/// Logs are written to the host's per-app log file
355/// (`{log_dir}/{app_name}.log`) and forwarded to the host's tracing
356/// subscriber, so they appear in the daemon's main log output too.
357///
358/// This function is a no-op if the context was not provided during init
359/// or if the message contains invalid UTF-8.
360///
361/// # Arguments
362/// * `level` - Log level (0=trace, 1=debug, 2=info, 3=warn, 4=error). Use the [`log_level`] constants.
363/// * `message` - The log message (must not contain interior NUL bytes).
364///
365/// # Example
366/// ```ignore
367/// use node_app_sdk_rust::{log, log_level};
368///
369/// log(log_level::INFO, "App initialized successfully");
370/// log(log_level::ERROR, "Something went wrong!");
371/// ```
372///
373/// Most callers should use the convenience macros ([`log_info!`],
374/// [`log_error!`], etc.) which accept `format!`-style arguments.
375pub fn log(level: u32, message: &str) {
376    let ctx_ptr = APP_CONTEXT.load(Ordering::Acquire);
377    if ctx_ptr.is_null() {
378        return;
379    }
380
381    let c_message = match CString::new(message) {
382        Ok(s) => s,
383        Err(_) => return, // Invalid message (contains null byte)
384    };
385
386    // Safety: ctx_ptr is valid for the app's lifetime, and host_log is a valid function pointer
387    unsafe {
388        let ctx = &*ctx_ptr;
389        (ctx.host_log)(ctx.host_data, level, c_message.as_ptr());
390    }
391}
392
393/// Invoke a capability on the host via the capability router.
394///
395/// This is the primary mechanism for app-to-app communication. The host
396/// resolves the capability name to the providing app, dispatches the
397/// request, and returns the response. The provider may itself be a
398/// different app, the host kernel, or a remote node (transparent to the
399/// caller).
400///
401/// Trace context is propagated automatically: if this call happens
402/// inside `handle_capability` and the inbound request carried a
403/// `trace_id`, the same trace ID is injected on outbound calls. Callers
404/// may override this by setting `trace_id` explicitly on the request.
405///
406/// # Errors
407///
408/// Returns [`NodeAppError::CapabilityError`] when:
409/// - The host context is not available (called before init).
410/// - The host's invoke callback is not wired up.
411/// - The capability is not registered, the provider rejects the call,
412///   or the response cannot be deserialized.
413pub fn invoke_capability(request: &CapabilityRequest) -> Result<CapabilityResponse, NodeAppError> {
414    let ctx_ptr = APP_CONTEXT.load(Ordering::Acquire);
415    if ctx_ptr.is_null() {
416        return Err(NodeAppError::CapabilityError(
417            "Host context not available".into(),
418        ));
419    }
420
421    // Propagate distributed trace context if a parent span is active on this thread.
422    // Only inject when the caller hasn't already set trace fields.
423    let active_invocation_context =
424        CURRENT_INVOCATION_CONTEXT.with(|current| current.borrow().clone());
425    let active_trace = if request.trace_id.is_none() {
426        CURRENT_TRACE.with(|current| current.borrow().clone())
427    } else {
428        None
429    };
430    let effective_request: std::borrow::Cow<CapabilityRequest> =
431        if active_trace.is_some() || active_invocation_context.is_some() {
432            let mut injected = request.clone();
433            if let Some(trace) = active_trace {
434                injected.trace_id = Some(trace.trace_id);
435                injected.span_id = Some(trace.span_id);
436                injected.parent_span_id = None;
437                injected.trace_depth = Some(trace.depth);
438            }
439            if let Some(context_id) = active_invocation_context {
440                injected.invocation_context_id = Some(context_id);
441            }
442            std::borrow::Cow::Owned(injected)
443        } else {
444            std::borrow::Cow::Borrowed(request)
445        };
446
447    // Serialize the request to JSON
448    let request_json = serde_json::to_vec(effective_request.as_ref())?;
449
450    // Safety: ctx_ptr is valid for the app's lifetime
451    unsafe {
452        let ctx = &*ctx_ptr;
453
454        // Check if the callback is available
455        if ctx.host_invoke_capability as usize == 0 {
456            return Err(NodeAppError::CapabilityError(
457                "host_invoke_capability callback not available".into(),
458            ));
459        }
460
461        let result = (ctx.host_invoke_capability)(
462            ctx.host_data,
463            request_json.as_ptr(),
464            request_json.len(),
465        );
466
467        if result.success && !result.data.is_null() && result.data_len > 0 {
468            let response_slice = std::slice::from_raw_parts(result.data, result.data_len);
469            let response: CapabilityResponse = serde_json::from_slice(response_slice)
470                .map_err(|e| NodeAppError::CapabilityError(format!("Response deserialization error: {}", e)))?;
471            // Free the host-allocated data
472            // Note: The host is responsible for freeing this memory via its own allocator
473            Ok(response)
474        } else if !result.success {
475            Err(NodeAppError::CapabilityError(format!(
476                "Host capability invocation failed with error code {}",
477                result.error_code
478            )))
479        } else {
480            Err(NodeAppError::CapabilityError(
481                "Empty response from host".into(),
482            ))
483        }
484    }
485}
486
487/// Maximum event name length in bytes (256).
488///
489/// Names that exceed this limit are rejected by [`publish_event`] with
490/// [`NodeAppError::EventFailed`]. The host enforces the same limit on
491/// the receiving side.
492pub const MAX_EVENT_NAME_LEN: usize = 256;
493
494/// Maximum event data length in bytes (64 KiB).
495///
496/// Payloads that exceed this limit are rejected by [`publish_event`]
497/// with [`NodeAppError::EventFailed`]. For larger artifacts, store them
498/// (e.g. via `core.storage.insert`) and emit an event referencing the
499/// storage key instead.
500pub const MAX_EVENT_DATA_LEN: usize = 64 * 1024;
501
502/// Publish a domain event to the host event bus.
503///
504/// The event is queued asynchronously (fire-and-forget). The event
505/// name **must** be namespaced with the app name prefix
506/// (e.g. `lightning.payment_received`, `my-app.user_created`); the host
507/// rejects events whose name does not start with `{app_name}.`.
508///
509/// # Errors
510///
511/// Returns [`NodeAppError::EventFailed`] if the host context is not
512/// available, the event name or data exceeds size limits
513/// ([`MAX_EVENT_NAME_LEN`] / [`MAX_EVENT_DATA_LEN`]), or the host
514/// rejects the event.
515pub fn publish_event(name: &str, data: &serde_json::Value) -> Result<(), NodeAppError> {
516    let ctx_ptr = APP_CONTEXT.load(Ordering::Acquire);
517    if ctx_ptr.is_null() {
518        return Err(NodeAppError::EventFailed(
519            "Host context not available".into(),
520        ));
521    }
522
523    let name_bytes = name.as_bytes();
524    if name_bytes.len() > MAX_EVENT_NAME_LEN {
525        return Err(NodeAppError::EventFailed(format!(
526            "Event name exceeds {} byte limit (got {})",
527            MAX_EVENT_NAME_LEN,
528            name_bytes.len()
529        )));
530    }
531
532    let data_json = serde_json::to_vec(data)?;
533    if data_json.len() > MAX_EVENT_DATA_LEN {
534        return Err(NodeAppError::EventFailed(format!(
535            "Event data exceeds {} byte limit (got {})",
536            MAX_EVENT_DATA_LEN,
537            data_json.len()
538        )));
539    }
540
541    unsafe {
542        let ctx = &*ctx_ptr;
543
544        if ctx.host_publish_event as usize == 0 {
545            return Err(NodeAppError::EventFailed(
546                "host_publish_event callback not available".into(),
547            ));
548        }
549
550        let result = (ctx.host_publish_event)(
551            ctx.host_data,
552            name_bytes.as_ptr(),
553            name_bytes.len(),
554            data_json.as_ptr(),
555            data_json.len(),
556        );
557
558        if result == 0 {
559            Ok(())
560        } else {
561            Err(NodeAppError::EventFailed(format!(
562                "host_publish_event returned error code {}",
563                result
564            )))
565        }
566    }
567}
568
569/// Log a TRACE-level message using `format!`-style arguments.
570///
571/// No-op when the host context has not been wired up yet. See [`log`]
572/// for details about delivery semantics.
573#[macro_export]
574macro_rules! log_trace {
575    ($($arg:tt)*) => {
576        $crate::log($crate::log_level::TRACE, &format!($($arg)*))
577    };
578}
579
580/// Log a DEBUG-level message using `format!`-style arguments.
581#[macro_export]
582macro_rules! log_debug {
583    ($($arg:tt)*) => {
584        $crate::log($crate::log_level::DEBUG, &format!($($arg)*))
585    };
586}
587
588/// Log an INFO-level message using `format!`-style arguments.
589#[macro_export]
590macro_rules! log_info {
591    ($($arg:tt)*) => {
592        $crate::log($crate::log_level::INFO, &format!($($arg)*))
593    };
594}
595
596/// Log a WARN-level message using `format!`-style arguments.
597#[macro_export]
598macro_rules! log_warn {
599    ($($arg:tt)*) => {
600        $crate::log($crate::log_level::WARN, &format!($($arg)*))
601    };
602}
603
604/// Log an ERROR-level message using `format!`-style arguments.
605#[macro_export]
606macro_rules! log_error {
607    ($($arg:tt)*) => {
608        $crate::log($crate::log_level::ERROR, &format!($($arg)*))
609    };
610}
611
612/// Store the context pointer for use by host helper functions.
613///
614/// Called automatically by [`declare_node_app!`] during init. App code
615/// must not call this directly.
616///
617/// Uses `AtomicPtr::store` so the pointer is replaced on every init —
618/// required on macOS where `dlclose` never unloads user libraries and the
619/// same app instance is reused across hot-reloads.
620#[doc(hidden)]
621pub fn __store_context(ctx: *const NodeAppContext) {
622    APP_CONTEXT.store(ctx as *mut NodeAppContext, Ordering::Release);
623}
624
625/// Get a configuration value from the host by key.
626///
627/// Common keys include `data_dir`, `host_port`, `app_name`, and
628/// `app_id`. The full set of available keys is determined by the host
629/// at app-init time.
630///
631/// Returns `None` if the key is not registered or the context is not
632/// available (called before init).
633pub fn get_config(key: &str) -> Option<String> {
634    let ctx_ptr = APP_CONTEXT.load(Ordering::Acquire);
635    if ctx_ptr.is_null() {
636        return None;
637    }
638    let c_key = CString::new(key).ok()?;
639    unsafe {
640        let ctx = &*ctx_ptr;
641        let result = (ctx.host_get_config)(ctx.host_data, c_key.as_ptr());
642        if result.is_null() {
643            return None;
644        }
645        Some(std::ffi::CStr::from_ptr(result).to_string_lossy().into_owned())
646    }
647}
648
649/// Set a storage value scoped to this app.
650///
651/// Storage is persisted by the host across app restarts and is isolated
652/// per-app: other apps cannot read or write keys you set here. Useful
653/// for small bits of configuration or state; for larger or relational
654/// data, prefer the `core.storage.*` capabilities.
655///
656/// No-op when the host context is not available or `key`/`value`
657/// contain interior NUL bytes.
658pub fn set_storage(key: &str, value: &str) {
659    let ctx_ptr = APP_CONTEXT.load(Ordering::Acquire);
660    if ctx_ptr.is_null() {
661        return;
662    }
663    let c_key = match CString::new(key) {
664        Ok(s) => s,
665        Err(_) => return,
666    };
667    let c_value = match CString::new(value) {
668        Ok(s) => s,
669        Err(_) => return,
670    };
671    unsafe {
672        let ctx = &*ctx_ptr;
673        (ctx.host_set_storage)(ctx.host_data, c_key.as_ptr(), c_value.as_ptr());
674    }
675}
676
677/// Get a storage value by key, scoped to this app.
678///
679/// See [`set_storage`] for scoping semantics. Returns `None` if the key
680/// is not found or the context is not available.
681pub fn get_storage(key: &str) -> Option<String> {
682    let ctx_ptr = APP_CONTEXT.load(Ordering::Acquire);
683    if ctx_ptr.is_null() {
684        return None;
685    }
686    let c_key = CString::new(key).ok()?;
687    unsafe {
688        let ctx = &*ctx_ptr;
689        let result = (ctx.host_get_storage)(ctx.host_data, c_key.as_ptr());
690        if result.is_null() {
691            return None;
692        }
693        Some(std::ffi::CStr::from_ptr(result).to_string_lossy().into_owned())
694    }
695}
696
697/// Error type returned by app operations.
698///
699/// Variants surface specific failure modes from each lifecycle hook plus
700/// transport-level errors (serialization, capability dispatch). The
701/// `#[from] serde_json::Error` impl on [`NodeAppError::SerializationError`]
702/// allows the `?` operator to propagate JSON errors directly.
703#[derive(Debug, thiserror::Error)]
704pub enum NodeAppError {
705    /// Returned from [`NodeApp::init`] when initialization failed.
706    #[error("Initialization failed: {0}")]
707    InitFailed(String),
708    /// Returned from [`NodeApp::handle_request`] when handling failed.
709    #[error("Request handling failed: {0}")]
710    RequestFailed(String),
711    /// Returned from [`NodeApp::handle_event`] or [`publish_event`].
712    #[error("Event handling failed: {0}")]
713    EventFailed(String),
714    /// Returned from [`NodeApp::shutdown`] when graceful shutdown failed.
715    #[error("Shutdown failed: {0}")]
716    ShutdownFailed(String),
717    /// JSON (de)serialization error — auto-converted from
718    /// [`serde_json::Error`] via the `?` operator.
719    #[error("Serialization error: {0}")]
720    SerializationError(#[from] serde_json::Error),
721    /// Returned from [`NodeApp::handle_capability`] or [`invoke_capability`].
722    #[error("Capability error: {0}")]
723    CapabilityError(String),
724}
725
726/// Trait implemented by node-app plugins.
727///
728/// Pair this with [`declare_node_app!`] to generate all the FFI
729/// boilerplate. Implementors must be `Default + Send + Sync + 'static`
730/// because the macro stores the singleton instance behind a static
731/// `OnceLock<Mutex<T>>`.
732///
733/// All trait methods have sensible default implementations; override
734/// only the hooks the app uses. Apps that opt into a hook must also
735/// declare the matching capability in their manifest, e.g.
736/// `"http_handler"` for HTTP request routing.
737pub trait NodeApp: Default + Send + Sync + 'static {
738    /// Return app metadata (name, version, author, description, capabilities).
739    ///
740    /// Called once when the host loads the shared library. The values
741    /// are cached for the lifetime of the process.
742    fn metadata() -> NodeAppInfo;
743
744    /// Initialize the app with the host context.
745    ///
746    /// Called once after loading. The default impl is a no-op — override
747    /// to set up state (DB connections, caches, etc.). The context may
748    /// be `None` if the host has not wired up callbacks yet (very early
749    /// in bootstrap or in unit tests).
750    fn init(&mut self, _ctx: Option<&NodeAppContext>) -> Result<(), NodeAppError> {
751        Ok(())
752    }
753
754    /// Start shutting down: called at once when the host stops the app,
755    /// while earlier calls may still be running.
756    ///
757    /// Runs under a shared borrow and does not wait for in-flight calls.
758    /// From this point the SDK refuses new calls with
759    /// [`ERROR_SHUTTING_DOWN`]. Override it to make in-flight work return
760    /// soon: cancel a token, set a flag, close a channel. Keep it short and
761    /// do not block in it. [`NodeApp::shutdown`] runs after it, once every
762    /// in-flight call has returned. Default is a no-op. See ADR-083.
763    fn begin_shutdown(&self) {}
764
765    /// Shut down the app gracefully.
766    ///
767    /// Called before unloading, after [`NodeApp::begin_shutdown`] and after
768    /// every in-flight call has returned, so it has exclusive access.
769    /// Default is a no-op. Override to flush pending writes, close
770    /// connections, etc.
771    fn shutdown(&mut self) -> Result<(), NodeAppError> {
772        Ok(())
773    }
774
775    /// Handle an incoming HTTP request proxied from the host.
776    ///
777    /// Only invoked if the app declared the `http_handler` capability.
778    /// Default returns `501 Not Implemented`.
779    fn handle_request(&self, _request: AppRequest) -> Result<AppResponse, NodeAppError> {
780        Ok(AppResponse {
781            status: 501,
782            headers: Default::default(),
783            body: serde_json::json!({"error": "Not implemented"}),
784        })
785    }
786
787    /// Handle a domain event from the host event bus.
788    ///
789    /// Only invoked if the app declared the `event_listener` capability
790    /// and subscribed to the event's namespace via the manifest.
791    fn handle_event(&self, _event: AppEvent) -> Result<(), NodeAppError> {
792        Ok(())
793    }
794
795    /// Return the list of service capabilities this app provides.
796    ///
797    /// Override to declare capabilities for the host's capability
798    /// registry. Default returns an empty list (no capabilities
799    /// provided). Each entry must use the namespace `{app_name}.{domain}.{action}`
800    /// — the `core.*` namespace is reserved for first-party apps.
801    fn provided_capabilities() -> Vec<ProvidedCapability> {
802        Vec::new()
803    }
804
805    /// Handle a capability invocation from another app via the
806    /// capability router.
807    ///
808    /// Override to implement capability handling logic. The default
809    /// returns [`NodeAppError::CapabilityError`] indicating the
810    /// capability is not implemented.
811    ///
812    /// Trace context is automatically captured into thread-local state
813    /// for the duration of this call so [`invoke_capability`] can
814    /// propagate it on outbound calls without explicit threading.
815    fn handle_capability(
816        &self,
817        _request: CapabilityRequest,
818    ) -> Result<CapabilityResponse, NodeAppError> {
819        Err(NodeAppError::CapabilityError(
820            "Capability handling not implemented".into(),
821        ))
822    }
823}
824
825/// Generate all FFI boilerplate for a [`NodeApp`] implementation.
826///
827/// This macro creates:
828/// - A `OnceLock<RwLock<T>>` instance for the app (sound concurrent access).
829/// - A call gate: the shutdown entry calls [`NodeApp::begin_shutdown`] at
830///   once, refuses new calls with [`ERROR_SHUTTING_DOWN`], and takes the
831///   write lock for [`NodeApp::shutdown`] only after in-flight calls return.
832/// - The `_node_app_entry` export symbol returning a [`NodeAppVTable`].
833/// - FFI wrapper functions for `init`, `shutdown`, `handle_request`,
834///   `handle_event`, `handle_capability`, and `free`.
835/// - `catch_unwind` guards on every FFI boundary to prevent UB from
836///   panics crossing the Rust/C ABI.
837/// - Thread-local trace span management around `handle_capability`.
838///
839/// Invoke once at the crate root after defining the app type:
840///
841/// ```ignore
842/// declare_node_app!(MyApp);
843/// ```
844#[macro_export]
845macro_rules! declare_node_app {
846    ($app_type:ty) => {
847        // Feature 463 B5 fix — RwLock instead of Mutex so capability calls
848        // can run CONCURRENTLY (read-locked). The previous Mutex serialized
849        // every call, causing a deadlock when one capability (e.g. agent.prompt)
850        // invokes another capability on the SAME app re-entrantly through the
851        // graph engine (agent.prompt → graph.engine.send_event → graph node →
852        // graph.tasks.storage.* → core.agent_storage.* → BACK to agent app).
853        // init still takes the write lock. Shutdown takes it only after the
854        // call gate has drained, so begin_shutdown runs at once (ADR-083).
855        static APP_INSTANCE: std::sync::OnceLock<std::sync::RwLock<$app_type>> =
856            std::sync::OnceLock::new();
857        static CALL_GATE: $crate::CallGate = $crate::CallGate::new();
858        static VTABLE: std::sync::OnceLock<$crate::NodeAppVTable> =
859            std::sync::OnceLock::new();
860
861        // OnceLock-backed CStrings for metadata (valid for process lifetime)
862        static META_NAME: std::sync::OnceLock<std::ffi::CString> = std::sync::OnceLock::new();
863        static META_VERSION: std::sync::OnceLock<std::ffi::CString> = std::sync::OnceLock::new();
864        static META_AUTHOR: std::sync::OnceLock<std::ffi::CString> = std::sync::OnceLock::new();
865        static META_DESCRIPTION: std::sync::OnceLock<std::ffi::CString> =
866            std::sync::OnceLock::new();
867
868        unsafe extern "C" fn __node_app_init(
869            ctx: *const std::os::raw::c_void,
870        ) -> $crate::FfiResult {
871            match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
872                let ctx_opt = if ctx.is_null() {
873                    None
874                } else {
875                    // Store context pointer for use by log() helper
876                    let ctx_typed = ctx as *const $crate::NodeAppContext;
877                    $crate::__store_context(ctx_typed);
878                    Some(unsafe { &*ctx_typed })
879                };
880                let app = APP_INSTANCE.get_or_init(|| {
881                    std::sync::RwLock::new(<$app_type>::default())
882                });
883                let mut guard = match app.write() {
884                    Ok(g) => g,
885                    Err(e) => {
886                        eprintln!("[node-app] rwlock poisoned in init: {}", e);
887                        return $crate::FfiResult::error(-10);
888                    }
889                };
890                // A cached library is initialized again after a shutdown
891                // (reload), so accept calls again.
892                CALL_GATE.open();
893                match guard.init(ctx_opt) {
894                    Ok(()) => $crate::FfiResult::ok(),
895                    Err(e) => {
896                        let msg = format!("init error: {}", e);
897                        $crate::log($crate::log_level::ERROR, &msg);
898                        eprintln!("[node-app] {}", msg);
899                        $crate::FfiResult::error(-1)
900                    }
901                }
902            })) {
903                Ok(result) => result,
904                Err(_) => {
905                    eprintln!("[node-app] panic in init");
906                    $crate::FfiResult::error(-99)
907                }
908            }
909        }
910
911        unsafe extern "C" fn __node_app_shutdown() -> $crate::FfiResult {
912            match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
913                if let Some(app) = APP_INSTANCE.get() {
914                    // ADR-083: refuse new calls, tell the app at once, then
915                    // take the write lock only after in-flight calls return.
916                    // Never return before they do: the host may unload the
917                    // library when this entry returns.
918                    CALL_GATE.close();
919                    let hook = match app.read() {
920                        Ok(app) => std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
921                            <$app_type as $crate::NodeApp>::begin_shutdown(&app)
922                        })),
923                        Err(e) => {
924                            eprintln!("[node-app] rwlock poisoned in shutdown: {}", e);
925                            return $crate::FfiResult::error(-10);
926                        }
927                    };
928                    // A panicking hook must not skip the drain below.
929                    if hook.is_err() {
930                        let msg = "panic in begin_shutdown; still waiting for in-flight calls";
931                        $crate::log($crate::log_level::ERROR, msg);
932                        eprintln!("[node-app] {}", msg);
933                    }
934                    let in_flight =
935                        CALL_GATE.wait_drained(Some($crate::SHUTDOWN_DRAIN_WARN_AFTER));
936                    if in_flight > 0 {
937                        let msg = format!(
938                            "shutdown: {} call(s) still in flight after {} ms; waiting for them to return",
939                            in_flight,
940                            $crate::SHUTDOWN_DRAIN_WARN_AFTER.as_millis()
941                        );
942                        $crate::log($crate::log_level::WARN, &msg);
943                        eprintln!("[node-app] {}", msg);
944                        CALL_GATE.wait_drained(None);
945                    }
946                    let mut guard = match app.write() {
947                        Ok(g) => g,
948                        Err(e) => {
949                            eprintln!("[node-app] rwlock poisoned in shutdown: {}", e);
950                            return $crate::FfiResult::error(-10);
951                        }
952                    };
953                    match guard.shutdown() {
954                        Ok(()) => $crate::FfiResult::ok(),
955                        Err(e) => {
956                            let msg = format!("shutdown error: {}", e);
957                            $crate::log($crate::log_level::ERROR, &msg);
958                            eprintln!("[node-app] {}", msg);
959                            $crate::FfiResult::error(-1)
960                        }
961                    }
962                } else {
963                    $crate::FfiResult::ok()
964                }
965            })) {
966                Ok(result) => result,
967                Err(_) => {
968                    eprintln!("[node-app] panic in shutdown");
969                    $crate::FfiResult::error(-99)
970                }
971            }
972        }
973
974        unsafe extern "C" fn __node_app_handle_request(
975            request_json: *const u8,
976            request_len: usize,
977        ) -> $crate::FfiResult {
978            match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
979                let json_slice =
980                    unsafe { std::slice::from_raw_parts(request_json, request_len) };
981                let request: $crate::AppRequest = match serde_json::from_slice(json_slice) {
982                    Ok(r) => r,
983                    Err(e) => {
984                        eprintln!("[node-app] request deserialization error: {}", e);
985                        return $crate::FfiResult::error(-2);
986                    }
987                };
988
989                let app = match APP_INSTANCE.get() {
990                    Some(a) => a,
991                    None => return $crate::FfiResult::error(-3),
992                };
993                // ADR-083 — refused once shutdown starts; counted in flight
994                // until this permit drops (after the read guard below).
995                let _call = match CALL_GATE.enter() {
996                    Some(call) => call,
997                    None => return $crate::FfiResult::error($crate::ERROR_SHUTTING_DOWN),
998                };
999                // Feature 463 B5 — read-lock allows concurrent calls
1000                let guard = match app.read() {
1001                    Ok(g) => g,
1002                    Err(e) => {
1003                        eprintln!("[node-app] rwlock poisoned in handle_request: {}", e);
1004                        return $crate::FfiResult::error(-10);
1005                    }
1006                };
1007
1008                match guard.handle_request(request) {
1009                    Ok(response) => match serde_json::to_vec(&response) {
1010                        Ok(bytes) => {
1011                            let len = bytes.len();
1012                            let boxed = bytes.into_boxed_slice();
1013                            let ptr = Box::into_raw(boxed) as *mut u8;
1014                            $crate::FfiResult {
1015                                success: true,
1016                                error_code: 0,
1017                                data: ptr,
1018                                data_len: len,
1019                            }
1020                        }
1021                        Err(e) => {
1022                            eprintln!("[node-app] response serialization error: {}", e);
1023                            $crate::FfiResult::error(-4)
1024                        }
1025                    },
1026                    Err(e) => {
1027                        eprintln!("[node-app] handle_request error: {}", e);
1028                        $crate::FfiResult::error(-5)
1029                    }
1030                }
1031            })) {
1032                Ok(result) => result,
1033                Err(_) => {
1034                    eprintln!("[node-app] panic in handle_request");
1035                    $crate::FfiResult::error(-99)
1036                }
1037            }
1038        }
1039
1040        unsafe extern "C" fn __node_app_handle_event(
1041            event_json: *const u8,
1042            event_len: usize,
1043        ) -> $crate::FfiResult {
1044            match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1045                let json_slice =
1046                    unsafe { std::slice::from_raw_parts(event_json, event_len) };
1047                let event: $crate::AppEvent = match serde_json::from_slice(json_slice) {
1048                    Ok(e) => e,
1049                    Err(e) => {
1050                        eprintln!("[node-app] event deserialization error: {}", e);
1051                        return $crate::FfiResult::error(-2);
1052                    }
1053                };
1054
1055                let app = match APP_INSTANCE.get() {
1056                    Some(a) => a,
1057                    None => return $crate::FfiResult::error(-3),
1058                };
1059                // ADR-083 — refused once shutdown starts; declared before the
1060                // read guard so the guard drops first.
1061                let _call = match CALL_GATE.enter() {
1062                    Some(call) => call,
1063                    None => return $crate::FfiResult::error($crate::ERROR_SHUTTING_DOWN),
1064                };
1065                // Feature 463 B5 — read-lock allows concurrent calls
1066                let guard = match app.read() {
1067                    Ok(g) => g,
1068                    Err(e) => {
1069                        eprintln!("[node-app] rwlock poisoned in handle_event: {}", e);
1070                        return $crate::FfiResult::error(-10);
1071                    }
1072                };
1073
1074                match guard.handle_event(event) {
1075                    Ok(()) => $crate::FfiResult::ok(),
1076                    Err(e) => {
1077                        eprintln!("[node-app] handle_event error: {}", e);
1078                        $crate::FfiResult::error(-5)
1079                    }
1080                }
1081            })) {
1082                Ok(result) => result,
1083                Err(_) => {
1084                    eprintln!("[node-app] panic in handle_event");
1085                    $crate::FfiResult::error(-99)
1086                }
1087            }
1088        }
1089
1090        unsafe extern "C" fn __node_app_handle_capability(
1091            request_json: *const u8,
1092            request_len: usize,
1093        ) -> $crate::FfiResult {
1094            match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1095                let json_slice =
1096                    unsafe { std::slice::from_raw_parts(request_json, request_len) };
1097                let request: $crate::CapabilityRequest = match serde_json::from_slice(json_slice) {
1098                    Ok(r) => r,
1099                    Err(e) => {
1100                        eprintln!("[node-app] capability request deserialization error: {}", e);
1101                        return $crate::FfiResult::error(-2);
1102                    }
1103                };
1104
1105                let app = match APP_INSTANCE.get() {
1106                    Some(a) => a,
1107                    None => return $crate::FfiResult::error(-3),
1108                };
1109                // ADR-083 — refused once shutdown starts; declared before the
1110                // read guard so the guard drops first.
1111                let _call = match CALL_GATE.enter() {
1112                    Some(call) => call,
1113                    None => return $crate::FfiResult::error($crate::ERROR_SHUTTING_DOWN),
1114                };
1115                // Feature 463 B5 — read-lock allows the SAME app to handle
1116                // re-entrant capability calls (e.g. agent.prompt outer call
1117                // and an inner core.agent_storage.write from a graph node)
1118                // concurrently. handle_capability is &self so this is sound.
1119                let guard = match app.read() {
1120                    Ok(g) => g,
1121                    Err(e) => {
1122                        eprintln!("[node-app] rwlock poisoned in handle_capability: {}", e);
1123                        return $crate::FfiResult::error(-10);
1124                    }
1125                };
1126
1127                // Set thread-local trace context for duration of this capability call.
1128                // This allows invoke_capability() to propagate trace automatically.
1129                $crate::CURRENT_TRACE.with(|tl| {
1130                    *tl.borrow_mut() = if let Some(ref trace_id) = request.trace_id {
1131                        Some($crate::CurrentTrace {
1132                            trace_id: trace_id.clone(),
1133                            span_id: request.span_id.clone().unwrap_or_default(),
1134                            depth: request.trace_depth.unwrap_or(0),
1135                        })
1136                    } else {
1137                        None
1138                    };
1139                });
1140
1141                let _invocation_context_guard =
1142                    $crate::CurrentInvocationContextGuard::enter(
1143                        request.invocation_context_id.clone(),
1144                    );
1145
1146                let cap_result = guard.handle_capability(request);
1147
1148                // Clear trace context after call completes.
1149                $crate::CURRENT_TRACE.with(|tl| {
1150                    *tl.borrow_mut() = None;
1151                });
1152
1153                match cap_result {
1154                    Ok(response) => match serde_json::to_vec(&response) {
1155                        Ok(bytes) => {
1156                            // Enforce 16MB response limit
1157                            if bytes.len() > $crate::MAX_CAPABILITY_RESPONSE_SIZE {
1158                                eprintln!(
1159                                    "[node-app] capability response exceeds 16MB limit ({} bytes)",
1160                                    bytes.len()
1161                                );
1162                                return $crate::FfiResult::error(-6);
1163                            }
1164                            let len = bytes.len();
1165                            let boxed = bytes.into_boxed_slice();
1166                            let ptr = Box::into_raw(boxed) as *mut u8;
1167                            $crate::FfiResult {
1168                                success: true,
1169                                error_code: 0,
1170                                data: ptr,
1171                                data_len: len,
1172                            }
1173                        }
1174                        Err(e) => {
1175                            eprintln!("[node-app] capability response serialization error: {}", e);
1176                            $crate::FfiResult::error(-4)
1177                        }
1178                    },
1179                    Err(e) => {
1180                        eprintln!("[node-app] handle_capability error: {}", e);
1181                        $crate::FfiResult::error(-5)
1182                    }
1183                }
1184            })) {
1185                Ok(result) => result,
1186                Err(_) => {
1187                    eprintln!("[node-app] panic in handle_capability");
1188                    $crate::FfiResult::error(-99)
1189                }
1190            }
1191        }
1192
1193        unsafe extern "C" fn __node_app_free(ptr: *mut u8, len: usize) {
1194            if !ptr.is_null() && len > 0 {
1195                let _ = unsafe { Box::from_raw(std::slice::from_raw_parts_mut(ptr, len)) };
1196            }
1197        }
1198
1199        #[no_mangle]
1200        pub unsafe extern "C" fn _node_app_entry(
1201            _ctx: *const std::os::raw::c_void,
1202        ) -> *const $crate::NodeAppVTable {
1203            match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1204                let info = <$app_type as $crate::NodeApp>::metadata();
1205                let caps = info.capability_flags();
1206
1207                let name = META_NAME
1208                    .get_or_init(|| std::ffi::CString::new(info.name).unwrap_or_default());
1209                let version = META_VERSION
1210                    .get_or_init(|| std::ffi::CString::new(info.version).unwrap_or_default());
1211                let author = META_AUTHOR
1212                    .get_or_init(|| std::ffi::CString::new(info.author).unwrap_or_default());
1213                let description = META_DESCRIPTION
1214                    .get_or_init(|| std::ffi::CString::new(info.description).unwrap_or_default());
1215
1216                let metadata = $crate::NodeAppMetadata {
1217                    api_version: $crate::API_VERSION,
1218                    name: name.as_ptr(),
1219                    version: version.as_ptr(),
1220                    author: author.as_ptr(),
1221                    description: description.as_ptr(),
1222                    capabilities: caps.bits(),
1223                };
1224
1225                VTABLE.get_or_init(|| $crate::NodeAppVTable {
1226                    metadata,
1227                    init: __node_app_init,
1228                    shutdown: __node_app_shutdown,
1229                    handle_request: __node_app_handle_request,
1230                    handle_event: __node_app_handle_event,
1231                    handle_capability: __node_app_handle_capability,
1232                    free: __node_app_free,
1233                }) as *const $crate::NodeAppVTable
1234            })) {
1235                Ok(ptr) => ptr,
1236                Err(_) => {
1237                    eprintln!("[node-app] panic in _node_app_entry");
1238                    std::ptr::null()
1239                }
1240            }
1241        }
1242    };
1243}
1244
1245#[cfg(test)]
1246mod invocation_context_tests {
1247    use std::ffi::{c_char, c_void};
1248    use std::sync::Mutex;
1249
1250    use super::*;
1251
1252    pub(super) static TEST_LOCK: Mutex<()> = Mutex::new(());
1253
1254    unsafe extern "C" fn noop_log(_: *const c_void, _: u32, _: *const c_char) {}
1255    unsafe extern "C" fn no_config(_: *const c_void, _: *const c_char) -> *const c_char {
1256        std::ptr::null()
1257    }
1258    unsafe extern "C" fn noop_storage(_: *const c_void, _: *const c_char, _: *const c_char) {}
1259    unsafe extern "C" fn no_storage(_: *const c_void, _: *const c_char) -> *const c_char {
1260        std::ptr::null()
1261    }
1262    unsafe extern "C" fn noop_publish(
1263        _: *const c_void,
1264        _: *const u8,
1265        _: usize,
1266        _: *const u8,
1267        _: usize,
1268    ) -> i32 {
1269        0
1270    }
1271    unsafe extern "C" fn capture_invoke(
1272        host_data: *const c_void,
1273        request_json: *const u8,
1274        request_len: usize,
1275    ) -> FfiResult {
1276        let requests = &*(host_data as *const Mutex<Vec<CapabilityRequest>>);
1277        let bytes = std::slice::from_raw_parts(request_json, request_len);
1278        requests
1279            .lock()
1280            .unwrap()
1281            .push(serde_json::from_slice(bytes).unwrap());
1282        let response = CapabilityResponse {
1283            id: "nested-1".into(),
1284            success: true,
1285            payload: serde_json::json!({"ok": true}),
1286        };
1287        let bytes = serde_json::to_vec(&response).unwrap().into_boxed_slice();
1288        let len = bytes.len();
1289        let data = Box::into_raw(bytes) as *mut u8;
1290        FfiResult {
1291            success: true,
1292            error_code: 0,
1293            data,
1294            data_len: len,
1295        }
1296    }
1297
1298    #[test]
1299    fn nested_invocation_propagates_opaque_context() {
1300        let _serial = TEST_LOCK.lock().unwrap();
1301        let captured = Mutex::new(Vec::<CapabilityRequest>::new());
1302        let context = NodeAppContext {
1303            host_data: &captured as *const _ as *const c_void,
1304            host_log: noop_log,
1305            host_get_config: no_config,
1306            host_set_storage: noop_storage,
1307            host_get_storage: no_storage,
1308            host_invoke_capability: capture_invoke,
1309            host_publish_event: noop_publish,
1310        };
1311        __store_context(&context);
1312        CURRENT_TRACE.with(|current| {
1313            *current.borrow_mut() = Some(CurrentTrace {
1314                trace_id: "trace-1".into(),
1315                span_id: "span-1".into(),
1316                depth: 4,
1317            });
1318        });
1319        let _invocation =
1320            CurrentInvocationContextGuard::enter(Some("opaque-context-1".into()));
1321
1322        invoke_capability(&CapabilityRequest {
1323            id: "nested-1".into(),
1324            capability: "example.nested".into(),
1325            payload: serde_json::json!({}),
1326            ..Default::default()
1327        })
1328        .unwrap();
1329
1330        let requests = captured.lock().unwrap();
1331        assert_eq!(requests.len(), 1);
1332        assert_eq!(
1333            requests[0].invocation_context_id.as_deref(),
1334            Some("opaque-context-1")
1335        );
1336        assert_eq!(requests[0].trace_id.as_deref(), Some("trace-1"));
1337        assert_eq!(requests[0].span_id.as_deref(), Some("span-1"));
1338        drop(requests);
1339        CURRENT_TRACE.with(|current| *current.borrow_mut() = None);
1340        __store_context(std::ptr::null());
1341    }
1342}
1343
1344#[cfg(test)]
1345mod shutdown_gate_tests {
1346    use std::sync::atomic::AtomicBool;
1347    use std::sync::{Condvar, Mutex, PoisonError};
1348    use std::time::Duration;
1349
1350    use super::*;
1351
1352    struct Signal {
1353        set: Mutex<bool>,
1354        changed: Condvar,
1355    }
1356
1357    impl Signal {
1358        const fn new() -> Self {
1359            Self {
1360                set: Mutex::new(false),
1361                changed: Condvar::new(),
1362            }
1363        }
1364
1365        fn set(&self) {
1366            *self.set.lock().unwrap() = true;
1367            self.changed.notify_all();
1368        }
1369
1370        fn reset(&self) {
1371            *self.set.lock().unwrap() = false;
1372        }
1373
1374        fn is_set(&self) -> bool {
1375            *self.set.lock().unwrap()
1376        }
1377
1378        fn wait(&self, timeout: Duration) -> bool {
1379            let set = self.set.lock().unwrap();
1380            *self
1381                .changed
1382                .wait_timeout_while(set, timeout, |set| !*set)
1383                .unwrap()
1384                .0
1385        }
1386    }
1387
1388    // Tests that use the app reset these first and hold TEST_LOCK: the macro
1389    // allows one app per binary, so they share it.
1390    static CALL_ENTERED: Signal = Signal::new();
1391    static RELEASE_CALL: Signal = Signal::new();
1392    static HOOK_RAN: Signal = Signal::new();
1393    static SHUTDOWN_RAN: Signal = Signal::new();
1394    static PANIC_IN_HOOK: AtomicBool = AtomicBool::new(false);
1395
1396    #[derive(Default)]
1397    struct SlowApp;
1398
1399    impl NodeApp for SlowApp {
1400        fn metadata() -> NodeAppInfo {
1401            NodeAppInfo {
1402                name: "slow".into(),
1403                version: "0.0.0".into(),
1404                author: "test".into(),
1405                description: "holds a capability call open".into(),
1406                capabilities: Vec::new(),
1407            }
1408        }
1409
1410        fn begin_shutdown(&self) {
1411            HOOK_RAN.set();
1412            if PANIC_IN_HOOK.load(Ordering::SeqCst) {
1413                panic!("begin_shutdown test panic");
1414            }
1415        }
1416
1417        fn shutdown(&mut self) -> Result<(), NodeAppError> {
1418            SHUTDOWN_RAN.set();
1419            Ok(())
1420        }
1421
1422        fn handle_capability(
1423            &self,
1424            request: CapabilityRequest,
1425        ) -> Result<CapabilityResponse, NodeAppError> {
1426            CALL_ENTERED.set();
1427            // Holds the app's read lock until the test lets go. The bound only
1428            // keeps a failing test from hanging.
1429            if !RELEASE_CALL.wait(Duration::from_secs(10)) {
1430                return Err(NodeAppError::CapabilityError("never released".into()));
1431            }
1432            Ok(CapabilityResponse {
1433                id: request.id,
1434                success: true,
1435                payload: serde_json::json!({}),
1436            })
1437        }
1438    }
1439
1440    crate::declare_node_app!(SlowApp);
1441
1442    /// 0 on success, else the wrapper's error code. Frees any response.
1443    fn code(result: FfiResult) -> i32 {
1444        if !result.success {
1445            return result.error_code;
1446        }
1447        if !result.data.is_null() {
1448            unsafe { __node_app_free(result.data, result.data_len) };
1449        }
1450        0
1451    }
1452
1453    fn call() -> i32 {
1454        let request = serde_json::to_vec(&CapabilityRequest {
1455            id: "call-1".into(),
1456            capability: "slow.wait".into(),
1457            payload: serde_json::json!({}),
1458            ..Default::default()
1459        })
1460        .unwrap();
1461        code(unsafe { __node_app_handle_capability(request.as_ptr(), request.len()) })
1462    }
1463
1464    fn request() -> i32 {
1465        let request = serde_json::to_vec(&AppRequest {
1466            id: "request-1".into(),
1467            method: "GET".into(),
1468            path: "/".into(),
1469            headers: Default::default(),
1470            body: serde_json::Value::Null,
1471            caller: None,
1472        })
1473        .unwrap();
1474        code(unsafe { __node_app_handle_request(request.as_ptr(), request.len()) })
1475    }
1476
1477    fn event() -> i32 {
1478        let event = serde_json::to_vec(&AppEvent {
1479            name: "slow.tick".into(),
1480            data: serde_json::Value::Null,
1481        })
1482        .unwrap();
1483        code(unsafe { __node_app_handle_event(event.as_ptr(), event.len()) })
1484    }
1485
1486    /// Stop the app while a capability call holds its read lock.
1487    fn stop_while_a_call_is_in_flight(panic_in_hook: bool) {
1488        for signal in [&CALL_ENTERED, &RELEASE_CALL, &HOOK_RAN, &SHUTDOWN_RAN] {
1489            signal.reset();
1490        }
1491        PANIC_IN_HOOK.store(panic_in_hook, Ordering::SeqCst);
1492        assert!(unsafe { __node_app_init(std::ptr::null()) }.success);
1493
1494        let in_flight = std::thread::spawn(call);
1495        assert!(
1496            CALL_ENTERED.wait(Duration::from_secs(5)),
1497            "call never started"
1498        );
1499
1500        let shutdown = std::thread::spawn(|| unsafe { __node_app_shutdown() }.success);
1501        // The in-flight call cannot return before RELEASE_CALL, so the hook
1502        // runs while that call still holds the read lock.
1503        assert!(
1504            HOOK_RAN.wait(Duration::from_secs(1)),
1505            "begin_shutdown waited behind the in-flight call"
1506        );
1507        assert_eq!(call(), ERROR_SHUTTING_DOWN, "capability call not refused");
1508        assert_eq!(request(), ERROR_SHUTTING_DOWN, "request not refused");
1509        assert_eq!(event(), ERROR_SHUTTING_DOWN, "event not refused");
1510        assert!(
1511            !SHUTDOWN_RAN.is_set(),
1512            "shutdown(&mut self) ran while a call was in flight"
1513        );
1514
1515        RELEASE_CALL.set();
1516        assert_eq!(in_flight.join().unwrap(), 0, "in-flight call failed");
1517        assert!(shutdown.join().unwrap(), "shutdown entry failed");
1518        assert!(SHUTDOWN_RAN.is_set());
1519    }
1520
1521    #[test]
1522    fn shutdown_hook_runs_while_a_call_is_in_flight() {
1523        let _serial = super::invocation_context_tests::TEST_LOCK
1524            .lock()
1525            .unwrap_or_else(PoisonError::into_inner);
1526        stop_while_a_call_is_in_flight(false);
1527
1528        // A cached library is initialized again on reload; it takes calls again.
1529        assert!(unsafe { __node_app_init(std::ptr::null()) }.success);
1530        assert_eq!(call(), 0, "capability call refused after init");
1531        assert_eq!(request(), 0, "request refused after init");
1532        assert_eq!(event(), 0, "event refused after init");
1533    }
1534
1535    #[test]
1536    fn panicking_shutdown_hook_still_waits_for_in_flight_calls() {
1537        let _serial = super::invocation_context_tests::TEST_LOCK
1538            .lock()
1539            .unwrap_or_else(PoisonError::into_inner);
1540        stop_while_a_call_is_in_flight(true);
1541    }
1542
1543    #[test]
1544    fn call_gate_counts_each_call_until_its_permit_drops() {
1545        let gate = CallGate::new();
1546        let first = gate.enter().expect("open gate admits");
1547        let second = gate.enter().expect("open gate admits");
1548        gate.close();
1549        assert!(gate.enter().is_none(), "closed gate admitted a call");
1550        assert_eq!(gate.wait_drained(Some(Duration::from_millis(10))), 2);
1551        drop(first);
1552        assert_eq!(gate.wait_drained(Some(Duration::from_millis(10))), 1);
1553        drop(second);
1554        assert_eq!(gate.wait_drained(None), 0);
1555        gate.open();
1556        assert!(gate.enter().is_some(), "reopened gate refused a call");
1557    }
1558}