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
14 changes: 3 additions & 11 deletions crates/tracedecay-session-temporal-store/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ use self::hydration::GlobalDbTemporalHydrationPort;
use self::participant_freeze::{
freeze_participants, freeze_prepared_candidate_participants, root_readiness,
};
use self::retrieval::{GlobalDbPreparedCandidatePort, GlobalDbTemporalReadPort};
use self::retrieval::GlobalDbTemporalReadPort;
use self::sql::TemporalSqlRead;
use tracedecay_lcm::payload::read_verified_payload_content_with_checkpoint;

Expand Down Expand Up @@ -893,16 +893,8 @@ impl<'db, D: SessionTemporalRegisteredDb + Sync>
request.direct_anchor(),
request.snapshot_request().semantic_filter().goals,
);
let preparation = GlobalDbPreparedCandidatePort::new(
&candidate_read,
request.snapshot_request(),
&plan,
);
let prepared =
tracedecay_temporal_query::ports::prepare_temporal_candidate_cohort(
request.snapshot_request(),
&preparation,
)
let prepared = candidate_read
.prepare_root_candidate_cohort(request.snapshot_request(), &plan)
.await
.map_err(map_control_error)?;
let readiness = root_readiness(&temporal_read, request).await?;
Expand Down
109 changes: 64 additions & 45 deletions crates/tracedecay-session-temporal-store/src/retrieval.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,12 @@ use tracedecay_domain::{
use tracedecay_runtime_core::db::engine;
use tracedecay_temporal_query::candidates::{CandidateChannel, CandidatePlan};
use tracedecay_temporal_query::ports::{
CandidatePageSink, MeasuredTemporalValue, PageRequest, PageStatus, PortFuture,
TemporalCandidateFilterV1, TemporalCandidatePreparationPort, TemporalExecutionSnapshot,
TemporalMessageTypeFilterV1, TemporalPortError, TemporalReadPort, TemporalRecordPageSink,
CANDIDATE_READ_BUDGET, CandidateFieldCaps, CandidatePageSink, CandidateReadState,
MeasuredTemporalValue, PageLimits, PageRequest, PageStatus, PortFuture,
TemporalCandidateFilterV1, TemporalExecutionSnapshot, TemporalMessageTypeFilterV1,
TemporalPortError, TemporalPreparedCandidateCohort, TemporalReadPort, TemporalRecordPageSink,
TemporalRetrievalScope, TemporalSessionScopeFilterV1, TemporalSnapshotRequest,
await_controlled, begin_prepared_candidate_pull, commit_prepared_candidate_pull,
};
use tracedecay_temporal_query::ranking::RankingCandidate;

Expand Down Expand Up @@ -182,48 +184,6 @@ pub struct GlobalDbTemporalReadPort<'a> {
git_scope_session_ids: Option<&'a BTreeSet<(String, String)>>,
}

pub(super) struct GlobalDbPreparedCandidatePort<'port, 'db, 'request> {
read_port: &'port GlobalDbTemporalReadPort<'db>,
request: &'request TemporalSnapshotRequest,
plan: &'request CandidatePlan,
}

impl<'port, 'db, 'request> GlobalDbPreparedCandidatePort<'port, 'db, 'request> {
#[hotpath::skip]
pub(super) const fn new(
read_port: &'port GlobalDbTemporalReadPort<'db>,
request: &'request TemporalSnapshotRequest,
plan: &'request CandidatePlan,
) -> Self {
Self {
read_port,
request,
plan,
}
}
}

impl TemporalCandidatePreparationPort for GlobalDbPreparedCandidatePort<'_, '_, '_> {
fn produce_prepared_candidate_page<'a>(
&'a self,
request: PageRequest,
sink: &'a mut CandidatePageSink<'_>,
) -> PortFuture<'a, PageStatus> {
Box::pin(async move {
self.read_port
.produce_candidates_from_request(
&TemporalRetrievalScope::AllSessionsInAuthorizedRoot,
self.request,
1,
self.plan,
&request,
sink,
)
.await
})
}
}

struct SessionReadRelationAuthority<'a> {
scope: &'a SessionRelationScope,
store: SessionRelationGraphStore,
Expand Down Expand Up @@ -276,6 +236,65 @@ impl<'a> GlobalDbTemporalReadPort<'a> {
}
}

