use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use vti_common::error::AppError;
use vti_common::store::KeyspaceHandle;
const PENDING_PREFIX: &[u8] = b"pending:";
const GRANT_PREFIX: &[u8] = b"grant:";
const WIRE_INDEX_PREFIX: &[u8] = b"wire:";
const DIGEST_DOMAIN: &[u8] = b"vta/task-consent/v1\0";
pub fn payload_digest(type_uri: &str, payload: &serde_json::Value) -> Result<String, AppError> {
digest_with(type_uri, payload, None)
}
pub fn wire_digest(
type_uri: &str,
payload: &serde_json::Value,
challenge: &str,
) -> Result<String, AppError> {
digest_with(type_uri, payload, Some(challenge))
}
fn digest_with(
type_uri: &str,
payload: &serde_json::Value,
challenge: Option<&str>,
) -> Result<String, AppError> {
let canonical = serde_json_canonicalizer::to_string(payload)
.map_err(|e| AppError::Internal(format!("payload JCS canonicalization failed: {e}")))?;
let mut h = Sha256::new();
h.update(DIGEST_DOMAIN);
h.update((type_uri.len() as u64).to_be_bytes());
h.update(type_uri.as_bytes());
h.update((canonical.len() as u64).to_be_bytes());
h.update(canonical.as_bytes());
if let Some(c) = challenge {
h.update(c.as_bytes());
}
Ok(encode_digest_multibase(&h.finalize()))
}
const MULTIHASH_SHA2_256_32: [u8; 2] = [0x12, 0x20];
fn encode_digest_multibase(digest: &[u8]) -> String {
let mut buf = Vec::with_capacity(MULTIHASH_SHA2_256_32.len() + digest.len());
buf.extend_from_slice(&MULTIHASH_SHA2_256_32);
buf.extend_from_slice(digest);
multibase::encode(multibase::Base::Base58Btc, &buf)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PendingTaskConsent {
pub digest: String,
pub wire_digest: String,
pub type_uri: String,
pub requester_did: String,
pub approver_set: String,
pub min_approvals: u32,
pub exclude_requester: bool,
pub challenge: String,
#[serde(default)]
pub correlator: String,
pub approvals: Vec<String>,
#[serde(default)]
pub state_pin: Option<crate::effects::StatePin>,
#[serde(default)]
pub guards: vti_common::guards::Guards,
#[serde(default)]
pub subject_context: Option<String>,
#[serde(default = "default_true")]
pub requester_authorized: bool,
pub created_at: u64,
pub expires_at: u64,
}
fn default_true() -> bool {
true
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TaskConsentGrant {
pub digest: String,
pub requester_did: String,
pub type_uri: String,
pub approvers: Vec<String>,
#[serde(default)]
pub state_pin: Option<crate::effects::StatePin>,
#[serde(default)]
pub guards: vti_common::guards::Guards,
#[serde(default)]
pub delegated_contexts: Vec<String>,
pub granted_at: u64,
pub expires_at: u64,
}
fn pending_key(digest: &str) -> Vec<u8> {
[PENDING_PREFIX, digest.as_bytes()].concat()
}
fn wire_index_key(wire_digest: &str) -> Vec<u8> {
[WIRE_INDEX_PREFIX, wire_digest.as_bytes()].concat()
}
fn grant_key(requester_did: &str, digest: &str) -> Vec<u8> {
[
GRANT_PREFIX,
digest.as_bytes(),
b":",
requester_did.as_bytes(),
]
.concat()
}
fn decode<T: for<'de> Deserialize<'de>>(bytes: &[u8]) -> Result<T, AppError> {
serde_json::from_slice(bytes)
.map_err(|e| AppError::Internal(format!("task-consent decode: {e}")))
}
pub async fn store_pending(ks: &KeyspaceHandle, p: &PendingTaskConsent) -> Result<(), AppError> {
let key = String::from_utf8(pending_key(&p.digest))
.map_err(|e| AppError::Internal(format!("pending key not utf-8: {e}")))?;
ks.insert(key, p).await?;
ks.insert_raw(
String::from_utf8(wire_index_key(&p.wire_digest))
.map_err(|e| AppError::Internal(format!("wire index key not utf-8: {e}")))?,
p.digest.as_bytes().to_vec(),
)
.await
}
pub async fn pending_by_wire_digest(
ks: &KeyspaceHandle,
wire_digest: &str,
now: u64,
) -> Result<Option<PendingTaskConsent>, AppError> {
let Some(bytes) = ks.get_raw(wire_index_key(wire_digest)).await? else {
return Ok(None);
};
let internal = String::from_utf8(bytes)
.map_err(|e| AppError::Internal(format!("wire index value not utf-8: {e}")))?;
get_pending(ks, &internal, now).await
}
pub async fn get_pending(
ks: &KeyspaceHandle,
digest: &str,
now: u64,
) -> Result<Option<PendingTaskConsent>, AppError> {
let key = pending_key(digest);
let Some(b) = ks.get_raw(key.clone()).await? else {
return Ok(None);
};
let p: PendingTaskConsent = decode(&b)?;
if p.expires_at <= now {
ks.remove(key).await?;
ks.remove(wire_index_key(&p.wire_digest)).await?;
return Ok(None);
}
Ok(Some(p))
}
pub async fn delete_pending(ks: &KeyspaceHandle, p: &PendingTaskConsent) -> Result<(), AppError> {
ks.remove(pending_key(&p.digest)).await?;
ks.remove(wire_index_key(&p.wire_digest)).await
}
pub async fn add_approval(
ks: &KeyspaceHandle,
digest: &str,
approver_did: &str,
now: u64,
) -> Result<Option<PendingTaskConsent>, AppError> {
let Some(mut p) = get_pending(ks, digest, now).await? else {
return Ok(None);
};
if !p.approvals.iter().any(|a| a == approver_did) {
p.approvals.push(approver_did.to_string());
store_pending(ks, &p).await?;
}
Ok(Some(p))
}
pub async fn sweep_expired(ks: &KeyspaceHandle, now: u64) -> Result<usize, AppError> {
let mut pruned = 0usize;
for (key, value) in ks.prefix_iter_raw("pending:").await? {
match serde_json::from_slice::<PendingTaskConsent>(&value) {
Ok(p) if p.expires_at <= now => {
ks.remove(key).await?;
ks.remove(wire_index_key(&p.wire_digest)).await?;
pruned += 1;
}
Ok(_) => {}
Err(e) => tracing::debug!(error = %e, "task-consent sweeper: unreadable pending row"),
}
}
for (key, value) in ks.prefix_iter_raw("grant:").await? {
match serde_json::from_slice::<TaskConsentGrant>(&value) {
Ok(g) if g.expires_at <= now => {
ks.remove(key).await?;
pruned += 1;
}
Ok(_) => {}
Err(e) => tracing::debug!(error = %e, "task-consent sweeper: unreadable grant row"),
}
}
Ok(pruned)
}
pub async fn store_grant(ks: &KeyspaceHandle, g: &TaskConsentGrant) -> Result<(), AppError> {
let key = String::from_utf8(grant_key(&g.requester_did, &g.digest))
.map_err(|e| AppError::Internal(format!("grant key not utf-8: {e}")))?;
ks.insert(key, g).await
}
pub async fn consume_grant(
ks: &KeyspaceHandle,
requester_did: &str,
type_uri: &str,
digest: &str,
now: u64,
) -> Result<Option<TaskConsentGrant>, AppError> {
let key = grant_key(requester_did, digest);
let Some(bytes) = ks.get_raw(key.clone()).await? else {
return Ok(None);
};
let grant: TaskConsentGrant = decode(&bytes)?;
ks.remove(key).await?;
if grant.expires_at <= now {
return Ok(None);
}
if grant.type_uri != type_uri {
return Err(AppError::Internal(format!(
"task-consent grant type mismatch: granted for '{}', presented for '{type_uri}'",
grant.type_uri
)));
}
Ok(Some(grant))
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use vta_config::StoreConfig;
use vti_common::store::Store;
async fn temp_ks() -> (KeyspaceHandle, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
let store = Store::open(&StoreConfig {
data_dir: dir.path().to_path_buf(),
})
.unwrap();
(store.keyspace(vta_keyspaces::TASK_CONSENT).unwrap(), dir)
}
const T_UPDATE: &str = "https://trusttasks.org/spec/webvh/dids/update/1.0";
const T_ROTATE: &str = "https://trusttasks.org/spec/webvh/dids/rotate-keys/1.0";
#[test]
fn digest_is_deterministic_and_key_order_independent() {
let a = payload_digest(T_UPDATE, &json!({ "b": 2, "a": 1 })).unwrap();
let b = payload_digest(T_UPDATE, &json!({ "a": 1, "b": 2 })).unwrap();
assert_eq!(a, b, "JCS canonicalization must ignore key order");
assert_ne!(
a,
payload_digest(T_UPDATE, &json!({ "a": 1, "b": 3 })).unwrap()
);
assert!(a.starts_with('z'), "must be multibase base58btc, got {a}");
let (base, bytes) = multibase::decode(&a).expect("valid multibase");
assert_eq!(base, multibase::Base::Base58Btc);
assert_eq!(
&bytes[..2],
&[0x12, 0x20],
"must carry the sha2-256 multihash prefix in-band"
);
assert_eq!(bytes.len(), 34, "2-byte multihash prefix + 32-byte digest");
}
#[test]
fn digest_binds_the_type_uri() {
let payload = json!({ "did": "did:webvh:example.com:abc", "contextId": "default" });
assert_ne!(
payload_digest(T_UPDATE, &payload).unwrap(),
payload_digest(T_ROTATE, &payload).unwrap(),
"same payload under a different task URI must not collide"
);
}
#[test]
fn wire_digest_is_salted_and_distinct_from_the_internal_one() {
let p = json!({ "did": "did:webvh:example.com:acme" });
let internal = payload_digest(T_UPDATE, &p).unwrap();
let a = wire_digest(T_UPDATE, &p, "challenge-aaaa").unwrap();
let b = wire_digest(T_UPDATE, &p, "challenge-bbbb").unwrap();
assert_ne!(
internal, a,
"the wire digest must not equal the storage key"
);
assert_ne!(
a, b,
"a different challenge must yield a different wire digest"
);
assert_eq!(a, wire_digest(T_UPDATE, &p, "challenge-aaaa").unwrap());
}
#[test]
fn wire_digest_binds_the_type_uri() {
let p = json!({ "did": "did:webvh:example.com:acme" });
assert_ne!(
wire_digest(T_UPDATE, &p, "c").unwrap(),
wire_digest(T_ROTATE, &p, "c").unwrap()
);
}
#[tokio::test]
async fn a_decision_resolves_its_pending_by_wire_digest() {
let (ks, _d) = temp_ks().await;
let p = pending("deadbeef", 1);
store_pending(&ks, &p).await.unwrap();
let found = pending_by_wire_digest(&ks, &p.wire_digest, 200)
.await
.unwrap()
.expect("resolved via the wire index");
assert_eq!(found.digest, "deadbeef");
assert!(
pending_by_wire_digest(&ks, "not-a-digest", 200)
.await
.unwrap()
.is_none()
);
delete_pending(&ks, &p).await.unwrap();
assert!(
pending_by_wire_digest(&ks, &p.wire_digest, 200)
.await
.unwrap()
.is_none()
);
}
#[test]
fn digest_uri_payload_boundary_is_unambiguous() {
assert_ne!(
payload_digest("ab", &json!("c")).unwrap(),
payload_digest("a", &json!("bc")).unwrap(),
);
}
fn pending(digest: &str, min: u32) -> PendingTaskConsent {
PendingTaskConsent {
digest: digest.into(),
wire_digest: format!("wire-{digest}"),
correlator: "urn:uuid:test-correlator".into(),
state_pin: None,
guards: Default::default(),
type_uri: "https://…/dids/update/1.0".into(),
requester_did: "did:key:zReq".into(),
approver_set: "operators".into(),
min_approvals: min,
exclude_requester: true,
challenge: "nonce123".into(),
approvals: vec![],
subject_context: None,
requester_authorized: true,
created_at: 100,
expires_at: 1000,
}
}
#[tokio::test]
async fn approvals_accumulate_idempotently() {
let (ks, _d) = temp_ks().await;
store_pending(&ks, &pending("deadbeef", 2)).await.unwrap();
let p = add_approval(&ks, "deadbeef", "did:key:zA", 200)
.await
.unwrap()
.unwrap();
assert_eq!(p.approvals.len(), 1);
let p = add_approval(&ks, "deadbeef", "did:key:zA", 200)
.await
.unwrap()
.unwrap();
assert_eq!(p.approvals.len(), 1);
let p = add_approval(&ks, "deadbeef", "did:key:zB", 200)
.await
.unwrap()
.unwrap();
assert_eq!(p.approvals.len(), 2);
assert!(p.approvals.len() as u32 >= p.min_approvals);
assert!(
add_approval(&ks, "nope", "did:key:zA", 200)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn expired_pending_reads_as_absent_and_cannot_be_approved() {
let (ks, _d) = temp_ks().await;
store_pending(&ks, &pending("deadbeef", 1)).await.unwrap();
assert!(get_pending(&ks, "deadbeef", 999).await.unwrap().is_some());
assert!(get_pending(&ks, "deadbeef", 1000).await.unwrap().is_none());
assert!(
add_approval(&ks, "deadbeef", "did:key:zA", 1001)
.await
.unwrap()
.is_none(),
"an expired pending must not accept approvals"
);
assert!(get_pending(&ks, "deadbeef", 500).await.unwrap().is_none());
}
#[tokio::test]
async fn sweeper_prunes_lapsed_pendings_and_grants() {
let (ks, _d) = temp_ks().await;
store_pending(&ks, &pending("d1", 1)).await.unwrap(); store_grant(&ks, &grant("d2", 500)).await.unwrap();
assert_eq!(
sweep_expired(&ks, 400).await.unwrap(),
0,
"nothing lapsed yet"
);
assert_eq!(sweep_expired(&ks, 2000).await.unwrap(), 2, "both lapsed");
assert_eq!(sweep_expired(&ks, 2000).await.unwrap(), 0, "idempotent");
}
fn grant(digest: &str, expires_at: u64) -> TaskConsentGrant {
TaskConsentGrant {
digest: digest.into(),
state_pin: None,
guards: Default::default(),
requester_did: "did:key:zReq".into(),
type_uri: T_UPDATE.into(),
approvers: vec!["did:key:zA".into()],
delegated_contexts: vec![],
granted_at: 100,
expires_at,
}
}
#[tokio::test]
async fn grant_is_single_use_and_expiry_checked() {
let (ks, _d) = temp_ks().await;
let g = grant("d1", 500);
store_grant(&ks, &g).await.unwrap();
assert!(
consume_grant(&ks, "did:key:zReq", T_UPDATE, "d1", 200)
.await
.unwrap()
.is_some()
);
assert!(
consume_grant(&ks, "did:key:zReq", T_UPDATE, "d1", 200)
.await
.unwrap()
.is_none()
);
store_grant(&ks, &g).await.unwrap();
assert!(
consume_grant(&ks, "did:key:zReq", T_UPDATE, "d1", 999)
.await
.unwrap()
.is_none()
);
assert!(
consume_grant(&ks, "did:key:zReq", T_UPDATE, "d1", 200)
.await
.unwrap()
.is_none()
);
store_grant(&ks, &g).await.unwrap();
assert!(
consume_grant(&ks, "did:key:zOther", T_UPDATE, "d1", 200)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn grant_refuses_a_different_task_uri() {
let (ks, _d) = temp_ks().await;
store_grant(&ks, &grant("d1", 500)).await.unwrap();
let err = consume_grant(&ks, "did:key:zReq", T_ROTATE, "d1", 200)
.await
.unwrap_err();
assert!(
err.to_string().contains("type mismatch"),
"expected a type-mismatch refusal, got: {err}"
);
}
}