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 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
214pub 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}