#[hotpath::measure(
future = true,
label = "session_temporal.query.prepare_root_candidates"
)]
pub(super) async fn prepare_root_candidate_cohort(
&self,
request: &TemporalSnapshotRequest,
plan: &CandidatePlan,
) -> Result<TemporalPreparedCandidateCohort, TemporalPortError> {
request.execution_control().checkpoint()?;
let limits = request.limits();
let candidate_page_items = limits.candidate_limit.min(64);
let candidate_limits = PageLimits::new(
limits.candidate_limit,
limits.candidate_total_bytes,
limits.candidate_item_bytes,
candidate_page_items,
)?;
let mut state = CandidateReadState::new(candidate_limits);
let mut candidates = Vec::with_capacity(limits.candidate_limit.min(256));
let scope = TemporalRetrievalScope::AllSessionsInAuthorizedRoot;
loop {
let limits = begin_prepared_candidate_pull(request, &mut state)?;
let control = request.execution_control();
let field_caps = CandidateFieldCaps::new(
limits.candidate_stable_id_bytes,
limits.candidate_anchor_id_bytes,
limits.candidate_metadata_field_bytes,
);
let page_request = state.request(limits.candidate_key_bytes, Some(field_caps));
let mut sink = state.begin_page(
control,
limits.candidate_key_bytes,
Some(field_caps),
CANDIDATE_READ_BUDGET,
);
let status = await_controlled(
control,
Box::pin(self.produce_candidates_from_request(
&scope,
request,
1,
plan,
&page_request,
&mut sink,
)),
)
.await?;
let page = sink.finish(status)?;
let page = commit_prepared_candidate_pull(&mut state, page)?;
let status = page.status();
candidates.extend(page.into_items());
if status == PageStatus::Complete {
break;
}
}
TemporalPreparedCandidateCohort::new(candidates)
}

#[hotpath::skip]
async fn candidate_matches_filter(
&self,
Expand Down
124 changes: 122 additions & 2 deletions crates/tracedecay-session-temporal-store/src/retrieval/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,15 @@ use tracedecay_runtime_core::db::{
engine::{Connection, Executor, TestConnection, Value as SqlValue},
};
use tracedecay_temporal_query::candidates::CandidateChannel;
use tracedecay_temporal_query::plan_temporal_candidates;
use tracedecay_temporal_query::ports::{
BindingDigest, ExecutionControl, KernelVersions, PageRequest, TemporalAuthorizedRoot,
BindingDigest, CANDIDATE_READ_BUDGET, CandidateFieldCaps, CandidateReadState, ExecutionControl,
ExecutionLimits, KernelVersions, PageLimits, PageRequest, PageStatus, TemporalAuthorizedRoot,
TemporalExecutionSnapshot, TemporalParticipantAuthorization, TemporalParticipantGeneration,
TemporalParticipantManifest, TemporalPortError, TemporalPreparedCandidateCohort,
TemporalRecord, TemporalRetrievalScope, TemporalSnapshotRequest, TemporalSourceAccess,
TemporalWatermarks,
TemporalWatermarks, await_controlled, begin_prepared_candidate_pull,
commit_prepared_candidate_pull,
};
use tracedecay_temporal_query::ranking::RankingCandidate;
use tracedecay_temporal_query::resolution::{SummarySourceState, ValidatedAuthorization};
Expand Down Expand Up @@ -234,6 +237,123 @@ fn prepared_root_candidates_bind_each_frozen_participant_generation() {
);
}

fn root_preparation_request(mode: TemporalModeV1) -> TemporalSnapshotRequest {
root_snapshot_with_mode(1, None, mode).request().clone()
}

