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
Original file line number Diff line number Diff line change
Expand Up @@ -466,6 +466,14 @@ fn serving_swap_gate() -> &'static Mutex<BTreeMap<PathBuf, WorkerStepGateV1>> {
GATE.get_or_init(|| Mutex::new(BTreeMap::new()))
}

/// Holds a complete-seat probe right after its read missed.
#[cfg(test)]
fn complete_seat_probe_miss_gate() -> &'static Mutex<BTreeMap<PathBuf, WorkerStepGateV1>> {
static GATE: std::sync::OnceLock<Mutex<BTreeMap<PathBuf, WorkerStepGateV1>>> =
std::sync::OnceLock::new();
GATE.get_or_init(|| Mutex::new(BTreeMap::new()))
}

/// Holds a worker's graph prepare before and after it decodes the active
/// generation.
#[cfg(test)]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -863,30 +863,42 @@ impl CodeIndexSchedulerRegistryV1 {
.clone()
};
let root = &root;
let slots = || async move {
let mounted = self.mounted.lock().await;
let worktree = mounted.get(root)?;
let seated = worktree
.serving_generation
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_some();
let text_published = worktree
.text_generation
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_some();
Some((seated, text_published))
};
let Ok(waited) = self
.wait_for_seat(root, deadline, |_| async move {
// Only a seat that was already there when the read declined it
// fails to serve this scope. The read's own demand can land
// the seat after the miss; its change signal re-probes.
let Some((seated_before_read, _)) = slots().await else {
return Ok(Some(CodeIndexSeatWaitV1::Cancelled));
};
if let Some(latest) = self.latest_complete_fresh_for_scope(scope).await {
return Ok(Some(CodeIndexSeatWaitV1::Seated(latest)));
}
let mounted = self.mounted.lock().await;
let Some(worktree) = mounted.get(root) else {
#[cfg(test)]
Self::wait_for_complete_seat_probe_miss_gate(root).await;
let Some((seated, text_published)) = slots().await else {
return Ok(Some(CodeIndexSeatWaitV1::Cancelled));
};
let seated = worktree
.serving_generation
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_some();
let text_published = worktree
.text_generation
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_some();
Ok::<_, Infallible>(if seated {
Ok::<_, Infallible>(if seated_before_read && seated {
Some(CodeIndexSeatWaitV1::Parked(
CodeIndexSeatParkV1::SeatNotServable,
))
} else if text_published {
} else if seated || text_published {
None
} else {
Some(CodeIndexSeatWaitV1::Parked(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,8 @@ use super::{
ColdMountPostCheckTestControlV1, PendingWakeDropGateTestV1, PendingWakeV1,
QueryAdmissionTestControlV1, ServingGenerationInstallationV1,
ServingGenerationRollbackOutcomeV1, WorkerStepGateV1, cold_mount_admission_barriers,
cold_mount_open_controls, cold_mount_post_check_controls, graph_decode_gate,
published_text_projection_gate, query_admission_controls, serving_swap_gate,
cold_mount_open_controls, cold_mount_post_check_controls, complete_seat_probe_miss_gate,
graph_decode_gate, published_text_projection_gate, query_admission_controls, serving_swap_gate,
unique_mounted_for_scope, wait_notified_if_unset,
};
use tracedecay_runtime_core::path_safety::canonical_existing_identity;
Expand Down Expand Up @@ -86,6 +86,39 @@ impl CodeIndexSchedulerRegistryV1 {
Self::pass_worker_step_gate(gate).await;
}

/// Hold the next complete-seat probe for `project_root` right after its
/// read missed. The first receiver resolves once the probe waits there;
/// sending on the returned sender releases it.
#[cfg(test)]
pub fn pause_next_complete_seat_probe_miss(
&self,
project_root: PathBuf,
) -> (
tokio::sync::oneshot::Receiver<()>,
tokio::sync::oneshot::Sender<()>,
) {
let (entered, entered_observed) = tokio::sync::oneshot::channel();
let (released, release) = tokio::sync::oneshot::channel();
let replaced = complete_seat_probe_miss_gate()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(project_root, WorkerStepGateV1 { entered, release });
assert!(
replaced.is_none(),
"one complete-seat probe gate per worktree"
);
(entered_observed, released)
}

#[cfg(test)]
pub(super) async fn wait_for_complete_seat_probe_miss_gate(project_root: &Path) {
let gate = complete_seat_probe_miss_gate()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(project_root);
Self::pass_worker_step_gate(gate).await;
}

/// Hold the next graph prepare of the worker for `project_root` twice:
/// right before it decodes the active generation, and right after the
/// decode returned. Each pair is the entered receiver and release sender.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ use std::time::{Duration, Instant};
use tempfile::TempDir;
use tracedecay_contracts::code_index_freshness::CodeIndexStalenessStateV1;
use tracedecay_domain::{CodeGenerationId, ProjectId, WorktreeId};
use tracedecay_runtime_core::path_safety::canonical_existing_identity;
use tracedecay_runtime_core::resident_memory::{
ProcessResidentMemoryV1, ProcessResidentSampleV1, ProcessSharedMemoryReservationV1,
ResidentHoldingV1, ResidentMemoryComponentIdV1, ResidentMemoryPressureV1, ResidentOwnerBytesV1,
Expand All @@ -16,9 +17,10 @@ use tracedecay_runtime_core::resident_memory::{
use super::super::{
CodeIndexCadenceTelemetryV1, CodeIndexCadenceTriggerV1, CodeIndexWorkerPhaseV1,
};

use super::{
CodeIndexSchedulerRegistryV1, GitFixture, core_search_request, git,
mounted_core_query_worktree_at, mounted_core_query_worktree_in, test_project_id,
CodeIndexSchedulerRegistryV1, GitFixture, SERVING_SEAT_FAILURE_CEILING, core_search_request,
git, mounted_core_query_worktree_at, mounted_core_query_worktree_in, test_project_id,
wait_for_generation_change, wait_for_live_complete_generation,
wait_for_queryable_text_generation, wait_for_settled_owner, wait_for_worker_phase,
with_untouched_fillers,
Expand Down Expand Up @@ -127,6 +129,60 @@ async fn an_idle_worktree_gives_back_its_decode_and_search_still_answers_fresh()
registry.shutdown().await;
}

/// The decode a complete read demands can seat between that read's miss and
/// its look at the slot. The read serves that seat instead of refusing it.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_complete_read_serves_the_seat_its_demand_lands_after_its_miss() {
let fixture = GitFixture::new(&[("src/main.rs", "fn main() {}\n")]);
let store = TempDir::new().expect("store root");
let owners = Arc::new(ResidentOwnersV1::new(IDLE_WINDOW));
let (registry, scope) = mounted_core_query_worktree_in(
CodeIndexSchedulerRegistryV1::new(1).with_resident_owners(Arc::clone(&owners)),
&fixture,
&store,
)
.await;
wait_for_settled_owner(&registry, fixture.path()).await;
assert_eq!(owners.release_idle(Instant::now() + IDLE_WINDOW).len(), 1);
assert!(
registry
.latest_complete_serving_for_test(fixture.path())
.await
.is_none(),
"only the text owner stays seated"
);

let root = canonical_existing_identity(fixture.path()).expect("canonical fixture");
let (missed, release_probe) = registry.pause_next_complete_seat_probe_miss(root);
let read = tokio::spawn({
let registry = registry.clone();
let scope = scope.clone();
async move {
registry
.latest_complete_fresh_for_scope_awaiting_seat(
&scope,
tokio::time::Instant::now() + SERVING_SEAT_FAILURE_CEILING,
)
.await
}
});
missed
.await
.expect("the read misses before the decode seats");
let seated = wait_for_live_complete_generation(&registry, fixture.path()).await;
release_probe.send(()).expect("the probe waits at its gate");

let served = read.await.expect("read task").expect(
"a seat that landed after the miss serves the read instead of ending it unavailable",
);
assert_eq!(
served.generation.manifest().generation_id,
seated.generation.manifest().generation_id
);

registry.shutdown().await;
}

/// A module with documented items, a struct with methods, cross-module calls
/// and a clone-sized body, so a decode holds every kind of page record.
fn linked_module(ordinal: usize) -> (String, String) {
Expand Down
Loading