diff --git a/crates/tracedecay-code-index-runtime/src/git_transactions/owner.rs b/crates/tracedecay-code-index-runtime/src/git_transactions/owner.rs index 1d6c015981..0be39c784e 100644 --- a/crates/tracedecay-code-index-runtime/src/git_transactions/owner.rs +++ b/crates/tracedecay-code-index-runtime/src/git_transactions/owner.rs @@ -19,7 +19,7 @@ use tracedecay_domain::{ ProjectId, RepositoryIndexStateV1, RepositoryWorkingTreeStateV1, UtcMicros, canonical_sha256, }; use tracedecay_policy::{GitConflictRiskV1, GitEffectAuthorizationV1, GitEffectClassifierV1}; -use tracedecay_tool_catalog::CapabilityId; +use tracedecay_tool_catalog::{CapabilityId, CatalogSnapshotV1}; use super::{ CurrentGitIndexPolicyStateV1, DaemonGitIndexTransactionService, @@ -27,7 +27,7 @@ use super::{ GitIndexTransactionStoreRegistry, RepositoryMutationQueue, SharedDaemonGitIndexTransactionStore, canonicalize_repository_root, }; -use crate::ports::ApplicationCatalogProviderV1; +use crate::ports::ApplicationCatalogSnapshotErrorV1; use tracedecay_application::ProjectSourceAccessSnapshot; use tracedecay_application::configuration::ConfigurationControlStore; use tracedecay_global_db::RegisteredGlobalDbLeaseV1; @@ -37,6 +37,8 @@ const GIT_POLICY_REVISION: u64 = 2; type ProfiledStdRwLock = hotpath::rw_locks::RwLock; type ProfiledTokioMutex = hotpath::wrap::tokio::sync::Mutex; +type ApplicationCatalogComposer = + Arc Result + Send + Sync>; #[derive(Clone, Debug)] pub struct DaemonGitAuthorityStateV1 { @@ -71,7 +73,7 @@ pub trait DaemonGitAuthoritySource: Send + Sync { struct ProductionDaemonGitAuthoritySource { access: ProjectSourceAccessSnapshot, configuration: OwnedGlobalDbConfigurationControlStore, - catalog: ApplicationCatalogProviderV1, + catalog: ApplicationCatalogComposer, runtime: tokio::runtime::Handle, } @@ -158,10 +160,8 @@ impl DaemonGitAuthoritySource for ProductionDaemonGitAuthoritySource { if !effective_capabilities.contains(capability_id) { return Err(GitIndexTransactionPortError::PolicyDenied); } - let catalog = self - .catalog - .snapshot() - .map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?; + let catalog = + (self.catalog)().map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?; let manifest = catalog .capability(capability_id) .ok_or(GitIndexTransactionPortError::PolicyDenied)?; @@ -440,7 +440,7 @@ impl ServiceKey { /// worktrees share one session store actor without sharing native executors or /// mutation authority. pub struct DaemonGitIndexTransactionServiceRegistry { - catalog: ApplicationCatalogProviderV1, + catalog: ApplicationCatalogComposer, stores: GitIndexTransactionStoreRegistry, mutation_queue: Arc, services: ProfiledTokioMutex>, @@ -453,9 +453,14 @@ impl DaemonGitIndexTransactionServiceRegistry { /// Root supplies the catalog composer here: every owner this registry /// mounts resolves capability manifests through it, so there is no window /// in which an owner exists without one. - pub fn new(catalog: ApplicationCatalogProviderV1) -> Self { + pub fn new( + catalog: impl Fn() -> Result + + Send + + Sync + + 'static, + ) -> Self { Self { - catalog, + catalog: Arc::new(catalog), stores: GitIndexTransactionStoreRegistry::default(), mutation_queue: Arc::new(RepositoryMutationQueue::default()), services: hotpath::mutex!( diff --git a/crates/tracedecay-code-index-runtime/src/git_transactions/tests.rs b/crates/tracedecay-code-index-runtime/src/git_transactions/tests.rs index 69eeba3f37..5aa3d92b16 100644 --- a/crates/tracedecay-code-index-runtime/src/git_transactions/tests.rs +++ b/crates/tracedecay-code-index-runtime/src/git_transactions/tests.rs @@ -42,14 +42,13 @@ use tracedecay_global_db::tests::harness::RegisteredGlobalDbHarness; /// Registry owners take their catalog composer by construction, so these /// fixtures compose a real (contribution-free) snapshot rather than relying on /// whatever some other test installed first. -fn test_catalog_provider() -> crate::ports::ApplicationCatalogProviderV1 { - crate::ports::ApplicationCatalogProviderV1::new(|| { - tracedecay_tool_catalog::CatalogSnapshotBuilderV1::new() - .build() - .map_err(|error| { - crate::ports::ApplicationCatalogSnapshotErrorV1::new(error.to_string()) - }) - }) +fn test_catalog_snapshot() -> Result< + tracedecay_tool_catalog::CatalogSnapshotV1, + crate::ports::ApplicationCatalogSnapshotErrorV1, +> { + tracedecay_tool_catalog::CatalogSnapshotBuilderV1::new() + .build() + .map_err(|error| crate::ports::ApplicationCatalogSnapshotErrorV1::new(error.to_string())) } #[test] @@ -865,7 +864,7 @@ fn quarantine_clears_only_after_a_proven_recovery_receipt() { async fn daemon_owner_reuses_one_service_for_the_same_project_database() { let directory = tempfile::tempdir().expect("project directory"); let database = RegisteredGlobalDbHarness::open("git-index-owner-singleton").await; - let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_provider()); + let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_snapshot); let project_id = id::("project.singleton.fixture"); let first = registry @@ -894,7 +893,7 @@ async fn daemon_owner_reuses_one_service_for_the_same_project_database() { async fn daemon_owner_shutdown_fences_admission_clears_services_and_joins_store_actor() { let directory = tempfile::tempdir().expect("project directory"); let database = RegisteredGlobalDbHarness::open("git-index-owner-shutdown").await; - let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_provider()); + let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_snapshot); let project_id = id::("project.shutdown.fixture"); let service = registry .ensure( @@ -945,7 +944,7 @@ async fn daemon_owner_isolates_worktrees_sharing_a_project_database() { let directory = tempfile::tempdir().expect("project directory"); let alternate = tempfile::tempdir().expect("alternate worktree directory"); let database = RegisteredGlobalDbHarness::open("git-index-owner-worktrees").await; - let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_provider()); + let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_snapshot); let project_id = id::("project.singleton.fixture"); let primary = registry .ensure( @@ -1003,7 +1002,7 @@ async fn daemon_owner_resolves_symlink_alias_to_the_canonical_mounted_root() { let alias = alias_parent.path().join("worktree-alias"); symlink(directory.path(), &alias).expect("worktree root symlink alias"); let database = RegisteredGlobalDbHarness::open("git-index-owner-alias").await; - let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_provider()); + let registry = DaemonGitIndexTransactionServiceRegistry::new(test_catalog_snapshot); let project_id = id::("project.alias.fixture"); let first = registry .ensure( diff --git a/crates/tracedecay-code-index-runtime/src/lib.rs b/crates/tracedecay-code-index-runtime/src/lib.rs index ebe0795176..a0fc45395f 100644 --- a/crates/tracedecay-code-index-runtime/src/lib.rs +++ b/crates/tracedecay-code-index-runtime/src/lib.rs @@ -76,9 +76,8 @@ pub use code_graph_seat::{ pub use code_index_scheduler::CodeIndexSchedulerRegistryV1; pub use code_index_scheduler::identity::resolved_scope_for_project; pub use ports::{ - AdmissionParkLeaseV1, ApplicationCatalogProviderV1, ApplicationCatalogSnapshotErrorV1, - CONNECTION_ADMISSION, GitWatchMaintenanceWakeV1, GitWatchSyncConfigV1, - PreparedQueryActivationViewV1, park_admission, + AdmissionParkLeaseV1, ApplicationCatalogSnapshotErrorV1, CONNECTION_ADMISSION, + GitWatchMaintenanceWakeV1, GitWatchSyncConfigV1, PreparedQueryActivationViewV1, park_admission, }; pub use semantic_evaluation_shutdown::{ SemanticEvaluationShutdownJoinV1, SemanticEvaluationShutdownReceiptV1, diff --git a/crates/tracedecay-code-index-runtime/src/ports.rs b/crates/tracedecay-code-index-runtime/src/ports.rs index f69a1ef5fb..5fbb4aef40 100644 --- a/crates/tracedecay-code-index-runtime/src/ports.rs +++ b/crates/tracedecay-code-index-runtime/src/ports.rs @@ -8,7 +8,6 @@ use tokio::time::{Duration, timeout}; use tracedecay_contracts::ResolvedScope; use tracedecay_domain::configuration::ConfigurationRevisionId; use tracedecay_query::retrieval::QueryAuthorityV1; -use tracedecay_tool_catalog::CatalogSnapshotV1; /// Scheduler-facing view of a prepared query activation. pub struct PreparedQueryActivationViewV1 { @@ -85,34 +84,6 @@ impl Default for GitWatchMaintenanceWakeV1 { } } -/// Catalog snapshot provider the git-transaction owner consults for capability -/// manifests. Root hands one to each owner at construction; there is no ambient -/// registration, so an owner cannot exist without a composer. -#[derive(Clone)] -pub struct ApplicationCatalogProviderV1 { - compose: - Arc Result + Send + Sync>, -} - -impl ApplicationCatalogProviderV1 { - pub fn new( - compose: impl Fn() -> Result - + Send - + Sync - + 'static, - ) -> Self { - Self { - compose: Arc::new(compose), - } - } - - /// Composes the current snapshot. Root composes lazily per call, so this is - /// deliberately not a captured snapshot. - pub fn snapshot(&self) -> Result { - (self.compose)() - } -} - #[derive(Clone, Debug, PartialEq, Eq)] pub struct ApplicationCatalogSnapshotErrorV1 { pub message: String, diff --git a/crates/tracedecay/src/daemon/branch_admin.rs b/crates/tracedecay/src/daemon/branch_admin.rs index 45821edf2b..044f81c95e 100644 --- a/crates/tracedecay/src/daemon/branch_admin.rs +++ b/crates/tracedecay/src/daemon/branch_admin.rs @@ -545,9 +545,7 @@ impl Default for StoreAdministration { ), git_index_transaction_services: Arc::new( DaemonGitIndexTransactionServiceRegistry::new( - tracedecay_code_index_runtime::ApplicationCatalogProviderV1::new( - crate::runtime_ports::compose_application_catalog_snapshot, - ), + crate::runtime_ports::compose_application_catalog_snapshot, ), ), native_integration_services: Arc::new( diff --git a/crates/tracedecay/src/daemon/callable_code_authorization.rs b/crates/tracedecay/src/daemon/callable_code_authorization.rs index 9b59aa4b12..3aaae43ff2 100644 --- a/crates/tracedecay/src/daemon/callable_code_authorization.rs +++ b/crates/tracedecay/src/daemon/callable_code_authorization.rs @@ -1,6 +1,4 @@ -use std::future::Future; use std::path::PathBuf; -use std::pin::Pin; use std::sync::Arc; use tracedecay_contracts::{ @@ -21,64 +19,56 @@ use tracedecay_configuration::{ }; use tracedecay_graph_query::CodeGraphReadError; -type CurrentAccessFuture<'a> = Pin< - Box> + Send + 'a>, ->; - -trait CurrentCallableCodeAccessPort: Send + Sync { - fn current_access(&self, observed_at: UtcMicros) -> CurrentAccessFuture<'_>; -} - -struct ProductionCallableCodeAccessPort { - project_root: PathBuf, - scope: ResolvedScope, - configuration: Arc, -} - -impl CurrentCallableCodeAccessPort for ProductionCallableCodeAccessPort { - fn current_access(&self, observed_at: UtcMicros) -> CurrentAccessFuture<'_> { - Box::pin(async move { - let current = self - .configuration - .configuration_store() - .current() - .await - .map_err(configuration_current_problem)?; - let configuration = tracedecay_configuration::config::PinnedRuntimeConfiguration::new( - self.configuration.configuration_target().clone(), - current.revision_id, - current.snapshot, - ) - .map_err(|_| concealed())?; - crate::daemon::project_open_owners::daemon_owned_project_source_access_at( - &self.scope, - &self.project_root, - &configuration, - observed_at, - ) - .map_err(|_| concealed()) - }) - } -} +type CurrentCallableCodeAccess = + dyn Fn(UtcMicros) -> CurrentCallableCodeAccessFuture<'static> + Send + Sync; #[derive(Clone)] pub(super) struct DaemonCallableCodeAuthorizationSource { - access: Arc, + access: Arc, } impl DaemonCallableCodeAuthorizationSource { + fn new( + access: impl Fn(UtcMicros) -> CurrentCallableCodeAccessFuture<'static> + Send + Sync + 'static, + ) -> Self { + Self { + access: Arc::new(access), + } + } + pub(super) fn production( project_root: PathBuf, scope: ResolvedScope, configuration: Arc, ) -> Self { - Self { - access: Arc::new(ProductionCallableCodeAccessPort { - project_root, - scope, - configuration, - }), - } + let project_root = Arc::new(project_root); + let scope = Arc::new(scope); + Self::new(move |observed_at| { + let project_root = Arc::clone(&project_root); + let scope = Arc::clone(&scope); + let configuration = Arc::clone(&configuration); + Box::pin(async move { + let current = configuration + .configuration_store() + .current() + .await + .map_err(configuration_current_problem)?; + let configuration = + tracedecay_configuration::config::PinnedRuntimeConfiguration::new( + configuration.configuration_target().clone(), + current.revision_id, + current.snapshot, + ) + .map_err(|_| concealed())?; + crate::daemon::project_open_owners::daemon_owned_project_source_access_at( + &scope, + &project_root, + &configuration, + observed_at, + ) + .map_err(|_| concealed()) + }) + }) } #[hotpath::skip] @@ -86,7 +76,7 @@ impl DaemonCallableCodeAuthorizationSource { &self, observed_at: UtcMicros, ) -> Result { - self.access.current_access(observed_at).await + (self.access)(observed_at).await } pub(super) fn authorize( @@ -102,7 +92,7 @@ impl DaemonCallableCodeAuthorizationSource { impl CallableCodeAuthorizationSourcePort for DaemonCallableCodeAuthorizationSource { fn current(&self, observed_at: UtcMicros) -> CurrentCallableCodeAccessFuture<'_> { - Box::pin(async move { self.access.current_access(observed_at).await }) + (self.access)(observed_at) } fn authorize( @@ -442,20 +432,18 @@ mod tests { use super::*; - struct MutableAccess { - current: Mutex, - } - - impl CurrentCallableCodeAccessPort for MutableAccess { - fn current_access(&self, _observed_at: UtcMicros) -> CurrentAccessFuture<'_> { + fn mutable_source( + current: Arc>, + ) -> DaemonCallableCodeAuthorizationSource { + DaemonCallableCodeAuthorizationSource::new(move |_observed_at| { + let current = Arc::clone(¤t); Box::pin(async move { - Ok(self - .current + Ok(current .lock() .unwrap_or_else(|_| panic!("mutable access lock")) .clone()) }) - } + }) } fn access(operation: &ApplicationOperation) -> ProjectSourceAccessSnapshot { @@ -526,12 +514,8 @@ mod tests { .get(CallableCodeOperationKind::ExactOccurrence) .clone(); let mounted = access(&operation); - let mutable = Arc::new(MutableAccess { - current: Mutex::new(mounted.clone()), - }); - let source = DaemonCallableCodeAuthorizationSource { - access: mutable.clone(), - }; + let mutable = Arc::new(Mutex::new(mounted.clone())); + let source = mutable_source(Arc::clone(&mutable)); let authorization = source.authorize(mounted.clone()); let context = context(&mounted, &operation); let admission = authorization @@ -541,7 +525,6 @@ mod tests { { let mut current = mutable - .current .lock() .unwrap_or_else(|_| panic!("mutable access lock")); current.configuration_revision = @@ -577,11 +560,7 @@ mod tests { let observed_at = tracedecay_contracts::now_micros(); let mut mounted = access(&operation); mounted.grant_expires_at = UtcMicros(observed_at.0.saturating_add(60_000_000)); - let source = DaemonCallableCodeAuthorizationSource { - access: Arc::new(MutableAccess { - current: Mutex::new(mounted.clone()), - }), - }; + let source = mutable_source(Arc::new(Mutex::new(mounted.clone()))); let admission = DaemonCodeGraphReadAdmission::new(mounted.scope.clone(), source); let request_id = RequestId::new("request.graph-read-admission").expect("request id"); let cancellation = diff --git a/crates/tracedecay/src/daemon/project_open_owners/git_catalog_tests.rs b/crates/tracedecay/src/daemon/project_open_owners/git_catalog_tests.rs index ea5aecfba5..16561d3d07 100644 --- a/crates/tracedecay/src/daemon/project_open_owners/git_catalog_tests.rs +++ b/crates/tracedecay/src/daemon/project_open_owners/git_catalog_tests.rs @@ -1,7 +1,6 @@ use super::{GRANT_HORIZON, daemon_owned_project_source_access_at}; use crate::runtime_ports::compose_application_catalog_snapshot; use tracedecay_application::git_intelligence::NativeGitIntelligence; -use tracedecay_code_index_runtime::ApplicationCatalogProviderV1; use tracedecay_code_index_runtime::git_transactions::DaemonGitIndexTransactionServiceRegistry; use tracedecay_contracts::git::GitIndexTransactionPortError; use tracedecay_contracts::{ @@ -80,9 +79,8 @@ async fn git_owner_uses_explicit_canonical_catalog_and_rechecks_authorization() .mount_registered_project_sessions(project_id.clone()) .await .unwrap(); - let registry = DaemonGitIndexTransactionServiceRegistry::new( - ApplicationCatalogProviderV1::new(compose_application_catalog_snapshot), - ); + let registry = + DaemonGitIndexTransactionServiceRegistry::new(compose_application_catalog_snapshot); registry .ensure( database.clone(), @@ -119,9 +117,7 @@ async fn git_owner_uses_explicit_canonical_catalog_and_rechecks_authorization() // A separate owner can refuse its own provider without replacing the // canonical dependency already retained by the first owner. - let independent = DaemonGitIndexTransactionServiceRegistry::new( - ApplicationCatalogProviderV1::new(unavailable_catalog), - ); + let independent = DaemonGitIndexTransactionServiceRegistry::new(unavailable_catalog); independent .ensure( database.clone(),