Skip to main content

lenso_service/
delivery_failure_recovery.rs

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}