Skip to main content

feagi_evolutionary/genome/migration/
chain.rs

1// Copyright 2025 Neuraville Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Chain runner: walks a genome from its detected schema version to a
5//! target version, invoking the registered migrators in order.
6//!
7//! See `feagi-core/docs/GENOME_SCHEMA_VERSIONING.md` for the runner
8//! contract: between hops it runs the per-version validator as advisory;
9//! at the final hop it runs the target validator as blocking.
10
11use serde_json::{json, Value};
12
13use super::{ChainRegistry, ChainResult, MigrationError};
14use crate::genome::schema::{detect_schema_version, GenomeSchemaVersion};
15
16/// Walks a JSON genome through the chain of registered migrators.
17pub struct ChainRunner<'a> {
18    registry: &'a ChainRegistry,
19}
20
21impl<'a> ChainRunner<'a> {
22    pub fn new(registry: &'a ChainRegistry) -> Self {
23        Self { registry }
24    }
25
26    /// Run the chain, migrating `genome` from its detected version up to
27    /// `target`. The genome is mutated in place. The returned `ChainResult`
28    /// captures every hop and the final blocking validator's verdict.
29    ///
30    /// Errors:
31    /// - `DetectionFailed` if `detect_schema_version` rejects the input.
32    /// - `DowngradeRefused` if the genome is already past `target`.
33    /// - `MissingMigrator` if the chain has a gap somewhere in the range.
34    /// - `StepFailed` if any migrator returns an error.
35    pub fn run_to(
36        &self,
37        genome: &mut Value,
38        target: GenomeSchemaVersion,
39    ) -> Result<ChainResult, MigrationError> {
40        let from = detect_schema_version(genome)
41            .map_err(|e| MigrationError::DetectionFailed(e.to_string()))?;
42
43        if from > target {
44            return Err(MigrationError::DowngradeRefused { from, target });
45        }
46
47        let mut migrators_applied: Vec<&'static str> = Vec::new();
48        let mut normalizers_applied: Vec<&'static str> = Vec::new();
49        let mut per_step_diagnostics = Vec::new();
50        let mut per_normalizer_diagnostics = Vec::new();
51        let mut advisory_warnings: Vec<String> = Vec::new();
52
53        // If no migration is needed, the chain still owes the caller a
54        // normalize+validate pass at the starting (== target) version.
55        // The post-hop branch below handles the migration case symmetrically.
56        if from == target {
57            run_normalizer_if_present(
58                self.registry,
59                target,
60                genome,
61                &mut normalizers_applied,
62                &mut per_normalizer_diagnostics,
63            )?;
64        }
65
66        let mut current = from;
67        while current < target {
68            let migrator =
69                self.registry
70                    .migrator_for(current)
71                    .ok_or(MigrationError::MissingMigrator {
72                        from: current,
73                        target,
74                    })?;
75
76            let next = migrator.to_version();
77            debug_assert_eq!(
78                next.as_u32(),
79                current.as_u32() + 1,
80                "registry should have rejected non-contiguous migrators at registration"
81            );
82
83            let diagnostics = migrator.migrate(genome)?;
84            migrators_applied.push(migrator.name());
85
86            // The runner is the single source of truth for the schema
87            // version on the wire. Migrators should not touch this field;
88            // we stamp it after the step so if a migrator forgets, we
89            // still produce a self-consistent output genome.
90            stamp_schema_version(genome, next);
91            per_step_diagnostics.push(diagnostics);
92
93            // Normalizer (if any) cleans bad values within the new
94            // version before that version's validator inspects the genome.
95            run_normalizer_if_present(
96                self.registry,
97                next,
98                genome,
99                &mut normalizers_applied,
100                &mut per_normalizer_diagnostics,
101            )?;
102
103            // Advisory validation between hops. Errors are demoted to
104            // warnings here; only the final-target validator is blocking.
105            let is_final_hop = next == target;
106            if !is_final_hop {
107                let intermediate = self.registry.run_validator(next, genome);
108                for w in intermediate.warnings {
109                    advisory_warnings.push(format!("v{}: {}", next.as_u32(), w));
110                }
111                for e in intermediate.errors {
112                    advisory_warnings.push(format!("v{} (advisory): {}", next.as_u32(), e));
113                }
114            }
115
116            current = next;
117        }
118
119        // Final blocking validation at the target version. Runs once,
120        // exactly here, regardless of whether we got here via migration
121        // hops or were already at target.
122        let final_report = self.registry.run_validator(target, genome);
123        let blocking_errors = final_report.errors;
124        for w in final_report.warnings {
125            advisory_warnings.push(format!("v{}: {}", target.as_u32(), w));
126        }
127
128        Ok(ChainResult {
129            from_version: from,
130            to_version: target,
131            migrators_applied,
132            normalizers_applied,
133            per_step_diagnostics,
134            per_normalizer_diagnostics,
135            advisory_warnings,
136            blocking_errors,
137        })
138    }
139}
140
141/// Run the normalizer at `version` if one is registered, recording its
142/// name and diagnostics. No-op when no normalizer is registered for that
143/// version.
144fn run_normalizer_if_present(
145    registry: &ChainRegistry,
146    version: GenomeSchemaVersion,
147    genome: &mut Value,
148    applied: &mut Vec<&'static str>,
149    diagnostics: &mut Vec<crate::genome::normalizers::NormalizationDiagnostics>,
150) -> Result<(), MigrationError> {
151    if let Some(normalizer) = registry.normalizer_for(version) {
152        let diag = normalizer.normalize(genome)?;
153        applied.push(normalizer.name());
154        diagnostics.push(diag);
155    }
156    Ok(())
157}
158
159/// Write the integer `genome_schema_version` field on the genome.
160fn stamp_schema_version(genome: &mut Value, version: GenomeSchemaVersion) {
161    if let Some(obj) = genome.as_object_mut() {
162        obj.insert("genome_schema_version".to_string(), json!(version.as_u32()));
163    }
164}
165
166#[cfg(test)]
167mod tests {
168    use super::*;
169    use crate::genome::migration::test_support::{make_failing, make_ok};
170    use crate::genome::schema::CURRENT_SCHEMA_VERSION;
171    use serde_json::json;
172
173    fn registry_with_chain(from: u32, to_inclusive: u32) -> ChainRegistry {
174        let mut reg = ChainRegistry::new();
175        for v in from..to_inclusive {
176            let name = match v {
177                2 => "v2_to_v3_test",
178                100 => "v100_to_v101_test",
179                101 => "v101_to_v102_test",
180                _ => "synthetic",
181            };
182            reg.register_migrator(make_ok(v, name)).unwrap();
183        }
184        reg
185    }
186
187    #[test]
188    fn no_op_when_already_at_target() {
189        let reg = ChainRegistry::new();
190        let runner = ChainRunner::new(&reg);
191        let mut genome = json!({ "genome_schema_version": 3 });
192
193        let result = runner.run_to(&mut genome, CURRENT_SCHEMA_VERSION).unwrap();
194
195        assert_eq!(result.from_version, CURRENT_SCHEMA_VERSION);
196        assert_eq!(result.to_version, CURRENT_SCHEMA_VERSION);
197        assert!(result.migrators_applied.is_empty());
198        assert!(result.per_step_diagnostics.is_empty());
199        assert!(result.is_blocking_clean());
200    }
201
202    #[test]
203    fn single_hop_runs_one_migrator_and_stamps_version() {
204        let reg = registry_with_chain(2, 3);
205        let runner = ChainRunner::new(&reg);
206        let mut genome = json!({ "version": "2.0" });
207
208        let result = runner.run_to(&mut genome, GenomeSchemaVersion(3)).unwrap();
209
210        assert_eq!(result.from_version, GenomeSchemaVersion(2));
211        assert_eq!(result.to_version, GenomeSchemaVersion(3));
212        assert_eq!(result.migrators_applied, vec!["v2_to_v3_test"]);
213        assert_eq!(result.per_step_diagnostics.len(), 1);
214        assert_eq!(genome["genome_schema_version"], json!(3));
215        assert_eq!(genome["step_count"], json!(1));
216    }
217
218    #[test]
219    fn multi_hop_runs_each_migrator_in_order() {
220        let reg = registry_with_chain(100, 103);
221        let runner = ChainRunner::new(&reg);
222        let mut genome = json!({ "genome_schema_version": 100 });
223
224        let result = runner
225            .run_to(&mut genome, GenomeSchemaVersion(103))
226            .unwrap();
227
228        assert_eq!(result.from_version, GenomeSchemaVersion(100));
229        assert_eq!(result.to_version, GenomeSchemaVersion(103));
230        assert_eq!(
231            result.migrators_applied,
232            vec!["v100_to_v101_test", "v101_to_v102_test", "synthetic"]
233        );
234        assert_eq!(result.per_step_diagnostics.len(), 3);
235        assert_eq!(genome["genome_schema_version"], json!(103));
236        assert_eq!(genome["step_count"], json!(3));
237    }
238
239    #[test]
240    fn missing_migrator_in_range_errors() {
241        let mut reg = ChainRegistry::new();
242        // Register only the second hop, not the first.
243        reg.register_migrator(make_ok(101, "v101_to_v102_test"))
244            .unwrap();
245        let runner = ChainRunner::new(&reg);
246        let mut genome = json!({ "genome_schema_version": 100 });
247
248        let err = runner
249            .run_to(&mut genome, GenomeSchemaVersion(102))
250            .unwrap_err();
251        assert!(matches!(
252            err,
253            MigrationError::MissingMigrator { from, target }
254                if from == GenomeSchemaVersion(100) && target == GenomeSchemaVersion(102)
255        ));
256    }
257
258    #[test]
259    fn step_failure_aborts_chain() {
260        let mut reg = ChainRegistry::new();
261        reg.register_migrator(make_ok(100, "v100_to_v101_test"))
262            .unwrap();
263        reg.register_migrator(make_failing(101, "v101_to_v102_test"))
264            .unwrap();
265        let runner = ChainRunner::new(&reg);
266        let mut genome = json!({ "genome_schema_version": 100 });
267
268        let err = runner
269            .run_to(&mut genome, GenomeSchemaVersion(102))
270            .unwrap_err();
271        assert!(matches!(err, MigrationError::StepFailed { .. }));
272
273        // The first hop did succeed and stamped its version, even though
274        // the chain ultimately failed. This is intentional: failure leaves
275        // partially-migrated state visible for diagnostics.
276        assert_eq!(genome["genome_schema_version"], json!(101));
277    }
278
279    #[test]
280    fn downgrade_is_refused() {
281        let reg = ChainRegistry::new();
282        let runner = ChainRunner::new(&reg);
283        let mut genome = json!({ "genome_schema_version": 5 });
284        let err = runner
285            .run_to(&mut genome, GenomeSchemaVersion(3))
286            .unwrap_err();
287        assert!(matches!(
288            err,
289            MigrationError::DowngradeRefused { from, target }
290                if from == GenomeSchemaVersion(5) && target == GenomeSchemaVersion(3)
291        ));
292    }
293
294    #[test]
295    fn detection_failure_propagates() {
296        let reg = ChainRegistry::new();
297        let runner = ChainRunner::new(&reg);
298        let mut genome = json!({});
299        let err = runner
300            .run_to(&mut genome, GenomeSchemaVersion(3))
301            .unwrap_err();
302        assert!(matches!(err, MigrationError::DetectionFailed(_)));
303    }
304
305    /// Synthetic normalizer that adds a `was_normalized_at` array entry
306    /// recording the version it was invoked at, so tests can assert the
307    /// runner invoked it the right number of times in the right order.
308    struct SyntheticNormalizer {
309        version: GenomeSchemaVersion,
310        name: &'static str,
311    }
312
313    impl crate::genome::normalizers::Normalizer for SyntheticNormalizer {
314        fn schema_version(&self) -> GenomeSchemaVersion {
315            self.version
316        }
317        fn name(&self) -> &'static str {
318            self.name
319        }
320        fn normalize(
321            &self,
322            genome: &mut Value,
323        ) -> Result<crate::genome::normalizers::NormalizationDiagnostics, MigrationError> {
324            let mut diag = crate::genome::normalizers::NormalizationDiagnostics::new(self.version);
325            let arr = genome
326                .as_object_mut()
327                .expect("test genome must be a JSON object")
328                .entry("was_normalized_at".to_string())
329                .or_insert_with(|| json!([]));
330            arr.as_array_mut()
331                .expect("was_normalized_at must be array")
332                .push(json!(self.version.as_u32()));
333            diag.record(format!("normalized at v{}", self.version.as_u32()));
334            Ok(diag)
335        }
336    }
337
338    #[test]
339    fn normalizer_runs_when_already_at_target() {
340        let mut reg = ChainRegistry::new();
341        reg.register_normalizer(Box::new(SyntheticNormalizer {
342            version: GenomeSchemaVersion(3),
343            name: "v3_norm",
344        }))
345        .unwrap();
346        let runner = ChainRunner::new(&reg);
347        let mut genome = json!({ "genome_schema_version": 3 });
348
349        let result = runner.run_to(&mut genome, GenomeSchemaVersion(3)).unwrap();
350
351        assert_eq!(result.normalizers_applied, vec!["v3_norm"]);
352        assert_eq!(result.per_normalizer_diagnostics.len(), 1);
353        assert_eq!(genome["was_normalized_at"], json!([3]));
354    }
355
356    #[test]
357    fn normalizer_runs_after_each_hop() {
358        let mut reg = registry_with_chain(100, 103);
359        for v in 101..=103 {
360            reg.register_normalizer(Box::new(SyntheticNormalizer {
361                version: GenomeSchemaVersion(v),
362                name: match v {
363                    101 => "v101_norm",
364                    102 => "v102_norm",
365                    103 => "v103_norm",
366                    _ => unreachable!(),
367                },
368            }))
369            .unwrap();
370        }
371        let runner = ChainRunner::new(&reg);
372        let mut genome = json!({ "genome_schema_version": 100 });
373
374        let result = runner
375            .run_to(&mut genome, GenomeSchemaVersion(103))
376            .unwrap();
377
378        assert_eq!(
379            result.normalizers_applied,
380            vec!["v101_norm", "v102_norm", "v103_norm"]
381        );
382        assert_eq!(genome["was_normalized_at"], json!([101, 102, 103]));
383    }
384
385    #[test]
386    fn missing_normalizer_is_silent() {
387        // No normalizer registered for v3; runner should still succeed.
388        let reg = registry_with_chain(2, 3);
389        let runner = ChainRunner::new(&reg);
390        let mut genome = json!({ "version": "2.0" });
391
392        let result = runner.run_to(&mut genome, GenomeSchemaVersion(3)).unwrap();
393        assert!(result.normalizers_applied.is_empty());
394        assert!(result.per_normalizer_diagnostics.is_empty());
395    }
396
397    #[test]
398    fn registry_rejects_duplicate_normalizer() {
399        let mut reg = ChainRegistry::new();
400        reg.register_normalizer(Box::new(SyntheticNormalizer {
401            version: GenomeSchemaVersion(3),
402            name: "first",
403        }))
404        .unwrap();
405        let err = reg
406            .register_normalizer(Box::new(SyntheticNormalizer {
407                version: GenomeSchemaVersion(3),
408                name: "second",
409            }))
410            .unwrap_err();
411        assert!(matches!(err, MigrationError::InvalidRegistry(_)));
412    }
413}