use super::helpers::{
TRANSPORT_TRUST_TASK, TrustTaskOutcome, app_error_to_reject, parse_payload, success_response,
};
use serde_json::Value;
use trust_tasks_rs::TrustTask;
use trust_tasks_rs::specs::vta::services as spec;
use crate::auth::AuthClaims;
use crate::error::AppError;
use crate::operations::protocol::list::ListServicesError;
use crate::operations::protocol::{OpContext, ServiceOpDeps};
use crate::server::AppState;
use vta_sdk::protocol::services::ServiceState as SdkState;
fn resolver(state: &AppState) -> Result<affinidi_did_resolver_cache_sdk::DIDCacheClient, AppError> {
state
.did_resolver
.as_ref()
.cloned()
.ok_or_else(|| AppError::Internal("DID resolver not available".into()))
}
fn setup_acl_from_ext<T: serde::Serialize>(ext: Option<&T>) -> Result<bool, AppError> {
let Some(ext) = ext else { return Ok(false) };
let value = serde_json::to_value(ext)
.map_err(|e| AppError::Validation(format!("invalid service extension: {e}")))?;
match value
.get("org.openvtc")
.and_then(|namespace| namespace.get("setupAcl"))
{
None => Ok(false),
Some(Value::Bool(enabled)) => Ok(*enabled),
Some(_) => Err(AppError::Validation(
"ext.org.openvtc.setupAcl must be a boolean".into(),
)),
}
}
#[cfg(feature = "didcomm")]
async fn live_didcomm_prover(
state: &AppState,
) -> Option<crate::messaging::live_prover::DIDCommServiceProver> {
let vta_did = {
let config = state.config.read().await;
config.vta_did.clone()?
};
crate::messaging::live_prover::try_build_from_parts(
&state.didcomm_bridge,
&vta_did,
state.secrets_resolver.as_ref()?,
state.signing_vm_id.as_ref()?,
state.ka_vm_id.as_ref()?,
)
.await
}
#[cfg(feature = "didcomm")]
async fn run_first_enable_handshake(
state: &AppState,
resolver: &affinidi_did_resolver_cache_sdk::DIDCacheClient,
mediator_did: &str,
timeout: std::time::Duration,
setup_acl: bool,
) -> Result<(), crate::messaging::handshake::HandshakeError> {
use crate::messaging::handshake::{HandshakeError, HandshakeOptions, HandshakeStage};
use crate::messaging::transient_handshake::{
TransientHandshakeContext, run_transient_handshake,
};
use affinidi_tdk::common::config::TDKConfig;
use affinidi_tdk::secrets_resolver::SecretsResolver;
let missing = |name: &str| HandshakeError::Failed {
stage: HandshakeStage::Connect,
cause: format!("first-enable handshake prerequisite missing: {name}"),
};
let secrets_resolver = state
.secrets_resolver
.as_ref()
.ok_or_else(|| missing("secrets resolver"))?;
let signing_vm_id = state
.signing_vm_id
.as_ref()
.ok_or_else(|| missing("signing verification method"))?;
let ka_vm_id = state
.ka_vm_id
.as_ref()
.ok_or_else(|| missing("key-agreement verification method"))?;
let vta_did = state
.config
.read()
.await
.vta_did
.clone()
.ok_or_else(|| missing("VTA DID"))?;
let mut secrets = Vec::with_capacity(2);
if let Some(secret) = secrets_resolver.get_secret(signing_vm_id).await {
secrets.push(secret);
}
if let Some(secret) = secrets_resolver.get_secret(ka_vm_id).await {
secrets.push(secret);
}
if secrets.is_empty() {
return Err(missing("VTA signing and key-agreement secrets"));
}
let mut resolver = resolver.clone();
crate::server::preload_self_did_document(&mut resolver, &vta_did, Some(&state.webvh_ks)).await;
let tdk_config = TDKConfig::builder()
.with_did_resolver(resolver.clone())
.with_load_environment(false)
.build()
.map_err(|e| HandshakeError::Failed {
stage: HandshakeStage::Connect,
cause: format!("build resolver-backed TDK config: {e}"),
})?;
run_transient_handshake(
TransientHandshakeContext {
vta_did,
secrets,
tdk_config: Some(tdk_config),
},
&resolver,
&state.telemetry,
mediator_did,
HandshakeOptions {
timeout,
setup_acl,
channel: TRANSPORT_TRUST_TASK.to_string(),
force: false,
},
)
.await
.map(|_| ())
}
fn state(s: &SdkState) -> spec::list::v1_0::ServiceState {
use spec::list::v1_0::{ServiceKind as K, ServiceState as Out};
let (kind, enabled, mediator_did, url) = match s {
SdkState::Tsp {
enabled,
mediator_did,
} => (K::Tsp, *enabled, mediator_did.clone(), None),
SdkState::Rest { enabled, url } => (K::Rest, *enabled, None, url.clone()),
SdkState::Didcomm {
enabled,
mediator_did,
..
} => (K::Didcomm, *enabled, mediator_did.clone(), None),
SdkState::Webauthn { enabled, url } => (K::Webauthn, *enabled, None, url.clone()),
};
Out::builder()
.kind(kind)
.enabled(enabled)
.mediator_did(mediator_did)
.url(url)
.try_into()
.expect("every required member of the service state is set above")
}
fn built<T, B>(builder: B) -> Result<T, AppError>
where
B: TryInto<T>,
B::Error: std::fmt::Display,
{
builder
.try_into()
.map_err(|e| AppError::Internal(format!("couldn't build a services response: {e}")))
}
fn kind_of(s: &SdkState) -> spec::list::v1_0::ServiceKind {
state(s).kind
}
fn list_error(e: ListServicesError) -> AppError {
match e {
ListServicesError::Auth(msg) => AppError::Forbidden(msg),
ListServicesError::VtaDidNotConfigured => {
AppError::Conflict("VTA DID is not configured — run `vta setup` first".into())
}
other => AppError::Internal(other.to_string()),
}
}
pub(super) async fn handle_list(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let _req: spec::list::v1_0::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
match crate::operations::protocol::list::list_services(&state_.config, &state_.webvh_ks, auth)
.await
{
Ok(body) => match built::<spec::list::v1_0::Response, _>(
spec::list::v1_0::Response::builder()
.services(body.services.iter().map(state).collect::<Vec<_>>()),
) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
},
Err(e) => app_error_to_reject(&doc, list_error(e)),
}
}
pub(super) async fn handle_get(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: spec::get::v1_0::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
match crate::operations::protocol::list::list_services(&state_.config, &state_.webvh_ks, auth)
.await
{
Ok(body) => {
let want = match req.service {
spec::get::v1_0::ServiceKind::Didcomm => spec::list::v1_0::ServiceKind::Didcomm,
spec::get::v1_0::ServiceKind::Rest => spec::list::v1_0::ServiceKind::Rest,
spec::get::v1_0::ServiceKind::Tsp => spec::list::v1_0::ServiceKind::Tsp,
spec::get::v1_0::ServiceKind::Webauthn => spec::list::v1_0::ServiceKind::Webauthn,
_ => {
return app_error_to_reject(
&doc,
AppError::Validation(
"services/get: unrecognised service kind — this VTA does not serve \
a transport added to the registry after it was built"
.into(),
),
);
}
};
match body.services.iter().find(|s| kind_of(s) == want) {
Some(found) => match built::<spec::get::v1_0::Response, _>(
spec::get::v1_0::Response::builder().state(
serde_json::from_value::<spec::get::v1_0::ServiceState>(
serde_json::to_value(state(found)).unwrap_or(Value::Null),
)
.expect("the two generated ServiceState shapes are identical"),
),
) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
},
None => app_error_to_reject(
&doc,
crate::error::AppError::NotFound(format!(
"no {:?} transport is configured on this agent",
req.service
)),
),
}
}
Err(e) => app_error_to_reject(&doc, list_error(e)),
}
}
use crate::operations::protocol::disable_didcomm::{
DisableDidcommParams, DisableTransport, disable_didcomm,
};
use crate::operations::protocol::disable_rest::{DisableRestParams, disable_rest};
use crate::operations::protocol::disable_tsp::{DisableTspParams, disable_tsp};
use crate::operations::protocol::disable_webauthn::{DisableWebauthnParams, disable_webauthn};
use crate::operations::protocol::drain_cancel::{DrainCancelParams, drain_cancel};
use crate::operations::protocol::enable_rest::{EnableRestParams, enable_rest};
use crate::operations::protocol::enable_tsp::{EnableTspParams, enable_tsp};
use crate::operations::protocol::enable_webauthn::{EnableWebauthnParams, enable_webauthn};
use crate::operations::protocol::list_drain::list_drain;
use crate::operations::protocol::rollback_rest::{RollbackRestParams, rollback_rest};
use crate::operations::protocol::rollback_tsp::{RollbackTspParams, rollback_tsp};
use crate::operations::protocol::rollback_webauthn::{RollbackWebauthnParams, rollback_webauthn};
use crate::operations::protocol::update_rest::{UpdateRestParams, update_rest};
use crate::operations::protocol::update_tsp::{UpdateTspParams, update_tsp};
use crate::operations::protocol::update_webauthn::{UpdateWebauthnParams, update_webauthn};
fn arrival_transport() -> DisableTransport {
match super::transport::current() {
super::transport::TransportConfidentiality::EndToEnd => DisableTransport::Didcomm,
super::transport::TransportConfidentiality::HopByHop => DisableTransport::Rest,
}
}
fn need_url(url: Option<String>) -> Result<String, AppError> {
url.ok_or_else(|| AppError::Conflict("this service needs `config.url`; none was sent".into()))
}
fn need_mediator(mediator_did: Option<String>) -> Result<String, AppError> {
mediator_did.ok_or_else(|| {
AppError::Conflict("this service needs `config.mediatorDid`; none was sent".into())
})
}
fn mutation(
log_entry_version_id: String,
vta_did: String,
serverless: bool,
drain_until: Option<chrono::DateTime<chrono::Utc>>,
draining_mediator: Option<String>,
) -> Result<spec::enable::v1_0::ServiceMutationResult, AppError> {
built(
spec::enable::v1_0::ServiceMutationResult::builder()
.log_entry_version_id(
spec::enable::v1_0::ServiceMutationResultLogEntryVersionId::try_from(
log_entry_version_id,
)
.map_err(|_| {
AppError::Internal("operation returned an empty log entry id".into())
})?,
)
.effective_at(chrono::Utc::now())
.drain_until(drain_until)
.draining_mediator(draining_mediator)
.vta_did((!vta_did.is_empty()).then_some(vta_did))
.serverless(serverless),
)
}
fn retype<T: serde::de::DeserializeOwned>(v: impl serde::Serialize) -> Result<T, AppError> {
serde_json::to_value(v)
.and_then(serde_json::from_value)
.map_err(|e| AppError::Internal(format!("service result re-type: {e}")))
}
macro_rules! op {
($doc:expr, $call:expr) => {
match $call.await {
Ok(r) => r,
Err(e) => return app_error_to_reject($doc, AppError::Conflict(e.to_string())),
}
};
}
pub(super) async fn handle_enable(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let req: spec::enable::v1_0::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let resolver = match resolver(state_) {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
let deps = ServiceOpDeps::from_app_state(state_, &resolver);
use spec::enable::v1_0::ServiceKind as K;
let setup_acl = match setup_acl_from_ext(req.ext.as_ref()) {
Ok(value) => value,
Err(e) => return app_error_to_reject(&doc, e),
};
let result = match req.service {
K::Rest => {
let url = match need_url(req.config.url.clone().map(String::from)) {
Ok(u) => u,
Err(e) => return app_error_to_reject(&doc, e),
};
let r = op!(
&doc,
enable_rest(
&deps,
auth,
EnableRestParams { url },
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Webauthn => {
let url = match need_url(req.config.url.clone().map(String::from)) {
Ok(u) => u,
Err(e) => return app_error_to_reject(&doc, e),
};
let r = op!(
&doc,
enable_webauthn(
&deps,
auth,
EnableWebauthnParams { url },
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Tsp => {
let mediator_did =
match need_mediator(req.config.mediator_did.clone().map(String::from)) {
Ok(m) => m,
Err(e) => return app_error_to_reject(&doc, e),
};
let r = op!(
&doc,
enable_tsp(
&deps,
auth,
EnableTspParams { mediator_did },
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Didcomm => {
let mediator_did =
match need_mediator(req.config.mediator_did.clone().map(String::from)) {
Ok(m) => m,
Err(e) => return app_error_to_reject(&doc, e),
};
let force = req.config.force.unwrap_or(false);
let handshake_timeout = std::time::Duration::from_secs(
req.config.handshake_timeout_secs.map_or(10, u64::from),
);
#[cfg(feature = "didcomm")]
if !force
&& let Err(error) = run_first_enable_handshake(
state_,
&resolver,
&mediator_did,
handshake_timeout,
setup_acl,
)
.await
{
return app_error_to_reject(&doc, AppError::Conflict(error.to_string()));
}
let params = crate::operations::protocol::enable_didcomm::EnableDidcommParams {
mediator_did,
setup_acl,
force,
handshake_timeout,
};
let prover = crate::messaging::handshake::AlwaysOkProver;
let r = op!(
&doc,
crate::operations::protocol::enable_didcomm::enable_didcomm(
&deps,
&prover,
auth,
params,
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
_ => {
return app_error_to_reject(
&doc,
AppError::Validation(
"unrecognised service kind — this VTA does not serve a \
transport added to the registry after it was built"
.into(),
),
);
}
};
match result {
Ok(result) => match built::<spec::enable::v1_0::Response, _>(
spec::enable::v1_0::Response::builder().result(result),
) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
},
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_update(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let req: spec::update::v1_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
use spec::update::v1_1::ServiceKind as K;
if req.drain_ttl_secs.is_some() && matches!(req.service, K::Rest | K::Webauthn) {
return app_error_to_reject(
&doc,
AppError::Validation(
"drainTtlSecs applies only to a mediated transport (didcomm, tsp)".into(),
),
);
}
let resolver = match resolver(state_) {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
let deps = ServiceOpDeps::from_app_state(state_, &resolver);
let setup_acl = match setup_acl_from_ext(req.ext.as_ref()) {
Ok(value) => value,
Err(e) => return app_error_to_reject(&doc, e),
};
let result = match req.service {
K::Rest => {
let url = match need_url(req.config.url.clone().map(String::from)) {
Ok(u) => u,
Err(e) => return app_error_to_reject(&doc, e),
};
let r = op!(
&doc,
update_rest(
&deps,
auth,
UpdateRestParams { url },
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Webauthn => {
let url = match need_url(req.config.url.clone().map(String::from)) {
Ok(u) => u,
Err(e) => return app_error_to_reject(&doc, e),
};
let r = op!(
&doc,
update_webauthn(
&deps,
auth,
UpdateWebauthnParams { url },
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Tsp => {
let mediator_did =
match need_mediator(req.config.mediator_did.clone().map(String::from)) {
Ok(m) => m,
Err(e) => return app_error_to_reject(&doc, e),
};
let r = op!(
&doc,
update_tsp(
&deps,
auth,
UpdateTspParams { mediator_did },
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Didcomm => {
let mediator_did =
match need_mediator(req.config.mediator_did.clone().map(String::from)) {
Ok(m) => m,
Err(e) => return app_error_to_reject(&doc, e),
};
let force = req.config.force.unwrap_or(false);
#[cfg(feature = "didcomm")]
let live_prover = live_didcomm_prover(state_).await;
let always_ok = crate::messaging::handshake::AlwaysOkProver;
let prover: &(dyn crate::messaging::handshake::ListenerProver + Send + Sync) = {
#[cfg(feature = "didcomm")]
if force {
&always_ok
} else if let Some(prover) = live_prover.as_ref() {
prover
} else {
return app_error_to_reject(
&doc,
AppError::ServiceError {
status: axum::http::StatusCode::BAD_GATEWAY,
message: "DIDComm messaging is not running; cannot prove the candidate mediator"
.into(),
},
);
}
#[cfg(not(feature = "didcomm"))]
if force {
&always_ok
} else {
return app_error_to_reject(
&doc,
AppError::ServiceError {
status: axum::http::StatusCode::BAD_GATEWAY,
message: "DIDComm support is not compiled in; cannot prove the candidate mediator"
.into(),
},
);
}
};
let floor = crate::operations::protocol::disable_didcomm::MIN_DRAIN_TTL_OVER_DIDCOMM;
let transport = arrival_transport();
let drain_ttl = match req.drain_ttl_secs {
None => floor,
Some(secs) => {
let asked = std::time::Duration::from_secs(secs);
match transport {
crate::operations::protocol::disable_didcomm::DisableTransport::Didcomm => {
asked.max(floor)
}
_ => asked,
}
}
};
let params = crate::operations::protocol::update_didcomm::UpdateDidcommParams {
new_mediator_did: mediator_did,
drain_ttl,
setup_acl,
force,
handshake_timeout: std::time::Duration::from_secs(
req.config.handshake_timeout_secs.map_or(10, u64::from),
),
audit_kind: crate::operations::protocol::update_didcomm::MigrateAuditKind::Forward,
transport,
};
let r = op!(
&doc,
crate::operations::protocol::update_didcomm::update_didcomm(
&deps,
prover,
auth,
params,
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(
r.new_version_id,
r.vta_did,
r.serverless,
Some(r.drains_until),
Some(r.prior_mediator_did),
)
}
_ => {
return app_error_to_reject(
&doc,
AppError::Validation(
"unrecognised service kind — this VTA does not serve a \
transport added to the registry after it was built"
.into(),
),
);
}
};
match result {
Ok(result) => match retype::<spec::update::v1_0::ServiceMutationResult>(result) {
Ok(result) => {
match built::<spec::update::v1_0::Response, _>(
spec::update::v1_0::Response::builder().result(result),
) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
}
}
Err(e) => app_error_to_reject(&doc, e),
},
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_disable(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let req: spec::disable::v1_0::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let resolver = match resolver(state_) {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
let deps = ServiceOpDeps::from_app_state(state_, &resolver);
use spec::disable::v1_0::ServiceKind as K;
let result = match req.service {
K::Rest => {
let r = op!(
&doc,
disable_rest(
&deps,
auth,
DisableRestParams,
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Webauthn => {
let r = op!(
&doc,
disable_webauthn(
&deps,
auth,
DisableWebauthnParams {},
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Tsp => {
let r = op!(
&doc,
disable_tsp(
&deps,
auth,
DisableTspParams,
OpContext::Direct,
TRANSPORT_TRUST_TASK
)
);
mutation(r.new_version_id, r.vta_did, r.serverless, None, None)
}
K::Didcomm => {
let params = DisableDidcommParams {
drain_ttl: std::time::Duration::from_secs(req.drain_ttl_secs.unwrap_or(0)),
transport: arrival_transport(),
};
let r = op!(
&doc,
disable_didcomm(&deps, auth, params, OpContext::Direct, TRANSPORT_TRUST_TASK)
);
let draining = r.drains_until.map(|_| r.prior_mediator_did.clone());
mutation(
r.new_version_id,
r.vta_did,
r.serverless,
r.drains_until,
draining,
)
}
_ => {
return app_error_to_reject(
&doc,
AppError::Validation(
"unrecognised service kind — this VTA does not serve a \
transport added to the registry after it was built"
.into(),
),
);
}
};
match result {
Ok(result) => match retype::<spec::disable::v1_0::ServiceMutationResult>(result) {
Ok(result) => {
match built::<spec::disable::v1_0::Response, _>(
spec::disable::v1_0::Response::builder().result(result),
) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
}
}
Err(e) => app_error_to_reject(&doc, e),
},
Err(e) => app_error_to_reject(&doc, e),
}
}
macro_rules! rollback_kind {
($m:ident, $k:expr) => {{
use crate::operations::protocol::$m::RollbackKind as K;
use spec::rollback::v1_0::RollbackResultKind as Out;
match $k {
K::Disabled => Out::Disabled,
K::Enabled => Out::Enabled,
K::Updated => Out::Updated,
K::NoOp => Out::NoOp,
}
}};
}
fn rollback_result(
kind: spec::rollback::v1_0::RollbackResultKind,
new_version_id: Option<String>,
vta_did: String,
serverless: bool,
draining_mediator: Option<String>,
) -> spec::rollback::v1_0::RollbackResult {
let wrote_an_entry = new_version_id.is_some();
spec::rollback::v1_0::RollbackResult::builder()
.kind(kind)
.log_entry_version_id(new_version_id)
.effective_at(wrote_an_entry.then(chrono::Utc::now))
.draining_mediator(draining_mediator)
.vta_did((!vta_did.is_empty()).then_some(vta_did))
.serverless(serverless)
.try_into()
.expect("every required member of the rollback result is set above")
}
pub(super) async fn handle_rollback(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let req: spec::rollback::v1_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
use spec::rollback::v1_1::ServiceKind as K;
if req.drain_ttl_secs.is_some() && matches!(req.service, K::Rest | K::Webauthn) {
return app_error_to_reject(
&doc,
AppError::Validation(
"drainTtlSecs applies only to a mediated transport (didcomm, tsp)".into(),
),
);
}
let resolver = match resolver(state_) {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
let deps = ServiceOpDeps::from_app_state(state_, &resolver);
let result = match req.service {
K::Rest => {
let r = op!(
&doc,
rollback_rest(&deps, auth, RollbackRestParams, TRANSPORT_TRUST_TASK)
);
rollback_result(
rollback_kind!(rollback_rest, r.kind),
r.new_version_id,
r.vta_did,
r.serverless,
None,
)
}
K::Webauthn => {
let r = op!(
&doc,
rollback_webauthn(&deps, auth, RollbackWebauthnParams, TRANSPORT_TRUST_TASK)
);
rollback_result(
rollback_kind!(rollback_webauthn, r.kind),
r.new_version_id,
r.vta_did,
r.serverless,
None,
)
}
K::Tsp => {
let r = op!(
&doc,
rollback_tsp(&deps, auth, RollbackTspParams, TRANSPORT_TRUST_TASK)
);
rollback_result(
rollback_kind!(rollback_tsp, r.kind),
r.new_version_id,
r.vta_did,
r.serverless,
None,
)
}
K::Didcomm => {
let floor = crate::operations::protocol::disable_didcomm::MIN_DRAIN_TTL_OVER_DIDCOMM;
let transport = arrival_transport();
let drain_ttl = match req.drain_ttl_secs {
None => floor,
Some(secs) => {
let asked = std::time::Duration::from_secs(secs);
match transport {
crate::operations::protocol::disable_didcomm::DisableTransport::Didcomm => {
asked.max(floor)
}
_ => asked,
}
}
};
let params = crate::operations::protocol::rollback_didcomm::RollbackDidcommParams {
drain_ttl,
transport,
};
let prover = crate::messaging::handshake::AlwaysOkProver;
let r = op!(
&doc,
crate::operations::protocol::rollback_didcomm::rollback_didcomm(
&deps,
&prover,
auth,
params,
TRANSPORT_TRUST_TASK
)
);
let draining = r.draining_mediator.clone();
rollback_result(
rollback_kind!(rollback_didcomm, r.kind),
r.new_version_id,
r.vta_did,
r.serverless,
draining,
)
}
_ => {
return app_error_to_reject(
&doc,
AppError::Validation(
"unrecognised service kind — this VTA does not serve a \
transport added to the registry after it was built"
.into(),
),
);
}
};
match built::<spec::rollback::v1_0::Response, _>(
spec::rollback::v1_0::Response::builder().result(result),
) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_drain_list(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let _req: spec::drain::list::v1_0::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
match list_drain(&state_.config, &state_.drains_ks, auth).await {
Ok(body) => {
let entries = body
.entries
.into_iter()
.map(|e| {
built(
spec::drain::list::v1_0::DrainEntry::builder()
.mediator_did(
spec::drain::list::v1_0::DrainEntryMediatorDid::try_from(
e.mediator_did,
)
.map_err(|_| {
AppError::Internal(
"drain entry has an empty mediator did".into(),
)
})?,
)
.endpoint(e.endpoint)
.drains_until(
e.drains_until
.parse::<chrono::DateTime<chrono::Utc>>()
.map_err(|_| {
AppError::Internal(
"drain entry has an unparseable deadline".into(),
)
})?,
),
)
})
.collect::<Result<Vec<_>, AppError>>();
match entries {
Ok(entries) => match built::<spec::drain::list::v1_0::Response, _>(
spec::drain::list::v1_0::Response::builder().entries(entries),
) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
},
Err(e) => app_error_to_reject(&doc, e),
}
}
Err(e) => app_error_to_reject(&doc, AppError::Conflict(e.to_string())),
}
}
pub(super) async fn handle_drain_cancel(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let req: spec::drain::cancel::v1_0::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let params = DrainCancelParams {
mediator_did: req.mediator_did.to_string(),
};
match drain_cancel(
&state_.config,
&state_.drains_ks,
&state_.mediator_registry,
&state_.telemetry,
auth,
params,
TRANSPORT_TRUST_TASK,
)
.await
{
Ok(r) => match built::<spec::drain::cancel::v1_0::Response, _>(
spec::drain::cancel::v1_0::Response::builder().mediator_did(r.mediator_did),
) {
Ok(resp) => success_response(&doc, resp),
Err(e) => app_error_to_reject(&doc, e),
},
Err(e) => app_error_to_reject(&doc, AppError::Conflict(e.to_string())),
}
}
pub(super) async fn handle_report(
state_: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let req: spec::report::v0_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
if let (Some(since), Some(until)) = (req.since, req.until)
&& since > until
{
return super::helpers::reject_declared(
&doc,
spec::report::v0_1::error_codes::INVALID_WINDOW,
"`since` is after `until`",
);
}
let report = match crate::operations::protocol::report::mediator_report(
&state_.telemetry,
auth,
crate::operations::protocol::report::ReportParams {
since: req.since,
until: req.until,
},
)
.await
{
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, AppError::Internal(e.to_string())),
};
let ts =
|t: chrono::DateTime<chrono::Utc>| t.to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let body = serde_json::json!({
"since": report.since.map(ts),
"until": ts(report.until),
"mediators": report.mediators.iter().map(|m| serde_json::json!({
"mediatorDid": m.mediator_did,
"inboundCount": m.inbound_count,
"firstSeen": ts(m.first_seen),
"lastSeen": ts(m.last_seen),
})).collect::<Vec<_>>(),
"senders": report.senders.iter().map(|s| serde_json::json!({
"senderDid": s.sender_did,
"lastSeenMediator": s.last_seen_mediator,
"lastSeenAt": ts(s.last_seen_at),
})).collect::<Vec<_>>(),
});
let mut body = body;
if report.since.is_none()
&& let Some(obj) = body.as_object_mut()
{
obj.remove("since");
}
match serde_json::from_value::<spec::report::v0_1::Response>(body) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(
&doc,
AppError::Internal(format!("report does not match its schema: {e}")),
),
}
}
#[cfg(test)]
mod setup_acl_extension_tests {
use super::*;
#[test]
fn reads_setup_acl_from_the_openvtc_extension() {
let payload: spec::update::v1_1::Payload = serde_json::from_value(serde_json::json!({
"service": "didcomm",
"config": { "mediatorDid": "did:web:mediator.example" },
"ext": { "org.openvtc": { "setupAcl": true } }
}))
.unwrap();
assert!(setup_acl_from_ext(payload.ext.as_ref()).unwrap());
}
#[test]
fn rejects_a_non_boolean_setup_acl_extension() {
let payload: spec::update::v1_1::Payload = serde_json::from_value(serde_json::json!({
"service": "didcomm",
"config": { "mediatorDid": "did:web:mediator.example" },
"ext": { "org.openvtc": { "setupAcl": "yes" } }
}))
.unwrap();
assert!(setup_acl_from_ext(payload.ext.as_ref()).is_err());
}
}