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 @@ -37,8 +37,7 @@ use tracedecay_query::retrieval::rerank::RerankExecutionControlV1;
use tracedecay_query::retrieval::semantic::{
SemanticAbstentionDispositionV1, SemanticAbstentionV1, SemanticCompositionExecutionAuthorityV1,
SemanticCompositionExecutionOutcomeV1, SemanticExecutionControl, SemanticQueryModeV1,
SemanticQueryServiceError, SemanticRerankExecutionPortV1, SemanticRerankReadinessV1,
SemanticRetrievalRequestV1,
SemanticQueryServiceError, SemanticRetrievalRequestV1, apply_bounded_rerank_outcome,
};

#[derive(Clone)]
Expand Down Expand Up @@ -582,38 +581,11 @@ impl CodeIndexSchedulerRegistryV1 {
label = "daemon.query.semantic.vector_and_lane"
)
.await?;
let mut rerank_executor = authority
.rerank
.as_ref()
.and_then(|configured| {
configured
.mounted
.as_ref()
.filter(|rerank| rerank.compatibility() == &configured.pins)
})
.map(|rerank| SemanticRerankExecutorV1 {
rerank,
code_generation,
query_view,
control,
});
let rerank_readiness = if authority.execution.rerank_policy().is_none() {
None
} else {
Some(match rerank_executor.as_mut() {
Some(executor) => SemanticRerankReadinessV1::Ready(executor),
None => SemanticRerankReadinessV1::Unavailable(
tracedecay_domain::SanitizedStageFailure::AuthorityUnavailable,
),
})
};
let outcome = hotpath::measure_block!("daemon.query.semantic.compose", {
authority.execution.execute(
base,
authorized_query,
outcome,
semantic_abstention_disposition(mode),
rerank_readiness,
)
})?;
match outcome {
Expand All @@ -624,7 +596,22 @@ impl CodeIndexSchedulerRegistryV1 {
abstention,
fallback,
}),
SemanticCompositionExecutionOutcomeV1::Augmented(executed) => {
SemanticCompositionExecutionOutcomeV1::Augmented(mut executed) => {
if authorized_query
.request_cursor
.as_ref()
.and_then(|cursor| cursor.semantic.as_ref())
.is_none()
{
executed.rerank = apply_configured_semantic_rerank(
&authority,
code_generation,
query_view,
base,
&mut executed.composition,
control,
);
}
let mut composition = executed.composition;
let Some(query_authority) = hotpath::future!(
self.query_authority_for_scope(scope),
Expand Down Expand Up @@ -704,33 +691,45 @@ where
}
}

struct SemanticRerankExecutorV1<'a, C: ?Sized> {
rerank: &'a ProductionCodeRerankAuthorityV1,
code_generation: &'a CodeIndexPublishedGenerationV1,
query_view: &'a EphemeralSanitizedQueryViewV1,
control: &'a C,
fn mounted_compatible_rerank(
configured: Option<&ConfiguredRerankAuthorityV1>,
) -> Option<&ProductionCodeRerankAuthorityV1> {
configured.and_then(|configured| {
configured
.mounted
.as_ref()
.filter(|rerank| rerank.compatibility() == &configured.pins)
})
}

impl<C> SemanticRerankExecutionPortV1 for SemanticRerankExecutorV1<'_, C>
fn apply_configured_semantic_rerank<C>(
authority: &SemanticQueryAuthorityV1,
code_generation: &CodeIndexPublishedGenerationV1,
query_view: &EphemeralSanitizedQueryViewV1,
request: &RetrievalRequest,
composition: &mut CompositionOutputV1,
control: &C,
) -> OptionalStagePublicStatus
where
C: SemanticExecutionControl + ?Sized,
{
fn execute_rerank(
&mut self,
request: &RetrievalRequest,
policy: &tracedecay_domain::RerankPolicy,
pre_rerank: &[tracedecay_domain::RankedCandidate],
) -> tracedecay_query::retrieval::rerank::BoundedRerankOutcomeV1 {
let rerank_control = SemanticRerankControlV1(self.control);
self.rerank.execute(
self.code_generation,
self.query_view,
request,
policy,
pre_rerank,
&rerank_control,
)
}
let Some(policy) = authority.execution.rerank_policy() else {
return OptionalStagePublicStatus::NotRequested;
};
let Some(rerank) = mounted_compatible_rerank(authority.rerank.as_ref()) else {
return OptionalStagePublicStatus::Unavailable(
tracedecay_domain::SanitizedStageFailure::AuthorityUnavailable,
);
};
let outcome = rerank.execute(
code_generation,
query_view,
request,
policy,
&composition.ranked_candidates,
&SemanticRerankControlV1(control),
);
apply_bounded_rerank_outcome(composition, outcome)
}

fn paginate_semantic_composition(
Expand Down Expand Up @@ -900,6 +899,11 @@ mod tests {

use super::*;
use tracedecay_query::retrieval::fusion::RetrievalCursorKeyringV1;
use tracedecay_query::retrieval::rerank::{
AdmittedNativeRerankExecutorV1, DeterministicLocalRerankExecutorV1, LocalRerankFailureV1,
LocalRerankInputV1, LocalRerankPermitV1,
};
use tracedecay_semantic_contracts::RerankCompatibilityPinsV1;

fn id<T>(value: &str) -> T
where
Expand Down Expand Up @@ -1772,4 +1776,83 @@ mod tests {
} if selected == generation
));
}

