Setting the file. One moment. Prompt Cleanup · Foundry Iq · microsoft/azure-skills · Skills DocsFile Cu Canary
helpers/prompt_cleanup.py
Python·621 lines·23 KB
import
_cleanup_dependencies
as
cleanup_dependencies
12 from ._common import (
13 MANAGEMENT_AUDIENCE,
14 HelperFailure,
15 TokenProvider,
16 Transport,
17 azure_cli_token,
18 blocked_result,
19 digest,
20 emit_result,
21 http_request,
22 is_ambiguous_mutation_failure,
23 is_ambiguous_sdk_error,
24 load_approved_input,
25 reject_secrets,
26 require_allowed_fields,
27 sdk_error_status,
28 sdk_error_metadata,
29 )
30 from .prompt_connect import _connection_url, _load_connection_sdk as _load_sdk, _project_identity
31except ImportError:
32 import _cleanup_dependencies as cleanup_dependencies
33 from _common import ( # type: ignore[no-redef]
34 MANAGEMENT_AUDIENCE,
35 HelperFailure,
36 TokenProvider,
37 Transport,
38 azure_cli_token,
39 blocked_result,
40 digest,
41 emit_result,
42 http_request,
43 is_ambiguous_mutation_failure,
44 is_ambiguous_sdk_error,
45 load_approved_input,
46 reject_secrets,
47 require_allowed_fields,
48 sdk_error_status,
49 sdk_error_metadata,
50 )
51 from prompt_connect import ( # type: ignore[no-redef]
52 _connection_url,
53 _load_connection_sdk as _load_sdk,
54 _project_identity,
55 )
56
57
58def load_cleanup_sdk():
59 try:
60 installed = version("azure-ai-projects").split(".")
61 if int(installed[0]) != 2 or int(installed[1]) < 4:
62 raise ValueError
63 except (PackageNotFoundError, ValueError, IndexError) as exc:
64 raise HelperFailure("sdk-version-invalid", "Cleanup requires azure-ai-projects >=2.4,<3 for complete draft inventories.",
65 blocked_at="execution") from exc
66 return _load_sdk()
67
68
69def _validate_owned_resource(
70 resource: Any,
71 *,
72 label: str,
73 identity_fields: tuple[str, ...],
74) -> dict[str, Any] | None:
75 if resource is None:
76 return None
77 if (
78 not isinstance(resource, dict)
79 or resource.get("run_owned") is not True
80 or not isinstance(resource.get("owned_definition_digest"), str)
81 or not all(isinstance(resource.get(field), str) and resource[field] for field in identity_fields)
82 ):
83 raise HelperFailure(
84 "ownership-unproven",
85 f"{label} cleanup requires exact run-owned identity and definition digest.",
86 blocked_at="reconciliation",
87 )
88 return resource
89
90
91def _validate_plan(
92 plan: dict[str, Any],
93) -> tuple[dict[str, Any] | None, dict[str, Any] | None]:
94 reject_secrets(plan)
95 require_allowed_fields(
96 plan,
97 {
98 "operation",
99 "outcome",
100 "plan_kind",
101 "cleanup_approved",
102 "sdk_major",
103 "project_resource_id",
104 "project_endpoint",
105 "agent",
106 "connection",
107 "owner",
108 "dependency_guard",
109 },
110 label="Prompt cleanup plan",
111 )
112 cleanup_dependencies.validate_guard(plan)
113 if (
114 plan.get("operation") != "delete"
115 or plan.get("plan_kind") != "cleanup"
116 or plan.get("cleanup_approved") is not True
117 or plan.get("sdk_major") != 2
118 ):
119 raise HelperFailure(
120 "cleanup-approval-mismatch",
121 "Prompt cleanup requires a separate approved cleanup plan and SDK major 2.",
122 blocked_at="confirmation",
123 )
124 _project_identity(plan)
125 agent = _validate_owned_resource(
126 plan.get("agent"),
127 label="Agent version",
128 identity_fields=("name", "version"),
129 )
130 connection = _validate_owned_resource(
131 plan.get("connection"),
132 label="Project connection",
133 identity_fields=("name",),
134 )
135 if agent is not None:
136 require_allowed_fields(
137 agent,
138 {"name", "version", "run_owned", "owned_definition_digest", "owned_version_identity"},
139 label="Prompt cleanup agent",
140 )
141 if "owned_version_identity" in agent:
142 try:
143 from . import _cleanup_receipts
144 except ImportError:
145 import _cleanup_receipts
146 identity = agent["owned_version_identity"]
147 if (not isinstance(identity, dict) or set(identity) != {"id", "created_at"}
148 or _cleanup_receipts.version_identity(identity) != identity):
149 raise HelperFailure("ownership-unproven", "Native version birth identity must be exact.", blocked_at="reconciliation")
150 if connection is not None:
151 require_allowed_fields(
152 connection,
153 {
154 "name",
155 "run_owned",
156 "owned_definition_digest",
157 "expected_etag",
158 },
159 label="Prompt cleanup connection",
160 )
161 if connection is not None and not isinstance(connection.get("expected_etag"), str):
162 raise HelperFailure(
163 "ownership-unproven",
164 "Project connection cleanup requires the approved current ETag.",
165 blocked_at="reconciliation",
166 )
167 if agent is None and connection is None:
168 raise HelperFailure(
169 "ownership-unproven",
170 "Prompt cleanup must identify at least one run-owned resource.",
171 blocked_at="reconciliation",
172 )
173 if connection is not None:
174 _connection_url({**plan, "connection": connection})
175 return agent, connection
176
177
178def _delete_agent(
179 plan: dict[str, Any],
180 agent: dict[str, Any],
181 *,
182 sdk_loader: Callable[[], tuple[Any, Any, Any, Any, Any]],
183) -> tuple[str, dict[str, Any]]:
184 AIProjectClient, _, _, _, extras = sdk_loader()
185 AzureCliCredential, AzureError = extras
186 client = AIProjectClient(
187 endpoint=_project_identity(plan)[1],
188 credential=AzureCliCredential(),
189 )
190 identity = {"type": "prompt-agent-version", "name": agent["name"], "version": agent["version"]}
191 version_options = {"include_drafts": True}
192 try:
193 signature = inspect.signature(client.agents.list_versions)
194 if ("include_drafts" not in signature.parameters
195 and not any(p.kind == inspect.Parameter.VAR_KEYWORD for p in signature.parameters.values())):
196 raise HelperFailure("sdk-version-invalid", "Cleanup requires draft-inclusive version listing.",
197 blocked_at="execution")
198 versions = list(client.agents.list_versions(agent_name=agent["name"], **version_options))
199 if agent["version"] not in {str(item.version) for item in versions}:
200 return "skipped", identity
201 current = client.agents.get_version(
202 agent_name=agent["name"],
203 agent_version=agent["version"],
204 )
205 if "owned_version_identity" in agent:
206 try:
207 from . import _cleanup_receipts
208 except ImportError:
209 import _cleanup_receipts
210 if _cleanup_receipts.version_identity(current) != agent["owned_version_identity"]:
211 raise HelperFailure("definition-drift", "Native version identity changed since creation.", blocked_at="reconciliation")
212 if digest(current.definition.as_dict()) != agent["owned_definition_digest"]:
213 raise HelperFailure(
214 "definition-drift",
215 "The agent version definition changed after cleanup approval.",
216 blocked_at="reconciliation",
217 )
218 try:
219 client.agents.delete_version(
220 agent_name=agent["name"],
221 agent_version=agent["version"],
222 )
223 except AzureError as exc:
224 status = sdk_error_status(exc)
225 if status != 404 and not is_ambiguous_sdk_error(exc):
226 raise HelperFailure(
227 message="The Prompt Agent SDK delete operation failed.",
228 blocked_at="execution",
229 **sdk_error_metadata(exc, "agent-cleanup-failed"),
230 ) from exc
231 try:
232 remaining = list(
233 client.agents.list_versions(agent_name=agent["name"], **version_options)
234 )
235 except AzureError as readback_exc:
236 raise HelperFailure(
237 "agent-delete-outcome-ambiguous",
238 "Agent deletion and same-identity readback are ambiguous.",
239 blocked_at="verification",
240 resources_remaining=[identity],
241 partial=True,
242 **sdk_error_metadata(exc),
243 ) from readback_exc
244 if agent["version"] not in {
245 str(item.version) for item in remaining
246 }:
247 return ("skipped" if status == 404 else "deleted"), identity
248 if status == 404:
249 raise HelperFailure(
250 "agent-cleanup-failed",
251 "Agent delete returned not found but same-identity readback still found the version.",
252 blocked_at="verification",
253 resources_remaining=[identity],
254 **sdk_error_metadata(exc),
255 ) from exc
256 raise HelperFailure(
257 "agent-delete-outcome-ambiguous",
258 "Same-identity readback still found the agent version after an ambiguous delete.",
259 blocked_at="verification",
260 resources_remaining=[identity],
261 partial=True,
262 **sdk_error_metadata(exc),
263 ) from exc
264 try:
265 remaining = list(
266 client.agents.list_versions(agent_name=agent["name"], **version_options)
267 )
268 except AzureError as exc:
269 raise HelperFailure(
270 "agent-delete-outcome-ambiguous",
271 "Agent delete completed but absence readback is ambiguous.",
272 blocked_at="verification",
273 writes=[{"action": "deleted", **identity}],
274 resources_remaining=[identity],
275 partial=True,
276 **sdk_error_metadata(exc),
277 ) from exc
278 if agent["version"] in {str(item.version) for item in remaining}:
279 raise HelperFailure(
280 "absence-unverified",
281 "The exact agent version still exists after delete.",
282 blocked_at="verification",
283 writes=[{"action": "deleted", **identity}],
284 resources_remaining=[identity],
285 partial=True,
286 )
287 return "deleted", identity
288 except AzureError as exc:
289 raise HelperFailure(
290 message="The Prompt Agent SDK cleanup operation failed.",
291 blocked_at="execution",
292 partial=False,
293 **sdk_error_metadata(exc, "agent-cleanup-failed"),
294 ) from exc
295 finally:
296 client.close()
297
298
299def _get_connection(
300 url: str,
301 token: str,
302 *,
303 transport: Transport,
304) -> tuple[dict[str, Any] | None, str | None]:
305 try:
306 result = transport("GET", url, token)
307 except HelperFailure as failure:
308 if failure.http_status == 404:
309 return None, failure.request_id
310 raise
311 if result.status != 200 or not isinstance(result.body, dict):
312 raise HelperFailure(
313 "connection-readback-invalid",
314 "Project connection readback was not one JSON object.",
315 blocked_at="reconciliation",
316 request_id=result.request_id,
317 status=result.status,
318 )
319 return result.body, result.request_id
320
321
322def _delete_connection(
323 plan: dict[str, Any],
324 connection: dict[str, Any],
325 token: str,
326 *,
327 transport: Transport,
328 sdk_loader: Callable[[], tuple[Any, Any, Any, Any, Any]] = load_cleanup_sdk,
329) -> tuple[str, dict[str, Any], list[str]]:
330 url = _connection_url({**plan, "connection": connection})
331 current, initial_request_id = _get_connection(url, token, transport=transport)
332 identity = {"type": "project-connection", "name": connection["name"]}
333 request_ids = [initial_request_id] if initial_request_id else []
334 if current is None:
335 return "skipped", identity, request_ids
336 if digest(current) != connection["owned_definition_digest"]:
337 raise HelperFailure(
338 "definition-drift",
339 "The project connection changed after cleanup approval.",
340 blocked_at="reconciliation",
341 )
342 current_etag = current.get("etag") or current.get("@odata.etag")
343 if current_etag != connection["expected_etag"]:
344 raise HelperFailure(
345 "definition-drift",
346 "The project connection ETag changed after cleanup approval.",
347 blocked_at="reconciliation",
348 )
349 if plan.get("dependency_guard") is not None:
350 cleanup_dependencies.verify_prompt(plan, current, sdk_loader=sdk_loader, allow_selected=False)
351 refreshed, _ = _get_connection(url, token, transport=transport)
352 if refreshed != current:
353 raise HelperFailure("definition-drift", "Connection changed during final consumer discovery.", blocked_at="reconciliation")
354 try:
355 result = transport(
356 "DELETE",
357 url,
358 token,
359 headers={"If-Match": connection["expected_etag"]},
360 )
361 except HelperFailure as failure:
362 if failure.http_status == 404:
363 return "skipped", identity, request_ids
364 if not is_ambiguous_mutation_failure(failure):
365 raise
366 return _recover_ambiguous_connection_delete(
367 url,
368 token,
369 identity,
370 request_ids,
371 failure,
372 transport=transport,
373 )
374 if result.status not in {200, 202, 204}:
375 if result.status == 404:
376 return "skipped", identity, request_ids
377 if result.status in {408, 429} or result.status >= 500:
378 return _recover_ambiguous_connection_delete(
379 url,
380 token,
381 identity,
382 request_ids,
383 HelperFailure(
384 "connection-delete-outcome-ambiguous",
385 f"Project connection delete returned ambiguous HTTP {result.status}.",
386 blocked_at="execution",
387 request_id=result.request_id,
388 status=result.status,
389 partial=True,
390 ),
391 transport=transport,
392 )
393 raise HelperFailure(
394 "connection-cleanup-failed",
395 f"Project connection delete returned HTTP {result.status}.",
396 blocked_at="execution",
397 request_id=result.request_id,
398 status=result.status,
399 )
400 if result.request_id:
401 request_ids.append(result.request_id)
402 write = {"action": "deleted", **identity}
403 try:
404 after, verify_request_id = _get_connection(
405 url, token, transport=transport
406 )
407 except HelperFailure as failure:
408 raise HelperFailure(
409 failure.code,
410 failure.message,
411 blocked_at=failure.blocked_at,
412 writes=[write, *failure.writes],
413 resources_remaining=[identity],
414 request_id=failure.request_id,
415 status=failure.http_status,
416 partial=True,
417 ) from failure
418 if verify_request_id:
419 request_ids.append(verify_request_id)
420 if after is not None:
421 raise HelperFailure(
422 "absence-unverified",
423 "The project connection still exists after delete.",
424 blocked_at="verification",
425 writes=[write],
426 resources_remaining=[identity],
427 request_id=verify_request_id,
428 partial=True,
429 )
430 return "deleted", identity, request_ids
431
432
433def _recover_ambiguous_connection_delete(
434 url: str,
435 token: str,
436 identity: dict[str, Any],
437 request_ids: list[str],
438 failure: HelperFailure,
439 *,
440 transport: Transport,
441) -> tuple[str, dict[str, Any], list[str]]:
442 if failure.request_id:
443 request_ids.append(failure.request_id)
444 try:
445 after, verify_request_id = _get_connection(
446 url, token, transport=transport
447 )
448 except HelperFailure as readback_failure:
449 raise HelperFailure(
450 "connection-delete-outcome-ambiguous",
451 "Connection deletion and same-identity readback are ambiguous.",
452 blocked_at="verification",
453 resources_remaining=[identity],
454 request_id=readback_failure.request_id or failure.request_id,
455 status=failure.http_status,
456 partial=True,
457 ) from readback_failure
458 if verify_request_id:
459 request_ids.append(verify_request_id)
460 if after is not None:
461 raise HelperFailure(
462 "connection-delete-outcome-ambiguous",
463 "Same-identity readback still found the connection after an ambiguous delete.",
464 blocked_at="verification",
465 resources_remaining=[identity],
466 request_id=verify_request_id or failure.request_id,
467 status=failure.http_status,
468 partial=True,
469 )
470 return "deleted", identity, request_ids
471
472
473def execute(
474 document: dict[str, Any],
475 *,
476 token_provider: TokenProvider = azure_cli_token,
477 transport: Transport = http_request,
478 sdk_loader: Callable[[], tuple[Any, Any, Any, Any, Any]] = load_cleanup_sdk,
479) -> dict[str, Any]:
480 plan = document["plan"]
481 fingerprint = document["_computed_fingerprint"]
482 agent, connection = _validate_plan(plan)
483 resources: dict[str, list[dict[str, Any]]] = {
484 "created": [],
485 "reused": [],
486 "updated": [],
487 "skipped": [],
488 "deleted": [],
489 }
490 writes: list[dict[str, Any]] = []
491 request_ids: list[str] = []
492
493 if connection is not None and plan.get("dependency_guard") is not None:
494 current, _ = _get_connection(
495 _connection_url({**plan, "connection": connection}),
496 token_provider(MANAGEMENT_AUDIENCE), transport=transport,
497 )
498 if current is not None:
499 cleanup_dependencies.verify_prompt(plan, current, sdk_loader=sdk_loader)
500 refreshed, _ = _get_connection(
501 _connection_url({**plan, "connection": connection}),
502 token_provider(MANAGEMENT_AUDIENCE), transport=transport,
503 )
504 if refreshed != current:
505 raise HelperFailure("definition-drift", "Connection changed before the first cleanup write.", blocked_at="reconciliation")
506
507 if agent is not None:
508 action, identity = _delete_agent(plan, agent, sdk_loader=sdk_loader)
509 resources[action].append(identity)
510 if action == "deleted":
511 writes.append({"action": action, **identity})
512
513 if connection is not None:
514 try:
515 token = token_provider(MANAGEMENT_AUDIENCE)
516 action, identity, connection_request_ids = _delete_connection(
517 plan,
518 connection,
519 token,
520 transport=transport,
521 sdk_loader=sdk_loader,
522 )
523 except HelperFailure as failure:
524 raise HelperFailure(
525 failure.code,
526 failure.message,
527 blocked_at=failure.blocked_at,
528 writes=writes + failure.writes,
529 resources_remaining=(
530 [
531 {
532 "type": "project-connection",
533 "name": connection["name"],
534 }
535 ]
536 + [
537 item
538 for item in failure.resources_remaining
539 if item
540 != {
541 "type": "project-connection",
542 "name": connection["name"],
543 }
544 ]
545 ),
546 request_id=failure.request_id,
547 status=failure.http_status,
548 partial=bool(writes or failure.writes or failure.partial),
549 ) from failure
550 resources[action].append(identity)
551 request_ids.extend(connection_request_ids)
552 if action == "deleted":
553 writes.append({"action": action, **identity})
554
555 return {
556 "status": "completed",
557 "outcome": str(plan.get("outcome") or "cleanup-prompt-connection"),
558 "approved_plan": {"fingerprint": fingerprint, "confirmed": True},
559 "resources": resources,
560 "api_contracts": [
561 {
562 "operation": "prompt-agent-version-delete",
563 "version": "azure-ai-projects-2.x",
564 "preview": True,
565 },
566 {
567 "operation": "project-connection-delete",
568 "version": "2025-10-01-preview",
569 "preview": True,
570 },
571 ],
572 "data_movement": {"boundary": "none", "result": "none"},
573 "auth": {"mode": "entra-user", "principals": []},
574 "rbac": {"assignments": []},
575 "network": {"posture": "preserved", "evidence": None},
576 "verification": {
577 "absence": True,
578 "request_ids": request_ids,
579 "idempotency": "already absent run-owned resources are zero-write",
580 },
581 "warnings": [],
582 "ownership": {
583 "run_owned": [],
584 "reused_not_owned": [],
585 "owner": plan.get("owner"),
586 },
587 "cleanup": {
588 "status": "completed",
589 "separate_confirmation_required": True,
590 },
591 }
592
593
594def main(argv: list[str] | None = None) -> int:
595 parser = argparse.ArgumentParser()
596 parser.add_argument("--input", type=Path, required=True)
597 args = parser.parse_args(argv)
598 fingerprint: str | None = None
599 owner: Any = None
600 outcome = "cleanup-prompt-connection"
601 try:
602 document, plan, fingerprint = load_approved_input(args.input)
603 document["_computed_fingerprint"] = fingerprint
604 owner = plan.get("owner")
605 outcome = str(plan.get("outcome") or outcome)
606 result = execute(document)
607 except HelperFailure as failure:
608 result = blocked_result(
609 failure,
610 outcome=outcome,
611 fingerprint=fingerprint,
612 owner=owner,
613 )
614 emit_result(result)
615 return 3 if result["status"] == "partial" else 2
616 emit_result(result)
617 return 0
618
619
620if __name__ == "__main__":
621 sys.exit(main())