use std::sync::Arc;
use std::time::Duration;
use affinidi_did_resolver_cache_sdk::DIDCacheClient;
use affinidi_messaging_core::{Inbound, InboundKind, MessageTransport, Protocol, ReceivedMessage};
#[cfg(feature = "didcomm")]
use affinidi_messaging_delivery::Delivery;
use affinidi_messaging_delivery::{MessagingService, OutboxStore};
#[cfg(feature = "didcomm")]
use affinidi_messaging_didcomm::Message;
use affinidi_tdk::common::TDKSharedState;
use affinidi_tdk::common::config::TDKConfig;
use affinidi_tdk::messaging::config::ATMConfig;
use affinidi_tdk::messaging::profiles::ATMProfile;
use affinidi_tdk::messaging::{ATM, DidCommTransport};
use affinidi_tdk::secrets_resolver::SecretsResolver;
use affinidi_tdk::secrets_resolver::secrets::Secret;
use futures_util::StreamExt;
use tokio_util::sync::CancellationToken;
use tracing::{info, warn};
use vti_common::outbox_store::VtiOutboxStore;
#[cfg(feature = "didcomm")]
use crate::messaging::router::{self, VtaState};
use crate::messaging::sender_order::{FrameOrder, SenderOrder};
#[cfg(feature = "didcomm")]
use crate::messaging::shim::{DIDCommResponse, ProblemReport, ServiceProblemReport};
use crate::server::AppState;
use crate::store::KeyspaceHandle;
pub struct VtaMessaging {
pub service: Arc<MessagingService>,
pub atm: Arc<ATM>,
pub profile: Arc<ATMProfile>,
#[cfg(feature = "tsp")]
pub tsp: crate::messaging::tsp_transport::TspTransport,
}
pub async fn build_messaging(
secrets: Vec<Secret>,
vta_did: &str,
mediator_did: &str,
outbox_ks: KeyspaceHandle,
pushes_ks: KeyspaceHandle,
relationships_ks: KeyspaceHandle,
relationship_drop_counter: std::sync::Arc<std::sync::atomic::AtomicU64>,
did_resolver: Option<&DIDCacheClient>,
resolver_url: Option<&str>,
) -> Result<VtaMessaging, String> {
let mut builder = TDKConfig::builder().with_load_environment(false);
if let Some(dr) = did_resolver {
builder = builder.with_did_resolver(dr.clone());
} else {
let resolver_config = vta_sdk::resolver::build_verifier_did_cache_config(
resolver_url,
vti_common::config::DID_CACHE_TTL_DEFAULT_SECS,
vti_common::config::DID_CACHE_CAPACITY_DEFAULT,
);
builder = builder.with_did_resolver_config(resolver_config);
}
let tdk_config = builder
.build()
.map_err(|e| format!("build TDK config: {e}"))?;
let tdk = TDKSharedState::new(tdk_config)
.await
.map_err(|e| format!("create TDK shared state: {e}"))?;
for secret in secrets {
tdk.secrets_resolver().insert(secret).await;
}
let atm_config_builder = ATMConfig::builder();
#[cfg(feature = "tsp")]
let atm_config_builder = atm_config_builder
.with_relationship_store(Arc::new(
affinidi_messaging_sdk::PersistentRelationshipStore::new(
crate::messaging::tsp_relationship_store::KeyspaceRelationshipKv::new(
relationships_ks,
),
),
))
.with_relationship_drop_counter(relationship_drop_counter.clone());
#[cfg(not(feature = "tsp"))]
let _ = (relationships_ks, &relationship_drop_counter);
let atm = Arc::new(
ATM::new(
atm_config_builder
.build()
.map_err(|e| format!("build ATM config: {e}"))?,
Arc::new(tdk),
)
.await
.map_err(|e| format!("create ATM: {e}"))?,
);
let profile = ATMProfile::new(
&atm,
None,
vta_did.to_string(),
Some(mediator_did.to_string()),
)
.await
.map_err(|e| format!("create ATM profile: {e}"))?;
let profile = atm
.profile_add(&profile, false)
.await
.map_err(|e| format!("register ATM profile: {e}"))?;
#[cfg(feature = "tsp")]
let tsp = match crate::messaging::tsp_transport::TspTransport::new(
(*atm).clone(),
profile.clone(),
) {
Some(t) => t,
None => {
atm.graceful_shutdown().await;
return Err(
"the profile registered for this session carries no mediator, so it can neither send nor unseal over TSP"
.to_string(),
);
}
};
let transport = match connect_transport(&atm, &profile).await {
Ok(transport) => transport,
Err(e) => {
atm.graceful_shutdown().await;
return Err(e);
}
};
let outbox: Arc<dyn OutboxStore> = Arc::new(VtiOutboxStore::new(outbox_ks));
let service = Arc::new(MessagingService::new(transport.clone(), outbox.clone()));
tokio::spawn(affinidi_messaging_delivery::drain_loop(
outbox.clone(),
service.primary_handle(),
Duration::from_secs(2),
));
tokio::spawn(affinidi_messaging_delivery::outbox_drain_loop(
service.primary_handle(),
outbox.clone(),
Duration::from_secs(10),
));
tokio::spawn(affinidi_messaging_delivery::confirmation_loop(
outbox.clone(),
Duration::from_secs(30),
));
crate::messaging::push::register_transports(
&service,
outbox.clone(),
pushes_ks,
&atm,
&profile,
mediator_did,
);
Ok(VtaMessaging {
service,
atm,
profile,
#[cfg(feature = "tsp")]
tsp,
})
}
async fn connect_transport(
atm: &Arc<ATM>,
profile: &Arc<ATMProfile>,
) -> Result<Arc<dyn MessageTransport>, String> {
match tokio::time::timeout(
Duration::from_secs(30),
atm.profile_enable_websocket(profile),
)
.await
{
Ok(res) => res.map_err(|e| format!("enable websocket: {e}"))?,
Err(_) => {
return Err(
"timeout enabling websocket to mediator after 30s — mediator may be unreachable"
.to_string(),
);
}
}
Ok(Arc::new(
DidCommTransport::new((**atm).clone(), profile.clone())
.await
.map_err(|e| format!("bind DidComm transport: {e}"))?,
))
}
pub async fn run_inbound_loop(
messaging: Arc<VtaMessaging>,
app_state: AppState,
vta_did: String,
shutdown: CancellationToken,
) {
#[cfg(feature = "didcomm")]
let vta_state = Arc::new(VtaState::from(&app_state));
let mut stream = messaging.service.subscribe();
info!("VTA messaging connected to mediator — inbound messages will be processed");
const MAX_INFLIGHT_INBOUND: usize = 32;
let inflight = Arc::new(tokio::sync::Semaphore::new(MAX_INFLIGHT_INBOUND));
let sender_order = SenderOrder::new();
loop {
tokio::select! {
maybe = stream.next() => {
let Some(inbound) = maybe else {
warn!("VTA inbound stream ended — messaging dispatcher stopping");
break;
};
let permit = match Arc::clone(&inflight).acquire_owned().await {
Ok(p) => p,
Err(_) => {
warn!("inbound concurrency semaphore closed — stopping");
break;
}
};
let ticket = sender_order
.admit(inbound.message.sender.as_deref(), frame_order(&inbound));
let messaging = Arc::clone(&messaging);
let app_state = app_state.clone();
#[cfg(feature = "didcomm")]
let vta_state = Arc::clone(&vta_state);
let vta_did = vta_did.clone();
tokio::spawn(async move {
let _permit = permit;
ticket.ready().await;
let _ticket = ticket;
#[cfg(feature = "didcomm")]
handle_inbound(inbound, &messaging, &app_state, &vta_state, &vta_did).await;
#[cfg(not(feature = "didcomm"))]
handle_inbound(inbound, &messaging, &app_state, &vta_did).await;
});
}
_ = shutdown.cancelled() => {
info!("VTA messaging stopping (shutdown signalled)");
break;
}
}
}
info!("VTA messaging stopped");
}
fn frame_order(inbound: &Inbound) -> FrameOrder {
match inbound.message.protocol {
Protocol::TSP => match inbound.kind {
InboundKind::RelationshipControl { .. } => FrameOrder::Barrier,
_ => FrameOrder::Follower,
},
_ => FrameOrder::Unordered,
}
}
async fn handle_inbound(
inbound: Inbound,
messaging: &Arc<VtaMessaging>,
app_state: &AppState,
#[cfg(feature = "didcomm")] vta_state: &Arc<VtaState>,
vta_did: &str,
) {
match inbound.message.protocol {
#[cfg(feature = "didcomm")]
Protocol::DIDComm => {
handle_didcomm(inbound, messaging, app_state, vta_state, vta_did).await;
}
#[cfg(not(feature = "didcomm"))]
Protocol::DIDComm => {
let _ = (messaging, app_state, vta_did);
warn!(
"received an inbound DIDComm frame but the `didcomm` feature is disabled — dropping"
);
}
#[cfg(feature = "tsp")]
Protocol::TSP => {
handle_tsp(inbound, messaging, app_state).await;
}
#[cfg(not(feature = "tsp"))]
Protocol::TSP => {
warn!("received an inbound TSP frame but the `tsp` feature is disabled — dropping");
}
Protocol::DIDCommV1 => {
warn!("received an inbound DIDComm v1 frame; this VTA speaks v2.1 only — dropping");
}
other => {
warn!(
protocol = ?other,
"received an inbound frame in a protocol this VTA does not implement — dropping"
);
}
}
}
#[cfg(feature = "didcomm")]
#[derive(Debug, Clone, PartialEq, Eq)]
enum InboundGate {
NotEncrypted,
Unauthenticated,
Authenticated(String),
}
#[cfg(feature = "didcomm")]
impl InboundGate {
fn authenticated_sender(&self) -> Option<&str> {
match self {
InboundGate::Authenticated(did) => Some(did.as_str()),
_ => None,
}
}
}
#[cfg(feature = "didcomm")]
fn inbound_gate(message: &ReceivedMessage) -> InboundGate {
if !message.encrypted {
return InboundGate::NotEncrypted;
}
match message.sender.as_deref() {
Some(did) if message.verified => InboundGate::Authenticated(did.to_string()),
_ => InboundGate::Unauthenticated,
}
}
#[cfg(feature = "didcomm")]
async fn gate_refusal(
gate: &InboundGate,
msg: &Message,
app_state: &AppState,
) -> Option<DIDCommResponse> {
let (why, report): (&str, fn(&str) -> ProblemReport) = match gate {
InboundGate::NotEncrypted => ("DIDComm message must be encrypted", |m| {
ProblemReport::bad_request(m)
}),
InboundGate::Unauthenticated => (
"DIDComm message must be authenticated (authcrypt) with a non-anonymous sender",
|m| ProblemReport::unauthorized(m),
),
InboundGate::Authenticated(_) => return None,
};
if msg.typ != trust_tasks_didcomm::ENVELOPE_TYPE {
return Some(DIDCommResponse::problem_report(report(why)));
}
if matches!(
vta_sdk::inbound::classify(&msg.body),
vta_sdk::inbound::Inbound::Error
) {
warn!(
reason = why,
"refused an inbound trust-task error at the DIDComm gate — terminal, not answered"
);
return None;
}
let body = serde_json::to_vec(&msg.body).unwrap_or_default();
let outcome = crate::trust_tasks::sign_response(
app_state,
crate::trust_tasks::reject_trust_task(
&body,
trust_tasks_rs::RejectReason::PermissionDenied {
reason: why.to_string(),
},
),
)
.await;
match serde_json::from_slice::<serde_json::Value>(&outcome.body) {
Ok(doc) => {
Some(DIDCommResponse::new(trust_tasks_didcomm::ENVELOPE_TYPE, doc).thid(msg.id.clone()))
}
Err(e) => {
warn!(error = %e, "could not read back a gate refusal document");
Some(DIDCommResponse::problem_report(report(why)))
}
}
}
#[cfg(feature = "didcomm")]
async fn handle_didcomm(
inbound: Inbound,
messaging: &Arc<VtaMessaging>,
app_state: &AppState,
vta_state: &Arc<VtaState>,
vta_did: &str,
) {
let mut msg: Message = match serde_json::from_slice(&inbound.message.payload) {
Ok(m) => m,
Err(e) => {
warn!(error = %e, "failed to parse inbound DIDComm message — dropping");
return;
}
};
let plaintext_from = msg.from.clone();
let gate = inbound_gate(&inbound.message);
let auth_sender = gate.authenticated_sender().map(str::to_string);
msg.from = auth_sender.clone();
let reply_to = auth_sender.clone().or(plaintext_from);
let msg_id = msg.id.clone();
let message_type = msg.typ.clone();
let start = std::time::Instant::now();
let reply = match &gate {
InboundGate::NotEncrypted | InboundGate::Unauthenticated => {
gate_refusal(&gate, &msg, app_state).await
}
InboundGate::Authenticated(_) => {
let ctx = crate::messaging::shim::HandlerContext {
sender_did: auth_sender.clone(),
};
router::dispatch(msg, ctx, vta_state.clone(), app_state.clone()).await
}
};
info!(
target: "didcomm_server::request",
message_type = %message_type,
sender = %auth_sender.as_deref().unwrap_or("<anon>"),
status = if reply.is_some() { "ok(response)" } else { "ok(empty)" },
latency = ?start.elapsed(),
"Request processed"
);
let Some(reply) = reply else {
return;
};
let Some(to) = reply_to else {
warn!(
reply_type = %reply.type_,
"computed a DIDComm reply but the inbound message had no sender/from to reply to — dropping"
);
return;
};
let reply_id = uuid::Uuid::new_v4().to_string();
let thid = reply.thid.unwrap_or(msg_id);
let reply_msg = Message::build(reply_id, reply.type_, reply.body)
.from(vta_did.to_string())
.to(to.clone())
.thid(thid)
.finalize();
match messaging
.atm
.pack_encrypted(&reply_msg, &to, Some(vta_did), Some(vta_did))
.await
{
Ok((packed, _)) => {
if let Err(e) = messaging
.service
.send(&to, packed.into_bytes(), Delivery::BestEffort)
.await
{
warn!(recipient = %to, error = %e, "failed to send DIDComm reply");
}
}
Err(e) => warn!(recipient = %to, error = %e, "failed to pack DIDComm reply"),
}
}
#[cfg(feature = "tsp")]
async fn handle_tsp(inbound: Inbound, messaging: &Arc<VtaMessaging>, app_state: &AppState) {
let Some(sender_vid) = inbound.message.sender.clone() else {
warn!("inbound TSP frame has no authenticated sender VID — dropping");
return;
};
if let affinidi_messaging_core::InboundKind::RelationshipControl {
request,
thread_digest,
reply_expected,
introduces,
} = &inbound.kind
{
handle_tsp_control(
&sender_vid,
*request,
*thread_digest,
*reply_expected,
introduces.as_deref(),
messaging,
)
.await;
return;
}
let reply = crate::messaging::tsp_inbound::dispatch_one(
app_state,
&inbound.message.payload,
&sender_vid,
)
.await;
if reply.is_empty() {
return;
}
if let Err(e) = messaging.tsp.send_to(&sender_vid, &reply).await {
warn!(recipient = %sender_vid, error = %e, "failed to send TSP reply");
}
}
#[cfg(feature = "tsp")]
async fn handle_tsp_control(
sender_vid: &str,
request: affinidi_messaging_core::RelationshipRequest,
thread_digest: [u8; 32],
reply_expected: bool,
introduces: Option<&str>,
messaging: &Arc<VtaMessaging>,
) {
use crate::messaging::tsp_inbound::{ControlDecision, decide_control};
let atm = messaging.atm.clone();
let profile = messaging.profile.clone();
match decide_control(request, reply_expected) {
ControlDecision::Accept => {
match atm
.tsp()
.accept_relationship(&profile, sender_vid, thread_digest)
.await
{
Ok(state) => info!(
sender = %sender_vid, ?request, ?state, introduced = ?introduces,
"accepted an inbound TSP relationship request",
),
Err(e) => warn!(
sender = %sender_vid, error = %e,
"could not send a TSP relationship accept; the relationship stays recorded, \
so traffic still flows, but the peer sees no answer",
),
}
}
ControlDecision::Cancel(why) => {
match atm
.tsp()
.answer_cancellation(&profile, sender_vid, thread_digest)
.await
{
Ok(_) => info!(
sender = %sender_vid, ?request, reason = %why,
"answered an inbound TSP relationship cancellation (§7.3) after the \
transport's own answer failed",
),
Err(e) => warn!(
sender = %sender_vid, reason = %why, error = %e,
"could not send the §7.3 answer to a TSP relationship cancellation; \
the relationship is forgotten on this side, but the peer sees no answer",
),
}
}
ControlDecision::Nothing => info!(
sender = %sender_vid, ?request,
"recorded an inbound TSP relationship request; no answer is due",
),
}
}
#[cfg(test)]
mod tests {
use super::{InboundGate, inbound_gate};
use affinidi_messaging_core::{Protocol, ReceivedMessage};
const SENDER: &str = "did:key:z6MkSenderUnderTest";
fn received(encrypted: bool, sender: Option<&str>, verified: bool) -> ReceivedMessage {
ReceivedMessage {
id: "urn:uuid:test".to_string(),
sender: sender.map(str::to_string),
recipient: "did:key:z6MkRecipient".to_string(),
payload: Vec::new(),
protocol: Protocol::DIDComm,
verified,
encrypted,
}
}
#[test]
fn encrypted_and_verified_sender_is_authenticated() {
assert_eq!(
inbound_gate(&received(true, Some(SENDER), true)),
InboundGate::Authenticated(SENDER.to_string()),
);
}
#[test]
fn unverified_sender_is_rejected() {
assert_eq!(
inbound_gate(&received(true, Some(SENDER), false)),
InboundGate::Unauthenticated,
);
assert_eq!(
inbound_gate(&received(true, Some(SENDER), false)).authenticated_sender(),
None,
);
}
#[test]
fn anonymous_encrypted_sender_is_rejected() {
assert_eq!(
inbound_gate(&received(true, None, false)),
InboundGate::Unauthenticated,
);
assert_eq!(
inbound_gate(&received(true, None, true)),
InboundGate::Unauthenticated,
);
}
#[test]
fn plaintext_frame_is_rejected() {
assert_eq!(
inbound_gate(&received(false, Some(SENDER), false)),
InboundGate::NotEncrypted,
);
assert_eq!(
inbound_gate(&received(false, Some(SENDER), true)),
InboundGate::NotEncrypted,
);
assert_eq!(
inbound_gate(&received(false, Some(SENDER), true)).authenticated_sender(),
None,
);
}
}
#[cfg(test)]
mod frame_order_tests {
use super::frame_order;
use crate::messaging::sender_order::FrameOrder;
use affinidi_messaging_core::{
Inbound, InboundAck, InboundKind, Protocol, ReceivedMessage, RelationshipRequest,
};
fn inbound(protocol: Protocol) -> Inbound {
Inbound::new(
ReceivedMessage {
id: "urn:uuid:test".to_string(),
sender: Some("did:peer:2.phone".to_string()),
recipient: "did:key:z6MkRecipient".to_string(),
payload: Vec::new(),
protocol,
verified: true,
encrypted: true,
},
None,
InboundAck("ack".into()),
)
}
#[test]
fn a_tsp_relationship_request_is_a_barrier() {
let xrfi = inbound(Protocol::TSP).with_kind(InboundKind::RelationshipControl {
request: RelationshipRequest::Invite,
thread_digest: [7u8; 32],
reply_expected: false,
introduces: None,
});
assert_eq!(frame_order(&xrfi), FrameOrder::Barrier);
}
#[test]
fn tsp_traffic_follows() {
assert_eq!(frame_order(&inbound(Protocol::TSP)), FrameOrder::Follower);
}
#[test]
fn didcomm_is_unordered() {
assert_eq!(
frame_order(&inbound(Protocol::DIDComm)),
FrameOrder::Unordered
);
}
}
#[cfg(all(test, feature = "didcomm"))]
mod keyring_vti_27_gate {
use super::{InboundGate, gate_refusal};
use crate::test_support::build_signing_test_app_state;
use affinidi_messaging_didcomm::Message;
use serde_json::{Value, json};
use trust_tasks_didcomm::ENVELOPE_TYPE;
fn message(typ: &str, body: Value) -> Message {
Message::build(
format!("urn:uuid:{}", uuid::Uuid::new_v4()),
typ.to_string(),
body,
)
.finalize()
}
fn request() -> Value {
json!({
"id": format!("urn:uuid:{}", uuid::Uuid::new_v4()),
"type": "https://trusttasks.org/spec/keys/list/0.1",
"issuer": "did:key:z6MkForgedSender",
"payload": {},
})
}
#[tokio::test]
async fn vti_27_a_gate_refused_trust_task_is_a_signed_trust_task_error() {
let (app_state, _dir) = build_signing_test_app_state().await;
for gate in [InboundGate::Unauthenticated, InboundGate::NotEncrypted] {
let req = request();
let msg = message(ENVELOPE_TYPE, req.clone());
let reply = gate_refusal(&gate, &msg, &app_state)
.await
.expect("a refused request is answered");
assert_eq!(reply.type_, ENVELOPE_TYPE, "{gate:?}: {:?}", reply.body);
assert_eq!(reply.thid.as_deref(), Some(msg.id.as_str()), "{gate:?}");
let doc = &reply.body;
assert_eq!(
doc["type"].as_str(),
Some(
crate::trust_tasks::framework_error_type_uri()
.to_string()
.as_str()
),
"{gate:?}: {doc}"
);
assert_eq!(
doc["payload"]["code"], "permissionDenied",
"{gate:?}: {doc}"
);
assert_eq!(doc["threadId"], req["id"], "{gate:?}: {doc}");
assert!(doc.get("proof").is_some(), "{gate:?}: unsigned: {doc}");
}
}
#[tokio::test]
async fn a_gate_refused_non_trust_task_is_still_a_problem_report() {
let (app_state, _dir) = build_signing_test_app_state().await;
for (gate, code) in [
(InboundGate::Unauthenticated, "e.p.msg.unauthorized"),
(InboundGate::NotEncrypted, "e.p.msg.bad-request"),
] {
let msg = message("https://didcomm.org/trust-ping/2.0/ping", json!({}));
let reply = gate_refusal(&gate, &msg, &app_state)
.await
.expect("answered");
assert_eq!(reply.type_, vta_sdk::protocols::PROBLEM_REPORT_TYPE);
assert_eq!(reply.body["code"], code, "{gate:?}");
}
}
#[tokio::test]
async fn a_gate_refused_trust_task_error_is_not_answered() {
let (app_state, _dir) = build_signing_test_app_state().await;
let error = json!({
"id": format!("urn:uuid:{}", uuid::Uuid::new_v4()),
"threadId": "urn:uuid:11111111-1111-1111-1111-111111111111",
"type": "https://trusttasks.org/spec/trust-task-error/0.5",
"payload": { "code": "permissionDenied", "message": "no" },
});
let msg = message(ENVELOPE_TYPE, error);
assert!(
gate_refusal(&InboundGate::Unauthenticated, &msg, &app_state)
.await
.is_none()
);
}
#[tokio::test]
async fn an_authenticated_frame_gets_no_gate_refusal() {
let (app_state, _dir) = build_signing_test_app_state().await;
let msg = message(ENVELOPE_TYPE, request());
assert!(
gate_refusal(
&InboundGate::Authenticated("did:key:z6MkSender".into()),
&msg,
&app_state
)
.await
.is_none()
);
}
}