Skip to main content

rings_node/extension/ext/
registry.rs

1//! Router + capability core.
2//!
3//! [`Extensions`] registers `(Protocol, Interpret)` pairs by namespace. Each interpreter is
4//! handed a namespace-scoped [`Scope`] (overlay `send` / `did` / self-`inject`, confined to its
5//! own namespace) — *not* the router-internal `Core`. Protocols that authenticate and select an
6//! exact physical next hop may additionally use crate-private `send_direct`; it preserves the
7//! namespace envelope while bypassing Chord route selection. `Core` is the crate-private capability
8//! that also routes an inbound [`Envelope`] to its protocol and drives the bounded re-injection
9//! fixpoint; the registry stays uniform (everything erased to the internal `Handler`) while each
10//! extension's shell is its own. Extension authors see only `Protocol` / `Interpret` / `Scope` /
11//! `Transition` / `Extensions`.
12
13use std::collections::HashMap;
14use std::collections::VecDeque;
15use std::ops::Deref;
16use std::sync::Arc;
17use std::sync::Mutex;
18use std::sync::RwLock;
19
20use bytes::Bytes;
21use futures::lock::Mutex as AsyncMutex;
22use rings_core::dht::Did;
23
24use super::Ctx;
25use super::Envelope;
26use super::Inbound;
27use super::Interpret;
28use super::MaybeSend;
29use super::Protocol;
30use super::Reject;
31use super::Transition;
32use super::Wire;
33use crate::error::Error;
34use crate::error::Result;
35use crate::processor::Processor;
36use crate::sync_lock::lock;
37
38/// Upper bound on re-injection iterations per inbound message, so a misbehaving
39/// protocol/effect cycle cannot diverge.
40const MAX_FIXPOINT_STEPS: u32 = 1024;
41
42/// Type-erased handler stored in the registry: native is `Send + Sync`, browser not.
43#[cfg(rings_native)]
44pub(crate) type DynHandler = dyn Handler + Send + Sync;
45/// Type-erased handler stored in the registry.
46#[cfg(rings_browser)]
47pub(crate) type DynHandler = dyn Handler;
48
49type HandlerMap = RwLock<HashMap<String, Arc<DynHandler>>>;
50
51/// Erased, runtime-facing handler — the router-internal ABI. Implemented once, generically, by
52/// `Runner`; protocol authors never name it (they write `Protocol` + `Interpret`).
53#[cfg_attr(rings_browser, async_trait::async_trait(?Send))]
54#[cfg_attr(rings_native, async_trait::async_trait)]
55pub(crate) trait Handler {
56    /// Decode → step (pure, committed) → run the protocol's effects, returning re-injected
57    /// messages. `handle : (from, payload) → IO [Inbound]`.
58    async fn handle(&self, core: &Core, from: Did, payload: Bytes) -> Result<Vec<Inbound>>;
59}
60
61/// The router-internal capability. Cloneable and `'static` so a long-running engine task can
62/// keep a copy and feed events back via [`inject`](Core::inject). Not handed to extension
63/// shells — they get a namespace-scoped [`Scope`] instead.
64#[derive(Clone)]
65pub(crate) struct Core {
66    processor: Arc<Processor>,
67    handlers: Arc<HandlerMap>,
68}
69
70impl Core {
71    /// This node's DID.
72    pub fn did(&self) -> Did {
73        self.processor.did()
74    }
75
76    /// Put a message on the overlay to `to` under `namespace`.
77    pub async fn send(&self, to: Did, namespace: &str, payload: Bytes) -> Result<()> {
78        let envelope = Envelope::new(namespace, payload);
79        self.processor.send_envelope(to, &envelope).await?;
80        Ok(())
81    }
82
83    /// Put a message on an already-connected transport edge without Chord route selection.
84    async fn send_direct(&self, to: Did, namespace: &str, payload: Bytes) -> Result<()> {
85        let envelope = Envelope::new(namespace, payload);
86        self.processor.send_direct_envelope(to, &envelope).await?;
87        Ok(())
88    }
89
90    /// Re-enter the router with a *self*-addressed message (`from = this node`): a locally
91    /// injected command, or an engine task feeding a lifecycle event back to its protocol.
92    pub async fn inject(&self, namespace: &str, payload: Bytes) -> Result<()> {
93        self.dispatch(self.did(), Envelope::new(namespace, payload))
94            .await
95    }
96
97    /// Route an inbound [`Envelope`] to its protocol and drive the bounded re-injection
98    /// fixpoint. Unknown namespaces are logged and dropped (non-fatal).
99    ///
100    /// This is the **authenticated ingress** capability: the caller chooses `from`, so a
101    /// protocol's `decode` will attribute the resulting event to that DID (for the relay, a
102    /// `from != me` envelope becomes a peer `Frame`). It is therefore `pub(crate)` — only the
103    /// router path may call it, and only [`Backend`](crate::extension::Backend) does, with
104    /// `from` taken from the message's verified signer. Extension code reaches the router only
105    /// through [`inject`](Core::inject) (self-addressed, `from = self.did()`); it can never
106    /// forge a remote `from`.
107    pub(crate) async fn dispatch(&self, from: Did, envelope: Envelope) -> Result<()> {
108        let mut queue: VecDeque<Inbound> = VecDeque::new();
109        queue.push_back(Inbound {
110            namespace: envelope.namespace,
111            from,
112            payload: envelope.payload,
113        });
114
115        let mut budget = MAX_FIXPOINT_STEPS;
116        while let Some(Inbound {
117            namespace,
118            from,
119            payload,
120        }) = queue.pop_front()
121        {
122            if budget == 0 {
123                return Err(Error::ExtensionError(format!(
124                    "fixpoint budget ({MAX_FIXPOINT_STEPS}) exhausted; last namespace {namespace:?}"
125                )));
126            }
127            budget -= 1;
128
129            match self.handler(namespace.as_str()) {
130                Some(handler) => queue.extend(handler.handle(self, from, payload).await?),
131                None => tracing::debug!(
132                    "no protocol registered for namespace {:?}, dropping",
133                    namespace
134                ),
135            }
136        }
137        Ok(())
138    }
139
140    fn handler(&self, namespace: &str) -> Option<Arc<DynHandler>> {
141        self.handlers.read().ok()?.get(namespace).map(Arc::clone)
142    }
143}
144
145/// A **namespace-scoped** capability handed to an [`Interpret`] shell — the effectful
146/// counterpart of the pure side's read-only [`Ctx`]. Every action is confined to the
147/// interpreter's own namespace: it may `send` to peers and self-`inject` only there, and it
148/// can neither reach another namespace nor forge a remote `from`. This is what keeps the
149/// capability honest: an extension shell cannot use the router as a generic
150/// inject-any-namespace bus (e.g. manufacture another extension's lifecycle events). Cloneable
151/// and `'static`, so a long-running engine task can keep a copy.
152#[derive(Clone)]
153pub struct Scope {
154    core: Core,
155    namespace: String,
156}
157
158impl Scope {
159    /// Confine `core` to `namespace`.
160    pub(crate) fn new(core: Core, namespace: String) -> Self {
161        Self { core, namespace }
162    }
163
164    /// This node's DID.
165    pub fn did(&self) -> Did {
166        self.core.did()
167    }
168
169    /// The namespace this scope is confined to.
170    pub fn namespace(&self) -> &str {
171        self.namespace.as_str()
172    }
173
174    /// Put a message on the overlay to `to`, under this interpreter's own namespace.
175    pub async fn send(&self, to: Did, payload: Bytes) -> Result<()> {
176        self.core.send(to, self.namespace.as_str(), payload).await
177    }
178
179    /// Send to an exact, already-connected peer under this scope's namespace.
180    ///
181    /// This is intentionally crate-private: only protocols that already authenticate their own
182    /// hop selection, such as onion circuits, may bypass the overlay routing decision.
183    pub(crate) async fn send_direct(&self, to: Did, payload: Bytes) -> Result<()> {
184        self.core
185            .send_direct(to, self.namespace.as_str(), payload)
186            .await
187    }
188
189    /// Self-inject `payload` into this interpreter's **own** namespace (`from = this node`).
190    ///
191    /// `pub(crate)`: this is the **long-lived lifecycle sink** for an extension's own engine
192    /// (e.g. the relay's spawned socket tasks reporting `Accepted`/`Untrack` later), and it
193    /// starts a **fresh** [`dispatch`](Core::dispatch) fixpoint with its own
194    /// `MAX_FIXPOINT_STEPS` budget — it is *not* part of the bounded feedback fixpoint that
195    /// drives a single inbound. The synchronous per-effect feedback path is the `Vec<Bytes>`
196    /// returned from [`Interpret::run`], which the runner reduces within its current ordered
197    /// turn and budget.
198    /// A third-party shell therefore gets only that bounded return path, never this re-entrant
199    /// sink, so it cannot recurse `inject` to escape the budget.
200    pub(crate) async fn inject(&self, payload: Bytes) -> Result<()> {
201        self.core.inject(self.namespace.as_str(), payload).await
202    }
203}
204
205/// Capability available while an interpreter applies one committed effect.
206///
207/// It can send overlay messages and return synchronous feedback, but cannot re-enter its own
208/// reducer. Long-lived engines receive a [`Scope`] explicitly through the crate-private
209/// `EffectScope::lifecycle` handoff, making ownership visible at the effect boundary.
210pub struct EffectScope {
211    scope: Scope,
212}
213
214impl EffectScope {
215    pub(crate) fn new(scope: Scope) -> Self {
216        Self { scope }
217    }
218
219    /// This node's DID.
220    pub fn did(&self) -> Did {
221        self.scope.did()
222    }
223
224    /// The namespace this effect is confined to.
225    pub fn namespace(&self) -> &str {
226        self.scope.namespace()
227    }
228
229    /// Put a message on the overlay under this effect's namespace.
230    pub async fn send(&self, to: Did, payload: Bytes) -> Result<()> {
231        self.scope.send(to, payload).await
232    }
233
234    /// Hand the lifecycle capability to an explicitly long-lived engine task.
235    ///
236    /// The interpreter itself must return same-turn feedback from [`Interpret::run`], rather
237    /// than await [`Scope::inject`] while the ordered effect turn is active.
238    pub(crate) fn lifecycle(&self) -> Scope {
239        self.scope.clone()
240    }
241}
242
243/// Adapter binding a pure [`Protocol`] to its [`Interpret`] shell and owned state; erased to
244/// [`Handler`]. Protocol authors never write this.
245struct Runner<P: Protocol, I> {
246    protocol: P,
247    interpret: I,
248    state: Mutex<P::State>,
249    transition_gate: AsyncMutex<()>,
250    #[cfg(all(test, rings_native))]
251    after_decode_for_test: Option<Arc<dyn Fn() + Send + Sync>>,
252    #[cfg(all(test, rings_native))]
253    after_commit_for_test: Option<Arc<dyn Fn() + Send + Sync>>,
254    #[cfg(all(test, rings_native))]
255    before_gate_wait_for_test: Option<Arc<dyn Fn(bool) + Send + Sync>>,
256}
257
258#[cfg_attr(rings_browser, async_trait::async_trait(?Send))]
259#[cfg_attr(rings_native, async_trait::async_trait)]
260impl<P, I> Handler for Runner<P, I>
261where
262    P: Protocol + MaybeSend + 'static,
263    P::State: MaybeSend + 'static,
264    P::Effect: MaybeSend,
265    I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
266{
267    async fn handle(&self, core: &Core, from: Did, payload: Bytes) -> Result<Vec<Inbound>> {
268        // Boundary: decode raw bytes to a typed event. An undecodable/foreign message is an
269        // explicit drop here, not a silent `Transition::pure` deep in `step`.
270        let event = match self.protocol.decode(Wire {
271            from,
272            me: core.did(),
273            payload: payload.as_ref(),
274        }) {
275            Ok(event) => event,
276            Err(Reject(why)) => {
277                tracing::debug!("drop on {}: {why}", self.protocol.namespace());
278                return Ok(Vec::new());
279            }
280        };
281
282        #[cfg(all(test, rings_native))]
283        if let Some(observe) = self.after_decode_for_test.as_ref() {
284            observe();
285        }
286
287        // The transition gate establishes the protocol's linearization order for state commit
288        // and its resulting effect trace. Decode remains outside that order: it has no state or
289        // effects by contract. Holding the gate through interpretation preserves:
290        // commit(A) < commit(B) => applying A's effects ends before applying B's effects begins.
291        // The state mutex itself remains synchronous and never crosses an await.
292        #[cfg(all(test, rings_native))]
293        if let Some(observe) = self.before_gate_wait_for_test.as_ref() {
294            // Witness the real synchronization boundary: false means this task could acquire
295            // the gate immediately; true means another transition owns it at this exact point.
296            observe(self.transition_gate.try_lock().is_none());
297        }
298        let _transition_turn = self.transition_gate.lock().await;
299
300        // Impure region: the state lock is released, while the transition gate keeps this
301        // protocol's effect trace and synchronous feedback fixpoint in commit order. Returned
302        // payloads are reduced before this turn releases the gate, so a later inbound cannot
303        // observe state that predates an effect's own feedback.
304        let namespace = self.protocol.namespace().to_string();
305        let scope = EffectScope::new(Scope::new(core.clone(), namespace.clone()));
306        let mut feedback = VecDeque::new();
307        feedback.push_back(event);
308        let mut feedback_budget = MAX_FIXPOINT_STEPS;
309        while let Some(event) = feedback.pop_front() {
310            if feedback_budget == 0 {
311                return Err(Error::ExtensionError(format!(
312                    "feedback fixpoint budget ({MAX_FIXPOINT_STEPS}) exhausted on {namespace:?}"
313                )));
314            }
315            feedback_budget -= 1;
316
317            // Pure region: a brief synchronous state fold. No await crosses `state`; the
318            // feedback queue makes every same-turn effect result part of this fold before a
319            // competing inbound can begin its own transition.
320            let effects = {
321                let mut guard = lock(&self.state)?;
322                let Transition { state, effects } = self.protocol.step(
323                    Ctx {
324                        did: core.did(),
325                        state: guard.deref(),
326                    },
327                    event,
328                );
329                *guard = state;
330                effects
331            };
332
333            #[cfg(all(test, rings_native))]
334            if let Some(observe) = self.after_commit_for_test.as_ref() {
335                observe();
336            }
337
338            for effect in effects {
339                for payload in self.interpret.run(&scope, effect).await? {
340                    match self.protocol.decode(Wire {
341                        from: core.did(),
342                        me: core.did(),
343                        payload: payload.as_ref(),
344                    }) {
345                        Ok(event) => feedback.push_back(event),
346                        Err(Reject(why)) => {
347                            tracing::debug!("drop feedback on {}: {why}", self.protocol.namespace())
348                        }
349                    }
350                }
351            }
352        }
353        Ok(Vec::new())
354    }
355}
356
357/// Registry of `(Protocol, Interpret)` pairs by namespace, plus the router-internal `Core`.
358/// Cheaply cloneable and shared (interior mutability) so the
359/// [`Provider`](crate::provider::Provider) and the inbound callback see the same table.
360#[derive(Clone)]
361pub struct Extensions {
362    core: Core,
363}
364
365impl Extensions {
366    /// Empty registry over a processor (the source of overlay `send` / `did`).
367    pub fn new(processor: Arc<Processor>) -> Self {
368        Self {
369            core: Core {
370                processor,
371                handlers: Arc::new(RwLock::new(HashMap::new())),
372            },
373        }
374    }
375
376    /// The capability handle (overlay `send` / `did` / self-addressed `inject`). `pub(crate)`:
377    /// public holders of an `Extensions` get **registration only** (`register` / `replace` /
378    /// `contains` / `register_many`), never a raw [`Core`]. An extension's local injection is
379    /// exposed through its own typed handle (e.g. `RelayHandle`), so application code cannot use
380    /// a generic inject-any-namespace bus to forge engine-lifecycle commands like the relay's
381    /// `Accepted` / `Untrack`.
382    pub(crate) fn core(&self) -> Core {
383        self.core.clone()
384    }
385
386    /// Register a protocol together with its interpreter under the protocol's namespace.
387    /// Errors if the namespace is already taken — use [`replace`](Extensions::replace) for
388    /// intentional replacement (no more silent overwrite).
389    pub fn register<P, I>(&self, protocol: P, interpret: I) -> Result<()>
390    where
391        P: Protocol + MaybeSend + 'static,
392        P::State: MaybeSend + 'static,
393        P::Effect: MaybeSend,
394        I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
395    {
396        self.insert(protocol, interpret, false)
397    }
398
399    /// Like [`register`](Extensions::register) but replaces an existing protocol on the same
400    /// namespace instead of erroring. For deliberate hot-swaps.
401    pub fn replace<P, I>(&self, protocol: P, interpret: I) -> Result<()>
402    where
403        P: Protocol + MaybeSend + 'static,
404        P::State: MaybeSend + 'static,
405        P::Effect: MaybeSend,
406        I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
407    {
408        self.insert(protocol, interpret, true)
409    }
410
411    /// Register several protocols **atomically**: build every runner, then under a single write
412    /// lock verify that none of their namespaces is taken (by an existing registration or by a
413    /// duplicate within the batch) and insert them all — or change nothing and return `Err`. The
414    /// pairs share the type `P`/`I` (e.g. the relay's TCP + UDP `Relay<T>` instances), so a
415    /// partial install can never leave one namespace claimed while the caller gets no handle.
416    pub fn register_many<P, I>(&self, items: Vec<(P, I)>) -> Result<()>
417    where
418        P: Protocol + MaybeSend + 'static,
419        P::State: MaybeSend + 'static,
420        P::Effect: MaybeSend,
421        I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
422    {
423        // Build (namespace, runner) outside the lock.
424        let prepared: Vec<(String, Arc<DynHandler>, Vec<&'static str>)> = items
425            .into_iter()
426            .map(|(protocol, interpret)| {
427                let capabilities = protocol.capabilities().to_vec();
428                let namespace = protocol.namespace().to_string();
429                let state = Mutex::new(protocol.init());
430                let runner: Arc<DynHandler> = Arc::new(Runner {
431                    protocol,
432                    interpret,
433                    state,
434                    transition_gate: AsyncMutex::new(()),
435                    #[cfg(all(test, rings_native))]
436                    after_decode_for_test: None,
437                    #[cfg(all(test, rings_native))]
438                    after_commit_for_test: None,
439                    #[cfg(all(test, rings_native))]
440                    before_gate_wait_for_test: None,
441                });
442                (namespace, runner, capabilities)
443            })
444            .collect();
445
446        let mut handlers = self.core.handlers.write().map_err(|_| Error::Lock)?;
447        // Check-all (existing table + intra-batch duplicates) before mutating anything.
448        for (index, (namespace, _, _)) in prepared.iter().enumerate() {
449            let duplicate_in_batch = prepared
450                .iter()
451                .take(index)
452                .any(|(seen, _, _)| seen == namespace);
453            if duplicate_in_batch || handlers.contains_key(namespace) {
454                return Err(Error::ExtensionError(format!(
455                    "namespace {namespace:?} is already registered"
456                )));
457            }
458        }
459        self.core.processor.add_online_node_capabilities(
460            prepared
461                .iter()
462                .flat_map(|(_, _, capabilities)| capabilities.iter().copied()),
463        )?;
464        // All free: insert the whole batch.
465        for (namespace, runner, _) in prepared {
466            handlers.insert(namespace, runner);
467        }
468        Ok(())
469    }
470
471    fn insert<P, I>(&self, protocol: P, interpret: I, replace: bool) -> Result<()>
472    where
473        P: Protocol + MaybeSend + 'static,
474        P::State: MaybeSend + 'static,
475        P::Effect: MaybeSend,
476        I: Interpret<Effect = P::Effect> + MaybeSend + 'static,
477    {
478        let capabilities = protocol.capabilities();
479        let namespace = protocol.namespace().to_string();
480        let state = Mutex::new(protocol.init());
481        let runner: Arc<DynHandler> = Arc::new(Runner {
482            protocol,
483            interpret,
484            state,
485            transition_gate: AsyncMutex::new(()),
486            #[cfg(all(test, rings_native))]
487            after_decode_for_test: None,
488            #[cfg(all(test, rings_native))]
489            after_commit_for_test: None,
490            #[cfg(all(test, rings_native))]
491            before_gate_wait_for_test: None,
492        });
493        let mut handlers = self.core.handlers.write().map_err(|_| Error::Lock)?;
494        if !replace && handlers.contains_key(&namespace) {
495            return Err(Error::ExtensionError(format!(
496                "namespace {namespace:?} is already registered"
497            )));
498        }
499        self.core
500            .processor
501            .add_online_node_capabilities(capabilities.iter().copied())?;
502        handlers.insert(namespace, runner);
503        Ok(())
504    }
505
506    /// Whether a namespace is registered.
507    pub fn contains(&self, namespace: &str) -> bool {
508        self.core
509            .handlers
510            .read()
511            .map(|h| h.contains_key(namespace))
512            .unwrap_or(false)
513    }
514
515    /// Route a decoded envelope (inbound entry point). `pub(crate)`: the authenticated ingress
516    /// belongs to the router path ([`Backend`](crate::extension::Backend)), not to public
517    /// holders of an `Extensions` value (which get registration + a self-addressed
518    /// [`Core`](Core::inject), never the ability to forge a remote `from`). See
519    /// [`Core::dispatch`].
520    pub(crate) async fn dispatch(&self, from: Did, envelope: Envelope) -> Result<()> {
521        self.core.dispatch(from, envelope).await
522    }
523}
524
525#[cfg(all(test, rings_native))]
526mod tests {
527    use std::collections::HashMap;
528    use std::net::SocketAddr;
529    use std::sync::Arc;
530    use std::sync::Mutex;
531
532    use async_trait::async_trait;
533    use rings_core::ecc::SecretKey;
534    use rings_core::session::SessionSk;
535    use tokio::sync::Notify;
536
537    use super::*;
538    use crate::extension::protocols::relay::ControlSendTestHook;
539    use crate::extension::protocols::relay::NativeRelay;
540    use crate::extension::protocols::relay::Relay;
541    use crate::extension::protocols::relay::RelayCommand;
542    use crate::extension::protocols::relay::RelayEffect;
543    use crate::extension::protocols::relay::TCP;
544    use crate::extension::transport::engine::TransportSessions;
545    use crate::extension::transport::Frame;
546    use crate::extension::transport::Initiator;
547    use crate::extension::transport::SessionId;
548    use crate::extension::transport::SessionKey;
549    use crate::processor::ProcessorBuilder;
550    use crate::processor::ProcessorConfig;
551
552    struct OrderedProtocol;
553
554    impl Protocol for OrderedProtocol {
555        type State = u8;
556        type Event = u8;
557        type Effect = u8;
558
559        fn namespace(&self) -> &str {
560            "ordered-effects"
561        }
562
563        fn init(&self) -> Self::State {
564            0
565        }
566
567        fn decode(&self, wire: Wire<'_>) -> std::result::Result<Self::Event, Reject> {
568            let event = wire
569                .payload
570                .first()
571                .copied()
572                .ok_or_else(|| Reject("missing effect value".to_string()))?;
573            Ok(event)
574        }
575
576        fn step(
577            &self,
578            ctx: Ctx<'_, Self::State>,
579            event: Self::Event,
580        ) -> Transition<Self::State, Self::Effect> {
581            Transition::with(ctx.state.saturating_add(1), vec![event])
582        }
583    }
584
585    #[derive(Default)]
586    struct BlockingOrderedInterpreter {
587        first_effect_started: Notify,
588        release_first_effect: Notify,
589        observed: Mutex<Vec<u8>>,
590    }
591
592    #[async_trait]
593    impl Interpret for Arc<BlockingOrderedInterpreter> {
594        type Effect = u8;
595
596        async fn run(&self, _scope: &EffectScope, effect: Self::Effect) -> Result<Vec<Bytes>> {
597            if effect == 1 {
598                self.first_effect_started.notify_one();
599                self.release_first_effect.notified().await;
600            }
601            lock(&self.observed)?.push(effect);
602            Ok(Vec::new())
603        }
604    }
605
606    #[derive(Default)]
607    struct RelayFeedbackInterpreter {
608        first_effect_started: Notify,
609        release_first_effect: Notify,
610        first_connect_seen: Mutex<bool>,
611        observed_connects: Mutex<Vec<SessionId>>,
612    }
613
614    #[async_trait]
615    impl Interpret for Arc<RelayFeedbackInterpreter> {
616        type Effect = RelayEffect<SocketAddr>;
617
618        async fn run(&self, _scope: &EffectScope, effect: Self::Effect) -> Result<Vec<Bytes>> {
619            match effect {
620                RelayEffect::Connect { key, .. } => {
621                    let first_connect = {
622                        let mut seen = lock(&self.first_connect_seen)?;
623                        let first_connect = !*seen;
624                        *seen = true;
625                        first_connect
626                    };
627                    if first_connect {
628                        self.first_effect_started.notify_one();
629                        self.release_first_effect.notified().await;
630                        let feedback = RelayCommand::<SocketAddr>::Untrack {
631                            peer: key.peer,
632                            session: key.session,
633                            initiator: key.initiator,
634                        };
635                        return rings_codec::serialize(&feedback)
636                            .map(Bytes::from)
637                            .map(|payload| vec![payload])
638                            .map_err(|_| Error::EncodeError);
639                    }
640                    lock(&self.observed_connects)?.push(key.session);
641                    Ok(Vec::new())
642                }
643                _ => Ok(Vec::new()),
644            }
645        }
646    }
647
648    #[derive(Default)]
649    struct FailingOrderedInterpreter {
650        observed: Mutex<Vec<u8>>,
651    }
652
653    #[async_trait]
654    impl Interpret for Arc<FailingOrderedInterpreter> {
655        type Effect = u8;
656
657        async fn run(&self, _scope: &EffectScope, effect: Self::Effect) -> Result<Vec<Bytes>> {
658            lock(&self.observed)?.push(effect);
659            if effect == 1 {
660                return Err(Error::ExtensionError(
661                    "intentional effect failure".to_string(),
662                ));
663            }
664            Ok(Vec::new())
665        }
666    }
667
668    fn extensions() -> Result<Extensions> {
669        let session = SessionSk::new_with_seckey(&SecretKey::random())?;
670        let config = ProcessorConfig::new(1, String::new(), session, 1);
671        let processor = ProcessorBuilder::from_config(&config)?
672            .advertise_presence(false)
673            .build()?;
674        Ok(Extensions::new(Arc::new(processor)))
675    }
676
677    #[tokio::test]
678    async fn test_unknown_legacy_namespace_is_a_nonfatal_drop() -> Result<()> {
679        let extensions = extensions()?;
680        let from = extensions.core().did();
681
682        extensions
683            .dispatch(
684                from,
685                Envelope::new("snark", Bytes::from_static(b"legacy-task")),
686            )
687            .await?;
688
689        assert!(extensions.core.handler("snark").is_none());
690        Ok(())
691    }
692
693    #[tokio::test]
694    async fn test_committed_transitions_execute_effects_in_commit_order() -> Result<()> {
695        // Invariant: while A's effect is blocked, B cannot commit or emit; releasing A
696        // produces the unique effect trace [A, B] for the protocol's state-transition order.
697        let extensions = extensions()?;
698        let interpreter = Arc::new(BlockingOrderedInterpreter::default());
699        let gate_wait = Arc::new(Notify::new());
700        let gate_contention = Arc::new(Mutex::new(Vec::new()));
701        let gate_observer = {
702            let gate_wait = Arc::clone(&gate_wait);
703            let gate_contention = Arc::clone(&gate_contention);
704            Arc::new(move |contended| {
705                gate_contention
706                    .lock()
707                    .expect("test gate witness lock")
708                    .push(contended);
709                gate_wait.notify_one();
710            }) as Arc<dyn Fn(bool) + Send + Sync>
711        };
712        let committed = Arc::new(Mutex::new(0_u8));
713        let commit_observer = {
714            let committed = Arc::clone(&committed);
715            Arc::new(move || {
716                *committed.lock().expect("test commit witness lock") += 1;
717            }) as Arc<dyn Fn() + Send + Sync>
718        };
719        let runner: Arc<DynHandler> = Arc::new(Runner {
720            protocol: OrderedProtocol,
721            interpret: Arc::clone(&interpreter),
722            state: Mutex::new(0),
723            transition_gate: AsyncMutex::new(()),
724            after_decode_for_test: None,
725            after_commit_for_test: Some(commit_observer),
726            before_gate_wait_for_test: Some(gate_observer),
727        });
728        extensions
729            .core
730            .handlers
731            .write()
732            .map_err(|_| Error::Lock)?
733            .insert("ordered-effects".to_string(), runner);
734        let from = extensions.core().did();
735
736        let first_extensions = extensions.clone();
737        let first = tokio::spawn(async move {
738            first_extensions
739                .dispatch(
740                    from,
741                    Envelope::new("ordered-effects", Bytes::from_static(&[1])),
742                )
743                .await
744        });
745        interpreter.first_effect_started.notified().await;
746        gate_wait.notified().await;
747
748        let second_extensions = extensions.clone();
749        let second = tokio::spawn(async move {
750            second_extensions
751                .dispatch(
752                    from,
753                    Envelope::new("ordered-effects", Bytes::from_static(&[2])),
754                )
755                .await
756        });
757        gate_wait.notified().await;
758        assert_eq!(*lock(&gate_contention)?, vec![false, true]);
759        assert!(!second.is_finished());
760        assert!(lock(&interpreter.observed)?.is_empty());
761        assert_eq!(*lock(&committed)?, 1);
762
763        interpreter.release_first_effect.notify_one();
764        first
765            .await
766            .map_err(|error| Error::ExtensionError(error.to_string()))??;
767        second
768            .await
769            .map_err(|error| Error::ExtensionError(error.to_string()))??;
770        assert_eq!(*lock(&committed)?, 2);
771        assert_eq!(*lock(&interpreter.observed)?, vec![1, 2]);
772        Ok(())
773    }
774
775    #[tokio::test]
776    async fn test_failed_effect_releases_ordered_turn_for_later_transition() -> Result<()> {
777        // Law: failure is an outcome of the committed transition, not a leaked gate. The next
778        // transition can run after the failed application has ended.
779        let extensions = extensions()?;
780        let interpreter = Arc::new(FailingOrderedInterpreter::default());
781        extensions.register(OrderedProtocol, Arc::clone(&interpreter))?;
782        let from = extensions.core().did();
783
784        let failed = extensions
785            .dispatch(
786                from,
787                Envelope::new("ordered-effects", Bytes::from_static(&[1])),
788            )
789            .await;
790        assert!(matches!(failed, Err(Error::ExtensionError(_))));
791        extensions
792            .dispatch(
793                from,
794                Envelope::new("ordered-effects", Bytes::from_static(&[2])),
795            )
796            .await?;
797
798        assert_eq!(*lock(&interpreter.observed)?, vec![1, 2]);
799        Ok(())
800    }
801
802    #[tokio::test]
803    async fn test_returned_feedback_precedes_a_waiting_transition() -> Result<()> {
804        // Invariant: the first inbound Open returns a real RelayCommand::Untrack feedback before
805        // a duplicate Open may inspect relay state. Therefore the duplicate is accepted and emits
806        // its own Connect; the old queue-after-gate implementation dropped it as a live duplicate.
807        let extensions = extensions()?;
808        let interpreter = Arc::new(RelayFeedbackInterpreter::default());
809        let decoded = Arc::new(Notify::new());
810        let observer = {
811            let decoded = Arc::clone(&decoded);
812            Arc::new(move || decoded.notify_one()) as Arc<dyn Fn() + Send + Sync>
813        };
814        let protocol = Relay::tcp(HashMap::from([(
815            "web".to_string(),
816            "127.0.0.1:80"
817                .parse::<SocketAddr>()
818                .map_err(|error| Error::ExtensionError(error.to_string()))?,
819        )]));
820        let state = protocol.init();
821        let runner: Arc<DynHandler> = Arc::new(Runner {
822            protocol,
823            interpret: Arc::clone(&interpreter),
824            state: Mutex::new(state),
825            transition_gate: AsyncMutex::new(()),
826            after_decode_for_test: Some(observer),
827            after_commit_for_test: None,
828            before_gate_wait_for_test: None,
829        });
830        extensions
831            .core
832            .handlers
833            .write()
834            .map_err(|_| Error::Lock)?
835            .insert(TCP.to_string(), runner);
836        let from: Did = SecretKey::random().address().into();
837        let open = rings_codec::serialize(&Frame::Open {
838            session: SessionId(0),
839            service: "web".to_string(),
840        })
841        .map(Bytes::from)
842        .map_err(|_| Error::EncodeError)?;
843        let first_open = open.clone();
844
845        let first_extensions = extensions.clone();
846        let first = tokio::spawn(async move {
847            first_extensions
848                .dispatch(from, Envelope::new(TCP, first_open))
849                .await
850        });
851        interpreter.first_effect_started.notified().await;
852        decoded.notified().await;
853
854        let second_extensions = extensions.clone();
855        let second = tokio::spawn(async move {
856            second_extensions
857                .dispatch(from, Envelope::new(TCP, open))
858                .await
859        });
860        decoded.notified().await;
861        assert!(!second.is_finished());
862
863        interpreter.release_first_effect.notify_one();
864        let timeout = std::time::Duration::from_secs(1);
865        tokio::time::timeout(timeout, first)
866            .await
867            .map_err(|_| Error::ExtensionError("first feedback turn timed out".to_string()))?
868            .map_err(|error| Error::ExtensionError(error.to_string()))??;
869        tokio::time::timeout(timeout, second)
870            .await
871            .map_err(|_| Error::ExtensionError("second feedback turn timed out".to_string()))?
872            .map_err(|error| Error::ExtensionError(error.to_string()))??;
873        assert_eq!(*lock(&interpreter.observed_connects)?, vec![SessionId(0)]);
874        Ok(())
875    }
876
877    #[tokio::test]
878    async fn test_missing_open_accepted_resource_returns_synchronous_untrack() -> Result<()> {
879        let extensions = extensions()?;
880        let effect_scope = EffectScope::new(Scope::new(extensions.core(), TCP.to_string()));
881        let interpreter = NativeRelay::new(Arc::new(TransportSessions::new()));
882        let peer: Did = SecretKey::random().address().into();
883        let key = SessionKey::new(peer, TCP, SessionId(9), Initiator::Local);
884
885        let feedback = interpreter
886            .run(&effect_scope, RelayEffect::OpenAccepted {
887                token: 77,
888                key: key.clone(),
889                service: "missing-pending-resource".to_string(),
890            })
891            .await?;
892
893        assert_eq!(feedback.len(), 1);
894        assert!(matches!(
895            rings_codec::deserialize::<RelayCommand<SocketAddr>>(feedback[0].as_ref()),
896            Ok(RelayCommand::Untrack {
897                peer: actual_peer,
898                session: SessionId(9),
899                initiator: Initiator::Local,
900            }) if actual_peer == peer
901        ));
902        Ok(())
903    }
904
905    #[tokio::test]
906    async fn test_terminal_relay_control_effect_does_not_await_overlay_send() -> Result<()> {
907        let extensions = extensions()?;
908        let hook = Arc::new(ControlSendTestHook::default());
909        let interpreter = Arc::new(NativeRelay::new_with_control_send_test_hook(
910            Arc::new(TransportSessions::new()),
911            Arc::clone(&hook),
912        ));
913        let peer: Did = SecretKey::random().address().into();
914        let core = extensions.core();
915        let application = tokio::spawn(async move {
916            let effect_scope = EffectScope::new(Scope::new(core, TCP.to_string()));
917            interpreter
918                .run(&effect_scope, RelayEffect::SendClose {
919                    to: peer,
920                    session: SessionId(5),
921                    from_opener: false,
922                })
923                .await
924        });
925
926        // Await a real outbox-worker suspension before observing completion. The interpreter
927        // must already have returned; otherwise this join times out deterministically while the
928        // hook remains held.
929        tokio::time::timeout(std::time::Duration::from_secs(1), hook.wait_until_blocked())
930            .await
931            .map_err(|_| {
932                Error::ExtensionError("control outbox did not reach test gate".to_string())
933            })?;
934        let applied = tokio::time::timeout(std::time::Duration::from_secs(1), application)
935            .await
936            .map_err(|_| {
937                Error::ExtensionError("terminal control effect held the gate".to_string())
938            })?
939            .map_err(|error| Error::ExtensionError(error.to_string()))??;
940
941        assert!(applied.is_empty());
942        hook.release();
943        Ok(())
944    }
945
946    #[tokio::test]
947    async fn test_saturated_peer_control_lane_does_not_block_another_peer() -> Result<()> {
948        let extensions = extensions()?;
949        let hook = Arc::new(ControlSendTestHook::default());
950        let interpreter = NativeRelay::new_with_control_send_test_hook(
951            Arc::new(TransportSessions::new()),
952            Arc::clone(&hook),
953        );
954        let blocked_peer: Did = SecretKey::random().address().into();
955        let independent_peer: Did = SecretKey::random().address().into();
956        let effect_scope = EffectScope::new(Scope::new(extensions.core(), TCP.to_string()));
957
958        interpreter
959            .run(&effect_scope, RelayEffect::SendClose {
960                to: blocked_peer,
961                session: SessionId(0),
962                from_opener: false,
963            })
964            .await?;
965        tokio::time::timeout(std::time::Duration::from_secs(1), hook.wait_until_blocked())
966            .await
967            .map_err(|_| {
968                Error::ExtensionError("first peer control lane did not block".to_string())
969            })?;
970
971        let mut saturated = false;
972        for session in 1..=8 {
973            let result = interpreter
974                .run(&effect_scope, RelayEffect::SendClose {
975                    to: blocked_peer,
976                    session: SessionId(session),
977                    from_opener: false,
978                })
979                .await;
980            if result.is_err() {
981                saturated = true;
982                break;
983            }
984        }
985        assert!(saturated, "the blocked peer must have a finite lane budget");
986
987        interpreter
988            .run(&effect_scope, RelayEffect::SendClose {
989                to: independent_peer,
990                session: SessionId(9),
991                from_opener: false,
992            })
993            .await?;
994        tokio::time::timeout(
995            std::time::Duration::from_secs(1),
996            hook.wait_until_completed(independent_peer),
997        )
998        .await
999        .map_err(|_| {
1000            Error::ExtensionError("independent peer control lane was blocked".to_string())
1001        })??;
1002
1003        hook.release();
1004        tokio::time::timeout(
1005            std::time::Duration::from_secs(1),
1006            hook.wait_until_completed(blocked_peer),
1007        )
1008        .await
1009        .map_err(|_| Error::ExtensionError("blocked peer lane did not resume".to_string()))??;
1010        Ok(())
1011    }
1012}