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
76 changes: 44 additions & 32 deletions crates/tracedecay-daemon-service/src/callable_code_authorization.rs
Original file line number Diff line number Diff line change
@@ -1,18 +1,21 @@
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;

use tracedecay_contracts::{
ApplicationContractError, ApplicationOperation, ApplicationProblem, ApplicationProblemKind,
AuthorityReceipt, CallableCodeAuthorizationAdmission, CallableCodeAuthorizationFuture,
CallableCodeAuthorizationPort, RequestAdmission, RequestContext, ResolvedScope, RetryDirective,
SafeDiagnostic,
};
use tracedecay_domain::{ComponentVersion, UtcMicros};
use tracedecay_domain::{ActorId, ComponentVersion, UtcMicros};

use crate::callable_code_request_context;
use crate::project_owner_registration::project_owner_capabilities;
use tracedecay_application::project_open_authorization::ProjectOpenSourceAccessAuthorityV1;
use tracedecay_application::{
CallableCodeAuthorizationSourcePort, CurrentCallableCodeAccessFuture,
ProjectSourceAccessSnapshot,
ProjectSourceAccessSnapshot, ProjectSourceAccessSnapshotPort,
};
use tracedecay_configuration::config::PinnedRuntimeConfiguration;
use tracedecay_configuration::{
Expand All @@ -23,6 +26,37 @@ use tracedecay_graph_query::CodeGraphReadError;
type CurrentCallableCodeAccess =
dyn Fn(UtcMicros) -> CurrentCallableCodeAccessFuture<'static> + Send + Sync;

const DAEMON_REQUESTER: &str = "actor.tracedecay-daemon.project-open";
pub const GRANT_HORIZON: Duration = Duration::from_hours(24);

pub fn project_open_source_access_authority()
-> Result<ProjectOpenSourceAccessAuthorityV1, ApplicationContractError> {
let requester = ActorId::new(DAEMON_REQUESTER.to_owned()).map_err(|_| {
ApplicationContractError::Inconsistent {
field: "project-open requester",
}
})?;
Ok(ProjectOpenSourceAccessAuthorityV1::new(
requester,
project_owner_capabilities()?,
GRANT_HORIZON,
))
}

pub fn daemon_owned_project_source_access_at(
scope: &ResolvedScope,
project_root: &Path,
configuration: &PinnedRuntimeConfiguration,
observed_at: UtcMicros,
) -> Result<ProjectSourceAccessSnapshot, ApplicationContractError> {
project_open_source_access_authority()?.source_access_at(
scope,
project_root,
configuration,
observed_at,
)
}

#[derive(Clone)]
pub struct DaemonCallableCodeAuthorizationSource {
access: Arc<CurrentCallableCodeAccess>,
Expand All @@ -41,25 +75,13 @@ impl DaemonCallableCodeAuthorizationSource {
project_root: PathBuf,
scope: ResolvedScope,
configuration: Arc<ProjectConfigurationRuntime>,
source_access_at: impl Fn(
&ResolvedScope,
&Path,
&PinnedRuntimeConfiguration,
UtcMicros,
)
-> Result<ProjectSourceAccessSnapshot, ApplicationContractError>
+ Send
+ Sync
+ 'static,
) -> Self {
let project_root = Arc::new(project_root);
let scope = Arc::new(scope);
let source_access_at = Arc::new(source_access_at);
Self::new(move |observed_at| {
let project_root = Arc::clone(&project_root);
let scope = Arc::clone(&scope);
let configuration = Arc::clone(&configuration);
let source_access_at = Arc::clone(&source_access_at);
Box::pin(async move {
let current = configuration
.configuration_store()
Expand All @@ -72,8 +94,13 @@ impl DaemonCallableCodeAuthorizationSource {
current.snapshot,
)
.map_err(|_| concealed())?;
source_access_at(&scope, &project_root, &configuration, observed_at)
.map_err(|_| concealed())
daemon_owned_project_source_access_at(
&scope,
&project_root,
&configuration,
observed_at,
)
.map_err(|_| concealed())
})
})
}
Expand Down Expand Up @@ -124,25 +151,10 @@ impl DaemonCodeGraphReadAdmission {
project_root: PathBuf,
scope: ResolvedScope,
configuration: Arc<ProjectConfigurationRuntime>,
source_access_at: impl Fn(
&ResolvedScope,
&Path,
&PinnedRuntimeConfiguration,
UtcMicros,
)
-> Result<ProjectSourceAccessSnapshot, ApplicationContractError>
+ Send
+ Sync
+ 'static,
) -> Self {
Self::new(
scope.clone(),
DaemonCallableCodeAuthorizationSource::production(
project_root,
scope,
configuration,
source_access_at,
),
DaemonCallableCodeAuthorizationSource::production(project_root, scope, configuration),
)
}

Expand Down
5 changes: 3 additions & 2 deletions crates/tracedecay-daemon-service/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,14 +67,15 @@ pub mod project_runtime;
pub mod query_authority_provider;
pub mod query_mcp_admission;
pub mod remote_http_transport;
pub mod remote_protocol;
mod remote_protocol;
pub mod request_cancellation;
mod shutdown_coordination;

mod multi_root;

pub use callable_code_authorization::{
DaemonCallableCodeAuthorizationSource, DaemonCodeGraphReadAdmission,
DaemonCallableCodeAuthorizationSource, DaemonCodeGraphReadAdmission, GRANT_HORIZON,
daemon_owned_project_source_access_at, project_open_source_access_authority,
};
pub use invocation::semantic_evaluation::SemanticInvocationControlV1;
#[cfg(any(test, feature = "test-helpers"))]
Expand Down
114 changes: 112 additions & 2 deletions crates/tracedecay-daemon-service/src/remote_protocol.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,6 @@ use tracedecay_store_runtime::{

mod observability;

pub use observability::remote_query_result_observation;

use tracedecay_store_runtime::{DaemonRemoteCredentialAuthorityV1, DaemonRemoteCredentialLookupV1};

struct DaemonRemoteEnrollmentProtocolPortV1 {
Expand Down Expand Up @@ -459,3 +457,115 @@ mod recovery_control_tests {
);
}
}

#[cfg(test)]
mod observation_tests {
use tracedecay_contracts::remote::composition::{
AuthenticityClaimV1, AuthorizationClaimV1, IntegrityClaimV1, PendingLocalEvidenceV1,
PendingLocalObservationsV1, QueryManifestBindingV1, RemoteCompletenessV1,
RemoteFreshnessV1, RemoteQueryCompositionV1, ShardCoverageStateV1,
ShardQueryContributionV1,
};
use tracedecay_contracts::remote::query::{
RemoteExactObservationResultV1, RemoteQueryResultV1,
};
use tracedecay_domain::{CoverageStateV1, ObservedTernaryV1};

use super::observability::remote_query_result_observation;

fn remote_query_result(
coverage: ShardCoverageStateV1,
pending_local: PendingLocalEvidenceV1,
) -> RemoteQueryResultV1 {
RemoteQueryResultV1 {
composition: RemoteQueryCompositionV1 {
contributions: vec![ShardQueryContributionV1 {
manifest: QueryManifestBindingV1 {
brain_id: "brain.remote-coverage".to_owned(),
shard_id: "shard.remote-coverage".to_owned(),
generation_id: "generation.remote-coverage".to_owned(),
schema_digest: [1; 32],
watermark_sequence: 1,
placement_revision: 1,
authority_epoch: 1,
cache_age_millis: 0,
cache_lag_commits: 0,
},
integrity: IntegrityClaimV1::Verified,
authenticity: AuthenticityClaimV1::Authenticated,
freshness: RemoteFreshnessV1::Current,
completeness: RemoteCompletenessV1::Complete,
authorization: AuthorizationClaimV1::Authorized,
coverage,
authority_receipt: None,
value: None,
reason_code: (coverage != ShardCoverageStateV1::Complete)
.then(|| "remote_shard_degraded".to_owned()),
}],
pending_local,
coverage,
},
observation: RemoteExactObservationResultV1::NotFound,
}
}

#[test]
fn remote_query_coverage_preserves_real_shard_and_pending_counts() {
let result = remote_query_result(
ShardCoverageStateV1::Stale,
PendingLocalObservationsV1 {
count: 3,
oldest_age_millis: Some(9),
has_sequence_gap: false,
has_quarantined: false,
}
.into(),
);
result.validate().expect("valid stale remote query result");

let observation = remote_query_result_observation(
"request.remote-coverage",
1,
&result,
ObservedTernaryV1::Yes,
);

assert_eq!(observation.expected_shards, Some(1));
assert_eq!(observation.observed_shards, Some(1));
assert_eq!(observation.pending_local_evidence, Some(3));
assert_eq!(observation.terminal_succeeded, ObservedTernaryV1::Yes);
assert_eq!(observation.coverage, CoverageStateV1::Stale);
assert_eq!(
observation.unavailable_reason.as_deref(),
Some("pending_local_evidence")
);
}

#[test]
fn remote_query_coverage_does_not_fabricate_unavailable_pending_count() {
let result = remote_query_result(
ShardCoverageStateV1::Unknown,
PendingLocalEvidenceV1::Unavailable {
reason: tracedecay_contracts::remote::composition::PendingLocalUnavailableReasonV1::AuthorityUnavailable,
},
);
result
.validate()
.expect("valid unavailable remote query result");

let observation = remote_query_result_observation(
"request.remote-coverage-unavailable",
1,
&result,
ObservedTernaryV1::Unknown,
);

assert_eq!(observation.pending_local_evidence, None);
assert_eq!(observation.terminal_succeeded, ObservedTernaryV1::Unknown);
assert_eq!(observation.coverage, CoverageStateV1::Unknown);
assert_eq!(
observation.unavailable_reason.as_deref(),
Some("pending_local_authority_unavailable")
);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ pub(super) fn record_remote_query_response(
);
}

pub fn remote_query_result_observation(
pub(super) fn remote_query_result_observation(
operation_ref: &str,
expected_shards: usize,
result: &RemoteQueryResultV1,
Expand Down
4 changes: 3 additions & 1 deletion crates/tracedecay-graph-query/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,9 +45,11 @@ pub use queries::{
pub use source_authority::{
CodeGraphSourceAuthorityPort, CodeGraphSourceBindFuture, CodeGraphSourceBindRequest,
};
#[cfg(any(test, feature = "test-helpers"))]
pub use verified_query::admitted_verified_graph_query_port;
pub use verified_query::{
AdmittedVerifiedGraphQueryPort, VerifiedGraphQuery, VerifiedGraphQueryFuture,
VerifiedGraphQueryPort, VerifiedGraphQueryRequest, admitted_verified_graph_query_port,
VerifiedGraphQueryPort, VerifiedGraphQueryRequest,
admitted_verified_graph_query_port_with_source, open_verified_graph_query,
};

Expand Down
1 change: 1 addition & 0 deletions crates/tracedecay-graph-query/src/verified_query.rs
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ impl VerifiedGraphQueryPort for AdmittedVerifiedGraphQueryPort {
}
}

#[cfg(any(test, feature = "test-helpers"))]
#[must_use]
pub fn admitted_verified_graph_query_port(
admission: Arc<dyn CodeGraphReadAdmissionPort>,
Expand Down
3 changes: 0 additions & 3 deletions crates/tracedecay/src/daemon/invocation_tests/types_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -995,7 +995,6 @@ async fn feedback_admission_conflicts_construct_zero_losing_producers() {
project.path().to_path_buf(),
scope.clone(),
Arc::clone(graph.configuration_runtime()),
crate::daemon::project_open_owners::daemon_owned_project_source_access_at,
)),
),
)
Expand Down Expand Up @@ -1042,7 +1041,6 @@ async fn feedback_admission_conflicts_construct_zero_losing_producers() {
project.path().to_path_buf(),
scope.clone(),
Arc::clone(graph.configuration_runtime()),
crate::daemon::project_open_owners::daemon_owned_project_source_access_at,
)),
)
.await;
Expand Down Expand Up @@ -1130,7 +1128,6 @@ async fn feedback_admission_conflicts_construct_zero_losing_producers() {
publisher_root,
publisher_scope,
publisher_configuration,
crate::daemon::project_open_owners::daemon_owned_project_source_access_at,
)),
)
.await
Expand Down
7 changes: 4 additions & 3 deletions crates/tracedecay/src/daemon/project_composition.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,9 @@
use super::*;
use tracedecay_code_index_runtime::code_index_scheduler;
use tracedecay_daemon_identity::profile_identity;
use tracedecay_daemon_service::DaemonSemanticRuntimeRegistrationError;
use tracedecay_daemon_service::{
DaemonSemanticRuntimeRegistrationError, daemon_owned_project_source_access_at,
};
use tracedecay_runtime_core::logging::log_daemon_event;
use tracedecay_semantic_contracts::SemanticResourceCeilings;
use tracedecay_session_runtime::session_sync::DaemonSessionSyncConfig;
Expand Down Expand Up @@ -1194,7 +1196,7 @@ impl ProjectOpenInputs<'_> {
user_session_db.clone(),
])
.await;
let delivery_access = project_open_owners::daemon_owned_project_source_access_at(
let delivery_access = daemon_owned_project_source_access_at(
&code_index.scope,
self.canonical_project_path,
runtime_configuration,
Expand Down Expand Up @@ -1914,7 +1916,6 @@ fn project_code_index_authorities(
canonical_project_path.to_path_buf(),
scope.clone(),
Arc::clone(cg.configuration_runtime()),
crate::daemon::project_open_owners::daemon_owned_project_source_access_at,
),
);
let search_admission = tracedecay_daemon_service::admit_query_mcp_read(
Expand Down
Loading
Loading