Setting the file. One moment.
Cleanup Plan · Foundry Iq · microsoft/azure-skills · Skills Docs
ContentsBack to the top of the page File Cu Canary
— line 271
This file
Number 10.41
Position 41 of 77
Type Python
Size 42 KB
Lines 735 helpers/ cleanup_plan.py
Python · 735 lines · 42 KB
.
import
prompt_cleanup, prompt_connect, search_reconcile, _cleanup_dependencies
as
dependencies, _cleanup_receipts
as
receipts
13 from ._common import (
14 SEARCH_AUDIENCE , MANAGEMENT_AUDIENCE , HelperFailure, TokenProvider, Transport, azure_cli_token,
15 blocked_result, digest, emit_result, http_request, load_approved_input,
16 reject_secrets, require_allowed_fields, sdk_error_metadata, sdk_error_status,
17 validate_search_endpoint,
18 )
19 except ImportError :
20 import prompt_cleanup
21 import prompt_connect
22 import search_reconcile
23 import _cleanup_dependencies as dependencies
24 import _cleanup_receipts as receipts
25 from _common import (
26 SEARCH_AUDIENCE , MANAGEMENT_AUDIENCE , HelperFailure, TokenProvider, Transport, azure_cli_token,
27 blocked_result, digest, emit_result, http_request, load_approved_input,
28 reject_secrets, require_allowed_fields, sdk_error_metadata, sdk_error_status,
29 validate_search_endpoint,
30 )
31
32
33 RETAIN = [
34 "outside-plan resources and consumers" , "original local and Storage documents" ,
35 "Search service" , "accounts" , "projects" , "models" , "role assignments" ,
36 ]
37 REQUEST_FIELDS = {
38 "schema_version" , "owner" , "target" , "creation_input_file" ,
39 "creation_result_file" , "creation_response_file" ,
40 "creation_receipt_file" , "inventory_limits" , "agent_version" , "agent_creation_response_file" , "agent_creation_receipt_file" ,
41 }
42 SEARCH_FIELDS = { "type" , "endpoint" , "api_version" , "name" }
43 PROMPT_FIELDS = { "type" , "project_resource_id" , "project_endpoint" , "name" , "version" }
44 RESPONSE_FIELDS = { "schema_version" , "target" , "operation" , "status" , "request_id" , "body" , "generated_resources" }
45
46
47 def _failure (code: str , message: str ) -> HelperFailure:
48 return HelperFailure(code, message, blocked_at = "cleanup-planning" )
49
50
51 def _object (value: Any, label: str ) -> dict[ str , Any]:
52 if not isinstance (value, dict ):
53 raise _failure( "input-schema-invalid" , f " { label } must be an object." )
54 return value
55
56
57 def _text (value: Any) -> bool :
58 return isinstance (value, str ) and bool (value.strip())
59
60
61 def _read_json (path: Path) -> dict[ str , Any]:
62 try :
63 result = json.loads(path.read_text( encoding = "utf-8" ))
64 json.dumps(result, allow_nan = False )
65 except ( OSError , UnicodeError , ValueError ) as exc:
66 raise _failure( "input-unreadable" , "Select readable retained UTF-8 JSON records." ) from exc
67 return _object(result, "Record" )
68
69
70 def _target (value: Any) -> dict[ str , Any]:
71 target = copy.deepcopy(_object(value, "Cleanup target" ))
72 kind = target.get( "type" )
73 if not isinstance (kind, str ):
74 raise _failure( "target-selection-required" , "Select one typed cleanup target." )
75 if kind == "hosted" :
76 require_allowed_fields(target, { "type" }, label = "Hosted cleanup target" )
77 raise _failure( "hosted-cleanup-unsupported" , "Hosted teardown remains unsupported; retain the deployment and toolbox." )
78 if kind in { "knowledge-base" , "knowledge-source" }:
79 require_allowed_fields(target, SEARCH_FIELDS , label = "Search cleanup target" )
80 search_reconcile.resource_url({ ** target, "resource_type" : kind})
81 target[ "endpoint" ] = validate_search_endpoint(target[ "endpoint" ])
82 elif kind in { "prompt-agent-version" , "project-connection" }:
83 fields = PROMPT_FIELDS if kind == "prompt-agent-version" else PROMPT_FIELDS - { "version" }
84 require_allowed_fields(target, fields, label = "Prompt cleanup target" )
85 project_id, endpoint = prompt_connect._project_identity(target)
86 target[ "project_resource_id" ] = project_id.casefold()
87 target[ "project_endpoint" ] = endpoint
88 if not _text(target.get( "name" )):
89 raise _failure( "target-selection-required" , "Select one exact agent or connection name." )
90 if kind == "prompt-agent-version" and (
91 not isinstance (target.get( "version" ), str )
92 or re.fullmatch( r " [ 1-9 ][ 0-9 ] * " , target[ "version" ]) is None
93 ):
94 raise _failure( "target-selection-required" , "Select one exact numeric Prompt version, not latest or a list." )
95 else :
96 raise _failure( "target-selection-required" , "Select one supported exact cleanup target." )
97 return target
98
99
100 def _path (request: dict[ str , Any], key: str , base_dir: Path) -> Path:
101 value = request.get(key)
102 if not _text(value):
103 raise _failure( "ownership-unproven" , "Retained approval, result and definitive create response files are required." )
104 path = Path(value)
105 return path if path.is_absolute() else base_dir / path
106
107
108 def _records (
109 request: dict[ str , Any], target: dict[ str , Any], base_dir: Path,
110 ) -> tuple[dict[ str , Any], dict[ str , Any], dict[ str , Any]]:
111 if "creation_receipt_file" in request:
112 if target[ "type" ] != "knowledge-source" or "creation_response_file" in request:
113 raise _failure( "input-schema-invalid" , "Select either an original Blob checkpoint or native result/response records." )
114 try :
115 from . import blob_recheck
116 except ImportError :
117 import blob_recheck
118 original = _read_json(_path(request, "creation_result_file" , base_dir))
119 reject_secrets(original)
120 retained = original.get( "recheck_checkpoint" )
121 if (original.get( "outcome" ) != "create-blob-knowledge-source" or original.get( "status" ) not in ( "completed" , "partial" )
122 or not isinstance (retained, dict ) or retained.get( "status" ) != "retained" ):
123 raise _failure(
124 "generated-creation-evidence-unavailable" ,
125 "Require the original creation result retaining its generated checkpoint; ACK/recovery GETs cannot prove original child versions." ,
126 )
127 prior, receipt = blob_recheck._load(
128 _path(request, "creation_input_file" , base_dir), _path(request, "creation_receipt_file" , base_dir),
129 )
130 if prior[ "source" ][ "action" ] != "create" or prior[ "owner" ] != request[ "owner" ]:
131 raise _failure( "ownership-unproven" , "Only an original run-owned creation checkpoint is eligible, never captured reuse." )
132 original_ownership = _object(
133 original.get( "resources_remaining" if original[ "status" ] == "partial" else "ownership" ),
134 "Original creation ownership" ,
135 )
136 original_owner = original.get( "owner" ) if original[ "status" ] == "partial" else original_ownership.get( "owner" )
137 original_owned = original_ownership.get( "run_owned" )
138 if (retained.get( "evidence_digest" ) != receipt[ "integrity" ]
139 or retained.get( "operation_id" ) != receipt[ "operation_id" ]
140 or original.get( "approved_plan" ) != { "confirmed" : True , "fingerprint" : receipt[ "plan_digest" ]}
141 or original_owner != prior[ "owner" ] or not isinstance (original_owned, list )
142 or receipt[ "creation" ][ "ownership" ][ "run_owned" ][ 0 ] not in original_owned):
143 raise _failure( "generated-ownership-unproven" , "The original creation result must bind this exact checkpoint and approval; never substitute recovery snapshots." )
144 result = { "status" : "completed" , ** copy.deepcopy(receipt[ "creation" ])}
145 generated = receipt[ "creation" ][ "source" ][ "generated" ]
146 body = copy.deepcopy(prior[ "source" ][ "desired" ])
147 body[ "@odata.etag" ] = receipt[ "creation" ][ "verification" ][ "readback" ][ "etag" ]
148 body[ "azureBlobParameters" ][ "createdResources" ] = {item[ "type" ]: item[ "name" ] for item in generated}
149 snapshots = [
150 { "type" : item[ "type" ], "name" : item[ "name" ], "etag" : receipt[ "configuration" ][item[ "type" ]][ "etag" ],
151 "definition_digest" : receipt[ "configuration" ][item[ "type" ]][ "digest" ]}
152 for item in generated
153 ]
154 return prior, result, { "body" : body, "_checkpoint_validated" : True , "_generated_snapshots" : snapshots}
155 paths = [_path(request, field, base_dir) for field in (
156 "creation_input_file" , "creation_result_file" , "creation_response_file" ,
157 )]
158 try :
159 _, prior, fingerprint = load_approved_input(paths[ 0 ])
160 except UnicodeError as exc:
161 raise _failure( "input-unreadable" , "Retained approval must be UTF-8 JSON." ) from exc
162 result, response = _read_json(paths[ 1 ]), _read_json(paths[ 2 ])
163 reject_secrets(result)
164 reject_secrets(response)
165 require_allowed_fields(response, RESPONSE_FIELDS , label = "Definitive create response record" )
166 if (
167 response.get( "schema_version" ) != "1.0"
168 or not _text(response.get( "request_id" ))
169 or _target(response.get( "target" )) != target
170 ):
171 raise _failure( "ownership-unproven" , "Original create response must bind the exact scope and native request ID." )
172 if (
173 result.get( "status" ) != "completed"
174 or result.get( "approved_plan" ) != { "confirmed" : True , "fingerprint" : fingerprint}
175 or _object(result.get( "ownership" ), "Creation ownership" ).get( "owner" ) != prior.get( "owner" )
176 or prior.get( "owner" ) != request[ "owner" ]
177 ):
178 raise _failure( "ownership-unproven" , "Retained completed result, approval and accountable owner must agree." )
179 _object(response.get( "body" ), "Create response body" )
180 return prior, result, response
181
182
183 def _owned_entry (result: dict[ str , Any], target: dict[ str , Any]) -> dict[ str , Any]:
184 identity = {key: target[key] for key in ( "type" , "name" , "version" ) if key in target}
185
186 def matches (item: Any) -> bool :
187 return isinstance (item, dict ) and all (item.get(key) == value for key, value in identity.items())
188
189 resources = _object(result.get( "resources" ), "Creation resources" )
190 ownership = _object(result.get( "ownership" ), "Creation ownership" )
191 for container, keys in (
192 (resources, ( "created" , "reused" , "updated" , "skipped" )),
193 (ownership, ( "run_owned" , "reused_not_owned" )),
194 ):
195 for key in keys:
196 if not isinstance (container.get(key, []), list ):
197 raise _failure( "ownership-unproven" , "Creation resource records must be complete lists." )
198 created = [item for item in resources.get( "created" , []) if matches(item)]
199 owned = [item for item in ownership.get( "run_owned" , []) if matches(item)]
200 if (
201 len (created) != 1 or len (owned) != 1 or created[ 0 ] != owned[ 0 ]
202 or any (matches(item) for key in ( "reused" , "updated" , "skipped" ) for item in resources.get(key, []))
203 or any (matches(item) for item in ownership.get( "reused_not_owned" , []))
204 ):
205 raise _failure( "ownership-unproven" , "Only one definitively created run-owned resource is eligible; never reused or updated." )
206 return created[ 0 ]
207
208
209 def _search_prior (prior: dict[ str , Any], target: dict[ str , Any]) -> dict[ str , Any]:
210 source = prior
211 if prior.get( "operation" ) in ( "reconcile-and-ingest" , "reconcile-and-monitor" ):
212 source = _object(prior.get( "source" ), "Prior source" )
213 if _object(source.get( "desired" ), "Prior source definition" ).get( "kind" ) == "file" :
214 try :
215 from . import file_source
216 except ImportError :
217 import file_source
218 file_source._validate_plan(prior)
219 else :
220 try :
221 from . import blob_source
222 except ImportError :
223 import blob_source
224 blob_source._validate_plan(prior)
225 search_reconcile._validate_plan(source)
226 if (
227 source.get( "operation" ) != "reconcile" or source.get( "action" ) != "create"
228 or search_reconcile.resource_url(source) != search_reconcile.resource_url(
229 { ** target, "resource_type" : target[ "type" ]}
230 )
231 ):
232 raise _failure( "ownership-unproven" , "Creation approval must select this exact source/base with action create." )
233 return source
234
235
236 def _summary (target: dict[ str , Any]) -> dict[ str , Any]:
237 is_agent = target[ "type" ] == "prompt-agent-version"
238 summary = {
239 "delete" : [copy.deepcopy(target)],
240 "retain" : RETAIN + (
241 [ "agent container" , "prior and other agent versions" , "project connection" , "KB" , "KS" ]
242 if is_agent else [ "all KS and generated indexes/pipelines" , "agent versions" , "project connections" ]
243 ),
244 "blocked" : [],
245 "order" : [ "Only this exact target; stop after failure. No dependent cleanup is chained." ],
246 "impact" : (
247 "The selected version and its tool binding become unavailable; the project connection remains."
248 if is_agent else "The selected KB retrieval endpoint becomes unavailable; retained consumers are not detached."
249 ),
250 "ownership" : "Original successful-create evidence required; owner metadata and equivalent GET are not proof or RBAC." ,
251 "approval" : "Separate cleanup consent: change only approval.confirmed after review." ,
252 "hosted_cleanup" : "unsupported" ,
253 }
254 if target[ "type" ] == "knowledge-source" :
255 summary.update(
256 retain = RETAIN + [ "all KBs" , "other KS/pipelines" , "agent versions" , "project connections" ],
257 impact = "The selected source, its uploaded File copies/indexed content and exact generated objects are deleted; originals remain." ,
258 order = [ "Referencing KBs must already be absent or no longer reference this source; no detach is performed." ,
259 "Delete only the source through Search; the service owns the exact approved cascade." ],
260 )
261 elif target[ "type" ] == "project-connection" :
262 summary.update(
263 retain = RETAIN + [ "agent containers" , "prior/other agent versions" , "KBs" , "KS/pipelines" ],
264 impact = "The selected project connection is removed; no retained agent version or tool is edited." ,
265 order = [ "Verify all consumers; delete and verify the explicitly selected owned version first, if any." ,
266 "Rescan protected consumers and conditionally delete only this connection; stop on failure." ],
267 )
268 return summary
269
270
271 def _original_generated (response):
272 names = dependencies.generated(response[ "body" ])
273 if response.get( "_checkpoint_validated" ) is True :
274 snapshots = response[ "_generated_snapshots" ]
275 else :
276 records = response.get( "generated_resources" )
277 if not isinstance (records, list ):
278 raise _failure(
279 "generated-creation-evidence-unavailable" ,
280 "Original generated-object GET snapshots with ETags are required: source ETag and createdResources names "
281 "cannot distinguish a replaced child. Use the retained Blob checkpoint or original generated-object audit records." ,
282 )
283 snapshots = []
284 for record in records:
285 if not isinstance (record, dict ) or set (record) != { "type" , "body" }:
286 raise _failure( "input-schema-invalid" , "Original generated records require type and complete native body." )
287 if not isinstance (record[ "type" ], str ) or record[ "type" ] not in names:
288 raise _failure( "generated-ownership-unproven" , "Unexpected original generated resource type." )
289 body = _object(record[ "body" ], "Original generated body" )
290 snapshots.append({
291 "type" : record[ "type" ], "name" : body.get( "name" ), "etag" : body.get( "@odata.etag" ),
292 "definition_digest" : digest(body),
293 })
294 if (
295 len (snapshots) != len (names)
296 or any (item.get( "type" ) not in names or item.get( "name" ) != names[item[ "type" ]]
297 or not _text(item.get( "etag" )) for item in snapshots)
298 or len ({item[ "type" ] for item in snapshots}) != len (names)
299 ):
300 raise _failure( "generated-ownership-unproven" , "Original child snapshots must match the exact complete acknowledged cascade." )
301 return sorted (snapshots, key =lambda item: item[ "type" ])
302
303
304 def _absent (target: dict[ str , Any], request_id: str | None = None ) -> dict[ str , Any]:
305 summary = _summary(target)
306 summary[ "delete" ] = []
307 if target[ "type" ] == "knowledge-source" :
308 summary[ "impact" ] = "Source already absent; generated-object absence is unverified. No independent child deletion is authorized."
309 summary[ "retain" ].append( "unverified generated objects" )
310 return {
311 "status" : "already-absent" , "outcome" : "plan-cleanup" ,
312 "execution_required" : False , "mutation_approval_required" : False ,
313 "approval_summary" : summary,
314 "verification" : { "absence" : True , "request_ids" : [request_id] if request_id else []},
315 }
316
317
318 def _planned (
319 target: dict[ str , Any], plan: dict[ str , Any], executor: str ,
320 retained_targets: list[dict[ str , Any]],
321 ) -> dict[ str , Any]:
322 fingerprint = digest(plan)
323 summary = _summary(target)
324 summary[ "retained_targets" ] = retained_targets
325 return {
326 "status" : "planned" , "outcome" : "plan-cleanup" , "executor" : executor,
327 "execution_required" : True , "mutation_approval_required" : True ,
328 "plan_fingerprint" : fingerprint,
329 "execution_input" : {
330 "schema_version" : "1.0" , "plan" : plan,
331 "approval" : { "confirmed" : False , "fingerprint" : fingerprint},
332 },
333 "approval_summary" : summary,
334 }
335
336
337 def _plan_search (
338 request: dict[ str , Any], target: dict[ str , Any], prior: dict[ str , Any],
339 result: dict[ str , Any], response: dict[ str , Any], * ,
340 token_provider: TokenProvider, transport: Transport,
341 ) -> dict[ str , Any]:
342 source = _search_prior(prior, target)
343 owned = _owned_entry(result, target)
344 body = response[ "body" ]
345 owned_digest = digest(search_reconcile._definition(body))
346 etag = body.get( "@odata.etag" )
347 if (
348 (response.get( "_checkpoint_validated" ) is not True and (
349 response.get( "operation" ) != "search-create" or type (response.get( "status" )) is not int or response[ "status" ] != 201
350 )) or not _text(etag)
351 or body.get( "name" ) != target[ "name" ]
352 or not search_reconcile.definitions_match(source[ "desired" ], body)
353 or owned.get( "definition_digest" ) != owned_digest or owned.get( "etag" ) != etag
354 ):
355 raise _failure( "ownership-unproven" , "Require the original HTTP 201 create body and matching retained owned readback, not GET recovery." )
356 plan = {
357 "operation" : "delete" , "outcome" : "cleanup-search-resource" , "plan_kind" : "cleanup" ,
358 "cleanup_approved" : True , "owner" : request[ "owner" ],
359 "resource_type" : target[ "type" ], "endpoint" : target[ "endpoint" ],
360 "name" : target[ "name" ], "api_version" : target[ "api_version" ],
361 "owned_definition_digest" : owned_digest, "expected_etag" : etag,
362 }
363 token = token_provider( SEARCH_AUDIENCE )
364 current, request_id = search_reconcile.read_resource(search_reconcile.resource_url(plan), token, transport = transport)
365 if current is None :
366 return _absent(target, request_id)
367 if (
368 digest(search_reconcile._definition(current)) != owned_digest
369 or current.get( "@odata.etag" ) != etag
370 ):
371 raise _failure( "definition-drift" , "Current definition or creation ETag changed; same-name replacement is not owned." )
372 if target[ "type" ] == "knowledge-source" :
373 original = _original_generated(response)
374 if dependencies.generated(current) != dependencies.generated(body):
375 raise _failure( "generated-ownership-unproven" , "Generated identities differ from the acknowledged source creation." )
376 snapshot = dependencies.search_snapshot(
377 plan, current, token, transport = transport, bounds = dependencies.limits(request.get( "inventory_limits" )),
378 )
379 if snapshot[ "generated" ] != original:
380 raise _failure( "generated-incarnation-drift" , "Generated definitions/ETags differ from original retained ownership evidence." )
381 plan[ "dependency_guard" ] = snapshot
382 search_reconcile._validate_plan(plan)
383 planned = _planned(target, plan, "helpers/search_reconcile.py" , [])
384 planned[ "approval_summary" ][ "delete" ].extend([
385 { "type" : item[ "type" ], "name" : item[ "name" ], "endpoint" : target[ "endpoint" ], "service_managed" : True }
386 for item in original
387 ])
388 return planned
389 reject_secrets(current)
390 plan[ "desired" ] = copy.deepcopy(current)
391 search_reconcile._validate_plan(plan)
392 retained = [
393 { ** target, "type" : "knowledge-source" , "name" : source[ "name" ]}
394 for source in current[ "knowledgeSources" ]
395 ]
396 return _planned(target, plan, "helpers/search_reconcile.py" , retained)
397
398
399 def _plan_prompt (
400 request: dict[ str , Any], target: dict[ str , Any], prior: dict[ str , Any],
401 result: dict[ str , Any], response: dict[ str , Any], * ,
402 sdk_loader: Callable[[], tuple[Any, Any, Any, Any, Any]],
403 ) -> dict[ str , Any]:
404 prompt_connect._validate_plan(prior)
405 owned = _owned_entry(result, target)
406 prior_project, prior_endpoint = prompt_connect._project_identity(prior)
407 body = response[ "body" ]
408 require_allowed_fields(body, { "name" , "version" , "definition" , "id" , "created_at" }, label = "SDK create-version response" )
409 definition = _object(body.get( "definition" ), "Created Prompt definition" )
410 if (
411 response.get( "operation" ) != "agents.create_version" or response.get( "status" ) != "succeeded"
412 or prior_project.casefold() != target[ "project_resource_id" ]
413 or prior_endpoint != target[ "project_endpoint" ]
414 or prior[ "agent" ][ "name" ] != target[ "name" ] or prior[ "agent" ][ "version" ] == target[ "version" ]
415 or body.get( "name" ) != target[ "name" ] or body.get( "version" ) != target[ "version" ]
416 or definition.get( "kind" ) != "prompt"
417 or owned.get( "definition_digest" ) != digest(definition)
418 or receipts.version_identity(body) is None
419 ):
420 raise _failure( "ownership-unproven" , "Require the original successful SDK new-version return, never the baseline or recovered equivalent." )
421 AIProjectClient, _, _, _, extras = sdk_loader()
422 AzureCliCredential, AzureError = extras
423 client = AIProjectClient( endpoint = target[ "project_endpoint" ], credential = AzureCliCredential())
424 try :
425 try :
426 current = client.agents.get_version( agent_name = target[ "name" ], agent_version = target[ "version" ])
427 except AzureError as exc:
428 if sdk_error_status(exc) == 404 :
429 return _absent(target, sdk_error_metadata(exc).get( "request_id" ))
430 raise HelperFailure(
431 message = "Exact Prompt version readback failed; no inventory or deletion attempted." ,
432 blocked_at = "cleanup-planning" , ** sdk_error_metadata(exc, "agent-readback-failed" ),
433 ) from exc
434 try :
435 actual = current.definition.as_dict()
436 identity_matches = current.name == target[ "name" ] and str (current.version) == target[ "version" ]
437 except ( AttributeError , TypeError , ValueError ) as exc:
438 raise _failure( "agent-readback-invalid" , "Exact version readback is incomplete." ) from exc
439 if ( not identity_matches or not isinstance (actual, dict ) or digest(actual) != owned[ "definition_digest" ]
440 or receipts.version_identity(current) != receipts.version_identity(body)):
441 raise _failure( "definition-drift" , "Exact Prompt version identity or definition changed." )
442 finally :
443 client.close()
444 plan = {
445 "operation" : "delete" , "outcome" : "cleanup-prompt-version" , "plan_kind" : "cleanup" ,
446 "cleanup_approved" : True , "sdk_major" : 2 , "owner" : request[ "owner" ],
447 "project_resource_id" : target[ "project_resource_id" ], "project_endpoint" : target[ "project_endpoint" ],
448 "agent" : {
449 "name" : target[ "name" ], "version" : target[ "version" ], "run_owned" : True ,
450 "owned_definition_digest" : owned[ "definition_digest" ],
451 "owned_version_identity" : receipts.version_identity(body),
452 },
453 }
454 prompt_cleanup._validate_plan(plan)
455 return _planned(target, plan, "helpers/prompt_cleanup.py" , [
456 { ** target, "version" : prior[ "agent" ][ "version" ]},
457 {
458 "type" : "project-connection" , "project_resource_id" : target[ "project_resource_id" ],
459 "project_endpoint" : target[ "project_endpoint" ], "name" : prior[ "connection" ][ "name" ],
460 },
461 ])
462
463
464 def _plan_connection (request, target, prior, result, response, * , base_dir, token_provider, transport, sdk_loader):
465 prompt_connect._validate_plan(prior)
466 project, endpoint = prompt_connect._project_identity(prior)
467 body = response[ "body" ]
468 created = { "type" : "project-connection" , "name" : target[ "name" ]}
469 owned_write = { "action" : "created" , "connection" : target[ "name" ]}
470 resources = _object(result.get( "resources" ), "Creation resources" )
471 ownership = _object(result.get( "ownership" ), "Creation ownership" )
472 verification = _object(result.get( "verification" ), "Creation verification" )
473 readback = _object(verification.get( "connection_readback" ), "Creation connection readback" )
474 def matches (item):
475 return isinstance (item, dict ) and (
476 (item.get( "type" ) == "project-connection" and item.get( "name" ) == target[ "name" ])
477 or item.get( "connection" ) == target[ "name" ]
478 )
479
480 if any ( not isinstance (container.get(key, []), list ) for container, keys in (
481 (resources, ( "created" , "reused" , "updated" , "skipped" )), (ownership, ( "run_owned" , "reused_not_owned" )),
482 ) for key in keys):
483 raise _failure( "ownership-unproven" , "Creation resource and ownership lists must be complete." )
484 etag = body.get( "etag" ) or body.get( "@odata.etag" )
485 if (
486 project.casefold() != target[ "project_resource_id" ] or endpoint != target[ "project_endpoint" ]
487 or prior[ "connection" ][ "name" ] != target[ "name" ] or prior[ "connection" ][ "action" ] != "create"
488 or response.get( "operation" ) != "project-connection-create" or type (response.get( "status" )) is not int
489 or response[ "status" ] != 201 or not _text(etag)
490 or [item for item in resources.get( "created" , []) if matches(item)] != [created]
491 or any (matches(item) for key in ( "reused" , "updated" , "skipped" ) for item in resources.get(key, []))
492 or [item for item in ownership.get( "run_owned" , []) if matches(item)] != [owned_write]
493 or any (matches(item) for item in ownership.get( "reused_not_owned" , []))
494 or readback.get( "name" ) != target[ "name" ] or readback.get( "definition_digest" ) != digest(body)
495 ):
496 raise _failure( "ownership-unproven" , "Retain the exact acknowledged created project connection; updated/reused/recovered GETs are not create evidence." )
497 if not prompt_connect._connection_readback(prior, body)[ 0 ]:
498 raise _failure( "ownership-unproven" , "The original connection body must match its approved project and configuration." )
499 plan = {
500 "operation" : "delete" , "outcome" : "cleanup-prompt-connection" , "plan_kind" : "cleanup" ,
501 "cleanup_approved" : True , "sdk_major" : 2 , "owner" : request[ "owner" ],
502 "project_resource_id" : target[ "project_resource_id" ], "project_endpoint" : target[ "project_endpoint" ],
503 "connection" : { "name" : target[ "name" ], "run_owned" : True , "expected_etag" : etag, "owned_definition_digest" : digest(body)},
504 }
505 url = prompt_connect._connection_url(plan)
506 token = token_provider( MANAGEMENT_AUDIENCE )
507 current, request_id = prompt_cleanup._get_connection(url, token, transport = transport)
508 if current is not None and current != body:
509 raise _failure( "definition-drift" , "The project connection changed since acknowledged creation." )
510 selected = None
511 if "agent_version" in request:
512 selected_target = {
513 ** target, "type" : "prompt-agent-version" , "name" : prior[ "agent" ][ "name" ], "version" : request[ "agent_version" ],
514 }
515 selected_request = {
516 "schema_version" : "1.0" , "owner" : request[ "owner" ], "target" : selected_target,
517 "creation_input_file" : request[ "creation_input_file" ], "creation_result_file" : request[ "creation_result_file" ],
518 "creation_response_file" : request[ "agent_creation_response_file" ],
519 }
520 selected = plan_cleanup(selected_request, base_dir = base_dir, token_provider = token_provider, transport = transport, sdk_loader = sdk_loader)
521 if selected[ "status" ] == "planned" :
522 plan[ "agent" ] = selected[ "execution_input" ][ "plan" ][ "agent" ]
523 if current is None :
524 if selected is not None and selected[ "status" ] == "planned" :
525 selected[ "approval_summary" ][ "already_absent" ] = [target]
526 return selected
527 return _absent(target, request_id)
528 plan[ "dependency_guard" ] = dependencies.prompt_snapshot(
529 plan, current, sdk_loader = sdk_loader, bounds = dependencies.limits(request.get( "inventory_limits" )),
530 )
531 refreshed, _ = prompt_cleanup._get_connection(url, token, transport = transport)
532 if refreshed != current:
533 raise _failure( "definition-drift" , "Connection changed during project consumer discovery." )
534 prompt_cleanup._validate_plan(plan)
535 retained = [
536 { ** target, "type" : "prompt-agent-version" , "name" : item[ "name" ], "version" : item[ "version" ]}
537 for item in plan[ "dependency_guard" ][ "versions" ]
538 if not plan.get( "agent" ) or (item[ "name" ], item[ "version" ]) != (plan[ "agent" ][ "name" ], plan[ "agent" ][ "version" ])
539 ]
540 planned = _planned(target, plan, "helpers/prompt_cleanup.py" , retained)
541 if selected is not None and selected[ "status" ] == "planned" :
542 planned[ "approval_summary" ][ "delete" ].insert( 0 , selected[ "approval_summary" ][ "delete" ][ 0 ])
543 return planned
544
545
546 def _plan_protected (request, target, prior, record, * , base_dir, token_provider, transport, sdk_loader):
547 snapshot = record[ "snapshot" ]
548 outcome = ( "cleanup-search-resource" if target[ "type" ] in ( "knowledge-base" , "knowledge-source" )
549 else "cleanup-prompt-version" if target[ "type" ] == "prompt-agent-version" else "cleanup-prompt-connection" )
550 plan = { "operation" : "delete" , "outcome" : outcome, "plan_kind" : "cleanup" ,
551 "cleanup_approved" : True , "owner" : request[ "owner" ]}
552 if record[ "owner" ] != request[ "owner" ]:
553 raise _failure( "ownership-unproven" , "Original producer owner and cleanup accountable owner differ." )
554 if target[ "type" ] in ( "knowledge-base" , "knowledge-source" ):
555 source = _search_prior(prior, target)
556 if snapshot[ "definition_digest" ] != digest(search_reconcile._definition(source[ "desired" ])):
557 raise _failure( "ownership-unproven" , "Receipt does not bind the original approved Search definition." )
558 plan.update( resource_type = target[ "type" ], ** {k: target[k] for k in ( "endpoint" , "name" , "api_version" )},
559 owned_definition_digest = snapshot[ "definition_digest" ], expected_etag = snapshot[ "etag" ])
560 token = token_provider( SEARCH_AUDIENCE )
561 current, request_id = search_reconcile.read_resource(search_reconcile.resource_url(plan), token, transport = transport)
562 if current is None :
563 return _absent(target, request_id)
564 if digest(search_reconcile._definition(current)) != snapshot[ "definition_digest" ] or current.get( "@odata.etag" ) != snapshot[ "etag" ]:
565 raise _failure( "definition-drift" , "Current Search state differs from the original producer snapshot." )
566 if target[ "type" ] == "knowledge-source" :
567 if dependencies.generated(current) != record[ "acknowledgement" ][ "generated" ]:
568 raise _failure( "generated-ownership-unproven" , "Current generated identities differ from the native create acknowledgement." )
569 plan[ "dependency_guard" ] = dependencies.search_snapshot(
570 plan, current, token, transport = transport, bounds = dependencies.limits(request.get( "inventory_limits" )))
571 if plan[ "dependency_guard" ][ "generated" ] != snapshot[ "generated" ]:
572 raise _failure( "generated-incarnation-drift" , "Generated resources differ from original producer snapshots." )
573 else :
574 reject_secrets(current)
575 plan[ "desired" ] = current
576 search_reconcile._validate_plan(plan)
577 planned = _planned(target, plan, "helpers/search_reconcile.py" , [])
578 if target[ "type" ] == "knowledge-source" :
579 planned[ "approval_summary" ][ "delete" ].extend(
580 { "type" : child[ "type" ], "name" : child[ "name" ], "endpoint" : target[ "endpoint" ], "service_managed" : True }
581 for child in snapshot[ "generated" ])
582 else :
583 planned[ "approval_summary" ][ "retained_targets" ] = [
584 { ** target, "type" : "knowledge-source" , "name" : item[ "name" ]} for item in current[ "knowledgeSources" ]]
585 return planned
586 if prior.get( "operation" ) == "create-initial-prompt-agent" :
587 try :
588 from . import _initial_prompt
589 except ImportError :
590 import _initial_prompt
591 _initial_prompt.validate(prior)
592 if (target[ "type" ] != "prompt-agent-version" or record[ "acknowledgement" ][ "operation" ] != "agents.create"
593 or snapshot[ "definition_digest" ] != digest(prior[ "agent" ][ "definition" ])):
594 raise _failure( "ownership-unproven" , "Initial receipts authorize only the acknowledged version, never a connection/container." )
595 else :
596 prompt_connect._validate_plan(prior)
597 if target[ "type" ] == "prompt-agent-version" and record[ "acknowledgement" ][ "operation" ] != "agents.create_version" :
598 raise _failure( "ownership-unproven" , "Connect receipts must come from the native new-version operation." )
599 project, endpoint = prompt_connect._project_identity(prior)
600 if project.casefold() != target[ "project_resource_id" ] or endpoint != target[ "project_endpoint" ]:
601 raise _failure( "ownership-unproven" , "Receipt project differs from original creation approval." )
602 plan.update( sdk_major = 2 , project_resource_id = target[ "project_resource_id" ], project_endpoint = target[ "project_endpoint" ])
603 if target[ "type" ] == "prompt-agent-version" :
604 if prior[ "agent" ][ "name" ] != target[ "name" ] or prior[ "agent" ].get( "version" ) == target[ "version" ]:
605 raise _failure( "ownership-unproven" , "Only the acknowledged new version is eligible, never the baseline." )
606 Client, _, _, _, extras = sdk_loader()
607 Credential, AzureError = extras
608 client = Client( endpoint = endpoint, credential = Credential())
609 try :
610 current = client.agents.get_version( agent_name = target[ "name" ], agent_version = target[ "version" ])
611 if (current.name != target[ "name" ] or str (current.version) != target[ "version" ]
612 or digest(current.definition.as_dict()) != snapshot[ "definition_digest" ]
613 or receipts.version_identity(current) != snapshot[ "version_identity" ]):
614 raise _failure( "definition-drift" , "Selected version differs from the original SDK return." )
615 except AzureError as exc:
616 if sdk_error_status(exc) == 404 :
617 return _absent(target)
618 raise HelperFailure( message = "Exact version readback failed." , blocked_at = "cleanup-planning" ,
619 ** sdk_error_metadata(exc, "agent-readback-failed" )) from exc
620 finally :
621 client.close()
622 plan[ "agent" ] = { "name" : target[ "name" ], "version" : target[ "version" ], "run_owned" : True ,
623 "owned_definition_digest" : snapshot[ "definition_digest" ], "owned_version_identity" : snapshot[ "version_identity" ]}
624 prompt_cleanup._validate_plan(plan)
625 return _planned(target, plan, "helpers/prompt_cleanup.py" , [])
626 if prior[ "connection" ][ "name" ] != target[ "name" ] or prior[ "connection" ][ "action" ] != "create" :
627 raise _failure( "ownership-unproven" , "Connection receipt must bind original create intent." )
628 plan[ "connection" ] = { "name" : target[ "name" ], "run_owned" : True , "expected_etag" : snapshot[ "etag" ],
629 "owned_definition_digest" : snapshot[ "definition_digest" ]}
630 token = token_provider( MANAGEMENT_AUDIENCE )
631 url = prompt_connect._connection_url(plan)
632 current, request_id = prompt_cleanup._get_connection(url, token, transport = transport)
633 if current is not None and (digest(current) != snapshot[ "definition_digest" ] or not prompt_connect._connection_readback(prior, current)[ 0 ]):
634 raise _failure( "definition-drift" , "Connection differs from its original acknowledged producer state." )
635 selected = None
636 if "agent_version" in request:
637 selected = plan_cleanup({
638 "schema_version" : "1.0" , "owner" : request[ "owner" ],
639 "target" : { ** target, "type" : "prompt-agent-version" , "name" : prior[ "agent" ][ "name" ], "version" : request[ "agent_version" ]},
640 "creation_input_file" : request[ "creation_input_file" ],
641 "creation_receipt_file" : request[ "agent_creation_receipt_file" ],
642 }, base_dir = base_dir, token_provider = token_provider, transport = transport, sdk_loader = sdk_loader)
643 if selected[ "status" ] == "planned" :
644 plan[ "agent" ] = selected[ "execution_input" ][ "plan" ][ "agent" ]
645 if current is None :
646 if selected and selected[ "status" ] == "planned" :
647 selected[ "approval_summary" ][ "already_absent" ] = [target]
648 return selected
649 return _absent(target, request_id)
650 plan[ "dependency_guard" ] = dependencies.prompt_snapshot(
651 plan, current, sdk_loader = sdk_loader, bounds = dependencies.limits(request.get( "inventory_limits" )))
652 if prompt_cleanup._get_connection(url, token, transport = transport)[ 0 ] != current:
653 raise _failure( "definition-drift" , "Connection changed during consumer discovery." )
654 prompt_cleanup._validate_plan(plan)
655 planned = _planned(target, plan, "helpers/prompt_cleanup.py" , [])
656 if selected and selected[ "status" ] == "planned" :
657 planned[ "approval_summary" ][ "delete" ].insert( 0 , selected[ "approval_summary" ][ "delete" ][ 0 ])
658 return planned
659
660
661 def plan_cleanup (
662 request: dict[ str , Any], * , base_dir: Path = Path( "." ),
663 token_provider: TokenProvider = azure_cli_token, transport: Transport = http_request,
664 sdk_loader: Callable[[], tuple[Any, Any, Any, Any, Any]] = prompt_cleanup.load_cleanup_sdk,
665 ) -> dict[ str , Any]:
666 _object(request, "Cleanup planning request" )
667 reject_secrets(request)
668 require_allowed_fields(request, REQUEST_FIELDS , label = "Cleanup planning request" )
669 if request.get( "schema_version" ) != "1.0" or not _text(request.get( "owner" )):
670 raise _failure( "input-schema-invalid" , "schema_version 1.0 and accountable owner are required." )
671 target = _target(request.get( "target" ))
672 dependencies.limits(request.get( "inventory_limits" ))
673 selection_fields = { "agent_version" , "agent_creation_response_file" , "agent_creation_receipt_file" } & set (request)
674 evidence_field = "agent_creation_receipt_file" if "agent_creation_receipt_file" in request else "agent_creation_response_file"
675 if selection_fields and (selection_fields != { "agent_version" , evidence_field} or target[ "type" ] != "project-connection" ):
676 raise _failure( "input-schema-invalid" , "Only connection cleanup can explicitly select both agent_version and its original create response." )
677 if selection_fields and (
678 not isinstance (request[ "agent_version" ], str ) or re.fullmatch( r " [ 1-9 ][ 0-9 ] * " , request[ "agent_version" ]) is None
679 or not _text(request[evidence_field])
680 ):
681 raise _failure( "target-selection-required" , "Select an exact numeric created agent version and original SDK response file." )
682 if "creation_receipt_file" in request:
683 try :
684 from .blob_recheck import read_private
685 except ImportError :
686 from blob_recheck import read_private
687 path = _path(request, "creation_receipt_file" , base_dir)
688 if read_private(path).get( "kind" ) == "cleanup-creation-receipt" :
689 if ( "creation_response_file" in request or "creation_result_file" in request
690 or (selection_fields and evidence_field != "agent_creation_receipt_file" )):
691 raise _failure( "input-schema-invalid" , "Protected receipts cannot be mixed with manual response adapters." )
692 prior, record = receipts.load(_path(request, "creation_input_file" , base_dir), path, target)
693 return _plan_protected(request, target, prior, record, base_dir = base_dir,
694 token_provider = token_provider, transport = transport, sdk_loader = sdk_loader)
695 if evidence_field == "agent_creation_receipt_file" :
696 raise _failure( "input-schema-invalid" , "Selected producer receipts require a protected connection receipt." )
697 prior, result, response = _records(request, target, base_dir)
698 if target[ "type" ] in { "knowledge-base" , "knowledge-source" }:
699 return _plan_search(
700 request, target, prior, result, response, token_provider = token_provider, transport = transport,
701 )
702 if target[ "type" ] == "project-connection" :
703 return _plan_connection(
704 request, target, prior, result, response, base_dir = base_dir,
705 token_provider = token_provider, transport = transport, sdk_loader = sdk_loader,
706 )
707 return _plan_prompt(request, target, prior, result, response, sdk_loader = sdk_loader)
708
709
710 def main (argv: list[ str ] | None = None ) -> int :
711 parser = argparse.ArgumentParser()
712 parser.add_argument( "--plan" , type = Path, required = True )
713 args = parser.parse_args(argv)
714 request: dict[ str , Any] = {}
715 try :
716 request = _read_json(args.plan)
717 result = plan_cleanup(request, base_dir = args.plan.resolve().parent)
718 emit_result(result, preserve_unapproved_input = result[ "status" ] == "planned" )
719 except HelperFailure as failure:
720 result = blocked_result(failure, outcome = "plan-cleanup" , fingerprint = None , owner = request.get( "owner" ))
721 result[ "approval_summary" ] = {
722 "delete" : [], "retain" : RETAIN + [ "selected target and all its dependencies" ],
723 "blocked" : [failure.code], "order" : [], "hosted_cleanup" : "unsupported" ,
724 }
725 try :
726 result[ "approval_summary" ][ "retained_targets" ] = [_target(request.get( "target" ))]
727 except HelperFailure:
728 pass
729 emit_result(result)
730 return 2
731 return 0
732
733
734 if __name__ == "__main__" :
735 sys.exit(main())