Skip to main content

lenso_service/
extraction_plan.rs

1use crate::{
2    CommonContextRequirement, CompatibilityCategory, EXTRACTION_READINESS_REPORT_PROTOCOL,
3    ExtractionContractDirection, ExtractionContractKind, ExtractionCursorEvidence,
4    ExtractionDataEvidenceSource, ExtractionReadinessIssueCode, ExtractionReadinessReport,
5    ExtractionReadinessSurfaceSummary, ModuleManifest, ServiceTenancyMode, system_v2_graph,
6};
7use schemars::JsonSchema;
8use serde::{Deserialize, Serialize};
9use serde_json::{Value, json};
10use sha2::{Digest, Sha256};
11use std::collections::{BTreeMap, BTreeSet};
12use std::fmt::{self, Write as _};
13
14pub const EXTRACTION_PLAN_PROTOCOL: &str = "lenso.extraction-plan.v1";
15pub const EXTRACTION_PLAN_GENERATOR_VERSION: &str = "lenso.extraction-plan-generator.v1";
16const EXTRACTION_PLAN_SCHEMA_ID: &str =
17    "https://contracts.lenso.local/extraction/lenso.extraction-plan.v1.schema.json";
18
19#[derive(
20    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
21)]
22#[serde(rename_all = "snake_case")]
23pub enum ExtractionAuthorityKind {
24    LinkedHost,
25    AutonomousService,
26}
27
28#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
29#[serde(rename_all = "camelCase")]
30pub struct ExtractionExpectedAuthority {
31    pub kind: ExtractionAuthorityKind,
32    pub owner_id: String,
33    pub revision: String,
34}
35
36#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
37#[serde(rename_all = "camelCase")]
38pub struct ExtractionPlanContractVersion {
39    pub contract_id: String,
40    pub version: String,
41    pub kind: ExtractionContractKind,
42    pub direction: ExtractionContractDirection,
43    pub artifact_reference: String,
44    pub artifact_digest: String,
45    pub artifact_format: ExtractionContractArtifactFormat,
46    pub tenancy_mode: ServiceTenancyMode,
47    #[serde(default)]
48    pub required_context: Vec<CommonContextRequirement>,
49    #[serde(default, skip_serializing_if = "Option::is_none")]
50    pub producer_id: Option<String>,
51    #[serde(default)]
52    pub consumer_ids: Vec<String>,
53}
54
55#[derive(
56    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
57)]
58#[serde(rename_all = "snake_case")]
59pub enum ExtractionContractArtifactFormat {
60    Openapi,
61    Protobuf,
62    JsonSchema,
63}
64
65#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
66#[serde(rename_all = "camelCase")]
67pub struct ExtractionEvidenceDigest {
68    pub reference: String,
69    pub digest: String,
70}
71
72#[derive(Debug, Clone, Serialize, Deserialize)]
73pub struct ExtractionPlanInputs {
74    pub readiness_report: ExtractionReadinessReport,
75    pub module: ModuleManifest,
76    pub system: Value,
77    pub contract_versions: Vec<ExtractionPlanContractVersion>,
78    pub expected_authority: ExtractionExpectedAuthority,
79    pub evidence_digests: Vec<ExtractionEvidenceDigest>,
80}
81
82#[derive(
83    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
84)]
85#[serde(rename_all = "snake_case")]
86pub enum ExtractionInputPinKind {
87    ReadinessEvidence,
88    ModuleDeclaration,
89    ContractVersion,
90    SystemGraph,
91    AnalyzerVersion,
92    DataMapping,
93    AuthorityRevision,
94    Evidence,
95}
96
97#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
98#[serde(rename_all = "camelCase")]
99pub struct ExtractionInputPin {
100    pub kind: ExtractionInputPinKind,
101    pub subject: String,
102    pub digest: String,
103}
104
105#[derive(
106    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
107)]
108#[serde(rename_all = "snake_case")]
109pub enum ExtractionCopyMode {
110    OnlineCheckpointed,
111    BoundedWritePause,
112}
113
114#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
115#[serde(rename_all = "camelCase")]
116pub struct ExtractionTableMapping {
117    pub source_table: String,
118    pub destination_table: String,
119    pub destination_store: String,
120    pub owner_module: String,
121    pub copy_mode: ExtractionCopyMode,
122    #[serde(default)]
123    pub evidence_sources: Vec<ExtractionDataEvidenceSource>,
124    #[serde(default)]
125    pub cursors: Vec<ExtractionCursorEvidence>,
126}
127
128#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
129#[serde(rename_all = "camelCase")]
130pub struct ExtractionMigrationMapping {
131    pub source_migration: String,
132    pub source_reference: String,
133    pub source_digest: String,
134    pub destination_store: String,
135    pub owner_module: String,
136}
137
138#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
139#[serde(rename_all = "camelCase")]
140pub struct ExtractionDataMapping {
141    pub store_engine: String,
142    pub destination_store: String,
143    #[serde(default)]
144    pub tables: Vec<ExtractionTableMapping>,
145    #[serde(default)]
146    pub migrations: Vec<ExtractionMigrationMapping>,
147}
148
149#[derive(
150    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
151)]
152#[serde(rename_all = "snake_case")]
153pub enum ExtractionWorkloadRole {
154    Api,
155    Worker,
156    Migration,
157}
158
159#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
160#[serde(rename_all = "camelCase")]
161pub struct ExtractionWorkloadPlan {
162    pub workload_id: String,
163    pub role: ExtractionWorkloadRole,
164}
165
166#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
167#[serde(rename_all = "camelCase")]
168pub struct ExtractionStorePlan {
169    pub store_id: String,
170    pub engine: String,
171    pub isolated: bool,
172}
173
174#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
175#[serde(rename_all = "camelCase")]
176pub struct ExtractionServiceReferencePlan {
177    pub reference_id: String,
178    pub contract_id: String,
179    pub version: String,
180    pub direction: ExtractionContractDirection,
181    pub target_service_id: String,
182}
183
184#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
185#[serde(rename_all = "camelCase")]
186pub struct ExtractionGeneratedClientPlan {
187    pub client_id: String,
188    pub owner_id: String,
189    pub contract_id: String,
190    pub version: String,
191    pub artifact_reference: String,
192}
193
194#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
195#[serde(rename_all = "camelCase")]
196pub struct ExtractionServicePlan {
197    pub service_id: String,
198    pub module_id: String,
199    pub workloads: Vec<ExtractionWorkloadPlan>,
200    pub store: ExtractionStorePlan,
201    pub contract_versions: Vec<ExtractionPlanContractVersion>,
202    pub service_references: Vec<ExtractionServiceReferencePlan>,
203    pub generated_clients: Vec<ExtractionGeneratedClientPlan>,
204    pub preserved_capabilities: Vec<String>,
205    pub preserved_surfaces: ExtractionReadinessSurfaceSummary,
206}
207
208#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema)]
209#[serde(rename_all = "camelCase")]
210pub struct ExtractionPlanDiffEntry {
211    pub subject: String,
212    #[serde(default, skip_serializing_if = "Option::is_none")]
213    pub before: Option<String>,
214    #[serde(default, skip_serializing_if = "Option::is_none")]
215    pub after: Option<String>,
216}
217
218#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
219#[serde(rename_all = "camelCase")]
220pub struct ExtractionPlanDiff {
221    pub entries: Vec<ExtractionPlanDiffEntry>,
222}
223
224#[derive(
225    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
226)]
227#[serde(rename_all = "snake_case")]
228pub enum ExtractionPlanIssueCode {
229    PlanIntegrityInvalid,
230    ReadinessEvidenceChanged,
231    ModuleDeclarationChanged,
232    ContractVersionChanged,
233    SystemGraphChanged,
234    AnalyzerVersionChanged,
235    DataMappingChanged,
236    AuthorityRevisionChanged,
237    InputEvidenceChanged,
238    ScaffoldConflict,
239    DestinationExpansionFailed,
240    BackfillCheckpointStale,
241    ReconciliationMismatch,
242    DrainIncomplete,
243    ProvisionalCutoverFailed,
244    VerificationFailed,
245    RollbackRequired,
246    FinalApprovalRequired,
247    TerminalEvidenceIncomplete,
248}
249
250#[derive(
251    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
252)]
253#[serde(rename_all = "snake_case")]
254pub enum ExtractionPlanPhaseKind {
255    Analysis,
256    Scaffold,
257    DestinationExpansion,
258    Backfill,
259    Reconciliation,
260    Drain,
261    ProvisionalCutover,
262    Verification,
263    RollbackOrCommit,
264    TerminalEvidence,
265}
266
267#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
268#[serde(rename_all = "camelCase")]
269pub struct ExtractionApprovalBoundary {
270    pub boundary_id: String,
271    pub phase_id: String,
272    pub action: String,
273    pub reason: String,
274    pub required_pins: Vec<ExtractionInputPinKind>,
275}
276
277#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
278#[serde(rename_all = "camelCase")]
279pub struct ExtractionPlanPhase {
280    pub phase_id: String,
281    pub order: u16,
282    pub kind: ExtractionPlanPhaseKind,
283    #[serde(default)]
284    pub prerequisite_phase_ids: Vec<String>,
285    pub prerequisites: Vec<String>,
286    pub intended_mutations: Vec<String>,
287    pub expected_evidence: Vec<String>,
288    pub rollback_conditions: Vec<String>,
289    pub issue_codes: Vec<ExtractionPlanIssueCode>,
290    pub next_actions: Vec<String>,
291    #[serde(default, skip_serializing_if = "Option::is_none")]
292    pub approval_boundary: Option<ExtractionApprovalBoundary>,
293}
294
295#[allow(clippy::struct_excessive_bools)]
296#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
297#[serde(rename_all = "camelCase")]
298pub struct ExtractionPlanEffects {
299    pub writes_repository_files: bool,
300    pub starts_workloads: bool,
301    pub copies_data: bool,
302    pub changes_authority: bool,
303}
304
305#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
306#[serde(rename_all = "camelCase")]
307pub struct ExtractionPlan {
308    pub protocol: String,
309    pub generator_version: String,
310    pub plan_id: String,
311    pub plan_digest: String,
312    pub target_module: String,
313    pub source_system_id: String,
314    pub readiness_classification: CompatibilityCategory,
315    pub readiness_issue_codes: Vec<ExtractionReadinessIssueCode>,
316    pub expected_authority: ExtractionExpectedAuthority,
317    pub pinned_inputs: Vec<ExtractionInputPin>,
318    pub data_mapping: ExtractionDataMapping,
319    pub proposed_service: ExtractionServicePlan,
320    pub diff: ExtractionPlanDiff,
321    pub phases: Vec<ExtractionPlanPhase>,
322    pub approval_boundaries: Vec<ExtractionApprovalBoundary>,
323    pub effects: ExtractionPlanEffects,
324}
325
326#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
327#[serde(rename_all = "snake_case")]
328pub enum ExtractionPlanGenerationIssueCode {
329    ReadinessNotReady,
330    ReadinessTargetMismatch,
331    SystemEvidenceInvalid,
332    AuthorityMismatch,
333    ContractVersionsMissing,
334    InputInvalid,
335}
336
337#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
338#[serde(rename_all = "camelCase")]
339pub struct ExtractionPlanGenerationError {
340    pub code: ExtractionPlanGenerationIssueCode,
341    pub message: String,
342    pub next_actions: Vec<String>,
343}
344
345impl fmt::Display for ExtractionPlanGenerationError {
346    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
347        formatter.write_str(&self.message)
348    }
349}
350
351impl std::error::Error for ExtractionPlanGenerationError {}
352
353#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
354#[serde(rename_all = "camelCase")]
355pub struct ExtractionStaleInput {
356    pub kind: ExtractionInputPinKind,
357    pub subject: String,
358    #[serde(default, skip_serializing_if = "Option::is_none")]
359    pub planned_digest: Option<String>,
360    #[serde(default, skip_serializing_if = "Option::is_none")]
361    pub current_digest: Option<String>,
362}
363
364#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
365#[serde(rename_all = "camelCase")]
366pub struct ExtractionPlanRejection {
367    pub plan_id: String,
368    pub message: String,
369    pub issue_codes: Vec<ExtractionPlanIssueCode>,
370    pub stale_inputs: Vec<ExtractionStaleInput>,
371    pub next_actions: Vec<String>,
372    pub effects: ExtractionPlanEffects,
373}
374
375impl fmt::Display for ExtractionPlanRejection {
376    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
377        formatter.write_str(&self.message)
378    }
379}
380
381impl std::error::Error for ExtractionPlanRejection {}
382
383#[derive(Serialize)]
384#[serde(rename_all = "camelCase")]
385struct ExtractionPlanContent<'a> {
386    protocol: &'a str,
387    generator_version: &'a str,
388    target_module: &'a str,
389    source_system_id: &'a str,
390    readiness_classification: CompatibilityCategory,
391    readiness_issue_codes: &'a [ExtractionReadinessIssueCode],
392    expected_authority: &'a ExtractionExpectedAuthority,
393    pinned_inputs: &'a [ExtractionInputPin],
394    data_mapping: &'a ExtractionDataMapping,
395    proposed_service: &'a ExtractionServicePlan,
396    diff: &'a ExtractionPlanDiff,
397    phases: &'a [ExtractionPlanPhase],
398    approval_boundaries: &'a [ExtractionApprovalBoundary],
399    effects: ExtractionPlanEffects,
400}
401
402#[must_use]
403pub fn extraction_input_digest(bytes: impl AsRef<[u8]>) -> String {
404    let digest = Sha256::digest(bytes.as_ref());
405    let mut rendered = String::with_capacity(7 + digest.len() * 2);
406    rendered.push_str("sha256:");
407    for byte in digest {
408        write!(&mut rendered, "{byte:02x}").expect("writing to String cannot fail");
409    }
410    rendered
411}
412
413pub fn generate_extraction_plan(
414    inputs: &ExtractionPlanInputs,
415) -> Result<ExtractionPlan, ExtractionPlanGenerationError> {
416    validate_generation_inputs(inputs)?;
417    let graph = system_v2_graph(&inputs.system).map_err(|issues| {
418        generation_error(
419            ExtractionPlanGenerationIssueCode::SystemEvidenceInvalid,
420            format!(
421                "System evidence is invalid: {}",
422                issues
423                    .iter()
424                    .map(|issue| issue.code.as_str())
425                    .collect::<Vec<_>>()
426                    .join(", ")
427            ),
428            "Correct the lenso.system.v2 graph and regenerate the Extraction Plan.",
429        )
430    })?;
431    let target_owner = graph
432        .nodes
433        .iter()
434        .find(|node| node.kind == "module" && node.id == inputs.module.module_id)
435        .and_then(|node| node.owner.as_deref());
436    if target_owner != Some(inputs.expected_authority.owner_id.as_str()) {
437        return Err(generation_error(
438            ExtractionPlanGenerationIssueCode::AuthorityMismatch,
439            "The current System graph does not assign the target Module to the expected linked Host authority.",
440            "Refresh readiness and authority evidence from the current System graph.",
441        ));
442    }
443    let source_system_id = graph.system_id;
444    if inputs.readiness_report.system_id.as_deref() != Some(source_system_id.as_str()) {
445        return Err(generation_error(
446            ExtractionPlanGenerationIssueCode::SystemEvidenceInvalid,
447            "The readiness report and current System graph identify different Systems.",
448            "Regenerate readiness evidence from the current lenso.system.v2 artifact.",
449        ));
450    }
451    let proposed_service = proposed_service(inputs);
452    let data_mapping = data_mapping(inputs, &proposed_service.store.store_id)?;
453    let pinned_inputs = pinned_inputs(inputs, &data_mapping)?;
454    let diff = plan_diff(inputs, &proposed_service);
455    let phases = plan_phases(&proposed_service, &data_mapping);
456    let approval_boundaries = phases
457        .iter()
458        .filter_map(|phase| phase.approval_boundary.clone())
459        .collect::<Vec<_>>();
460    let effects = ExtractionPlanEffects::default();
461    let content = ExtractionPlanContent {
462        protocol: EXTRACTION_PLAN_PROTOCOL,
463        generator_version: EXTRACTION_PLAN_GENERATOR_VERSION,
464        target_module: &inputs.module.module_id,
465        source_system_id: &source_system_id,
466        readiness_classification: inputs.readiness_report.classification,
467        readiness_issue_codes: &inputs.readiness_report.issue_codes,
468        expected_authority: &inputs.expected_authority,
469        pinned_inputs: &pinned_inputs,
470        data_mapping: &data_mapping,
471        proposed_service: &proposed_service,
472        diff: &diff,
473        phases: &phases,
474        approval_boundaries: &approval_boundaries,
475        effects,
476    };
477    let plan_digest = digest_serializable(&content)?;
478    Ok(ExtractionPlan {
479        protocol: EXTRACTION_PLAN_PROTOCOL.to_owned(),
480        generator_version: EXTRACTION_PLAN_GENERATOR_VERSION.to_owned(),
481        plan_id: format!("extraction-plan:{plan_digest}"),
482        plan_digest,
483        target_module: inputs.module.module_id.clone(),
484        source_system_id,
485        readiness_classification: inputs.readiness_report.classification,
486        readiness_issue_codes: inputs.readiness_report.issue_codes.clone(),
487        expected_authority: inputs.expected_authority.clone(),
488        pinned_inputs,
489        data_mapping,
490        proposed_service,
491        diff,
492        phases,
493        approval_boundaries,
494        effects,
495    })
496}
497
498pub fn dry_run_extraction_plan(
499    inputs: &ExtractionPlanInputs,
500) -> Result<ExtractionPlan, ExtractionPlanGenerationError> {
501    generate_extraction_plan(inputs)
502}
503
504#[must_use]
505pub fn extraction_plan_integrity_is_valid(plan: &ExtractionPlan) -> bool {
506    if plan.protocol != EXTRACTION_PLAN_PROTOCOL
507        || plan.generator_version != EXTRACTION_PLAN_GENERATOR_VERSION
508        || plan.plan_id != format!("extraction-plan:{}", plan.plan_digest)
509    {
510        return false;
511    }
512    let content = ExtractionPlanContent {
513        protocol: &plan.protocol,
514        generator_version: &plan.generator_version,
515        target_module: &plan.target_module,
516        source_system_id: &plan.source_system_id,
517        readiness_classification: plan.readiness_classification,
518        readiness_issue_codes: &plan.readiness_issue_codes,
519        expected_authority: &plan.expected_authority,
520        pinned_inputs: &plan.pinned_inputs,
521        data_mapping: &plan.data_mapping,
522        proposed_service: &plan.proposed_service,
523        diff: &plan.diff,
524        phases: &plan.phases,
525        approval_boundaries: &plan.approval_boundaries,
526        effects: plan.effects,
527    };
528    digest_serializable(&content).is_ok_and(|digest| digest == plan.plan_digest)
529}
530
531#[allow(clippy::result_large_err)]
532pub fn ensure_extraction_plan_fresh(
533    plan: &ExtractionPlan,
534    current_inputs: &ExtractionPlanInputs,
535) -> Result<(), ExtractionPlanRejection> {
536    if !extraction_plan_integrity_is_valid(plan) {
537        return Err(ExtractionPlanRejection {
538            plan_id: plan.plan_id.clone(),
539            message: "Extraction Plan integrity validation failed before mutation.".to_owned(),
540            issue_codes: vec![ExtractionPlanIssueCode::PlanIntegrityInvalid],
541            stale_inputs: Vec::new(),
542            next_actions: vec![
543                "Discard the modified plan and generate a new content-addressed Extraction Plan."
544                    .to_owned(),
545            ],
546            effects: ExtractionPlanEffects::default(),
547        });
548    }
549
550    let destination_store = format!("{}-service-store", current_inputs.module.module_id);
551    let current_mapping = data_mapping(current_inputs, &destination_store).map_err(|error| {
552        ExtractionPlanRejection {
553            plan_id: plan.plan_id.clone(),
554            message:
555                "Current migration evidence could not be pinned; the plan is stale before mutation."
556                    .to_owned(),
557            issue_codes: vec![ExtractionPlanIssueCode::DataMappingChanged],
558            stale_inputs: Vec::new(),
559            next_actions: error.next_actions,
560            effects: ExtractionPlanEffects::default(),
561        }
562    })?;
563    let current_pins = pinned_inputs(current_inputs, &current_mapping).map_err(|error| {
564        ExtractionPlanRejection {
565            plan_id: plan.plan_id.clone(),
566            message:
567                "Current extraction inputs could not be pinned; the plan is stale before mutation."
568                    .to_owned(),
569            issue_codes: vec![ExtractionPlanIssueCode::InputEvidenceChanged],
570            stale_inputs: Vec::new(),
571            next_actions: error.next_actions,
572            effects: ExtractionPlanEffects::default(),
573        }
574    })?;
575    let planned = plan
576        .pinned_inputs
577        .iter()
578        .map(|pin| ((pin.kind, pin.subject.as_str()), pin.digest.as_str()))
579        .collect::<BTreeMap<_, _>>();
580    let current = current_pins
581        .iter()
582        .map(|pin| ((pin.kind, pin.subject.as_str()), pin.digest.as_str()))
583        .collect::<BTreeMap<_, _>>();
584    let keys = planned
585        .keys()
586        .chain(current.keys())
587        .copied()
588        .collect::<BTreeSet<_>>();
589    let mut stale_inputs = keys
590        .into_iter()
591        .filter_map(|(kind, subject)| {
592            let planned_digest = planned.get(&(kind, subject)).copied();
593            let current_digest = current.get(&(kind, subject)).copied();
594            (planned_digest != current_digest).then(|| ExtractionStaleInput {
595                kind,
596                subject: subject.to_owned(),
597                planned_digest: planned_digest.map(str::to_owned),
598                current_digest: current_digest.map(str::to_owned),
599            })
600        })
601        .collect::<Vec<_>>();
602    stale_inputs
603        .sort_by(|left, right| (&left.kind, &left.subject).cmp(&(&right.kind, &right.subject)));
604    if stale_inputs.is_empty() {
605        return Ok(());
606    }
607    let mut issue_codes = stale_inputs
608        .iter()
609        .map(|input| stale_issue_code(input.kind))
610        .collect::<Vec<_>>();
611    issue_codes.sort();
612    issue_codes.dedup();
613    Err(ExtractionPlanRejection {
614        plan_id: plan.plan_id.clone(),
615        message: "Pinned extraction inputs changed; the stale plan was rejected before mutation."
616            .to_owned(),
617        issue_codes,
618        stale_inputs,
619        next_actions: vec![
620            "Rerun readiness analysis and generate a new Extraction Plan from the current inputs."
621                .to_owned(),
622        ],
623        effects: ExtractionPlanEffects::default(),
624    })
625}
626
627fn validate_generation_inputs(
628    inputs: &ExtractionPlanInputs,
629) -> Result<(), ExtractionPlanGenerationError> {
630    if inputs.readiness_report.protocol != EXTRACTION_READINESS_REPORT_PROTOCOL
631        || inputs.readiness_report.analyzer_version.trim().is_empty()
632    {
633        return Err(generation_error(
634            ExtractionPlanGenerationIssueCode::InputInvalid,
635            "The Extraction Readiness Report protocol or analyzer version is invalid.",
636            "Regenerate readiness evidence with a supported public analyzer.",
637        ));
638    }
639    if !inputs.readiness_report.ready
640        || matches!(
641            inputs.readiness_report.classification,
642            CompatibilityCategory::Breaking | CompatibilityCategory::Blocked
643        )
644    {
645        return Err(generation_error(
646            ExtractionPlanGenerationIssueCode::ReadinessNotReady,
647            "Only a ready linked Module can produce an Extraction Plan.",
648            "Resolve the readiness findings and rerun the public readiness command.",
649        ));
650    }
651    if inputs.readiness_report.target_module != inputs.module.module_id {
652        return Err(generation_error(
653            ExtractionPlanGenerationIssueCode::ReadinessTargetMismatch,
654            "The readiness report does not describe the requested Module declaration.",
655            "Regenerate readiness evidence for exactly the target Module.",
656        ));
657    }
658    if inputs.expected_authority.kind != ExtractionAuthorityKind::LinkedHost
659        || inputs.expected_authority.owner_id.trim().is_empty()
660        || inputs.expected_authority.revision.trim().is_empty()
661        || inputs.readiness_report.target_owner.as_deref()
662            != Some(inputs.expected_authority.owner_id.as_str())
663    {
664        return Err(generation_error(
665            ExtractionPlanGenerationIssueCode::AuthorityMismatch,
666            "Expected authority must pin the linked Host owner and a non-empty revision.",
667            "Read the current linked authority and regenerate the plan with its exact revision.",
668        ));
669    }
670    if inputs.contract_versions.is_empty() {
671        return Err(generation_error(
672            ExtractionPlanGenerationIssueCode::ContractVersionsMissing,
673            "Extraction Plan inputs must include the relevant authoritative Contract Versions.",
674            "Resolve every provided or consumed Contract artifact and supply its version and digest.",
675        ));
676    }
677    validate_contracts(&inputs.contract_versions)?;
678    validate_relevant_contracts(inputs)?;
679    validate_evidence_digests(&inputs.evidence_digests)?;
680    Ok(())
681}
682
683fn validate_relevant_contracts(
684    inputs: &ExtractionPlanInputs,
685) -> Result<(), ExtractionPlanGenerationError> {
686    let planned = inputs.contract_versions.iter().fold(
687        BTreeMap::<&str, Vec<&ExtractionPlanContractVersion>>::new(),
688        |mut map, contract| {
689            map.entry(contract.contract_id.as_str())
690                .or_default()
691                .push(contract);
692            map
693        },
694    );
695    for evidence in &inputs.readiness_report.contract_evidence {
696        if evidence.status != crate::ExtractionEvidenceStatus::Present {
697            continue;
698        }
699        let Some(contract_id) = evidence.contract_id.as_deref() else {
700            continue;
701        };
702        let matches = planned
703            .get(contract_id)
704            .map(Vec::as_slice)
705            .unwrap_or_default();
706        if !matches.iter().any(|contract| {
707            contract.kind == evidence.kind && contract.direction == evidence.direction
708        }) {
709            return Err(generation_error(
710                ExtractionPlanGenerationIssueCode::ContractVersionsMissing,
711                format!(
712                    "Readiness evidence requires Contract `{contract_id}` with kind {:?} and direction {:?}, but the plan input does not pin it.",
713                    evidence.kind, evidence.direction
714                ),
715                "Supply the exact authoritative Contract Version and artifact digest used by readiness.",
716            ));
717        }
718    }
719    Ok(())
720}
721
722fn validate_contracts(
723    contracts: &[ExtractionPlanContractVersion],
724) -> Result<(), ExtractionPlanGenerationError> {
725    let mut identities = BTreeSet::new();
726    for contract in contracts {
727        if contract.contract_id.trim().is_empty()
728            || contract.version.trim().is_empty()
729            || contract.artifact_reference.trim().is_empty()
730            || !valid_sha256_digest(&contract.artifact_digest)
731        {
732            return Err(generation_error(
733                ExtractionPlanGenerationIssueCode::InputInvalid,
734                "Contract Version inputs require stable identities, artifact references, and SHA-256 digests.",
735                "Correct the Contract Version input and regenerate the Extraction Plan.",
736            ));
737        }
738        if contract.kind == ExtractionContractKind::Service
739            && contract.direction == ExtractionContractDirection::Consumes
740            && contract
741                .producer_id
742                .as_deref()
743                .is_none_or(|producer| producer.trim().is_empty())
744        {
745            return Err(generation_error(
746                ExtractionPlanGenerationIssueCode::InputInvalid,
747                "A consumed Service Contract must identify its producing Service.",
748                "Resolve the producer from the current System graph before planning a generated client.",
749            ));
750        }
751        if !matches!(
752            (contract.kind, contract.artifact_format),
753            (
754                ExtractionContractKind::Service,
755                ExtractionContractArtifactFormat::Openapi
756                    | ExtractionContractArtifactFormat::Protobuf
757            ) | (
758                ExtractionContractKind::Event,
759                ExtractionContractArtifactFormat::JsonSchema
760                    | ExtractionContractArtifactFormat::Protobuf
761            )
762        ) {
763            return Err(generation_error(
764                ExtractionPlanGenerationIssueCode::InputInvalid,
765                "A Contract Version uses an artifact format that does not match its Service or Event kind.",
766                "Use OpenAPI or Protobuf for Service Contracts and JSON Schema or Protobuf for Event Contracts.",
767            ));
768        }
769        if contract.tenancy_mode == ServiceTenancyMode::Required
770            && !contract
771                .required_context
772                .contains(&CommonContextRequirement::Tenant)
773        {
774            return Err(generation_error(
775                ExtractionPlanGenerationIssueCode::InputInvalid,
776                "A tenant-required Contract Version must require tenant context.",
777                "Preserve the authoritative Contract context requirements before generating the plan.",
778            ));
779        }
780        let mut required_context = contract.required_context.clone();
781        required_context.sort();
782        required_context.dedup();
783        if required_context != contract.required_context {
784            return Err(generation_error(
785                ExtractionPlanGenerationIssueCode::InputInvalid,
786                "Contract context requirements must be unique and deterministically ordered.",
787                "Sort and deduplicate the authoritative Contract context requirements.",
788            ));
789        }
790        let identity = (contract.contract_id.as_str(), contract.version.as_str());
791        if !identities.insert(identity) {
792            return Err(generation_error(
793                ExtractionPlanGenerationIssueCode::InputInvalid,
794                "A Contract Version input is duplicated.",
795                "Supply one authoritative entry per Contract Version, kind, and direction.",
796            ));
797        }
798    }
799    Ok(())
800}
801
802fn validate_evidence_digests(
803    evidence: &[ExtractionEvidenceDigest],
804) -> Result<(), ExtractionPlanGenerationError> {
805    if evidence.is_empty() {
806        return Err(generation_error(
807            ExtractionPlanGenerationIssueCode::InputInvalid,
808            "Extraction Plan inputs must include digests for the readiness evidence bundle.",
809            "Digest the analyzer and Store evidence consumed by readiness before planning.",
810        ));
811    }
812    let mut references = BTreeMap::new();
813    for item in evidence {
814        if item.reference.trim().is_empty() || !valid_sha256_digest(&item.digest) {
815            return Err(generation_error(
816                ExtractionPlanGenerationIssueCode::InputInvalid,
817                "Evidence inputs require a stable reference and SHA-256 digest.",
818                "Digest every analyzer, topology, Contract, and Store evidence input before planning.",
819            ));
820        }
821        if references
822            .insert(item.reference.as_str(), item.digest.as_str())
823            .is_some()
824        {
825            return Err(generation_error(
826                ExtractionPlanGenerationIssueCode::InputInvalid,
827                "An evidence input reference is duplicated.",
828                "Supply exactly one digest for every evidence input reference.",
829            ));
830        }
831    }
832    Ok(())
833}
834
835fn proposed_service(inputs: &ExtractionPlanInputs) -> ExtractionServicePlan {
836    let service_id = format!("{}-service", inputs.module.module_id);
837    let store_id = format!("{service_id}-store");
838    let mut contracts = inputs.contract_versions.clone();
839    for contract in &mut contracts {
840        normalize_strings(&mut contract.consumer_ids);
841        contract.required_context.sort();
842        contract.required_context.dedup();
843    }
844    contracts.sort();
845    let mut service_references = contracts
846        .iter()
847        .filter(|contract| contract.kind == ExtractionContractKind::Service)
848        .map(|contract| ExtractionServiceReferencePlan {
849            reference_id: format!(
850                "{}-{}-{}-reference",
851                service_id,
852                stable_slug(&contract.contract_id),
853                stable_slug(&contract.version)
854            ),
855            contract_id: contract.contract_id.clone(),
856            version: contract.version.clone(),
857            direction: contract.direction,
858            target_service_id: if contract.direction == ExtractionContractDirection::Provides {
859                service_id.clone()
860            } else {
861                contract
862                    .producer_id
863                    .clone()
864                    .expect("consumed Service Contracts require a producer")
865            },
866        })
867        .collect::<Vec<_>>();
868    service_references.sort();
869    service_references.dedup();
870
871    let mut generated_clients = Vec::new();
872    for contract in contracts
873        .iter()
874        .filter(|contract| contract.kind == ExtractionContractKind::Service)
875    {
876        if contract.direction == ExtractionContractDirection::Consumes {
877            generated_clients.push(ExtractionGeneratedClientPlan {
878                client_id: format!(
879                    "{}-{}-client",
880                    service_id,
881                    stable_slug(&contract.contract_id)
882                ),
883                owner_id: service_id.clone(),
884                contract_id: contract.contract_id.clone(),
885                version: contract.version.clone(),
886                artifact_reference: contract.artifact_reference.clone(),
887            });
888        } else {
889            generated_clients.extend(contract.consumer_ids.iter().map(|consumer_id| {
890                ExtractionGeneratedClientPlan {
891                    client_id: format!(
892                        "{}-{}-client",
893                        stable_slug(consumer_id),
894                        stable_slug(&contract.contract_id)
895                    ),
896                    owner_id: consumer_id.clone(),
897                    contract_id: contract.contract_id.clone(),
898                    version: contract.version.clone(),
899                    artifact_reference: contract.artifact_reference.clone(),
900                }
901            }));
902        }
903    }
904    generated_clients.sort();
905    generated_clients.dedup();
906    let mut capabilities = inputs.module.capabilities.clone();
907    normalize_strings(&mut capabilities);
908    ExtractionServicePlan {
909        service_id: service_id.clone(),
910        module_id: inputs.module.module_id.clone(),
911        workloads: vec![
912            ExtractionWorkloadPlan {
913                workload_id: format!("{service_id}-api"),
914                role: ExtractionWorkloadRole::Api,
915            },
916            ExtractionWorkloadPlan {
917                workload_id: format!("{service_id}-worker"),
918                role: ExtractionWorkloadRole::Worker,
919            },
920            ExtractionWorkloadPlan {
921                workload_id: format!("{service_id}-migration"),
922                role: ExtractionWorkloadRole::Migration,
923            },
924        ],
925        store: ExtractionStorePlan {
926            store_id,
927            engine: "postgres".to_owned(),
928            isolated: true,
929        },
930        contract_versions: contracts,
931        service_references,
932        generated_clients,
933        preserved_capabilities: capabilities,
934        preserved_surfaces: inputs.readiness_report.surfaces.clone(),
935    }
936}
937
938fn data_mapping(
939    inputs: &ExtractionPlanInputs,
940    destination_store: &str,
941) -> Result<ExtractionDataMapping, ExtractionPlanGenerationError> {
942    #[derive(Default)]
943    struct TableEvidence {
944        sources: BTreeSet<ExtractionDataEvidenceSource>,
945        cursors: BTreeSet<ExtractionCursorEvidence>,
946    }
947    let mut tables = BTreeMap::<(String, String), TableEvidence>::new();
948    for table in &inputs.readiness_report.service_data.tables {
949        let owner = table
950            .owner_module
951            .clone()
952            .unwrap_or_else(|| "unresolved".to_owned());
953        let item = tables.entry((table.table.clone(), owner)).or_default();
954        item.sources.insert(table.source.clone());
955        if let Some(cursor) = &table.cursor {
956            item.cursors.insert(cursor.clone());
957        }
958    }
959    let tables = tables
960        .into_iter()
961        .map(|((table, owner), evidence)| {
962            let cursors = evidence.cursors.into_iter().collect::<Vec<_>>();
963            ExtractionTableMapping {
964                source_table: table.clone(),
965                destination_table: table,
966                destination_store: destination_store.to_owned(),
967                owner_module: owner,
968                copy_mode: if cursors.iter().any(|cursor| cursor.trustworthy) {
969                    ExtractionCopyMode::OnlineCheckpointed
970                } else {
971                    ExtractionCopyMode::BoundedWritePause
972                },
973                evidence_sources: evidence.sources.into_iter().collect(),
974                cursors,
975            }
976        })
977        .collect::<Vec<_>>();
978    let evidence_digests = inputs
979        .evidence_digests
980        .iter()
981        .map(|evidence| (evidence.reference.as_str(), evidence.digest.as_str()))
982        .collect::<BTreeMap<_, _>>();
983    let mut migrations = Vec::new();
984    for migration in &inputs.readiness_report.service_data.migrations {
985        let references = migration
986            .evidence_references
987            .iter()
988            .filter_map(|reference| {
989                evidence_digests
990                    .get(reference.as_str())
991                    .map(|digest| (reference, *digest))
992            })
993            .collect::<Vec<_>>();
994        let [(source_reference, source_digest)] = references.as_slice() else {
995            return Err(generation_error(
996                ExtractionPlanGenerationIssueCode::InputInvalid,
997                format!(
998                    "Migration `{}` must resolve to exactly one digest-pinned source artifact.",
999                    migration.migration
1000                ),
1001                "Record the authoritative migration path and content digest as Extraction Plan evidence.",
1002            ));
1003        };
1004        migrations.push(ExtractionMigrationMapping {
1005            source_migration: migration.migration.clone(),
1006            source_reference: (*source_reference).clone(),
1007            source_digest: (*source_digest).to_owned(),
1008            destination_store: destination_store.to_owned(),
1009            owner_module: migration
1010                .owner_module
1011                .clone()
1012                .unwrap_or_else(|| "unresolved".to_owned()),
1013        });
1014    }
1015    migrations.sort();
1016    migrations.dedup();
1017    Ok(ExtractionDataMapping {
1018        store_engine: "postgres".to_owned(),
1019        destination_store: destination_store.to_owned(),
1020        tables,
1021        migrations,
1022    })
1023}
1024
1025fn pinned_inputs(
1026    inputs: &ExtractionPlanInputs,
1027    data_mapping: &ExtractionDataMapping,
1028) -> Result<Vec<ExtractionInputPin>, ExtractionPlanGenerationError> {
1029    let source_system_id = inputs
1030        .system
1031        .get("systemId")
1032        .and_then(Value::as_str)
1033        .unwrap_or("unknown");
1034    let mut pins = vec![
1035        serializable_pin(
1036            ExtractionInputPinKind::ReadinessEvidence,
1037            &inputs.module.module_id,
1038            &inputs.readiness_report,
1039        )?,
1040        serializable_pin(
1041            ExtractionInputPinKind::ModuleDeclaration,
1042            &inputs.module.module_id,
1043            &inputs.module,
1044        )?,
1045        serializable_pin(
1046            ExtractionInputPinKind::SystemGraph,
1047            source_system_id,
1048            &inputs.system,
1049        )?,
1050        serializable_pin(
1051            ExtractionInputPinKind::AnalyzerVersion,
1052            &inputs.readiness_report.analyzer_version,
1053            &inputs.readiness_report.analyzer_version,
1054        )?,
1055        serializable_pin(
1056            ExtractionInputPinKind::DataMapping,
1057            &inputs.module.module_id,
1058            data_mapping,
1059        )?,
1060        serializable_pin(
1061            ExtractionInputPinKind::AuthorityRevision,
1062            &inputs.expected_authority.owner_id,
1063            &inputs.expected_authority,
1064        )?,
1065    ];
1066    let mut contracts = inputs.contract_versions.clone();
1067    for contract in &mut contracts {
1068        normalize_strings(&mut contract.consumer_ids);
1069        contract.required_context.sort();
1070        contract.required_context.dedup();
1071    }
1072    contracts.sort();
1073    pins.extend(
1074        contracts
1075            .iter()
1076            .map(|contract| {
1077                serializable_pin(
1078                    ExtractionInputPinKind::ContractVersion,
1079                    &format!("{}@{}", contract.contract_id, contract.version),
1080                    contract,
1081                )
1082            })
1083            .collect::<Result<Vec<_>, _>>()?,
1084    );
1085    pins.extend(
1086        inputs
1087            .evidence_digests
1088            .iter()
1089            .map(|evidence| ExtractionInputPin {
1090                kind: ExtractionInputPinKind::Evidence,
1091                subject: evidence.reference.clone(),
1092                digest: evidence.digest.clone(),
1093            }),
1094    );
1095    pins.sort();
1096    pins.dedup();
1097    Ok(pins)
1098}
1099
1100fn serializable_pin<T: Serialize>(
1101    kind: ExtractionInputPinKind,
1102    subject: &str,
1103    value: &T,
1104) -> Result<ExtractionInputPin, ExtractionPlanGenerationError> {
1105    Ok(ExtractionInputPin {
1106        kind,
1107        subject: subject.to_owned(),
1108        digest: digest_serializable(value)?,
1109    })
1110}
1111
1112fn plan_diff(inputs: &ExtractionPlanInputs, service: &ExtractionServicePlan) -> ExtractionPlanDiff {
1113    let mut entries = vec![
1114        ExtractionPlanDiffEntry {
1115            subject: format!("system.host.modules.{}", inputs.module.module_id),
1116            before: Some(inputs.expected_authority.owner_id.clone()),
1117            after: None,
1118        },
1119        ExtractionPlanDiffEntry {
1120            subject: format!("system.autonomousServices.{}", service.service_id),
1121            before: None,
1122            after: Some(format!("module={}", inputs.module.module_id)),
1123        },
1124        ExtractionPlanDiffEntry {
1125            subject: format!("authority.module.{}", inputs.module.module_id),
1126            before: Some(format!(
1127                "linked_host:{}@{}",
1128                inputs.expected_authority.owner_id, inputs.expected_authority.revision
1129            )),
1130            after: Some(format!("autonomous_service:{}", service.service_id)),
1131        },
1132        ExtractionPlanDiffEntry {
1133            subject: format!("store.{}", service.store.store_id),
1134            before: None,
1135            after: Some("postgres:isolated".to_owned()),
1136        },
1137    ];
1138    entries.extend(
1139        service
1140            .workloads
1141            .iter()
1142            .map(|workload| ExtractionPlanDiffEntry {
1143                subject: format!("workload.{}", workload.workload_id),
1144                before: None,
1145                after: Some(format!("{}:{:?}", service.service_id, workload.role).to_lowercase()),
1146            }),
1147    );
1148    entries.extend(
1149        service
1150            .service_references
1151            .iter()
1152            .map(|reference| ExtractionPlanDiffEntry {
1153                subject: format!("serviceReference.{}", reference.reference_id),
1154                before: None,
1155                after: Some(format!(
1156                    "{}@{}:{}",
1157                    reference.contract_id, reference.version, reference.target_service_id
1158                )),
1159            }),
1160    );
1161    entries.extend(
1162        service
1163            .generated_clients
1164            .iter()
1165            .map(|client| ExtractionPlanDiffEntry {
1166                subject: format!("generatedClient.{}", client.client_id),
1167                before: None,
1168                after: Some(format!("{}@{}", client.contract_id, client.version)),
1169            }),
1170    );
1171    entries.extend(
1172        service
1173            .contract_versions
1174            .iter()
1175            .filter(|contract| contract.direction == ExtractionContractDirection::Provides)
1176            .map(|contract| ExtractionPlanDiffEntry {
1177                subject: format!(
1178                    "contractProducer.{}@{}",
1179                    contract.contract_id, contract.version
1180                ),
1181                before: Some(format!("host:{}", inputs.expected_authority.owner_id)),
1182                after: Some(format!("autonomous_service:{}", service.service_id)),
1183            }),
1184    );
1185    entries.sort();
1186    entries.dedup();
1187    ExtractionPlanDiff { entries }
1188}
1189
1190#[allow(clippy::too_many_lines)]
1191fn plan_phases(
1192    service: &ExtractionServicePlan,
1193    data_mapping: &ExtractionDataMapping,
1194) -> Vec<ExtractionPlanPhase> {
1195    let full_pause_tables = data_mapping
1196        .tables
1197        .iter()
1198        .filter(|table| table.copy_mode == ExtractionCopyMode::BoundedWritePause)
1199        .map(|table| table.source_table.as_str())
1200        .collect::<Vec<_>>();
1201    let backfill_action = if full_pause_tables.is_empty() {
1202        "Copy ordered idempotent batches through pinned trustworthy cursors and durable checkpoints.".to_owned()
1203    } else {
1204        format!(
1205            "Keep full-copy tables blocked until the bounded write pause: {}.",
1206            full_pause_tables.join(", ")
1207        )
1208    };
1209    let approval = ExtractionApprovalBoundary {
1210        boundary_id: "commit-extraction-authority".to_owned(),
1211        phase_id: "09-rollback-or-commit".to_owned(),
1212        action: "commit_authority_to_autonomous_service".to_owned(),
1213        reason: "Final ownership transfer and write reopening are irreversible without a separately reviewed reverse-migration plan.".to_owned(),
1214        required_pins: vec![
1215            ExtractionInputPinKind::ReadinessEvidence,
1216            ExtractionInputPinKind::ContractVersion,
1217            ExtractionInputPinKind::SystemGraph,
1218            ExtractionInputPinKind::DataMapping,
1219            ExtractionInputPinKind::AuthorityRevision,
1220            ExtractionInputPinKind::Evidence,
1221        ],
1222    };
1223    let provisional_cutover_action = format!(
1224        "Route declared verification traffic to candidate Service `{}` without admitting authoritative mutations.",
1225        service.service_id
1226    );
1227    vec![
1228        phase(
1229            1,
1230            ExtractionPlanPhaseKind::Analysis,
1231            Vec::new(),
1232            vec!["The target Module is linked, ready, and owned by the pinned Host authority."],
1233            Vec::new(),
1234            vec!["Fresh input digests and a content-addressed Extraction Plan."],
1235            vec!["No mutation occurs; regenerate the plan if any pinned input changes."],
1236            vec![
1237                ExtractionPlanIssueCode::ReadinessEvidenceChanged,
1238                ExtractionPlanIssueCode::ModuleDeclarationChanged,
1239                ExtractionPlanIssueCode::ContractVersionChanged,
1240                ExtractionPlanIssueCode::SystemGraphChanged,
1241                ExtractionPlanIssueCode::AnalyzerVersionChanged,
1242                ExtractionPlanIssueCode::DataMappingChanged,
1243                ExtractionPlanIssueCode::AuthorityRevisionChanged,
1244                ExtractionPlanIssueCode::InputEvidenceChanged,
1245            ],
1246            vec!["Review the exact plan, diff, risks, and Approval Boundary."],
1247            None,
1248        ),
1249        phase(
1250            2,
1251            ExtractionPlanPhaseKind::Scaffold,
1252            vec!["01-analysis"],
1253            vec!["The exact plan is fresh and the deterministic scaffold patch has been reviewed."],
1254            vec![
1255                "Write the candidate API, Worker, and Migration Workload scaffold plus generated bindings and clients.",
1256            ],
1257            vec![
1258                "A deterministic patch, file digests, compile evidence, and identity-preservation evidence.",
1259            ],
1260            vec!["Remove only plan-owned generated files; refuse changed or unrecognized files."],
1261            vec![ExtractionPlanIssueCode::ScaffoldConflict],
1262            vec!["Apply the scaffold without changing linked authority."],
1263            None,
1264        ),
1265        phase(
1266            3,
1267            ExtractionPlanPhaseKind::DestinationExpansion,
1268            vec!["02-scaffold"],
1269            vec![
1270                "The candidate Migration Workload is valid and the destination Store is isolated.",
1271            ],
1272            vec![
1273                "Create the isolated Service Store and apply expand-first destination schema changes.",
1274            ],
1275            vec!["Idempotent migration receipts and candidate health evidence."],
1276            vec!["Discard the candidate Store; never contract or delete source schema."],
1277            vec![ExtractionPlanIssueCode::DestinationExpansionFailed],
1278            vec!["Verify destination schema compatibility before copying Service Data."],
1279            None,
1280        ),
1281        phase(
1282            4,
1283            ExtractionPlanPhaseKind::Backfill,
1284            vec!["03-destination-expansion"],
1285            vec![
1286                "Destination expansion succeeded and every online table has a pinned trustworthy cursor.",
1287            ],
1288            vec![backfill_action.as_str()],
1289            vec![
1290                "Durable source high-water marks, destination checkpoints, counts, and batch digests.",
1291            ],
1292            vec![
1293                "Stop copying and retain the linked implementation as the sole authoritative writer.",
1294            ],
1295            vec![ExtractionPlanIssueCode::BackfillCheckpointStale],
1296            vec!["Resume only from validated checkpoints, then reconcile the copied state."],
1297            None,
1298        ),
1299        phase(
1300            5,
1301            ExtractionPlanPhaseKind::Reconciliation,
1302            vec!["04-backfill"],
1303            vec!["Backfill checkpoint and source high-water mark are stable."],
1304            Vec::new(),
1305            vec![
1306                "Matching identities, counts, field digests, relationships, and declared business invariants.",
1307            ],
1308            vec![
1309                "Keep linked authority and remediate or repeat backfill when reconciliation differs.",
1310            ],
1311            vec![ExtractionPlanIssueCode::ReconciliationMismatch],
1312            vec!["Record reconciliation evidence bound to the exact plan and checkpoint."],
1313            None,
1314        ),
1315        phase(
1316            6,
1317            ExtractionPlanPhaseKind::Drain,
1318            vec!["05-reconciliation"],
1319            vec!["Candidate readiness and pre-pause reconciliation passed."],
1320            vec![
1321                "Pause new Module mutations, drain requests, Inbox, Outbox, schedules, and Workflows, then copy the final delta.",
1322            ],
1323            vec![
1324                "Source quiescence, drained-work evidence, final checkpoint, and final reconciliation.",
1325            ],
1326            vec![
1327                "Reopen linked writes before provisional routing if drain or final reconciliation fails.",
1328            ],
1329            vec![ExtractionPlanIssueCode::DrainIncomplete],
1330            vec!["Proceed only while external authoritative mutations remain paused."],
1331            None,
1332        ),
1333        phase(
1334            7,
1335            ExtractionPlanPhaseKind::ProvisionalCutover,
1336            vec!["06-drain"],
1337            vec!["The source is quiescent, work is drained, and the final delta reconciles."],
1338            vec![provisional_cutover_action.as_str()],
1339            vec!["Candidate health, routing, compatibility, and provisional authority evidence."],
1340            vec!["Restore linked routing and authority without reverse data movement."],
1341            vec![ExtractionPlanIssueCode::ProvisionalCutoverFailed],
1342            vec!["Run behavior, Contract, policy, health, and Runtime Story verification."],
1343            None,
1344        ),
1345        phase(
1346            8,
1347            ExtractionPlanPhaseKind::Verification,
1348            vec!["07-provisional-cutover"],
1349            vec![
1350                "Provisional routing is active and all external authoritative mutations remain paused.",
1351            ],
1352            Vec::new(),
1353            vec![
1354                "Compatibility, policy, business scenario, durable state, Event, Workflow, and Runtime Story comparison evidence.",
1355            ],
1356            vec!["Rollback provisional routing on any mismatch or stale input."],
1357            vec![ExtractionPlanIssueCode::VerificationFailed],
1358            vec!["Choose rollback on failure or request exact final commit approval on success."],
1359            None,
1360        ),
1361        phase(
1362            9,
1363            ExtractionPlanPhaseKind::RollbackOrCommit,
1364            vec!["08-verification"],
1365            vec![
1366                "Verification is terminal, the plan is fresh, and the authority revision still matches.",
1367            ],
1368            vec![
1369                "Either restore linked routing and writes, or compare-and-set authority and topology to the candidate before reopening writes.",
1370            ],
1371            vec![
1372                "Rollback evidence, or verified approval plus one-owner authority and topology evidence.",
1373            ],
1374            vec![
1375                "Before commit, restore linked authority; after new Autonomous writes, block fast rollback without a reviewed reverse plan.",
1376            ],
1377            vec![
1378                ExtractionPlanIssueCode::RollbackRequired,
1379                ExtractionPlanIssueCode::FinalApprovalRequired,
1380            ],
1381            vec!["Stop at the human Approval Boundary before final ownership transfer."],
1382            Some(approval),
1383        ),
1384        phase(
1385            10,
1386            ExtractionPlanPhaseKind::TerminalEvidence,
1387            vec!["09-rollback-or-commit"],
1388            vec!["Rollback or commit reached one unambiguous authoritative owner."],
1389            Vec::new(),
1390            vec![
1391                "Terminal plan, phase receipts, evidence, authority, topology, rollback constraints, and next actions.",
1392            ],
1393            vec!["Do not erase source data, linked recovery state, or audit evidence."],
1394            vec![ExtractionPlanIssueCode::TerminalEvidenceIncomplete],
1395            vec![
1396                "Publish the versioned terminal extraction evidence for operators and automation.",
1397            ],
1398            None,
1399        ),
1400    ]
1401}
1402
1403#[allow(clippy::too_many_arguments)]
1404fn phase(
1405    order: u16,
1406    kind: ExtractionPlanPhaseKind,
1407    prerequisite_phase_ids: Vec<&str>,
1408    prerequisites: Vec<&str>,
1409    intended_mutations: Vec<&str>,
1410    expected_evidence: Vec<&str>,
1411    rollback_conditions: Vec<&str>,
1412    issue_codes: Vec<ExtractionPlanIssueCode>,
1413    next_actions: Vec<&str>,
1414    approval_boundary: Option<ExtractionApprovalBoundary>,
1415) -> ExtractionPlanPhase {
1416    let label = match kind {
1417        ExtractionPlanPhaseKind::Analysis => "analysis",
1418        ExtractionPlanPhaseKind::Scaffold => "scaffold",
1419        ExtractionPlanPhaseKind::DestinationExpansion => "destination-expansion",
1420        ExtractionPlanPhaseKind::Backfill => "backfill",
1421        ExtractionPlanPhaseKind::Reconciliation => "reconciliation",
1422        ExtractionPlanPhaseKind::Drain => "drain",
1423        ExtractionPlanPhaseKind::ProvisionalCutover => "provisional-cutover",
1424        ExtractionPlanPhaseKind::Verification => "verification",
1425        ExtractionPlanPhaseKind::RollbackOrCommit => "rollback-or-commit",
1426        ExtractionPlanPhaseKind::TerminalEvidence => "terminal-evidence",
1427    };
1428    ExtractionPlanPhase {
1429        phase_id: format!("{order:02}-{label}"),
1430        order,
1431        kind,
1432        prerequisite_phase_ids: owned(prerequisite_phase_ids),
1433        prerequisites: owned(prerequisites),
1434        intended_mutations: owned(intended_mutations),
1435        expected_evidence: owned(expected_evidence),
1436        rollback_conditions: owned(rollback_conditions),
1437        issue_codes,
1438        next_actions: owned(next_actions),
1439        approval_boundary,
1440    }
1441}
1442
1443fn owned(values: Vec<&str>) -> Vec<String> {
1444    values.into_iter().map(str::to_owned).collect()
1445}
1446
1447fn stale_issue_code(kind: ExtractionInputPinKind) -> ExtractionPlanIssueCode {
1448    match kind {
1449        ExtractionInputPinKind::ReadinessEvidence => {
1450            ExtractionPlanIssueCode::ReadinessEvidenceChanged
1451        }
1452        ExtractionInputPinKind::ModuleDeclaration => {
1453            ExtractionPlanIssueCode::ModuleDeclarationChanged
1454        }
1455        ExtractionInputPinKind::ContractVersion => ExtractionPlanIssueCode::ContractVersionChanged,
1456        ExtractionInputPinKind::SystemGraph => ExtractionPlanIssueCode::SystemGraphChanged,
1457        ExtractionInputPinKind::AnalyzerVersion => ExtractionPlanIssueCode::AnalyzerVersionChanged,
1458        ExtractionInputPinKind::DataMapping => ExtractionPlanIssueCode::DataMappingChanged,
1459        ExtractionInputPinKind::AuthorityRevision => {
1460            ExtractionPlanIssueCode::AuthorityRevisionChanged
1461        }
1462        ExtractionInputPinKind::Evidence => ExtractionPlanIssueCode::InputEvidenceChanged,
1463    }
1464}
1465
1466fn digest_serializable<T: Serialize>(value: &T) -> Result<String, ExtractionPlanGenerationError> {
1467    serde_json::to_vec(value)
1468        .map(extraction_input_digest)
1469        .map_err(|error| {
1470            generation_error(
1471                ExtractionPlanGenerationIssueCode::InputInvalid,
1472                format!("Extraction input could not be serialized deterministically: {error}"),
1473                "Correct the structured input and regenerate the Extraction Plan.",
1474            )
1475        })
1476}
1477
1478fn generation_error(
1479    code: ExtractionPlanGenerationIssueCode,
1480    message: impl Into<String>,
1481    next_action: impl Into<String>,
1482) -> ExtractionPlanGenerationError {
1483    ExtractionPlanGenerationError {
1484        code,
1485        message: message.into(),
1486        next_actions: vec![next_action.into()],
1487    }
1488}
1489
1490fn valid_sha256_digest(value: &str) -> bool {
1491    value.strip_prefix("sha256:").is_some_and(|digest| {
1492        digest.len() == 64
1493            && digest
1494                .bytes()
1495                .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1496    })
1497}
1498
1499fn stable_slug(value: &str) -> String {
1500    let mut slug = String::new();
1501    let mut separator = false;
1502    for character in value.chars() {
1503        if character.is_ascii_alphanumeric() {
1504            if separator && !slug.is_empty() {
1505                slug.push('-');
1506            }
1507            slug.push(character.to_ascii_lowercase());
1508            separator = false;
1509        } else {
1510            separator = true;
1511        }
1512    }
1513    slug
1514}
1515
1516fn normalize_strings(values: &mut Vec<String>) {
1517    values.retain(|value| !value.trim().is_empty());
1518    values.sort();
1519    values.dedup();
1520}
1521
1522#[must_use]
1523pub fn render_extraction_plan(plan: &ExtractionPlan) -> String {
1524    let mut output = vec![
1525        format!("Extraction plan: {}", plan.target_module),
1526        format!("Plan ID: {}", plan.plan_id),
1527        format!("System: {}", plan.source_system_id),
1528        format!(
1529            "Expected authority: {}:{}@{}",
1530            serialized_label(plan.expected_authority.kind),
1531            plan.expected_authority.owner_id,
1532            plan.expected_authority.revision
1533        ),
1534        format!("Candidate Service: {}", plan.proposed_service.service_id),
1535        format!(
1536            "Workloads: {}",
1537            plan.proposed_service
1538                .workloads
1539                .iter()
1540                .map(|workload| {
1541                    format!(
1542                        "{} ({})",
1543                        workload.workload_id,
1544                        serialized_label(workload.role)
1545                    )
1546                })
1547                .collect::<Vec<_>>()
1548                .join(", ")
1549        ),
1550        format!(
1551            "Store: {} (postgres, isolated)",
1552            plan.proposed_service.store.store_id
1553        ),
1554        "Effects: dry-run; writesRepositoryFiles=false; startsWorkloads=false; copiesData=false; changesAuthority=false".to_owned(),
1555        "Pinned inputs:".to_owned(),
1556    ];
1557    output.extend(plan.pinned_inputs.iter().map(|pin| {
1558        format!(
1559            "- {} {}: {}",
1560            serialized_label(pin.kind),
1561            pin.subject,
1562            pin.digest
1563        )
1564    }));
1565    output.push("Diff:".to_owned());
1566    output.extend(plan.diff.entries.iter().map(|entry| {
1567        format!(
1568            "- {}: {} -> {}",
1569            entry.subject,
1570            entry.before.as_deref().unwrap_or("absent"),
1571            entry.after.as_deref().unwrap_or("absent")
1572        )
1573    }));
1574    output.push("Phases:".to_owned());
1575    for phase in &plan.phases {
1576        output.push(format!(
1577            "- {} {}",
1578            phase.phase_id,
1579            serialized_label(phase.kind)
1580        ));
1581        for prerequisite in &phase.prerequisites {
1582            output.push(format!("  prerequisite: {prerequisite}"));
1583        }
1584        if phase.intended_mutations.is_empty() {
1585            output.push("  mutation: none".to_owned());
1586        } else {
1587            for mutation in &phase.intended_mutations {
1588                output.push(format!("  mutation: {mutation}"));
1589            }
1590        }
1591        for evidence in &phase.expected_evidence {
1592            output.push(format!("  evidence: {evidence}"));
1593        }
1594        for condition in &phase.rollback_conditions {
1595            output.push(format!("  rollback: {condition}"));
1596        }
1597        output.push(format!(
1598            "  issueCodes: {}",
1599            phase
1600                .issue_codes
1601                .iter()
1602                .map(|code| serialized_label(*code))
1603                .collect::<Vec<_>>()
1604                .join(", ")
1605        ));
1606        for action in &phase.next_actions {
1607            output.push(format!("  next: {action}"));
1608        }
1609        if let Some(boundary) = &phase.approval_boundary {
1610            output.push(format!("  approvalBoundary: {}", boundary.boundary_id));
1611        }
1612    }
1613    output.push(String::new());
1614    output.join("\n")
1615}
1616
1617fn serialized_label<T: Serialize>(value: T) -> String {
1618    serde_json::to_value(value)
1619        .ok()
1620        .and_then(|value| value.as_str().map(str::to_owned))
1621        .unwrap_or_else(|| "unknown".to_owned())
1622}
1623
1624pub fn extraction_plan_json(plan: &ExtractionPlan) -> Result<String, serde_json::Error> {
1625    serde_json::to_string_pretty(plan).map(|rendered| format!("{rendered}\n"))
1626}
1627
1628#[must_use]
1629pub fn extraction_plan_schema() -> Value {
1630    let mut schema = serde_json::to_value(schemars::schema_for!(ExtractionPlan))
1631        .expect("Extraction Plan schema must serialize");
1632    let object = schema
1633        .as_object_mut()
1634        .expect("Extraction Plan schema must be an object");
1635    object.insert(
1636        "$id".to_owned(),
1637        Value::String(EXTRACTION_PLAN_SCHEMA_ID.to_owned()),
1638    );
1639    object.insert(
1640        "title".to_owned(),
1641        Value::String("Lenso Extraction Plan v1".to_owned()),
1642    );
1643    schema["properties"]["protocol"] = json!({
1644        "type": "string",
1645        "const": EXTRACTION_PLAN_PROTOCOL
1646    });
1647    schema["properties"]["generatorVersion"] = json!({
1648        "type": "string",
1649        "const": EXTRACTION_PLAN_GENERATOR_VERSION
1650    });
1651    schema["properties"]["planId"] = json!({
1652        "type": "string",
1653        "pattern": "^extraction-plan:sha256:[0-9a-f]{64}$"
1654    });
1655    schema["properties"]["planDigest"] = json!({
1656        "type": "string",
1657        "pattern": "^sha256:[0-9a-f]{64}$"
1658    });
1659    for field in [
1660        "writesRepositoryFiles",
1661        "startsWorkloads",
1662        "copiesData",
1663        "changesAuthority",
1664    ] {
1665        schema["$defs"]["ExtractionPlanEffects"]["properties"][field] = json!({
1666            "type": "boolean",
1667            "const": false
1668        });
1669    }
1670    schema
1671}
1672
1673#[cfg(test)]
1674mod tests {
1675    use super::*;
1676    use crate::{
1677        EXTRACTION_READINESS_ANALYZER_VERSION, EXTRACTION_READINESS_REPORT_PROTOCOL,
1678        ExtractionContractEvidence, ExtractionDataTableEvidence, ExtractionEvidenceStatus,
1679        ExtractionReadinessEffects, ExtractionServiceDataEvidence,
1680    };
1681    use lenso_contracts::ModuleManifest;
1682
1683    fn inputs() -> ExtractionPlanInputs {
1684        let module = ModuleManifest::builder("acme/support-ticket")
1685            .capabilities(vec!["support.tickets.read".to_owned()])
1686            .build();
1687        let report = ExtractionReadinessReport {
1688            protocol: EXTRACTION_READINESS_REPORT_PROTOCOL.to_owned(),
1689            analyzer_version: EXTRACTION_READINESS_ANALYZER_VERSION.to_owned(),
1690            target_module: module.module_id.clone(),
1691            system_id: Some("support-system".to_owned()),
1692            target_owner: Some("support-host".to_owned()),
1693            classification: CompatibilityCategory::Safe,
1694            ready: true,
1695            issue_codes: Vec::new(),
1696            contract_evidence: Vec::new(),
1697            active_consumers: Vec::new(),
1698            surfaces: ExtractionReadinessSurfaceSummary::default(),
1699            service_data: ExtractionServiceDataEvidence {
1700                complete: true,
1701                ..ExtractionServiceDataEvidence::default()
1702            },
1703            findings: Vec::new(),
1704            effects: ExtractionReadinessEffects::default(),
1705        };
1706        ExtractionPlanInputs {
1707            readiness_report: report,
1708            module,
1709            system: json!({
1710                "protocol": "lenso.system.v2",
1711                "systemId": "support-system",
1712                "host": { "hostId": "support-host", "modules": ["acme/support-ticket"] },
1713                "providers": [{
1714                    "providerId": "notification-provider",
1715                    "modules": ["notification-gateway"]
1716                }],
1717                "autonomousServices": [{
1718                    "serviceId": "support-sla-service",
1719                    "modules": ["support-sla"],
1720                    "workloads": [{ "workloadId": "support-sla-api", "role": "api" }]
1721                }],
1722                "contracts": [{
1723                    "contractId": "support.sla-updated.v1",
1724                    "version": "v1",
1725                    "producerKind": "autonomous_service",
1726                    "producerId": "support-sla-service",
1727                    "artifact": {
1728                        "format": "json_schema",
1729                        "path": "contracts/events/support.sla-updated.v1.schema.json"
1730                    },
1731                    "tenancyMode": "required"
1732                }],
1733                "consumers": [{
1734                    "consumerId": "support-ticket-sla-updates",
1735                    "ownerKind": "host",
1736                    "ownerId": "support-host",
1737                    "contractId": "support.sla-updated.v1",
1738                    "tenancyMode": "required"
1739                }]
1740            }),
1741            contract_versions: vec![ExtractionPlanContractVersion {
1742                contract_id: "support-ticket-http.v1".to_owned(),
1743                version: "v1".to_owned(),
1744                kind: ExtractionContractKind::Service,
1745                direction: ExtractionContractDirection::Provides,
1746                artifact_reference: "contracts/openapi/support-ticket.v1.yaml".to_owned(),
1747                artifact_digest: extraction_input_digest(b"support-ticket-http-v1"),
1748                artifact_format: ExtractionContractArtifactFormat::Openapi,
1749                tenancy_mode: ServiceTenancyMode::Required,
1750                required_context: vec![CommonContextRequirement::Tenant],
1751                producer_id: None,
1752                consumer_ids: vec!["support-portal".to_owned()],
1753            }],
1754            expected_authority: ExtractionExpectedAuthority {
1755                kind: ExtractionAuthorityKind::LinkedHost,
1756                owner_id: "support-host".to_owned(),
1757                revision: "authority-7".to_owned(),
1758            },
1759            evidence_digests: vec![ExtractionEvidenceDigest {
1760                reference: "analyzer:rust/support-ticket".to_owned(),
1761                digest: extraction_input_digest(b"boundary-clean"),
1762            }],
1763        }
1764    }
1765
1766    #[test]
1767    fn plan_is_content_addressed_and_dry_run_is_exact() {
1768        let inputs = inputs();
1769        let plan = generate_extraction_plan(&inputs).expect("plan should generate");
1770        let dry_run = dry_run_extraction_plan(&inputs).expect("dry run should generate");
1771
1772        assert_eq!(plan, dry_run);
1773        assert!(extraction_plan_integrity_is_valid(&plan));
1774        assert_eq!(plan.effects, ExtractionPlanEffects::default());
1775        assert_eq!(
1776            plan.proposed_service
1777                .workloads
1778                .iter()
1779                .map(|workload| workload.role)
1780                .collect::<Vec<_>>(),
1781            vec![
1782                ExtractionWorkloadRole::Api,
1783                ExtractionWorkloadRole::Worker,
1784                ExtractionWorkloadRole::Migration,
1785            ]
1786        );
1787        assert!(plan.proposed_service.store.isolated);
1788    }
1789
1790    #[test]
1791    fn plan_has_the_complete_ordered_phase_protocol() {
1792        let plan = generate_extraction_plan(&inputs()).expect("plan should generate");
1793        assert_eq!(
1794            plan.phases
1795                .iter()
1796                .map(|phase| phase.kind)
1797                .collect::<Vec<_>>(),
1798            vec![
1799                ExtractionPlanPhaseKind::Analysis,
1800                ExtractionPlanPhaseKind::Scaffold,
1801                ExtractionPlanPhaseKind::DestinationExpansion,
1802                ExtractionPlanPhaseKind::Backfill,
1803                ExtractionPlanPhaseKind::Reconciliation,
1804                ExtractionPlanPhaseKind::Drain,
1805                ExtractionPlanPhaseKind::ProvisionalCutover,
1806                ExtractionPlanPhaseKind::Verification,
1807                ExtractionPlanPhaseKind::RollbackOrCommit,
1808                ExtractionPlanPhaseKind::TerminalEvidence,
1809            ]
1810        );
1811        assert!(plan.phases.iter().all(|phase| {
1812            !phase.prerequisites.is_empty()
1813                && !phase.expected_evidence.is_empty()
1814                && !phase.rollback_conditions.is_empty()
1815                && !phase.issue_codes.is_empty()
1816                && !phase.next_actions.is_empty()
1817        }));
1818        assert_eq!(plan.approval_boundaries.len(), 1);
1819        assert_eq!(
1820            plan.approval_boundaries[0].action,
1821            "commit_authority_to_autonomous_service"
1822        );
1823    }
1824
1825    #[test]
1826    fn input_order_does_not_change_plan_identity() {
1827        let mut left = inputs();
1828        left.contract_versions.push(ExtractionPlanContractVersion {
1829            contract_id: "support-sla-grpc.v1".to_owned(),
1830            version: "v1".to_owned(),
1831            kind: ExtractionContractKind::Service,
1832            direction: ExtractionContractDirection::Consumes,
1833            artifact_reference: "contracts/services/support-sla.v1.proto".to_owned(),
1834            artifact_digest: extraction_input_digest(b"support-sla-grpc-v1"),
1835            artifact_format: ExtractionContractArtifactFormat::Protobuf,
1836            tenancy_mode: ServiceTenancyMode::Required,
1837            required_context: vec![CommonContextRequirement::Tenant],
1838            producer_id: Some("support-sla-service".to_owned()),
1839            consumer_ids: Vec::new(),
1840        });
1841        left.evidence_digests.push(ExtractionEvidenceDigest {
1842            reference: "store:host-postgres".to_owned(),
1843            digest: extraction_input_digest(b"read-only-observation"),
1844        });
1845        let mut right = left.clone();
1846        right.contract_versions.reverse();
1847        right.evidence_digests.reverse();
1848
1849        let left = generate_extraction_plan(&left).expect("left plan");
1850        let right = generate_extraction_plan(&right).expect("right plan");
1851        assert_eq!(left.plan_id, right.plan_id);
1852        assert_eq!(left, right);
1853    }
1854
1855    #[test]
1856    fn authority_or_evidence_drift_rejects_before_mutation() {
1857        let inputs = inputs();
1858        let plan = generate_extraction_plan(&inputs).expect("plan should generate");
1859        let mut changed = inputs.clone();
1860        changed.expected_authority.revision = "authority-8".to_owned();
1861        changed.evidence_digests[0].digest = extraction_input_digest(b"changed-evidence");
1862
1863        let rejection = ensure_extraction_plan_fresh(&plan, &changed)
1864            .expect_err("changed inputs must reject the stale plan");
1865        assert_eq!(rejection.effects, ExtractionPlanEffects::default());
1866        assert_eq!(
1867            rejection.issue_codes,
1868            vec![
1869                ExtractionPlanIssueCode::AuthorityRevisionChanged,
1870                ExtractionPlanIssueCode::InputEvidenceChanged,
1871            ]
1872        );
1873        assert_eq!(
1874            rejection
1875                .stale_inputs
1876                .iter()
1877                .map(|input| input.kind)
1878                .collect::<Vec<_>>(),
1879            vec![
1880                ExtractionInputPinKind::AuthorityRevision,
1881                ExtractionInputPinKind::Evidence,
1882            ]
1883        );
1884    }
1885
1886    #[test]
1887    fn modified_plan_fails_integrity_before_freshness() {
1888        let inputs = inputs();
1889        let mut plan = generate_extraction_plan(&inputs).expect("plan should generate");
1890        plan.diff.entries.clear();
1891
1892        let rejection = ensure_extraction_plan_fresh(&plan, &inputs)
1893            .expect_err("modified plans must fail integrity");
1894        assert_eq!(
1895            rejection.issue_codes,
1896            vec![ExtractionPlanIssueCode::PlanIntegrityInvalid]
1897        );
1898        assert!(rejection.stale_inputs.is_empty());
1899    }
1900
1901    #[test]
1902    fn every_pinned_input_category_rejects_drift() {
1903        let original = inputs();
1904        let plan = generate_extraction_plan(&original).expect("plan should generate");
1905
1906        let mut readiness = original.clone();
1907        readiness.readiness_report.system_id = Some("changed-evidence-system".to_owned());
1908        assert_stale_code(
1909            &plan,
1910            &readiness,
1911            ExtractionPlanIssueCode::ReadinessEvidenceChanged,
1912        );
1913
1914        let mut module = original.clone();
1915        module
1916            .module
1917            .capabilities
1918            .push("support.tickets.write".to_owned());
1919        assert_stale_code(
1920            &plan,
1921            &module,
1922            ExtractionPlanIssueCode::ModuleDeclarationChanged,
1923        );
1924
1925        let mut contract = original.clone();
1926        contract.contract_versions[0].artifact_digest =
1927            extraction_input_digest(b"changed-contract");
1928        assert_stale_code(
1929            &plan,
1930            &contract,
1931            ExtractionPlanIssueCode::ContractVersionChanged,
1932        );
1933
1934        let mut topology = original.clone();
1935        topology.system["host"]["modules"] = json!(["auth", "support-ticket"]);
1936        assert_stale_code(
1937            &plan,
1938            &topology,
1939            ExtractionPlanIssueCode::SystemGraphChanged,
1940        );
1941
1942        let mut analyzer = original.clone();
1943        analyzer.readiness_report.analyzer_version = "lenso.extraction-readiness.v3".to_owned();
1944        assert_stale_code(
1945            &plan,
1946            &analyzer,
1947            ExtractionPlanIssueCode::AnalyzerVersionChanged,
1948        );
1949
1950        let mut data = original.clone();
1951        data.readiness_report
1952            .service_data
1953            .tables
1954            .push(ExtractionDataTableEvidence {
1955                table: "support.tickets".to_owned(),
1956                owner_module: Some("support-ticket".to_owned()),
1957                source: ExtractionDataEvidenceSource::StaticDeclaration,
1958                volume: None,
1959                cursor: None,
1960                evidence_references: vec!["changed:data-mapping".to_owned()],
1961            });
1962        assert_stale_code(&plan, &data, ExtractionPlanIssueCode::DataMappingChanged);
1963
1964        let mut authority = original.clone();
1965        authority.expected_authority.revision = "authority-8".to_owned();
1966        assert_stale_code(
1967            &plan,
1968            &authority,
1969            ExtractionPlanIssueCode::AuthorityRevisionChanged,
1970        );
1971
1972        let mut evidence = original;
1973        evidence.evidence_digests[0].digest = extraction_input_digest(b"changed-evidence");
1974        assert_stale_code(
1975            &plan,
1976            &evidence,
1977            ExtractionPlanIssueCode::InputEvidenceChanged,
1978        );
1979    }
1980
1981    fn assert_stale_code(
1982        plan: &ExtractionPlan,
1983        inputs: &ExtractionPlanInputs,
1984        expected: ExtractionPlanIssueCode,
1985    ) {
1986        let rejection = ensure_extraction_plan_fresh(plan, inputs)
1987            .expect_err("changed pinned input must reject the plan");
1988        assert!(rejection.issue_codes.contains(&expected));
1989        assert_eq!(rejection.effects, ExtractionPlanEffects::default());
1990    }
1991
1992    #[test]
1993    fn plan_schema_accepts_public_json_and_v1_reader_ignores_future_fields() {
1994        let plan = generate_extraction_plan(&inputs()).expect("plan should generate");
1995        let value = serde_json::to_value(&plan).expect("plan should serialize");
1996        let validator = jsonschema::validator_for(&extraction_plan_schema())
1997            .expect("plan schema should compile");
1998        assert!(validator.is_valid(&value));
1999
2000        let mut future = value;
2001        future["futureField"] = json!(true);
2002        let decoded: ExtractionPlan =
2003            serde_json::from_value(future).expect("v1 reader should ignore future fields");
2004        assert_eq!(decoded.plan_id, plan.plan_id);
2005    }
2006
2007    #[test]
2008    fn blocked_readiness_cannot_generate_a_plan() {
2009        let mut inputs = inputs();
2010        inputs.readiness_report.ready = false;
2011        inputs.readiness_report.classification = CompatibilityCategory::Blocked;
2012
2013        let error = generate_extraction_plan(&inputs).expect_err("blocked readiness must fail");
2014        assert_eq!(
2015            error.code,
2016            ExtractionPlanGenerationIssueCode::ReadinessNotReady
2017        );
2018    }
2019
2020    #[test]
2021    fn every_contract_used_by_readiness_must_be_pinned() {
2022        let mut inputs = inputs();
2023        inputs
2024            .readiness_report
2025            .contract_evidence
2026            .push(ExtractionContractEvidence {
2027                subject: "event-handler:apply_sla_update".to_owned(),
2028                kind: ExtractionContractKind::Event,
2029                direction: ExtractionContractDirection::Consumes,
2030                status: ExtractionEvidenceStatus::Present,
2031                contract_id: Some("support.sla-updated.v1".to_owned()),
2032                evidence_references: vec![
2033                    "contracts/events/support.sla-updated.v1.schema.json".to_owned(),
2034                ],
2035            });
2036
2037        let error = generate_extraction_plan(&inputs)
2038            .expect_err("readiness Contract Versions must be pinned");
2039        assert_eq!(
2040            error.code,
2041            ExtractionPlanGenerationIssueCode::ContractVersionsMissing
2042        );
2043    }
2044}