Skip to main content

graphforge_storage/
embedding_freshness.rs

1//! Deterministic freshness policy for verified embedding generations.
2
3use crate::{
4    EmbeddingCompatibilityId, EmbeddingGenerationId, EmbeddingGenerationManifest,
5    EmbeddingSourceFingerprint, EmbeddingSourceState, SearchArtifactError,
6};
7
8/// Changed-UUID percentage that makes a generation substantially stale.
9pub const EMBEDDING_SUBSTANTIAL_CHANGED_PERCENT: u64 = 5;
10/// Relevant committed mutation batches that make a generation substantially stale.
11pub const EMBEDDING_SUBSTANTIAL_MUTATION_BATCHES: u64 = 128;
12
13/// Durable mutation evidence observed after one generation source snapshot.
14#[derive(Clone, Copy, Debug, PartialEq, Eq)]
15pub struct EmbeddingMutationObservation {
16    current_source: EmbeddingSourceState,
17    changed_distinct_uuids: u64,
18    relevant_committed_batches: u64,
19    structural_mutation: bool,
20    scope_proven: bool,
21}
22
23impl EmbeddingMutationObservation {
24    /// Validate one bounded freshness observation.
25    ///
26    /// # Errors
27    /// Rejects changed UUID or structural evidence without a committed batch.
28    pub fn new(
29        current_source: EmbeddingSourceState,
30        changed_distinct_uuids: u64,
31        relevant_committed_batches: u64,
32        structural_mutation: bool,
33        scope_proven: bool,
34    ) -> Result<Self, SearchArtifactError> {
35        if relevant_committed_batches == 0 && (changed_distinct_uuids != 0 || structural_mutation) {
36            return Err(invalid(
37                "embedding mutation observation",
38                "changed UUID and structural evidence require a committed relevant batch",
39            ));
40        }
41        Ok(Self {
42            current_source,
43            changed_distinct_uuids,
44            relevant_committed_batches,
45            structural_mutation,
46            scope_proven,
47        })
48    }
49
50    /// Current durable dependency projection.
51    #[must_use]
52    pub const fn current_source(self) -> EmbeddingSourceState {
53        self.current_source
54    }
55
56    /// Distinct UUIDs changed across relevant committed batches.
57    #[must_use]
58    pub const fn changed_distinct_uuids(self) -> u64 {
59        self.changed_distinct_uuids
60    }
61
62    /// Relevant committed mutation batches since the recorded source.
63    #[must_use]
64    pub const fn relevant_committed_batches(self) -> u64 {
65        self.relevant_committed_batches
66    }
67
68    /// Whether any relevant topology mutation occurred.
69    #[must_use]
70    pub const fn structural_mutation(self) -> bool {
71        self.structural_mutation
72    }
73
74    /// Whether mutation scope was proven from durable metadata.
75    #[must_use]
76    pub const fn scope_proven(self) -> bool {
77        self.scope_proven
78    }
79}
80
81/// Freshness state for one verified, compatible, complete generation.
82#[derive(Clone, Copy, Debug, PartialEq, Eq)]
83pub enum EmbeddingFreshnessState {
84    /// Recorded and current dependency projections are exactly identical.
85    Fresh,
86    /// A bounded relevant mutation occurred below every substantial threshold.
87    Stale,
88    /// Ordinary search must refresh successfully before serving.
89    SubstantiallyStale,
90}
91
92/// Stable reason for a non-fresh state.
93#[derive(Clone, Copy, Debug, PartialEq, Eq)]
94pub enum EmbeddingFreshnessReason {
95    /// Relevant scope cannot be proven from durable metadata.
96    UnprovenMutationScope,
97    /// A structural space observed a relevant topology mutation.
98    StructuralMutation,
99    /// Any relevant mutation affected a recorded empty generation.
100    EmptyGenerationMutation,
101    /// At least five percent of recorded eligible UUIDs changed.
102    ChangedUuidFraction,
103    /// At least 128 relevant mutation batches accumulated.
104    MutationBatchLimit,
105    /// A known relevant mutation remains below substantial thresholds.
106    RelevantMutation,
107}
108
109impl EmbeddingFreshnessReason {
110    /// Stable bounded diagnostic token.
111    #[must_use]
112    pub const fn as_str(self) -> &'static str {
113        match self {
114            Self::UnprovenMutationScope => "unproven_mutation_scope",
115            Self::StructuralMutation => "structural_mutation",
116            Self::EmptyGenerationMutation => "empty_generation_mutation",
117            Self::ChangedUuidFraction => "changed_uuid_fraction",
118            Self::MutationBatchLimit => "mutation_batch_limit",
119            Self::RelevantMutation => "relevant_mutation",
120        }
121    }
122}
123
124/// Classified freshness plus the exact current durable source state.
125#[derive(Clone, Copy, Debug, PartialEq, Eq)]
126pub struct EmbeddingFreshness {
127    state: EmbeddingFreshnessState,
128    reason: Option<EmbeddingFreshnessReason>,
129    current_source: EmbeddingSourceState,
130}
131
132impl EmbeddingFreshness {
133    /// Classified state.
134    #[must_use]
135    pub const fn state(self) -> EmbeddingFreshnessState {
136        self.state
137    }
138
139    /// Stable reason, absent only for `fresh`.
140    #[must_use]
141    pub const fn reason(self) -> Option<EmbeddingFreshnessReason> {
142        self.reason
143    }
144
145    /// Current durable dependency projection used for classification.
146    #[must_use]
147    pub const fn current_source(self) -> EmbeddingSourceState {
148        self.current_source
149    }
150}
151
152/// Stable observable diagnostic emitted only for an explicit forced stale read.
153#[derive(Clone, Copy, Debug, PartialEq, Eq)]
154pub struct EmbeddingForcedStaleDiagnostic {
155    /// Compatibility lineage selected by the caller.
156    pub compatibility_id: EmbeddingCompatibilityId,
157    /// Last complete immutable generation served.
158    pub generation_id: EmbeddingGenerationId,
159    /// Source fingerprint recorded by that generation.
160    pub recorded_source: EmbeddingSourceFingerprint,
161    /// Current durable dependency fingerprint.
162    pub current_source: EmbeddingSourceFingerprint,
163    /// Deterministic substantial-staleness reason.
164    pub reason: EmbeddingFreshnessReason,
165}
166
167impl EmbeddingForcedStaleDiagnostic {
168    /// Stable, bounded, content-free diagnostic representation.
169    #[must_use]
170    pub fn stable_message(self) -> String {
171        format!(
172            "embedding_force_stale:v1 compatibility_id={} generation_id={} recorded_source={} current_source={} reason={}",
173            self.compatibility_id,
174            self.generation_id,
175            self.recorded_source,
176            self.current_source,
177            self.reason.as_str(),
178        )
179    }
180}
181
182/// Serving decision after compatibility and primary-data validation succeeded.
183#[derive(Clone, Copy, Debug, PartialEq, Eq)]
184pub enum EmbeddingReadDecision {
185    /// Serve the verified fresh generation.
186    ServeFresh,
187    /// Serve a mildly stale generation while refresh is queued.
188    ServeStale {
189        /// Stable non-substantial reason.
190        reason: EmbeddingFreshnessReason,
191    },
192    /// Refresh must succeed before ordinary search may serve.
193    RefreshRequired {
194        /// Stable substantial-staleness reason.
195        reason: EmbeddingFreshnessReason,
196    },
197    /// Explicitly serve the last complete substantially stale generation.
198    ServeForcedStale {
199        /// Stable observable forced-read diagnostic.
200        diagnostic: EmbeddingForcedStaleDiagnostic,
201    },
202}
203
204/// Classify one verified generation against current durable mutation evidence.
205///
206/// Precedence is deterministic: unproven scope, structural mutation, empty
207/// generation, changed UUID fraction, mutation-batch limit, then ordinary
208/// relevant mutation. Wall-clock timestamps never participate.
209///
210/// # Errors
211/// Rejects observations that predate the generation, impossible batch counts,
212/// unexplained fingerprint changes, and checked threshold arithmetic overflow.
213pub fn classify_embedding_freshness(
214    manifest: &EmbeddingGenerationManifest,
215    observation: EmbeddingMutationObservation,
216) -> Result<EmbeddingFreshness, SearchArtifactError> {
217    let recorded = manifest.source();
218    let current = observation.current_source;
219    if current.graph_generation() < recorded.graph_generation() {
220        return Err(invalid(
221            "embedding freshness observation",
222            "current graph generation predates the recorded generation",
223        ));
224    }
225    let generation_delta = current.graph_generation() - recorded.graph_generation();
226    if observation.relevant_committed_batches > generation_delta {
227        return Err(invalid(
228            "embedding freshness observation",
229            "relevant batch count exceeds the committed graph generation delta",
230        ));
231    }
232    let no_relevant_evidence = observation.relevant_committed_batches == 0
233        && observation.changed_distinct_uuids == 0
234        && !observation.structural_mutation
235        && observation.scope_proven;
236    if recorded.fingerprint() == current.fingerprint() && no_relevant_evidence {
237        return Ok(EmbeddingFreshness {
238            state: EmbeddingFreshnessState::Fresh,
239            reason: None,
240            current_source: current,
241        });
242    }
243    if recorded.fingerprint() != current.fingerprint()
244        && no_relevant_evidence
245        && observation.scope_proven
246    {
247        return Err(invalid(
248            "embedding freshness observation",
249            "source fingerprint changed without a relevant mutation or unproven scope",
250        ));
251    }
252
253    let reason = if !observation.scope_proven {
254        EmbeddingFreshnessReason::UnprovenMutationScope
255    } else if observation.structural_mutation {
256        EmbeddingFreshnessReason::StructuralMutation
257    } else if recorded.eligible_uuid_count() == 0 {
258        EmbeddingFreshnessReason::EmptyGenerationMutation
259    } else if changed_fraction_is_substantial(
260        observation.changed_distinct_uuids,
261        recorded.eligible_uuid_count(),
262    )? {
263        EmbeddingFreshnessReason::ChangedUuidFraction
264    } else if observation.relevant_committed_batches >= EMBEDDING_SUBSTANTIAL_MUTATION_BATCHES {
265        EmbeddingFreshnessReason::MutationBatchLimit
266    } else {
267        EmbeddingFreshnessReason::RelevantMutation
268    };
269    let state = if reason == EmbeddingFreshnessReason::RelevantMutation {
270        EmbeddingFreshnessState::Stale
271    } else {
272        EmbeddingFreshnessState::SubstantiallyStale
273    };
274    Ok(EmbeddingFreshness {
275        state,
276        reason: Some(reason),
277        current_source: current,
278    })
279}
280
281/// Select the serving boundary for a classified verified generation.
282#[must_use]
283pub fn decide_embedding_read(
284    manifest: &EmbeddingGenerationManifest,
285    freshness: EmbeddingFreshness,
286    force_stale: bool,
287) -> EmbeddingReadDecision {
288    match freshness.state {
289        EmbeddingFreshnessState::Fresh => EmbeddingReadDecision::ServeFresh,
290        EmbeddingFreshnessState::Stale => EmbeddingReadDecision::ServeStale {
291            reason: freshness
292                .reason
293                .expect("a stale classification always has a reason"),
294        },
295        EmbeddingFreshnessState::SubstantiallyStale if force_stale => {
296            EmbeddingReadDecision::ServeForcedStale {
297                diagnostic: EmbeddingForcedStaleDiagnostic {
298                    compatibility_id: manifest.compatibility_id(),
299                    generation_id: manifest.generation_id(),
300                    recorded_source: manifest.source().fingerprint(),
301                    current_source: freshness.current_source.fingerprint(),
302                    reason: freshness
303                        .reason
304                        .expect("a substantially stale classification always has a reason"),
305                },
306            }
307        }
308        EmbeddingFreshnessState::SubstantiallyStale => EmbeddingReadDecision::RefreshRequired {
309            reason: freshness
310                .reason
311                .expect("a substantially stale classification always has a reason"),
312        },
313    }
314}
315
316fn changed_fraction_is_substantial(
317    changed: u64,
318    recorded_eligible: u64,
319) -> Result<bool, SearchArtifactError> {
320    let changed_percent = changed.checked_mul(100).ok_or_else(arithmetic_overflow)?;
321    let threshold = recorded_eligible
322        .max(1)
323        .checked_mul(EMBEDDING_SUBSTANTIAL_CHANGED_PERCENT)
324        .ok_or_else(arithmetic_overflow)?;
325    Ok(changed_percent >= threshold)
326}
327
328fn arithmetic_overflow() -> SearchArtifactError {
329    SearchArtifactError::ResourceExhausted {
330        resource: "embedding_freshness_arithmetic",
331        limit: u64::MAX,
332    }
333}
334
335fn invalid(field: &'static str, reason: impl Into<String>) -> SearchArtifactError {
336    SearchArtifactError::InvalidSelector {
337        field,
338        reason: reason.into(),
339    }
340}
341
342#[cfg(test)]
343mod tests {
344    use super::*;
345    use crate::{
346        EmbeddingContentDigest, EmbeddingGenerationManifestInput, EmbeddingPublicationFingerprint,
347    };
348
349    fn source(generation: u64, eligible: u64, marker: u8) -> EmbeddingSourceState {
350        EmbeddingSourceState::new(
351            generation,
352            [marker; 32],
353            [marker.wrapping_add(1); 32],
354            eligible,
355        )
356    }
357
358    fn manifest(recorded: EmbeddingSourceState, generated_at: i64) -> EmbeddingGenerationManifest {
359        EmbeddingGenerationManifest::new(EmbeddingGenerationManifestInput {
360            compatibility_id: EmbeddingCompatibilityId::from_hex(&"11".repeat(32)).unwrap(),
361            source: recorded,
362            content_digest: EmbeddingContentDigest::from_hex(&"22".repeat(32)).unwrap(),
363            vector_count: recorded.eligible_uuid_count(),
364            dimension: 2,
365            generated_at_micros: generated_at,
366            committed_at_micros: generated_at + 1,
367            publication_fingerprint: EmbeddingPublicationFingerprint::from_hex(&"33".repeat(32))
368                .unwrap(),
369        })
370        .unwrap()
371    }
372
373    fn observation(
374        current: EmbeddingSourceState,
375        changed: u64,
376        batches: u64,
377        structural: bool,
378        proven: bool,
379    ) -> EmbeddingMutationObservation {
380        EmbeddingMutationObservation::new(current, changed, batches, structural, proven).unwrap()
381    }
382
383    #[test]
384    fn exact_source_without_mutation_evidence_is_fresh_and_time_independent() {
385        let recorded = source(7, 100, 1);
386        let observation = observation(recorded, 0, 0, false, true);
387        let early = classify_embedding_freshness(&manifest(recorded, 10), observation).unwrap();
388        let late = classify_embedding_freshness(&manifest(recorded, 9_999), observation).unwrap();
389        assert_eq!(early, late);
390        assert_eq!(early.state(), EmbeddingFreshnessState::Fresh);
391        assert_eq!(early.reason(), None);
392        assert_eq!(
393            decide_embedding_read(&manifest(recorded, 10), early, true),
394            EmbeddingReadDecision::ServeFresh
395        );
396    }
397
398    #[test]
399    fn substantial_reason_precedence_is_deterministic() {
400        let recorded = source(1, 100, 1);
401        let current = source(130, 100, 2);
402        let cases = [
403            (
404                observation(current, 100, 129, true, false),
405                EmbeddingFreshnessReason::UnprovenMutationScope,
406            ),
407            (
408                observation(current, 100, 129, true, true),
409                EmbeddingFreshnessReason::StructuralMutation,
410            ),
411            (
412                observation(current, 5, 127, false, true),
413                EmbeddingFreshnessReason::ChangedUuidFraction,
414            ),
415            (
416                observation(current, 4, 128, false, true),
417                EmbeddingFreshnessReason::MutationBatchLimit,
418            ),
419        ];
420        for (observation, expected) in cases {
421            let freshness =
422                classify_embedding_freshness(&manifest(recorded, 10), observation).unwrap();
423            assert_eq!(
424                freshness.state(),
425                EmbeddingFreshnessState::SubstantiallyStale
426            );
427            assert_eq!(freshness.reason(), Some(expected));
428        }
429    }
430
431    #[test]
432    fn any_relevant_change_to_an_empty_generation_is_substantial() {
433        let recorded = source(1, 0, 1);
434        let freshness = classify_embedding_freshness(
435            &manifest(recorded, 10),
436            observation(source(2, 1, 2), 1, 1, false, true),
437        )
438        .unwrap();
439        assert_eq!(
440            freshness.reason(),
441            Some(EmbeddingFreshnessReason::EmptyGenerationMutation)
442        );
443    }
444
445    #[test]
446    fn changed_uuid_fraction_uses_exact_five_percent_boundary() {
447        let recorded = source(1, 100, 1);
448        let below = classify_embedding_freshness(
449            &manifest(recorded, 10),
450            observation(source(2, 100, 2), 4, 1, false, true),
451        )
452        .unwrap();
453        let exact = classify_embedding_freshness(
454            &manifest(recorded, 10),
455            observation(source(2, 100, 2), 5, 1, false, true),
456        )
457        .unwrap();
458        assert_eq!(below.state(), EmbeddingFreshnessState::Stale);
459        assert_eq!(
460            exact.reason(),
461            Some(EmbeddingFreshnessReason::ChangedUuidFraction)
462        );
463    }
464
465    #[test]
466    fn mutation_batch_limit_uses_exact_128_boundary() {
467        let recorded = source(1, 1_000, 1);
468        let below = classify_embedding_freshness(
469            &manifest(recorded, 10),
470            observation(source(128, 1_000, 2), 0, 127, false, true),
471        )
472        .unwrap();
473        let exact = classify_embedding_freshness(
474            &manifest(recorded, 10),
475            observation(source(129, 1_000, 2), 0, 128, false, true),
476        )
477        .unwrap();
478        assert_eq!(below.state(), EmbeddingFreshnessState::Stale);
479        assert_eq!(
480            exact.reason(),
481            Some(EmbeddingFreshnessReason::MutationBatchLimit)
482        );
483    }
484
485    #[test]
486    fn forced_stale_is_explicit_deterministic_and_content_free() {
487        let recorded = source(1, 100, 1);
488        let manifest = manifest(recorded, 10);
489        let freshness = classify_embedding_freshness(
490            &manifest,
491            observation(source(2, 100, 2), 100, 1, false, true),
492        )
493        .unwrap();
494        assert_eq!(
495            decide_embedding_read(&manifest, freshness, false),
496            EmbeddingReadDecision::RefreshRequired {
497                reason: EmbeddingFreshnessReason::ChangedUuidFraction,
498            }
499        );
500        let EmbeddingReadDecision::ServeForcedStale { diagnostic } =
501            decide_embedding_read(&manifest, freshness, true)
502        else {
503            panic!("force must select the last complete generation");
504        };
505        let message = diagnostic.stable_message();
506        assert_eq!(message, diagnostic.stable_message());
507        assert!(message.starts_with("embedding_force_stale:v1 compatibility_id="));
508        assert!(message.ends_with("reason=changed_uuid_fraction"));
509        assert!(!message.contains("vector"));
510        assert!(!message.contains("body"));
511    }
512
513    #[test]
514    fn inconsistent_observations_and_overflow_are_structured() {
515        let recorded = source(5, 100, 1);
516        assert!(EmbeddingMutationObservation::new(recorded, 1, 0, false, true).is_err());
517        for observation in [
518            observation(source(4, 100, 2), 0, 0, false, false),
519            observation(source(6, 100, 2), 0, 2, false, true),
520            observation(source(6, 100, 2), 0, 0, false, true),
521        ] {
522            assert!(classify_embedding_freshness(&manifest(recorded, 10), observation).is_err());
523        }
524
525        let overflow_recorded = source(1, u64::MAX, 1);
526        let error = classify_embedding_freshness(
527            &manifest(overflow_recorded, 10),
528            observation(source(2, u64::MAX, 2), u64::MAX, 1, false, true),
529        )
530        .unwrap_err();
531        assert!(matches!(
532            error,
533            SearchArtifactError::ResourceExhausted {
534                resource: "embedding_freshness_arithmetic",
535                ..
536            }
537        ));
538    }
539}