use std::sync::Arc;
use affinidi_messaging_didcomm::Message;
use tokio::sync::RwLock;
use affinidi_did_resolver_cache_sdk::DIDCacheClient;
use crate::config::AppConfig;
use crate::didcomm_bridge::DIDCommBridge;
use crate::keys::seed_store::SeedStore;
use crate::messaging::shim::{DIDCommResponse, DIDCommServiceError, HandlerContext, ProblemReport};
#[cfg(feature = "didcomm")]
use crate::messaging::shim::{Extension, ServiceProblemReport};
use crate::server::AppState;
use crate::store::KeyspaceHandle;
#[cfg(feature = "didcomm")]
use super::handlers;
#[cfg(feature = "didcomm")]
use vta_sdk::protocols;
#[cfg(feature = "didcomm")]
const TRUST_PING_TYPE: &str = "https://didcomm.org/trust-ping/2.0/ping";
#[cfg(feature = "didcomm")]
const TRUST_PONG_TYPE: &str = "https://didcomm.org/trust-ping/2.0/ping-response";
pub(crate) const MESSAGE_PICKUP_STATUS_TYPE: &str = "https://didcomm.org/messagepickup/3.0/status";
#[derive(Clone)]
pub struct VtaState {
pub keys_ks: KeyspaceHandle,
pub acl_ks: KeyspaceHandle,
pub sessions_ks: KeyspaceHandle,
pub contexts_ks: KeyspaceHandle,
pub did_templates_ks: KeyspaceHandle,
pub audit_ks: KeyspaceHandle,
pub audit_sink: vta_audit::SharedAuditSink,
pub imported_ks: KeyspaceHandle,
pub internal_ks: KeyspaceHandle,
pub service_state_ks: KeyspaceHandle,
#[cfg(feature = "webvh")]
pub webvh_ks: KeyspaceHandle,
pub issued_credentials_ks: KeyspaceHandle,
pub sealed_nonces_ks: KeyspaceHandle,
#[cfg(feature = "webvh")]
pub drains_ks: KeyspaceHandle,
#[cfg(feature = "webvh")]
pub snapshot_ks: KeyspaceHandle,
#[cfg(feature = "webvh")]
pub mediator_registry: Arc<crate::messaging::registry::MediatorListenerRegistry>,
#[cfg(feature = "webvh")]
pub drain_sweeper: Arc<crate::messaging::drain_sweeper::DrainSweeper>,
pub telemetry: vti_common::telemetry::SharedTelemetrySink,
pub seed_store: Arc<dyn SeedStore>,
pub config: Arc<RwLock<AppConfig>>,
pub did_resolver: Option<DIDCacheClient>,
pub didcomm_bridge: Arc<DIDCommBridge>,
#[cfg(feature = "didcomm")]
pub secrets_resolver: Option<Arc<affinidi_tdk::secrets_resolver::ThreadedSecretsResolver>>,
#[cfg(feature = "didcomm")]
pub signing_vm_id: Option<String>,
#[cfg(feature = "didcomm")]
pub ka_vm_id: Option<String>,
#[cfg(feature = "tee")]
pub tee_state: Option<crate::tee::TeeState>,
pub restart_tx: tokio::sync::watch::Sender<bool>,
pub store: vti_common::store::Store,
pub storage_encryption_key: Option<[u8; 32]>,
pub in_enclave: bool,
}
impl VtaState {
pub fn backup_access(&self) -> crate::restore::BackupAccess<'_> {
crate::restore::BackupAccess {
store: &self.store,
storage_key: self.storage_encryption_key,
in_enclave: self.in_enclave,
seed_store: self.seed_store.as_ref(),
config: &self.config,
}
}
}
#[cfg(feature = "webvh")]
impl From<&VtaState> for crate::operations::provision_integration::ProvisionIntegrationDeps {
fn from(state: &VtaState) -> Self {
Self {
keys_ks: state.keys_ks.clone(),
acl_ks: state.acl_ks.clone(),
audit: std::sync::Arc::clone(&state.audit_sink),
contexts_ks: state.contexts_ks.clone(),
did_templates_ks: state.did_templates_ks.clone(),
imported_ks: state.imported_ks.clone(),
webvh_ks: state.webvh_ks.clone(),
sealed_nonces_ks: state.sealed_nonces_ks.clone(),
seed_store: state.seed_store.clone(),
config: state.config.clone(),
did_resolver: state.did_resolver.clone(),
didcomm_bridge: state.didcomm_bridge.clone(),
}
}
}
impl From<&AppState> for VtaState {
fn from(state: &AppState) -> Self {
Self {
keys_ks: state.keys_ks.clone(),
acl_ks: state.acl_ks.clone(),
sessions_ks: state.sessions_ks.clone(),
contexts_ks: state.contexts_ks.clone(),
did_templates_ks: state.did_templates_ks.clone(),
audit_ks: state.audit_ks.clone(),
audit_sink: std::sync::Arc::clone(&state.audit_sink),
imported_ks: state.imported_ks.clone(),
internal_ks: state.internal_ks.clone(),
service_state_ks: state.service_state_ks.clone(),
#[cfg(feature = "webvh")]
webvh_ks: state.webvh_ks.clone(),
issued_credentials_ks: state.issued_credentials_ks.clone(),
sealed_nonces_ks: state.sealed_nonces_ks.clone(),
#[cfg(feature = "webvh")]
drains_ks: state.drains_ks.clone(),
#[cfg(feature = "webvh")]
snapshot_ks: state.snapshot_ks.clone(),
#[cfg(feature = "webvh")]
mediator_registry: Arc::clone(&state.mediator_registry),
#[cfg(feature = "webvh")]
drain_sweeper: Arc::clone(&state.drain_sweeper),
telemetry: Arc::clone(&state.telemetry),
seed_store: state.seed_store.clone(),
config: Arc::clone(&state.config),
did_resolver: state.did_resolver.clone(),
didcomm_bridge: Arc::clone(&state.didcomm_bridge),
#[cfg(feature = "didcomm")]
secrets_resolver: state.secrets_resolver.clone(),
#[cfg(feature = "didcomm")]
signing_vm_id: state.signing_vm_id.clone(),
#[cfg(feature = "didcomm")]
ka_vm_id: state.ka_vm_id.clone(),
#[cfg(feature = "tee")]
tee_state: state.tee.as_ref().map(|tc| tc.state.clone()),
restart_tx: state.restart_tx.clone(),
store: state.store.clone(),
storage_encryption_key: state.storage_encryption_key,
in_enclave: state.tee.is_some(),
}
}
}
#[cfg(feature = "didcomm")]
type HandlerResult = Result<Option<DIDCommResponse>, DIDCommServiceError>;
#[cfg(feature = "didcomm")]
fn finish(result: HandlerResult) -> Option<DIDCommResponse> {
match result {
Ok(opt) => opt,
Err(e) => Some(DIDCommResponse::problem_report(
ProblemReport::internal_error(e.to_string()),
)),
}
}
#[cfg(feature = "didcomm")]
fn trust_ping_reply(msg: &Message, sender_did: Option<&str>) -> Option<DIDCommResponse> {
#[derive(serde::Deserialize)]
struct PingBody {
#[serde(default = "default_true")]
response_requested: bool,
}
fn default_true() -> bool {
true
}
let body: PingBody = serde_json::from_value(msg.body.clone()).unwrap_or(PingBody {
response_requested: true,
});
if !body.response_requested {
return None;
}
sender_did?;
Some(DIDCommResponse::new(TRUST_PONG_TYPE, serde_json::Value::Null).thid(msg.id.clone()))
}
#[cfg(feature = "didcomm")]
pub async fn dispatch(
msg: Message,
ctx: HandlerContext,
_vta_state: Arc<VtaState>,
app_state: AppState,
) -> Option<DIDCommResponse> {
let t = msg.typ.clone();
let t = t.as_str();
if t == MESSAGE_PICKUP_STATUS_TYPE {
return None;
}
if t == TRUST_PING_TYPE {
return trust_ping_reply(&msg, ctx.sender_did.as_deref());
}
if t == trust_tasks_didcomm::ENVELOPE_TYPE {
return finish(handlers::handle_trust_task(ctx, msg, Extension(app_state)).await);
}
if t == protocols::PROBLEM_REPORT_TYPE {
return finish(handlers::handle_problem_report(ctx, msg).await);
}
finish(handlers::handle_unknown(ctx, msg).await)
}
#[cfg(all(test, feature = "didcomm"))]
mod envelope_only_carriage {
use super::*;
use serde_json::json;
use trust_tasks_didcomm::ENVELOPE_TYPE;
const STRANGER: &str = "did:key:z6MkhaXgBZDvotDkL5257faiztiGiC2QtKLGpbnnEGta2doK";
fn problem_comment(resp: &DIDCommResponse) -> Option<&str> {
(resp.type_ == vta_sdk::protocols::PROBLEM_REPORT_TYPE)
.then(|| resp.body.get("comment").and_then(|c| c.as_str()))
.flatten()
}
#[tokio::test]
async fn every_dispatched_uri_is_served_in_the_envelope_and_refused_typed_as_itself() {
let (app_state, _dir) = crate::test_support::build_signing_test_app_state().await;
let vta_state = Arc::new(VtaState::from(&app_state));
let uris = crate::trust_tasks::dispatched_uris();
assert!(!uris.is_empty(), "the dispatcher serves nothing?");
let mut typed_arm: Vec<&str> = Vec::new();
for uri in uris {
let doc = json!({
"id": format!("urn:uuid:{}", uuid::Uuid::new_v4()),
"type": uri,
"issuer": STRANGER,
"payload": {},
});
let req_id = format!("urn:uuid:{}", uuid::Uuid::new_v4());
let msg = Message::build(req_id.clone(), ENVELOPE_TYPE.to_string(), doc.clone())
.from(STRANGER.to_string())
.finalize();
let ctx = HandlerContext {
sender_did: Some(STRANGER.to_string()),
};
let resp = dispatch(msg, ctx, vta_state.clone(), app_state.clone())
.await
.unwrap_or_else(|| panic!("`{uri}` in the envelope got no reply at all"));
assert!(
!problem_comment(&resp).is_some_and(|c| c.contains("unsupported message type")),
"`{uri}` is dispatched, but the router refused its envelope: {:?}",
resp.body
);
assert_eq!(
resp.type_, ENVELOPE_TYPE,
"`{uri}`: a reply to an enveloped request rides the envelope (binding §5)"
);
let req_id = format!("urn:uuid:{}", uuid::Uuid::new_v4());
let msg = Message::build(req_id.clone(), uri.to_string(), doc)
.from(STRANGER.to_string())
.finalize();
let ctx = HandlerContext {
sender_did: Some(STRANGER.to_string()),
};
let resp = dispatch(msg, ctx, vta_state.clone(), app_state.clone())
.await
.unwrap_or_else(|| panic!("`{uri}` typed as itself got no reply at all"));
let refused = problem_comment(&resp).is_some_and(|c| c.contains(ENVELOPE_TYPE));
if !refused {
typed_arm.push(uri);
continue;
}
assert_eq!(
resp.thid.as_deref(),
Some(req_id.as_str()),
"`{uri}`: unthreaded"
);
}
assert!(
typed_arm.is_empty(),
"served URIs the router answers typed as the task: {typed_arm:?} — carry them in \
the envelope instead"
);
}
#[test]
fn a_non_trust_task_type_is_not_told_about_the_envelope() {
assert!(
handlers::trust_task_needs_envelope("https://example.com/protocols/x/1.0/y").is_none()
);
}
}
#[cfg(all(test, feature = "didcomm"))]
mod keyring_vti_09_27 {
use super::*;
use crate::acl::{AclEntry, Role, store_acl_entry};
use crate::auth::session::now_epoch;
use serde_json::{Value, json};
use trust_tasks_didcomm::ENVELOPE_TYPE;
const KEYS_LIST: &str = "https://trusttasks.org/spec/keys/list/0.1";
async fn send(app_state: &AppState, sender: &str, body: Value) -> DIDCommResponse {
let vta_state = Arc::new(VtaState::from(app_state));
let msg = Message::build(
format!("urn:uuid:{}", uuid::Uuid::new_v4()),
ENVELOPE_TYPE.to_string(),
body,
)
.from(sender.to_string())
.finalize();
let ctx = HandlerContext {
sender_did: Some(sender.to_string()),
};
dispatch(msg, ctx, vta_state, app_state.clone())
.await
.expect("an enveloped Trust Task is always answered")
}
fn signed_request(seed: u8, type_uri: &str) -> Value {
let (did, _) = crate::test_support::did_for_seed(seed);
let mut doc = vta_sdk::trust_task_sign::build_unsigned(
type_uri,
json!({}),
&did,
"did:key:z6MkfMo6gxqdBhaHMNnmfhgZFBjpCDTkmJMJLoypsBZS9PwD",
)
.expect("a well-formed document");
crate::test_support::sign_as(seed, &mut doc);
serde_json::to_value(doc).expect("serialise")
}
fn a_signed_trust_task_error(reply: &DIDCommResponse) -> (String, String) {
assert_eq!(
reply.type_, ENVELOPE_TYPE,
"a refusal of an enveloped Trust Task rides the envelope, not a problem-report: {:?}",
reply.body
);
let doc = &reply.body;
assert_eq!(
doc["type"].as_str(),
Some(
crate::trust_tasks::framework_error_type_uri()
.to_string()
.as_str()
),
"{doc}"
);
assert!(doc.get("proof").is_some(), "the refusal is signed: {doc}");
(
doc["payload"]["code"]
.as_str()
.unwrap_or_default()
.to_string(),
doc["payload"]["message"]
.as_str()
.unwrap_or_default()
.to_string(),
)
}
#[tokio::test]
async fn vti_09_a_bare_payload_over_didcomm_is_refused_as_malformed() {
let (app_state, _dir) = crate::test_support::build_signing_test_app_state().await;
let (admin, _) = crate::test_support::did_for_seed(0x61);
store_acl_entry(
&app_state.acl_ks,
&AclEntry::new(&admin, Role::Admin, "test"),
)
.await
.unwrap();
let (stranger, _) = crate::test_support::did_for_seed(0x62);
for sender in [&admin, &stranger] {
let reply = send(&app_state, sender, json!({ "contextId": "ctx1" })).await;
let (code, message) = a_signed_trust_task_error(&reply);
assert_eq!(code, "malformedRequest", "{sender}: {message}");
assert!(
message.contains("missing field `id`"),
"{sender}: the refusal names what is missing, as the VTC's does: {message}"
);
}
}
#[tokio::test]
async fn vti_27_an_acl_refusal_over_didcomm_is_a_signed_trust_task_error() {
let (app_state, _dir) = crate::test_support::build_signing_test_app_state().await;
let (lapsed, _) = crate::test_support::did_for_seed(0x64);
store_acl_entry(
&app_state.acl_ks,
&AclEntry::new(&lapsed, Role::Admin, "test")
.with_created_at(now_epoch().saturating_sub(7200))
.with_expires_at(Some(now_epoch().saturating_sub(60))),
)
.await
.unwrap();
for seed in [0x63u8, 0x64] {
let (sender, _) = crate::test_support::did_for_seed(seed);
let request = signed_request(seed, KEYS_LIST);
let reply = send(&app_state, &sender, request.clone()).await;
let (code, message) = a_signed_trust_task_error(&reply);
assert_eq!(code, "permissionDenied", "{sender}: {message}");
assert_eq!(
reply.body["threadId"], request["id"],
"{sender}: the refusal threads to the request it refuses"
);
}
}
#[tokio::test]
async fn vti_27_an_inbound_trust_task_error_over_didcomm_is_not_answered() {
let (app_state, _dir) = crate::test_support::build_signing_test_app_state().await;
let vta_state = Arc::new(VtaState::from(&app_state));
let (peer, _) = crate::test_support::did_for_seed(0x65);
let error = json!({
"id": format!("urn:uuid:{}", uuid::Uuid::new_v4()),
"threadId": "urn:uuid:11111111-1111-1111-1111-111111111111",
"type": crate::trust_tasks::framework_error_type_uri().to_string(),
"payload": { "code": "permissionDenied", "message": "no" },
});
let msg = Message::build(
format!("urn:uuid:{}", uuid::Uuid::new_v4()),
ENVELOPE_TYPE.to_string(),
error,
)
.from(peer.clone())
.finalize();
let ctx = HandlerContext {
sender_did: Some(peer),
};
let reply = dispatch(msg, ctx, vta_state, app_state.clone()).await;
assert!(
reply.is_none(),
"an inbound error is never answered: {:?}",
reply.map(|r| (r.type_, r.body))
);
}
}