vgi-bridge 0.16.0

The per-community VGI bridge: holds the community's forge credentials, takes git-ns/bridge jobs from its VTC over TSP or DIDComm, runs the forge adapters, reports results, events and drift, and posts the commit-trust check where no forge feature can be trusted to.
//! The bridge's channel to the Trust Registry for the bridge-posted check.
//!
//! When the registry's DID document advertises TSP or DIDComm, the bridge
//! queries it **as its own DID** — the stable, VTA-provisioned identity the
//! community already knows — over the mediator session it already holds for
//! the VTC ([`VtcLink`]): one channel per transport, TSP preferred, chosen by
//! verify-trust's route discovery with no fallback once chosen. It does not
//! open a second session: the mediator permits one websocket per DID, and a
//! second would take the job link's slot.
//!
//! Replies come back on the same inbound stream as the VTC's jobs.
//! [`Bridge::handle_inbound`](crate::Bridge::handle_inbound) offers each
//! inbound document to [`RegistryReplies::route`] first; a document is taken
//! only when the transport proved it came from the registry's DID, it came
//! over the transport the query went out on, **and** it answers a query in
//! flight (`threadId`). Everything else — including a registry that is also
//! the VTC, sending a job — carries on to the job path untouched.
//!
//! The answer's integrity rests on the registry's key, not on the bridge's:
//! a reply the transport did not authenticate as the registry is never
//! delivered to a waiting query, which then times out as
//! `registryUnavailable`.
//!
//! The authenticated sender is taken from the messaging SDK, which hands the
//! bridge each message already unpacked — the envelope is not available here
//! for verify-trust's own header check. Over TSP it is the sender VID the
//! message's signature verified against (there is no separate plaintext
//! claim). Over DIDComm it is believed because the SDK
//! (affinidi-messaging-sdk 0.27.2 / affinidi-messaging-didcomm 0.15.9 and
//! later, which this workspace requires) binds an authcrypt message's
//! reported sender to the key its key agreement actually used. On top of it,
//! the sender must be a *verified* one, the reply's thread must be a query in
//! flight — every query carries a fresh random id (verify-trust's channel
//! transport) — and the first failure is latched for the rest of the check.
//! A registry whose mediator refuses the bridge's DID (access-list mode
//! `ExplicitAllow` without the bridge on the list) fails the same way —
//! closed.

use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::Duration;

use async_trait::async_trait;
use serde_json::Value;
use tokio::sync::oneshot;
use trql_client::{TransportKind, TrqlError};
use verify_trust::RegistryChannel;

use crate::transport::{InboundDoc, Via, VtcLink};

/// How long a registry query waits for its reply.
const REPLY_TIMEOUT: Duration = Duration::from_secs(30);

/// A query in flight: the transport it went out on, and its waiter.
type Waiter = (Via, oneshot::Sender<Value>);

/// Registry queries in flight, keyed by request document id.
#[derive(Debug)]
pub struct RegistryReplies {
    registry_did: String,
    pending: Mutex<HashMap<String, Waiter>>,
}

impl RegistryReplies {
    /// For the registry at `registry_did`.
    pub fn new(registry_did: impl Into<String>) -> Self {
        RegistryReplies {
            registry_did: registry_did.into(),
            pending: Mutex::new(HashMap::new()),
        }
    }

    fn lock(&self) -> std::sync::MutexGuard<'_, HashMap<String, Waiter>> {
        self.pending.lock().unwrap_or_else(|p| p.into_inner())
    }

    /// Register a query in flight. The entry lives as long as the returned
    /// guard: however the wait ends — answer, timeout, send failure, or the
    /// exchange future being dropped (abort, shutdown) — it is removed, so a
    /// late reply to it is never taken and the map cannot grow.
    fn register(&self, id: &str, via: Via) -> (Pending<'_>, oneshot::Receiver<Value>) {
        let (tx, rx) = oneshot::channel();
        self.lock().insert(id.to_string(), (via, tx));
        (
            Pending {
                replies: self,
                id: id.to_string(),
            },
            rx,
        )
    }

    #[cfg(test)]
    fn in_flight(&self) -> usize {
        self.lock().len()
    }

    /// Take `inbound` if it is the registry's answer to a query in flight.
    /// Returns whether it was taken; if not, it belongs to the job path.
    pub fn route(&self, inbound: &InboundDoc) -> bool {
        let Some(sender) = inbound.authenticated_sender.as_deref() else {
            return false;
        };
        let sender = sender.split_once('#').map_or(sender, |(did, _)| did);
        if sender != self.registry_did {
            return false;
        }
        let Some(thread) = inbound.doc.get("threadId").and_then(Value::as_str) else {
            return false;
        };
        let waiter = {
            let mut pending = self.lock();
            // Over the transport the query went out on, or not an answer:
            // the binding was chosen, and nothing falls back from it.
            if pending
                .get(thread)
                .is_none_or(|(via, _)| *via != inbound.via)
            {
                return false;
            }
            match pending.remove(thread) {
                Some((_, waiter)) => waiter,
                None => return false,
            }
        };
        // A waiter that gave up (timed out) still means the document was a
        // registry reply, not a job.
        let _ = waiter.send(inbound.doc.clone());
        true
    }
}

