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
3 changes: 2 additions & 1 deletion crates/tracedecay-daemon-service/src/invocation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -233,7 +233,8 @@ pub use work_routing::DaemonWorkProposalRoutingAuthorityV1;
pub use feedback::{
DaemonAdvisoryCycleInvocationFuture, DaemonAdvisoryCycleInvocationOwner,
DaemonAdvisoryCycleInvocationPort, DaemonAdvisoryCycleInvocationRequest,
DaemonAdvisoryCycleMountFuture, DaemonAdvisoryCycleMountV1, DaemonFeedbackInvocationOwner,
DaemonAdvisoryCycleMountChangedFuture, DaemonAdvisoryCycleMountFuture,
DaemonAdvisoryCycleMountV1, DaemonFeedbackInvocationOwner,
DaemonFeedbackProximityInvocationFuture, DaemonFeedbackProximityInvocationRequest,
advisory_cycle_invocation_result, daemon_operation_event_authority,
feedback_proximity_invocation_result,
Expand Down
7 changes: 6 additions & 1 deletion crates/tracedecay-daemon-service/src/invocation/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -513,7 +513,12 @@ impl DaemonInvocationService {
cancellation,
} => {
let advisory_cycle = match advisory_cycle {
Some(owner) if owner.service.mount().await == DaemonAdvisoryCycleMountV1::Answers => {
Some(owner)
if matches!(
owner.service.mount().await,
DaemonAdvisoryCycleMountV1::Answers
) =>
{
Some(owner)
}
_ => {
Expand Down
46 changes: 31 additions & 15 deletions crates/tracedecay-daemon-service/src/invocation/feedback.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,13 +109,16 @@ pub type DaemonFeedbackProximityInvocationFuture<'a> = Pin<
>;

/// Whether an advisory-cycle owner answers for its project now, or stands in
/// for a full cycle a ready sealed generation is mounting.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
/// for a full cycle its mount is still publishing.
pub enum DaemonAdvisoryCycleMountV1 {
Answers,
Mounting,
/// Resolves when the owner's own mount state next moves, so a waiting
/// request re-reads it instead of holding out for its deadline.
Mounting(DaemonAdvisoryCycleMountChangedFuture),
}

pub type DaemonAdvisoryCycleMountChangedFuture = Pin<Box<dyn Future<Output = ()> + Send>>;

pub type DaemonAdvisoryCycleMountFuture<'a> =
Pin<Box<dyn Future<Output = DaemonAdvisoryCycleMountV1> + Send + 'a>>;

Expand Down Expand Up @@ -960,11 +963,12 @@ impl DaemonInvocationService {
/// The owner that answers an advisory-cycle request for `project_root`.
///
/// A reopened project serves its first requests while project open is
/// still publishing owners, and then while a ready sealed generation
/// mounts the full cycle behind a placeholder. Those requests wait for the
/// publication within their own deadline instead of failing a call a
/// retry would answer. A finished publication without an owner, or an
/// owner that answers, returns at once.
/// still publishing owners, and then while its mount publishes the full
/// cycle behind a placeholder. Those requests wait for the publication
/// within their own deadline instead of failing a call a retry would
/// answer, and re-read the placeholder whenever its mount state moves. A
/// finished publication without an owner, or an owner that answers,
/// returns at once.
#[hotpath::measure(label = "daemon.service.feedback.advisory_owner_wait", future = true)]
pub(super) async fn answering_advisory_cycle_owner(
&self,
Expand All @@ -975,18 +979,30 @@ impl DaemonInvocationService {
loop {
let (owner, publication, mut changed) =
self.project_runtimes.advisory_cycle_view(project_root);
let waits = match &owner {
Some(owner) => owner.service.mount().await == DaemonAdvisoryCycleMountV1::Mounting,
None => publication == Some(ProjectRuntimePublicationStateV1::Warming),
let mount_changed: DaemonAdvisoryCycleMountChangedFuture = match &owner {
Some(current) => match current.service.mount().await {
DaemonAdvisoryCycleMountV1::Answers => return owner,
DaemonAdvisoryCycleMountV1::Mounting(mount_changed) => mount_changed,
},
None if publication == Some(ProjectRuntimePublicationStateV1::Warming) => {
Box::pin(std::future::pending())
}
None => return owner,
};
let remaining_micros = deadline.expires_at.0.saturating_sub(now_micros().0);
if !waits || remaining_micros <= 0 {
if remaining_micros <= 0 {
return owner;
}
let remaining = Duration::from_micros(remaining_micros.unsigned_abs());
match tokio::time::timeout(remaining, changed.changed()).await {
Ok(Ok(())) => {}
Ok(Err(_)) | Err(_) => return owner,
let woke = tokio::time::timeout(remaining, async {
tokio::select! {
published = changed.changed() => published.is_ok(),
() = mount_changed => true,
}
})
.await;
if !matches!(woke, Ok(true)) {
return owner;
}
}
}
Expand Down
3 changes: 2 additions & 1 deletion crates/tracedecay-daemon-service/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,8 @@ pub use invocation::{
BoundedHookOrchestratorV1, ConfigurationRuntimeRefreshFuture, ConfigurationRuntimeRefreshPort,
DaemonAdvisoryCycleInvocationFuture, DaemonAdvisoryCycleInvocationOwner,
DaemonAdvisoryCycleInvocationPort, DaemonAdvisoryCycleInvocationRequest,
DaemonAdvisoryCycleMountFuture, DaemonAdvisoryCycleMountV1, DaemonAdvisoryRuntimeRegistrar,
DaemonAdvisoryCycleMountChangedFuture, DaemonAdvisoryCycleMountFuture,
DaemonAdvisoryCycleMountV1, DaemonAdvisoryRuntimeRegistrar,
DaemonAdvisoryRuntimeRegistrationError, DaemonConfigurationGrantAuthority,
DaemonConfigurationRuntimeRegistrar, DaemonContextScoutRuntimeRegistrar,
DaemonContextScoutRuntimeRegistrationError, DaemonFeedbackInvocationOwner,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -156,6 +156,108 @@ async fn first_advisory_cycle_after_a_reopen_answers_without_a_retry() {
harness.shutdown().await;
}

/// A reopened checkout whose deferred advisory mount fails terminally names
/// that failure as soon as the mount gives up, instead of holding the request
/// to its deadline and then answering the retryable pre-mount state.
#[tokio::test(flavor = "multi_thread")]
#[hotpath::skip]
async fn reopened_checkout_whose_advisory_mount_failed_names_the_failure() {
let isolation = tempfile::TempDir::new().expect("production harness isolation");
let project = isolation.path().join("project");
std::fs::create_dir_all(project.join("src")).expect("project source dir");
std::fs::write(
project.join("src/lib.rs"),
"pub fn failed_mount(value: u32) -> u32 {\n value * 2\n}\n",
)
.expect("rust source");
git(&project, &["init", "--quiet", "-b", "feature/failed-mount"]);
git(&project, &["add", "."]);
git(&project, &["commit", "--quiet", "-m", "seed failed mount"]);

let harness = ProductionProjectCompositionHarnessV1::open(isolation.path(), [project.clone()])
.await
.expect("first production composition");
let (refused, settled) = settled_advisory_cycle(&harness, &project, "src/lib.rs").await;
assert!(!refused, "the first open must mount the cycle: {settled}");
let project_id = tracedecay_domain::ProjectId::new(
harness
.project_id(&project)
.await
.expect("registered project id"),
)
.expect("project id");
harness.shutdown().await;

// A foreign live hook-notice queue for this checkout refuses the reopened
// advisory owner its registration on every attempt.
let canonical_project = std::fs::canonicalize(&project).expect("canonical project");
let scope =
tracedecay_code_index_runtime::resolved_scope_for_project(&canonical_project, &project_id)
.expect("resolved scope");
let (hook_project, hook_worktree) = tracedecay_agent_hosts::hooks::hook_scope_locators(&scope);
let foreign = tracedecay_application::advisory::AdvisoryHookNoticeQueueV1::new(
tracedecay_application::feedback::resolve_project_feedback_scope_v1(
&canonical_project,
&scope,
)
.expect("feedback scope"),
);
assert!(
tracedecay_application::advisory::register_advisory_hook_notice_queue(
hook_project,
hook_worktree,
&foreign,
)
);

let harness = ProductionProjectCompositionHarnessV1::open(isolation.path(), [project.clone()])
.await
.expect("reopened production composition");
let document_uri = url::Url::from_file_path(project.join("src/lib.rs"))
.expect("document uri")
.to_string();
let request_deadline = tracedecay_mcp::tools::binding::canonical_tool_dispatch_ceiling(
"tracedecay_feedback_advisory_cycle",
)
.expect("advisory cycle dispatch deadline");
let started = Instant::now();
let response = harness
.call_tool(
&project,
"tracedecay_feedback_advisory_cycle",
json!({"document_uri": document_uri, "format": "json"}),
)
.await
.expect("advisory cycle call");
let elapsed = started.elapsed();
let (refused, answer) = tool_answer(&response);
assert!(refused, "a failed mount cannot run the cycle: {answer}");
let problem = &answer["problem"];
assert_eq!(
problem["code"],
json!("feedback.advisory-cycle.mount-failed"),
"{answer}"
);
assert_eq!(problem["kind"], json!("unavailable"), "{answer}");
assert_eq!(
problem["legal_actions"],
json!(["contact_administrator"]),
"{answer}"
);
assert!(
elapsed < request_deadline / 2,
"the failure was held for {elapsed:?} of a {request_deadline:?} deadline: {answer}"
);
harness.shutdown().await;
assert!(
tracedecay_application::advisory::unregister_advisory_hook_notice_queue(
hook_project,
hook_worktree,
&foreign,
)
);
}

#[tokio::test(flavor = "multi_thread")]
#[hotpath::skip]
async fn checkout_without_indexable_source_names_why_the_advisory_cycle_cannot_run() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ use std::sync::Arc;
use std::time::{Duration, Instant};

use sha2::{Digest, Sha256};
use tokio::sync::watch;
use tracedecay_application::advisory::github_runtime::{
ConfiguredGitHubSourceAccessAuthorityV1, GitHubDiscoveryControlV1,
GitHubExactCommitDiscoveryOutcomeV1, GitHubProviderLifecycleV1, GitHubSourceAccessAuthorityV1,
Expand Down Expand Up @@ -115,6 +116,7 @@ mod deferred;
mod model;

pub(super) use deferred::wait_for_generation_change;
use deferred::{AdvisoryMountPublisherV1, AdvisoryMountStateV1};
pub(crate) use model::ProjectOpenDependentOwnerState;
use model::advisory_monotonic_deadline;
#[cfg(test)]
Expand Down Expand Up @@ -1358,12 +1360,14 @@ pub(in crate::daemon) async fn register_project_open_dependent_owners(
return Ok(());
}
register_project_delivery_read_authority(invocation, project_root, &state).await?;
let (mount_publisher, mount_state) = AdvisoryMountPublisherV1::channel();
// Proximity is a typed read over session/git correlation and the current
// code graph. Like Delivery, it must be available before a sealed
// generation mounts the full advisory cycle, otherwise HTTP
// `/api/feedback/proximity` answers `feedback.proximity.unavailable`
// while Work/application are already serving.
register_project_proximity_read_authority(invocation, project_root, &state).await?;
register_project_proximity_read_authority(invocation, project_root, &state, mount_state)
.await?;
let indexed_generation =
selected_feedback_generation(invocation, project_root, &state.scope).await;
if let (Some(lsp_session_factory), Some(indexed_generation)) =
Expand Down Expand Up @@ -1398,9 +1402,11 @@ pub(in crate::daemon) async fn register_project_open_dependent_owners(
invocation.clone(),
project_root.to_path_buf(),
state,
mount_publisher,
);
return Ok(());
}
mount_publisher.publish(AdvisoryMountStateV1::Mounted);
tracing::info!(
event = "project_open_owner_phase",
project = %project_root.display(),
Expand Down Expand Up @@ -1450,6 +1456,7 @@ pub(in crate::daemon) async fn register_project_open_dependent_owners(
invocation.clone(),
project_root.to_path_buf(),
state,
mount_publisher,
);
Ok(())
}
Expand Down Expand Up @@ -1984,6 +1991,7 @@ async fn register_project_proximity_read_authority(
invocation: &DaemonInvocationState,
project_root: &Path,
state: &ProjectOpenDependentOwnerState,
mount_state: watch::Receiver<AdvisoryMountStateV1>,
) -> Result<()> {
let feedback_scope = match resolve_project_feedback_scope_v1(project_root, &state.scope) {
Ok(scope) => scope,
Expand Down Expand Up @@ -2025,6 +2033,7 @@ async fn register_project_proximity_read_authority(
feedback_scope,
proximity_read,
code_index_schedulers: invocation.code_index_schedulers.clone(),
mount_state,
}) as Arc<dyn DaemonAdvisoryCycleInvocationPort>,
);
invocation
Expand Down Expand Up @@ -2053,6 +2062,7 @@ struct ProjectOpenProximityReadOwnerV1 {
feedback_scope: FeedbackScopeV1,
proximity_read: FeedbackProximityReadRuntimeV1,
code_index_schedulers: CodeIndexSchedulerRegistryV1,
mount_state: watch::Receiver<AdvisoryMountStateV1>,
}

impl DaemonAdvisoryCycleInvocationPort for ProjectOpenProximityReadOwnerV1 {
Expand All @@ -2078,6 +2088,9 @@ impl DaemonAdvisoryCycleInvocationPort for ProjectOpenProximityReadOwnerV1 {
LegalAction::CorrectRequest,
));
}
if let AdvisoryMountStateV1::Failed(failure) = *self.mount_state.borrow() {
return Err(failure.problem());
}
if self
.code_index_schedulers
.reconciled_without_generation_for_scope(&self.scope)
Expand All @@ -2096,35 +2109,27 @@ impl DaemonAdvisoryCycleInvocationPort for ProjectOpenProximityReadOwnerV1 {
})
}

/// With a ready sealed generation the deferred mount is already upgrading
/// this owner, so a request waits for that publication instead of taking
/// the retryable warming answer. A restart serves its retained generation
/// before its first pass proves the checkout still matches it; that
/// proof, or the successor it finds owed, is what upgrades this owner, so
/// a retained owner waits as well.
// ponytail: a deferred mount that fails terminally leaves this owner in
// place, so such a request waits out its own deadline before the warming
// answer; surfacing the terminal mount failure here would end it early.
/// While the project-open mount is upgrading this owner a request waits
/// for that publication instead of taking the retryable warming answer;
/// a mount still awaiting its first generation, or one that failed,
/// answers at once.
fn mount(&self) -> DaemonAdvisoryCycleMountFuture<'_> {
let mut mount_state = self.mount_state.clone();
Box::pin(async move {
if code_index_disabled_for_scope(&self.code_index_schedulers, &self.scope) {
return DaemonAdvisoryCycleMountV1::Answers;
}
if self
.code_index_schedulers
.retained_text_owner_freshness_for_scope(&self.scope)
.await
.is_some()
|| self
.code_index_schedulers
.latest_feedback_generation_for_scope(&self.project_root, &self.scope)
.await
.is_some()
{
DaemonAdvisoryCycleMountV1::Mounting
} else {
DaemonAdvisoryCycleMountV1::Answers
let state = *mount_state.borrow_and_update();
match state {
AdvisoryMountStateV1::AwaitingGeneration | AdvisoryMountStateV1::Failed(_) => {
return DaemonAdvisoryCycleMountV1::Answers;
}
AdvisoryMountStateV1::Mounting | AdvisoryMountStateV1::Mounted => {}
}
DaemonAdvisoryCycleMountV1::Mounting(Box::pin(async move {
// Only a `Mounted` publisher closes without a later state; the
// replacement owner's publication wakes that waiter instead.
if mount_state.changed().await.is_err() {
std::future::pending::<()>().await;
}
}))
})
}

Expand Down
Loading
Loading