Skip to main content

ahri_tre_app/projections/
lifecycle.rs

1use super::*;
2use crate::service::{GovernedDeleteOutcome, GovernedDeletePlan};
3use ahri_tre_protocol::{
4    PublicUuid, lifecycle as wire, refs::ObjectKind, session::SessionStatusPayload,
5};
6
7fn lifecycle_plan(ds: PublicUuid, plan: &crate::ArchiveDeletePlan) -> serde_json::Value {
8    let target = match plan.target {
9        crate::DeleteTarget::Study(id) => object_ref(ds, ObjectKind::Study, id.0),
10        crate::DeleteTarget::Asset(id) => object_ref(ds, ObjectKind::Asset, id.0),
11        crate::DeleteTarget::AssetVersion(id) => object_ref(ds, ObjectKind::AssetVersion, id.0),
12    };
13    json!({
14        "target": target,
15        "resolved_version": plan.resolved_version.map(|id| object_ref(ds, ObjectKind::AssetVersion, id.0)),
16        "delete_all_versions": plan.delete_all_versions,
17        "archive_mode": plan.archive_mode, "archive_required": plan.archive_required,
18        "cascade": plan.cascade, "force": plan.force, "reason": plan.reason,
19        "actor": plan.actor, "authorization": plan.authorization,
20        "dependency_summary": {"status": plan.dependency_summary.status,
21            "items": plan.dependency_summary.items.iter().map(|item| json!({"kind": item.kind, "action": item.action})).collect::<Vec<_>>()},
22        "archive_preview": plan.archive_preview,
23        "metadata_cleanup_preview": plan.metadata_cleanup_preview,
24        "lake_cleanup_preview": plan.lake_cleanup_preview, "retry_state": plan.retry_state,
25    })
26}
27
28fn planned_delete(ds: PublicUuid, resolved: GovernedDeletePlan) -> serde_json::Value {
29    json!({
30        "accepted": resolved.planning.accepted,
31        "plan": resolved.planning.plan.as_ref().map(|plan| lifecycle_plan(ds, plan)),
32        "validation_errors": resolved.planning.validation_errors.into_iter().map(|error| json!({
33            "code": error.code, "message": "Deletion plan was rejected",
34        })).collect::<Vec<_>>()
35    })
36}
37
38fn result_summary(kind: wire::DeleteWorkflowResultKind) -> wire::DeleteWorkflowResultSummary {
39    wire::DeleteWorkflowResultSummary {
40        kind,
41        asset_name: None,
42        reason: None,
43        withdrawn_version: None,
44        lake_table_drop: None,
45        archive: None,
46        tombstone: None,
47        dependency_summary: None,
48        provenance: None,
49        removed_link_count: None,
50        promoted_latest_version: None,
51        deleted_assets: 0,
52        deleted_versions: 0,
53        skipped_versions: 0,
54        failed_versions: 0,
55        already_withdrawn: false,
56    }
57}
58
59fn semantic_plan(ds: PublicUuid, plan: &crate::SemanticDeletePlan) -> serde_json::Value {
60    use crate::SemanticDeleteTargetKind as Kind;
61    let (kind, scope) = match plan.target_identity.kind {
62        Kind::Domain => (ObjectKind::Domain, "domain"),
63        Kind::Variable => (ObjectKind::Variable, "variable"),
64        Kind::Vocabulary => (ObjectKind::Vocabulary, "vocabulary"),
65        Kind::Entity => (ObjectKind::Entity, "entity"),
66        Kind::EntityRelation => (ObjectKind::Relation, "relation"),
67    };
68    let target = ahri_tre_protocol::refs::ObjectRef {
69        datastore_id: ds,
70        kind,
71        id: ahri_tre_protocol::refs::encode_scoped_integer_ref(
72            scope,
73            plan.target_identity
74                .id
75                .parse()
76                .expect("resolved semantic ID"),
77        ),
78    };
79    json!({"target": target, "name": plan.target_identity.name, "domain": plan.target_identity.domain,
80        "authorization": plan.authorization, "reason": plan.reason, "actor": plan.actor,
81        "disposition": plan.disposition, "dependency_summary": semantic_dependencies(&plan.dependency_summary)})
82}
83
84fn semantic_dependencies(summary: &crate::SemanticDeleteDependencySummary) -> serde_json::Value {
85    // Dependency rows include composite links and unaddressable storage IDs.
86    // The description carries the qualified object names needed to resolve a blocker.
87    let values = |items: &[crate::SemanticDeleteDependency]| {
88        items
89            .iter()
90            .map(|item| json!({"kind": item.kind, "description": item.description}))
91            .collect::<Vec<_>>()
92    };
93    json!({"status": summary.status, "active_blockers": values(&summary.active_blockers),
94        "historical_blockers": values(&summary.historical_blockers), "unsupported": values(&summary.unsupported),
95        "delete_candidates": values(&summary.delete_candidates)})
96}
97
98pub fn governed_semantic_delete_response(
99    session: SessionStatusPayload,
100    outcome: crate::service::GovernedSemanticDeleteOutcome,
101) -> wire::DeleteWorkflowResponse {
102    use crate::service::GovernedSemanticDeleteOutcome as Outcome;
103    let ds = session
104        .datastore_id
105        .expect("semantic Session retains its Datastore binding");
106    let (status, planning, result, dependencies) = match outcome {
107        Outcome::Planned(planning) => (
108            wire::DeleteWorkflowStatus::Planned,
109            Some(json!({
110                "accepted": planning.accepted, "plan": planning.plan.as_ref().map(|p| semantic_plan(ds, p)),
111                "validation_errors": planning.validation_errors,
112            })),
113            None,
114            None,
115        ),
116        Outcome::Blocked(plan) => (
117            wire::DeleteWorkflowStatus::Blocked,
118            Some(json!({
119                "accepted": true, "plan": semantic_plan(ds, &plan), "validation_errors": [],
120            })),
121            None,
122            Some(semantic_dependencies(&plan.dependency_summary)),
123        ),
124        Outcome::Executed(executed) => {
125            let mut result = result_summary(wire::DeleteWorkflowResultKind::Semantic);
126            result.provenance = Some(semantic_plan(ds, &executed.plan));
127            result.tombstone = Some(wire::TombstoneSummary {
128                metadata_removed: executed.deleted,
129                lake_removed: false,
130                retry_state: None,
131            });
132            (
133                wire::DeleteWorkflowStatus::Deleted,
134                None,
135                Some(result),
136                None,
137            )
138        }
139    };
140    let blocked = status == wire::DeleteWorkflowStatus::Blocked;
141    wire::DeleteWorkflowResponse {
142        session,
143        status,
144        blocked,
145        blocked_reason: blocked.then(|| "Semantic deletion is dependency-blocked".into()),
146        next_step: blocked
147            .then(|| "Remove or migrate the listed dependencies and inspect a fresh plan".into()),
148        dependency_summary: dependencies,
149        planning,
150        result,
151        warnings: Vec::new(),
152    }
153}
154
155fn tombstone(
156    metadata_removed: bool,
157    lake_removed: bool,
158    retry_state: crate::ArchiveDeleteRetryState,
159) -> wire::TombstoneSummary {
160    wire::TombstoneSummary {
161        metadata_removed,
162        lake_removed,
163        retry_state: Some(json!(retry_state)),
164    }
165}
166
167fn asset_lake_removed(result: &crate::DeleteAssetResult) -> bool {
168    result.failed_versions == 0
169        && !result.version_results.is_empty()
170        && result.version_results.iter().all(|version| match version {
171            crate::DeleteAssetVersionResult::Dataset(result) => {
172                result.lake_cleanup == crate::LakeTableDropDisposition::DroppedIfExisted
173            }
174            crate::DeleteAssetVersionResult::Datafile(result) => matches!(
175                result.managed_file_cleanup,
176                crate::ArchiveDeleteCleanupDisposition::Completed
177                    | crate::ArchiveDeleteCleanupDisposition::Skipped
178            ),
179            crate::DeleteAssetVersionResult::Skipped { .. } => false,
180        })
181}
182
183fn withdrawal_summary(
184    ds: PublicUuid,
185    result: crate::WithdrawDatasetVersionResult,
186) -> wire::DeleteWorkflowResultSummary {
187    let mut summary = result_summary(wire::DeleteWorkflowResultKind::DatasetWithdrawal);
188    summary.asset_name = Some(result.catalog.asset.name.as_str().into());
189    summary.reason = Some(result.withdrawal.reason);
190    summary.withdrawn_version = Some(json!(safe_asset_version_summary(
191        ds,
192        result.withdrawn_version
193    )));
194    summary.promoted_latest_version = result
195        .promoted_latest_version
196        .map(|version| json!(safe_asset_version_summary(ds, version)));
197    summary.lake_table_drop = Some(
198        match result.lake_table_drop {
199            crate::LakeTableDropDisposition::NotRequested => "not_requested",
200            crate::LakeTableDropDisposition::DroppedIfExisted => "dropped_if_existed",
201            crate::LakeTableDropDisposition::AlreadyWithdrawn => "already_withdrawn",
202        }
203        .into(),
204    );
205    summary.tombstone = Some(wire::TombstoneSummary {
206        metadata_removed: false,
207        lake_removed: result.lake_table_drop == crate::LakeTableDropDisposition::DroppedIfExisted,
208        retry_state: None,
209    });
210    summary.already_withdrawn = result.already_withdrawn;
211    summary
212}
213
214/// Allowlisted lifecycle evidence. Archive manifests, raw records, physical
215/// locations and artifact bytes never enter a user-controlled response.
216pub fn governed_delete_response(
217    session: SessionStatusPayload,
218    result: GovernedDeleteOutcome,
219) -> wire::DeleteWorkflowResponse {
220    use crate::ArchiveDeleteCleanupDisposition::Completed;
221    let ds = session
222        .datastore_id
223        .expect("lifecycle Session retains its Datastore binding");
224    let mut response = wire::DeleteWorkflowResponse {
225        session,
226        status: wire::DeleteWorkflowStatus::Deleted,
227        blocked: false,
228        blocked_reason: None,
229        dependency_summary: None,
230        next_step: None,
231        planning: None,
232        result: None,
233        warnings: Vec::new(),
234    };
235    response.result = Some(match result {
236        GovernedDeleteOutcome::Planned(plan) => {
237            response.status = wire::DeleteWorkflowStatus::Planned;
238            response.planning = Some(planned_delete(ds, *plan));
239            return response;
240        }
241        GovernedDeleteOutcome::Study(result) => {
242            let mut summary = result_summary(wire::DeleteWorkflowResultKind::Study);
243            summary.reason = Some(result.plan.reason.clone());
244            summary.provenance = Some(lifecycle_plan(ds, &result.plan));
245            summary.archive = Some(wire::ArchiveSummary {
246                archived: true,
247                asset_name: Some(result.archive.archive_asset.name.as_str().into()),
248                version: Some(object_ref(
249                    ds,
250                    ObjectKind::AssetVersion,
251                    result.archive.archive_version.version_id.0,
252                )),
253                summary: None,
254            });
255            summary.deleted_assets = result
256                .asset_results
257                .iter()
258                .filter(|asset| asset.metadata_cleanup == Completed)
259                .count();
260            summary.deleted_versions = result.deleted_versions;
261            summary.skipped_versions = result.skipped_versions;
262            summary.failed_versions = result.failed_versions;
263            summary.tombstone = Some(tombstone(
264                result.metadata_cleanup == Completed,
265                !result.asset_results.is_empty()
266                    && result.asset_results.iter().all(asset_lake_removed),
267                result.retry_state,
268            ));
269            summary
270        }
271        GovernedDeleteOutcome::Asset(result) => {
272            let mut summary = result_summary(wire::DeleteWorkflowResultKind::Asset);
273            summary.asset_name = Some(result.asset.name.as_str().into());
274            summary.reason = Some(result.plan.reason.clone());
275            summary.provenance = Some(lifecycle_plan(ds, &result.plan));
276            summary.deleted_versions = result.deleted_versions;
277            summary.skipped_versions = result.skipped_versions;
278            summary.failed_versions = result.failed_versions;
279            summary.deleted_assets = usize::from(result.metadata_cleanup == Completed);
280            summary.tombstone = Some(tombstone(
281                result.metadata_cleanup == Completed,
282                asset_lake_removed(&result),
283                result.retry_state,
284            ));
285            summary
286        }
287        GovernedDeleteOutcome::DataFile(result) => {
288            let mut summary = result_summary(wire::DeleteWorkflowResultKind::DataFile);
289            summary.asset_name = Some(result.asset.name.as_str().into());
290            summary.reason = Some(result.plan.reason.clone());
291            summary.provenance = Some(lifecycle_plan(ds, &result.plan));
292            summary.deleted_versions = usize::from(result.metadata_cleanup == Completed);
293            summary.promoted_latest_version = result
294                .promoted_latest_version
295                .map(|version| json!(safe_asset_version_summary(ds, version)));
296            summary.tombstone = Some(tombstone(
297                result.metadata_cleanup == Completed,
298                result.managed_file_cleanup == Completed,
299                result.retry_state,
300            ));
301            summary
302        }
303        GovernedDeleteOutcome::Dataset(result) => {
304            let mut summary = withdrawal_summary(ds, result.deletion);
305            summary.kind = wire::DeleteWorkflowResultKind::Dataset;
306            summary.provenance = Some(lifecycle_plan(ds, &result.plan));
307            summary.deleted_versions = 1;
308            summary
309        }
310        GovernedDeleteOutcome::Withdrawal(result) => {
311            response.status = wire::DeleteWorkflowStatus::Withdrawn;
312            withdrawal_summary(ds, *result)
313        }
314    });
315    response
316}