/// A query's entry in [`RegistryReplies`], removed when this is dropped.
struct Pending<'a> {
    replies: &'a RegistryReplies,
    id: String,
}

impl Drop for Pending<'_> {
    fn drop(&mut self) {
        self.replies.lock().remove(&self.id);
    }
}

/// [`RegistryChannel`] over the bridge's mediator link, as the bridge's DID,
/// on one transport.
///
/// It keeps [`RegistryChannel::exchange`]'s security contract through
/// [`RegistryReplies::route`]: only a reply whose transport-verified sender
/// is the registry DID, over this channel's transport, on the thread of a
/// query in flight, is returned.
pub struct BridgeRegistryChannel {
    link: Arc<dyn VtcLink>,
    did: String,
    replies: Arc<RegistryReplies>,
    via: Via,
    timeout: Duration,
}

impl BridgeRegistryChannel {
    /// Send as `did` over `link`, on `via` only; replies arrive through
    /// `replies`.
    pub fn new(
        link: Arc<dyn VtcLink>,
        did: impl Into<String>,
        replies: Arc<RegistryReplies>,
        via: Via,
    ) -> Self {
        BridgeRegistryChannel {
            link,
            did: did.into(),
            replies,
            via,
            timeout: REPLY_TIMEOUT,
        }
    }

    /// Wait at most `timeout` for each reply (tests).
    #[doc(hidden)]
    pub fn with_timeout(mut self, timeout: Duration) -> Self {
        self.timeout = timeout;
        self
    }
}

#[async_trait]
impl RegistryChannel for BridgeRegistryChannel {
    fn kind(&self) -> TransportKind {
        match self.via {
            Via::Tsp => TransportKind::Tsp,
            Via::Didcomm => TransportKind::Didcomm,
        }
    }

    fn sender_did(&self) -> &str {
        &self.did
    }

