#![cfg(any(test, feature = "test-support"))]
pub mod response_conformance;
use std::collections::BTreeMap;
use std::path::PathBuf;
use std::sync::Arc;
use chrono::Duration;
use ed25519_dalek::SigningKey;
use serde_json::Value;
use tokio::sync::RwLock;
use vti_common::slip10::{DerivationPath, ExtendedSigningKey};
use affinidi_did_resolver_cache_sdk::{DIDCacheClient, config::DIDCacheConfigBuilder};
use crate::acl::Role;
use crate::auth::AuthClaims;
use crate::config::{AppConfig, StoreConfig};
use crate::didcomm_bridge::DIDCommBridge;
use crate::keys::seed_store::PlaintextSeedStore;
use crate::keys::{KeyType as SdkKeyType, save_key_record};
use crate::operations::provision_integration::ProvisionIntegrationDeps;
use crate::store::{KeyspaceHandle, Store};
use vta_sdk::did_key::ed25519_multibase_pubkey;
use vta_sdk::provision_integration::{
BootstrapAsk, BootstrapRequest, DidTemplateRef, TemplateBootstrapAsk, VerifiedBootstrapRequest,
};
pub struct TestStore {
_dir: tempfile::TempDir,
_store: Store,
pub contexts_ks: KeyspaceHandle,
pub did_templates_ks: KeyspaceHandle,
pub keys_ks: KeyspaceHandle,
pub acl_ks: KeyspaceHandle,
pub audit_ks: KeyspaceHandle,
pub audit: vta_audit::SharedAuditSink,
pub imported_ks: KeyspaceHandle,
pub internal_ks: KeyspaceHandle,
pub webvh_ks: KeyspaceHandle,
pub sealed_nonces_ks: KeyspaceHandle,
pub drains_ks: KeyspaceHandle,
pub snapshot_ks: KeyspaceHandle,
pub service_state_ks: KeyspaceHandle,
pub data_dir: PathBuf,
}
pub async fn open_test_store() -> TestStore {
let dir = tempfile::tempdir().expect("temp dir");
let data_dir = dir.path().to_path_buf();
let store = Store::open(&StoreConfig {
data_dir: data_dir.clone(),
})
.expect("open store");
TestStore {
contexts_ks: store
.keyspace(crate::keyspaces::CONTEXTS)
.expect("contexts ks"),
did_templates_ks: store
.keyspace(crate::keyspaces::DID_TEMPLATES)
.expect("did_templates ks"),
keys_ks: store.keyspace(crate::keyspaces::KEYS).expect("keys ks"),
acl_ks: store.keyspace(crate::keyspaces::ACL).expect("acl ks"),
audit_ks: store.keyspace(crate::keyspaces::AUDIT).expect("audit ks"),
audit: vta_audit::shared_keyspace_sink(
store.keyspace(crate::keyspaces::AUDIT).expect("audit ks"),
),
imported_ks: store
.keyspace(crate::keyspaces::IMPORTED_SECRETS)
.expect("imported ks"),
internal_ks: store
.keyspace(crate::keyspaces::INTERNAL_KEYS)
.expect("internal keys ks"),
webvh_ks: store.keyspace(crate::keyspaces::WEBVH).expect("webvh ks"),
sealed_nonces_ks: store
.keyspace(crate::keyspaces::SEALED_NONCES)
.expect("nonces ks"),
drains_ks: store.keyspace(crate::keyspaces::DRAINS).expect("drains ks"),
snapshot_ks: store
.keyspace(crate::operations::protocol::snapshot::KEYSPACE_NAME)
.expect("snapshot ks"),
service_state_ks: store
.keyspace(crate::keyspaces::SERVICE_STATE)
.expect("service_state ks"),
_dir: dir,
_store: store,
data_dir,
}
}
pub fn test_app_config(data_dir: PathBuf) -> AppConfig {
AppConfig {
trusted_presentation_verifiers: Vec::new(),
credential_holder_did: None,
vta_did: None,
vta_name: None,
public_url: None,
resolver_url: None,
server: Default::default(),
log: Default::default(),
store: StoreConfig { data_dir },
messaging: None,
mediator_readiness: Default::default(),
services: Default::default(),
auth: Default::default(),
audit: Default::default(),
vault: Default::default(),
app_state: Default::default(),
policy: Default::default(),
secrets: Default::default(),
#[cfg(feature = "tee")]
tee: Default::default(),
hardened: Default::default(),
config_path: PathBuf::new(),
unknown_keys: Vec::new(),
effective_config_digest: None,
effective_config_view: None,
}
}
pub fn test_deps(ts: &TestStore) -> ProvisionIntegrationDeps {
ProvisionIntegrationDeps {
keys_ks: ts.keys_ks.clone(),
acl_ks: ts.acl_ks.clone(),
audit: vta_audit::shared_keyspace_sink(ts.audit_ks.clone()),
contexts_ks: ts.contexts_ks.clone(),
did_templates_ks: ts.did_templates_ks.clone(),
imported_ks: ts.imported_ks.clone(),
webvh_ks: ts.webvh_ks.clone(),
sealed_nonces_ks: ts.sealed_nonces_ks.clone(),
seed_store: Arc::new(PlaintextSeedStore::new(&ts.data_dir)),
config: Arc::new(RwLock::new(test_app_config(ts.data_dir.clone()))),
did_resolver: None,
didcomm_bridge: Arc::new(DIDCommBridge::placeholder()),
webvh_auth_locks: crate::operations::did_webvh::WebvhAuthLocks::new(),
}
}
pub const TEST_ADMIN_SEED: [u8; 32] = [0x7A; 32];
pub const TEST_VTA_SEED: [u8; 32] = [0x5A; 32];
pub fn test_vta_did() -> (String, String) {
let (did, _vm) = did_for_seed(TEST_VTA_SEED[0]);
let vm = format!("{did}#key-0");
(did, vm)
}
pub(crate) fn test_vta_signer() -> VtaOwnSigner {
use ed25519_dalek::SigningKey;
let sk = SigningKey::from_bytes(&[TEST_VTA_SEED[0]; 32]);
let (_did, vm_id) = test_vta_did();
VtaOwnSigner {
vm_id: vm_id.clone(),
secret: secret_from_ed25519(&sk, vm_id),
}
}
pub async fn sign_response_document(
secret: &affinidi_secrets_resolver::secrets::Secret,
doc: &mut serde_json::Value,
) -> bool {
crate::trust_tasks::attach_proof_in_place(secret, doc).await
}
pub(crate) fn secret_from_ed25519(
sk: &ed25519_dalek::SigningKey,
id: String,
) -> affinidi_tdk::secrets_resolver::secrets::Secret {
let mut prefixed = vec![0x80u8, 0x26];
prefixed.extend_from_slice(sk.to_bytes().as_slice());
let mb = multibase::encode(multibase::Base::Base58Btc, prefixed);
let mut secret = affinidi_tdk::secrets_resolver::secrets::Secret::from_multibase(&mb, None)
.expect("construct an ed25519 Secret");
secret.id = id;
secret
}
pub fn test_admin_did() -> (String, String) {
did_for_seed(TEST_ADMIN_SEED[0])
}
pub fn did_for_seed(seed: u8) -> (String, String) {
use ed25519_dalek::SigningKey;
let sk = SigningKey::from_bytes(&[seed; 32]);
let mut mc = vec![0xed, 0x01];
mc.extend_from_slice(sk.verifying_key().as_bytes());
let mb = multibase::encode(multibase::Base::Base58Btc, mc);
let did = format!("did:key:{mb}");
let vm = format!("{did}#{mb}");
(did, vm)
}
pub fn sign_as_test_admin(doc: &mut trust_tasks_rs::TrustTask<serde_json::Value>) {
sign_as(TEST_ADMIN_SEED[0], doc)
}
pub fn sign_as(seed: u8, doc: &mut trust_tasks_rs::TrustTask<serde_json::Value>) {
use affinidi_data_integrity::DataIntegrityProof;
use affinidi_data_integrity::crypto_suites::CryptoSuite;
use affinidi_data_integrity::prepare_sign_input;
use ed25519_dalek::{Signer, SigningKey};
let (_did, vm) = did_for_seed(seed);
let sk = SigningKey::from_bytes(&[seed; 32]);
let mut di = DataIntegrityProof::new(
CryptoSuite::EddsaJcs2022,
vm,
"assertionMethod".to_string(),
None,
Some(
chrono::Utc::now()
.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
.to_string(),
),
None,
);
doc.proof = None;
let input = prepare_sign_input(&*doc, &di, CryptoSuite::EddsaJcs2022)
.expect("the test document prepares for signing");
di.proof_value = Some(multibase::encode(
multibase::Base::Base58Btc,
sk.sign(&input).to_bytes(),
));
doc.proof = Some(
serde_json::from_value(serde_json::to_value(&di).expect("proof serialises"))
.expect("proof round-trips into the framework type"),
);
}
pub fn super_admin_claims() -> AuthClaims {
AuthClaims {
did: test_admin_did().0,
role: Role::Admin,
allowed_contexts: Vec::new(),
session_id: "test-session".into(),
access_expires_at: 0,
issued_at: 0,
amr: Vec::new(),
acr: String::new(),
}
}
pub fn admin_claims_for_context(context: &str) -> AuthClaims {
AuthClaims {
allowed_contexts: vec![context.to_string()],
..super_admin_claims()
}
}
pub async fn signed_request(template_name: &str, context_hint: &str) -> VerifiedBootstrapRequest {
signed_request_with_vars(template_name, context_hint, BTreeMap::new()).await
}
pub async fn signed_request_with_vars(
template_name: &str,
context_hint: &str,
vars: BTreeMap<String, Value>,
) -> VerifiedBootstrapRequest {
let seed = [7u8; 32];
let signing = SigningKey::from_bytes(&seed);
let pub_bytes: [u8; 32] = signing.verifying_key().to_bytes();
let client_did = affinidi_crypto::did_key::ed25519_pub_to_did_key(&pub_bytes);
let ask = BootstrapAsk::TemplateBootstrap(TemplateBootstrapAsk {
context_hint: Some(context_hint.into()),
template: DidTemplateRef {
name: template_name.into(),
vars,
},
admin_template: None,
note: None,
});
let req = BootstrapRequest::sign(
&seed,
&client_did,
[0u8; 16],
Duration::minutes(10),
None,
ask,
)
.await
.expect("sign bootstrap request");
req.verify().expect("verify bootstrap request")
}
pub async fn signed_admin_rotation_request(
admin_template_name: &str,
context_hint: &str,
) -> VerifiedBootstrapRequest {
use vta_sdk::provision_integration::AdminRotationAsk;
let seed = [7u8; 32];
let signing = SigningKey::from_bytes(&seed);
let pub_bytes: [u8; 32] = signing.verifying_key().to_bytes();
let client_did = affinidi_crypto::did_key::ed25519_pub_to_did_key(&pub_bytes);
let ask = BootstrapAsk::AdminRotation(AdminRotationAsk {
context_hint: Some(context_hint.into()),
admin_template: DidTemplateRef {
name: admin_template_name.into(),
vars: BTreeMap::new(),
},
note: None,
});
let req = BootstrapRequest::sign(
&seed,
&client_did,
[0u8; 16],
Duration::minutes(10),
None,
ask,
)
.await
.expect("sign bootstrap request");
req.verify().expect("verify bootstrap request")
}
async fn provision_vta_signing_identity(
keys_ks: &KeyspaceHandle,
data_dir: &std::path::Path,
vta_did_override: Option<&str>,
) -> (String, Arc<PlaintextSeedStore>, VtaOwnSigner) {
use crate::keys::seeds::{SeedRecord, save_seed_record, set_active_seed_id};
let raw_seed = [0xC5u8; 64];
let seed_store = PlaintextSeedStore::new(data_dir);
crate::keys::seed_store::SeedStore::set(&seed_store, &raw_seed)
.await
.expect("write test seed to plaintext store");
let now = chrono::Utc::now();
save_seed_record(
keys_ks,
&SeedRecord {
id: 0,
seed_hex: None,
seed_enc: None,
created_at: now,
retired_at: None,
},
)
.await
.expect("save seed record");
set_active_seed_id(keys_ks, 0)
.await
.expect("set active seed id");
let vta_base_path = "m/26'/1'/0'";
let root = ExtendedSigningKey::from_seed(&raw_seed).expect("bip-32 root");
let dp: DerivationPath = vta_base_path.parse().expect("derivation path");
let derived = root.derive(&dp).expect("derive VTA key");
let signing = ed25519_dalek::SigningKey::from_bytes(derived.signing_key.as_bytes());
let pub_bytes = signing.verifying_key().to_bytes();
let multibase = ed25519_multibase_pubkey(&pub_bytes);
let vta_did = match vta_did_override {
Some(did) => did.to_string(),
None => format!("did:key:{multibase}"),
};
let key_id = format!("{vta_did}#key-0");
let own_signer = {
let mut prefixed = vec![0x80u8, 0x26];
prefixed.extend_from_slice(signing.to_bytes().as_slice());
let private_key_multibase = multibase::encode(multibase::Base::Base58Btc, prefixed);
let mut secret = affinidi_tdk::secrets_resolver::secrets::Secret::from_multibase(
&private_key_multibase,
None,
)
.expect("construct the VTA's own signing secret");
secret.id = key_id.clone();
VtaOwnSigner {
vm_id: key_id.clone(),
secret,
}
};
save_key_record(
keys_ks,
&key_id,
vta_base_path,
SdkKeyType::Ed25519,
&multibase,
"VTA signing key",
None,
Some(0),
)
.await
.expect("save VTA key record");
let st_base_path = "m/26'/1'/1'";
let st_dp: DerivationPath = st_base_path.parse().expect("st derivation path");
let st_derived = root.derive(&st_dp).expect("derive VTA sealed-transfer key");
let st_signing = ed25519_dalek::SigningKey::from_bytes(st_derived.signing_key.as_bytes());
let st_pub_bytes = st_signing.verifying_key().to_bytes();
let st_multibase = ed25519_multibase_pubkey(&st_pub_bytes);
save_key_record(
keys_ks,
&format!("{vta_did}#sealed-transfer-0"),
st_base_path,
SdkKeyType::Ed25519,
&st_multibase,
"VTA sealed-transfer producer-assertion key",
None,
Some(0),
)
.await
.expect("save VTA sealed-transfer key record");
(
vta_did,
Arc::new(PlaintextSeedStore::new(data_dir)),
own_signer,
)
}
pub async fn bootstrap_test_vta(ts: &TestStore) -> (String, ProvisionIntegrationDeps) {
let (vta_did, _seed_store, _own_signer) =
provision_vta_signing_identity(&ts.keys_ks, &ts.data_dir, None).await;
let mut config = test_app_config(ts.data_dir.clone());
config.vta_did = Some(vta_did.clone());
config.public_url = Some("https://vta.test".into());
let resolver = DIDCacheClient::new(DIDCacheConfigBuilder::default().build())
.await
.expect("DID resolver");
let deps = ProvisionIntegrationDeps {
keys_ks: ts.keys_ks.clone(),
acl_ks: ts.acl_ks.clone(),
audit: vta_audit::shared_keyspace_sink(ts.audit_ks.clone()),
contexts_ks: ts.contexts_ks.clone(),
did_templates_ks: ts.did_templates_ks.clone(),
imported_ks: ts.imported_ks.clone(),
webvh_ks: ts.webvh_ks.clone(),
sealed_nonces_ks: ts.sealed_nonces_ks.clone(),
seed_store: Arc::new(PlaintextSeedStore::new(&ts.data_dir)),
config: Arc::new(RwLock::new(config)),
did_resolver: Some(resolver),
didcomm_bridge: Arc::new(DIDCommBridge::placeholder()),
webvh_auth_locks: crate::operations::did_webvh::WebvhAuthLocks::new(),
};
(vta_did, deps)
}
pub async fn build_signing_test_app_state() -> (crate::server::AppState, tempfile::TempDir) {
build_signing_test_app_state_with_sink(None).await
}
pub async fn build_signing_test_app_state_with_sink(
audit_sink: Option<vta_audit::SharedAuditSink>,
) -> (crate::server::AppState, tempfile::TempDir) {
use crate::server::{AppStateParts, build_app_state};
use tokio::sync::watch;
init_jwt_provider();
let dir = tempfile::tempdir().expect("temp dir");
let store = Store::open(&StoreConfig {
data_dir: dir.path().to_path_buf(),
})
.expect("open store");
let keys_ks = store.keyspace(crate::keyspaces::KEYS).expect("keys ks");
let (vta_did, seed_store, _own_signer) =
provision_vta_signing_identity(&keys_ks, dir.path(), None).await;
let seed_store: Arc<dyn crate::keys::seed_store::SeedStore> = seed_store;
let mut config = test_app_config(dir.path().to_path_buf());
config.vta_did = Some(vta_did);
config.public_url = Some("https://vta.test".into());
let (restart_tx, _rx) = watch::channel(false);
let state = build_app_state(
config,
&store,
seed_store,
None,
None,
restart_tx,
AppStateParts {
audit_sink,
..Default::default()
},
)
.await
.expect("build signing app state");
(state, dir)
}
pub const PROVISIONABLE_CONTEXT: &str = "provisionable-ctx";
pub async fn bootstrap_provisionable_test_vta(
ts: &TestStore,
) -> (String, ProvisionIntegrationDeps) {
let (vta_did, deps) = bootstrap_test_vta(ts).await;
crate::contexts::create_context(
&ts.contexts_ks,
PROVISIONABLE_CONTEXT,
"Provisionable Context",
)
.await
.expect("create provisionable context");
(vta_did, deps)
}
pub fn provisionable_mediator_vars() -> BTreeMap<String, Value> {
let mut vars = BTreeMap::new();
vars.insert("URL".into(), Value::String("https://mediator.test".into()));
vars.insert(
"WS_URL".into(),
Value::String("wss://mediator.test/ws".into()),
);
vars.insert("ROUTING_KEYS".into(), Value::Array(vec![]));
vars
}
pub fn init_jwt_provider() {
use std::sync::Once;
static INIT: Once = Once::new();
INIT.call_once(|| {
let _ = jsonwebtoken::crypto::aws_lc::DEFAULT_PROVIDER.install_default();
});
}
pub struct TestSeedStore(pub Vec<u8>);
impl crate::keys::seed_store::SeedStore for TestSeedStore {
fn get(
&self,
) -> std::pin::Pin<
Box<
dyn std::future::Future<Output = Result<Option<Vec<u8>>, crate::error::AppError>>
+ Send
+ '_,
>,
> {
let v = self.0.clone();
Box::pin(async move { Ok(Some(v)) })
}
fn set(
&self,
_seed: &[u8],
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<(), crate::error::AppError>> + Send + '_>,
> {
Box::pin(async { Ok(()) })
}
}
pub struct TestAppContext {
pub jwt_keys: Arc<vti_common::auth::jwt::JwtKeys>,
pub sessions_ks: KeyspaceHandle,
pub acl_ks: KeyspaceHandle,
pub keys_ks: KeyspaceHandle,
pub vault_ks: KeyspaceHandle,
pub backup_bundles_ks: KeyspaceHandle,
pub backup_blob_dir: std::path::PathBuf,
#[cfg(feature = "webvh")]
pub webvh_ks: KeyspaceHandle,
pub policy_ks: KeyspaceHandle,
pub contexts_ks: KeyspaceHandle,
pub vta_did: String,
pub config: Arc<RwLock<AppConfig>>,
pub outbox_ks: KeyspaceHandle,
pub relationships_ks: KeyspaceHandle,
pub state: crate::server::AppState,
pub _dir: tempfile::TempDir,
}
impl TestAppContext {
pub async fn mint_signing_identity(
&self,
seed: u8,
role: &str,
contexts: Vec<String>,
vta_did: &str,
) -> (vta_sdk::client::ClientIdentity, String) {
use ed25519_dalek::SigningKey;
let sk = SigningKey::from_bytes(&[seed; 32]);
let mut mc = vec![0xed, 0x01];
mc.extend_from_slice(sk.verifying_key().as_bytes());
let did = format!(
"did:key:{}",
multibase::encode(multibase::Base::Base58Btc, mc)
);
let token = self.mint_token(&did, role, contexts).await;
let private_key_multibase =
multibase::encode(multibase::Base::Base58Btc, sk.to_bytes().as_slice());
(
vta_sdk::client::ClientIdentity {
client_did: did,
private_key_multibase,
vta_did: vta_did.to_string(),
verification_method: None,
},
token,
)
}
pub async fn mint_token(&self, did: &str, role: &str, contexts: Vec<String>) -> String {
use vti_common::auth::session::{Session, SessionState, store_session};
let session_id = format!("sess-{}", uuid::Uuid::new_v4());
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let session = Session {
session_id: session_id.clone(),
did: did.to_string(),
challenge: String::new(),
state: SessionState::Authenticated,
created_at: now,
last_seen: now,
refresh_token: None,
refresh_expires_at: None,
tee_attested: false,
amr: Vec::new(),
acr: String::new(),
acr_expires_at: None,
token_id: None,
session_pubkey_b58btc: None,
};
store_session(&self.sessions_ks, &session)
.await
.expect("store session");
let claims = self.jwt_keys.new_claims(
did.to_string(),
session_id,
role.to_string(),
contexts,
900,
false,
);
self.jwt_keys.encode(&claims).expect("encode jwt")
}
}
#[derive(Default)]
pub struct TestAppOptions {
pub provisionable_vta: bool,
pub preseed_did_docs: Vec<(String, serde_json::Value)>,
#[cfg(feature = "webvh")]
pub webvh_servers: Vec<(String, String)>,
pub atm: Option<affinidi_tdk::messaging::ATM>,
pub vta_transport: Option<VtaTransportIdentity>,
}
#[derive(Clone)]
pub struct VtaTransportIdentity {
pub did: String,
pub secrets: Vec<affinidi_tdk::secrets_resolver::secrets::Secret>,
pub mediator_did: String,
}
pub(crate) struct VtaOwnSigner {
pub(crate) vm_id: String,
pub(crate) secret: affinidi_tdk::secrets_resolver::secrets::Secret,
}
#[derive(Default)]
struct TransportState {
secrets_resolver: Option<Arc<affinidi_secrets_resolver::ThreadedSecretsResolver>>,
signing_vm_id: Option<String>,
ka_vm_id: Option<String>,
atm: Option<affinidi_tdk::messaging::ATM>,
}
async fn build_transport_state(
identity: Option<&VtaTransportIdentity>,
did_resolver: Option<&DIDCacheClient>,
) -> TransportState {
use affinidi_secrets_resolver::SecretsResolver as _;
use affinidi_tdk::common::TDKSharedState;
use affinidi_tdk::common::config::TDKConfig;
use affinidi_tdk::messaging::config::ATMConfig;
let Some(identity) = identity else {
return TransportState::default();
};
let (secrets_resolver, _task) =
affinidi_secrets_resolver::ThreadedSecretsResolver::new(None).await;
secrets_resolver.insert_vec(&identity.secrets).await;
let vm_id = |suffix: &str| -> Option<String> {
identity
.secrets
.iter()
.map(|s| s.id.clone())
.find(|id| id.ends_with(suffix))
};
let mut builder = TDKConfig::builder().with_load_environment(false);
if let Some(dr) = did_resolver {
builder = builder.with_did_resolver(dr.clone());
}
let tdk = TDKSharedState::new(builder.build().expect("TDK config"))
.await
.expect("TDK shared state");
for secret in &identity.secrets {
tdk.secrets_resolver().insert(secret.clone()).await;
}
let atm = affinidi_tdk::messaging::ATM::new(
ATMConfig::builder().build().expect("ATM config"),
Arc::new(tdk),
)
.await
.expect("transport ATM");
TransportState {
secrets_resolver: Some(Arc::new(secrets_resolver)),
signing_vm_id: vm_id("#key-1"),
ka_vm_id: vm_id("#key-2"),
atm: Some(atm),
}
}
pub async fn build_offline_atm() -> affinidi_tdk::messaging::ATM {
use affinidi_tdk::common::TDKSharedState;
use affinidi_tdk::common::config::TDKConfig;
use affinidi_tdk::messaging::config::ATMConfig;
let tdk = TDKSharedState::new(TDKConfig::builder().build().expect("TDK config"))
.await
.expect("TDK shared state");
affinidi_tdk::messaging::ATM::new(
ATMConfig::builder().build().expect("ATM config"),
Arc::new(tdk),
)
.await
.expect("offline ATM")
}
pub async fn build_test_app() -> (axum::Router, TestAppContext) {
build_test_app_with(TestAppOptions::default()).await
}
pub async fn build_provisionable_test_app() -> (axum::Router, TestAppContext) {
build_test_app_with(TestAppOptions {
provisionable_vta: true,
..Default::default()
})
.await
}
pub async fn build_test_app_with(opts: TestAppOptions) -> (axum::Router, TestAppContext) {
use base64::Engine;
use base64::engine::general_purpose::URL_SAFE_NO_PAD as BASE64;
use tokio::sync::watch;
init_jwt_provider();
let dir = tempfile::tempdir().expect("temp dir");
let store_config = StoreConfig {
data_dir: dir.path().to_path_buf(),
};
let store = Store::open(&store_config).expect("open store");
let keys_ks = store.keyspace(crate::keyspaces::KEYS).unwrap();
let sessions_ks = store.keyspace(crate::keyspaces::SESSIONS).unwrap();
let acl_ks = store.keyspace(crate::keyspaces::ACL).unwrap();
let contexts_ks = store.keyspace(crate::keyspaces::CONTEXTS).unwrap();
{
use chrono::Utc;
let now = Utc::now();
crate::contexts::store_context(
&contexts_ks,
&crate::contexts::ContextRecord {
id: "ctx1".into(),
name: "ctx1".into(),
did: None,
description: None,
parent: None,
base_path: "m/26'/2'/0'".into(),
index: 0,
created_at: now,
updated_at: now,
context_policy: None,
},
)
.await
.expect("seed ctx1");
}
let audit_ks = store.keyspace(crate::keyspaces::AUDIT).unwrap();
let cache_ks = store.keyspace(crate::keyspaces::CACHE).unwrap();
let vault_ks = store.keyspace(crate::keyspaces::VAULT).unwrap();
let vault_ks_ctx = vault_ks.clone();
let service_state_ks = store.keyspace(crate::keyspaces::SERVICE_STATE).unwrap();
let imported_ks = store.keyspace(crate::keyspaces::IMPORTED_SECRETS).unwrap();
let sealed_nonces_ks = store.keyspace(crate::keyspaces::SEALED_NONCES).unwrap();
let backup_bundles_ks = store.keyspace(crate::keyspaces::BACKUP_BUNDLES).unwrap();
let backup_blob_dir = dir.path().join("backups");
let did_templates_ks = store.keyspace(crate::keyspaces::DID_TEMPLATES).unwrap();
#[cfg(feature = "webvh")]
let webvh_ks = store.keyspace(crate::keyspaces::WEBVH).unwrap();
#[cfg(feature = "webvh")]
for (id, did) in &opts.webvh_servers {
seed_webvh_server(&webvh_ks, id, did).await;
}
#[cfg(feature = "webvh")]
let passkey_vms_ks = store.keyspace(crate::keyspaces::PASSKEY_VMS).unwrap();
#[cfg(feature = "webvh")]
let drains_ks = store.keyspace(crate::keyspaces::DRAINS).unwrap();
#[cfg(feature = "webvh")]
let snapshot_ks = store
.keyspace(crate::operations::protocol::snapshot::KEYSPACE_NAME)
.unwrap();
let jwt_seed = [0x42u8; 32];
let jwt_keys = Arc::new(
vti_common::auth::jwt::JwtKeys::from_ed25519_bytes(&jwt_seed, "VTA").expect("jwt keys"),
);
let (vta_did, seed_store, own_signer): (
String,
Arc<dyn crate::keys::seed_store::SeedStore>,
Option<VtaOwnSigner>,
) = if opts.provisionable_vta {
let (did, ps, signer) = provision_vta_signing_identity(
&keys_ks,
dir.path(),
opts.vta_transport.as_ref().map(|t| t.did.as_str()),
)
.await;
let store: Arc<dyn crate::keys::seed_store::SeedStore> = ps;
(did, store, Some(signer))
} else {
let store: Arc<dyn crate::keys::seed_store::SeedStore> =
Arc::new(TestSeedStore(vec![0xABu8; 32]));
let (did, _vm) = test_vta_did();
(did, store, Some(test_vta_signer()))
};
let mut config: AppConfig = toml::from_str(&format!(
r#"
vta_did = "{vta_did}"
[store]
data_dir = "{}"
[auth]
jwt_signing_key = "{}"
"#,
dir.path().display(),
BASE64.encode(jwt_seed),
))
.expect("parse config");
config.config_path = dir.path().join("config.toml");
let (restart_tx, _rx) = watch::channel(false);
let telemetry: vti_common::telemetry::SharedTelemetrySink =
Arc::new(vti_common::telemetry::RingBufferTelemetry::new());
#[cfg(feature = "webvh")]
let mediator_registry = Arc::new(crate::messaging::registry::MediatorListenerRegistry::new(
Arc::clone(&telemetry),
));
#[cfg(feature = "webvh")]
let drain_sweeper = {
let (tx, _rx) = crate::messaging::drain_sweeper::teardown_channel(8);
Arc::new(crate::messaging::drain_sweeper::DrainSweeper::new(
Arc::clone(&mediator_registry),
drains_ks.clone(),
tx,
))
};
let config = Arc::new(RwLock::new(config));
let did_resolver = {
let mut resolver = DIDCacheClient::new(DIDCacheConfigBuilder::default().build())
.await
.ok();
if let Some(client) = resolver.as_mut() {
for (did, doc_json) in &opts.preseed_did_docs {
let doc = serde_json::from_value(doc_json.clone())
.expect("preseed DID document must deserialize into a resolver Document");
client.add_did_document(did, doc).await;
}
}
resolver
};
let mut transport =
build_transport_state(opts.vta_transport.as_ref(), did_resolver.as_ref()).await;
if transport.signing_vm_id.is_none()
&& let Some(signer) = own_signer
{
use affinidi_secrets_resolver::SecretsResolver as _;
let (resolver, _task) = affinidi_secrets_resolver::ThreadedSecretsResolver::new(None).await;
resolver.insert(signer.secret).await;
transport.secrets_resolver = Some(Arc::new(resolver));
transport.signing_vm_id = Some(signer.vm_id);
}
let policy_ks = store.keyspace(crate::keyspaces::POLICY).unwrap();
let state = crate::server::AppState {
audit_sink: vta_audit::shared_keyspace_sink(
store.keyspace(crate::keyspaces::AUDIT).unwrap(),
),
internal_ks: store.keyspace(crate::keyspaces::INTERNAL_KEYS).unwrap(),
idempotency_ks: store.keyspace(crate::keyspaces::IDEMPOTENCY).unwrap(),
mdoc_trust: std::sync::Arc::new(
vta_vault::mdoc_trust::IacaTrustAnchors::from_pem(&[])
.expect("an empty anchor set always parses"),
),
keys_ks: keys_ks.clone(),
sessions_ks: sessions_ks.clone(),
acl_ks: acl_ks.clone(),
contexts_ks: contexts_ks.clone(),
did_templates_ks,
audit_ks,
imported_ks,
cache_ks,
vault_ks,
consent_ks: store.keyspace(crate::keyspaces::CONSENT).unwrap(),
consent_approvers_ks: store.keyspace(crate::keyspaces::CONSENT_APPROVERS).unwrap(),
issued_credentials_ks: store
.keyspace(crate::keyspaces::ISSUED_CREDENTIALS)
.unwrap(),
memory_ks: store.keyspace(crate::keyspaces::MEMORY).unwrap(),
room_groups_ks: store.keyspace(crate::keyspaces::ROOM_GROUPS).unwrap(),
room_invitations_ks: store.keyspace(crate::keyspaces::ROOM_INVITATIONS).unwrap(),
app_state_ks: store.keyspace(crate::keyspaces::APP_STATE).unwrap(),
persona_ks: store.keyspace(crate::keyspaces::PERSONA).unwrap(),
persona_correlation_key: [0u8; 32],
app_state_locks: crate::operations::app_state::NamespaceLocks::default(),
policy_ks: policy_ks.clone(),
task_consent_ks: store.keyspace(crate::keyspaces::TASK_CONSENT).unwrap(),
service_state_ks,
sealed_nonces_ks,
backup_bundles_ks: backup_bundles_ks.clone(),
backup_blob_dir: backup_blob_dir.clone(),
#[cfg(feature = "webvh")]
webvh_ks: webvh_ks.clone(),
#[cfg(feature = "webvh")]
passkey_vms_ks,
#[cfg(feature = "webvh")]
drains_ks,
#[cfg(feature = "webvh")]
snapshot_ks,
#[cfg(feature = "webvh")]
mediator_registry,
#[cfg(feature = "webvh")]
drain_sweeper,
#[cfg(feature = "webvh")]
webvh_auth_locks: crate::operations::did_webvh::WebvhAuthLocks::new(),
telemetry,
wrapping_cache: crate::keys::wrapping::WrappingKeyCache::new(),
config: config.clone(),
seed_store,
did_resolver,
status_list_resolver: None,
secrets_resolver: transport.secrets_resolver,
signing_vm_id: transport.signing_vm_id,
#[cfg(feature = "didcomm")]
ka_vm_id: transport.ka_vm_id,
#[cfg(feature = "didcomm")]
didcomm_bridge: Arc::new(DIDCommBridge::placeholder()),
#[cfg(feature = "tsp")]
tsp_reach: Arc::new(crate::messaging::tsp_reach::TspReachability::new()),
pending_replies: crate::trust_tasks::pending_replies::PendingReplies::new(),
#[cfg(feature = "tsp")]
tsp_recovery: std::sync::Arc::new(affinidi_messaging_sdk::RecoveryCoordinator::new(
affinidi_messaging_sdk::BackoffPolicy::default(),
)),
jwt_keys: Some(jwt_keys.clone()),
atm: transport.atm.or(opts.atm),
tee: None,
restart_tx,
metrics_handle: None,
};
let state_for_ctx = state.clone();
let router = crate::routes::router_with_cors(
&[],
&["127.0.0.1/32".parse().unwrap()],
crate::routes::QuotaSource::Live(state.config.clone()),
)
.with_state(state.clone())
.merge(crate::routes::health_router().with_state(state))
.layer(axum::middleware::from_fn(
vti_common::rate_limit::insert_default_connect_info_if_missing,
));
let ctx = TestAppContext {
jwt_keys,
sessions_ks,
acl_ks,
keys_ks,
vault_ks: vault_ks_ctx,
backup_bundles_ks,
backup_blob_dir,
#[cfg(feature = "webvh")]
webvh_ks,
policy_ks,
contexts_ks: contexts_ks.clone(),
vta_did,
config,
outbox_ks: store.keyspace(crate::keyspaces::OUTBOX).unwrap(),
relationships_ks: store.keyspace(crate::keyspaces::RELATIONSHIPS).unwrap(),
state: state_for_ctx,
_dir: dir,
};
(router, ctx)
}
#[cfg(feature = "webvh")]
pub async fn seed_webvh_server(webvh_ks: &KeyspaceHandle, id: &str, server_did: &str) {
use chrono::Utc;
let now = Utc::now();
let record = vta_sdk::webvh::WebvhServerRecord {
id: id.to_string(),
did: server_did.to_string(),
label: Some(format!("test server {id}")),
created_at: now,
updated_at: now,
};
crate::webvh_store::store_server(webvh_ks, &record)
.await
.expect("seed webvh server");
}
pub async fn seed_acl_entry(
acl_ks: &KeyspaceHandle,
did: &str,
role: crate::acl::Role,
contexts: Vec<String>,
) {
let entry = crate::acl::AclEntry::new(did, role, "test-support").with_contexts(contexts);
crate::acl::store_acl_entry(acl_ks, &entry)
.await
.expect("seed acl entry");
}
#[cfg(feature = "webvh")]
pub const STUB_WEBVH_DID_URL: &str = "https://webvh-host.test/dids/persona/did.jsonl";
#[cfg(feature = "webvh")]
pub struct StubWebvhHost {
base_url: String,
fail_puts: std::sync::Arc<std::sync::atomic::AtomicUsize>,
deletes: std::sync::Arc<std::sync::atomic::AtomicUsize>,
shutdown: Option<tokio::sync::oneshot::Sender<()>>,
handle: Option<tokio::task::JoinHandle<()>>,
}
#[cfg(feature = "webvh")]
impl StubWebvhHost {
pub async fn start() -> StubWebvhHost {
use axum::routing::post;
use serde_json::json;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
let fail_puts = Arc::new(AtomicUsize::new(0));
let deletes = Arc::new(AtomicUsize::new(0));
fn require_bearer(headers: &axum::http::HeaderMap) -> Result<(), axum::http::StatusCode> {
let ok = headers
.get(axum::http::header::AUTHORIZATION)
.and_then(|v| v.to_str().ok())
.is_some_and(|v| v.starts_with("Bearer ") && v.len() > "Bearer ".len());
if ok {
Ok(())
} else {
Err(axum::http::StatusCode::UNAUTHORIZED)
}
}
async fn tokens() -> axum::Json<serde_json::Value> {
axum::Json(json!({
"session": {
"id": "stub-session",
"subject": "did:webvh:stub:vta",
"issuedAt": chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
"expiresAt": "2026-01-02T00:00:00Z",
},
"tokens": {
"accessToken": "stub-access-token",
"refreshToken": "stub-refresh-token",
"tokenType": "Bearer",
"expiresIn": 9_999_999u64,
"refreshExpiresIn": 9_999_999u64,
}
}))
}
let router = axum::Router::new()
.route(
"/api/auth/challenge",
post(|| async {
axum::Json(json!({
"challenge": "stub-challenge-0000000000000000",
"sessionId": "stub-session",
"expiresAt": "2099-01-01T00:00:00Z"
}))
}),
)
.route("/api/auth/", post(tokens))
.route("/api/auth/refresh", post(tokens))
.route(
"/api/agent-names/update",
post(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(json!({ "record": {} })))
}),
)
.route(
"/api/agent-names/remove",
post(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(json!({ "record": {} })))
}),
)
.route(
"/api/agent-names/check",
post(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(json!({
"name": "coverage-agent",
"domain": "webvh-host.test",
"available": true,
"reserved": false,
})))
}),
)
.route(
"/api/dids",
axum::routing::get(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(json!([{
"mnemonic": "cov-orphan-slot",
"domain": "webvh-host.test",
"disabled": false,
"updatedAt": 1_767_225_600u64,
}])))
}),
)
.route(
"/api/me/domains",
axum::routing::get(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(
json!({
"domains": [{
"name": "webvh-host.test",
"defaultDomain": true,
"status": "active",
"createdAt": 1_767_225_600u64,
}],
"default": "webvh-host.test",
}),
))
}),
)
.route(
"/api/dids",
post(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(
json!({ "didUrl": STUB_WEBVH_DID_URL, "mnemonic": "stub-mnemonic" }),
))
}),
)
.route(
"/api/dids/register",
post(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(
json!({ "didUrl": STUB_WEBVH_DID_URL, "mnemonic": "stub-mnemonic" }),
))
}),
)
.route(
"/api/dids/check",
post(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(json!({ "available": true })))
}),
)
.route(
"/api/dids/{mnemonic}",
axum::routing::get(|headers: axum::http::HeaderMap| async move {
require_bearer(&headers)?;
Ok::<_, axum::http::StatusCode>(axum::Json(json!({
"agentNames": [{
"name": "coverage-agent",
"enabled": true,
"createdAt": 1_767_225_600u64,
}],
})))
})
.put({
let fail_puts = fail_puts.clone();
move |headers: axum::http::HeaderMap| {
let fail_puts = fail_puts.clone();
async move {
require_bearer(&headers)?;
if fail_puts
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |n| {
n.checked_sub(1)
})
.is_ok()
{
return Err(axum::http::StatusCode::INTERNAL_SERVER_ERROR);
}
Ok::<_, axum::http::StatusCode>(axum::http::StatusCode::OK)
}
}
})
.delete({
let deletes = deletes.clone();
move |headers: axum::http::HeaderMap| {
let deletes = deletes.clone();
async move {
require_bearer(&headers)?;
deletes.fetch_add(1, Ordering::SeqCst);
Ok::<_, axum::http::StatusCode>(axum::http::StatusCode::OK)
}
}
}),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind stub webvh host port");
let addr = listener.local_addr().expect("stub host local addr");
let base_url = format!("http://{addr}");
let (tx, rx) = tokio::sync::oneshot::channel::<()>();
let handle = tokio::spawn(async move {
let _ = axum::serve(listener, router)
.with_graceful_shutdown(async move {
let _ = rx.await;
})
.await;
});
StubWebvhHost {
base_url,
fail_puts,
deletes,
shutdown: Some(tx),
handle: Some(handle),
}
}
pub fn base_url(&self) -> &str {
&self.base_url
}
pub fn fail_next_publishes(&self, n: usize) {
self.fail_puts.store(n, std::sync::atomic::Ordering::SeqCst);
}
pub fn deletes(&self) -> usize {
self.deletes.load(std::sync::atomic::Ordering::SeqCst)
}
}
#[cfg(feature = "webvh")]
impl Drop for StubWebvhHost {
fn drop(&mut self) {
if let Some(tx) = self.shutdown.take() {
let _ = tx.send(());
}
if let Some(handle) = self.handle.take() {
handle.abort();
}
}
}
pub struct MockVta {
base_url: String,
pub ctx: TestAppContext,
shutdown: Option<tokio::sync::oneshot::Sender<()>>,
handle: Option<tokio::task::JoinHandle<()>>,
#[cfg(feature = "webvh")]
webvh_host: Option<StubWebvhHost>,
#[cfg(feature = "transport-harness")]
transports: Option<MockVtaTransports>,
}
#[cfg(feature = "transport-harness")]
struct MockVtaTransports {
mediator: affinidi_messaging_test_mediator::TestMediatorHandle,
mediator_did: String,
shutdown: tokio_util::sync::CancellationToken,
loop_handle: Option<tokio::task::JoinHandle<()>>,
atm: Arc<affinidi_tdk::messaging::ATM>,
profile: Arc<affinidi_tdk::messaging::profiles::ATMProfile>,
}
impl MockVta {
pub async fn start() -> MockVta {
Self::serve(build_test_app().await).await
}
pub async fn start_provisionable() -> MockVta {
Self::serve(build_provisionable_test_app().await).await
}
#[cfg(feature = "webvh")]
pub const WEBVH_SERVER_ID: &'static str = "stub-webvh";
#[cfg(feature = "webvh")]
pub async fn start_with_webvh_host() -> MockVta {
use serde_json::json;
let host = StubWebvhHost::start().await;
let server_did = "did:webvh:stubscid0000000000000000:webvh-host.test".to_string();
let server_doc = json!({
"@context": ["https://www.w3.org/ns/did/v1"],
"id": server_did,
"service": [{
"id": format!("{server_did}#webvh"),
"type": "WebVHHosting",
"serviceEndpoint": host.base_url(),
}]
});
let opts = TestAppOptions {
provisionable_vta: true,
preseed_did_docs: vec![(server_did.clone(), server_doc)],
webvh_servers: vec![(Self::WEBVH_SERVER_ID.to_string(), server_did)],
atm: None,
vta_transport: None,
};
let mut mock = Self::serve(build_test_app_with(opts).await).await;
mock.webvh_host = Some(host);
mock
}
#[cfg(feature = "transport-harness")]
pub async fn start_with_transports() -> MockVta {
use affinidi_messaging_test_mediator::TestMediator;
use affinidi_tdk::dids::{
OneOrMany, PeerService, PeerServiceEndpoint, PeerServiceEndpointLong,
};
affinidi_messaging_test_mediator::install_default_crypto_provider();
let mediator = TestMediator::builder()
.local_direct_delivery(true, false)
.spawn()
.await
.expect("spawn test mediator");
let mediator_did = mediator.did().to_string();
let didcomm = crate::operations::did_peer::mediator_did_didcomm_service(
&mediator_did,
vec![],
vec![],
);
let tsp = PeerService {
type_: "TSPTransport".into(),
endpoint: PeerServiceEndpoint::Long(OneOrMany::One(PeerServiceEndpointLong {
uri: mediator.endpoint().to_string(),
accept: vec![],
routing_keys: vec![],
})),
id: None,
};
let (vta_did, secrets) = crate::operations::did_peer::mint_did_peer_with_services(
didcomm.into_iter().chain(std::iter::once(tsp)).collect(),
)
.expect("mint VTA did:peer");
const MAX_DID_BYTES: usize = 1_000;
assert!(
vta_did.len() < MAX_DID_BYTES,
"the VTA's did:peer is {} bytes, over the {MAX_DID_BYTES}-byte stock resolver limit \
({}-byte mediator DID). Nothing with a default resolver — this process, the \
mediator, or a consumer — can resolve it. Shed service metadata, or embed the \
mediator DID in fewer services.",
vta_did.len(),
mediator_did.len(),
);
mediator
.register_local_did(&vta_did)
.await
.expect("register the VTA as a local mediator account");
let (router, ctx) = build_test_app_with(TestAppOptions {
provisionable_vta: true,
vta_transport: Some(VtaTransportIdentity {
did: vta_did.clone(),
secrets: secrets.clone(),
mediator_did: mediator_did.clone(),
}),
..Default::default()
})
.await;
let messaging = crate::messaging::service::build_messaging(
secrets,
&vta_did,
&mediator_did,
ctx.outbox_ks.clone(),
ctx.relationships_ks.clone(),
std::sync::Arc::new(std::sync::atomic::AtomicU64::new(0)),
ctx.state.did_resolver.as_ref(),
None,
)
.await
.expect("build VTA messaging over the test mediator");
let atm = messaging.atm.clone();
let messaging_profile = messaging.profile.clone();
ctx.state.didcomm_bridge.set_messaging(
messaging.service.clone(),
(*messaging.atm).clone(),
messaging.profile.clone(),
vta_did.clone(),
);
let shutdown = tokio_util::sync::CancellationToken::new();
let loop_handle = tokio::spawn({
let (messaging, state, vta_did, shutdown) = (
Arc::new(messaging),
ctx.state.clone(),
vta_did.clone(),
shutdown.clone(),
);
async move {
crate::messaging::service::run_inbound_loop(messaging, state, vta_did, shutdown)
.await;
}
});
let mut mock = Self::serve((router, ctx)).await;
mock.transports = Some(MockVtaTransports {
mediator,
mediator_did,
shutdown,
loop_handle: Some(loop_handle),
atm,
profile: messaging_profile,
});
mock
}
#[cfg(feature = "transport-harness")]
pub async fn forget_tsp_relationship(&self, peer_did: &str) {
let t = self
.transports
.as_ref()
.expect("forget_tsp_relationship() requires start_with_transports()");
t.atm
.tsp()
.reset_relationship(&t.profile, peer_did)
.await
.expect("reset the VTA-side TSP relationship");
}
#[cfg(feature = "transport-harness")]
pub fn mediator_did(&self) -> &str {
&self
.transports
.as_ref()
.expect("mediator_did() requires start_with_transports()")
.mediator_did
}
#[cfg(feature = "transport-harness")]
pub async fn register_mediator_account(&self, did: &str) {
self.transports
.as_ref()
.expect("register_mediator_account() requires start_with_transports()")
.mediator
.register_local_did(did)
.await
.expect("register a local mediator account");
}
async fn serve((router, ctx): (axum::Router, TestAppContext)) -> MockVta {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind ephemeral loopback port");
let addr = listener.local_addr().expect("resolve local addr");
let base_url = format!("http://{addr}");
let (tx, rx) = tokio::sync::oneshot::channel::<()>();
let handle = tokio::spawn(async move {
let _ = axum::serve(
listener,
router.into_make_service_with_connect_info::<std::net::SocketAddr>(),
)
.with_graceful_shutdown(async move {
let _ = rx.await;
})
.await;
});
MockVta {
base_url,
ctx,
shutdown: Some(tx),
handle: Some(handle),
#[cfg(feature = "webvh")]
webvh_host: None,
#[cfg(feature = "transport-harness")]
transports: None,
}
}
pub fn base_url(&self) -> &str {
&self.base_url
}
pub fn vta_did(&self) -> &str {
&self.ctx.vta_did
}
pub async fn signing_client(
&self,
seed: u8,
role: &str,
contexts: Vec<String>,
) -> vta_sdk::client::VtaClient {
let vta_did = self.vta_did().to_string();
let (identity, token) = self
.ctx
.mint_signing_identity(seed, role, contexts, &vta_did)
.await;
let client = vta_sdk::client::VtaClient::new(self.base_url()).with_identity(identity);
client.set_token_async(token).await;
client
}
#[cfg(feature = "webvh")]
pub fn fail_next_publishes(&self, n: usize) {
if let Some(host) = &self.webvh_host {
host.fail_next_publishes(n);
}
}
#[cfg(feature = "webvh")]
pub fn webvh_host_deletes(&self) -> usize {
self.webvh_host.as_ref().map_or(0, StubWebvhHost::deletes)
}
#[cfg(feature = "webvh")]
pub async fn corrupt_supersede_keys(&self, scid: &str, version_id: &str) {
crate::operations::did_webvh::webvh_keys::supersede_keys_for_version(
&self.ctx.keys_ks,
scid,
version_id,
)
.await
.expect("supersede keys for test corruption");
}
#[cfg(feature = "webvh")]
pub async fn seed_webvh_server(&self, id: &str, server_did: &str) {
seed_webvh_server(&self.ctx.webvh_ks, id, server_did).await;
}
pub async fn authorize_did(&self, did: &str, role: crate::acl::Role, contexts: Vec<String>) {
seed_acl_entry(&self.ctx.acl_ks, did, role, contexts).await;
}
pub async fn grant_super_admin(&self, did: &str) {
self.authorize_did(did, crate::acl::Role::Admin, Vec::new())
.await;
}
pub async fn shutdown(mut self) {
#[cfg(feature = "transport-harness")]
if let Some(mut transports) = self.transports.take() {
transports.shutdown.cancel();
if let Some(handle) = transports.loop_handle.take() {
let _ = handle.await;
}
transports.atm.graceful_shutdown().await;
transports.mediator.shutdown();
}
if let Some(tx) = self.shutdown.take() {
let _ = tx.send(());
}
if let Some(handle) = self.handle.take() {
let _ = handle.await;
}
}
}
impl Drop for MockVta {
fn drop(&mut self) {
#[cfg(feature = "transport-harness")]
if let Some(mut transports) = self.transports.take() {
transports.shutdown.cancel();
if let Some(handle) = transports.loop_handle.take() {
handle.abort();
}
transports.mediator.shutdown();
}
if let Some(tx) = self.shutdown.take() {
let _ = tx.send(());
}
if let Some(handle) = self.handle.take() {
handle.abort();
}
}
}
#[cfg(all(test, feature = "transport-harness"))]
mod transport_harness_tests {
use super::*;
use affinidi_did_resolver_cache_sdk::{DIDCacheClient, config::DIDCacheConfigBuilder};
#[tokio::test]
async fn a_two_service_did_peer_encodes_both_transports() {
use affinidi_tdk::dids::{
OneOrMany, PeerService, PeerServiceEndpoint, PeerServiceEndpointLong,
};
let mediator = "did:peer:2.Ez6LSmediator";
let didcomm = crate::operations::did_peer::mediator_did_didcomm_service(
mediator,
vec!["didcomm/v2".to_string()],
vec![],
);
let tsp = PeerService {
type_: "TSPTransport".into(),
endpoint: PeerServiceEndpoint::Long(OneOrMany::One(PeerServiceEndpointLong {
uri: mediator.to_string(),
accept: vec![],
routing_keys: vec![],
})),
id: None,
};
let (did, _secrets) = crate::operations::did_peer::mint_did_peer_with_services(
didcomm.into_iter().chain(std::iter::once(tsp)).collect(),
)
.expect("mint two-service did:peer");
let resolver = DIDCacheClient::new(DIDCacheConfigBuilder::default().build())
.await
.expect("local DID cache");
let resolved = resolver.resolve(&did).await.expect("did:peer resolves");
let types: Vec<String> = resolved
.doc
.service
.iter()
.map(|s| s.type_.clone().into_iter().collect::<Vec<_>>().join(","))
.collect();
let ids: Vec<String> = resolved
.doc
.service
.iter()
.map(|s| format!("{:?}", s.id))
.collect();
println!("service types: {types:?}");
println!("service ids: {ids:?}");
assert!(
types.iter().any(|t| t.contains("DIDCommMessaging")),
"types {types:?}"
);
assert!(
types.iter().any(|t| t.contains("TSPTransport")),
"types {types:?}"
);
}
#[tokio::test]
async fn the_mock_advertises_both_transports_to_a_foreign_resolver() {
let _ = tracing_subscriber::fmt()
.with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
.with_test_writer()
.try_init();
let mock = MockVta::start_with_transports().await;
let resolver = DIDCacheClient::new(DIDCacheConfigBuilder::default().build())
.await
.expect("local DID cache");
let resolved = resolver
.resolve(mock.vta_did())
.await
.expect("a did:peer resolves offline in any resolver");
let types: Vec<String> = resolved
.doc
.service
.iter()
.map(|s| s.type_.clone().into_iter().collect::<Vec<_>>().join(","))
.collect();
assert!(
types.iter().any(|t| t.contains("DIDCommMessaging")),
"expected a DIDCommMessaging service, got {types:?}"
);
assert!(
types.iter().any(|t| t.contains("TSPTransport")),
"expected a TSPTransport service, got {types:?}"
);
mock.shutdown().await;
}
#[tokio::test]
async fn the_vta_can_initiate_a_tsp_send_not_only_answer_one() {
let mock = MockVta::start_with_transports().await;
let transport = mock
.ctx
.state
.tsp_transport()
.expect("a mediator-connected VTA has a TSP transport to send on");
assert_eq!(
transport.mediator_did(),
mock.mediator_did(),
"the transport must route through this session's own mediator"
);
let body = vta_sdk::tsp_binding::wrap_envelope(br#"{"probe":true}"#);
let sent = transport.send_to(mock.vta_did(), &body).await;
assert!(
sent.is_ok(),
"the VTA could not put a TSP frame on the wire: {:?}",
sent.err()
);
mock.shutdown().await;
}
fn orphan_did(seed: u8) -> String {
let sk = ed25519_dalek::SigningKey::from_bytes(&[seed; 32]);
let pk = sk.verifying_key().to_bytes();
format!(
"did:key:{}",
vta_sdk::did_key::ed25519_multibase_pubkey(&pk)
)
}
fn did_key_from_seed(seed: u8) -> (String, String) {
let sk = ed25519_dalek::SigningKey::from_bytes(&[seed; 32]);
let pk = sk.verifying_key().to_bytes();
let did = format!(
"did:key:{}",
vta_sdk::did_key::ed25519_multibase_pubkey(&pk)
);
let mut priv_bytes = vec![0x80, 0x26];
priv_bytes.extend_from_slice(&[seed; 32]);
let priv_mb = multibase::encode(multibase::Base::Base58Btc, &priv_bytes);
(did, priv_mb)
}
struct AnsweringPeer {
session: std::sync::Arc<vta_sdk::session::TspSession>,
loop_handle: tokio::task::JoinHandle<()>,
did: String,
log: std::sync::Arc<std::sync::Mutex<PeerLog>>,
}
#[derive(Default)]
struct PeerLog {
received: usize,
sent: usize,
faults: Vec<String>,
}
impl AnsweringPeer {
async fn spawn(mock: &MockVta, seed: u8) -> Self {
let (did, priv_mb) = did_key_from_seed(seed);
mock.register_mediator_account(&did).await;
mock.grant_super_admin(&did).await;
let session =
vta_sdk::session::TspSession::connect(&did, &priv_mb, mock.mediator_did())
.await
.expect("the answering peer connects its TSP socket");
session
.relate(mock.vta_did())
.await
.expect("the peer invites the VTA so its reply is admitted");
let session = std::sync::Arc::new(session);
let loop_session = session.clone();
let vta_did = mock.vta_did().to_string();
let mediator_did = mock.mediator_did().to_string();
let log = std::sync::Arc::new(std::sync::Mutex::new(PeerLog::default()));
let loop_log = log.clone();
let loop_handle = tokio::spawn(async move {
loop {
let polled = loop_session
.receive_next(1)
.await
.map_err(|e| e.to_string());
let doc = match polled {
Ok(Some(doc)) => doc,
Ok(None) => continue,
Err(e) => {
loop_log
.lock()
.expect("peer log")
.faults
.push(format!("the peer's socket ended: {e}"));
break;
}
};
loop_log.lock().expect("peer log").received += 1;
let Ok(request) =
serde_json::from_str::<trust_tasks_rs::TrustTask<serde_json::Value>>(&doc)
else {
loop_log
.lock()
.expect("peer log")
.faults
.push(format!("a frame that is not a Trust Task: {doc}"));
continue;
};
let reply = request.respond_with(
format!("urn:uuid:d6-reply-{}", request.id),
serde_json::json!({ "answered": true }),
);
if let Ok(bytes) = serde_json::to_vec(&reply) {
let sent = loop_session
.send_document(&vta_did, &mediator_did, &bytes)
.await
.map_err(|e| e.to_string());
let mut log = loop_log.lock().expect("peer log");
match sent {
Ok(()) => log.sent += 1,
Err(e) => log.faults.push(format!("the reply could not be sent: {e}")),
}
}
}
});
Self {
session,
loop_handle,
did,
log,
}
}
fn did(&self) -> &str {
&self.did
}
fn report(&self) -> String {
let log = self.log.lock().expect("peer log");
format!(
"the peer received {} document(s), sent {} reply(ies); faults: {}",
log.received,
log.sent,
if log.faults.is_empty() {
"none".to_string()
} else {
log.faults.join("; ")
}
)
}
async fn stop(self) {
self.loop_handle.abort();
let _ = self.loop_handle.await;
self.session.shutdown().await;
}
}
#[tokio::test]
async fn d6_drives_recovery_on_a_reply_timeout() {
let mock = MockVta::start_with_transports().await;
let orphan = orphan_did(0x9a);
mock.register_mediator_account(&orphan).await;
let tsp = crate::operations::outbound::TspSender::from_app_state(&mock.ctx.state)
.expect("a mediator-connected VTA has a TSP transport")
.with_reply_timeout(std::time::Duration::from_millis(600));
let framed = vta_sdk::tsp_binding::wrap_envelope(br#"{"probe":true}"#);
let out = tsp
.recover_for_test(
&orphan,
"thread-d6-1",
&framed,
vta_sdk::trust_tasks::TASK_ACL_GRANT_0_1,
)
.await;
assert!(
out.is_err(),
"the orphan never answers, so recovery cannot complete: {out:?}"
);
assert_eq!(
tsp.recovery().metrics().attempts,
1,
"the coordinator recorded exactly one recovery attempt"
);
mock.shutdown().await;
}
#[tokio::test]
async fn d6_coalesces_concurrent_recoveries() {
let mock = MockVta::start_with_transports().await;
let orphan = orphan_did(0x9b);
mock.register_mediator_account(&orphan).await;
let mk = || {
crate::operations::outbound::TspSender::from_app_state(&mock.ctx.state)
.expect("a mediator-connected VTA has a TSP transport")
.with_reply_timeout(std::time::Duration::from_millis(800))
};
let a = mk();
let b = mk();
let framed = vta_sdk::tsp_binding::wrap_envelope(br#"{"probe":true}"#);
let (ra, rb) = tokio::join!(
a.recover_for_test(
&orphan,
"t-a",
&framed,
vta_sdk::trust_tasks::TASK_ACL_GRANT_0_1
),
b.recover_for_test(
&orphan,
"t-b",
&framed,
vta_sdk::trust_tasks::TASK_ACL_GRANT_0_1
),
);
assert!(
ra.is_err() && rb.is_err(),
"neither recovery can complete against a silent peer"
);
assert_eq!(
a.recovery().metrics().attempts,
1,
"two concurrent recoveries for one peer coalesce onto a single attempt \
(single-flight); a second concurrent send gets InFlight"
);
mock.shutdown().await;
}
#[tokio::test]
async fn d6_recovers_a_reply_from_an_answering_peer() {
let mock = MockVta::start_with_transports().await;
let peer = AnsweringPeer::spawn(&mock, 0x9c).await;
let tsp = crate::operations::outbound::TspSender::from_app_state(&mock.ctx.state)
.expect("a mediator-connected VTA has a TSP transport")
.with_reply_timeout(std::time::Duration::from_secs(30));
let thread = "urn:uuid:d6-success-req-1";
let request = trust_tasks_rs::TrustTask::new(
thread,
vta_sdk::trust_tasks::TASK_ACL_GRANT_0_1
.parse::<trust_tasks_rs::TypeUri>()
.expect("acl/grant/0.1 is a valid Type URI"),
serde_json::json!({}),
);
let framed = vta_sdk::tsp_binding::wrap_envelope(
&serde_json::to_vec(&request).expect("serialise the request"),
);
let out = tsp
.recover_for_test(
peer.did(),
thread,
&framed,
vta_sdk::trust_tasks::TASK_ACL_GRANT_0_1,
)
.await;
assert!(
out.is_ok(),
"the answering peer replies, so recovery completes with the reply: {out:?} — {}",
peer.report()
);
let reply = out.expect("recovery returned the correlated reply");
assert_eq!(
reply.get("threadId").and_then(|t| t.as_str()),
Some(thread),
"the correlated document is the reply the peer threaded to the request"
);
let metrics = tsp.recovery().metrics();
assert_eq!(metrics.attempts, 1, "exactly one recovery attempt ran");
assert_eq!(
metrics.successes, 1,
"and it settled as a success — the relationship recovered and the reply arrived"
);
assert_eq!(metrics.give_ups, 0, "a recovered peer is not given up on");
peer.stop().await;
mock.shutdown().await;
}
#[tokio::test]
async fn tsp_is_not_offered_before_a_mediator_session_exists() {
let (_router, ctx) = build_test_app().await;
assert!(
ctx.state.tsp_transport().is_none(),
"an unconnected VTA must not claim a TSP transport"
);
assert!(
crate::operations::outbound::TspSender::from_app_state(&ctx.state).is_none(),
"TSP must drop out of transport selection until a session can carry it"
);
}
async fn poll_tsp_inbox(
env: &affinidi_messaging_test_mediator::TestEnvironment,
profile: &std::sync::Arc<affinidi_messaging_sdk::profiles::ATMProfile>,
) -> String {
use affinidi_messaging_sdk::messages::fetch::FetchOptions;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(15);
while std::time::Instant::now() < deadline {
let fetched = env
.atm
.fetch_messages(profile, &FetchOptions::default())
.await
.expect("fetch messages");
if let Some(msg) = fetched.success.first().and_then(|e| e.msg.as_ref()) {
return msg.clone();
}
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
}
panic!("no TSP message received within the deadline");
}
#[tokio::test]
async fn send_metadata_private_reaches_a_cross_mediator_peer() {
use affinidi_messaging_test_mediator::topology::TestTopology;
let topology = TestTopology::builder()
.mediators(2)
.spawn()
.await
.expect("spawn a two-mediator topology");
let mediator_b = topology
.mediator_did(1)
.expect("mediator B DID")
.to_string();
let alice = topology
.add_user(0, "alice")
.await
.expect("alice on mediator A");
let bob = topology
.add_user(1, "bob")
.await
.expect("bob on mediator B");
topology
.relate_directly(&alice, &bob)
.await
.expect("seed the TSP relationship both ways");
let transport = crate::messaging::tsp_transport::TspTransport::new(
topology.node(0).expect("node A").atm.clone(),
alice.profile.clone(),
)
.expect("alice's profile carries a mediator, so a transport can be built");
assert_ne!(
transport.mediator_did(),
mediator_b,
"the peer must be on a different mediator for this to exercise nesting"
);
let body = vta_sdk::tsp_binding::wrap_envelope(br#"{"probe":"cross-mediator"}"#);
transport
.send_metadata_private(&bob.did, Some(&mediator_b), &body)
.await
.expect("the nested routed send is accepted for delivery");
let bob_env = topology.node(1).expect("node B");
let stored = poll_tsp_inbox(bob_env, &bob.profile).await;
let (recovered, sender) = bob_env
.atm
.tsp()
.unpack(&bob.profile, &stored)
.await
.expect("bob unpacks the nested message");
assert_eq!(recovered, body, "bob receives the framed body intact");
assert_eq!(sender, alice.did, "and TSP proves alice as the sender");
topology.shutdown().await.expect("shutdown the topology");
}
#[tokio::test]
async fn send_metadata_private_stays_direct_for_a_same_mediator_peer() {
use affinidi_messaging_test_mediator::topology::TestTopology;
let topology = TestTopology::builder()
.mediators(1)
.spawn()
.await
.expect("spawn a one-mediator topology");
let mediator_a = topology
.mediator_did(0)
.expect("mediator A DID")
.to_string();
let alice = topology
.add_user(0, "alice")
.await
.expect("alice on mediator A");
let carol = topology
.add_user(0, "carol")
.await
.expect("carol on mediator A");
topology
.relate_directly(&alice, &carol)
.await
.expect("seed the TSP relationship both ways");
let transport = crate::messaging::tsp_transport::TspTransport::new(
topology.node(0).expect("node A").atm.clone(),
alice.profile.clone(),
)
.expect("alice's profile carries a mediator");
assert_eq!(
transport.mediator_did(),
mediator_a,
"alice and carol share this mediator"
);
let body = vta_sdk::tsp_binding::wrap_envelope(br#"{"probe":"same-mediator"}"#);
transport
.send_metadata_private(&carol.did, Some(&mediator_a), &body)
.await
.expect("the direct routed send is accepted for delivery");
let carol_env = topology.node(0).expect("node A");
let stored = poll_tsp_inbox(carol_env, &carol.profile).await;
let (recovered, sender) = carol_env
.atm
.tsp()
.unpack(&carol.profile, &stored)
.await
.expect("carol unpacks the direct message");
assert_eq!(recovered, body);
assert_eq!(sender, alice.did);
topology.shutdown().await.expect("shutdown the topology");
}
}
pub struct SoftAuthenticator {
signing_key: p256::ecdsa::SigningKey,
credential_id: Vec<u8>,
}
pub struct SoftRegistration {
pub credential_id: String,
pub public_key_multibase: String,
pub cose_algorithm: i64,
pub attestation_object: String,
pub client_data_json: String,
pub authenticator_data: String,
}
impl SoftAuthenticator {
pub fn new(seed: u8) -> Self {
let signing_key = p256::ecdsa::SigningKey::from_bytes(&[seed; 32].into())
.expect("a fixed 32-byte scalar is a valid P-256 key");
Self {
signing_key,
credential_id: vec![seed; 32],
}
}
pub fn register(&self, rp_id: &str, origin: &str, challenge: &str) -> SoftRegistration {
use base64::Engine;
use base64::engine::general_purpose::URL_SAFE_NO_PAD as B64URL;
use sha2::{Digest, Sha256};
let point = self.signing_key.verifying_key().to_sec1_point(false);
let cose = ciborium::value::Value::Map(vec![
(
ciborium::value::Value::Integer(1.into()),
ciborium::value::Value::Integer(2.into()),
),
(
ciborium::value::Value::Integer(3.into()),
ciborium::value::Value::Integer((-7).into()),
),
(
ciborium::value::Value::Integer((-1).into()),
ciborium::value::Value::Integer(1.into()),
),
(
ciborium::value::Value::Integer((-2).into()),
ciborium::value::Value::Bytes(point.x().expect("uncompressed point").to_vec()),
),
(
ciborium::value::Value::Integer((-3).into()),
ciborium::value::Value::Bytes(point.y().expect("uncompressed point").to_vec()),
),
]);
let mut cose_bytes = Vec::new();
ciborium::ser::into_writer(&cose, &mut cose_bytes).expect("COSE key serialises");
let mut auth_data = Vec::new();
auth_data.extend_from_slice(&Sha256::digest(rp_id.as_bytes()));
auth_data.push(0x45);
auth_data.extend_from_slice(&0u32.to_be_bytes());
auth_data.extend_from_slice(&[0u8; 16]); auth_data.extend_from_slice(&(self.credential_id.len() as u16).to_be_bytes());
auth_data.extend_from_slice(&self.credential_id);
auth_data.extend_from_slice(&cose_bytes);
let attestation = ciborium::value::Value::Map(vec![
(
ciborium::value::Value::Text("fmt".into()),
ciborium::value::Value::Text("none".into()),
),
(
ciborium::value::Value::Text("attStmt".into()),
ciborium::value::Value::Map(vec![]),
),
(
ciborium::value::Value::Text("authData".into()),
ciborium::value::Value::Bytes(auth_data.clone()),
),
]);
let mut attestation_bytes = Vec::new();
ciborium::ser::into_writer(&attestation, &mut attestation_bytes)
.expect("attestation object serialises");
let client_data = serde_json::json!({
"type": "webauthn.create",
"challenge": challenge,
"origin": origin,
"crossOrigin": false,
});
let client_data_json = serde_json::to_vec(&client_data).expect("client data serialises");
let (cose_algorithm, public_key_multibase) =
crate::operations::passkey_vms::multikey::cose_key_to_multikey(&cose_bytes)
.expect("the COSE key converts to a Multikey");
SoftRegistration {
credential_id: B64URL.encode(&self.credential_id),
public_key_multibase,
cose_algorithm,
attestation_object: B64URL.encode(&attestation_bytes),
client_data_json: B64URL.encode(&client_data_json),
authenticator_data: B64URL.encode(&auth_data),
}
}
}
pub async fn seed_held_credential(
vault_ks: &KeyspaceHandle,
vct: &str,
disclosable_claim: &str,
subject_did: &str,
) -> 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, SigningKey};
use vta_vault::mint::{MintRequest, mint_sd_jwt_vc};
struct EddsaSigner {
key: SigningKey,
kid: String,
}
impl JwtSigner for EddsaSigner {
fn algorithm(&self) -> &str {
"EdDSA"
}
fn key_id(&self) -> Option<&str> {
Some(&self.kid)
}
fn sign_jwt(
&self,
header: &serde_json::Value,
payload: &serde_json::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}");
let sig = self.key.sign(input.as_bytes());
Ok(format!(
"{input}.{}",
URL_SAFE_NO_PAD.encode(sig.to_bytes())
))
}
}
let issuer_key = SigningKey::from_bytes(&[0x71; 32]);
let issuer_did =
affinidi_crypto::did_key::ed25519_pub_to_did_key(issuer_key.verifying_key().as_bytes());
let signer = EddsaSigner {
key: issuer_key,
kid: format!("{issuer_did}#key-0"),
};
let compact = mint_sd_jwt_vc(
&MintRequest {
vct,
issuer_did: &issuer_did,
subject_did,
claims: &serde_json::json!({ disclosable_claim: "Alice" }),
disclosable: &[disclosable_claim],
iat: 1_700_000_000,
exp: Some(1_900_000_000),
},
&signer,
)
.expect("mint the SD-JWT-VC");
let body = serde_json::from_value(
serde_json::json!({ "credential_response": { "credential": compact } }),
)
.expect("issue body deserialises");
let cred = crate::operations::credential_exchange::receive_issued_credential(
vault_ks,
&body,
None,
None,
chrono::Utc::now(),
)
.await
.expect("receive the issued credential");
cred.id
}
pub async fn seed_holder_key(
state: &crate::server::AppState,
derivation_path: &str,
context_id: Option<&str>,
) -> String {
use vta_sdk::keys::{KeyOrigin, KeyRecord, KeyStatus, KeyType};
use vti_common::slip10::{DerivationPath, ExtendedSigningKey};
let seed = state
.seed_store
.get()
.await
.expect("read the seed")
.expect("an active seed");
let bip32 = ExtendedSigningKey::from_seed(&seed).expect("seed is a valid BIP-32 root");
let derived = bip32
.derive(&derivation_path.parse::<DerivationPath>().expect("a path"))
.expect("derive");
let did = affinidi_crypto::did_key::ed25519_pub_to_did_key(
derived.signing_key.verifying_key().as_bytes(),
);
let multibase = did.strip_prefix("did:key:").expect("a did:key");
let key_id = format!("{did}#{multibase}");
let record = KeyRecord {
key_id: key_id.clone(),
derivation_path: derivation_path.to_string(),
key_type: KeyType::Ed25519,
status: KeyStatus::Active,
public_key: multibase.to_string(),
label: None,
context_id: context_id.map(str::to_string),
seed_id: None,
exportable: None,
origin: KeyOrigin::Derived,
created_at: chrono::Utc::now(),
updated_at: chrono::Utc::now(),
};
state
.keys_ks
.insert(crate::keys::store_key(&key_id), &record)
.await
.expect("store the holder key record");
did
}