#![allow(clippy::result_large_err)]
use axum::http::StatusCode;
use axum::response::{IntoResponse, Response};
use serde_json::Value;
use trust_tasks_https::status_for_code;
use trust_tasks_rs::{
ErrorPayload, ErrorResponse, RejectReason, StandardCode, TrustTask, TrustTaskCode, TypeUri,
};
use uuid::Uuid;
use vta_sdk::protocols::trust_task_reject_reasons as reasons;
use crate::auth::AuthClaims;
use crate::error::AppError;
use crate::server::AppState;
use vti_common::acl::Capability;
use vta_sdk::protocols::trust_task_reject_details as details;
pub(super) const TRANSPORT_TRUST_TASK: &str = "trust-task";
pub(crate) struct TrustTaskOutcome {
pub(crate) status: StatusCode,
pub(crate) body: Vec<u8>,
}
impl IntoResponse for TrustTaskOutcome {
fn into_response(self) -> Response {
(
self.status,
[(axum::http::header::CONTENT_TYPE, "application/json")],
self.body,
)
.into_response()
}
}
pub(super) fn parse_payload<T: serde::de::DeserializeOwned>(
doc: &TrustTask<Value>,
) -> Result<T, TrustTaskOutcome> {
serde_json::from_value::<T>(doc.payload.clone()).map_err(|e| {
reject_with(
doc,
RejectReason::MalformedRequest {
reason: format!("payload parse: {e}"),
},
)
})
}
fn task_failed_because(message: String, reason: &str) -> RejectReason {
RejectReason::TaskFailed {
reason: message,
details: Some(serde_json::json!({ "reason": reason })),
}
}
pub(super) fn app_error_to_reject(doc: &TrustTask<Value>, err: AppError) -> TrustTaskOutcome {
let message = err.to_string();
let reason = match err {
AppError::Authentication(_) | AppError::Unauthorized(_) | AppError::Forbidden(_) => {
RejectReason::PermissionDenied { reason: message }
}
AppError::Validation(_) | AppError::TrustTaskMalformed(_) | AppError::InvalidCursor => {
RejectReason::MalformedRequest { reason: message }
}
AppError::NotFound(_) => task_failed_because(message, reasons::NOT_FOUND),
AppError::Conflict(_) => task_failed_because(message, reasons::CONFLICT),
AppError::Gone(_) => task_failed_because(message, reasons::GONE),
AppError::ServiceError { status, message }
if status == StatusCode::BAD_GATEWAY || status == StatusCode::GATEWAY_TIMEOUT =>
{
tracing::error!(cause = %message, "trust task failed: an upstream peer did not answer or refused");
task_failed_because(
UPSTREAM_UNAVAILABLE_MESSAGE.to_string(),
reasons::UPSTREAM_UNAVAILABLE,
)
}
AppError::Internal(cause) => {
tracing::error!(cause = %cause, "trust task failed with an internal error");
RejectReason::InternalError {
reason: OPAQUE_INTERNAL_ERROR.to_string(),
}
}
other => {
tracing::error!(cause = %other, "trust task failed with an internal error");
RejectReason::InternalError {
reason: OPAQUE_INTERNAL_ERROR.to_string(),
}
}
};
reject_with(doc, reason)
}
pub(super) const OPAQUE_INTERNAL_ERROR: &str =
"the consumer could not complete this task; the request itself was accepted";
pub(super) const UPSTREAM_UNAVAILABLE_MESSAGE: &str = "a service this VTA depends on did not answer or refused the request; \
the VTA's log names which one and why";
const DETAILS_MAX_JCS_BYTES: usize = 4096;
const DETAILS_MAX_MEMBERS: usize = 16;
fn bound_details(details: Option<Value>) -> Option<Value> {
let details = details?;
let too_many_members = details
.as_object()
.is_some_and(|o| o.len() > DETAILS_MAX_MEMBERS);
let too_large = serde_json_canonicalizer::to_string(&details)
.map(|jcs| jcs.len() > DETAILS_MAX_JCS_BYTES)
.unwrap_or(true);
if too_many_members || too_large {
tracing::warn!(
members = details.as_object().map(serde_json::Map::len),
"error `details` exceeds the framework bound and was dropped; the code still went out"
);
return None;
}
Some(details)
}
pub(super) async fn require_capability(
state: &AppState,
auth: &AuthClaims,
doc: &TrustTask<Value>,
cap: Capability,
what: &str,
) -> Result<(), TrustTaskOutcome> {
let allowed = match vti_common::acl::get_acl_entry(&state.acl_ks, &auth.did).await {
Ok(Some(entry)) => vti_common::acl::entry_has_capability(&entry, cap),
Ok(None) => vti_common::acl::role_has_capability(&auth.role, cap),
Err(e) => {
tracing::error!(
error = %e, did = %auth.did,
"could not read the ACL entry for a capability check; refusing"
);
false
}
};
if allowed {
return Ok(());
}
Err(reject_with(
doc,
RejectReason::PermissionDenied {
reason: format!(
"{what} denied: {} does not carry the {cap:?} capability",
auth.did
),
},
))
}
pub(super) fn reject_with(doc: &TrustTask<Value>, reason: RejectReason) -> TrustTaskOutcome {
let reason = match reason {
RejectReason::TaskFailed { reason, details } => RejectReason::TaskFailed {
reason,
details: bound_details(details),
},
other => other,
};
let routed = doc.reject_with(format!("urn:uuid:{}", Uuid::new_v4()), reason);
error_response(routed)
}
pub(super) fn reject_with_code(
doc: &TrustTask<Value>,
code: TrustTaskCode,
message: impl Into<String>,
details: Option<Value>,
) -> TrustTaskOutcome {
let mut payload = ErrorPayload::new(code).with_message(message);
if let Some(d) = bound_details(details) {
payload = payload.with_details(d);
}
let routed = doc.reject_with(format!("urn:uuid:{}", Uuid::new_v4()), payload);
error_response(routed)
}
pub(super) fn success_response<R: serde::Serialize>(
doc: &TrustTask<Value>,
payload: R,
) -> TrustTaskOutcome {
let response_doc = doc.respond_with(format!("urn:uuid:{}", Uuid::new_v4()), payload);
let body = match serde_json::to_vec(&response_doc) {
Ok(b) => b,
Err(e) => {
tracing::error!(error = %e, "failed to serialise success response doc");
return reject_with(
doc,
RejectReason::InternalError {
reason: format!("response serialisation: {e}"),
},
);
}
};
TrustTaskOutcome {
status: StatusCode::OK,
body,
}
}
#[allow(dead_code)]
pub(super) fn not_implemented_yet(doc: TrustTask<Value>, reason: &str) -> TrustTaskOutcome {
let reject = RejectReason::TaskFailed {
reason: reason.to_string(),
details: None,
};
let routed = doc.reject_with(format!("urn:uuid:{}", Uuid::new_v4()), reject);
error_response(routed)
}
fn unsupported_type(doc: TrustTask<Value>, type_uri: &str) -> TrustTaskOutcome {
reject_with_code(
&doc,
TrustTaskCode::Standard(StandardCode::UnsupportedType),
format!("unsupported type: {type_uri}"),
Some(serde_json::json!({ details::REQUESTED_TYPE: type_uri })),
)
}
fn task_family(type_uri: &str) -> Option<&str> {
type_uri.rsplit_once('/').map(|(family, _version)| family)
}
fn served_versions_of_family(type_uri: &str) -> Vec<&'static str> {
let Some(family) = task_family(type_uri) else {
return Vec::new();
};
let mut served: Vec<&'static str> = super::dispatched_uris()
.into_iter()
.filter(|uri| task_family(uri) == Some(family))
.collect();
served.sort_unstable();
served.dedup();
served
}
pub(super) fn method_not_found(doc: TrustTask<Value>, type_uri: &str) -> TrustTaskOutcome {
let served = served_versions_of_family(type_uri);
if served.is_empty() {
return unsupported_type(doc, type_uri);
}
reject_with_code(
&doc,
TrustTaskCode::Standard(StandardCode::UnsupportedVersion),
format!(
"unsupported version: {type_uri} — this VTA serves {}",
served.join(", ")
),
Some(serde_json::json!({
details::REQUESTED_TYPE: type_uri,
details::SERVED_VERSIONS: served,
})),
)
}
pub(super) fn error_response(err_doc: ErrorResponse) -> TrustTaskOutcome {
let status = StatusCode::from_u16(status_for_code(&err_doc.payload.code))
.unwrap_or(StatusCode::INTERNAL_SERVER_ERROR);
let body = serde_json::to_vec(&err_doc).unwrap_or_else(|_| Vec::new());
TrustTaskOutcome { status, body }
}
fn framework_error_type_uri() -> TypeUri {
"https://trusttasks.org/spec/trust-task-error/0.5"
.parse()
.expect("framework error Type URI parses")
}
pub(super) fn body_parse_error_response(reason: &str) -> TrustTaskOutcome {
malformed_request_response(format!(
"body did not parse as a Trust Task document: {reason}"
))
}
pub(crate) fn malformed_request_response(reason: String) -> TrustTaskOutcome {
let reject = RejectReason::MalformedRequest { reason };
let payload: ErrorPayload = reject.into();
let type_uri: TypeUri = framework_error_type_uri();
let err = ErrorResponse {
id: format!("urn:uuid:{}", Uuid::new_v4()),
thread_id: None,
parent_thread_id: None,
type_uri,
issuer: None,
recipient: None,
issued_at: Some(chrono::Utc::now()),
expires_at: None,
payload,
context: None,
ceremony: None,
proof: None,
extra: Default::default(),
};
error_response(err)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn doc() -> TrustTask<Value> {
let uri: TypeUri = vta_sdk::trust_tasks::TASK_WEBVH_DIDS_UPDATE_1_0
.parse()
.expect("update uri");
TrustTask::new("urn:uuid:test", uri, json!({}))
}
fn message_of(outcome: TrustTaskOutcome) -> String {
let doc: Value = serde_json::from_slice(&outcome.body).expect("error doc parses");
doc["payload"]["message"]
.as_str()
.expect("payload carries a message")
.to_string()
}
fn details_of(outcome: TrustTaskOutcome) -> Value {
let doc: Value = serde_json::from_slice(&outcome.body).expect("error doc parses");
doc["payload"]["details"].clone()
}
#[test]
fn a_not_found_is_discriminated_from_a_plain_task_failure() {
let outcome = app_error_to_reject(&doc(), AppError::NotFound("policy `x`".into()));
assert_eq!(details_of(outcome)["reason"], reasons::NOT_FOUND);
}
#[test]
fn a_conflict_is_discriminated_from_a_plain_task_failure() {
let outcome = app_error_to_reject(&doc(), AppError::Conflict("already exists".into()));
assert_eq!(details_of(outcome)["reason"], reasons::CONFLICT);
}
#[test]
fn a_gone_is_discriminated_from_a_plain_task_failure() {
let outcome = app_error_to_reject(&doc(), AppError::Gone("carve-out closed".into()));
assert_eq!(details_of(outcome)["reason"], reasons::GONE);
}
#[test]
fn the_reasons_are_distinct() {
let all = [
reasons::NOT_FOUND,
reasons::CONFLICT,
reasons::GONE,
reasons::UPSTREAM_UNAVAILABLE,
];
let mut seen = std::collections::BTreeSet::new();
for r in all {
assert!(seen.insert(r), "`{r}` is used for more than one outcome");
}
}
#[test]
fn task_family_strips_only_the_version_segment() {
assert_eq!(
task_family("scheme://host/spec/a/b/0.3"),
Some("scheme://host/spec/a/b")
);
assert_eq!(
task_family("scheme://host/spec/a/b/c-d/1.0"),
Some("scheme://host/spec/a/b/c-d")
);
assert_eq!(task_family("no-slashes-at-all"), None);
}
#[test]
fn served_versions_come_from_the_dispatch_table() {
let served = crate::trust_tasks::dispatched_uris();
let real = served
.first()
.expect("the dispatch table is not empty")
.to_string();
let family = task_family(&real).expect("a task URI has a family");
let bogus = format!("{family}/99.99");
let found = served_versions_of_family(&bogus);
assert!(
found.contains(&real.as_str()),
"{real} is dispatched, so a bogus version of its family should name it; got {found:?}"
);
}
#[test]
fn an_unknown_family_is_still_unsupported_type() {
let outcome = method_not_found(doc(), "https://trusttasks.org/spec/does-not-exist/9.9");
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(parsed["payload"]["code"], "unsupportedType");
assert!(parsed["payload"]["details"].get("servedVersions").is_none());
}
#[test]
fn a_known_family_at_an_unknown_version_names_what_is_served() {
let real = crate::trust_tasks::dispatched_uris()
.first()
.expect("the dispatch table is not empty")
.to_string();
let family = task_family(&real).expect("a task URI has a family");
let bogus = format!("{family}/99.99");
let outcome = method_not_found(doc(), &bogus);
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(parsed["payload"]["code"], "unsupportedVersion");
let message = parsed["payload"]["message"]
.as_str()
.expect("a message")
.to_string();
assert!(
message.contains(&real),
"the message must name the served version; got {message}"
);
assert_eq!(parsed["payload"]["details"]["servedVersions"][0], real);
}
#[test]
fn an_extended_code_survives_to_the_wire() {
let code: TrustTaskCode = "provision/integration:contextRequired"
.parse()
.expect("a legal extended code");
let outcome = reject_with_code(
&doc(),
code,
"which context?",
Some(json!({ "candidates": ["a", "b"] })),
);
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(
parsed["payload"]["code"],
"provision/integration:contextRequired"
);
assert_eq!(parsed["payload"]["details"]["candidates"][1], "b");
}
#[test]
fn an_extended_code_rejection_bounds_its_details_too() {
let code: TrustTaskCode = "provision/integration:contextRequired"
.parse()
.expect("a legal extended code");
let huge = json!({ "explanation": "x".repeat(DETAILS_MAX_JCS_BYTES + 1) });
let outcome = reject_with_code(&doc(), code, "which context?", Some(huge));
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(
parsed["payload"]["code"], "provision/integration:contextRequired",
"an oversized annex must never cost the code: {parsed}"
);
assert!(
parsed["payload"]["details"].is_null(),
"the oversized details should have been dropped: {parsed}"
);
}
#[test]
fn an_oversized_details_is_dropped_but_the_code_survives() {
let huge = serde_json::json!({ "explanation": "x".repeat(DETAILS_MAX_JCS_BYTES + 1) });
let outcome = reject_with(
&doc(),
RejectReason::TaskFailed {
reason: "policy denied".into(),
details: Some(huge),
},
);
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(
parsed["payload"]["code"], "taskFailed",
"an oversized annex must never cost the code: {parsed}"
);
assert!(
parsed["payload"].get("details").is_none_or(Value::is_null),
"the oversized details must not go out: {parsed}"
);
}
#[test]
fn a_details_with_too_many_members_is_dropped() {
let mut wide = serde_json::Map::new();
for i in 0..=DETAILS_MAX_MEMBERS {
wide.insert(format!("k{i}"), serde_json::json!(1));
}
let outcome = reject_with(
&doc(),
RejectReason::TaskFailed {
reason: "policy denied".into(),
details: Some(Value::Object(wide)),
},
);
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(parsed["payload"]["code"], "taskFailed");
assert!(parsed["payload"].get("details").is_none_or(Value::is_null));
}
#[test]
fn a_small_details_still_goes_out() {
let outcome = reject_with(
&doc(),
RejectReason::TaskFailed {
reason: "policy denied".into(),
details: Some(serde_json::json!({ "reason": "auth:consent_required" })),
},
);
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(
parsed["payload"]["details"]["reason"], "auth:consent_required",
"{parsed}"
);
}
#[test]
fn an_internal_error_reveals_no_internal_state() {
let secret = "log entry has no update_keys";
let message = message_of(app_error_to_reject(
&doc(),
AppError::Internal(secret.into()),
));
assert!(
!message.contains(secret),
"the cause must reach the operator's log, never the wire: {message}"
);
assert!(
message.contains(OPAQUE_INTERNAL_ERROR),
"the producer still needs to be told the failure was not its \
document's doing: {message}"
);
assert!(
!message.contains("internal error: internal error"),
"{message}"
);
}
#[test]
fn the_catch_all_arm_is_opaque_too() {
let message = message_of(app_error_to_reject(
&doc(),
AppError::SecretStore("vault backend at 10.0.0.7 refused the token".into()),
));
assert!(!message.contains("10.0.0.7"), "{message}");
assert!(!message.contains("vault backend"), "{message}");
assert!(message.contains(OPAQUE_INTERNAL_ERROR), "{message}");
}
#[test]
fn an_upstream_failure_is_named_as_one_without_its_cause() {
for status in [StatusCode::BAD_GATEWAY, StatusCode::GATEWAY_TIMEOUT] {
let err = || AppError::ServiceError {
status,
message: "bad gateway: `did:webvh:QmHost:dids.example` at \
https://10.0.0.7/trust-tasks did not answer"
.into(),
};
let parsed: Value = serde_json::from_slice(&app_error_to_reject(&doc(), err()).body)
.expect("error doc parses");
assert_eq!(parsed["payload"]["code"], "taskFailed", "{status}");
assert_eq!(
parsed["payload"]["details"]["reason"],
reasons::UPSTREAM_UNAVAILABLE,
"{status}"
);
let message = message_of(app_error_to_reject(&doc(), err()));
assert!(message.contains(UPSTREAM_UNAVAILABLE_MESSAGE), "{message}");
assert!(!message.contains("10.0.0.7"), "{message}");
assert!(!message.contains("QmHost"), "{message}");
assert!(!message.contains(OPAQUE_INTERNAL_ERROR), "{message}");
}
}
#[test]
fn other_service_errors_stay_internal() {
let message = message_of(app_error_to_reject(
&doc(),
AppError::ServiceError {
status: StatusCode::INTERNAL_SERVER_ERROR,
message: "key derivation failed at m/26'/2'".into(),
},
));
assert!(message.contains(OPAQUE_INTERNAL_ERROR), "{message}");
assert!(!message.contains("m/26'"), "{message}");
}
#[test]
fn unrouted_and_routed_errors_agree_on_the_type_uri() {
let routed = doc().reject_with(
"urn:uuid:routed",
RejectReason::InternalError {
reason: "probe".into(),
},
);
assert_eq!(
framework_error_type_uri(),
routed.type_uri,
"the unrouted body-parse error names a different document type than \
the framework stamps on a routed rejection"
);
}
#[test]
fn the_body_parse_error_goes_out_as_a_framework_error_document() {
let outcome = body_parse_error_response("not json");
let doc: Value = serde_json::from_slice(&outcome.body).expect("error doc parses");
assert_eq!(
doc["type"].as_str().expect("type present"),
framework_error_type_uri().to_string()
);
}
#[test]
fn a_not_found_keeps_the_cause_the_operator_needs() {
let message = message_of(app_error_to_reject(
&doc(),
AppError::NotFound("SCID QmNope not found".into()),
));
assert!(message.contains("SCID QmNope not found"), "{message}");
}
#[test]
fn a_gone_is_a_task_failure_not_an_internal_error() {
let message = message_of(app_error_to_reject(
&doc(),
AppError::Gone("carve-out has already been used".into()),
));
assert!(
message.contains("carve-out has already been used"),
"{message}"
);
assert!(
!message.starts_with("internal error"),
"a consumed resource must not report as a server fault: {message}"
);
}
}