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 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 Box::pin($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 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 params = crate::operations::protocol::enable_didcomm::EnableDidcommParams {
mediator_did,
force: req.config.force.unwrap_or(false),
handshake_timeout: std::time::Duration::from_secs(
req.config.handshake_timeout_secs.map_or(10, u64::from),
),
};
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_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::update::v1_0::ServiceKind as K;
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 prover = crate::messaging::handshake::AlwaysOkProver;
let params = crate::operations::protocol::update_didcomm::UpdateDidcommParams {
new_mediator_did: mediator_did,
drain_ttl: crate::operations::protocol::disable_didcomm::MIN_DRAIN_TTL_OVER_DIDCOMM,
force: req.config.force.unwrap_or(false),
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: arrival_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_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::rollback::v1_0::ServiceKind as K;
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 params = crate::operations::protocol::rollback_didcomm::RollbackDidcommParams {
drain_ttl: crate::operations::protocol::disable_didcomm::MIN_DRAIN_TTL_OVER_DIDCOMM,
transport: arrival_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())),
}
}