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}