Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 15 additions & 10 deletions crates/tracedecay-code-index-runtime/src/git_transactions/owner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,15 @@ 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,
DaemonProjectGitIndexPreviewAssembler, FixedDaemonGitIndexExecutor, GitIndexPolicyRecheckPort,
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;
Expand All @@ -37,6 +37,8 @@ const GIT_POLICY_REVISION: u64 = 2;

type ProfiledStdRwLock<T> = hotpath::rw_locks::RwLock<T>;
type ProfiledTokioMutex<T> = hotpath::wrap::tokio::sync::Mutex<T>;
type ApplicationCatalogComposer =
Arc<dyn Fn() -> Result<CatalogSnapshotV1, ApplicationCatalogSnapshotErrorV1> + Send + Sync>;

#[derive(Clone, Debug)]
pub struct DaemonGitAuthorityStateV1 {
Expand Down Expand Up @@ -71,7 +73,7 @@ pub trait DaemonGitAuthoritySource: Send + Sync {
struct ProductionDaemonGitAuthoritySource {
access: ProjectSourceAccessSnapshot,
configuration: OwnedGlobalDbConfigurationControlStore,
catalog: ApplicationCatalogProviderV1,
catalog: ApplicationCatalogComposer,
runtime: tokio::runtime::Handle,
}

Expand Down Expand Up @@ -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)?;
Expand Down Expand Up @@ -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<RepositoryMutationQueue>,
services: ProfiledTokioMutex<HashMap<ServiceKey, ServiceEntry>>,
Expand All @@ -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<CatalogSnapshotV1, ApplicationCatalogSnapshotErrorV1>
+ Send
+ Sync
+ 'static,
) -> Self {
Self {
catalog,
catalog: Arc::new(catalog),
stores: GitIndexTransactionStoreRegistry::default(),
mutation_queue: Arc::new(RepositoryMutationQueue::default()),
services: hotpath::mutex!(
Expand Down
23 changes: 11 additions & 12 deletions crates/tracedecay-code-index-runtime/src/git_transactions/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down Expand Up @@ -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::<ProjectId>("project.singleton.fixture");

let first = registry
Expand Down Expand Up @@ -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::<ProjectId>("project.shutdown.fixture");
let service = registry
.ensure(
Expand Down Expand Up @@ -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::<ProjectId>("project.singleton.fixture");
let primary = registry
.ensure(
Expand Down Expand Up @@ -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::<ProjectId>("project.alias.fixture");
let first = registry
.ensure(
Expand Down
5 changes: 2 additions & 3 deletions crates/tracedecay-code-index-runtime/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
29 changes: 0 additions & 29 deletions crates/tracedecay-code-index-runtime/src/ports.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<dyn Fn() -> Result<CatalogSnapshotV1, ApplicationCatalogSnapshotErrorV1> + Send + Sync>,
}

impl ApplicationCatalogProviderV1 {
pub fn new(
compose: impl Fn() -> Result<CatalogSnapshotV1, ApplicationCatalogSnapshotErrorV1>
+ 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<CatalogSnapshotV1, ApplicationCatalogSnapshotErrorV1> {
(self.compose)()
}
}

#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ApplicationCatalogSnapshotErrorV1 {
pub message: String,
Expand Down
4 changes: 1 addition & 3 deletions crates/tracedecay/src/daemon/branch_admin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
123 changes: 51 additions & 72 deletions crates/tracedecay/src/daemon/callable_code_authorization.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,4 @@
use std::future::Future;
use std::path::PathBuf;
use std::pin::Pin;
use std::sync::Arc;

use tracedecay_contracts::{
Expand All @@ -21,72 +19,64 @@ use tracedecay_configuration::{
};
use tracedecay_graph_query::CodeGraphReadError;

type CurrentAccessFuture<'a> = Pin<
Box<dyn Future<Output = Result<ProjectSourceAccessSnapshot, ApplicationProblem>> + Send + 'a>,
>;

trait CurrentCallableCodeAccessPort: Send + Sync {
fn current_access(&self, observed_at: UtcMicros) -> CurrentAccessFuture<'_>;
}

struct ProductionCallableCodeAccessPort {
project_root: PathBuf,
scope: ResolvedScope,
configuration: Arc<ProjectConfigurationRuntime>,
}

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<dyn CurrentCallableCodeAccessPort>,
access: Arc<CurrentCallableCodeAccess>,
}

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<ProjectConfigurationRuntime>,
) -> 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]
pub(super) async fn current(
&self,
observed_at: UtcMicros,
) -> Result<ProjectSourceAccessSnapshot, ApplicationProblem> {
self.access.current_access(observed_at).await
(self.access)(observed_at).await
}

pub(super) fn authorize(
Expand All @@ -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(
Expand Down Expand Up @@ -442,20 +432,18 @@ mod tests {

use super::*;

struct MutableAccess {
current: Mutex<ProjectSourceAccessSnapshot>,
}

impl CurrentCallableCodeAccessPort for MutableAccess {
fn current_access(&self, _observed_at: UtcMicros) -> CurrentAccessFuture<'_> {
fn mutable_source(
current: Arc<Mutex<ProjectSourceAccessSnapshot>>,
) -> DaemonCallableCodeAuthorizationSource {
DaemonCallableCodeAuthorizationSource::new(move |_observed_at| {
let current = Arc::clone(&current);
Box::pin(async move {
Ok(self
.current
Ok(current
.lock()
.unwrap_or_else(|_| panic!("mutable access lock"))
.clone())
})
}
})
}

fn access(operation: &ApplicationOperation) -> ProjectSourceAccessSnapshot {
Expand Down Expand Up @@ -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
Expand All @@ -541,7 +525,6 @@ mod tests {

{
let mut current = mutable
.current
.lock()
.unwrap_or_else(|_| panic!("mutable access lock"));
current.configuration_revision =
Expand Down Expand Up @@ -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 =
Expand Down
Loading
Loading