use super::helpers::TrustTaskOutcome;
use serde_json::{Value, json};
use trust_tasks_rs::{RejectReason, TrustTask};
use vta_sdk::protocols::backup_management::chunked::{
ALGORITHM_CHUNKED, finalize_import_1_1, get_chunk, initiate_export_1_1, initiate_import_1_1,
put_chunk,
};
use vta_sdk::protocols::backup_management::descriptors::{
AbortBundleBody, CompleteExportBody, FinalizeImportBody, InitiateExportBody, InitiateImportBody,
};
use vti_common::error::AppError;
use crate::auth::AuthClaims;
use crate::operations::backup::{chunked, descriptors};
use crate::server::AppState;
use super::helpers::{
TRANSPORT_TRUST_TASK, app_error_to_reject, parse_payload, reject_with, reject_with_code,
success_response,
};
const INITIATE_EXPORT_SLUG: &str = "vta/backup/initiate-export";
const INITIATE_IMPORT_SLUG: &str = "vta/backup/initiate-import";
fn transport_unavailable(doc: &TrustTask<Value>, slug: &str) -> TrustTaskOutcome {
tracing::warn!(
slug,
"backup descriptor refused: `public_url` is not configured, so the `stream` \
algorithm has no blob URL to publish"
);
let code = trust_tasks_rs::TrustTaskCode::new_extended(slug, "transportUnavailable")
.expect("backup extended code is grammar-valid");
reject_with_code(doc, code, descriptors::TRANSPORT_UNAVAILABLE_MESSAGE, None)
}
async fn initiate_precheck(
state: &AppState,
auth: &AuthClaims,
doc: &TrustTask<Value>,
slug: &str,
) -> Result<(), TrustTaskOutcome> {
auth.require_super_admin()
.map_err(|e| app_error_to_reject(doc, e))?;
if descriptors::blob_transport_base_url(&state.config)
.await
.is_none()
{
return Err(transport_unavailable(doc, slug));
}
Ok(())
}
async fn record_bundle_event(
state: &AppState,
auth: &AuthClaims,
action: &str,
bundle_id: &str,
detail: String,
) {
if let Err(e) = crate::audit::record_with_detail(
&state.audit_sink,
action,
&auth.did,
Some(bundle_id),
"success",
Some(TRANSPORT_TRUST_TASK),
None,
Some(&detail),
)
.await
{
tracing::warn!(error = %e, action, "audit record failed for {action}");
}
}
pub(super) async fn handle_initiate_export(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: InitiateExportBody = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
if let Err(resp) = initiate_precheck(state, auth, &doc, INITIATE_EXPORT_SLUG).await {
return resp;
}
let deps = crate::operations::descriptor_deps_from_app_state(state);
let include_audit = req.include_audit;
match descriptors::initiate_export(&deps, auth, req).await {
Ok(body) => {
record_bundle_event(
state,
auth,
"backup.initiate-export",
&body.descriptor.bundle_id,
format!(
"includeAudit={include_audit} bytes={} expires={}",
body.descriptor.expected_size_bytes, body.descriptor.expires_at
),
)
.await;
success_response(&doc, body)
}
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_complete_export(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: CompleteExportBody = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let deps = crate::operations::descriptor_deps_from_app_state(state);
match descriptors::complete_export(&deps, auth, req).await {
Ok(body) => {
record_bundle_event(
state,
auth,
"backup.complete-export",
&body.bundle_id,
format!("downloaded={}", body.downloaded),
)
.await;
success_response(&doc, body)
}
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_initiate_import(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: InitiateImportBody = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
if let Err(resp) = initiate_precheck(state, auth, &doc, INITIATE_IMPORT_SLUG).await {
return resp;
}
let deps = crate::operations::descriptor_deps_from_app_state(state);
match descriptors::initiate_import(&deps, auth, req).await {
Ok(body) => {
record_bundle_event(
state,
auth,
"backup.initiate-import",
&body.descriptor.bundle_id,
format!(
"sha256={} bytes={} expires={}",
body.descriptor.expected_sha256,
body.descriptor.expected_size_bytes,
body.descriptor.expires_at
),
)
.await;
success_response(&doc, body)
}
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_finalize_import(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: FinalizeImportBody = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
if let Err(e) = chunked::finalize_precheck(&state.backup_bundles_ks, auth, &req.bundle_id).await
{
return match e {
chunked::ChunkedError::App(app) => app_error_to_reject(&doc, app),
chunked::ChunkedError::NotFound => app_error_to_reject(
&doc,
AppError::NotFound(format!("bundle not found: {}", req.bundle_id)),
),
other => app_error_to_reject(&doc, AppError::Conflict(other.to_string())),
};
}
let deps = crate::operations::descriptor_deps_from_app_state(state);
match descriptors::finalize_import(&deps, auth, req).await {
Ok(body) => {
record_bundle_event(
state,
auth,
"backup.finalize-import",
&body.bundle_id,
format!(
"status={} source={} keys={} acls={} contexts={}",
body.status,
body.source_did.as_deref().unwrap_or("unknown"),
body.key_count,
body.acl_count,
body.context_count
),
)
.await;
success_response(&doc, body)
}
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_abort(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: AbortBundleBody = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let deps = crate::operations::descriptor_deps_from_app_state(state);
match descriptors::abort_bundle(&deps, auth, req).await {
Ok(body) => {
record_bundle_event(
state,
auth,
"backup.abort",
&body.bundle_id,
format!("aborted={}", body.aborted),
)
.await;
success_response(&doc, body)
}
Err(e) => app_error_to_reject(&doc, e),
}
}
fn backup_slug(doc: &TrustTask<Value>) -> String {
doc.type_uri
.to_string()
.strip_prefix("https://trusttasks.org/spec/")
.and_then(|rest| rest.rsplit_once('/'))
.map(|(slug, _ver)| slug.to_string())
.unwrap_or_else(|| "vta/backup".to_string())
}
fn chunked_reject(doc: &TrustTask<Value>, err: chunked::ChunkedError) -> TrustTaskOutcome {
use chunked::ChunkedError as E;
let slug = backup_slug(doc);
let message = err.to_string();
let (local, details) = match err {
E::App(e) => return app_error_to_reject(doc, e),
E::RateLimited { retry_after_secs } => {
return reject_with(
doc,
RejectReason::Unavailable {
retry_after: Some(
chrono::Utc::now() + chrono::Duration::seconds(retry_after_secs as i64),
),
},
);
}
E::NotFound => ("notFound", None),
E::TerminalState(_) => ("terminalState", None),
E::ChunkOutOfRange { .. } => ("chunkOutOfRange", None),
E::DigestMismatch {
expected_digest_multibase,
} => (
"digestMismatch",
Some(json!({ "expectedDigestMultibase": expected_digest_multibase })),
),
E::ChunkSizeMismatch { .. } => ("chunkSizeMismatch", None),
E::IncompleteUpload {
missing_count,
missing_indices,
} => (
"incompleteUpload",
Some(json!({ "missingCount": missing_count, "missingIndices": missing_indices })),
),
E::BundleDigestMismatch => ("bundleDigestMismatch", None),
E::BundleTooLarge { .. } => ("bundleTooLarge", None),
E::InvalidManifest(_) => ("invalidManifest", None),
};
let code = trust_tasks_rs::TrustTaskCode::new_extended(&slug, local)
.expect("backup extended code is grammar-valid");
reject_with_code(doc, code, message, details)
}
fn typed<R: serde::de::DeserializeOwned>(value: Value) -> Result<R, AppError> {
serde_json::from_value(value)
.map_err(|e| AppError::Internal(format!("backup response does not fit its schema: {e}")))
}
fn manifest_json(bundle: &chunked::ChunkedBundle) -> Value {
json!({
"bundleId": bundle.bundle_id.to_string(),
"algorithm": ALGORITHM_CHUNKED,
"chunks": {
"chunkSize": bundle.chunk_size,
"chunkCount": bundle.chunk_count,
"chunkDigests": bundle.digests,
},
"expectedSha256": bundle.expected_sha256,
"expectedSizeBytes": bundle.expected_size_bytes,
"expiresAt": bundle.expires_at,
})
}
fn is_chunked(algorithm: Option<&str>) -> bool {
algorithm == Some(ALGORITHM_CHUNKED)
}
pub(super) async fn handle_initiate_export_1_1(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: initiate_export_1_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let algorithm = req.algorithm.as_ref().map(|a| a.as_str());
if !is_chunked(algorithm) {
return handle_initiate_export(state, auth, doc).await;
}
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
let include_audit = req.include_audit.unwrap_or(false);
let deps = crate::operations::descriptor_deps_from_app_state(state);
let bundle = match chunked::initiate_export(
&deps,
auth,
req.password.as_str(),
include_audit,
req.max_chunk_size.map(|s| s.0.max(0) as u64),
)
.await
{
Ok(b) => b,
Err(e) => return chunked_reject(&doc, e),
};
record_bundle_event(
state,
auth,
"backup.initiate-export",
&bundle.bundle_id.to_string(),
format!(
"algorithm={ALGORITHM_CHUNKED} includeAudit={include_audit} bytes={} chunks={} expires={}",
bundle.expected_size_bytes, bundle.chunk_count, bundle.expires_at
),
)
.await;
match typed::<initiate_export_1_1::Response>(json!({
"descriptor": manifest_json(&bundle),
"completionHint": format!(
"Send get-chunk for indices 0 to {}, verify each, then send complete-export.",
bundle.chunk_count - 1
),
})) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_initiate_import_1_1(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: initiate_import_1_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let algorithm = req.algorithm.as_ref().map(|a| a.as_str());
if let Err(e) = auth.require_super_admin() {
return app_error_to_reject(&doc, e);
}
if !is_chunked(algorithm) {
if req.chunks.is_some() {
return chunked_reject(
&doc,
chunked::ChunkedError::InvalidManifest(
"`chunks` is only meaningful with algorithm chunkedTrustTask".into(),
),
);
}
return handle_initiate_import(state, auth, doc).await;
}
let Some(manifest) = req.chunks else {
return chunked_reject(
&doc,
chunked::ChunkedError::InvalidManifest(
"algorithm chunkedTrustTask requires a `chunks` manifest".into(),
),
);
};
let slot = match chunked::initiate_import(
&state.backup_bundles_ks,
auth,
req.expected_sha256.as_str(),
req.expected_size_bytes.0.get(),
manifest.chunk_size.0.max(0) as u64,
manifest.chunk_count.0.get(),
manifest
.chunk_digests
.iter()
.map(|d| d.as_str().to_string())
.collect(),
)
.await
{
Ok(s) => s,
Err(e) => return chunked_reject(&doc, e),
};
record_bundle_event(
state,
auth,
"backup.initiate-import",
&slot.bundle_id.to_string(),
format!(
"algorithm={ALGORITHM_CHUNKED} sha256={} bytes={} chunks={} expires={}",
slot.expected_sha256, slot.expected_size_bytes, slot.chunk_count, slot.expires_at
),
)
.await;
match typed::<initiate_import_1_1::Response>(json!({
"descriptor": manifest_json(&slot),
"completionHint": format!(
"Send put-chunk for indices 0 to {}, then send finalize-import.",
slot.chunk_count - 1
),
})) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_finalize_import_1_1(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: finalize_import_1_1::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
if let Err(e) =
chunked::finalize_precheck(&state.backup_bundles_ks, auth, req.bundle_id.as_str()).await
{
return chunked_reject(&doc, e);
}
handle_finalize_import(state, auth, doc).await
}
pub(super) async fn handle_get_chunk(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: get_chunk::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
let served = match chunked::get_chunk(
&state.backup_bundles_ks,
chunked::ChunkRateLimiter::global(),
auth,
req.bundle_id.as_str(),
req.index.0.max(0) as u64,
)
.await
{
Ok(c) => c,
Err(e) => return chunked_reject(&doc, e),
};
use base64::Engine;
match typed::<get_chunk::Response>(json!({
"bundleId": served.bundle_id.to_string(),
"index": served.index,
"digestMultibase": served.digest_multibase,
"data": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(&served.data),
"expiresAt": served.expires_at,
})) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
}
}
pub(super) async fn handle_put_chunk(
state: &AppState,
auth: &AuthClaims,
doc: TrustTask<Value>,
) -> TrustTaskOutcome {
let req: put_chunk::Payload = match parse_payload(&doc) {
Ok(r) => r,
Err(resp) => return resp,
};
use base64::Engine;
let data = match base64::engine::general_purpose::URL_SAFE_NO_PAD.decode(req.data.as_str()) {
Ok(d) => d,
Err(e) => {
return reject_with(
&doc,
RejectReason::MalformedRequest {
reason: format!("`data` is not unpadded base64url: {e}"),
},
);
}
};
let index = req.index.0.max(0) as u64;
let outcome = match chunked::put_chunk(
&state.backup_bundles_ks,
&state.backup_blob_dir,
chunked::ChunkRateLimiter::global(),
auth,
chunked::ChunkWrite {
bundle_id: req.bundle_id.as_str(),
index,
digest_multibase: req.digest_multibase.as_str(),
data: &data,
},
)
.await
{
Ok(o) => o,
Err(e) => return chunked_reject(&doc, e),
};
match typed::<put_chunk::Response>(json!({
"bundleId": req.bundle_id.as_str(),
"index": index,
"stored": outcome.stored,
"remainingCount": outcome.remaining_count,
"expiresAt": outcome.expires_at,
})) {
Ok(r) => success_response(&doc, r),
Err(e) => app_error_to_reject(&doc, e),
}
}
#[cfg(test)]
mod tests {
use super::*;
use trust_tasks_rs::TypeUri;
fn doc(uri: &str) -> TrustTask<Value> {
let uri: TypeUri = uri.parse().expect("backup uri");
TrustTask::new("urn:uuid:test", uri, serde_json::json!({}))
}
#[test]
fn transport_unavailable_uses_the_specified_extended_code() {
for (uri, slug) in [
(
vta_sdk::trust_tasks::TASK_BACKUP_INITIATE_EXPORT_1_0,
INITIATE_EXPORT_SLUG,
),
(
vta_sdk::trust_tasks::TASK_BACKUP_INITIATE_IMPORT_1_0,
INITIATE_IMPORT_SLUG,
),
] {
assert!(
uri.contains(slug),
"{uri}: the slug constant must match the dispatched URI"
);
let outcome = transport_unavailable(&doc(uri), slug);
let parsed: Value = serde_json::from_slice(&outcome.body).expect("error doc");
assert_eq!(
parsed["payload"]["code"],
format!("{slug}:transportUnavailable"),
"{parsed}"
);
assert_ne!(
outcome.status,
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
"transportUnavailable is not an internal error"
);
assert!(
!parsed["payload"]["message"]
.as_str()
.unwrap_or_default()
.contains("public_url"),
"the wire message must not name configuration: {parsed}"
);
}
}
}