    async fn exchange(&self, recipient: &str, request: Value) -> Result<Value, TrqlError> {
        let kind = self.kind();
        let transport = |detail: String| TrqlError::Transport { kind, detail };
        if recipient != self.replies.registry_did {
            return Err(TrqlError::Config(format!(
                "registry channel is for {}, not {recipient}",
                self.replies.registry_did
            )));
        }
        let id = request
            .get("id")
            .and_then(Value::as_str)
            .ok_or_else(|| TrqlError::Contract("request document has no id".to_string()))?
            .to_string();
        // Register before sending, so a fast reply cannot be lost. `_pending`
        // unregisters on every exit, including this future being dropped.
        let (_pending, reply) = self.replies.register(&id, self.via);
        if let Err(e) = self.link.send_via(recipient, &request, self.via).await {
            return Err(transport(format!("sending to the registry: {e:#}")));
        }
        match tokio::time::timeout(self.timeout, reply).await {
            Ok(Ok(doc)) => Ok(doc),
            Ok(Err(_)) => Err(transport("the reply channel closed".to_string())),
            Err(_) => Err(TrqlError::Timeout {
                kind,
                waited_secs: self.timeout.as_secs(),
            }),
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::transport::memory::ChannelLink;

    const REGISTRY: &str = "did:webvh:QmReg:registry.example";

    fn inbound(sender: Option<&str>, doc: Value) -> InboundDoc {
        InboundDoc {
            doc,
            authenticated_sender: sender.map(str::to_string),
            via: Via::Didcomm,
        }
    }

    #[tokio::test]
    async fn a_query_goes_out_as_the_bridge_and_its_proven_answer_comes_back() {
        let (link, mut sent) = ChannelLink::new();
        let replies = Arc::new(RegistryReplies::new(REGISTRY));
        let channel = BridgeRegistryChannel::new(
            Arc::new(link),
            "did:key:z6MkBridge",
            replies.clone(),
            Via::Didcomm,
        );

        let answer = tokio::spawn(async move {
            channel
                .exchange(
                    REGISTRY,
                    serde_json::json!({ "id": "urn:uuid:q", "issuer": "did:key:z6MkBridge" }),
                )
                .await
        });
        let (to, doc) = sent.recv().await.unwrap();
        assert_eq!(to, REGISTRY);
        assert_eq!(doc["issuer"], "did:key:z6MkBridge");

        // A correlated reply the transport did not prove is from the
        // registry is not an answer: it stays with the job path.
        let forged =
            serde_json::json!({ "threadId": "urn:uuid:q", "payload": { "authorized": true } });
        assert!(!replies.route(&inbound(Some("did:key:z6MkMallory"), forged.clone())));
        assert!(!replies.route(&inbound(None, forged)));

        let real =
            serde_json::json!({ "threadId": "urn:uuid:q", "payload": { "authorized": false } });
        assert!(replies.route(&inbound(Some(&format!("{REGISTRY}#key-2")), real)));
        let got = answer.await.unwrap().unwrap();
        assert_eq!(got["payload"]["authorized"], false);
    }

    #[tokio::test]
    async fn a_tsp_query_goes_out_over_tsp_and_only_a_tsp_answer_is_taken() {
        let (link, mut sent) = ChannelLink::new();
        let link = Arc::new(link);
        let replies = Arc::new(RegistryReplies::new(REGISTRY));
        let channel = BridgeRegistryChannel::new(
            link.clone(),
            "did:key:z6MkBridge",
            replies.clone(),
            Via::Tsp,
        );
        assert_eq!(channel.kind(), TransportKind::Tsp);
        let answer = tokio::spawn(async move {
            channel
                .exchange(REGISTRY, serde_json::json!({ "id": "urn:uuid:q" }))
                .await
        });
        sent.recv().await.unwrap();
        assert_eq!(
            link.pinned(),
            [Via::Tsp],
            "pinned to TSP, not the link's choice"
        );

        let doc =
            serde_json::json!({ "threadId": "urn:uuid:q", "payload": { "authorized": true } });
        // The registry's own key, but over DIDComm: not the chosen binding.
        assert!(!replies.route(&inbound(Some(REGISTRY), doc.clone())));
        // Over TSP, from another VID: not the registry.
        let mallory = InboundDoc {
            via: Via::Tsp,
            ..inbound(Some("did:key:z6MkMallory"), doc.clone())
        };
        assert!(!replies.route(&mallory));
        let real = InboundDoc {
            via: Via::Tsp,
            ..inbound(Some(REGISTRY), doc)
        };
        assert!(replies.route(&real));
        assert_eq!(
            answer.await.unwrap().unwrap()["payload"]["authorized"],
            true
        );
    }

    #[tokio::test]
    async fn a_registry_that_never_answers_times_out_and_later_mail_is_not_taken() {
        let (link, _sent) = ChannelLink::new();
        let replies = Arc::new(RegistryReplies::new(REGISTRY));
        let channel = BridgeRegistryChannel::new(
            Arc::new(link),
            "did:key:z6MkBridge",
            replies.clone(),
            Via::Didcomm,
        )
        .with_timeout(Duration::from_millis(20));
        let e = channel
            .exchange(REGISTRY, serde_json::json!({ "id": "urn:uuid:q" }))
            .await
            .unwrap_err();
        assert!(matches!(e, TrqlError::Timeout { .. }), "{e}");
        // The late reply — or a job from a registry that is also the VTC —
        // goes on to the job path.
        let late = serde_json::json!({ "threadId": "urn:uuid:q" });
        assert!(!replies.route(&inbound(Some(REGISTRY), late)));
    }

    #[tokio::test]
    async fn an_exchange_dropped_mid_wait_leaves_no_query_in_flight() {
        let (link, mut sent) = ChannelLink::new();
        let replies = Arc::new(RegistryReplies::new(REGISTRY));
        let channel = BridgeRegistryChannel::new(
            Arc::new(link),
            "did:key:z6MkBridge",
            replies.clone(),
            Via::Didcomm,
        );
        let waiting = tokio::spawn(async move {
            channel
                .exchange(REGISTRY, serde_json::json!({ "id": "urn:uuid:q" }))
                .await
        });
        sent.recv().await.unwrap();
        assert_eq!(replies.in_flight(), 1);
        waiting.abort();
        assert!(waiting.await.unwrap_err().is_cancelled());
        assert_eq!(replies.in_flight(), 0, "an aborted query must unregister");
        let late = serde_json::json!({ "threadId": "urn:uuid:q" });
        assert!(!replies.route(&inbound(Some(REGISTRY), late)));
    }

    #[tokio::test]
    async fn a_link_that_is_down_fails_the_query_rather_than_waiting() {
        let link = crate::transport::SupervisedLink::new(); // never connected
        let replies = Arc::new(RegistryReplies::new(REGISTRY));
        let channel =
            BridgeRegistryChannel::new(Arc::new(link), "did:key:z6MkBridge", replies, Via::Didcomm);
        let e = channel
            .exchange(REGISTRY, serde_json::json!({ "id": "urn:uuid:q" }))
            .await
            .unwrap_err();
        assert!(matches!(e, TrqlError::Transport { .. }), "{e}");
    }
}