Skip to main content

graphforge_storage/
embedding_refresh_config.rs

1//! Durable, content-free embedding refresh policy and terminal outcomes.
2
3use std::collections::{BTreeMap, BTreeSet};
4use std::fs::{File, OpenOptions};
5use std::io::Write;
6use std::path::{Path, PathBuf};
7use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
8
9use serde::{Deserialize, Serialize};
10use sha2::{Digest, Sha256};
11
12use crate::{
13    EmbeddingCompatibilityId, EmbeddingSourceFingerprint, SearchArtifactError,
14    SearchCoordinationLimits,
15};
16
17/// Refresh-control document schema implemented by this release.
18pub const EMBEDDING_REFRESH_CONFIG_VERSION: u32 = 1;
19/// Default process-local proactive refresh debounce.
20pub const DEFAULT_EMBEDDING_REFRESH_DEBOUNCE: Duration = Duration::from_millis(500);
21/// Default and maximum producer concurrency for one embedded project.
22pub const MAX_EMBEDDING_REFRESH_JOBS: usize = 2;
23/// Maximum accepted refresh-control bytes by default.
24pub const MAX_EMBEDDING_REFRESH_CONFIG_BYTES: usize = 256 * 1024;
25/// Maximum per-lineage policy/outcome records by default.
26pub const MAX_EMBEDDING_REFRESH_CONFIG_ENTRIES: usize = 1_024;
27
28const EMBEDDINGS_DIR: &str = "embeddings";
29const CONFIG_FILE: &str = "refresh.json";
30const CONFIG_LOCK: &str = ".refresh.lock";
31const CHECKSUM_DOMAIN: &[u8] = b"graphforge.embedding.refresh-config.v1\0";
32const MAX_DEBOUNCE: Duration = Duration::from_hours(1);
33
34/// Resource and coordination bounds for refresh-control access.
35#[derive(Clone, Copy, Debug)]
36pub struct EmbeddingRefreshConfigLimits {
37    /// Maximum canonical JSON bytes.
38    pub metadata_bytes: usize,
39    /// Maximum distinct compatibility-lineage records.
40    pub entries: usize,
41    /// Refresh-control writer lock and cleanup bounds.
42    pub coordination: SearchCoordinationLimits,
43}
44
45impl Default for EmbeddingRefreshConfigLimits {
46    fn default() -> Self {
47        Self {
48            metadata_bytes: MAX_EMBEDDING_REFRESH_CONFIG_BYTES,
49            entries: MAX_EMBEDDING_REFRESH_CONFIG_ENTRIES,
50            coordination: SearchCoordinationLimits::default(),
51        }
52    }
53}
54
55/// Project-wide refresh defaults.
56#[derive(Clone, Copy, Debug, PartialEq, Eq)]
57pub struct EmbeddingRefreshProjectPolicy {
58    /// Queue relevant mutations while this process is open.
59    pub proactive: bool,
60    /// Delay after the newest relevant mutation before work is ready.
61    pub debounce: Duration,
62    /// Maximum producer jobs active for this project.
63    pub max_concurrent_jobs: usize,
64}
65
66impl Default for EmbeddingRefreshProjectPolicy {
67    fn default() -> Self {
68        Self {
69            proactive: true,
70            debounce: DEFAULT_EMBEDDING_REFRESH_DEBOUNCE,
71            max_concurrent_jobs: MAX_EMBEDDING_REFRESH_JOBS,
72        }
73    }
74}
75
76/// Optional per-lineage overrides of project refresh defaults.
77#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
78pub struct EmbeddingRefreshSpacePolicy {
79    /// Override proactive enablement for this lineage.
80    pub proactive: Option<bool>,
81    /// Override mutation debounce for this lineage.
82    pub debounce: Option<Duration>,
83}
84
85/// Fully resolved policy for one compatibility lineage.
86#[derive(Clone, Copy, Debug, PartialEq, Eq)]
87pub struct ResolvedEmbeddingRefreshPolicy {
88    /// Effective proactive enablement.
89    pub proactive: bool,
90    /// Effective mutation debounce.
91    pub debounce: Duration,
92    /// Project-wide producer concurrency bound.
93    pub max_concurrent_jobs: usize,
94}
95
96/// Stable, content-free terminal refresh failure classification.
97#[derive(Clone, Copy, Debug, PartialEq, Eq)]
98pub enum EmbeddingRefreshFailureClass {
99    /// Provider or callback work failed.
100    Provider,
101    /// Caller input or producer output was invalid.
102    Validation,
103    /// A named time, memory, disk, queue, token, or cost bound was exhausted.
104    ResourceExhausted,
105    /// Durable storage or publication failed.
106    Storage,
107    /// The graph changed across both bounded attempts.
108    ConcurrentMutation,
109    /// Space identity or producer contract was incompatible.
110    Incompatible,
111    /// Primary or derived durable bytes were corrupt.
112    Corrupt,
113    /// The producer is unavailable in this process.
114    Unavailable,
115}
116
117/// Stable terminal interpretation of one refresh attempt.
118#[derive(Clone, Copy, Debug, PartialEq, Eq)]
119pub enum EmbeddingRefreshOutcomeStatus {
120    /// One complete generation published or was content-idempotently reused.
121    Succeeded,
122    /// Cooperative cancellation stopped private work.
123    Cancelled,
124    /// Refresh failed without publishing partial state.
125    Failed(EmbeddingRefreshFailureClass),
126}
127
128/// Content-free durable outcome for one compatibility lineage.
129#[derive(Clone, Copy, Debug, PartialEq, Eq)]
130pub struct EmbeddingRefreshOutcomeRecord {
131    /// Terminal completion classification.
132    pub status: EmbeddingRefreshOutcomeStatus,
133    /// Committed graph generation targeted by the attempt.
134    pub graph_generation: u64,
135    /// Exact source fingerprint targeted by the attempt.
136    pub source_fingerprint: EmbeddingSourceFingerprint,
137    /// Terminal UTC timestamp in microseconds since Unix epoch.
138    pub completed_at_micros: i64,
139}
140
141/// Content-free terminal outcome fields captured before durable lock acquisition.
142#[derive(Clone, Copy, Debug, PartialEq, Eq)]
143pub struct EmbeddingRefreshOutcomeAttempt {
144    /// Terminal completion classification.
145    pub status: EmbeddingRefreshOutcomeStatus,
146    /// Committed graph generation targeted by the attempt.
147    pub graph_generation: u64,
148    /// Exact source fingerprint targeted by the attempt.
149    pub source_fingerprint: EmbeddingSourceFingerprint,
150}
151
152/// One deterministic compatibility-lineage refresh record.
153#[derive(Clone, Copy, Debug, PartialEq, Eq)]
154pub struct EmbeddingRefreshSpaceState {
155    /// Exact embedding compatibility lineage.
156    pub compatibility_id: EmbeddingCompatibilityId,
157    /// Optional per-lineage override.
158    pub policy: Option<EmbeddingRefreshSpacePolicy>,
159    /// Optional latest content-free terminal outcome.
160    pub last_outcome: Option<EmbeddingRefreshOutcomeRecord>,
161}
162
163#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
164struct StoredSpaceState {
165    policy: Option<EmbeddingRefreshSpacePolicy>,
166    last_outcome: Option<EmbeddingRefreshOutcomeRecord>,
167}
168
169/// Fully validated durable refresh control state.
170#[derive(Clone, Debug, Default, PartialEq, Eq)]
171pub struct EmbeddingRefreshConfig {
172    project: EmbeddingRefreshProjectPolicy,
173    spaces: BTreeMap<EmbeddingCompatibilityId, StoredSpaceState>,
174}
175
176impl EmbeddingRefreshConfig {
177    /// Project-wide defaults.
178    #[must_use]
179    pub const fn project_policy(&self) -> EmbeddingRefreshProjectPolicy {
180        self.project
181    }
182
183    /// Number of lineages with an override or terminal outcome.
184    #[must_use]
185    pub fn len(&self) -> usize {
186        self.spaces.len()
187    }
188
189    /// Whether no per-lineage state is retained.
190    #[must_use]
191    pub fn is_empty(&self) -> bool {
192        self.spaces.is_empty()
193    }
194
195    /// Deterministic compatibility-identity-ordered lineage state.
196    #[must_use]
197    pub fn spaces(&self) -> Vec<EmbeddingRefreshSpaceState> {
198        self.spaces
199            .iter()
200            .map(|(compatibility_id, state)| EmbeddingRefreshSpaceState {
201                compatibility_id: *compatibility_id,
202                policy: state.policy,
203                last_outcome: state.last_outcome,
204            })
205            .collect()
206    }
207
208    /// Resolve one lineage's project defaults and optional overrides.
209    #[must_use]
210    pub fn resolved_policy(
211        &self,
212        compatibility_id: EmbeddingCompatibilityId,
213    ) -> ResolvedEmbeddingRefreshPolicy {
214        let override_policy = self
215            .spaces
216            .get(&compatibility_id)
217            .and_then(|state| state.policy);
218        ResolvedEmbeddingRefreshPolicy {
219            proactive: override_policy
220                .and_then(|policy| policy.proactive)
221                .unwrap_or(self.project.proactive),
222            debounce: override_policy
223                .and_then(|policy| policy.debounce)
224                .unwrap_or(self.project.debounce),
225            max_concurrent_jobs: self.project.max_concurrent_jobs,
226        }
227    }
228
229    fn to_canonical_json(
230        &self,
231        limits: EmbeddingRefreshConfigLimits,
232    ) -> Result<Vec<u8>, SearchArtifactError> {
233        validate_limits(limits)?;
234        if self.spaces.len() > limits.entries {
235            return Err(exhausted(
236                "embedding_refresh_config_entries",
237                limits.entries,
238            ));
239        }
240        validate_project_policy(self.project)?;
241        let spaces = self
242            .spaces
243            .iter()
244            .map(|(compatibility_id, state)| wire_space(*compatibility_id, *state))
245            .collect::<Result<Vec<_>, _>>()?;
246        let material = WireMaterial {
247            config_version: EMBEDDING_REFRESH_CONFIG_VERSION,
248            project: wire_project(self.project)?,
249            spaces,
250        };
251        let material_bytes = serde_json::to_vec(&material)
252            .map_err(|error| invalid("embedding refresh config", error.to_string()))?;
253        let bytes = serde_json::to_vec(&WireDocument {
254            config_version: material.config_version,
255            project: material.project,
256            spaces: material.spaces,
257            checksum: checksum(&material_bytes),
258        })
259        .map_err(|error| invalid("embedding refresh config", error.to_string()))?;
260        if bytes.len() > limits.metadata_bytes {
261            return Err(exhausted(
262                "embedding_refresh_config_bytes",
263                limits.metadata_bytes,
264            ));
265        }
266        Ok(bytes)
267    }
268}
269
270/// One explicit mutation of durable refresh control state.
271#[derive(Clone, Copy, Debug)]
272pub enum EmbeddingRefreshConfigUpdate {
273    /// Replace project-wide defaults.
274    SetProjectPolicy(EmbeddingRefreshProjectPolicy),
275    /// Set or clear one lineage override. Clearing retains its last outcome.
276    SetSpacePolicy {
277        /// Exact compatibility lineage.
278        compatibility_id: EmbeddingCompatibilityId,
279        /// Override to persist, or `None` to inherit project defaults.
280        policy: Option<EmbeddingRefreshSpacePolicy>,
281    },
282    /// Record one terminal content-free outcome.
283    RecordOutcome {
284        /// Exact compatibility lineage.
285        compatibility_id: EmbeddingCompatibilityId,
286        /// Newest terminal outcome.
287        outcome: EmbeddingRefreshOutcomeRecord,
288    },
289    /// Remove all policy and outcome state for one lineage.
290    RemoveSpace {
291        /// Exact compatibility lineage.
292        compatibility_id: EmbeddingCompatibilityId,
293    },
294}
295
296/// Read and validate durable embedding refresh controls.
297///
298/// A missing file returns pinned defaults without creating directories.
299///
300/// # Errors
301/// Returns structured cancellation, limit, corruption, incompatibility, or I/O errors.
302pub fn read_embedding_refresh_config<C>(
303    project_dir: &Path,
304    limits: EmbeddingRefreshConfigLimits,
305    mut checkpoint: C,
306) -> Result<EmbeddingRefreshConfig, SearchArtifactError>
307where
308    C: FnMut() -> Result<(), SearchArtifactError>,
309{
310    validate_limits(limits)?;
311    checkpoint()?;
312    read_config_file(&config_path(project_dir), limits)
313}
314
315/// Apply one refresh-control mutation under a bounded cross-process writer lock.
316///
317/// Exact-idempotent updates do not rewrite durable bytes. Publication uses a
318/// synchronized sibling temporary file and atomic replacement.
319///
320/// # Errors
321/// Returns structured validation, cancellation, lock, limit, corruption, or I/O errors.
322pub fn update_embedding_refresh_config<C>(
323    project_dir: &Path,
324    update: EmbeddingRefreshConfigUpdate,
325    limits: EmbeddingRefreshConfigLimits,
326    checkpoint: C,
327) -> Result<EmbeddingRefreshConfig, SearchArtifactError>
328where
329    C: FnMut() -> Result<(), SearchArtifactError>,
330{
331    mutate_embedding_refresh_config(project_dir, limits, checkpoint, |config, limits| {
332        apply_update(config, update, limits)
333    })
334}
335
336/// Record one terminal outcome with monotonic time assigned under the writer lock.
337///
338/// The latest durable same-lineage outcome is read only after cross-process lock
339/// acquisition. The new completion time is the current UTC microsecond or one
340/// microsecond after the prior outcome, whichever is greater.
341///
342/// # Errors
343/// Returns structured validation, cancellation, lock, clock-overflow, limit,
344/// corruption, or I/O errors. Failure leaves durable bytes unchanged.
345pub fn record_embedding_refresh_outcome<C>(
346    project_dir: &Path,
347    compatibility_id: EmbeddingCompatibilityId,
348    attempt: EmbeddingRefreshOutcomeAttempt,
349    limits: EmbeddingRefreshConfigLimits,
350    checkpoint: C,
351) -> Result<EmbeddingRefreshConfig, SearchArtifactError>
352where
353    C: FnMut() -> Result<(), SearchArtifactError>,
354{
355    record_embedding_refresh_outcome_with_clock(
356        project_dir,
357        compatibility_id,
358        attempt,
359        limits,
360        transaction_time_micros,
361        checkpoint,
362    )
363}
364
365fn record_embedding_refresh_outcome_with_clock<C, N>(
366    project_dir: &Path,
367    compatibility_id: EmbeddingCompatibilityId,
368    attempt: EmbeddingRefreshOutcomeAttempt,
369    limits: EmbeddingRefreshConfigLimits,
370    now: N,
371    checkpoint: C,
372) -> Result<EmbeddingRefreshConfig, SearchArtifactError>
373where
374    C: FnMut() -> Result<(), SearchArtifactError>,
375    N: FnOnce() -> i64,
376{
377    mutate_embedding_refresh_config(project_dir, limits, checkpoint, move |config, limits| {
378        let prior_time = config
379            .spaces
380            .get(&compatibility_id)
381            .and_then(|state| state.last_outcome)
382            .map_or(Ok(0), |outcome| {
383                outcome.completed_at_micros.checked_add(1).ok_or(
384                    SearchArtifactError::ResourceExhausted {
385                        resource: "embedding_refresh_outcome_timestamp",
386                        limit: u64::MAX,
387                    },
388                )
389            })?;
390        apply_update(
391            config,
392            EmbeddingRefreshConfigUpdate::RecordOutcome {
393                compatibility_id,
394                outcome: EmbeddingRefreshOutcomeRecord {
395                    status: attempt.status,
396                    graph_generation: attempt.graph_generation,
397                    source_fingerprint: attempt.source_fingerprint,
398                    completed_at_micros: now().max(prior_time),
399                },
400            },
401            limits,
402        )
403    })
404}
405
406fn mutate_embedding_refresh_config<C, M>(
407    project_dir: &Path,
408    limits: EmbeddingRefreshConfigLimits,
409    mut checkpoint: C,
410    mutate: M,
411) -> Result<EmbeddingRefreshConfig, SearchArtifactError>
412where
413    C: FnMut() -> Result<(), SearchArtifactError>,
414    M: FnOnce(
415        &mut EmbeddingRefreshConfig,
416        EmbeddingRefreshConfigLimits,
417    ) -> Result<bool, SearchArtifactError>,
418{
419    validate_limits(limits)?;
420    checkpoint()?;
421    let embeddings = project_dir.join(EMBEDDINGS_DIR);
422    ensure_owned_directory(&embeddings)?;
423    let _lock = ConfigWriterLock::acquire(&embeddings, limits.coordination, &mut checkpoint)?;
424    checkpoint()?;
425    cleanup_abandoned_temps(&embeddings, limits.coordination.cleanup_entries)?;
426    let path = embeddings.join(CONFIG_FILE);
427    let mut config = read_config_file(&path, limits)?;
428    if mutate(&mut config, limits)? {
429        checkpoint()?;
430        let bytes = config.to_canonical_json(limits)?;
431        checkpoint()?;
432        persist_synced_file(&path, &bytes)?;
433    }
434    Ok(config)
435}
436
437fn transaction_time_micros() -> i64 {
438    SystemTime::now()
439        .duration_since(UNIX_EPOCH)
440        .map_or(0, |duration| {
441            i64::try_from(duration.as_micros()).unwrap_or(i64::MAX)
442        })
443}
444
445fn apply_update(
446    config: &mut EmbeddingRefreshConfig,
447    update: EmbeddingRefreshConfigUpdate,
448    limits: EmbeddingRefreshConfigLimits,
449) -> Result<bool, SearchArtifactError> {
450    match update {
451        EmbeddingRefreshConfigUpdate::SetProjectPolicy(policy) => {
452            validate_project_policy(policy)?;
453            if config.project == policy {
454                Ok(false)
455            } else {
456                config.project = policy;
457                Ok(true)
458            }
459        }
460        EmbeddingRefreshConfigUpdate::SetSpacePolicy {
461            compatibility_id,
462            policy,
463        } => {
464            if let Some(policy) = policy {
465                validate_space_policy(policy)?;
466            }
467            if policy.is_none() && !config.spaces.contains_key(&compatibility_id) {
468                return Ok(false);
469            }
470            ensure_entry_capacity(config, compatibility_id, limits)?;
471            let current = config.spaces.entry(compatibility_id).or_default();
472            if current.policy == policy {
473                return Ok(false);
474            }
475            current.policy = policy;
476            if current.policy.is_none() && current.last_outcome.is_none() {
477                config.spaces.remove(&compatibility_id);
478            }
479            Ok(true)
480        }
481        EmbeddingRefreshConfigUpdate::RecordOutcome {
482            compatibility_id,
483            outcome,
484        } => {
485            validate_outcome(outcome)?;
486            ensure_entry_capacity(config, compatibility_id, limits)?;
487            let current = config.spaces.entry(compatibility_id).or_default();
488            if current.last_outcome == Some(outcome) {
489                return Ok(false);
490            }
491            if let Some(prior) = current.last_outcome {
492                validate_outcome_progress(prior, outcome)?;
493            }
494            current.last_outcome = Some(outcome);
495            Ok(true)
496        }
497        EmbeddingRefreshConfigUpdate::RemoveSpace { compatibility_id } => {
498            Ok(config.spaces.remove(&compatibility_id).is_some())
499        }
500    }
501}
502
503fn ensure_entry_capacity(
504    config: &EmbeddingRefreshConfig,
505    compatibility_id: EmbeddingCompatibilityId,
506    limits: EmbeddingRefreshConfigLimits,
507) -> Result<(), SearchArtifactError> {
508    if !config.spaces.contains_key(&compatibility_id) && config.spaces.len() >= limits.entries {
509        Err(exhausted(
510            "embedding_refresh_config_entries",
511            limits.entries,
512        ))
513    } else {
514        Ok(())
515    }
516}
517
518fn validate_project_policy(
519    policy: EmbeddingRefreshProjectPolicy,
520) -> Result<(), SearchArtifactError> {
521    validate_debounce(policy.debounce)?;
522    if !(1..=MAX_EMBEDDING_REFRESH_JOBS).contains(&policy.max_concurrent_jobs) {
523        return Err(invalid(
524            "embedding refresh max_concurrent_jobs",
525            "must be in 1..=2",
526        ));
527    }
528    Ok(())
529}
530
531fn validate_space_policy(policy: EmbeddingRefreshSpacePolicy) -> Result<(), SearchArtifactError> {
532    if policy.proactive.is_none() && policy.debounce.is_none() {
533        return Err(invalid(
534            "embedding refresh space policy",
535            "must override proactive or debounce",
536        ));
537    }
538    if let Some(debounce) = policy.debounce {
539        validate_debounce(debounce)?;
540    }
541    Ok(())
542}
543
544fn validate_debounce(debounce: Duration) -> Result<(), SearchArtifactError> {
545    if debounce.is_zero() || debounce > MAX_DEBOUNCE {
546        Err(invalid(
547            "embedding refresh debounce",
548            "must be in 1 millisecond..=1 hour",
549        ))
550    } else if !debounce.subsec_nanos().is_multiple_of(1_000_000) {
551        Err(invalid(
552            "embedding refresh debounce",
553            "must use whole milliseconds",
554        ))
555    } else {
556        Ok(())
557    }
558}
559
560fn validate_outcome(outcome: EmbeddingRefreshOutcomeRecord) -> Result<(), SearchArtifactError> {
561    if outcome.completed_at_micros < 0 {
562        Err(invalid(
563            "embedding refresh completed_at_micros",
564            "must be non-negative",
565        ))
566    } else {
567        Ok(())
568    }
569}
570
571fn validate_outcome_progress(
572    prior: EmbeddingRefreshOutcomeRecord,
573    next: EmbeddingRefreshOutcomeRecord,
574) -> Result<(), SearchArtifactError> {
575    if next.completed_at_micros < prior.completed_at_micros {
576        return Err(invalid(
577            "embedding refresh outcome",
578            "completion timestamp regressed",
579        ));
580    }
581    match next.graph_generation.cmp(&prior.graph_generation) {
582        std::cmp::Ordering::Less => Err(invalid(
583            "embedding refresh outcome",
584            "graph generation regressed",
585        )),
586        std::cmp::Ordering::Equal if next.source_fingerprint != prior.source_fingerprint => {
587            Err(invalid(
588                "embedding refresh outcome",
589                "source fingerprint conflicts at the same graph generation",
590            ))
591        }
592        std::cmp::Ordering::Equal | std::cmp::Ordering::Greater => Ok(()),
593    }
594}
595
596fn read_config_file(
597    path: &Path,
598    limits: EmbeddingRefreshConfigLimits,
599) -> Result<EmbeddingRefreshConfig, SearchArtifactError> {
600    if !path_exists(path)? {
601        return Ok(EmbeddingRefreshConfig::default());
602    }
603    ensure_regular_file(path)?;
604    let metadata = std::fs::metadata(path)
605        .map_err(|source| io("inspect embedding refresh config", path, source))?;
606    if metadata.len() > limits.metadata_bytes as u64 {
607        return Err(exhausted(
608            "embedding_refresh_config_bytes",
609            limits.metadata_bytes,
610        ));
611    }
612    let bytes =
613        std::fs::read(path).map_err(|source| io("read embedding refresh config", path, source))?;
614    let raw: RawDocument =
615        serde_json::from_slice(&bytes).map_err(|error| corrupt(path, error.to_string()))?;
616    if raw.config_version != u64::from(EMBEDDING_REFRESH_CONFIG_VERSION) {
617        return Err(SearchArtifactError::IncompatibleManifest {
618            path: path.to_path_buf(),
619            found: raw.config_version,
620            supported: EMBEDDING_REFRESH_CONFIG_VERSION,
621        });
622    }
623    if raw.spaces.len() > limits.entries {
624        return Err(exhausted(
625            "embedding_refresh_config_entries",
626            limits.entries,
627        ));
628    }
629    let project = parse_project(&raw.project).map_err(|error| corrupt(path, error.to_string()))?;
630    let mut spaces = BTreeMap::new();
631    let mut seen = BTreeSet::new();
632    for raw_space in raw.spaces {
633        let (compatibility_id, state) =
634            parse_space(raw_space).map_err(|error| corrupt(path, error.to_string()))?;
635        if !seen.insert(compatibility_id) {
636            return Err(corrupt(path, "duplicate refresh compatibility identity"));
637        }
638        spaces.insert(compatibility_id, state);
639    }
640    let config = EmbeddingRefreshConfig { project, spaces };
641    let material_bytes = material_bytes(&config)?;
642    if raw.checksum != checksum(&material_bytes) {
643        return Err(corrupt(path, "embedding refresh config checksum mismatch"));
644    }
645    let canonical = config
646        .to_canonical_json(limits)
647        .map_err(|error| corrupt(path, error.to_string()))?;
648    if canonical != bytes {
649        return Err(corrupt(
650            path,
651            "embedding refresh config bytes are not exact canonical JSON",
652        ));
653    }
654    Ok(config)
655}
656
657fn material_bytes(config: &EmbeddingRefreshConfig) -> Result<Vec<u8>, SearchArtifactError> {
658    let spaces = config
659        .spaces
660        .iter()
661        .map(|(compatibility_id, state)| wire_space(*compatibility_id, *state))
662        .collect::<Result<Vec<_>, _>>()?;
663    serde_json::to_vec(&WireMaterial {
664        config_version: EMBEDDING_REFRESH_CONFIG_VERSION,
665        project: wire_project(config.project)?,
666        spaces,
667    })
668    .map_err(|error| invalid("embedding refresh config", error.to_string()))
669}
670
671#[derive(Clone, Serialize)]
672struct WireProject {
673    proactive: bool,
674    debounce_millis: u64,
675    max_concurrent_jobs: usize,
676}
677
678#[derive(Clone, Serialize)]
679struct WireSpace {
680    compatibility_id: String,
681    policy: Option<WireSpacePolicy>,
682    last_outcome: Option<WireOutcome>,
683}
684
685#[derive(Clone, Serialize)]
686struct WireSpacePolicy {
687    proactive: Option<bool>,
688    debounce_millis: Option<u64>,
689}
690
691#[derive(Clone, Serialize)]
692struct WireOutcome {
693    status: &'static str,
694    failure_class: Option<&'static str>,
695    graph_generation: u64,
696    source_fingerprint: String,
697    completed_at_micros: i64,
698}
699
700#[derive(Serialize)]
701struct WireMaterial {
702    config_version: u32,
703    project: WireProject,
704    spaces: Vec<WireSpace>,
705}
706
707#[derive(Serialize)]
708struct WireDocument {
709    config_version: u32,
710    project: WireProject,
711    spaces: Vec<WireSpace>,
712    checksum: String,
713}
714
715#[derive(Deserialize)]
716#[serde(deny_unknown_fields)]
717struct RawDocument {
718    config_version: u64,
719    project: RawProject,
720    spaces: Vec<RawSpace>,
721    checksum: String,
722}
723
724#[derive(Deserialize)]
725#[serde(deny_unknown_fields)]
726struct RawProject {
727    proactive: bool,
728    debounce_millis: u64,
729    max_concurrent_jobs: usize,
730}
731
732#[derive(Deserialize)]
733#[serde(deny_unknown_fields)]
734struct RawSpace {
735    compatibility_id: String,
736    policy: Option<RawSpacePolicy>,
737    last_outcome: Option<RawOutcome>,
738}
739
740#[derive(Deserialize)]
741#[serde(deny_unknown_fields)]
742struct RawSpacePolicy {
743    proactive: Option<bool>,
744    debounce_millis: Option<u64>,
745}
746
747#[derive(Deserialize)]
748#[serde(deny_unknown_fields)]
749struct RawOutcome {
750    status: String,
751    failure_class: Option<String>,
752    graph_generation: u64,
753    source_fingerprint: String,
754    completed_at_micros: i64,
755}
756
757fn wire_project(policy: EmbeddingRefreshProjectPolicy) -> Result<WireProject, SearchArtifactError> {
758    validate_project_policy(policy)?;
759    Ok(WireProject {
760        proactive: policy.proactive,
761        debounce_millis: duration_millis(policy.debounce)?,
762        max_concurrent_jobs: policy.max_concurrent_jobs,
763    })
764}
765
766fn wire_space(
767    compatibility_id: EmbeddingCompatibilityId,
768    state: StoredSpaceState,
769) -> Result<WireSpace, SearchArtifactError> {
770    if state.policy.is_none() && state.last_outcome.is_none() {
771        return Err(invalid(
772            "embedding refresh space state",
773            "must retain a policy or outcome",
774        ));
775    }
776    let policy = state
777        .policy
778        .map(|policy| {
779            validate_space_policy(policy)?;
780            Ok(WireSpacePolicy {
781                proactive: policy.proactive,
782                debounce_millis: policy.debounce.map(duration_millis).transpose()?,
783            })
784        })
785        .transpose()?;
786    let last_outcome = state
787        .last_outcome
788        .map(|outcome| {
789            validate_outcome(outcome)?;
790            let (status, failure_class) = outcome_tokens(outcome.status);
791            Ok(WireOutcome {
792                status,
793                failure_class,
794                graph_generation: outcome.graph_generation,
795                source_fingerprint: outcome.source_fingerprint.to_hex(),
796                completed_at_micros: outcome.completed_at_micros,
797            })
798        })
799        .transpose()?;
800    Ok(WireSpace {
801        compatibility_id: compatibility_id.to_hex(),
802        policy,
803        last_outcome,
804    })
805}
806
807fn parse_project(raw: &RawProject) -> Result<EmbeddingRefreshProjectPolicy, SearchArtifactError> {
808    let policy = EmbeddingRefreshProjectPolicy {
809        proactive: raw.proactive,
810        debounce: Duration::from_millis(raw.debounce_millis),
811        max_concurrent_jobs: raw.max_concurrent_jobs,
812    };
813    validate_project_policy(policy)?;
814    Ok(policy)
815}
816
817fn parse_space(
818    raw: RawSpace,
819) -> Result<(EmbeddingCompatibilityId, StoredSpaceState), SearchArtifactError> {
820    let compatibility_id = EmbeddingCompatibilityId::from_hex(&raw.compatibility_id)?;
821    let policy = raw
822        .policy
823        .map(|raw| {
824            let policy = EmbeddingRefreshSpacePolicy {
825                proactive: raw.proactive,
826                debounce: raw.debounce_millis.map(Duration::from_millis),
827            };
828            validate_space_policy(policy)?;
829            Ok(policy)
830        })
831        .transpose()?;
832    let last_outcome = raw.last_outcome.as_ref().map(parse_outcome).transpose()?;
833    if policy.is_none() && last_outcome.is_none() {
834        return Err(invalid(
835            "embedding refresh space state",
836            "must retain a policy or outcome",
837        ));
838    }
839    Ok((
840        compatibility_id,
841        StoredSpaceState {
842            policy,
843            last_outcome,
844        },
845    ))
846}
847
848fn parse_outcome(raw: &RawOutcome) -> Result<EmbeddingRefreshOutcomeRecord, SearchArtifactError> {
849    let status = match (raw.status.as_str(), raw.failure_class.as_deref()) {
850        ("succeeded", None) => EmbeddingRefreshOutcomeStatus::Succeeded,
851        ("cancelled", None) => EmbeddingRefreshOutcomeStatus::Cancelled,
852        ("failed", Some(class)) => {
853            EmbeddingRefreshOutcomeStatus::Failed(parse_failure_class(class)?)
854        }
855        _ => {
856            return Err(invalid(
857                "embedding refresh outcome status",
858                "status and failure_class are inconsistent",
859            ));
860        }
861    };
862    let outcome = EmbeddingRefreshOutcomeRecord {
863        status,
864        graph_generation: raw.graph_generation,
865        source_fingerprint: EmbeddingSourceFingerprint::from_hex(&raw.source_fingerprint)?,
866        completed_at_micros: raw.completed_at_micros,
867    };
868    validate_outcome(outcome)?;
869    Ok(outcome)
870}
871
872fn outcome_tokens(status: EmbeddingRefreshOutcomeStatus) -> (&'static str, Option<&'static str>) {
873    match status {
874        EmbeddingRefreshOutcomeStatus::Succeeded => ("succeeded", None),
875        EmbeddingRefreshOutcomeStatus::Cancelled => ("cancelled", None),
876        EmbeddingRefreshOutcomeStatus::Failed(class) => ("failed", Some(failure_token(class))),
877    }
878}
879
880fn failure_token(class: EmbeddingRefreshFailureClass) -> &'static str {
881    match class {
882        EmbeddingRefreshFailureClass::Provider => "provider",
883        EmbeddingRefreshFailureClass::Validation => "validation",
884        EmbeddingRefreshFailureClass::ResourceExhausted => "resource_exhausted",
885        EmbeddingRefreshFailureClass::Storage => "storage",
886        EmbeddingRefreshFailureClass::ConcurrentMutation => "concurrent_mutation",
887        EmbeddingRefreshFailureClass::Incompatible => "incompatible",
888        EmbeddingRefreshFailureClass::Corrupt => "corrupt",
889        EmbeddingRefreshFailureClass::Unavailable => "unavailable",
890    }
891}
892
893fn parse_failure_class(value: &str) -> Result<EmbeddingRefreshFailureClass, SearchArtifactError> {
894    match value {
895        "provider" => Ok(EmbeddingRefreshFailureClass::Provider),
896        "validation" => Ok(EmbeddingRefreshFailureClass::Validation),
897        "resource_exhausted" => Ok(EmbeddingRefreshFailureClass::ResourceExhausted),
898        "storage" => Ok(EmbeddingRefreshFailureClass::Storage),
899        "concurrent_mutation" => Ok(EmbeddingRefreshFailureClass::ConcurrentMutation),
900        "incompatible" => Ok(EmbeddingRefreshFailureClass::Incompatible),
901        "corrupt" => Ok(EmbeddingRefreshFailureClass::Corrupt),
902        "unavailable" => Ok(EmbeddingRefreshFailureClass::Unavailable),
903        _ => Err(invalid(
904            "embedding refresh failure_class",
905            "is not a supported token",
906        )),
907    }
908}
909
910fn duration_millis(duration: Duration) -> Result<u64, SearchArtifactError> {
911    validate_debounce(duration)?;
912    u64::try_from(duration.as_millis())
913        .map_err(|_| exhausted("embedding_refresh_debounce_millis", usize::MAX))
914}
915
916fn checksum(material: &[u8]) -> String {
917    let mut hasher = Sha256::new();
918    hasher.update(CHECKSUM_DOMAIN);
919    hasher.update(material);
920    format!("{:x}", hasher.finalize())
921}
922
923struct ConfigWriterLock {
924    file: File,
925}
926
927impl ConfigWriterLock {
928    fn acquire<C>(
929        embeddings: &Path,
930        limits: SearchCoordinationLimits,
931        checkpoint: &mut C,
932    ) -> Result<Self, SearchArtifactError>
933    where
934        C: FnMut() -> Result<(), SearchArtifactError>,
935    {
936        let path = embeddings.join(CONFIG_LOCK);
937        if path_exists(&path)? {
938            ensure_regular_file(&path)?;
939        }
940        let file = OpenOptions::new()
941            .read(true)
942            .write(true)
943            .create(true)
944            .truncate(false)
945            .open(&path)
946            .map_err(|source| SearchArtifactError::Lock {
947                path: path.clone(),
948                reason: source.to_string(),
949            })?;
950        let started = Instant::now();
951        loop {
952            match file.try_lock() {
953                Ok(()) => return Ok(Self { file }),
954                Err(std::fs::TryLockError::WouldBlock) => {
955                    checkpoint()?;
956                    if started.elapsed() >= limits.lock_timeout {
957                        return Err(SearchArtifactError::Lock {
958                            path,
959                            reason: format!(
960                                "timed out after {} ms",
961                                limits.lock_timeout.as_millis()
962                            ),
963                        });
964                    }
965                    std::thread::sleep(limits.lock_poll_interval);
966                }
967                Err(std::fs::TryLockError::Error(source)) => {
968                    return Err(SearchArtifactError::Lock {
969                        path,
970                        reason: source.to_string(),
971                    });
972                }
973            }
974        }
975    }
976}
977
978impl Drop for ConfigWriterLock {
979    fn drop(&mut self) {
980        let _ = self.file.unlock();
981    }
982}
983
984fn cleanup_abandoned_temps(
985    embeddings: &Path,
986    max_entries: usize,
987) -> Result<usize, SearchArtifactError> {
988    if max_entries == 0 {
989        return Err(invalid(
990            "embedding refresh cleanup_entries",
991            "must be non-zero",
992        ));
993    }
994    let mut inspected = 0_usize;
995    let mut removed = 0_usize;
996    for entry in std::fs::read_dir(embeddings)
997        .map_err(|source| io("read embedding refresh directory", embeddings, source))?
998    {
999        inspected = inspected
1000            .checked_add(1)
1001            .ok_or_else(|| exhausted("embedding_refresh_cleanup_entries", max_entries))?;
1002        if inspected > max_entries {
1003            return Err(exhausted("embedding_refresh_cleanup_entries", max_entries));
1004        }
1005        let entry =
1006            entry.map_err(|source| io("read embedding refresh entry", embeddings, source))?;
1007        let file_type = entry
1008            .file_type()
1009            .map_err(|source| io("inspect embedding refresh entry", &entry.path(), source))?;
1010        if file_type.is_file() && is_refresh_temp_name(&entry.file_name()) {
1011            std::fs::remove_file(entry.path())
1012                .map_err(|source| io("remove embedding refresh temp", &entry.path(), source))?;
1013            removed += 1;
1014        }
1015    }
1016    Ok(removed)
1017}
1018
1019fn is_refresh_temp_name(name: &std::ffi::OsStr) -> bool {
1020    let Some(name) = name.to_str() else {
1021        return false;
1022    };
1023    let Some(random) = name
1024        .strip_prefix(".refresh.json.")
1025        .and_then(|name| name.strip_suffix(".tmp"))
1026    else {
1027        return false;
1028    };
1029    !random.is_empty() && random.bytes().all(|byte| byte.is_ascii_alphanumeric())
1030}
1031
1032fn validate_limits(limits: EmbeddingRefreshConfigLimits) -> Result<(), SearchArtifactError> {
1033    if limits.metadata_bytes == 0 || limits.entries == 0 || limits.coordination.cleanup_entries == 0
1034    {
1035        Err(invalid(
1036            "embedding refresh config limits",
1037            "must be non-zero",
1038        ))
1039    } else {
1040        Ok(())
1041    }
1042}
1043
1044fn config_path(project_dir: &Path) -> PathBuf {
1045    project_dir.join(EMBEDDINGS_DIR).join(CONFIG_FILE)
1046}
1047
1048fn ensure_owned_directory(path: &Path) -> Result<(), SearchArtifactError> {
1049    if path_exists(path)? {
1050        let metadata = std::fs::symlink_metadata(path)
1051            .map_err(|source| io("inspect embedding refresh directory", path, source))?;
1052        if metadata.file_type().is_symlink() || !metadata.is_dir() {
1053            return Err(corrupt(path, "expected an owned directory"));
1054        }
1055        return Ok(());
1056    }
1057    std::fs::create_dir_all(path)
1058        .map_err(|source| io("create embedding refresh directory", path, source))?;
1059    sync_directory(path.parent().unwrap_or(path))
1060}
1061
1062fn ensure_regular_file(path: &Path) -> Result<(), SearchArtifactError> {
1063    let metadata = std::fs::symlink_metadata(path)
1064        .map_err(|source| io("inspect embedding refresh config", path, source))?;
1065    if metadata.file_type().is_symlink() || !metadata.is_file() {
1066        return Err(corrupt(path, "expected a regular file"));
1067    }
1068    Ok(())
1069}
1070
1071fn persist_synced_file(path: &Path, bytes: &[u8]) -> Result<(), SearchArtifactError> {
1072    let parent = path
1073        .parent()
1074        .ok_or_else(|| invalid("embedding refresh config", "path has no parent"))?;
1075    let mut temp = tempfile::Builder::new()
1076        .prefix(".refresh.json.")
1077        .suffix(".tmp")
1078        .tempfile_in(parent)
1079        .map_err(|source| io("create embedding refresh temp", path, source))?;
1080    temp.write_all(bytes)
1081        .map_err(|source| io("write embedding refresh temp", path, source))?;
1082    temp.as_file()
1083        .sync_all()
1084        .map_err(|source| io("sync embedding refresh temp", path, source))?;
1085    temp.persist(path)
1086        .map_err(|error| io("publish embedding refresh config", path, error.error))?;
1087    sync_directory(parent)
1088}
1089
1090#[cfg(unix)]
1091fn sync_directory(path: &Path) -> Result<(), SearchArtifactError> {
1092    File::open(path)
1093        .and_then(|directory| directory.sync_all())
1094        .map_err(|source| io("sync embedding refresh directory", path, source))
1095}
1096
1097#[cfg(not(unix))]
1098fn sync_directory(_path: &Path) -> Result<(), SearchArtifactError> {
1099    Ok(())
1100}
1101
1102fn path_exists(path: &Path) -> Result<bool, SearchArtifactError> {
1103    match std::fs::symlink_metadata(path) {
1104        Ok(_) => Ok(true),
1105        Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(false),
1106        Err(source) => Err(io("inspect embedding refresh path", path, source)),
1107    }
1108}
1109
1110fn invalid(field: &'static str, reason: impl Into<String>) -> SearchArtifactError {
1111    SearchArtifactError::InvalidSelector {
1112        field,
1113        reason: reason.into(),
1114    }
1115}
1116
1117fn corrupt(path: &Path, reason: impl Into<String>) -> SearchArtifactError {
1118    SearchArtifactError::CorruptManifest {
1119        path: path.to_path_buf(),
1120        reason: reason.into(),
1121    }
1122}
1123
1124fn exhausted(resource: &'static str, limit: usize) -> SearchArtifactError {
1125    SearchArtifactError::ResourceExhausted {
1126        resource,
1127        limit: limit as u64,
1128    }
1129}
1130
1131fn io(operation: &'static str, path: &Path, source: std::io::Error) -> SearchArtifactError {
1132    SearchArtifactError::Io {
1133        operation,
1134        path: path.to_path_buf(),
1135        source,
1136    }
1137}
1138
1139#[cfg(test)]
1140mod tests {
1141    use std::cell::Cell;
1142    use std::sync::{Arc, Barrier};
1143
1144    use super::*;
1145
1146    fn id(value: u8) -> EmbeddingCompatibilityId {
1147        EmbeddingCompatibilityId::from_hex(&format!("{value:02x}").repeat(32)).unwrap()
1148    }
1149
1150    fn fingerprint(value: u8) -> EmbeddingSourceFingerprint {
1151        EmbeddingSourceFingerprint::from_hex(&format!("{value:02x}").repeat(32)).unwrap()
1152    }
1153
1154    fn read(project: &Path) -> EmbeddingRefreshConfig {
1155        read_embedding_refresh_config(project, EmbeddingRefreshConfigLimits::default(), || Ok(()))
1156            .unwrap()
1157    }
1158
1159    fn update(project: &Path, mutation: EmbeddingRefreshConfigUpdate) -> EmbeddingRefreshConfig {
1160        update_embedding_refresh_config(
1161            project,
1162            mutation,
1163            EmbeddingRefreshConfigLimits::default(),
1164            || Ok(()),
1165        )
1166        .unwrap()
1167    }
1168
1169    fn attempt(status: EmbeddingRefreshOutcomeStatus) -> EmbeddingRefreshOutcomeAttempt {
1170        EmbeddingRefreshOutcomeAttempt {
1171            status,
1172            graph_generation: 1,
1173            source_fingerprint: fingerprint(1),
1174        }
1175    }
1176
1177    #[test]
1178    fn locked_outcome_stamping_is_monotonic_under_same_lineage_contention() {
1179        let project = tempfile::tempdir().unwrap();
1180        let project = Arc::new(project.path().to_path_buf());
1181        let barrier = Arc::new(Barrier::new(2));
1182        let mut threads = Vec::new();
1183        for status in [
1184            EmbeddingRefreshOutcomeStatus::Succeeded,
1185            EmbeddingRefreshOutcomeStatus::Cancelled,
1186        ] {
1187            let project = Arc::clone(&project);
1188            let barrier = Arc::clone(&barrier);
1189            threads.push(std::thread::spawn(move || {
1190                barrier.wait();
1191                record_embedding_refresh_outcome_with_clock(
1192                    &project,
1193                    id(1),
1194                    attempt(status),
1195                    EmbeddingRefreshConfigLimits::default(),
1196                    || 5,
1197                    || Ok(()),
1198                )
1199                .unwrap()
1200                .spaces()
1201                .into_iter()
1202                .find(|state| state.compatibility_id == id(1))
1203                .unwrap()
1204                .last_outcome
1205                .unwrap()
1206                .completed_at_micros
1207            }));
1208        }
1209        let mut completed = threads
1210            .into_iter()
1211            .map(|thread| thread.join().unwrap())
1212            .collect::<Vec<_>>();
1213        completed.sort_unstable();
1214        assert_eq!(completed, [5, 6]);
1215        assert_eq!(
1216            read(&project)
1217                .spaces()
1218                .into_iter()
1219                .find(|state| state.compatibility_id == id(1))
1220                .unwrap()
1221                .last_outcome
1222                .unwrap()
1223                .completed_at_micros,
1224            6
1225        );
1226    }
1227
1228    #[test]
1229    fn locked_outcome_timestamp_overflow_is_structured_and_atomic() {
1230        let project = tempfile::tempdir().unwrap();
1231        update(
1232            project.path(),
1233            EmbeddingRefreshConfigUpdate::RecordOutcome {
1234                compatibility_id: id(1),
1235                outcome: EmbeddingRefreshOutcomeRecord {
1236                    completed_at_micros: i64::MAX,
1237                    ..outcome(1, 1, 1)
1238                },
1239            },
1240        );
1241        let stable = std::fs::read(config_path(project.path())).unwrap();
1242        assert!(matches!(
1243            record_embedding_refresh_outcome_with_clock(
1244                project.path(),
1245                id(1),
1246                attempt(EmbeddingRefreshOutcomeStatus::Succeeded),
1247                EmbeddingRefreshConfigLimits::default(),
1248                || 1,
1249                || Ok(())
1250            ),
1251            Err(SearchArtifactError::ResourceExhausted {
1252                resource: "embedding_refresh_outcome_timestamp",
1253                ..
1254            })
1255        ));
1256        assert_eq!(std::fs::read(config_path(project.path())).unwrap(), stable);
1257    }
1258
1259    fn outcome(
1260        generation: u64,
1261        fingerprint_value: u8,
1262        completed_at_micros: i64,
1263    ) -> EmbeddingRefreshOutcomeRecord {
1264        EmbeddingRefreshOutcomeRecord {
1265            status: EmbeddingRefreshOutcomeStatus::Succeeded,
1266            graph_generation: generation,
1267            source_fingerprint: fingerprint(fingerprint_value),
1268            completed_at_micros,
1269        }
1270    }
1271
1272    #[test]
1273    fn missing_config_returns_defaults_without_creating_files() {
1274        let project = tempfile::tempdir().unwrap();
1275        let config = read(project.path());
1276        assert_eq!(
1277            config.project_policy(),
1278            EmbeddingRefreshProjectPolicy::default()
1279        );
1280        assert!(config.is_empty());
1281        assert!(!project.path().join(EMBEDDINGS_DIR).exists());
1282    }
1283
1284    #[test]
1285    fn project_space_and_outcome_state_round_trip_canonically() {
1286        let project = tempfile::tempdir().unwrap();
1287        let compatibility_id = id(2);
1288        update(
1289            project.path(),
1290            EmbeddingRefreshConfigUpdate::SetProjectPolicy(EmbeddingRefreshProjectPolicy {
1291                proactive: false,
1292                debounce: Duration::from_millis(750),
1293                max_concurrent_jobs: 1,
1294            }),
1295        );
1296        update(
1297            project.path(),
1298            EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1299                compatibility_id,
1300                policy: Some(EmbeddingRefreshSpacePolicy {
1301                    proactive: Some(true),
1302                    debounce: Some(Duration::from_secs(2)),
1303                }),
1304            },
1305        );
1306        let config = update(
1307            project.path(),
1308            EmbeddingRefreshConfigUpdate::RecordOutcome {
1309                compatibility_id,
1310                outcome: EmbeddingRefreshOutcomeRecord {
1311                    status: EmbeddingRefreshOutcomeStatus::Failed(
1312                        EmbeddingRefreshFailureClass::Provider,
1313                    ),
1314                    graph_generation: 9,
1315                    source_fingerprint: fingerprint(3),
1316                    completed_at_micros: 10,
1317                },
1318            },
1319        );
1320        assert_eq!(config, read(project.path()));
1321        assert_eq!(
1322            config.resolved_policy(compatibility_id),
1323            ResolvedEmbeddingRefreshPolicy {
1324                proactive: true,
1325                debounce: Duration::from_secs(2),
1326                max_concurrent_jobs: 1,
1327            }
1328        );
1329        let state = config.spaces()[0];
1330        assert_eq!(state.compatibility_id, compatibility_id);
1331        assert!(matches!(
1332            state.last_outcome.unwrap().status,
1333            EmbeddingRefreshOutcomeStatus::Failed(EmbeddingRefreshFailureClass::Provider)
1334        ));
1335
1336        let path = config_path(project.path());
1337        let bytes = std::fs::read(&path).unwrap();
1338        assert_eq!(
1339            bytes,
1340            config
1341                .to_canonical_json(EmbeddingRefreshConfigLimits::default())
1342                .unwrap()
1343        );
1344        let text = String::from_utf8(bytes).unwrap();
1345        assert!(text.contains("\"checksum\""));
1346        for forbidden in [
1347            "alias",
1348            "credential",
1349            "payload",
1350            "vector",
1351            "source_text",
1352            "confidence",
1353        ] {
1354            assert!(!text.contains(forbidden));
1355        }
1356    }
1357
1358    #[test]
1359    fn idempotence_removal_limits_and_non_monotonic_outcomes_are_structured() {
1360        let project = tempfile::tempdir().unwrap();
1361        let compatibility_id = id(1);
1362        update(
1363            project.path(),
1364            EmbeddingRefreshConfigUpdate::RecordOutcome {
1365                compatibility_id,
1366                outcome: outcome(2, 2, 20),
1367            },
1368        );
1369        let path = config_path(project.path());
1370        let stable = std::fs::read(&path).unwrap();
1371        update(
1372            project.path(),
1373            EmbeddingRefreshConfigUpdate::RecordOutcome {
1374                compatibility_id,
1375                outcome: outcome(2, 2, 20),
1376            },
1377        );
1378        assert_eq!(std::fs::read(&path).unwrap(), stable);
1379
1380        for invalid_outcome in [outcome(1, 1, 21), outcome(2, 3, 21), outcome(3, 3, 19)] {
1381            assert!(matches!(
1382                update_embedding_refresh_config(
1383                    project.path(),
1384                    EmbeddingRefreshConfigUpdate::RecordOutcome {
1385                        compatibility_id,
1386                        outcome: invalid_outcome,
1387                    },
1388                    EmbeddingRefreshConfigLimits::default(),
1389                    || Ok(())
1390                ),
1391                Err(SearchArtifactError::InvalidSelector { .. })
1392            ));
1393            assert_eq!(std::fs::read(&path).unwrap(), stable);
1394        }
1395
1396        assert!(matches!(
1397            update_embedding_refresh_config(
1398                project.path(),
1399                EmbeddingRefreshConfigUpdate::SetProjectPolicy(EmbeddingRefreshProjectPolicy {
1400                    max_concurrent_jobs: 3,
1401                    ..EmbeddingRefreshProjectPolicy::default()
1402                }),
1403                EmbeddingRefreshConfigLimits::default(),
1404                || Ok(())
1405            ),
1406            Err(SearchArtifactError::InvalidSelector { .. })
1407        ));
1408        assert!(matches!(
1409            update_embedding_refresh_config(
1410                project.path(),
1411                EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1412                    compatibility_id: id(9),
1413                    policy: Some(EmbeddingRefreshSpacePolicy::default()),
1414                },
1415                EmbeddingRefreshConfigLimits::default(),
1416                || Ok(())
1417            ),
1418            Err(SearchArtifactError::InvalidSelector { .. })
1419        ));
1420
1421        let config = update(
1422            project.path(),
1423            EmbeddingRefreshConfigUpdate::RemoveSpace { compatibility_id },
1424        );
1425        assert!(config.is_empty());
1426    }
1427
1428    #[test]
1429    fn corruption_version_cancellation_and_interrupted_temp_fail_closed() {
1430        let project = tempfile::tempdir().unwrap();
1431        let compatibility_id = id(1);
1432        update(
1433            project.path(),
1434            EmbeddingRefreshConfigUpdate::RecordOutcome {
1435                compatibility_id,
1436                outcome: outcome(1, 1, 1),
1437            },
1438        );
1439        let path = config_path(project.path());
1440        let stable = std::fs::read(&path).unwrap();
1441        let calls = Cell::new(0_u8);
1442        assert!(matches!(
1443            update_embedding_refresh_config(
1444                project.path(),
1445                EmbeddingRefreshConfigUpdate::SetProjectPolicy(EmbeddingRefreshProjectPolicy {
1446                    proactive: false,
1447                    ..EmbeddingRefreshProjectPolicy::default()
1448                }),
1449                EmbeddingRefreshConfigLimits::default(),
1450                || {
1451                    let next = calls.get() + 1;
1452                    calls.set(next);
1453                    if next >= 4 {
1454                        Err(SearchArtifactError::Cancelled)
1455                    } else {
1456                        Ok(())
1457                    }
1458                }
1459            ),
1460            Err(SearchArtifactError::Cancelled)
1461        ));
1462        assert_eq!(std::fs::read(&path).unwrap(), stable);
1463
1464        let interrupted = project
1465            .path()
1466            .join(EMBEDDINGS_DIR)
1467            .join(".refresh.json.ABC123.tmp");
1468        std::fs::write(&interrupted, b"partial").unwrap();
1469        assert_eq!(read(project.path()).spaces().len(), 1);
1470        update(
1471            project.path(),
1472            EmbeddingRefreshConfigUpdate::SetProjectPolicy(EmbeddingRefreshProjectPolicy {
1473                proactive: false,
1474                ..EmbeddingRefreshProjectPolicy::default()
1475            }),
1476        );
1477        assert!(!interrupted.exists());
1478
1479        let valid = std::fs::read_to_string(&path).unwrap();
1480        let damaged = valid.replacen("\"checksum\":\"", "\"checksum\":\"0", 1);
1481        std::fs::write(&path, damaged).unwrap();
1482        assert!(matches!(
1483            read_embedding_refresh_config(
1484                project.path(),
1485                EmbeddingRefreshConfigLimits::default(),
1486                || Ok(())
1487            ),
1488            Err(SearchArtifactError::CorruptManifest { .. })
1489        ));
1490
1491        std::fs::write(
1492            &path,
1493            b"{\"config_version\":1,\"project\":{\"proactive\":true,\"debounce_millis\":500,\"max_concurrent_jobs\":2},\"spaces\":[{\"compatibility_id\":\"invalid\",\"policy\":{\"proactive\":false,\"debounce_millis\":null},\"last_outcome\":null}],\"checksum\":\"bad\"}",
1494        )
1495        .unwrap();
1496        assert!(matches!(
1497            read_embedding_refresh_config(
1498                project.path(),
1499                EmbeddingRefreshConfigLimits::default(),
1500                || Ok(())
1501            ),
1502            Err(SearchArtifactError::CorruptManifest { .. })
1503        ));
1504
1505        std::fs::write(
1506            &path,
1507            b"{\"config_version\":2,\"project\":{\"proactive\":true,\"debounce_millis\":500,\"max_concurrent_jobs\":2},\"spaces\":[],\"checksum\":\"bad\"}",
1508        )
1509        .unwrap();
1510        assert!(matches!(
1511            read_embedding_refresh_config(
1512                project.path(),
1513                EmbeddingRefreshConfigLimits::default(),
1514                || Ok(())
1515            ),
1516            Err(SearchArtifactError::IncompatibleManifest { .. })
1517        ));
1518    }
1519
1520    #[test]
1521    fn concurrent_writers_serialize_without_lost_lineages() {
1522        use std::sync::{Arc, Barrier};
1523
1524        let project = tempfile::tempdir().unwrap();
1525        let project_path = Arc::new(project.path().to_path_buf());
1526        let barrier = Arc::new(Barrier::new(3));
1527        let workers = [id(1), id(2)].map(|compatibility_id| {
1528            let project_path = Arc::clone(&project_path);
1529            let barrier = Arc::clone(&barrier);
1530            std::thread::spawn(move || {
1531                barrier.wait();
1532                update_embedding_refresh_config(
1533                    &project_path,
1534                    EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1535                        compatibility_id,
1536                        policy: Some(EmbeddingRefreshSpacePolicy {
1537                            proactive: Some(false),
1538                            debounce: None,
1539                        }),
1540                    },
1541                    EmbeddingRefreshConfigLimits::default(),
1542                    || Ok(()),
1543                )
1544                .unwrap();
1545            })
1546        });
1547        barrier.wait();
1548        for worker in workers {
1549            worker.join().unwrap();
1550        }
1551        let config = read(project.path());
1552        assert_eq!(
1553            config
1554                .spaces()
1555                .iter()
1556                .map(|state| state.compatibility_id)
1557                .collect::<Vec<_>>(),
1558            [id(1), id(2)]
1559        );
1560    }
1561
1562    #[test]
1563    fn entry_and_metadata_bounds_do_not_publish_partial_state() {
1564        let project = tempfile::tempdir().unwrap();
1565        let limits = EmbeddingRefreshConfigLimits {
1566            entries: 1,
1567            ..EmbeddingRefreshConfigLimits::default()
1568        };
1569        update_embedding_refresh_config(
1570            project.path(),
1571            EmbeddingRefreshConfigUpdate::RecordOutcome {
1572                compatibility_id: id(1),
1573                outcome: outcome(1, 1, 1),
1574            },
1575            limits,
1576            || Ok(()),
1577        )
1578        .unwrap();
1579        let path = config_path(project.path());
1580        let stable = std::fs::read(&path).unwrap();
1581        assert!(matches!(
1582            update_embedding_refresh_config(
1583                project.path(),
1584                EmbeddingRefreshConfigUpdate::RecordOutcome {
1585                    compatibility_id: id(2),
1586                    outcome: outcome(2, 2, 2),
1587                },
1588                limits,
1589                || Ok(())
1590            ),
1591            Err(SearchArtifactError::ResourceExhausted { .. })
1592        ));
1593        assert_eq!(std::fs::read(&path).unwrap(), stable);
1594
1595        let tiny = EmbeddingRefreshConfigLimits {
1596            metadata_bytes: 1,
1597            ..EmbeddingRefreshConfigLimits::default()
1598        };
1599        assert!(matches!(
1600            read_embedding_refresh_config(project.path(), tiny, || Ok(())),
1601            Err(SearchArtifactError::ResourceExhausted { .. })
1602        ));
1603    }
1604
1605    #[cfg(not(unix))]
1606    #[test]
1607    fn directory_sync_is_a_supported_noop() {
1608        let directory = tempfile::tempdir().unwrap();
1609        sync_directory(directory.path()).unwrap();
1610    }
1611}