struct IdentityRerankExecutorV1 {
digest: ManifestDigest,
}

impl DeterministicLocalRerankExecutorV1 for IdentityRerankExecutorV1 {
fn planned_model_invocations(
&self,
_candidate_count: u32,
) -> Result<u32, LocalRerankFailureV1> {
Ok(1)
}

fn rerank(
&self,
_policy: &tracedecay_domain::RerankPolicy,
inputs: &[LocalRerankInputV1<'_>],
_permit: LocalRerankPermitV1,
) -> Result<Vec<RetrievalAnchorId>, LocalRerankFailureV1> {
Ok(inputs
.iter()
.map(|input| input.candidate.candidate.anchor_id.clone())
.collect())
}
}

impl AdmittedNativeRerankExecutorV1 for IdentityRerankExecutorV1 {
fn artifact_manifest_digest(&self) -> &ManifestDigest {
&self.digest
}
}

fn rerank_pins(byte: char) -> RerankCompatibilityPinsV1 {
RerankCompatibilityPinsV1 {
implementation_revision: id("rerank.fastembed.production.v1"),
artifact_manifest_digest: digest(byte),
runtime_compatibility_digest: digest(byte),
}
}

#[test]
fn configured_rerank_is_unavailable_when_unmounted_or_pins_diverge() {
let pins = rerank_pins('a');
let unmounted = ConfiguredRerankAuthorityV1 {
pins: pins.clone(),
mounted: None,
};
assert!(mounted_compatible_rerank(Some(&unmounted)).is_none());
assert!(mounted_compatible_rerank(None).is_none());

let mounted = ProductionCodeRerankAuthorityV1::from_executor_for_test(
rerank_pins('b'),
Arc::new(IdentityRerankExecutorV1 {
digest: digest('b'),
}),
);
let mismatched = ConfiguredRerankAuthorityV1 {
pins,
mounted: Some(mounted),
};
assert!(mounted_compatible_rerank(Some(&mismatched)).is_none());
}

#[test]
fn configured_rerank_selects_the_mounted_authority_with_exact_pins() {
let pins = rerank_pins('c');
let mounted = ProductionCodeRerankAuthorityV1::from_executor_for_test(
pins.clone(),
Arc::new(IdentityRerankExecutorV1 {
digest: digest('c'),
}),
);
let configured = ConfiguredRerankAuthorityV1 {
pins: pins.clone(),
mounted: Some(mounted),
};
let selected = mounted_compatible_rerank(Some(&configured)).expect("compatible mount");
assert_eq!(selected.compatibility(), &pins);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -46,9 +46,9 @@ use tracedecay_application::semantic_runtime::{
use tracedecay_graph_db::NeverCancelled;
#[cfg(feature = "semantic-fastembed")]
use tracedecay_runtime_core::db::{Database, DatabaseAuthority, TestDatabaseRuntimeMode};
use tracedecay_semantic_contracts::SemanticFallbackReasonV1;
#[cfg(feature = "semantic-fastembed")]
use tracedecay_semantic_contracts::{DEFAULT_FASTEMBED_MODEL_ID, SemanticResourceCeilings};
use tracedecay_semantic_contracts::{RerankCompatibilityPinsV1, SemanticFallbackReasonV1};

use super::registry::{
ColdMountOpenEventV1, ServingGenerationInstallationOutcomeV1,
Expand All @@ -65,7 +65,9 @@ use crate::code_index::production::{
CodeIndexProductionErrorV1, CodeIndexPublicationStoreErrorV1,
UninterruptibleCodeIndexControlV1, VerifiedSealedLexicalPageReadV1,
};
use crate::semantic_code::rerank_adapter::GenerationBoundCodeRerankViewsV1;
use crate::semantic_code::rerank_adapter::{
GenerationBoundCodeRerankViewsV1, ProductionCodeRerankAuthorityV1,
};
use tracedecay_query::retrieval::QueryAuthorityV1;
use tracedecay_query::retrieval::exact::{
CentralExactAdmissionAuthorityV1, ExactAdmissionAuthority, ExactLaneRequest,
Expand All @@ -77,9 +79,10 @@ use tracedecay_query::retrieval::lexical::{
LexicalRouteKindV1, LexicalRoutingV1,
};
use tracedecay_query::retrieval::rerank::{
BoundedRerankRuntimeV1, DeterministicLocalRerankExecutorV1, LocalRerankFailureV1,
LocalRerankInputV1, LocalRerankPermitV1, RerankExecutionControlV1,
AdmittedNativeRerankExecutorV1, BoundedRerankRuntimeV1, DeterministicLocalRerankExecutorV1,
LocalRerankFailureV1, LocalRerankInputV1, LocalRerankPermitV1, RerankExecutionControlV1,
};
use tracedecay_query::retrieval::semantic::apply_bounded_rerank_outcome;
use tracedecay_query::retrieval::semantic::{
SemanticAbstentionV1, SemanticExecutionControl, SemanticQueryModeV1,
};
Expand Down Expand Up @@ -2496,6 +2499,15 @@ fn oversized_generations_still_produce_a_complete_retention_finding() {

struct MixedAnchorReverseRerankExecutorV1;

impl AdmittedNativeRerankExecutorV1 for MixedAnchorReverseRerankExecutorV1 {
fn artifact_manifest_digest(&self) -> &ManifestDigest {
static DIGEST: OnceLock<ManifestDigest> = OnceLock::new();
DIGEST.get_or_init(|| {
ManifestDigest::new(format!("sha256:{}", "a".repeat(64))).expect("artifact digest")
})
}
}

impl DeterministicLocalRerankExecutorV1 for MixedAnchorReverseRerankExecutorV1 {
fn planned_model_invocations(
&self,
Expand Down Expand Up @@ -2530,6 +2542,18 @@ impl RerankExecutionControlV1 for ReadyRerankControlV1 {
}
}

struct CancelledRerankControlV1;

impl RerankExecutionControlV1 for CancelledRerankControlV1 {
fn elapsed_micros(&self) -> u64 {
0
}

fn is_cancelled(&self) -> bool {
true
}
}

struct ReadySemanticControlV1;

impl SemanticExecutionControl for ReadySemanticControlV1 {
Expand Down Expand Up @@ -3803,18 +3827,74 @@ fn generation_bound_rerank_authorizes_mixed_symbol_and_chunk_anchors() {
deadline_micros: None,
};
let mut views = GenerationBoundCodeRerankViewsV1::new(&latest.generation, &query);
let outcome = BoundedRerankRuntimeV1::new(&mut views, &MixedAnchorReverseRerankExecutorV1)
.rerank(&request, &policy, &candidates, &ReadyRerankControlV1);
let runtime_outcome = BoundedRerankRuntimeV1::new(
&mut views,
&MixedAnchorReverseRerankExecutorV1,
)
.rerank(&request, &policy, &candidates, &ReadyRerankControlV1);
let pins = RerankCompatibilityPinsV1 {
implementation_revision: ComponentRevision::new("rerank.fastembed.production.v1")
.expect("implementation revision"),
artifact_manifest_digest: MixedAnchorReverseRerankExecutorV1
.artifact_manifest_digest()
.clone(),
runtime_compatibility_digest: ManifestDigest::new(format!("sha256:{}", "b".repeat(64)))
.expect("runtime digest"),
};
let authority = ProductionCodeRerankAuthorityV1::from_executor_for_test(
pins,
Arc::new(MixedAnchorReverseRerankExecutorV1),
);
let execute_outcome = authority.execute(
&latest.generation,
&query,
&request,
&policy,
&candidates,
&ReadyRerankControlV1,
);

assert_eq!(outcome.public_status, OptionalStagePublicStatus::Complete);
assert_eq!(execute_outcome, runtime_outcome);
assert_eq!(
execute_outcome.public_status,
OptionalStagePublicStatus::Complete
);
assert_eq!(
outcome
execute_outcome
.ordered_candidates
.iter()
.map(|candidate| candidate.candidate.anchor_id.clone())
.collect::<Vec<_>>(),
anchors.into_iter().rev().collect::<Vec<_>>()
);

let cancelled = authority.execute(
&latest.generation,
&query,
&request,
&policy,
&candidates,
&CancelledRerankControlV1,
);
assert_eq!(
cancelled.public_status,
OptionalStagePublicStatus::Cancelled
);
assert_eq!(cancelled.ordered_candidates, candidates);
let mut composition = tracedecay_query::retrieval::fusion::CompositionOutputV1 {
profile_id: request.profile_id.clone(),
ranked_candidates: candidates.clone(),
comparator_records: Vec::new(),
internal_lane_outcomes: BTreeMap::new(),
public_lane_statuses: BTreeMap::new(),
freshness: Vec::new(),
lane_checkpoints: Vec::new(),
dedupe_decisions: Vec::new(),
diversity_decisions: Vec::new(),
};
let status = apply_bounded_rerank_outcome(&mut composition, cancelled);
assert_eq!(status, OptionalStagePublicStatus::Cancelled);
assert_eq!(composition.ranked_candidates, candidates);
}

#[test]
Expand Down
2 changes: 1 addition & 1 deletion crates/tracedecay-query/src/retrieval/semantic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ mod service;
pub use execution_authority::{
ExecutedSemanticCompositionV1, SemanticCompositionAuthorityErrorV1,
SemanticCompositionExecutionAuthorityV1, SemanticCompositionExecutionOutcomeV1,
SemanticRerankExecutionPortV1, SemanticRerankReadinessV1, restore_frozen_semantic_order,
apply_bounded_rerank_outcome, restore_frozen_semantic_order,
};
pub use service::{
CalibratedSemanticQueryService, CompleteSemanticGenerationV1, SemanticAbstentionDispositionV1,
Expand Down
Loading
Loading