#![allow(clippy::result_large_err)]
use super::helpers::TrustTaskOutcome;
use serde_json::Value;
use trust_tasks_rs::TrustTask;
use vta_sdk::protocols::credential_exchange::{
self as cx, IssueBody, OfferBody, PendingApproveBody, PendingDenyBody, PendingDenyResponse,
PendingListResponse, PendingPresentationSummary, QueryBody, RequestedCredentialSummary,
};
use crate::acl::Role;
use crate::audit::audit;
use crate::auth::AuthClaims;
use crate::error::AppError;
use crate::operations::credential_exchange::{
ConsentPolicy, PresentOutcome, RequestedCredential, approve_pending_presentation,
build_credential_request_for_offer, defer_presentation, deny_pending_presentation, pending,
present_query, receive_issued_credential,
};
use crate::server::AppState;
use super::helpers::{acknowledge, app_error_to_reject, parse_payload, silence, success_response};
const EXCHANGE_DELIVER_BY: std::time::Duration = std::time::Duration::from_secs(24 * 3600);
pub(crate) fn is_counterparty_task(type_uri: &str) -> bool {
type_uri == cx::OFFER || type_uri == cx::ISSUE || type_uri == cx::QUERY
}
fn thread_of(doc: &TrustTask<Value>) -> String {
doc.thread_id.clone().unwrap_or_else(|| doc.id.clone())
}
async fn own_authority(state: &AppState) -> AuthClaims {
AuthClaims {
did: state
.config
.read()
.await
.vta_did
.clone()
.unwrap_or_else(|| "vta:self".into()),
role: Role::Admin,
allowed_contexts: Vec::new(),
..Default::default()
}
}
#[cfg_attr(
not(any(feature = "didcomm", feature = "tsp")),
allow(unused_variables)
)]
async fn push_step(
state: &AppState,
recipient: &str,
type_uri: &str,
payload: Value,
thread: &str,
) -> Result<String, AppError> {
#[cfg(any(feature = "didcomm", feature = "tsp"))]
{
let vta_did = state
.config
.read()
.await
.vta_did
.clone()
.ok_or_else(|| AppError::Internal("VTA DID not configured".into()))?;
let mut doc =
vti_common::capability_client::build_document(&vta_did, recipient, type_uri, payload);
doc.thread_id = Some(thread.to_string());
let id = doc.id.clone();
let mut doc_value = serde_json::to_value(&doc)
.map_err(|e| AppError::Internal(format!("serialise {type_uri} document: {e}")))?;
if !super::sign_outbound_request(state, &mut doc_value).await {
return Err(AppError::Internal(format!(
"{type_uri} could not be signed, so it was not sent"
)));
}
crate::messaging::push::push_trust_task(state, recipient, doc_value, EXCHANGE_DELIVER_BY)
.await?;
Ok(id)
}
#[cfg(not(any(feature = "didcomm", feature = "tsp")))]
Err(AppError::Internal(format!(
"{type_uri} cannot be sent: this build has no messaging transport"
)))
}
pub(super) async fn handle_offer(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let body: OfferBody = match parse_payload(&doc) {
Ok(b) => b,
Err(resp) => return resp,
};
let issuer = auth.did.clone();
let Some(subject_did) = state.config.read().await.credential_holder_did.clone() else {
tracing::info!(
from = %issuer,
"credential offer received but no credential_holder_did configured — declining"
);
return app_error_to_reject(
&doc,
AppError::Forbidden(
"this VTA does not accept unsolicited credential offers \
(no credential_holder_did configured)"
.into(),
),
);
};
let request = match build_credential_request_for_offer(
&state.keys_ks,
&state.contexts_ks,
&state.seed_store,
&state.audit_sink,
&own_authority(state).await,
&body.credential_offer,
&subject_did,
chrono::Utc::now(),
)
.await
{
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
let payload = match serde_json::to_value(&request) {
Ok(v) => v,
Err(e) => {
return app_error_to_reject(
&doc,
AppError::Internal(format!("request serialise: {e}")),
);
}
};
if let Err(e) = push_step(state, &issuer, cx::REQUEST, payload, &thread_of(&doc)).await {
return app_error_to_reject(&doc, e);
}
tracing::info!(from = %issuer, subject = %subject_did, "answered credential offer with a request");
acknowledge(&doc)
}
pub(super) async fn handle_issue(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let body: IssueBody = match parse_payload(&doc) {
Ok(b) => b,
Err(resp) => return resp,
};
let stored = match receive_issued_credential(
&state.vault_ks,
&body,
state.did_resolver.as_ref(),
Some(auth.did.clone()),
chrono::Utc::now(),
)
.await
{
Ok(s) => s,
Err(e) => return app_error_to_reject(&doc, e),
};
tracing::info!(
credential_id = %stored.id,
format = ?stored.format,
from = %auth.did,
"received issued credential into vault"
);
acknowledge(&doc)
}
pub(super) async fn handle_query(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let body: QueryBody = match parse_payload(&doc) {
Ok(b) => b,
Err(resp) => return resp,
};
let verifier_did = auth.did.clone();
let thread = thread_of(&doc);
let policy = ConsentPolicy::trusting(
state
.config
.read()
.await
.trusted_presentation_verifiers
.clone(),
);
let outcome = match present_query(
&state.vault_ks,
&state.keys_ks,
&state.contexts_ks,
&state.seed_store,
&state.audit_sink,
&own_authority(state).await,
&body,
&verifier_did,
&policy,
state.status_list_resolver.as_deref(),
chrono::Utc::now(),
)
.await
{
Ok(o) => o,
Err(e) => return app_error_to_reject(&doc, e),
};
match outcome {
PresentOutcome::Presented(present) => {
let payload = match serde_json::to_value(&present) {
Ok(v) => v,
Err(e) => {
return app_error_to_reject(
&doc,
AppError::Internal(format!("present serialise: {e}")),
);
}
};
if let Err(e) = push_step(state, &verifier_did, cx::PRESENT, payload, &thread).await {
return app_error_to_reject(&doc, e);
}
tracing::info!(verifier = %verifier_did, "presented a vp_token");
acknowledge(&doc)
}
PresentOutcome::ConsentRequired {
requested, purpose, ..
} => {
let requested_count = requested.len();
if let Err(e) = defer_presentation(
&state.vault_ks,
&thread,
&verifier_did,
requested,
&body,
chrono::Utc::now(),
)
.await
{
return app_error_to_reject(&doc, e);
}
tracing::info!(
verifier = %verifier_did,
approval_id = %thread,
requested = requested_count,
%purpose,
"credential query deferred — holder consent required (pending approval persisted)"
);
silence()
}
}
}
fn summarize(record: pending::PendingPresentation) -> PendingPresentationSummary {
PendingPresentationSummary {
id: record.id,
verifier_did: record.verifier_did,
requested: record
.requested
.into_iter()
.map(summarize_requested)
.collect(),
purpose: record.purpose,
created_at: record.created_at,
expires_at: record.expires_at,
}
}
fn summarize_requested(r: RequestedCredential) -> RequestedCredentialSummary {
RequestedCredentialSummary {
credential_query_id: r.credential_query_id,
credential_id: r.credential_id,
claims: r.claims,
}
}
pub(super) async fn handle_pending_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 records = match pending::list(&state.vault_ks).await {
Ok(r) => r,
Err(e) => return app_error_to_reject(&doc, e),
};
let now = chrono::Utc::now();
let pending_out: Vec<PendingPresentationSummary> = records
.into_iter()
.filter(|r| r.status == pending::PendingStatus::Pending && r.expires_at > now)
.map(summarize)
.collect();
audit!(
"credential-exchange.pending-list",
actor = &auth.did,
resource = "pending-present",
outcome = "success"
);
success_response(
&doc,
PendingListResponse {
pending: pending_out,
},
)
}
pub(super) async fn handle_pending_approve(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let body: PendingApproveBody = match parse_payload(&doc) {
Ok(b) => b,
Err(resp) => return resp,
};
let verifier_did = match pending::get(&state.vault_ks, &body.id).await {
Ok(Some(record)) => record.verifier_did,
Ok(None) => {
return app_error_to_reject(
&doc,
AppError::NotFound(format!("no pending presentation `{}`", body.id)),
);
}
Err(e) => return app_error_to_reject(&doc, e),
};
let present = match approve_pending_presentation(
&state.vault_ks,
&state.keys_ks,
&state.contexts_ks,
&state.seed_store,
&state.audit_sink,
auth,
&body.id,
state.status_list_resolver.as_deref(),
chrono::Utc::now(),
)
.await
{
Ok(p) => p,
Err(e) => return app_error_to_reject(&doc, e),
};
audit!(
"credential-exchange.pending-approve",
actor = &auth.did,
resource = &body.id,
outcome = "success"
);
match serde_json::to_value(&present) {
Ok(payload) => {
if let Err(e) = push_step(state, &verifier_did, cx::PRESENT, payload, &body.id).await {
tracing::warn!(verifier = %verifier_did, pending = %body.id, error = %e, "approved presentation could not be pushed to the verifier; the vp_token is in the response for relay");
}
}
Err(e) => {
tracing::warn!(error = %e, "approved presentation did not serialise for the push")
}
}
success_response(&doc, present)
}
pub(super) async fn handle_pending_deny(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let body: PendingDenyBody = match parse_payload(&doc) {
Ok(b) => b,
Err(resp) => return resp,
};
if let Err(e) = deny_pending_presentation(&state.vault_ks, &body.id).await {
return app_error_to_reject(&doc, e);
}
audit!(
"credential-exchange.pending-deny",
actor = &auth.did,
resource = &body.id,
outcome = "success"
);
success_response(
&doc,
PendingDenyResponse {
id: body.id,
status: "denied".to_string(),
},
)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
const ISSUER: &str = "did:key:zIssuerCounterparty";
fn doc(type_uri: &str, payload: Value) -> TrustTask<Value> {
serde_json::from_value(json!({
"id": format!("urn:uuid:{}", uuid::Uuid::new_v4()),
"type": type_uri,
"issuer": ISSUER,
"recipient": "did:example:vta",
"issuedAt": chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
"payload": payload,
}))
.expect("a Trust Task document")
}
fn body(outcome: &TrustTaskOutcome) -> Value {
serde_json::from_slice(&outcome.body).expect("a JSON document")
}
fn minted_credential() -> String {
use affinidi_sd_jwt::error::SdJwtError;
use affinidi_sd_jwt::signer::JwtSigner;
use base64::Engine;
use base64::engine::general_purpose::URL_SAFE_NO_PAD;
use ed25519_dalek::Signer;
struct Issuer {
key: ed25519_dalek::SigningKey,
kid: String,
}
impl JwtSigner for Issuer {
fn algorithm(&self) -> &str {
"EdDSA"
}
fn key_id(&self) -> Option<&str> {
Some(&self.kid)
}
fn sign_jwt(&self, header: &Value, payload: &Value) -> Result<String, SdJwtError> {
let h = URL_SAFE_NO_PAD.encode(serde_json::to_string(header)?.as_bytes());
let p = URL_SAFE_NO_PAD.encode(serde_json::to_string(payload)?.as_bytes());
let input = format!("{h}.{p}");
Ok(format!(
"{input}.{}",
URL_SAFE_NO_PAD.encode(self.key.sign(input.as_bytes()).to_bytes())
))
}
}
let key = ed25519_dalek::SigningKey::from_bytes(&[9u8; 32]);
let issuer_did =
affinidi_crypto::did_key::ed25519_pub_to_did_key(key.verifying_key().as_bytes());
let signer = Issuer {
key,
kid: format!("{issuer_did}#key-0"),
};
let subject_did = affinidi_crypto::did_key::ed25519_pub_to_did_key(
ed25519_dalek::SigningKey::from_bytes(&[7u8; 32])
.verifying_key()
.as_bytes(),
);
crate::vault::mint::mint_sd_jwt_vc(
&crate::vault::mint::MintRequest {
vct: "https://openvtc.org/credentials/MembershipCredential",
issuer_did: &issuer_did,
subject_did: &subject_did,
claims: &json!({ "givenName": "Alice" }),
disclosable: &["givenName"],
iat: 1_700_000_000,
exp: Some(1_900_000_000),
},
&signer,
)
.expect("mint an SD-JWT-VC")
}
#[tokio::test]
async fn an_issued_credential_is_stored_and_acknowledged() {
let (state, _dir) = crate::test_support::build_signing_test_app_state().await;
let issue = doc(
cx::ISSUE,
json!({ "credential_response": { "credential": minted_credential() } }),
);
let outcome = handle_issue(
&state,
&super::super::ceremony::ceremony_claims(ISSUER),
issue,
)
.await;
let ack = body(&outcome);
assert_eq!(ack["type"], format!("{}#response", cx::ISSUE), "{ack}");
assert_eq!(
ack["payload"],
json!({}),
"an acknowledgement carries nothing"
);
let stored = state
.vault_ks
.prefix_keys("cred:")
.await
.expect("scan the credential store");
assert_eq!(stored.len(), 1, "the credential is in the vault");
}
#[tokio::test]
async fn an_issue_carrying_no_credential_is_refused() {
let (state, _dir) = crate::test_support::build_signing_test_app_state().await;
let outcome = handle_issue(
&state,
&super::super::ceremony::ceremony_claims(ISSUER),
doc(cx::ISSUE, json!({})),
)
.await;
let error = body(&outcome);
assert!(
error["type"]
.as_str()
.is_some_and(|t| t.contains("trust-task-error")),
"{error}"
);
}
#[tokio::test]
async fn an_offer_is_declined_without_a_configured_holder() {
let (state, _dir) = crate::test_support::build_signing_test_app_state().await;
state.config.write().await.credential_holder_did = None;
let offer = doc(
cx::OFFER,
json!({ "credential_offer": {
"credential_issuer": ISSUER,
"credential_configuration_ids": ["VIC"],
"grants": { "urn:ietf:params:oauth:grant-type:pre-authorized_code":
{ "pre-authorized_code": "code" } },
}}),
);
let outcome = handle_offer(
&state,
&super::super::ceremony::ceremony_claims(ISSUER),
offer,
)
.await;
let error = body(&outcome);
assert_eq!(error["payload"]["code"], "permissionDenied", "{error}");
#[cfg(any(feature = "didcomm", feature = "tsp"))]
assert!(
crate::messaging::push::take_pushes(&state).is_empty(),
"a declined offer sends no request"
);
}
}