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 {
kind,
enabled,
mediator_did,
url,
drains_until: None,
ext: None,
}
}
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) => success_response(
&doc,
spec::list::v1_0::Response {
services: body.services.iter().map(state).collect(),
ext: None,
},
),
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,
};
match body.services.iter().find(|s| kind_of(s) == want) {
Some(found) => success_response(
&doc,
spec::get::v1_0::Response {
state: serde_json::from_value(
serde_json::to_value(state(found)).unwrap_or(Value::Null),
)
.expect("the two generated ServiceState shapes are identical"),
ext: None,
},
),
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> {
Ok(spec::enable::v1_0::ServiceMutationResult {
log_entry_version_id: log_entry_version_id
.try_into()
.map_err(|_| AppError::Internal("operation returned an empty log entry id".into()))?,
effective_at: chrono::Utc::now(),
drain_until,
draining_mediator,
vta_did: (!vta_did.is_empty()).then_some(vta_did),
serverless,
ext: None,
})
}
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)
}
};
match result {
Ok(result) => success_response(&doc, spec::enable::v1_0::Response { result, ext: None }),
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),
)
}
};
match result {
Ok(result) => match retype(result) {
Ok(result) => {
success_response(&doc, spec::update::v1_0::Response { result, ext: None })
}
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,
)
}
};
match result {
Ok(result) => match retype(result) {
Ok(result) => {
success_response(&doc, spec::disable::v1_0::Response { result, ext: None })
}
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 {
kind,
log_entry_version_id: new_version_id,
effective_at: wrote_an_entry.then(chrono::Utc::now),
drain_until: None,
draining_mediator,
vta_did: (!vta_did.is_empty()).then_some(vta_did),
serverless,
ext: None,
}
}
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,
)
}
};
success_response(&doc, spec::rollback::v1_0::Response { result, ext: None })
}
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| {
Ok(spec::drain::list::v1_0::DrainEntry {
mediator_did: e.mediator_did.try_into().map_err(|_| {
AppError::Internal("drain entry has an empty mediator did".into())
})?,
endpoint: e.endpoint,
drains_until: e.drains_until.parse().map_err(|_| {
AppError::Internal("drain entry has an unparseable deadline".into())
})?,
ext: None,
})
})
.collect::<Result<Vec<_>, AppError>>();
match entries {
Ok(entries) => success_response(
&doc,
spec::drain::list::v1_0::Response { entries, ext: None },
),
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) => success_response(
&doc,
spec::drain::cancel::v1_0::Response {
mediator_did: r.mediator_did,
ext: None,
},
),
Err(e) => app_error_to_reject(&doc, AppError::Conflict(e.to_string())),
}
}