use std::sync::Arc;
use chrono::{DateTime, Utc};
use tokio_util::sync::CancellationToken;
use tracing::warn;
use crate::audit::audit_logger::BaseAuditLogger;
use crate::audit::bounded::BoundedAuditLogger;
use crate::audit::model::AuditLogger;
use crate::capabilities::{
CapabilityDefinition, CapabilitySet, CapabilityStateStore, DispatcherHandle,
FileCapabilityStore,
};
use crate::configs::TrustRegistryConfig;
use crate::dedup::{MemoryMessageIdStore, MessageIdStore};
use crate::didcomm::listener::DidCommSource;
use crate::health::RegistryHealth;
use crate::http::application_routes;
use crate::storage::repository::{TrustRecordAdminRepository, TrustRecordRepository};
use crate::trust_tasks::{ReplySigner, TaskHandler};
use crate::{SharedData, server::ServerHandle};
type BoxError = Box<dyn std::error::Error + Send + Sync>;
pub struct TrustRegistry {
config: Arc<TrustRegistryConfig>,
repository: Arc<dyn TrustRecordAdminRepository>,
capabilities: Arc<CapabilitySet>,
verifier: Arc<dyn trust_tasks_rs::DynProofVerifier>,
dedup: Arc<dyn MessageIdStore>,
audit: Arc<BoundedAuditLogger>,
health: Arc<RegistryHealth>,
didcomm_source: DidCommSource,
shutdown: CancellationToken,
service_start_timestamp: DateTime<Utc>,
}
impl TrustRegistry {
pub fn builder(config: impl Into<Arc<TrustRegistryConfig>>) -> TrustRegistryBuilder {
TrustRegistryBuilder {
config: config.into(),
repository: None,
capabilities: None,
capability_store: None,
dedup: None,
verifier: None,
didcomm_source: None,
shutdown: None,
}
}
pub fn router(&self) -> axum::Router {
application_routes("", self.shared_data())
}
pub fn health_router(&self) -> axum::Router {
use axum::{Json, routing::get};
let health = self.health.clone();
axum::Router::new().route(
"/health",
get(move || {
let health = health.clone();
async move { Json(health.to_json()) }
}),
)
}
pub fn task_handler(&self) -> TaskHandler {
write_task_handler(
&self.config,
&self.capabilities,
&self.verifier,
&self.dedup,
&self.audit,
)
}
pub async fn route_didcomm_envelope(
&self,
body: serde_json::Value,
sender_did: &str,
) -> Option<Result<trust_tasks_rs::TrustTask<serde_json::Value>, trust_tasks_rs::ErrorResponse>>
{
crate::didcomm::handlers::trust_tasks::route_envelope_body(
&self.task_handler(),
body,
Some(sender_did),
)
.await
}
pub fn query_task_handler(&self) -> TaskHandler {
let handler = TaskHandler::new(
self.capabilities.query_dispatcher(),
self.config.didcomm_config.profile_config.did.clone(),
Vec::new(),
self.verifier.clone(),
)
.with_audit(self.audit.clone() as Arc<dyn AuditLogger>);
with_reply_signer(handler, &self.config)
}
pub fn dispatcher(&self) -> DispatcherHandle {
self.capabilities.dispatcher()
}
pub fn query_dispatcher(&self) -> DispatcherHandle {
self.capabilities.query_dispatcher()
}
pub fn capabilities(&self) -> &Arc<CapabilitySet> {
&self.capabilities
}
pub fn repository(&self) -> &Arc<dyn TrustRecordAdminRepository> {
&self.repository
}
pub fn config(&self) -> &Arc<TrustRegistryConfig> {
&self.config
}
pub fn health(&self) -> &Arc<RegistryHealth> {
&self.health
}
pub fn shutdown_token(&self) -> CancellationToken {
self.shutdown.clone()
}
fn shared_data(&self) -> SharedData<dyn TrustRecordRepository> {
SharedData {
config: self.config.clone(),
service_start_timestamp: self.service_start_timestamp,
repository: self.repository.clone() as Arc<dyn TrustRecordRepository>,
query_dispatcher: self.capabilities.query_dispatcher(),
verifier: self.verifier.clone(),
audit: self.audit.clone(),
}
}
pub async fn serve(self) -> Result<ServerHandle, BoxError> {
crate::server::serve_registry(self).await
}
pub(crate) fn into_parts(self) -> RegistryParts {
RegistryParts {
config: self.config,
repository: self.repository,
capabilities: self.capabilities,
verifier: self.verifier,
dedup: self.dedup,
audit: self.audit,
health: self.health,
didcomm_source: self.didcomm_source,
shutdown: self.shutdown,
service_start_timestamp: self.service_start_timestamp,
}
}
}
pub(crate) struct RegistryParts {
pub(crate) config: Arc<TrustRegistryConfig>,
pub(crate) repository: Arc<dyn TrustRecordAdminRepository>,
pub(crate) capabilities: Arc<CapabilitySet>,
pub(crate) verifier: Arc<dyn trust_tasks_rs::DynProofVerifier>,
pub(crate) dedup: Arc<dyn MessageIdStore>,
pub(crate) audit: Arc<BoundedAuditLogger>,
pub(crate) health: Arc<RegistryHealth>,
pub(crate) didcomm_source: DidCommSource,
pub(crate) shutdown: CancellationToken,
pub(crate) service_start_timestamp: DateTime<Utc>,
}
fn write_task_handler(
config: &TrustRegistryConfig,
capabilities: &Arc<CapabilitySet>,
verifier: &Arc<dyn trust_tasks_rs::DynProofVerifier>,
dedup: &Arc<dyn MessageIdStore>,
audit: &Arc<BoundedAuditLogger>,
) -> TaskHandler {
let admin_config = &config.didcomm_config.admin_config;
let handler = TaskHandler::new(
capabilities.dispatcher(),
config.didcomm_config.profile_config.did.clone(),
admin_config.admin_dids.clone(),
verifier.clone(),
)
.with_admin_authorities(admin_config.admin_authorities.clone())
.with_dedup(dedup.clone())
.with_capabilities(capabilities.clone())
.with_audit(audit.clone() as Arc<dyn AuditLogger>);
with_reply_signer(handler, config)
}
pub(crate) fn with_reply_signer(handler: TaskHandler, config: &TrustRegistryConfig) -> TaskHandler {
match ReplySigner::for_config(config) {
Some(signer) => handler.with_reply_signer(signer),
None => handler,
}
}
impl RegistryParts {
pub(crate) fn task_handler(&self) -> TaskHandler {
write_task_handler(
&self.config,
&self.capabilities,
&self.verifier,
&self.dedup,
&self.audit,
)
}
pub(crate) fn shared_data(&self) -> SharedData<dyn TrustRecordRepository> {
SharedData {
config: self.config.clone(),
service_start_timestamp: self.service_start_timestamp,
repository: self.repository.clone() as Arc<dyn TrustRecordRepository>,
query_dispatcher: self.capabilities.query_dispatcher(),
verifier: self.verifier.clone(),
audit: self.audit.clone(),
}
}
}
pub struct TrustRegistryBuilder {
config: Arc<TrustRegistryConfig>,
repository: Option<Arc<dyn TrustRecordAdminRepository>>,
capabilities: Option<Vec<CapabilityDefinition>>,
capability_store: Option<Box<dyn CapabilityStateStore>>,
dedup: Option<Arc<dyn MessageIdStore>>,
verifier: Option<Arc<dyn trust_tasks_rs::DynProofVerifier>>,
didcomm_source: Option<DidCommSource>,
shutdown: Option<CancellationToken>,
}
impl TrustRegistryBuilder {
pub fn repository(mut self, repository: Arc<dyn TrustRecordAdminRepository>) -> Self {
self.repository = Some(repository);
self
}
pub fn capabilities(mut self, capabilities: Vec<CapabilityDefinition>) -> Self {
self.capabilities = Some(capabilities);
self
}
pub fn capability_store(mut self, store: Box<dyn CapabilityStateStore>) -> Self {
self.capability_store = Some(store);
self
}
pub fn dedup_store(mut self, dedup: Arc<dyn MessageIdStore>) -> Self {
self.dedup = Some(dedup);
self
}
pub fn verifier(mut self, verifier: Arc<dyn trust_tasks_rs::DynProofVerifier>) -> Self {
self.verifier = Some(verifier);
self
}
pub fn didcomm_source(mut self, source: DidCommSource) -> Self {
self.didcomm_source = Some(source);
self
}
pub fn shutdown(mut self, shutdown: CancellationToken) -> Self {
self.shutdown = Some(shutdown);
self
}
pub async fn build(self) -> Result<TrustRegistry, BoxError> {
if ReplySigner::for_config(&self.config).is_none() {
warn!(
"the registry profile carries no signing key for {}: replies will go out \
unsigned, and a requester that checks proofs will not accept them",
self.config.didcomm_config.profile_config.did
);
}
let repository = self.repository.ok_or_else(|| {
BoxError::from(
"TrustRegistryBuilder needs a repository: call `.repository(...)`, or build one \
from the config with `storage::factory::TrustStorageRepoFactory`",
)
})?;
let available = self
.capabilities
.map(Ok)
.unwrap_or_else(|| default_capabilities(repository.clone()))?;
let store = self.capability_store.unwrap_or_else(|| {
Box::new(FileCapabilityStore::new(
self.config.server_config.capability_state_path(),
))
});
let base_repository = repository.clone();
let query_repository = repository.clone();
let capabilities = CapabilitySet::new(
available,
store,
Box::new(move || crate::trust_tasks::build_dispatcher(base_repository.clone())),
Box::new(move || crate::trust_tasks::build_query_dispatcher(query_repository.clone())),
)
.map_err(BoxError::from)?;
let verifier = match self.verifier {
Some(v) => v,
None => crate::trust_tasks::build_verifier().await,
};
let dedup = self.dedup.unwrap_or_else(|| {
if self.config.didcomm_config.is_enabled {
warn!(
"Message-id dedup is in-memory: duplicate writes are suppressed while this \
process lives, but a restart forgets them and a redelivery could re-apply. \
Pass a durable store with `TrustRegistryBuilder::dedup_store`."
);
}
Arc::new(MemoryMessageIdStore::default())
});
let audit = Arc::new(BoundedAuditLogger::new(Arc::new(BaseAuditLogger::new(
self.config.didcomm_config.admin_config.audit_config.clone(),
))));
let shutdown = self.shutdown.unwrap_or_default();
audit.spawn_flusher(shutdown.clone());
Ok(TrustRegistry {
audit,
shutdown,
health: Arc::new(RegistryHealth::new(self.config.didcomm_config.is_enabled)),
config: self.config,
repository,
capabilities,
verifier,
dedup,
didcomm_source: self.didcomm_source.unwrap_or_default(),
service_start_timestamp: Utc::now(),
})
}
}
fn default_capabilities(
repository: Arc<dyn TrustRecordAdminRepository>,
) -> Result<Vec<CapabilityDefinition>, BoxError> {
Ok(vec![
crate::capabilities::git_trust::definition(repository).map_err(BoxError::from)?,
])
}
#[cfg(test)]
mod tests {
use super::*;
use crate::capabilities::MemoryCapabilityStore;
use crate::storage::adapters::local_storage::LocalStorage;
async fn registry() -> TrustRegistry {
let repository: Arc<dyn TrustRecordAdminRepository> = Arc::new(LocalStorage::new());
TrustRegistry::builder(TrustRegistryConfig::embedded("/tmp/tr-embed-test"))
.repository(repository)
.capability_store(Box::new(MemoryCapabilityStore::default()))
.build()
.await
.expect("builds")
}
#[tokio::test]
async fn build_without_a_repository_explains_itself() {
let result = TrustRegistry::builder(TrustRegistryConfig::embedded("/tmp/tr-embed-test"))
.capability_store(Box::new(MemoryCapabilityStore::default()))
.build()
.await;
let err = match result {
Ok(_) => panic!("a repository is required"),
Err(e) => e,
};
assert!(
err.to_string().contains("repository"),
"unhelpful error: {err}"
);
}
#[tokio::test]
async fn capability_store_is_injectable() {
let _registry = registry().await;
assert!(!std::path::Path::new("/tmp/tr-embed-test/capabilities.json").exists());
}
#[tokio::test]
async fn didcomm_source_defaults_to_managed() {
let registry = registry().await;
assert!(matches!(
registry.into_parts().didcomm_source,
DidCommSource::Managed
));
}
#[tokio::test]
async fn host_driven_source_reaches_the_service_layer() {
let repository: Arc<dyn TrustRecordAdminRepository> = Arc::new(LocalStorage::new());
let registry = TrustRegistry::builder(TrustRegistryConfig::embedded("/tmp/tr-embed-test"))
.repository(repository)
.capability_store(Box::new(MemoryCapabilityStore::default()))
.didcomm_source(DidCommSource::HostDriven)
.build()
.await
.expect("builds");
assert!(matches!(
registry.into_parts().didcomm_source,
DidCommSource::HostDriven
));
}
#[tokio::test]
async fn route_didcomm_envelope_drops_an_unusable_body() {
let registry = registry().await;
let outcome = registry
.route_didcomm_envelope(
serde_json::json!({"not": "a trust task"}),
"did:example:peer",
)
.await;
assert!(
outcome.is_none(),
"an undecodable body has no thread or issuer to address an error to"
);
}
#[tokio::test]
async fn host_supplied_shutdown_token_is_used() {
let token = CancellationToken::new();
let repository: Arc<dyn TrustRecordAdminRepository> = Arc::new(LocalStorage::new());
let registry = TrustRegistry::builder(TrustRegistryConfig::embedded("/tmp/tr-embed-test"))
.repository(repository)
.capability_store(Box::new(MemoryCapabilityStore::default()))
.shutdown(token.clone())
.build()
.await
.expect("builds");
assert!(!registry.shutdown_token().is_cancelled());
token.cancel();
assert!(
registry.shutdown_token().is_cancelled(),
"the registry must observe the host's token, not one of its own"
);
}
#[tokio::test]
async fn mountable_router_excludes_health() {
use axum::body::Body;
use axum::http::{Request, StatusCode};
use tower::ServiceExt;
let registry = registry().await;
let response = registry
.router()
.oneshot(
Request::builder()
.uri("/health")
.body(Body::empty())
.expect("request"),
)
.await
.expect("router responds");
assert_eq!(response.status(), StatusCode::NOT_FOUND);
let response = registry
.health_router()
.oneshot(
Request::builder()
.uri("/health")
.body(Body::empty())
.expect("request"),
)
.await
.expect("health router responds");
assert_eq!(response.status(), StatusCode::OK);
}
#[tokio::test]
async fn router_mounts_under_a_host_prefix() {
use crate::domain::*;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use tower::ServiceExt;
let repository: Arc<dyn TrustRecordAdminRepository> = Arc::new(LocalStorage::new());
repository
.create(
TrustRecordBuilder::new()
.entity_id(EntityId::new("did:example:entity"))
.authority_id(AuthorityId::new("did:example:authority"))
.action(Action::new("issue"))
.resource(Resource::new("vc"))
.recognized(true)
.authorized(true)
.record_type(RecordType::Recognition)
.build()
.expect("valid record"),
)
.await
.expect("seeded");
let registry = TrustRegistry::builder(TrustRegistryConfig::embedded("/tmp/tr-embed-test"))
.repository(repository)
.capability_store(Box::new(MemoryCapabilityStore::default()))
.build()
.await
.expect("builds");
let app = axum::Router::new().nest("/registry", registry.router());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/registry/recognition")
.header("content-type", "application/json")
.body(Body::from(
serde_json::json!({
"entity_id": "did:example:entity",
"authority_id": "did:example:authority",
"action": "issue",
"resource": "vc",
})
.to_string(),
))
.expect("request"),
)
.await
.expect("nested router responds");
assert_eq!(
response.status(),
StatusCode::OK,
"recognition route should answer at the host's chosen prefix"
);
}
}