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, ¤t_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}