1use std::collections::BTreeSet;
2
3use schemars::JsonSchema;
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6use utoipa::ToSchema;
7
8use crate::{DeliveryEffects, extraction_input_digest};
9
10pub const DELIVERY_FAILURE_RECOVERY_PROTOCOL: &str = "lenso.delivery-failure-recovery-evidence.v1";
11
12#[derive(
13 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema, ToSchema,
14)]
15#[serde(rename_all = "snake_case")]
16pub enum DeliveryFailureCondition {
17 DeploymentAdapterRejected,
18 OperatorReconciliationFailed,
19 GatewayDrift,
20 InvalidConfigRevision,
21 SecretReferenceUnavailable,
22 MigrationFailed,
23}
24
25#[derive(
26 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema, ToSchema,
27)]
28#[serde(rename_all = "snake_case")]
29pub enum DeliveryFailureStage {
30 BeforeApply,
31 Reconciling,
32 PartiallyApplied,
33}
34
35#[derive(
36 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema, ToSchema,
37)]
38#[serde(rename_all = "snake_case")]
39pub enum DeliveryRecoveryScope {
40 Deterministic,
41 EnvironmentVerification,
42}
43
44#[derive(
45 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema, ToSchema,
46)]
47#[serde(rename_all = "snake_case")]
48pub enum DeliveryRecoveryDecision {
49 Passed,
50 Blocked,
51}
52
53#[derive(
54 Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize, JsonSchema, ToSchema,
55)]
56#[serde(rename_all = "snake_case")]
57pub enum DeliveryRecoveryIssueCode {
58 InputInvalid,
59 PreApplyMutation,
60 DesiredObservedStateConflated,
61 LastValidConfigurationLost,
62 MigrationEvidenceIncomplete,
63 MigrationEffectWouldRepeat,
64 EnvironmentEvidenceInvalid,
65 CleanupIncomplete,
66}
67
68impl DeliveryRecoveryIssueCode {
69 #[must_use]
70 pub const fn as_str(self) -> &'static str {
71 match self {
72 Self::InputInvalid => "delivery_recovery_input_invalid",
73 Self::PreApplyMutation => "delivery_recovery_pre_apply_mutation",
74 Self::DesiredObservedStateConflated => {
75 "delivery_recovery_desired_observed_state_conflated"
76 }
77 Self::LastValidConfigurationLost => "delivery_recovery_last_valid_configuration_lost",
78 Self::MigrationEvidenceIncomplete => "delivery_recovery_migration_evidence_incomplete",
79 Self::MigrationEffectWouldRepeat => "delivery_recovery_migration_effect_would_repeat",
80 Self::EnvironmentEvidenceInvalid => "delivery_recovery_environment_evidence_invalid",
81 Self::CleanupIncomplete => "delivery_recovery_cleanup_incomplete",
82 }
83 }
84}
85
86#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
87#[serde(rename_all = "camelCase")]
88pub struct DeliveryRecoveryIssue {
89 pub code: DeliveryRecoveryIssueCode,
90 pub message: String,
91 pub remediation: String,
92 pub next_actions: Vec<String>,
93}
94
95#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
96#[serde(rename_all = "camelCase")]
97pub struct DeliveryStateObservation {
98 pub observation_id: String,
99 pub source: String,
100 pub revision: u64,
101 pub desired_digest: String,
102 pub observed_digest: String,
103 pub fresh: bool,
104 pub drifted: bool,
105}
106
107#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
108#[serde(rename_all = "camelCase")]
109pub struct MigrationRecoveryEvidence {
110 pub migration_id: String,
111 pub completed_effects: Vec<String>,
112 pub remaining_steps: Vec<String>,
113 pub retry_steps: Vec<String>,
114 pub state_compatible: bool,
115 pub rollback_allowed: bool,
116 pub intervention_required: bool,
117}
118
119#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
120#[serde(rename_all = "camelCase")]
121pub struct KubernetesRecoveryObservation {
122 pub cluster_identity: String,
123 pub api_server_version: String,
124 pub operator_version: String,
125 pub gateway_adapter_version: String,
126 pub used_real_api: bool,
127 pub observed_resource_version: String,
128 pub evidence_digest: String,
129}
130
131#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
132#[serde(rename_all = "camelCase")]
133pub struct DeliveryFailureRecoveryInput {
134 pub scenario_id: String,
135 pub condition: DeliveryFailureCondition,
136 pub stage: DeliveryFailureStage,
137 pub scope: DeliveryRecoveryScope,
138 pub desired_state: DeliveryStateObservation,
139 pub observed_state: DeliveryStateObservation,
140 pub previous_valid_config_revision_id: String,
141 pub attempted_config_revision_id: String,
142 pub active_config_revision_id: String,
143 #[serde(default, skip_serializing_if = "Option::is_none")]
144 pub migration: Option<MigrationRecoveryEvidence>,
145 #[serde(default)]
146 pub infrastructure_mutations: Vec<String>,
147 #[serde(default, skip_serializing_if = "Option::is_none")]
148 pub kubernetes: Option<KubernetesRecoveryObservation>,
149 pub cleanup_complete: bool,
150}
151
152#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
153#[serde(rename_all = "camelCase")]
154pub struct DeliveryFailureRecoveryEvidence {
155 pub protocol: String,
156 pub evidence_id: String,
157 pub evidence_digest: String,
158 pub scenario_id: String,
159 pub condition: DeliveryFailureCondition,
160 pub stage: DeliveryFailureStage,
161 pub scope: DeliveryRecoveryScope,
162 pub desired_state: DeliveryStateObservation,
163 pub observed_state: DeliveryStateObservation,
164 pub retained_config_revision_id: String,
165 #[serde(default, skip_serializing_if = "Option::is_none")]
166 pub migration: Option<MigrationRecoveryEvidence>,
167 #[serde(default, skip_serializing_if = "Option::is_none")]
168 pub environment_observation: Option<KubernetesRecoveryObservation>,
169 pub decision: DeliveryRecoveryDecision,
170 pub issues: Vec<DeliveryRecoveryIssue>,
171 pub next_actions: Vec<String>,
172 pub effects: DeliveryEffects,
173}
174
175#[must_use]
176pub fn evaluate_delivery_failure_recovery(
177 mut input: DeliveryFailureRecoveryInput,
178) -> DeliveryFailureRecoveryEvidence {
179 input.infrastructure_mutations.sort();
180 if let Some(migration) = &mut input.migration {
181 migration.completed_effects.sort();
182 migration.completed_effects.dedup();
183 migration.remaining_steps.sort();
184 migration.remaining_steps.dedup();
185 migration.retry_steps.sort();
186 migration.retry_steps.dedup();
187 }
188
189 let mut issues = Vec::new();
190 if input.scenario_id.trim().is_empty()
191 || input.desired_state.observation_id.trim().is_empty()
192 || input.observed_state.observation_id.trim().is_empty()
193 || input.observed_state.source.trim().is_empty()
194 || input.observed_state.revision < input.desired_state.revision
195 || !valid_digest(&input.desired_state.desired_digest)
196 || !valid_digest(&input.observed_state.observed_digest)
197 {
198 issues.push(issue(
199 DeliveryRecoveryIssueCode::InputInvalid,
200 "Delivery recovery evidence is missing a stable identity, digest, or monotonic observation.",
201 "Collect exact desired and observed state from their authoritative boundaries.",
202 "Refresh the failure observation before choosing a recovery action.",
203 ));
204 }
205 if input.stage == DeliveryFailureStage::BeforeApply
206 && !input.infrastructure_mutations.is_empty()
207 {
208 issues.push(issue(
209 DeliveryRecoveryIssueCode::PreApplyMutation,
210 "A validation, trust, policy, compatibility, configuration, freshness, or approval failure mutated infrastructure before apply.",
211 "Keep every pre-apply failure zero-effect.",
212 "Restore the prior state and repeat validation before apply.",
213 ));
214 }
215 if input.condition == DeliveryFailureCondition::GatewayDrift
216 && (!input.observed_state.drifted
217 || input.desired_state.desired_digest == input.observed_state.observed_digest)
218 {
219 issues.push(issue(
220 DeliveryRecoveryIssueCode::DesiredObservedStateConflated,
221 "Gateway drift evidence does not preserve the difference between desired and observed state.",
222 "Record both states and their independent source revisions.",
223 "Refresh the gateway observation without overwriting newer observed state.",
224 ));
225 }
226 if input.active_config_revision_id != input.previous_valid_config_revision_id {
227 issues.push(issue(
228 DeliveryRecoveryIssueCode::LastValidConfigurationLost,
229 "The Service did not retain its last valid Config Revision after activation failed.",
230 "Leave the rejected revision inactive and keep the previous revision authoritative.",
231 "Restore the previous valid Config Revision before retrying activation.",
232 ));
233 }
234 if input.condition == DeliveryFailureCondition::MigrationFailed {
235 match &input.migration {
236 Some(migration)
237 if !migration.migration_id.trim().is_empty()
238 && !migration.completed_effects.is_empty()
239 && !migration.remaining_steps.is_empty()
240 && (migration.state_compatible || migration.intervention_required) =>
241 {
242 let completed = migration
243 .completed_effects
244 .iter()
245 .map(String::as_str)
246 .collect::<BTreeSet<_>>();
247 if migration
248 .retry_steps
249 .iter()
250 .any(|step| completed.contains(step.as_str()))
251 {
252 issues.push(issue(
253 DeliveryRecoveryIssueCode::MigrationEffectWouldRepeat,
254 "The recovery plan would repeat a completed Migration effect.",
255 "Resume only the remaining steps identified by durable receipts.",
256 "Remove completed effects from the retry set and rebuild the plan.",
257 ));
258 }
259 }
260 _ => issues.push(issue(
261 DeliveryRecoveryIssueCode::MigrationEvidenceIncomplete,
262 "Partial Migration evidence does not identify completed effects, remaining steps, compatibility, and intervention constraints.",
263 "Bind recovery to durable per-effect receipts and current schema compatibility.",
264 "Collect complete Migration evidence before resuming or rolling back.",
265 )),
266 }
267 }
268 if input.scope == DeliveryRecoveryScope::EnvironmentVerification {
269 let valid = input.kubernetes.as_ref().is_some_and(|observation| {
270 observation.used_real_api
271 && !observation.cluster_identity.trim().is_empty()
272 && !observation.api_server_version.trim().is_empty()
273 && !observation.operator_version.trim().is_empty()
274 && !observation.gateway_adapter_version.trim().is_empty()
275 && !observation.observed_resource_version.trim().is_empty()
276 && valid_digest(&observation.evidence_digest)
277 });
278 if !valid {
279 issues.push(issue(
280 DeliveryRecoveryIssueCode::EnvironmentEvidenceInvalid,
281 "Environment Verification lacks a real Kubernetes API, Operator, and gateway observation.",
282 "Use the pinned real-cluster verification lane rather than mock-only desired resources.",
283 "Repeat the scenario against the supported environment adapter set.",
284 ));
285 }
286 }
287 if !input.cleanup_complete {
288 issues.push(issue(
289 DeliveryRecoveryIssueCode::CleanupIncomplete,
290 "Delivery failure cleanup is incomplete.",
291 "Remove or isolate every disposable resource while preserving durable evidence.",
292 "Finish cleanup before accepting the recovery result.",
293 ));
294 }
295
296 let decision = if issues.is_empty() {
297 DeliveryRecoveryDecision::Passed
298 } else {
299 DeliveryRecoveryDecision::Blocked
300 };
301 let next_actions = if issues.is_empty() {
302 match input.condition {
303 DeliveryFailureCondition::GatewayDrift => {
304 vec!["Reconcile desired state against the newer observed revision.".to_owned()]
305 }
306 DeliveryFailureCondition::MigrationFailed => {
307 vec!["Resume only the receipt-backed remaining Migration steps.".to_owned()]
308 }
309 _ => vec!["Retry the idempotent operation from the preserved state.".to_owned()],
310 }
311 } else {
312 issues
313 .iter()
314 .flat_map(|issue| issue.next_actions.iter().cloned())
315 .collect()
316 };
317 let mut evidence = DeliveryFailureRecoveryEvidence {
318 protocol: DELIVERY_FAILURE_RECOVERY_PROTOCOL.to_owned(),
319 evidence_id: String::new(),
320 evidence_digest: String::new(),
321 scenario_id: input.scenario_id,
322 condition: input.condition,
323 stage: input.stage,
324 scope: input.scope,
325 desired_state: input.desired_state,
326 observed_state: input.observed_state,
327 retained_config_revision_id: input.active_config_revision_id,
328 migration: input.migration,
329 environment_observation: input.kubernetes,
330 decision,
331 issues,
332 next_actions,
333 effects: DeliveryEffects::default(),
334 };
335 evidence.evidence_digest = digest_without_identity(&evidence);
336 evidence.evidence_id = format!(
337 "delivery-failure-recovery:{}",
338 &evidence.evidence_digest[7..23]
339 );
340 evidence
341}
342
343#[must_use]
344pub fn delivery_failure_recovery_schema() -> Value {
345 let mut schema = serde_json::to_value(schemars::schema_for!(DeliveryFailureRecoveryEvidence))
346 .expect("delivery failure recovery schema serializes");
347 schema["$id"] = Value::String(
348 "https://contracts.lenso.local/ga/lenso.delivery-failure-recovery-evidence.v1.schema.json"
349 .to_owned(),
350 );
351 schema
352}
353
354fn issue(
355 code: DeliveryRecoveryIssueCode,
356 message: impl Into<String>,
357 remediation: impl Into<String>,
358 next_action: impl Into<String>,
359) -> DeliveryRecoveryIssue {
360 DeliveryRecoveryIssue {
361 code,
362 message: message.into(),
363 remediation: remediation.into(),
364 next_actions: vec![next_action.into()],
365 }
366}
367
368fn valid_digest(value: &str) -> bool {
369 value.strip_prefix("sha256:").is_some_and(|digest| {
370 digest.len() == 64
371 && digest
372 .bytes()
373 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
374 })
375}
376
377fn digest_without_identity(evidence: &DeliveryFailureRecoveryEvidence) -> String {
378 let mut canonical = evidence.clone();
379 canonical.evidence_id.clear();
380 canonical.evidence_digest.clear();
381 extraction_input_digest(
382 &serde_json::to_vec(&canonical).expect("delivery recovery evidence serializes"),
383 )
384}
385
386#[must_use]
387pub fn delivery_failure_recovery_integrity_is_valid(
388 evidence: &DeliveryFailureRecoveryEvidence,
389) -> bool {
390 valid_digest(&evidence.evidence_digest)
391 && evidence.evidence_digest == digest_without_identity(evidence)
392 && evidence.evidence_id
393 == format!(
394 "delivery-failure-recovery:{}",
395 &evidence.evidence_digest[7..23]
396 )
397}