use std::sync::Arc;
use std::time::Duration;
use affinidi_did_resolver_cache_sdk::{DIDCacheClient, config::DIDCacheConfigBuilder};
use affinidi_tdk::common::TDKSharedState;
use affinidi_tdk::common::config::TDKConfig;
use affinidi_tdk::messaging::ATM;
use affinidi_tdk::messaging::config::ATMConfig;
use affinidi_tdk::secrets_resolver::{SecretsResolver, ThreadedSecretsResolver};
use vti_common::slip10::ExtendedSigningKey;
use base64::Engine;
use base64::engine::general_purpose::URL_SAFE_NO_PAD as BASE64;
use crate::auth::AuthState;
use crate::auth::jwt::JwtKeys;
use crate::auth::session::cleanup_expired_sessions;
use crate::config::{AppConfig, AuthConfig};
#[cfg(any(feature = "didcomm", feature = "tsp"))]
use crate::didcomm_bridge::DIDCommBridge;
use crate::error::AppError;
use crate::keys::KeyRecord;
use crate::keys::derivation::Bip32Extension;
use crate::keys::seed_store::SeedStore;
use crate::keys::seeds::load_seed_bytes;
#[cfg(feature = "rest")]
use crate::routes;
use crate::store::{KeyspaceHandle, Store};
use tokio::sync::{RwLock, watch};
#[cfg(feature = "rest")]
use tower_http::trace::{DefaultMakeSpan, DefaultOnRequest, DefaultOnResponse, TraceLayer};
use tracing::Level;
use tracing::{debug, error, info, warn};
#[cfg(any(feature = "didcomm", feature = "tsp"))]
use tokio_util::sync::CancellationToken;
#[cfg(any(feature = "didcomm", feature = "tsp"))]
use vta_sdk::acl_setup;
#[derive(Clone)]
#[cfg(feature = "tee")]
pub struct TeeContext {
pub state: crate::tee::TeeState,
pub mnemonic_guard: Option<Arc<crate::tee::mnemonic_guard::MnemonicExportGuard>>,
}
#[derive(Clone)]
#[cfg(not(feature = "tee"))]
pub struct TeeContext(());
pub fn trigger_restart(restart_tx: &watch::Sender<bool>) {
let tx = restart_tx.clone();
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let _ = tx.send(true);
});
}
#[derive(Clone)]
#[non_exhaustive]
pub struct AppState {
pub keys_ks: KeyspaceHandle,
pub sessions_ks: KeyspaceHandle,
pub acl_ks: KeyspaceHandle,
pub contexts_ks: KeyspaceHandle,
pub did_templates_ks: KeyspaceHandle,
pub audit_ks: KeyspaceHandle,
pub audit_sink: vta_audit::SharedAuditSink,
pub imported_ks: KeyspaceHandle,
pub internal_ks: KeyspaceHandle,
pub cache_ks: KeyspaceHandle,
pub vault_ks: KeyspaceHandle,
pub service_state_ks: KeyspaceHandle,
pub sealed_nonces_ks: KeyspaceHandle,
pub idempotency_ks: KeyspaceHandle,
pub backup_bundles_ks: KeyspaceHandle,
pub backup_blob_dir: std::path::PathBuf,
#[cfg(feature = "webvh")]
pub webvh_ks: KeyspaceHandle,
pub mdoc_trust: Arc<vta_vault::mdoc_trust::IacaTrustAnchors>,
#[cfg(feature = "webvh")]
pub passkey_vms_ks: KeyspaceHandle,
pub consent_ks: KeyspaceHandle,
pub consent_approvers_ks: KeyspaceHandle,
pub issued_credentials_ks: KeyspaceHandle,
pub memory_ks: KeyspaceHandle,
pub room_groups_ks: KeyspaceHandle,
pub room_invitations_ks: KeyspaceHandle,
pub app_state_ks: KeyspaceHandle,
pub persona_ks: KeyspaceHandle,
pub persona_correlation_key: [u8; 32],
pub app_state_locks: crate::operations::app_state::NamespaceLocks,
pub policy_ks: KeyspaceHandle,
pub task_consent_ks: KeyspaceHandle,
#[cfg(feature = "webvh")]
pub drains_ks: KeyspaceHandle,
#[cfg(feature = "webvh")]
pub snapshot_ks: KeyspaceHandle,
#[cfg(all(feature = "webvh", feature = "didcomm"))]
pub mediator_registry: Arc<crate::messaging::registry::MediatorListenerRegistry>,
#[cfg(feature = "webvh")]
pub webvh_auth_locks: crate::operations::did_webvh::WebvhAuthLocks,
#[cfg(all(feature = "webvh", feature = "didcomm"))]
pub drain_sweeper: Arc<crate::messaging::drain_sweeper::DrainSweeper>,
pub telemetry: vti_common::telemetry::SharedTelemetrySink,
pub wrapping_cache: crate::keys::wrapping::WrappingKeyCache,
pub config: Arc<RwLock<AppConfig>>,
pub seed_store: Arc<dyn SeedStore>,
pub did_resolver: Option<DIDCacheClient>,
pub status_list_resolver: Option<Arc<dyn crate::vault::status::StatusListResolver>>,
pub secrets_resolver: Option<Arc<ThreadedSecretsResolver>>,
pub signing_vm_id: Option<String>,
#[cfg(any(feature = "didcomm", feature = "tsp"))]
pub ka_vm_id: Option<String>,
#[cfg(any(feature = "didcomm", feature = "tsp"))]
pub didcomm_bridge: Arc<DIDCommBridge>,
#[cfg(feature = "tsp")]
pub tsp_reach: Arc<crate::messaging::tsp_reach::TspReachability>,
pub jwt_keys: Option<Arc<JwtKeys>>,
pub atm: Option<ATM>,
#[cfg(feature = "tsp")]
pub tsp_profile: Option<std::sync::Arc<affinidi_tdk::messaging::profiles::ATMProfile>>,
pub tee: Option<TeeContext>,
pub restart_tx: watch::Sender<bool>,
#[cfg(feature = "rest")]
pub metrics_handle: Option<crate::metrics::PrometheusHandle>,
}
impl AppState {
pub fn trust_task_vm_resolver(&self) -> vti_common::auth::TrustTaskVmResolver {
vti_common::auth::TrustTaskVmResolver::from_optional(self.did_resolver.clone())
}
}
impl AuthState for AppState {
fn jwt_keys(&self) -> Option<&Arc<JwtKeys>> {
self.jwt_keys.as_ref()
}
fn sessions_ks(&self) -> &KeyspaceHandle {
&self.sessions_ks
}
}
#[derive(Default)]
#[non_exhaustive]
pub struct AppStateParts {
pub telemetry: Option<vti_common::telemetry::SharedTelemetrySink>,
pub audit_sink: Option<vta_audit::SharedAuditSink>,
#[cfg(all(feature = "webvh", feature = "didcomm"))]
pub mediator_registry: Option<Arc<crate::messaging::registry::MediatorListenerRegistry>>,
#[cfg(all(feature = "webvh", feature = "didcomm"))]
pub drain_sweeper: Option<Arc<crate::messaging::drain_sweeper::DrainSweeper>>,
#[cfg(any(feature = "didcomm", feature = "tsp"))]
pub didcomm_bridge: Option<Arc<DIDCommBridge>>,
#[cfg(feature = "rest")]
pub metrics_handle: Option<crate::metrics::PrometheusHandle>,
}
pub async fn build_app_state(
config: AppConfig,
store: &Store,
seed_store: Arc<dyn SeedStore>,
storage_encryption_key: Option<[u8; 32]>,
tee_context: Option<TeeContext>,
restart_tx: watch::Sender<bool>,
parts: AppStateParts,
) -> Result<AppState, AppError> {
let mdoc_trust = Arc::new(vta_vault::mdoc_trust::IacaTrustAnchors::from_pem(
&config.vault.mdoc_iaca_trust_anchors,
)?);
let apply_encryption = |ks: KeyspaceHandle| -> KeyspaceHandle {
if let Some(key) = storage_encryption_key {
ks.with_encryption(key)
} else {
ks
}
};
let keys_ks = apply_encryption(store.keyspace(crate::keyspaces::KEYS)?);
let sessions_ks = apply_encryption(store.keyspace(crate::keyspaces::SESSIONS)?);
let acl_ks = apply_encryption(store.keyspace(crate::keyspaces::ACL)?);
let contexts_ks = apply_encryption(store.keyspace(crate::keyspaces::CONTEXTS)?);
let did_templates_ks = apply_encryption(store.keyspace(crate::keyspaces::DID_TEMPLATES)?);
let audit_ks = apply_encryption(store.keyspace(crate::keyspaces::AUDIT)?);
let imported_ks = apply_encryption(store.keyspace(crate::keyspaces::IMPORTED_SECRETS)?);
let internal_ks = apply_encryption(store.keyspace(crate::keyspaces::INTERNAL_KEYS)?);
let cache_ks = apply_encryption(store.keyspace(crate::keyspaces::CACHE)?);
let vault_ks = apply_encryption(store.keyspace(crate::keyspaces::VAULT)?);
let service_state_ks = apply_encryption(store.keyspace(crate::keyspaces::SERVICE_STATE)?);
let sealed_nonces_ks = apply_encryption(store.keyspace(crate::keyspaces::SEALED_NONCES)?);
let idempotency_ks = apply_encryption(store.keyspace(crate::keyspaces::IDEMPOTENCY)?);
let backup_bundles_ks = apply_encryption(store.keyspace(crate::keyspaces::BACKUP_BUNDLES)?);
let backup_blob_dir = config.store.data_dir.join("backups");
#[cfg(feature = "webvh")]
let webvh_ks = apply_encryption(store.keyspace(crate::keyspaces::WEBVH)?);
#[cfg(feature = "webvh")]
let passkey_vms_ks = apply_encryption(store.keyspace(crate::keyspaces::PASSKEY_VMS)?);
let consent_ks = apply_encryption(store.keyspace(crate::keyspaces::CONSENT)?);
let consent_approvers_ks =
apply_encryption(store.keyspace(crate::keyspaces::CONSENT_APPROVERS)?);
let issued_credentials_ks =
apply_encryption(store.keyspace(crate::keyspaces::ISSUED_CREDENTIALS)?);
let memory_ks = apply_encryption(store.keyspace(crate::keyspaces::MEMORY)?);
let room_groups_ks = apply_encryption(store.keyspace(crate::keyspaces::ROOM_GROUPS)?);
let room_invitations_ks = apply_encryption(store.keyspace(crate::keyspaces::ROOM_INVITATIONS)?);
let app_state_ks = apply_encryption(store.keyspace(crate::keyspaces::APP_STATE)?);
let persona_ks = apply_encryption(store.keyspace(crate::keyspaces::PERSONA)?);
let persona_correlation_key: [u8; 32] = {
use sha2::{Digest, Sha256};
let mut h = Sha256::new();
h.update(b"vta-persona/correlation-index/v1");
h.update(storage_encryption_key.unwrap_or([0u8; 32]));
h.finalize().into()
};
let policy_ks = apply_encryption(store.keyspace(crate::keyspaces::POLICY)?);
let task_consent_ks = apply_encryption(store.keyspace(crate::keyspaces::TASK_CONSENT)?);
#[cfg(feature = "webvh")]
let drains_ks = apply_encryption(store.keyspace(crate::keyspaces::DRAINS)?);
#[cfg(feature = "webvh")]
let snapshot_ks =
apply_encryption(store.keyspace(crate::operations::protocol::snapshot::KEYSPACE_NAME)?);
let auth = init_auth(
&config,
&*seed_store,
&keys_ks,
#[cfg(feature = "webvh")]
Some(&webvh_ks),
#[cfg(not(feature = "webvh"))]
None,
)
.await;
let telemetry: vti_common::telemetry::SharedTelemetrySink = parts
.telemetry
.unwrap_or_else(|| Arc::new(vti_common::telemetry::RingBufferTelemetry::new()));
#[cfg(all(feature = "webvh", feature = "didcomm"))]
let mediator_registry = parts.mediator_registry.unwrap_or_else(|| {
Arc::new(crate::messaging::registry::MediatorListenerRegistry::new(
Arc::clone(&telemetry),
))
});
#[cfg(all(feature = "webvh", feature = "didcomm"))]
let drain_sweeper = parts.drain_sweeper.unwrap_or_else(|| {
let (tx, _rx) = crate::messaging::drain_sweeper::teardown_channel(
crate::messaging::drain_sweeper::DEFAULT_TEARDOWN_CHANNEL_CAPACITY,
);
Arc::new(crate::messaging::drain_sweeper::DrainSweeper::new(
Arc::clone(&mediator_registry),
drains_ks.clone(),
tx,
))
});
let audit_sink: vta_audit::SharedAuditSink = parts
.audit_sink
.unwrap_or_else(|| Arc::new(vta_audit::KeyspaceAuditSink::new(audit_ks.clone())));
Ok(AppState {
keys_ks,
sessions_ks,
acl_ks,
contexts_ks,
did_templates_ks,
audit_ks,
audit_sink,
imported_ks,
internal_ks,
cache_ks,
vault_ks,
service_state_ks,
sealed_nonces_ks,
idempotency_ks,
backup_bundles_ks,
backup_blob_dir,
#[cfg(feature = "webvh")]
webvh_ks,
#[cfg(feature = "webvh")]
passkey_vms_ks,
consent_ks,
consent_approvers_ks,
issued_credentials_ks,
memory_ks,
room_groups_ks,
room_invitations_ks,
app_state_ks,
persona_ks,
persona_correlation_key,
app_state_locks: crate::operations::app_state::NamespaceLocks::default(),
policy_ks,
task_consent_ks,
#[cfg(feature = "webvh")]
drains_ks,
#[cfg(feature = "webvh")]
snapshot_ks,
#[cfg(all(feature = "webvh", feature = "didcomm"))]
mediator_registry,
#[cfg(all(feature = "webvh", feature = "didcomm"))]
drain_sweeper,
#[cfg(feature = "webvh")]
webvh_auth_locks: crate::operations::did_webvh::WebvhAuthLocks::new(),
telemetry,
wrapping_cache: crate::keys::wrapping::WrappingKeyCache::new(),
mdoc_trust,
config: Arc::new(RwLock::new(config)),
seed_store,
did_resolver: auth.did_resolver.clone(),
status_list_resolver: crate::vault::status::default_status_resolver(auth.did_resolver),
secrets_resolver: auth.secrets_resolver,
signing_vm_id: auth.signing_vm_id,
#[cfg(any(feature = "didcomm", feature = "tsp"))]
ka_vm_id: auth.ka_vm_id,
#[cfg(any(feature = "didcomm", feature = "tsp"))]
didcomm_bridge: parts
.didcomm_bridge
.unwrap_or_else(|| Arc::new(DIDCommBridge::placeholder())),
#[cfg(feature = "tsp")]
tsp_reach: Arc::new(crate::messaging::tsp_reach::TspReachability::new()),
jwt_keys: auth.jwt_keys,
atm: auth.atm,
#[cfg(feature = "tsp")]
tsp_profile: auth.tsp_profile,
tee: tee_context,
restart_tx,
#[cfg(feature = "rest")]
metrics_handle: parts.metrics_handle,
})
}
#[cfg_attr(not(feature = "webvh"), allow(unused_mut))]
pub async fn run(
mut config: AppConfig,
store: Store,
seed_store: Arc<dyn SeedStore>,
storage_encryption_key: Option<[u8; 32]>,
tee_context: Option<TeeContext>,
allow_degraded: bool,
#[cfg_attr(not(feature = "didcomm"), allow(unused_variables))] flush_queues: bool,
) -> Result<(), AppError> {
config.validate()?;
{
let keys_ks_boot = {
let ks = store.keyspace(crate::keyspaces::KEYS)?;
match storage_encryption_key {
Some(key) => ks.with_encryption(key),
None => ks,
}
};
if keys_ks_boot
.get_raw(crate::operations::backup::IMPORT_IN_PROGRESS_KEY)
.await?
.is_some()
{
return Err(AppError::Internal(
"a previous backup import did not complete — the store is in a \
half-imported, inconsistent state. Re-run the import to restore a \
consistent snapshot before starting the VTA."
.into(),
));
}
}
let boot_service_state_ks = {
let ks = store.keyspace(crate::keyspaces::SERVICE_STATE)?;
match storage_encryption_key {
Some(key) => ks.with_encryption(key),
None => ks,
}
};
#[cfg(feature = "webvh")]
{
crate::operations::protocol::runtime_state::migrate_from_config(
&boot_service_state_ks,
&config,
)
.await?;
config.services.rest =
crate::operations::protocol::runtime_state::is_rest_enabled(&boot_service_state_ks)
.await?;
config.services.didcomm =
crate::operations::protocol::runtime_state::is_didcomm_enabled(&boot_service_state_ks)
.await?;
}
#[cfg(not(feature = "webvh"))]
{
let _ = &boot_service_state_ks;
}
{
let keys_ks_boot = {
let ks = store.keyspace(crate::keyspaces::KEYS)?;
match storage_encryption_key {
Some(key) => ks.with_encryption(key),
None => ks,
}
};
match crate::keys::seeds::reconcile_archive(&keys_ks_boot, &*seed_store).await {
Ok(0) => {}
Ok(n) => info!(rewritten = n, "seed archive reconciled at boot"),
Err(e) => warn!(error = %e, "seed archive reconcile failed — continuing"),
}
}
#[cfg(feature = "tee")]
if let Some(storage_key) = storage_encryption_key
&& let Some(kms) = config.tee.kms.as_ref()
{
let enc = |name: &str| -> Result<KeyspaceHandle, AppError> {
Ok(store.keyspace(name)?.with_encryption(storage_key))
};
let anchor: Option<Arc<dyn vti_common::integrity::AnchorCounter>> =
match (kms.anchor.as_ref(), config.vta_did.as_ref()) {
(Some(anchor_cfg), Some(vta_did)) => {
let writer = match anchor_cfg.writer_credential_ciphertext.as_ref() {
Some(b64) => {
let ct = base64::engine::general_purpose::STANDARD
.decode(b64)
.map_err(|e| {
AppError::Config(format!(
"tee.kms.anchor.writer_credential_ciphertext is not \
valid base64: {e}"
))
})?;
let pt = crate::tee::kms_bootstrap::attested_decrypt(kms, &ct).await?;
let creds: crate::tee::anchor::WriterCredentials =
serde_json::from_slice(&pt).map_err(|e| {
AppError::Config(format!(
"anchor writer credential did not decrypt to \
{{access_key_id, secret_access_key}}: {e}"
))
})?;
info!("anchor writer credential unsealed (attestation-gated, P0.2c)");
Some(creds)
}
None => None,
};
Some(Arc::new(
crate::tee::anchor::DynamoAnchorCounter::new(
&kms.region,
anchor_cfg.table_name.clone(),
vta_did.clone(),
writer,
)
.await,
))
}
(Some(_), None) => {
warn!(
"tee.kms.anchor is configured but vta_did is unset — booting \
manifest-only (P0.2a); the external rollback counter is disabled"
);
None
}
(None, _) => None,
};
let outcome = vti_common::integrity::boot_verify_and_install(
vti_common::integrity::derive_mac_key(&storage_key),
enc("keys")?,
store.keyspace(crate::keyspaces::BOOTSTRAP)?, enc("acl")?,
enc("contexts")?,
anchor,
kms.allow_anchor_init,
kms.allow_unanchored,
)
.await?;
info!(?outcome, "TEE anti-rollback anchor checked");
}
let rest_enabled = cfg!(feature = "rest") && config.services.rest;
let didcomm_enabled = cfg!(feature = "didcomm") && config.services.didcomm;
if !rest_enabled && !didcomm_enabled {
return Err(AppError::Config(
"no services enabled — enable at least one of REST or DIDComm \
(compile-time feature flags + `pnm services {kind} enable`)"
.into(),
));
}
#[cfg(feature = "rest")]
let std_listener = if rest_enabled {
let addr = format!("{}:{}", config.server.host, config.server.port);
let listener = std::net::TcpListener::bind(&addr).map_err(AppError::Io)?;
listener.set_nonblocking(true).map_err(AppError::Io)?;
info!("server listening addr={addr}");
Some(listener)
} else {
None
};
#[cfg(feature = "rest")]
let metrics_handle = if rest_enabled {
Some(crate::metrics::install())
} else {
None
};
loop {
let apply_encryption = |ks: KeyspaceHandle| -> KeyspaceHandle {
match storage_encryption_key {
Some(key) => ks.with_encryption(key),
None => ks,
}
};
let sessions_ks = apply_encryption(store.keyspace(crate::keyspaces::SESSIONS)?);
let acl_ks = apply_encryption(store.keyspace(crate::keyspaces::ACL)?);
let audit_ks = apply_encryption(store.keyspace(crate::keyspaces::AUDIT)?);
let consent_ks = apply_encryption(store.keyspace(crate::keyspaces::CONSENT)?);
let idempotency_ks = apply_encryption(store.keyspace(crate::keyspaces::IDEMPOTENCY)?);
let task_consent_ks = apply_encryption(store.keyspace(crate::keyspaces::TASK_CONSENT)?);
let vault_ks = apply_encryption(store.keyspace(crate::keyspaces::VAULT)?);
let backup_bundles_ks = apply_encryption(store.keyspace(crate::keyspaces::BACKUP_BUNDLES)?);
let backup_blob_dir = config.store.data_dir.join("backups");
#[cfg(all(feature = "webvh", feature = "didcomm"))]
let drains_ks = apply_encryption(store.keyspace(crate::keyspaces::DRAINS)?);
let telemetry: vti_common::telemetry::SharedTelemetrySink =
Arc::new(vti_common::telemetry::RingBufferTelemetry::new());
#[cfg(all(feature = "webvh", feature = "didcomm"))]
let mediator_registry = Arc::new(
crate::messaging::registry::MediatorListenerRegistry::new(Arc::clone(&telemetry)),
);
#[cfg(all(feature = "webvh", feature = "didcomm"))]
let (teardown_tx, teardown_rx) = crate::messaging::drain_sweeper::teardown_channel(
crate::messaging::drain_sweeper::DEFAULT_TEARDOWN_CHANNEL_CAPACITY,
);
#[cfg(all(feature = "webvh", feature = "didcomm"))]
let drain_sweeper = Arc::new(crate::messaging::drain_sweeper::DrainSweeper::new(
Arc::clone(&mediator_registry),
drains_ks.clone(),
teardown_tx,
));
#[cfg(all(feature = "webvh", feature = "didcomm"))]
match mediator_registry.replay_drains(&drains_ks).await {
Ok(live) => {
if !live.is_empty() {
info!(count = live.len(), "drain set replayed from keyspace");
}
drain_sweeper.arm_all(&live).await;
}
Err(e) => {
warn!(error = %e, "drain replay failed — starting with empty drain set");
}
}
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let (restart_tx, mut restart_rx) = watch::channel(false);
#[cfg(any(feature = "didcomm", feature = "tsp"))]
let didcomm_shutdown = CancellationToken::new();
tokio::spawn({
let shutdown_tx = shutdown_tx.clone();
#[cfg(any(feature = "didcomm", feature = "tsp"))]
let didcomm_shutdown = didcomm_shutdown.clone();
async move {
shutdown_signal().await;
info!("shutting down — press Ctrl-C again to force exit");
let _ = shutdown_tx.send(true);
#[cfg(any(feature = "didcomm", feature = "tsp"))]
didcomm_shutdown.cancel();
shutdown_signal().await;
eprintln!("\nForcing exit.");
std::process::exit(130);
}
});
let storage_store = store.clone();
let storage_sessions_ks = sessions_ks.clone();
let storage_audit_ks = audit_ks.clone();
let audit_sink: vta_audit::SharedAuditSink =
Arc::new(vta_audit::KeyspaceAuditSink::new(audit_ks.clone()));
let storage_audit_sink = Arc::clone(&audit_sink);
let storage_acl_ks = acl_ks.clone();
let storage_consent_ks = consent_ks.clone();
let storage_idempotency_ks = idempotency_ks.clone();
let storage_task_consent_ks = task_consent_ks.clone();
let storage_vault_ks = vault_ks.clone();
let storage_app_state_ks = apply_encryption(store.keyspace(crate::keyspaces::APP_STATE)?);
let storage_app_state_retention_days = config.app_state.tombstone_retention_days;
let storage_backup_bundles_ks = backup_bundles_ks.clone();
let storage_backup_blob_dir = backup_blob_dir.clone();
let storage_audit_config = config.audit.clone();
let storage_auth_config = config.auth.clone();
#[cfg(any(feature = "didcomm", feature = "tsp"))]
let didcomm_bridge: Arc<DIDCommBridge> = Arc::new(DIDCommBridge::new("vta-main"));
#[cfg(any(feature = "rest", feature = "didcomm"))]
let app_state = {
let parts = AppStateParts {
telemetry: Some(Arc::clone(&telemetry)),
audit_sink: Some(Arc::clone(&audit_sink)),
#[cfg(all(feature = "webvh", feature = "didcomm"))]
mediator_registry: Some(Arc::clone(&mediator_registry)),
#[cfg(all(feature = "webvh", feature = "didcomm"))]
drain_sweeper: Some(Arc::clone(&drain_sweeper)),
#[cfg(any(feature = "didcomm", feature = "tsp"))]
didcomm_bridge: Some(didcomm_bridge.clone()),
#[cfg(feature = "rest")]
metrics_handle: metrics_handle.clone(), };
build_app_state(
config.clone(),
&store,
seed_store.clone(),
storage_encryption_key,
tee_context.clone(),
restart_tx.clone(),
parts,
)
.await?
};
let storage_app_state_locks = app_state.app_state_locks.clone();
#[cfg(any(feature = "rest", feature = "didcomm"))]
app_state.wrapping_cache.clone().spawn_reaper();
#[cfg(any(feature = "rest", feature = "didcomm"))]
{
crate::policy::install_default_policy(
&app_state.policy_ks,
&chrono::Utc::now().to_rfc3339(),
)
.await?;
crate::policy::remove_stale_config_consent_policy(&app_state.policy_ks).await?;
{
let cfg = app_state.config.read().await;
crate::policy::seed_declarative_approvals(
&app_state.policy_ks,
&cfg.policy.approvals,
&cfg.policy.approver_sets,
&chrono::Utc::now().to_rfc3339(),
)
.await?;
}
}
#[cfg(any(feature = "rest", feature = "didcomm"))]
if app_state.jwt_keys.is_none() && !allow_degraded {
return Err(AppError::Config(missing_identity_message(&config)));
}
#[cfg(any(feature = "rest", feature = "didcomm"))]
let has_auth = app_state.jwt_keys.is_some();
#[cfg(not(any(feature = "rest", feature = "didcomm")))]
let has_auth = false;
#[cfg(feature = "tee")]
if config.tee.mode == crate::config::TeeMode::Required && !has_auth {
warn!(
"TEE mode is 'required' but authentication is not initialized \
(vta_did not configured). The VTA will start but authenticated \
endpoints will return 401."
);
}
#[cfg(feature = "rest")]
let rest_handle = if let Some(ref listener_ref) = std_listener {
let listener = listener_ref.try_clone().map_err(AppError::Io)?;
let state = app_state.clone();
let mut rest_shutdown_rx = shutdown_rx.clone();
Some(
std::thread::Builder::new()
.name("vta-rest".into())
.spawn(move || run_rest_thread(listener, state, &mut rest_shutdown_rx))
.map_err(|e| AppError::Internal(format!("failed to spawn REST thread: {e}")))?,
)
} else {
None
};
#[cfg(not(feature = "rest"))]
let rest_handle: Option<std::thread::JoinHandle<()>> = None;
#[cfg(any(feature = "didcomm", feature = "tsp"))]
let readiness_fatal = Arc::new(std::sync::atomic::AtomicBool::new(false));
#[cfg(any(feature = "didcomm", feature = "tsp"))]
if config.services.didcomm || config.services.tsp {
match (
&app_state.secrets_resolver,
&config.vta_did,
&config.messaging,
) {
(Some(_), Some(vta_did), Some(messaging_config)) => {
let outbox_ks = apply_encryption(store.keyspace(crate::keyspaces::OUTBOX)?);
let supervisor = MessagingConnect {
app_state: app_state.clone(),
vta_did: vta_did.clone(),
messaging_config: messaging_config.clone(),
readiness: config.mediator_readiness.clone(),
resolver_url: config.resolver_url.clone(),
outbox_ks,
flush_queues,
shutdown: didcomm_shutdown.clone(),
fatal_shutdown: shutdown_tx.clone(),
fatal_flag: readiness_fatal.clone(),
};
tokio::spawn(supervisor.run());
}
_ => {
info!("DIDComm not configured — service not started");
}
}
}
#[cfg(all(feature = "webvh", feature = "didcomm"))]
let _teardown_handle = {
let teardown_app_state = app_state.clone();
let mut teardown_rx = teardown_rx;
let mut shutdown_rx_for_teardown = shutdown_rx.clone();
tokio::spawn(async move {
loop {
tokio::select! {
biased;
_ = shutdown_rx_for_teardown.changed() => {
if *shutdown_rx_for_teardown.borrow() {
break;
}
}
msg = teardown_rx.recv() => {
match msg {
None => break,
Some(mediator_did) => {
if let Some(svc) =
teardown_app_state.didcomm_bridge.messaging_handle()
{
svc.remove_transport(&mediator_did);
info!(
mediator = %mediator_did,
"drain teardown: transport removed"
);
} else {
debug!(
mediator = %mediator_did,
"drain teardown: DIDComm not connected, skipping remove_transport"
);
}
}
}
}
}
}
debug!("teardown consumer task exiting");
})
};
#[cfg(not(all(feature = "webvh", feature = "didcomm")))]
let _teardown_handle: Option<tokio::task::JoinHandle<()>> = None;
let mut storage_shutdown_rx = shutdown_rx.clone();
let storage_handle = std::thread::Builder::new()
.name("vta-storage".into())
.spawn(move || {
run_storage_thread(
storage_store,
storage_sessions_ks,
storage_audit_ks,
storage_audit_sink,
storage_acl_ks,
storage_consent_ks,
storage_idempotency_ks,
storage_task_consent_ks,
storage_vault_ks,
storage_app_state_ks,
storage_app_state_locks,
storage_app_state_retention_days,
storage_backup_bundles_ks,
storage_backup_blob_dir,
storage_audit_config,
storage_auth_config,
has_auth,
&mut storage_shutdown_rx,
)
})
.map_err(|e| AppError::Internal(format!("failed to spawn storage thread: {e}")))?;
let mut any_panic = false;
let is_restart;
if let Some(handle) = rest_handle {
tokio::select! {
result = tokio::task::spawn_blocking(move || handle.join()) => {
match result {
Ok(Ok(())) => info!("REST thread stopped"),
Ok(Err(_panic)) => { error!("REST thread panicked"); any_panic = true; }
Err(e) => { error!("failed to join REST thread: {e}"); any_panic = true; }
}
is_restart = false;
}
_ = restart_rx.changed() => {
info!("soft restart requested — shutting down services");
let _ = shutdown_tx.send(true);
is_restart = true;
}
}
} else {
tokio::select! {
_ = async {
let mut wait_rx = shutdown_rx.clone();
let _ = wait_rx.changed().await;
} => {
is_restart = false;
}
_ = restart_rx.changed() => {
info!("soft restart requested — shutting down services");
let _ = shutdown_tx.send(true);
is_restart = true;
}
}
}
#[cfg(any(feature = "didcomm", feature = "tsp"))]
{
didcomm_shutdown.cancel();
info!("mediator messaging stopped");
}
if any_panic {
let _ = shutdown_tx.send(true);
}
match storage_handle.join() {
Ok(()) => info!("storage thread stopped"),
Err(_panic) => {
error!("storage thread panicked");
any_panic = true;
}
}
if any_panic {
return Err(AppError::Internal("one or more threads panicked".into()));
}
#[cfg(any(feature = "didcomm", feature = "tsp"))]
if readiness_fatal.load(std::sync::atomic::Ordering::Relaxed) {
return Err(AppError::Internal(
"mediator self-readiness gate timed out with on_timeout = \"fail\"".into(),
));
}
if !is_restart {
info!("server shut down");
return Ok(());
}
info!("soft restart: re-initializing services");
}
}
#[allow(clippy::too_many_arguments)]
fn run_storage_thread(
store: Store,
sessions_ks: KeyspaceHandle,
audit_ks: KeyspaceHandle,
audit_sink: vta_audit::SharedAuditSink,
acl_ks: KeyspaceHandle,
consent_ks: KeyspaceHandle,
idempotency_ks: KeyspaceHandle,
task_consent_ks: KeyspaceHandle,
vault_ks: KeyspaceHandle,
app_state_ks: KeyspaceHandle,
app_state_locks: crate::operations::app_state::NamespaceLocks,
app_state_retention_days: u32,
backup_bundles_ks_storage: KeyspaceHandle,
backup_blob_dir_storage: std::path::PathBuf,
audit_config: crate::config::AuditConfig,
auth_config: AuthConfig,
has_auth: bool,
shutdown_rx: &mut watch::Receiver<bool>,
) {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("failed to build storage runtime");
rt.block_on(async {
info!("storage thread started");
if has_auth {
let interval = Duration::from_secs(auth_config.session_cleanup_interval);
let mut timer = tokio::time::interval(interval);
timer.tick().await;
loop {
tokio::select! {
_ = timer.tick() => {
if let Err(e) = cleanup_expired_sessions(&sessions_ks, auth_config.challenge_ttl).await {
warn!("session cleanup error: {e}");
}
let audit_retention = audit_config.retention_days;
if let Err(e) = crate::audit::cleanup_expired_logs(&audit_ks, audit_retention).await {
warn!("audit cleanup error: {e}");
}
if let Err(e) =
crate::acl_sweeper::sweep_expired(&acl_ks, &audit_sink).await
{
warn!("acl sweeper error: {e}");
}
if let Err(e) =
crate::consent_sweeper::sweep_expired(&consent_ks, &audit_sink).await
{
warn!("consent sweeper error: {e}");
}
crate::idempotency_sweeper::sweep_expired_logged(&idempotency_ks).await;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
match crate::policy::consent::sweep_expired(&task_consent_ks, now).await {
Ok(n) if n > 0 => debug!("task-consent sweeper pruned {n} rows"),
Ok(_) => {}
Err(e) => warn!("task-consent sweeper error: {e}"),
}
if let Err(e) = crate::backup_bundle_sweeper::sweep_bundles(
&backup_bundles_ks_storage,
&backup_blob_dir_storage,
)
.await
{
warn!("backup bundle sweeper error: {e}");
}
match crate::operations::credential_exchange::pending::sweep(
&vault_ks,
chrono::Utc::now(),
)
.await
{
Ok(n) if n > 0 => {
info!(reclaimed = n, "pending-present sweeper")
}
Ok(_) => {}
Err(e) => warn!("pending-present sweeper error: {e}"),
}
if let Err(e) =
crate::vault_sweeper::sweep_expired(&vault_ks, &audit_sink).await
{
warn!("vault sweeper error: {e}");
}
if app_state_retention_days > 0 {
match crate::operations::app_state::sweep_expired_tombstones(
&app_state_ks,
&app_state_locks,
&audit_sink,
u64::from(app_state_retention_days) * 24 * 60 * 60,
)
.await
{
Ok(n) if n > 0 => {
info!(reaped = n, "app-state tombstone sweeper")
}
Ok(_) => {}
Err(e) => warn!("app-state tombstone sweeper error: {e}"),
}
}
}
_ = shutdown_rx.changed() => {
info!("storage thread shutting down");
break;
}
}
}
} else {
let _ = shutdown_rx.changed().await;
info!("storage thread shutting down");
}
if let Err(e) = store.persist().await {
error!("failed to persist store on shutdown: {e}");
} else {
info!("store persisted");
}
});
}
#[cfg(feature = "rest")]
fn run_rest_thread(
std_listener: std::net::TcpListener,
state: AppState,
shutdown_rx: &mut watch::Receiver<bool>,
) {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("failed to build REST runtime");
rt.block_on(async {
info!("REST thread started");
let listener = tokio::net::TcpListener::from_std(std_listener)
.expect("failed to convert std TcpListener to tokio TcpListener");
let (cors_origins, trust_xff, interval_secs, burst) = {
let cfg = state.config.read().await;
(
cfg.server.cors_origins.clone(),
cfg.server.trust_xff,
cfg.server.rate_limit_interval_secs,
cfg.server.rate_limit_burst,
)
};
let traced_routes =
routes::router_with_cors(&cors_origins, trust_xff, interval_secs, burst)
.with_state(state.clone())
.layer(axum::middleware::from_fn(crate::metrics::track_metrics))
.layer(
TraceLayer::new_for_http()
.make_span_with(DefaultMakeSpan::new().level(Level::INFO))
.on_request(DefaultOnRequest::new().level(Level::INFO))
.on_response(DefaultOnResponse::new().level(Level::INFO)),
);
let app =
traced_routes.merge(routes::health_router_with_cors(&cors_origins).with_state(state));
let shutdown_rx = shutdown_rx.clone();
axum::serve(
listener,
app.into_make_service_with_connect_info::<std::net::SocketAddr>(),
)
.with_graceful_shutdown(async move {
let mut rx = shutdown_rx;
let _ = rx.changed().await;
})
.await
.expect("axum serve failed");
info!("REST thread shutting down");
});
}
struct AuthInit {
did_resolver: Option<DIDCacheClient>,
secrets_resolver: Option<Arc<ThreadedSecretsResolver>>,
jwt_keys: Option<Arc<JwtKeys>>,
atm: Option<ATM>,
#[cfg(feature = "tsp")]
tsp_profile: Option<std::sync::Arc<affinidi_tdk::messaging::profiles::ATMProfile>>,
#[cfg_attr(not(feature = "didcomm"), allow(dead_code))]
signing_vm_id: Option<String>,
#[cfg_attr(not(feature = "didcomm"), allow(dead_code))]
ka_vm_id: Option<String>,
}
impl AuthInit {
fn empty() -> Self {
Self {
did_resolver: None,
secrets_resolver: None,
jwt_keys: None,
atm: None,
#[cfg(feature = "tsp")]
tsp_profile: None,
signing_vm_id: None,
ka_vm_id: None,
}
}
}
fn missing_identity_message(config: &AppConfig) -> String {
let cause = if config.vta_did.is_none() {
"vta_did is not configured — this VTA has no identity. Run `vta setup` \
to provision one"
.to_string()
} else if config.auth.jwt_signing_key.is_none() {
"auth.jwt_signing_key is not configured — the VTA can't issue access \
tokens. Run `vta setup`, or restore the key to config.toml"
.to_string()
} else {
format!(
"vta_did is set ({}) but its signing identity could not be loaded — \
the VTA key records may be missing from the store or the seed \
backend may be unreachable. Check the store data_dir and the \
secrets backend",
config.vta_did.as_deref().unwrap_or_default()
)
};
format!(
"refusing to start: {cause}. A VTA without a usable signing identity \
boots but answers every authenticated request with 401. To start \
anyway (e.g. to inspect or finish provisioning a half-set-up \
instance), pass `--allow-degraded`."
)
}
async fn init_auth(
config: &AppConfig,
seed_store: &dyn SeedStore,
keys_ks: &KeyspaceHandle,
webvh_ks: Option<&KeyspaceHandle>,
) -> AuthInit {
let vta_did = match &config.vta_did {
Some(did) => did.clone(),
None => {
warn!("vta_did not configured — auth endpoints will not work (run setup first)");
return AuthInit::empty();
}
};
let (signing_path, ka_path, vta_seed_id) = match find_vta_key_paths(&vta_did, keys_ks).await {
Ok(paths) => paths,
Err(e) => {
warn!(
"failed to find VTA key records: {e} — auth endpoints will not work (run setup first)"
);
return AuthInit::empty();
}
};
let seed = match load_seed_bytes(keys_ks, seed_store, vta_seed_id).await {
Ok(s) => s,
Err(e) => {
warn!("failed to load seed: {e} — auth endpoints will not work");
return AuthInit::empty();
}
};
let root = match ExtendedSigningKey::from_seed(&seed) {
Ok(r) => r,
Err(e) => {
warn!("failed to create BIP-32 root key: {e} — auth endpoints will not work");
return AuthInit::empty();
}
};
let resolver_config = {
let mut builder = DIDCacheConfigBuilder::default();
if let Some(ref url) = config.resolver_url {
info!(url = %url, "DID resolver using network mode (remote resolver)");
builder = builder.with_network_mode(url);
} else {
info!("DID resolver using local mode");
}
builder.build()
};
let mut did_resolver = match DIDCacheClient::new(resolver_config).await {
Ok(r) => r,
Err(e) => {
warn!("failed to create DID resolver: {e} — auth endpoints will not work");
return AuthInit::empty();
}
};
preload_self_did_document(&mut did_resolver, &vta_did, webvh_ks).await;
let (secrets_resolver, _handle) = ThreadedSecretsResolver::new(None).await;
let mut signing_vm_id: Option<String> = None;
let mut ka_vm_id: Option<String> = None;
if vta_did.starts_with("did:key:") {
let dp: vti_common::slip10::DerivationPath = match signing_path.parse() {
Ok(p) => p,
Err(e) => {
warn!("invalid signing derivation path: {e}");
return AuthInit {
did_resolver: Some(did_resolver),
..AuthInit::empty()
};
}
};
match root.derive(&dp) {
Ok(derived) => {
let seed_bytes: &[u8; 32] = derived.signing_key.as_bytes();
match vta_sdk::did_key::secrets_from_did_key(&vta_did, seed_bytes) {
Ok(secrets) => {
signing_vm_id = Some(secrets.signing.id.clone());
ka_vm_id = Some(secrets.key_agreement.id.clone());
info!(signing_id = %secrets.signing.id, ka_id = %secrets.key_agreement.id, "did:key secrets loaded");
secrets_resolver.insert(secrets.signing).await;
secrets_resolver.insert(secrets.key_agreement).await;
}
Err(e) => {
warn!("failed to build did:key secrets: {e} — auth will not work");
return AuthInit {
did_resolver: Some(did_resolver),
..AuthInit::empty()
};
}
}
}
Err(e) => warn!("failed to derive VTA signing key: {e}"),
}
} else {
let ka_path = match ka_path {
Some(p) => p,
None => {
warn!(
"VTA key-agreement record missing — auth endpoints will not work (run setup first)"
);
return AuthInit {
did_resolver: Some(did_resolver),
..AuthInit::empty()
};
}
};
signing_vm_id = Some(format!("{vta_did}#key-0"));
ka_vm_id = Some(format!("{vta_did}#key-1"));
let stored_signing: Option<KeyRecord> = keys_ks
.get(crate::keys::store_key(&format!("{vta_did}#key-0")))
.await
.ok()
.flatten();
let stored_ka: Option<KeyRecord> = keys_ks
.get(crate::keys::store_key(&format!("{vta_did}#key-1")))
.await
.ok()
.flatten();
match root.derive_ed25519(&signing_path) {
Ok(mut signing_secret) => {
if let Some(ref record) = stored_signing {
match signing_secret.get_public_keymultibase() {
Ok(runtime_pub) if runtime_pub != record.public_key => {
error!(
key_id = %format!("{vta_did}#key-0"),
stored = %record.public_key,
runtime = %runtime_pub,
"SIGNING KEY MISMATCH: runtime-derived Ed25519 public key does not match \
the key stored in the key record (and published in the DID document). \
DIDComm message signing/verification will fail. \
This likely means the DID was created with different code or seed."
);
}
Ok(runtime_pub) => {
info!(key_id = %format!("{vta_did}#key-0"), pub_key = %runtime_pub, "signing key validated");
}
Err(e) => warn!("could not extract signing public key for validation: {e}"),
}
}
signing_secret.id = format!("{vta_did}#key-0");
secrets_resolver.insert(signing_secret).await;
}
Err(e) => warn!("failed to derive VTA signing key: {e}"),
}
match root.derive_x25519(&ka_path) {
Ok(mut ka_secret) => {
if let Some(ref record) = stored_ka {
match ka_secret.get_public_keymultibase() {
Ok(runtime_pub) if runtime_pub != record.public_key => {
error!(
key_id = %format!("{vta_did}#key-1"),
stored = %record.public_key,
runtime = %runtime_pub,
"KEY-AGREEMENT KEY MISMATCH: runtime-derived X25519 public key does not match \
the key stored in the key record (and published in the DID document). \
DIDComm encryption/decryption will fail. Others will encrypt to the DID \
document key but this VTA holds a different private key. \
The DID document must be updated or the VTA identity must be regenerated."
);
}
Ok(runtime_pub) => {
info!(key_id = %format!("{vta_did}#key-1"), pub_key = %runtime_pub, "key-agreement key validated");
}
Err(e) => warn!("could not extract KA public key for validation: {e}"),
}
}
ka_secret.id = format!("{vta_did}#key-1");
secrets_resolver.insert(ka_secret).await;
}
Err(e) => warn!("failed to derive VTA key-agreement key: {e}"),
}
}
let jwt_keys = match &config.auth.jwt_signing_key {
Some(b64) => match decode_jwt_key(b64) {
Ok(k) => k,
Err(e) => {
warn!("failed to load JWT signing key: {e} — auth endpoints will not work");
return AuthInit {
did_resolver: Some(did_resolver),
secrets_resolver: Some(Arc::new(secrets_resolver)),
signing_vm_id,
ka_vm_id,
..AuthInit::empty()
};
}
},
None => {
warn!(
"auth.jwt_signing_key not configured — auth endpoints will not work (run setup first)"
);
return AuthInit {
did_resolver: Some(did_resolver),
secrets_resolver: Some(Arc::new(secrets_resolver)),
signing_vm_id,
ka_vm_id,
..AuthInit::empty()
};
}
};
let secrets_resolver = Arc::new(secrets_resolver);
let atm = {
let tdk_config = TDKConfig::builder()
.with_did_resolver(did_resolver.clone())
.with_secrets_resolver((*secrets_resolver).clone())
.with_load_environment(false)
.build();
match tdk_config {
Ok(cfg) => match TDKSharedState::new(cfg).await {
Ok(tdk) => {
match ATM::new(ATMConfig::builder().build().unwrap(), Arc::new(tdk)).await {
Ok(a) => Some(a),
Err(e) => {
warn!("failed to create ATM for auth unpack: {e}");
None
}
}
}
Err(e) => {
warn!("failed to create TDK shared state: {e}");
None
}
},
Err(e) => {
warn!("failed to build TDK config: {e}");
None
}
}
};
#[cfg(feature = "tsp")]
let tsp_profile = if let Some(ref atm) = atm {
match affinidi_tdk::messaging::profiles::ATMProfile::new(
atm,
Some("VTA".to_string()),
vta_did.clone(),
None,
)
.await
{
Ok(profile) => match atm.profile_add(&profile, false).await {
Ok(arc) => {
info!("TSP profile registered for DID {vta_did}");
Some(arc)
}
Err(e) => {
warn!("failed to register TSP profile (TSP unseal disabled): {e}");
None
}
},
Err(e) => {
warn!("failed to build TSP profile (TSP unseal disabled): {e}");
None
}
}
} else {
warn!("ATM unavailable — TSP profile not built (TSP unseal disabled)");
None
};
info!("auth initialized for DID {vta_did}");
AuthInit {
did_resolver: Some(did_resolver),
secrets_resolver: Some(secrets_resolver),
jwt_keys: Some(Arc::new(jwt_keys)),
atm,
#[cfg(feature = "tsp")]
tsp_profile,
signing_vm_id,
ka_vm_id,
}
}
pub(crate) async fn preload_self_did_document(
did_resolver: &mut DIDCacheClient,
vta_did: &str,
webvh_ks: Option<&KeyspaceHandle>,
) {
#[cfg(not(feature = "webvh"))]
{
let _ = (&did_resolver, vta_did, webvh_ks);
}
#[cfg(feature = "webvh")]
{
if !vta_did.starts_with("did:webvh:") {
return;
}
let Some(webvh_ks) = webvh_ks else {
warn!(
did = %vta_did,
"webvh keyspace not available; self DID preload skipped"
);
return;
};
let Some(did_log) = (match crate::webvh_store::get_did_log(webvh_ks, vta_did).await {
Ok(log) => log,
Err(e) => {
warn!(did = %vta_did, error = %e, "failed to read local did.jsonl for resolver preload");
return;
}
}) else {
warn!(did = %vta_did, "no local did.jsonl found for resolver preload");
return;
};
let doc_value = match crate::operations::protocol::document::current_document_from_log(
&did_log,
) {
Ok(doc) => doc,
Err(e) => {
warn!(did = %vta_did, error = %e, "failed to parse local did.jsonl for resolver preload");
return;
}
};
let doc = match serde_json::from_value(doc_value) {
Ok(doc) => doc,
Err(e) => {
warn!(did = %vta_did, error = %e, "failed to decode DID document for resolver preload");
return;
}
};
did_resolver.add_did_document(vta_did, doc).await;
info!(did = %vta_did, "preloaded VTA DID into resolver cache from local did.jsonl");
}
}
async fn find_vta_key_paths(
vta_did: &str,
keys_ks: &KeyspaceHandle,
) -> Result<(String, Option<String>, Option<u32>), AppError> {
let signing_key_id = format!("{vta_did}#key-0");
let signing: KeyRecord = keys_ks
.get(crate::keys::store_key(&signing_key_id))
.await?
.ok_or_else(|| AppError::NotFound("VTA signing key not found".into()))?;
let ka_path = if vta_did.starts_with("did:key:") {
None
} else {
let ka_key_id = format!("{vta_did}#key-1");
let ka: KeyRecord = keys_ks
.get(crate::keys::store_key(&ka_key_id))
.await?
.ok_or_else(|| AppError::NotFound("VTA key-agreement key not found".into()))?;
Some(ka.derivation_path)
};
debug!(signing_path = %signing.derivation_path, ka_path = ?ka_path, "VTA key paths resolved");
Ok((signing.derivation_path, ka_path, signing.seed_id))
}
fn decode_jwt_key(b64: &str) -> Result<JwtKeys, AppError> {
let bytes = BASE64
.decode(b64)
.map_err(|e| AppError::Config(format!("invalid jwt_signing_key base64: {e}")))?;
let key_bytes: [u8; 32] = bytes
.try_into()
.map_err(|_| AppError::Config("jwt_signing_key must be exactly 32 bytes".into()))?;
let keys = JwtKeys::from_ed25519_bytes(&key_bytes, "VTA")?;
debug!("JWT signing key decoded successfully");
Ok(keys)
}
#[cfg(any(feature = "didcomm", feature = "tsp"))]
struct MessagingConnect {
app_state: AppState,
vta_did: String,
messaging_config: crate::config::MessagingConfig,
readiness: crate::config::MediatorReadinessConfig,
resolver_url: Option<String>,
outbox_ks: KeyspaceHandle,
flush_queues: bool,
shutdown: CancellationToken,
fatal_shutdown: watch::Sender<bool>,
fatal_flag: Arc<std::sync::atomic::AtomicBool>,
}
#[cfg(any(feature = "didcomm", feature = "tsp"))]
impl MessagingConnect {
async fn run(self) {
match crate::messaging::readiness::run_gate(
&self.vta_did,
&self.readiness,
self.resolver_url.as_deref(),
&self.shutdown,
)
.await
{
Ok(crate::messaging::readiness::GateDecision::Proceed) => {}
Ok(crate::messaging::readiness::GateDecision::Skip) => {
info!(
"DIDComm messaging not started this boot \
(mediator self-readiness gate: skip)"
);
return;
}
Err(e) => {
error!("mediator self-readiness gate failed, shutting down: {e}");
self.fatal_flag
.store(true, std::sync::atomic::Ordering::Relaxed);
let _ = self.fatal_shutdown.send(true);
return;
}
}
self.run_startup_recovery().await;
self.supervise().await;
}
async fn supervise(&self) {
let policy = crate::messaging::readiness::ReconnectPolicy::from_config(&self.readiness);
let mut window_started = tokio::time::Instant::now();
let mut attempt: u32 = 0;
loop {
if self.shutdown.is_cancelled() {
return;
}
if crate::messaging::readiness::self_did_network_resolvable(
&self.vta_did,
self.resolver_url.as_deref(),
)
.await
{
match self.connect_once().await {
Ok(messaging) => {
info!("DIDComm messaging started");
let session_started = tokio::time::Instant::now();
crate::messaging::service::run_inbound_loop(
messaging.clone(),
self.app_state.clone(),
self.vta_did.clone(),
self.messaging_config.mediator_did.clone(),
self.shutdown.clone(),
)
.await;
self.app_state.didcomm_bridge.clear_messaging();
messaging.atm.graceful_shutdown().await;
if self.shutdown.is_cancelled() {
return;
}
let lasted = session_started.elapsed();
warn!(
session_secs = lasted.as_secs(),
"mediator session ended unexpectedly; reconnecting"
);
if policy.session_was_healthy(lasted) {
attempt = 0;
window_started = tokio::time::Instant::now();
}
}
Err(e) => warn!("failed to start DIDComm messaging: {e}"),
}
} else {
warn!(
vta_did = %self.vta_did,
"VTA not yet self-resolvable over the network; deferring mediator connect"
);
}
let Some(sleep_for) = policy.next_backoff(attempt, window_started.elapsed()) else {
if self.readiness.reconnect {
warn!(
waited_secs = window_started.elapsed().as_secs(),
"mediator reconnect gave up after the configured horizon; \
a later restart will retry"
);
}
return;
};
debug!(
attempt = attempt + 1,
ceiling_secs = policy.ceiling_for(attempt).as_secs_f64(),
sleep_secs = sleep_for.as_secs_f64(),
"mediator connection not established; backing off before retry \
(exponential backoff + full jitter)"
);
tokio::select! {
_ = self.shutdown.cancelled() => {
info!("shutdown during mediator reconnect backoff; stopping supervisor");
return;
}
_ = tokio::time::sleep(sleep_for) => {}
}
attempt = attempt.saturating_add(1);
}
}
async fn run_startup_recovery(&self) {
let messaging_config = &self.messaging_config;
let vta_did = self.vta_did.as_str();
if messaging_config.drain_inbox_on_start || self.flush_queues {
match self.app_state.atm.as_ref() {
Some(atm) => {
let cleared =
drain_mediator_inbox(atm, &messaging_config.mediator_did, vta_did).await;
info!(
count = cleared,
mediator = %messaging_config.mediator_did,
"cleared queued mediator inbox before going live"
);
}
None => warn!("inbox drain requested but no ATM is available; skipping"),
}
}
if self.flush_queues {
match self.app_state.atm.as_ref() {
Some(atm) => {
let cleared =
flush_mediator_outbox(atm, &messaging_config.mediator_did, vta_did).await;
info!(
count = cleared,
mediator = %messaging_config.mediator_did,
"flush_queues: cleared outbound sender queue before going live"
);
}
None => {
warn!("--flush-queues set but no ATM is available; skipping outbox flush")
}
}
}
}
async fn connect_once(&self) -> Result<Arc<crate::messaging::service::VtaMessaging>, String> {
let app_state = &self.app_state;
let messaging_config = &self.messaging_config;
let vta_did = self.vta_did.as_str();
let sr = app_state
.secrets_resolver
.as_ref()
.ok_or_else(|| "no secrets resolver available".to_string())?;
let mut secrets = Vec::new();
if let Some(ref signing_id) = app_state.signing_vm_id
&& let Some(s) = sr.get_secret(signing_id).await
{
secrets.push(s);
}
if let Some(ref ka_id) = app_state.ka_vm_id
&& let Some(s) = sr.get_secret(ka_id).await
{
secrets.push(s);
}
let messaging = Arc::new(
crate::messaging::service::build_messaging(
secrets,
vta_did,
&messaging_config.mediator_did,
self.outbox_ks.clone(),
app_state.did_resolver.as_ref(),
self.resolver_url.as_deref(),
)
.await?,
);
app_state.didcomm_bridge.set_messaging(
messaging.service.clone(),
(*messaging.atm).clone(),
vta_did.to_string(),
);
#[cfg(all(feature = "webvh", feature = "didcomm"))]
app_state
.mediator_registry
.record_activate(crate::messaging::registry::MediatorBinding {
mediator_did: messaging_config.mediator_did.clone(),
endpoint: messaging_config.mediator_url.clone(),
})
.await;
if messaging_config.setup_acl {
if let Some(atm) = app_state.atm.as_ref() {
acl_setup::set_client_acl_on_connection(
atm,
vta_did,
messaging_config.mediator_did.as_str(),
"vta-main",
"vta",
)
.await;
} else {
warn!("setup_acl = true but no ATM available; skipping mediator ACL provisioning");
}
}
Ok(messaging)
}
}
#[cfg(any(feature = "didcomm", feature = "tsp"))]
async fn drain_mediator_inbox(atm: &ATM, mediator_did: &str, vta_did: &str) -> usize {
let profile = match affinidi_tdk::messaging::profiles::ATMProfile::new(
atm,
Some("VTA-drain".to_string()),
vta_did.to_string(),
Some(mediator_did.to_string()),
)
.await
{
Ok(p) => match atm.profile_add(&p, false).await {
Ok(arc) => arc,
Err(e) => {
warn!(error = %e, "drain: could not register mediator profile; skipping drain");
return 0;
}
},
Err(e) => {
warn!(error = %e, "drain: could not build mediator profile; skipping drain");
return 0;
}
};
let mut cleared = 0usize;
for _round in 0..200 {
let batch = match atm
.message_pickup()
.send_delivery_request(&profile, Some(20), true)
.await
{
Ok(b) => b,
Err(e) => {
warn!(error = %e, cleared, "drain: delivery-request failed; stopping — a queued message may be un-fetchable, inspect the mediator");
break;
}
};
if batch.is_empty() {
break;
}
let ids: Vec<String> = batch.iter().map(|(msg, _meta)| msg.id.clone()).collect();
for (msg, _meta) in &batch {
warn!(id = %msg.id, msg_type = %msg.typ, "drain: clearing queued mediator message");
}
match atm
.message_pickup()
.send_messages_received(&profile, &ids, true)
.await
{
Ok(_) => cleared += ids.len(),
Err(e) => {
warn!(error = %e, cleared, "drain: delete failed; stopping to avoid a loop");
break;
}
}
}
cleared
}
#[cfg(any(feature = "didcomm", feature = "tsp"))]
async fn flush_mediator_outbox(atm: &ATM, mediator_did: &str, vta_did: &str) -> usize {
use affinidi_tdk::messaging::messages::{DeleteMessageRequest, Folder};
let profile = match affinidi_tdk::messaging::profiles::ATMProfile::new(
atm,
Some("VTA-flush-outbox".to_string()),
vta_did.to_string(),
Some(mediator_did.to_string()),
)
.await
{
Ok(p) => match atm.profile_add(&p, false).await {
Ok(arc) => arc,
Err(e) => {
warn!(error = %e, "flush outbox: could not register mediator profile; skipping");
return 0;
}
},
Err(e) => {
warn!(error = %e, "flush outbox: could not build mediator profile; skipping");
return 0;
}
};
let list = match atm.list_messages(&profile, Folder::Outbox).await {
Ok(l) => l,
Err(e) => {
warn!(error = %e, "flush outbox: could not list outbound queue; skipping");
return 0;
}
};
let ids: Vec<String> = list.iter().map(|m| m.msg_id.clone()).collect();
if ids.is_empty() {
return 0;
}
let mut cleared = 0usize;
for chunk in ids.chunks(100) {
match atm
.delete_messages_direct(
&profile,
&DeleteMessageRequest {
message_ids: chunk.to_vec(),
},
)
.await
{
Ok(r) => cleared += r.success.len(),
Err(e) => {
warn!(error = %e, cleared, "flush outbox: delete batch failed; stopping");
break;
}
}
}
cleared
}
async fn shutdown_signal() {
let ctrl_c = async {
tokio::signal::ctrl_c()
.await
.expect("failed to install Ctrl+C handler");
};
#[cfg(unix)]
let terminate = async {
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("failed to install SIGTERM handler")
.recv()
.await;
};
#[cfg(not(unix))]
let terminate = std::future::pending::<()>();
tokio::select! {
() = ctrl_c => info!("received SIGINT"),
() = terminate => info!("received SIGTERM"),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::keys::{KeyType, save_key_record};
use crate::store::Store;
#[cfg(feature = "webvh")]
use crate::webvh_store;
use affinidi_did_resolver_cache_sdk::config::DIDCacheConfigBuilder;
use vti_common::config::StoreConfig;
fn temp_keys_ks() -> (Store, KeyspaceHandle, tempfile::TempDir) {
let dir = tempfile::tempdir().expect("temp dir");
let store = Store::open(&StoreConfig {
data_dir: dir.path().to_path_buf(),
})
.expect("store open");
let keys_ks = store
.keyspace(crate::keyspaces::KEYS)
.expect("keys keyspace");
(store, keys_ks, dir)
}
#[cfg(feature = "webvh")]
fn temp_webvh_ks() -> (Store, KeyspaceHandle, tempfile::TempDir) {
let dir = tempfile::tempdir().expect("temp dir");
let store = Store::open(&StoreConfig {
data_dir: dir.path().to_path_buf(),
})
.expect("store open");
let webvh_ks = store
.keyspace(crate::keyspaces::WEBVH)
.expect("webvh keyspace");
(store, webvh_ks, dir)
}
#[tokio::test]
async fn find_vta_key_paths_returns_none_ka_for_did_key() {
let (_store, keys_ks, _dir) = temp_keys_ks();
let did = "did:key:z6MkTestKey";
save_key_record(
&keys_ks,
&format!("{did}#key-0"),
"m/44'/0'/0'",
KeyType::Ed25519,
"z6MkSigningPub",
"VTA signing key",
Some("vta"),
Some(0),
)
.await
.unwrap();
let (signing_path, ka_path, seed_id) =
find_vta_key_paths(did, &keys_ks).await.expect("paths");
assert_eq!(signing_path, "m/44'/0'/0'");
assert!(ka_path.is_none(), "did:key must not require #key-1 lookup");
assert_eq!(seed_id, Some(0));
}
#[tokio::test]
async fn find_vta_key_paths_loads_ka_for_did_webvh() {
let (_store, keys_ks, _dir) = temp_keys_ks();
let did = "did:webvh:abc:example.com:vta";
save_key_record(
&keys_ks,
&format!("{did}#key-0"),
"m/44'/0'/0'",
KeyType::Ed25519,
"z6MkSigningPub",
"VTA signing key",
Some("vta"),
Some(0),
)
.await
.unwrap();
save_key_record(
&keys_ks,
&format!("{did}#key-1"),
"m/44'/0'/1'",
KeyType::X25519,
"z6LSKaPub",
"VTA key-agreement key",
Some("vta"),
Some(0),
)
.await
.unwrap();
let (signing_path, ka_path, _seed_id) =
find_vta_key_paths(did, &keys_ks).await.expect("paths");
assert_eq!(signing_path, "m/44'/0'/0'");
assert_eq!(ka_path.as_deref(), Some("m/44'/0'/1'"));
}
#[tokio::test]
async fn find_vta_key_paths_errors_when_did_webvh_missing_ka() {
let (_store, keys_ks, _dir) = temp_keys_ks();
let did = "did:webvh:abc:example.com:vta";
save_key_record(
&keys_ks,
&format!("{did}#key-0"),
"m/44'/0'/0'",
KeyType::Ed25519,
"z6MkSigningPub",
"VTA signing key",
Some("vta"),
Some(0),
)
.await
.unwrap();
let result = find_vta_key_paths(did, &keys_ks).await;
assert!(
matches!(result, Err(AppError::NotFound(_))),
"expected NotFound for did:webvh missing #key-1, got {result:?}"
);
}
fn cfg(toml_str: &str) -> AppConfig {
toml::from_str::<AppConfig>(toml_str).expect("parse test config")
}
#[test]
fn missing_identity_message_names_absent_vta_did() {
let msg = missing_identity_message(&cfg(""));
assert!(msg.contains("vta_did is not configured"), "{msg}");
assert!(msg.contains("vta setup"), "{msg}");
assert!(msg.contains("--allow-degraded"), "{msg}");
}
#[test]
fn missing_identity_message_names_absent_jwt_key() {
let msg = missing_identity_message(&cfg("vta_did = \"did:key:z6MkTest\"\n"));
assert!(msg.contains("auth.jwt_signing_key"), "{msg}");
assert!(!msg.contains("vta_did is not configured"), "{msg}");
assert!(msg.contains("--allow-degraded"), "{msg}");
}
#[test]
fn missing_identity_message_falls_back_to_key_material() {
let msg = missing_identity_message(&cfg(
"vta_did = \"did:key:z6MkTest\"\n[auth]\njwt_signing_key = \"AAAA\"\n",
));
assert!(msg.contains("could not be loaded"), "{msg}");
assert!(msg.contains("did:key:z6MkTest"), "{msg}");
assert!(msg.contains("--allow-degraded"), "{msg}");
}
#[cfg(feature = "webvh")]
#[tokio::test]
async fn preload_self_did_document_makes_vta_did_resolvable_from_cache() {
let (_store, webvh_ks, _dir) = temp_webvh_ks();
let did = "did:webvh:QmScid:vta.example.com:vta";
let log_line = serde_json::json!({
"versionId": "1-test",
"versionTime": "2026-05-06T00:00:00Z",
"parameters": {},
"state": {
"@context": ["https://www.w3.org/ns/did/v1"],
"id": did,
},
});
let log = serde_json::to_string(&log_line).expect("serialize log line");
webvh_store::store_did_log(&webvh_ks, did, &log)
.await
.expect("store did log");
let mut resolver = DIDCacheClient::new(DIDCacheConfigBuilder::default().build())
.await
.expect("resolver init");
preload_self_did_document(&mut resolver, did, Some(&webvh_ks)).await;
let resolved = resolver.resolve(did).await.expect("resolve preloaded did");
assert!(
resolved.cache_hit,
"preloaded DID was not served from cache"
);
let expected_value = crate::operations::protocol::document::current_document_from_log(&log)
.expect("extract current DID document from did.jsonl");
let expected_doc =
serde_json::from_value(expected_value).expect("deserialize expected DID document");
assert_eq!(
resolved.doc, expected_doc,
"resolved DID document should match local did.jsonl current state"
);
}
#[cfg(feature = "webvh")]
#[tokio::test]
async fn preload_self_did_document_ignores_malformed_local_log() {
let (_store, webvh_ks, _dir) = temp_webvh_ks();
let did = "did:webvh:QmBadScid:vta.example.com:vta";
webvh_store::store_did_log(&webvh_ks, did, "not-json")
.await
.expect("store malformed did log");
let mut resolver = DIDCacheClient::new(DIDCacheConfigBuilder::default().build())
.await
.expect("resolver init");
preload_self_did_document(&mut resolver, did, Some(&webvh_ks)).await;
let result = resolver.resolve(did).await;
assert!(
result.is_err(),
"malformed preload input should not seed resolver cache"
);
}
}