#[tokio::test]
async fn root_candidate_preparation_matches_direct_request_producer() {
let dir = tempdir().expect("temporary directory");
let runtime = HostAdmissionTestRuntimeV1::profile(dir.path())
.await
.expect("registered profile runtime");
runtime.seed_candidate_query_fixture_for_test().await;
let read = runtime.retrieval_read_for_test().await;
let adapter = read.adapter();
let request = root_preparation_request(TemporalModeV1::Current);
let plan = plan_temporal_candidates("needle candidate", None, false);
let direct = adapter
.prepare_root_candidate_cohort(&request, &plan)
.await
.expect("direct root candidate cohort");
let candidate_limits = PageLimits::new(
request.limits().candidate_limit,
request.limits().candidate_total_bytes,
request.limits().candidate_item_bytes,
request.limits().candidate_limit.min(64),
)
.expect("valid limits");
let mut state = CandidateReadState::new(candidate_limits);
let mut via_pages = Vec::new();
let scope = TemporalRetrievalScope::AllSessionsInAuthorizedRoot;
loop {
let limits = begin_prepared_candidate_pull(&request, &mut state).expect("pull limits");
let control = request.execution_control();
let field_caps = CandidateFieldCaps::new(
limits.candidate_stable_id_bytes,
limits.candidate_anchor_id_bytes,
limits.candidate_metadata_field_bytes,
);
let page_request = state.request(limits.candidate_key_bytes, Some(field_caps));
let mut sink = state.begin_page(
control,
limits.candidate_key_bytes,
Some(field_caps),
CANDIDATE_READ_BUDGET,
);
let status = await_controlled(
control,
Box::pin(adapter.produce_candidates_from_request(
&scope,
&request,
1,
&plan,
&page_request,
&mut sink,
)),
)
.await
.expect("page-driven root candidate cohort");
let page = sink.finish(status).expect("bounded page");
let page = commit_prepared_candidate_pull(&mut state, page).expect("committed page");
let status = page.status();
via_pages.extend(page.into_items());
if status == PageStatus::Complete {
break;
}
}
assert_eq!(
direct.candidates(),
via_pages.as_slice(),
"direct preparation must match the request producer path"
);
assert!(
!direct.candidates().is_empty(),
"root candidate preparation must return live fixture candidates"
);
}

#[tokio::test]
async fn root_candidate_preparation_preserves_item_byte_budget_failure() {
let dir = tempdir().expect("temporary directory");
let runtime = HostAdmissionTestRuntimeV1::profile(dir.path())
.await
.expect("registered profile runtime");
runtime.seed_candidate_query_fixture_for_test().await;
let read = runtime.retrieval_read_for_test().await;
let adapter = read.adapter();
let request = root_preparation_request(TemporalModeV1::Current).with_limits(ExecutionLimits {
candidate_limit: 8,
candidate_total_bytes: 64 * 1024,
candidate_item_bytes: 8,
..ExecutionLimits::default()
});
let plan = plan_temporal_candidates("needle candidate", None, false);
assert!(matches!(
adapter.prepare_root_candidate_cohort(&request, &plan).await,
Err(TemporalPortError::BudgetExceeded { .. })
));
}

#[tokio::test]
async fn root_candidate_preparation_preserves_live_cancellation() {
let dir = tempdir().expect("temporary directory");
let runtime = HostAdmissionTestRuntimeV1::profile(dir.path())
.await
.expect("registered profile runtime");
runtime.seed_candidate_query_fixture_for_test().await;
let read = runtime.retrieval_read_for_test().await;
let adapter = read.adapter();
let control = ExecutionControl::default();
control.cancel();
let request = root_preparation_request(TemporalModeV1::Current).with_execution_control(control);
let plan = plan_temporal_candidates("needle candidate", None, false);
assert_eq!(
adapter.prepare_root_candidate_cohort(&request, &plan).await,
Err(TemporalPortError::Cancelled)
);
}

fn record_request() -> PageRequest {
PageRequest::for_test(32, 64 * 1024, 8 * 1024, 32, 512)
}
Expand Down
25 changes: 25 additions & 0 deletions crates/tracedecay-temporal-query/src/ports/contracts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -181,6 +181,31 @@ pub async fn prepare_temporal_candidate_cohort(
TemporalPreparedCandidateCohort::new(candidates)
}

pub fn begin_prepared_candidate_pull(
request: &TemporalSnapshotRequest,
state: &mut CandidateReadState,
) -> Result<ExecutionLimits, TemporalPortError> {
begin_pull_request(
request,
state,
|limits| {
(
limits.candidate_limit,
limits.candidate_total_bytes,
limits.candidate_item_bytes,
)
},
CANDIDATE_READ_BUDGET,
)
}

pub fn commit_prepared_candidate_pull(
state: &mut CandidateReadState,
page: BoundedPage<RankingCandidate>,
) -> Result<BoundedPage<RankingCandidate>, TemporalPortError> {
commit_pulled_page(state, page, CANDIDATE_READ_BUDGET)
}

pub async fn pull_candidate_page(
port: &impl TemporalReadPort,
snapshot: &TemporalExecutionSnapshot,
Expand Down
2 changes: 1 addition & 1 deletion crates/tracedecay-temporal-query/src/ports/execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -238,7 +238,7 @@ impl ExecutionControl {
}
}

pub(crate) async fn await_controlled<T, E>(
pub async fn await_controlled<T, E>(
control: &ExecutionControl,
future: impl Future<Output = Result<T, E>>,
) -> Result<T, E>
Expand Down
Loading
Loading