1use 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
17pub const EMBEDDING_REFRESH_CONFIG_VERSION: u32 = 1;
19pub const DEFAULT_EMBEDDING_REFRESH_DEBOUNCE: Duration = Duration::from_millis(500);
21pub const MAX_EMBEDDING_REFRESH_JOBS: usize = 2;
23pub const MAX_EMBEDDING_REFRESH_CONFIG_BYTES: usize = 256 * 1024;
25pub 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#[derive(Clone, Copy, Debug)]
36pub struct EmbeddingRefreshConfigLimits {
37 pub metadata_bytes: usize,
39 pub entries: usize,
41 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#[derive(Clone, Copy, Debug, PartialEq, Eq)]
57pub struct EmbeddingRefreshProjectPolicy {
58 pub proactive: bool,
60 pub debounce: Duration,
62 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#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
78pub struct EmbeddingRefreshSpacePolicy {
79 pub proactive: Option<bool>,
81 pub debounce: Option<Duration>,
83}
84
85#[derive(Clone, Copy, Debug, PartialEq, Eq)]
87pub struct ResolvedEmbeddingRefreshPolicy {
88 pub proactive: bool,
90 pub debounce: Duration,
92 pub max_concurrent_jobs: usize,
94}
95
96#[derive(Clone, Copy, Debug, PartialEq, Eq)]
98pub enum EmbeddingRefreshFailureClass {
99 Provider,
101 Validation,
103 ResourceExhausted,
105 Storage,
107 ConcurrentMutation,
109 Incompatible,
111 Corrupt,
113 Unavailable,
115}
116
117#[derive(Clone, Copy, Debug, PartialEq, Eq)]
119pub enum EmbeddingRefreshOutcomeStatus {
120 Succeeded,
122 Cancelled,
124 Failed(EmbeddingRefreshFailureClass),
126}
127
128#[derive(Clone, Copy, Debug, PartialEq, Eq)]
130pub struct EmbeddingRefreshOutcomeRecord {
131 pub status: EmbeddingRefreshOutcomeStatus,
133 pub graph_generation: u64,
135 pub source_fingerprint: EmbeddingSourceFingerprint,
137 pub completed_at_micros: i64,
139}
140
141#[derive(Clone, Copy, Debug, PartialEq, Eq)]
143pub struct EmbeddingRefreshOutcomeAttempt {
144 pub status: EmbeddingRefreshOutcomeStatus,
146 pub graph_generation: u64,
148 pub source_fingerprint: EmbeddingSourceFingerprint,
150}
151
152#[derive(Clone, Copy, Debug, PartialEq, Eq)]
154pub struct EmbeddingRefreshSpaceState {
155 pub compatibility_id: EmbeddingCompatibilityId,
157 pub policy: Option<EmbeddingRefreshSpacePolicy>,
159 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#[derive(Clone, Debug, Default, PartialEq, Eq)]
171pub struct EmbeddingRefreshConfig {
172 project: EmbeddingRefreshProjectPolicy,
173 spaces: BTreeMap<EmbeddingCompatibilityId, StoredSpaceState>,
174}
175
176impl EmbeddingRefreshConfig {
177 #[must_use]
179 pub const fn project_policy(&self) -> EmbeddingRefreshProjectPolicy {
180 self.project
181 }
182
183 #[must_use]
185 pub fn len(&self) -> usize {
186 self.spaces.len()
187 }
188
189 #[must_use]
191 pub fn is_empty(&self) -> bool {
192 self.spaces.is_empty()
193 }
194
195 #[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 #[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#[derive(Clone, Copy, Debug)]
272pub enum EmbeddingRefreshConfigUpdate {
273 SetProjectPolicy(EmbeddingRefreshProjectPolicy),
275 SetSpacePolicy {
277 compatibility_id: EmbeddingCompatibilityId,
279 policy: Option<EmbeddingRefreshSpacePolicy>,
281 },
282 RecordOutcome {
284 compatibility_id: EmbeddingCompatibilityId,
286 outcome: EmbeddingRefreshOutcomeRecord,
288 },
289 RemoveSpace {
291 compatibility_id: EmbeddingCompatibilityId,
293 },
294}
295
296pub 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
315pub 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
336pub 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
1612 #[test]
1613 fn limits_and_cleanup_bounds_are_exact_and_side_effect_free() {
1614 for limits in [
1615 EmbeddingRefreshConfigLimits {
1616 metadata_bytes: 0,
1617 ..EmbeddingRefreshConfigLimits::default()
1618 },
1619 EmbeddingRefreshConfigLimits {
1620 entries: 0,
1621 ..EmbeddingRefreshConfigLimits::default()
1622 },
1623 EmbeddingRefreshConfigLimits {
1624 coordination: SearchCoordinationLimits {
1625 cleanup_entries: 0,
1626 ..SearchCoordinationLimits::default()
1627 },
1628 ..EmbeddingRefreshConfigLimits::default()
1629 },
1630 ] {
1631 assert!(validate_limits(limits).is_err());
1632 }
1633
1634 let config = EmbeddingRefreshConfig::default();
1635 assert!(matches!(
1636 config.to_canonical_json(EmbeddingRefreshConfigLimits {
1637 metadata_bytes: 1,
1638 ..EmbeddingRefreshConfigLimits::default()
1639 }),
1640 Err(SearchArtifactError::ResourceExhausted {
1641 resource: "embedding_refresh_config_bytes",
1642 ..
1643 })
1644 ));
1645
1646 let cleanup = tempfile::tempdir().unwrap();
1647 std::fs::write(cleanup.path().join("caller"), b"preserve").unwrap();
1648 assert!(matches!(
1649 cleanup_abandoned_temps(cleanup.path(), 0),
1650 Err(SearchArtifactError::InvalidSelector { .. })
1651 ));
1652 assert!(matches!(cleanup_abandoned_temps(cleanup.path(), 1), Ok(0)));
1653 std::fs::write(cleanup.path().join("caller-two"), b"preserve").unwrap();
1654 assert!(matches!(
1655 cleanup_abandoned_temps(cleanup.path(), 1),
1656 Err(SearchArtifactError::ResourceExhausted {
1657 resource: "embedding_refresh_cleanup_entries",
1658 ..
1659 })
1660 ));
1661 assert_eq!(
1662 std::fs::read(cleanup.path().join("caller")).unwrap(),
1663 b"preserve"
1664 );
1665 }
1666
1667 #[test]
1668 fn policy_and_wire_validation_matrix_is_exact_and_side_effect_free() {
1669 for debounce in [
1670 Duration::ZERO,
1671 MAX_DEBOUNCE + Duration::from_millis(1),
1672 Duration::from_nanos(1_500_000),
1673 ] {
1674 assert!(matches!(
1675 validate_debounce(debounce),
1676 Err(SearchArtifactError::InvalidSelector {
1677 field: "embedding refresh debounce",
1678 ..
1679 })
1680 ));
1681 }
1682 assert!(validate_debounce(Duration::from_millis(1)).is_ok());
1683 for max_concurrent_jobs in [0, MAX_EMBEDDING_REFRESH_JOBS + 1] {
1684 assert!(
1685 validate_project_policy(EmbeddingRefreshProjectPolicy {
1686 max_concurrent_jobs,
1687 ..EmbeddingRefreshProjectPolicy::default()
1688 })
1689 .is_err()
1690 );
1691 }
1692 assert!(
1693 validate_space_policy(EmbeddingRefreshSpacePolicy {
1694 proactive: None,
1695 debounce: None,
1696 })
1697 .is_err()
1698 );
1699 assert!(
1700 validate_outcome(EmbeddingRefreshOutcomeRecord {
1701 status: EmbeddingRefreshOutcomeStatus::Cancelled,
1702 graph_generation: 0,
1703 source_fingerprint: fingerprint(1),
1704 completed_at_micros: -1,
1705 })
1706 .is_err()
1707 );
1708
1709 let prior = outcome(3, 3, 30);
1710 for next in [outcome(3, 3, 29), outcome(2, 2, 30), outcome(3, 4, 30)] {
1711 assert!(validate_outcome_progress(prior, next).is_err());
1712 }
1713 assert!(validate_outcome_progress(prior, prior).is_ok());
1714
1715 let classes = [
1716 EmbeddingRefreshFailureClass::Provider,
1717 EmbeddingRefreshFailureClass::Validation,
1718 EmbeddingRefreshFailureClass::ResourceExhausted,
1719 EmbeddingRefreshFailureClass::Storage,
1720 EmbeddingRefreshFailureClass::ConcurrentMutation,
1721 EmbeddingRefreshFailureClass::Incompatible,
1722 EmbeddingRefreshFailureClass::Corrupt,
1723 EmbeddingRefreshFailureClass::Unavailable,
1724 ];
1725 for class in classes {
1726 let token = failure_token(class);
1727 assert_eq!(parse_failure_class(token).unwrap(), class);
1728 assert_eq!(
1729 outcome_tokens(EmbeddingRefreshOutcomeStatus::Failed(class)),
1730 ("failed", Some(token))
1731 );
1732 }
1733 assert!(parse_failure_class("unknown").is_err());
1734 assert_eq!(
1735 outcome_tokens(EmbeddingRefreshOutcomeStatus::Succeeded),
1736 ("succeeded", None)
1737 );
1738 assert_eq!(
1739 outcome_tokens(EmbeddingRefreshOutcomeStatus::Cancelled),
1740 ("cancelled", None)
1741 );
1742 }
1743
1744 #[test]
1745 fn direct_update_matrix_covers_idempotence_cleanup_capacity_and_policy_resolution() {
1746 let limits = EmbeddingRefreshConfigLimits {
1747 entries: 1,
1748 ..EmbeddingRefreshConfigLimits::default()
1749 };
1750 let mut config = EmbeddingRefreshConfig::default();
1751 assert!(config.is_empty());
1752 assert_eq!(config.len(), 0);
1753 let project_policy = config.project_policy();
1754 assert!(
1755 !apply_update(
1756 &mut config,
1757 EmbeddingRefreshConfigUpdate::SetProjectPolicy(project_policy),
1758 limits,
1759 )
1760 .unwrap()
1761 );
1762 assert!(
1763 !apply_update(
1764 &mut config,
1765 EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1766 compatibility_id: id(1),
1767 policy: None,
1768 },
1769 limits,
1770 )
1771 .unwrap()
1772 );
1773
1774 let policy = EmbeddingRefreshSpacePolicy {
1775 proactive: Some(false),
1776 debounce: Some(Duration::from_millis(25)),
1777 };
1778 assert!(
1779 apply_update(
1780 &mut config,
1781 EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1782 compatibility_id: id(1),
1783 policy: Some(policy),
1784 },
1785 limits,
1786 )
1787 .unwrap()
1788 );
1789 assert_eq!(config.len(), 1);
1790 assert_eq!(config.spaces()[0].policy, Some(policy));
1791 let resolved = config.resolved_policy(id(1));
1792 assert!(!resolved.proactive);
1793 assert_eq!(resolved.debounce, Duration::from_millis(25));
1794 assert!(
1795 !apply_update(
1796 &mut config,
1797 EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1798 compatibility_id: id(1),
1799 policy: Some(policy),
1800 },
1801 limits,
1802 )
1803 .unwrap()
1804 );
1805 assert!(matches!(
1806 apply_update(
1807 &mut config,
1808 EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1809 compatibility_id: id(2),
1810 policy: Some(policy),
1811 },
1812 limits,
1813 ),
1814 Err(SearchArtifactError::ResourceExhausted {
1815 resource: "embedding_refresh_config_entries",
1816 ..
1817 })
1818 ));
1819
1820 assert!(
1821 apply_update(
1822 &mut config,
1823 EmbeddingRefreshConfigUpdate::RecordOutcome {
1824 compatibility_id: id(1),
1825 outcome: outcome(1, 1, 10),
1826 },
1827 limits,
1828 )
1829 .unwrap()
1830 );
1831 assert!(
1832 apply_update(
1833 &mut config,
1834 EmbeddingRefreshConfigUpdate::SetSpacePolicy {
1835 compatibility_id: id(1),
1836 policy: None,
1837 },
1838 limits,
1839 )
1840 .unwrap()
1841 );
1842 assert_eq!(
1843 config.len(),
1844 1,
1845 "outcome retains the lineage after policy removal"
1846 );
1847 assert!(
1848 apply_update(
1849 &mut config,
1850 EmbeddingRefreshConfigUpdate::RemoveSpace {
1851 compatibility_id: id(1),
1852 },
1853 limits,
1854 )
1855 .unwrap()
1856 );
1857 assert!(
1858 !apply_update(
1859 &mut config,
1860 EmbeddingRefreshConfigUpdate::RemoveSpace {
1861 compatibility_id: id(1),
1862 },
1863 limits,
1864 )
1865 .unwrap()
1866 );
1867 assert!(config.is_empty());
1868
1869 let invalid_state = EmbeddingRefreshConfig {
1870 project: EmbeddingRefreshProjectPolicy::default(),
1871 spaces: BTreeMap::from([(id(1), StoredSpaceState::default())]),
1872 };
1873 assert!(invalid_state.to_canonical_json(limits).is_err());
1874 let too_many = EmbeddingRefreshConfig {
1875 project: EmbeddingRefreshProjectPolicy::default(),
1876 spaces: BTreeMap::from([
1877 (
1878 id(1),
1879 StoredSpaceState {
1880 policy: Some(policy),
1881 last_outcome: None,
1882 },
1883 ),
1884 (
1885 id(2),
1886 StoredSpaceState {
1887 policy: Some(policy),
1888 last_outcome: None,
1889 },
1890 ),
1891 ]),
1892 };
1893 assert!(matches!(
1894 too_many.to_canonical_json(limits),
1895 Err(SearchArtifactError::ResourceExhausted {
1896 resource: "embedding_refresh_config_entries",
1897 ..
1898 })
1899 ));
1900 }
1901
1902 #[test]
1903 fn wave10_raw_space_and_outcome_corruption_matrix_is_total() {
1904 let empty = RawSpace {
1905 compatibility_id: id(1).to_hex(),
1906 policy: None,
1907 last_outcome: None,
1908 };
1909 assert!(matches!(
1910 parse_space(empty),
1911 Err(SearchArtifactError::InvalidSelector {
1912 field: "embedding refresh space state",
1913 ..
1914 })
1915 ));
1916
1917 for (status, failure_class) in [
1918 ("succeeded", Some("provider")),
1919 ("cancelled", Some("storage")),
1920 ("failed", None),
1921 ("unknown", None),
1922 ] {
1923 let raw = RawOutcome {
1924 status: status.into(),
1925 failure_class: failure_class.map(str::to_owned),
1926 graph_generation: 1,
1927 source_fingerprint: fingerprint(1).to_hex(),
1928 completed_at_micros: 1,
1929 };
1930 assert!(matches!(
1931 parse_outcome(&raw),
1932 Err(SearchArtifactError::InvalidSelector {
1933 field: "embedding refresh outcome status",
1934 ..
1935 })
1936 ));
1937 }
1938 }
1939}