Setting the file. One moment.
Search Reconcile · Foundry Iq · microsoft/azure-skills · Skills Docs
ContentsBack to the top of the page File Cu Canary
def _creation_callback_response
— line 418
This file
Number 10.66
Position 66 of 77
Type Python
Size 67 KB
Lines 1,465 helpers/ search_reconcile.py
Python · 1,465 lines · 67 KB
urlencode
12
13 try :
14 from . import cu_ingestion_auth as file_cu_auth
15 from . import _cleanup_dependencies as cleanup_dependencies, _cleanup_receipts as cleanup_receipts
16 from ._common import (
17 SEARCH_AUDIENCE ,
18 HelperFailure,
19 HttpResult,
20 ReadRecovery,
21 TokenProvider,
22 Transport,
23 azure_cli_token,
24 blocked_result,
25 canonical_bytes,
26 digest,
27 emit_result,
28 http_request,
29 is_ambiguous_mutation_failure,
30 load_approved_input,
31 odata_name,
32 reject_secrets,
33 require_allowed_fields,
34 validate_search_endpoint,
35 RESOURCE_ID_CONNECTION ,
36 )
37 except ImportError :
38 import cu_ingestion_auth as file_cu_auth
39 import _cleanup_dependencies as cleanup_dependencies, _cleanup_receipts as cleanup_receipts
40 from _common import ( # type: ignore[no-redef]
41 SEARCH_AUDIENCE ,
42 HelperFailure,
43 HttpResult,
44 ReadRecovery,
45 TokenProvider,
46 Transport,
47 azure_cli_token,
48 blocked_result,
49 canonical_bytes,
50 digest,
51 emit_result,
52 http_request,
53 is_ambiguous_mutation_failure,
54 load_approved_input,
55 odata_name,
56 reject_secrets,
57 require_allowed_fields,
58 validate_search_endpoint,
59 RESOURCE_ID_CONNECTION ,
60 )
61
62
63 SUPPORTED_API_VERSIONS = { "2026-04-01" , "2026-08-01-preview" }
64 RESOURCE_SEGMENTS = {
65 "knowledge-source" : "knowledgesources" ,
66 "knowledge-base" : "knowledgebases" ,
67 }
68 DYNAMIC_FIELDS = {
69 "@odata.context" ,
70 "@odata.etag" ,
71 "currentSynchronizationState" ,
72 "lastSynchronizationState" ,
73 "synchronizationStatus" ,
74 "apiKey" ,
75 "createdResources" ,
76 # Azure Search always returns "<redacted>" for connectionString on GET
77 # (secret redaction), never the submitted value, so it can never be
78 # compared for exact equality against a desired definition.
79 "connectionString" ,
80 }
81 # Endpoint-URI fields where Azure Search's own readback normalization is
82 # inconsistent (resourceUri loses a trailing slash; aiServices.uri keeps
83 # whatever was submitted), so both are compared slash-insensitively.
84 URI_FIELDS = { "resourceUri" , "uri" }
85 SHA256 = re.compile( r " ^ sha256: [ a-f0-9 ] {64} $ " )
86 ENVIRONMENT_NAME = re.compile( r " ^[ A-Z ][ A-Z0-9_ ] * $ " )
87 PLAN_FIELDS = {
88 "operation" ,
89 "outcome" ,
90 "plan_kind" ,
91 "resource_type" ,
92 "endpoint" ,
93 "name" ,
94 "api_version" ,
95 "action" ,
96 "desired" ,
97 "source_evidence" ,
98 "verified_source" ,
99 "expected_etag" ,
100 "owned_definition_digest" ,
101 "ai_services_api_key_environment" ,
102 "ai_services_key_acquisition" ,
103 "ai_services_managed_identity" ,
104 "data_movement" ,
105 "rbac" ,
106 "network" ,
107 "owner" ,
108 "cleanup_approved" ,
109 "dependency_guard" ,
110 "kb_plan_version" ,
111 "kb_model" ,
112 }
113 KB_INTENT_FIELDS = {
114 "schema_version" , "endpoint" , "name" , "owner" , "api_version" , "source_name" ,
115 "reasoning_effort" , "output_mode" , "model" , "description" ,
116 "retrieval_instructions" , "answer_instructions" , "action" ,
117 "data_movement" , "rbac" , "network" ,
118 }
119 CHAT_MODELS = {
120 "gpt-4o" , "gpt-4o-mini" , "gpt-4.1" , "gpt-4.1-mini" , "gpt-4.1-nano" ,
121 "gpt-5" , "gpt-5-mini" , "gpt-5-nano" , "gpt-5.1" , "gpt-5.2" ,
122 "gpt-5.4" , "gpt-5.4-mini" , "gpt-5.4-nano" , "gpt-5.5" ,
123 "gpt-5.6-sol" , "gpt-5.6-terra" , "gpt-5.6-luna" ,
124 }
125
126
127 def _kb_failure (code: str , message: str , request_id: str | None = None ) -> HelperFailure:
128 return HelperFailure(code, message, blocked_at = "knowledge-base-planning" , request_id = request_id)
129
130
131 def _kb_model (value: Any, * , required: bool ) -> dict[ str , Any] | None :
132 if not required:
133 if value is not None :
134 raise _kb_failure( "kb-model-conflict" , "Minimal extractive planning must omit a KB model." )
135 return None
136 if not isinstance (value, dict ):
137 raise _kb_failure( "kb-model-required" , "Low/medium or synthesis requires a selected chat deployment." )
138 require_allowed_fields(value, { "endpoint" , "deployment" , "model" , "auth" , "prerequisites" },
139 label = "KB model choice" )
140 if (
141 not isinstance (value.get( "endpoint" ), str )
142 or re.fullmatch( r "https:// [ a-z0-9 ][ a-z0-9- ] {0,62} \. "
143 r " (?: openai \. azure \. com | services \. ai \. azure \. com | cognitiveservices \. azure \. com ) / ? " ,
144 value[ "endpoint" ]) is None
145 or not isinstance (value.get( "deployment" ), str )
146 or re.fullmatch( r " [ A-Za-z0-9 ][ A-Za-z0-9_.- ] {0,63} " , value[ "deployment" ]) is None
147 or not isinstance (value.get( "model" ), str ) or value[ "model" ] not in CHAT_MODELS
148 or value.get( "auth" ) != "system-assigned"
149 ):
150 raise _kb_failure( "kb-model-invalid" , "Select a supported chat model/deployment and exact Azure endpoint with system-assigned auth." )
151 evidence = value.get( "prerequisites" )
152 if not isinstance (evidence, dict ):
153 raise _kb_failure( "kb-model-evidence-missing" , "Provide owner-verified deployment, identity and network references." )
154 require_allowed_fields(evidence, { "deployment" , "identity" , "network" }, label = "KB model prerequisites" )
155 if any ( not isinstance (evidence.get(k), str ) or not evidence[k].strip() or len (evidence[k]) > 4096
156 for k in ( "deployment" , "identity" , "network" )):
157 raise _kb_failure( "kb-model-evidence-missing" , "Provide owner-verified deployment, identity and network references." )
158 result = copy.deepcopy(value)
159 result[ "endpoint" ] = result[ "endpoint" ].rstrip( "/" )
160 return result
161
162
163 def _kb_model_definition (choice: dict[ str , Any]) -> dict[ str , Any]:
164 try :
165 from .source_vector import model_definition
166 except ImportError :
167 from source_vector import model_definition
168 return model_definition(choice)
169
170
171 def _kb_guard (plan: dict[ str , Any], transport: Transport) -> Transport:
172 target = _resource_url(plan)
173
174 def guarded (method, url, token, ** kwargs):
175 result = transport(method, url, token, ** kwargs)
176 if method == "GET" and url == target and result.status == 200 and isinstance (result.body, dict ):
177 models = result.body.get( "models" )
178 for model in models if isinstance (models, list ) else []:
179 parameters = model.get( "azureOpenAIParameters" ) if isinstance (model, dict ) else None
180 if isinstance (parameters, dict ) and (
181 parameters.get( "apiKey" ) not in ( None , "" )
182 or parameters.get( "authIdentity" ) is not None
183 ):
184 raise _kb_failure( "kb-model-auth-conflict" ,
185 "KB readback does not prove keyless system-assigned model auth; details withheld." ,
186 result.request_id)
187 return result
188
189 return guarded
190
191
192 def plan_knowledge_base (
193 request: dict[ str , Any], * , token_provider: TokenProvider = azure_cli_token,
194 transport: Transport = http_request,
195 ) -> dict[ str , Any]:
196 """Compile resolved KB choices into an unapproved artifact using exact GETs."""
197 if not isinstance (request, dict ) or request.get( "schema_version" ) != "1.0" :
198 raise _kb_failure( "input-schema-invalid" , "KB planning requires an object with schema_version 1.0." )
199 try :
200 json.dumps(request, ensure_ascii = False , allow_nan = False ).encode( "utf-8" )
201 except ( TypeError , ValueError , UnicodeError , RecursionError ) as exc:
202 raise _kb_failure( "input-schema-invalid" , "KB choices must be finite UTF-8 JSON." ) from exc
203 reject_secrets(request)
204 require_allowed_fields(request, KB_INTENT_FIELDS , label = "KB planning input" )
205 for field in ( "name" , "source_name" ):
206 _odata_name(request.get(field))
207 if not isinstance (request.get( "owner" ), str ) or not request[ "owner" ].strip() or len (request[ "owner" ]) > 4096 :
208 raise _kb_failure( "input-schema-invalid" , "An explicit nonempty owner is required." )
209 action = request.get( "action" , "create-or-reuse" )
210 effort, output = request.get( "reasoning_effort" ), request.get( "output_mode" )
211 version = request.get( "api_version" )
212 if (
213 not isinstance (version, str ) or version not in SUPPORTED_API_VERSIONS
214 or action not in ( "create-or-reuse" , "agent-minimal-transition" )
215 or effort not in ( "minimal" , "low" , "medium" )
216 or output not in ( "extractiveData" , "answerSynthesis" )
217 or version == "2026-04-01" and (effort != "minimal" or output != "extractiveData" )
218 or action == "agent-minimal-transition" and (
219 version != "2026-08-01-preview" or effort != "minimal" or output != "extractiveData"
220 )
221 ):
222 raise _kb_failure( "mode-api-mismatch" , "Choose supported KB API/output/effort; transition is preview minimal/extractive only." )
223 model = _kb_model(request.get( "model" ), required = effort in ( "low" , "medium" ) or output == "answerSynthesis" )
224 desired: dict[ str , Any] = { "name" : request[ "name" ], "knowledgeSources" : [{ "name" : request[ "source_name" ]}]}
225 if version == "2026-08-01-preview" :
226 desired.update( outputMode = output, retrievalReasoningEffort = { "kind" : effort})
227 for key, wire in (( "description" , "description" ), ( "retrieval_instructions" , "retrievalInstructions" ),
228 ( "answer_instructions" , "answerInstructions" )):
229 if key in request:
230 value = request[key]
231 if value is not None and ( not isinstance (value, str ) or len (value) > 4096 ):
232 raise _kb_failure( "input-schema-invalid" , "KB description/instructions must be bounded text or null." )
233 if action == "agent-minimal-transition" or version == "2026-04-01" and key != "description" :
234 raise _kb_failure( "mode-api-mismatch" , "These optional fields are not supported in this planning mode." )
235 desired[wire] = value
236 if model is not None :
237 desired[ "models" ] = [_kb_model_definition(model)]
238 plan = {
239 "operation" : "reconcile" , "outcome" : "create-knowledge-base" , "resource_type" : "knowledge-base" ,
240 "endpoint" : validate_search_endpoint(request.get( "endpoint" )), "name" : request[ "name" ],
241 "api_version" : version, "owner" : request[ "owner" ], "cleanup_approved" : False ,
242 "action" : "create" , "desired" : desired, "kb_plan_version" : "1.0" , "kb_model" : model,
243 ** {k: copy.deepcopy(request[k]) for k in ( "data_movement" , "rbac" , "network" ) if k in request},
244 }
245 target = _resource_url(plan)
246 source_url = _resource_url({ ** plan, "resource_type" : "knowledge-source" , "name" : request[ "source_name" ]})
247 # Validate all local sections before authentication; no invented source proof.
248 for key, fields in (( "data_movement" , { "boundary" , "result" }), ( "rbac" , { "assignments" }),
249 ( "network" , { "posture" , "evidence" })):
250 if key in plan:
251 if not isinstance (plan[key], dict ):
252 raise _kb_failure( "input-schema-invalid" , "Optional controls must be objects." )
253 require_allowed_fields(plan[key], fields, label = key)
254 transport = _kb_guard(plan, transport)
255 token = token_provider( SEARCH_AUDIENCE )
256 source, source_id = _get(source_url, token, transport = transport)
257 if source is None or source.get( "name" ) != request[ "source_name" ] or source.get( "kind" ) not in ( "file" , "azureBlob" ):
258 raise _kb_failure( "source-unverified" , "Select an existing supported source with exact readback." , source_id)
259 parameters = source.get( "azureBlobParameters" )
260 if source[ "kind" ] == "azureBlob" and (
261 not isinstance (parameters, dict ) or not isinstance (parameters.get( "isADLSGen2" , False ), bool )
262 ):
263 raise _kb_failure( "source-unverified" , "Blob source type/boundary metadata is malformed." , source_id)
264 if version == "2026-04-01" and (
265 source[ "kind" ] != "azureBlob" or parameters.get( "isADLSGen2" , False )
266 ):
267 raise _kb_failure( "mode-api-mismatch" , "File and ADLS knowledge bases require supported preview." , source_id)
268 source_digest = digest(_definition(source))
269 plan[ "verified_source" ] = { "name" : request[ "source_name" ], "verified" : True , "definition_digest" : source_digest}
270 current, current_id = _get(target, token, transport = transport)
271 if action == "agent-minimal-transition" :
272 if current is None :
273 raise _kb_failure( "target-absent" , "Transition requires the exact existing KB." , current_id)
274 if (
275 current.get( "name" ) != plan[ "name" ]
276 or _definition(current.get( "knowledgeSources" )) != desired[ "knowledgeSources" ]
277 or current.get( "models" ) not in ( None , [])
278 or current.get( "outputMode" ) not in ( None , "extractiveData" )
279 or current.get( "retrievalReasoningEffort" ) not in ( None , { "kind" : "minimal" }, { "kind" : "low" })
280 ):
281 raise _kb_failure( "definition-conflict" , "Transition requires a model-free extractive KB with absent/minimal/low effort." , current_id)
282 desired = {k: copy.deepcopy(v) for k, v in current.items() if k not in ( "@odata.etag" , "@odata.context" )}
283 desired.update( outputMode = "extractiveData" , retrievalReasoningEffort = { "kind" : "minimal" })
284 plan[ "desired" ] = desired
285 plan[ "action" ] = "update"
286 if current is not None :
287 if not definitions_match(desired, current) and action != "agent-minimal-transition" :
288 raise _kb_failure( "definition-conflict" , "The exact KB differs; never overwrite or choose another name." , current_id)
289 etag = current.get( "@odata.etag" )
290 if not isinstance (etag, str ) or not etag.strip():
291 raise _kb_failure( "definition-evidence-missing" , "Existing KB readback requires an ETag." , current_id)
292 plan[ "expected_etag" ] = etag
293 if definitions_match(desired, current):
294 plan[ "action" ] = "reuse"
295 refreshed, refresh_id = _get(source_url, token, transport = transport)
296 if refreshed is None or digest(_definition(refreshed)) != source_digest:
297 raise _kb_failure( "source-drift" , "Source changed during planning; refresh its evidence before approval." , refresh_id)
298 request_ids = [source_id, current_id, refresh_id]
299 if current is not None :
300 after, after_id = _get(target, token, transport = transport)
301 request_ids.append(after_id)
302 if after is None or after.get( "@odata.etag" ) != plan[ "expected_etag" ] or not definitions_match(current, after):
303 raise _kb_failure( "definition-drift" , "KB changed during planning; discard this proposal." , after_id)
304 _validate_plan(plan)
305 fingerprint = digest(plan)
306 mutation_required = plan[ "action" ] != "reuse"
307 return {
308 "status" : "planned" , "outcome" : plan[ "outcome" ], "plan_fingerprint" : fingerprint,
309 "execution_input" : { "schema_version" : "1.0" , "plan" : plan,
310 "approval" : { "confirmed" : False , "fingerprint" : fingerprint}},
311 "approval_summary" : {
312 "target" : { "endpoint" : plan[ "endpoint" ], "name" : plan[ "name" ], "api_version" : version},
313 "source" : { "name" : request[ "source_name" ], "kind" : source[ "kind" ]},
314 "owner" : plan[ "owner" ], "action" : plan[ "action" ],
315 "execution_required" : mutation_required, "mutation_approval_required" : mutation_required,
316 "reasoning_effort" : effort, "output_mode" : output,
317 ** ({ "previous_mode" : {
318 "output_mode" : current.get( "outputMode" ),
319 "reasoning_effort" : current.get( "retrievalReasoningEffort" ),
320 }} if action == "agent-minimal-transition" else {}),
321 "model" : {k: model[k] for k in ( "endpoint" , "deployment" , "model" , "auth" )} if model else None ,
322 "cost_and_data" : "Selected KB chat use is billable; retrieved content may move to its endpoint." if model
323 else "No KB chat model; existing Search/source/storage charges still apply." ,
324 "remaining_checks" : [ "Source ingestion readiness" , "Effective access/network/model prerequisites" ,
325 "Supported and unrelated KB queries with original-source citations" ],
326 "cleanup" : "Not approved; retain reused/shared sources, models, roles and data." ,
327 },
328 "verification" : { "source_definition" : "verified" , "knowledge_base" : "reused" if not mutation_required else "unverified" ,
329 "ingestion" : "unverified" , "retrieval" : "unverified" ,
330 "request_ids" : [r for r in request_ids if r]},
331 "writes_performed" : [],
332 }
333
334
335 def _odata_name (name: Any) -> str :
336 return odata_name(name)
337
338
339 def _resource_url (plan: dict[ str , Any]) -> str :
340 endpoint = validate_search_endpoint(plan.get( "endpoint" ))
341 resource_type = plan.get( "resource_type" )
342 segment = RESOURCE_SEGMENTS .get(resource_type)
343 if segment is None :
344 raise HelperFailure(
345 "resource-type-invalid" ,
346 "resource_type must be knowledge-source or knowledge-base." ,
347 blocked_at = "input-resolution" ,
348 )
349 api_version = plan.get( "api_version" )
350 if api_version not in SUPPORTED_API_VERSIONS :
351 raise HelperFailure(
352 "api-version-invalid" ,
353 "API version must be 2026-04-01 or 2026-08-01-preview." ,
354 blocked_at = "input-resolution" ,
355 )
356 return (
357 f " { endpoint } / { segment } (' { _odata_name(plan.get( 'name' )) } ')?"
358 + urlencode({ "api-version" : api_version})
359 )
360
361
362 def _definition (value: Any, * , key: str | None = None ) -> Any:
363 if isinstance (value, dict ):
364 result = {
365 child_key: _definition(child, key = child_key)
366 for child_key, child in sorted (value.items())
367 if child_key not in DYNAMIC_FIELDS and child is not None
368 }
369 # Preview returns this documented default even when omitted on creation.
370 if value.get( "kind" ) == "azureBlob" and result.get( "resultsProcessing" ) == "rerank" :
371 result.pop( "resultsProcessing" )
372 return {
373 child_key: child
374 for child_key, child in result.items()
375 if child not in ({}, [])
376 }
377 if isinstance (value, list ):
378 return [_definition(child) for child in value]
379 # Azure Search silently strips a single trailing slash from
380 # azureOpenAIParameters.resourceUri on readback (while preserving it
381 # verbatim on aiServices.uri), so a byte-exact comparison would
382 # false-negative on functionally identical endpoints that differ only
383 # by a trailing slash. Normalize both known endpoint-URI field names.
384 if (
385 key in URI_FIELDS
386 and isinstance (value, str )
387 and value.endswith( "/" )
388 and len (value) > 1
389 ):
390 # Remove only one trailing slash; preserve intentional extra
391 # slashes (e.g. "https://example.com//") for exact comparison.
392 return value[: - 1 ]
393 return value
394
395
396 def response_etags (result: HttpResult) -> dict[ str , Any]:
397 return {
398 "body" : result.body.get( "@odata.etag" ) if isinstance (result.body, dict ) else None ,
399 "headers" : list (result.etag_values) if result.etag_values is not None else [
400 value for name, value in result.headers.items() if name.lower() == "etag"
401 ],
402 }
403
404
405 def resolve_etag (evidence: dict[ str , Any], request_id: str | None = None ) -> str | None :
406 values = [ * evidence[ "headers" ]]
407 if evidence[ "body" ] is not None :
408 values.append(evidence[ "body" ])
409 if any ( not isinstance (value, str ) or not value.strip() for value in values):
410 raise HelperFailure( "etag-invalid" , "Search returned malformed version evidence." ,
411 blocked_at = "verification" , request_id = request_id)
412 if len ( set (values)) > 1 :
413 raise HelperFailure( "etag-conflict" , "Search response header/body ETags conflict." ,
414 blocked_at = "verification" , request_id = request_id)
415 return values[ 0 ] if values else None
416
417
418 def _creation_callback_response (result: HttpResult) -> HttpResult:
419 """Project provenance only, not credential-bearing service bodies or headers."""
420 evidence = response_etags(result)
421 # Preserve invalidity without forwarding malformed containers that could contain credentials.
422 body_etag = evidence[ "body" ] if evidence[ "body" ] is None or isinstance (evidence[ "body" ], str ) else []
423 etags = tuple (value if isinstance (value, str ) else "" for value in evidence[ "headers" ])
424 parameters = result.body.get( "azureBlobParameters" ) if isinstance (result.body, dict ) else None
425 generated = parameters.get( "createdResources" ) if isinstance (parameters, dict ) else None
426 names = {
427 kind: name for kind, name in (generated.items() if isinstance (generated, dict ) else ())
428 if kind in { "datasource" , "dataSourceConnection" , "indexer" , "skillset" , "index" } and isinstance (name, str )
429 }
430 headers = { "request-id" : result.request_id} if isinstance (result.request_id, str ) else {}
431 return HttpResult(result.status, {
432 "@odata.etag" : body_etag, "azureBlobParameters" : { "createdResources" : names},
433 }, headers, etags)
434
435
436 def _get (
437 url: str ,
438 token: str ,
439 * ,
440 transport: Transport,
441 recovery: ReadRecovery | None = None ,
442 ) -> tuple[dict[ str , Any] | None , str | None ]:
443 try :
444 result = (recovery.get(url, token, transport = transport) if recovery is not None
445 else transport( "GET" , url, token))
446 except HelperFailure as failure:
447 if failure.http_status == 404 and failure.blocked_at != "local-persistence" :
448 return None , failure.request_id
449 raise
450 if result.status != 200 or not isinstance (result.body, dict ):
451 raise HelperFailure(
452 "readback-invalid" ,
453 "Search resource readback did not return one JSON object." ,
454 blocked_at = "reconciliation" ,
455 request_id = result.request_id,
456 status = result.status,
457 )
458 etag = resolve_etag(response_etags(result), result.request_id)
459 body = { ** result.body, "@odata.etag" : etag} if etag is not None else result.body
460 return body, result.request_id
461
462
463 def _validate_plan (plan: dict[ str , Any]) -> None :
464 reject_secrets(plan)
465 require_allowed_fields(plan, PLAN_FIELDS , label = "Search reconciliation plan" )
466 acquisition = plan.get( "ai_services_key_acquisition" )
467 mi = plan.get( "ai_services_managed_identity" )
468 if mi is not None and (
469 mi is not True or plan.get( "operation" ) != "reconcile"
470 or plan.get( "resource_type" ) != "knowledge-source"
471 or not isinstance (plan.get( "desired" ), dict ) or plan[ "desired" ].get( "kind" ) != "file"
472 or acquisition is not None or plan.get( "ai_services_api_key_environment" ) is not None
473 ):
474 raise file_cu_auth.failure( "cu-mi-plan-invalid" , "File MI metadata is exclusive to the verified File source workflow, without key channels." )
475 if acquisition is not None and (
476 plan.get( "operation" ) != "reconcile" or plan.get( "resource_type" ) != "knowledge-source"
477 or not isinstance (plan.get( "desired" ), dict ) or plan[ "desired" ].get( "kind" ) != "file"
478 or plan.get( "ai_services_api_key_environment" ) is not None
479 ):
480 raise file_cu_auth.failure( "cu-acquisition-invalid" , "Private acquisition is exclusive to Standard File source reconciliation, not other resources or auth channels." )
481 cleanup_dependencies.validate_guard(plan)
482 if "kb_plan_version" in plan or "kb_model" in plan:
483 if (
484 plan.get( "kb_plan_version" ) != "1.0" or "kb_model" not in plan
485 or plan.get( "resource_type" ) != "knowledge-base" or plan.get( "operation" ) != "reconcile"
486 ):
487 raise _kb_failure( "input-schema-invalid" , "KB planning metadata requires the versioned KB reconciliation contract." )
488 desired = plan.get( "desired" )
489 if not isinstance (desired, dict ):
490 raise _kb_failure( "desired-definition-invalid" , "KB desired definition must be an object." )
491 effort = desired.get( "retrievalReasoningEffort" )
492 required = desired.get( "outputMode" ) == "answerSynthesis" or (
493 isinstance (effort, dict ) and effort.get( "kind" ) in ( "low" , "medium" )
494 )
495 model = _kb_model(plan[ "kb_model" ], required = required)
496 expected = [_kb_model_definition(model)] if model else []
497 if _definition(desired.get( "models" , [])) != _definition(expected):
498 raise _kb_failure( "kb-model-conflict" , "KB wire models must match the selected independent chat configuration." )
499 for field, allowed in (
500 ( "data_movement" , { "boundary" , "result" }),
501 ( "rbac" , { "assignments" }),
502 ( "network" , { "posture" , "evidence" }),
503 ):
504 section = plan.get(field)
505 if section is not None :
506 if not isinstance (section, dict ):
507 raise HelperFailure(
508 "input-schema-invalid" ,
509 f " { field } must be an object." ,
510 blocked_at = "input-resolution" ,
511 )
512 require_allowed_fields(section, allowed, label = field)
513 operation = plan.get( "operation" )
514 if operation not in { "reconcile" , "delete" }:
515 raise HelperFailure(
516 "operation-invalid" ,
517 "operation must be reconcile or delete." ,
518 blocked_at = "input-resolution" ,
519 )
520 if operation == "reconcile" :
521 if plan.get( "cleanup_approved" ) is not False :
522 raise HelperFailure(
523 "cleanup-boundary-invalid" ,
524 "A reconciliation plan must set cleanup_approved to false." ,
525 blocked_at = "confirmation" ,
526 )
527 desired = plan.get( "desired" )
528 if not isinstance (desired, dict ) or desired.get( "name" ) != plan.get( "name" ):
529 raise HelperFailure(
530 "desired-definition-invalid" ,
531 "desired must be an object whose name matches the target name." ,
532 blocked_at = "input-resolution" ,
533 )
534 if plan.get( "resource_type" ) == "knowledge-source" :
535 kind = desired.get( "kind" )
536 if kind not in { "file" , "azureBlob" }:
537 raise HelperFailure(
538 "source-kind-invalid" ,
539 "Only file and azureBlob knowledge sources are supported." ,
540 blocked_at = "input-resolution" ,
541 )
542 if kind == "file" and plan.get( "api_version" ) != "2026-08-01-preview" :
543 raise HelperFailure(
544 "api-version-invalid" ,
545 "This helper uses 2026-08-01-preview for multipart metadata and 200-file inventories. "
546 "The service also supports 2026-05-01-preview minimal; use a compatible client or explicitly approve August." ,
547 blocked_at = "input-resolution" ,
548 )
549 if kind == "file" :
550 ingestion = (
551 desired.get( "fileParameters" , {}).get( "ingestionParameters" , {})
552 if isinstance (desired.get( "fileParameters" ), dict )
553 else {}
554 )
555 if not isinstance (ingestion, dict ):
556 raise HelperFailure(
557 "desired-definition-invalid" ,
558 "File ingestionParameters must be an object." ,
559 blocked_at = "input-resolution" ,
560 )
561 mode = ingestion.get( "contentExtractionMode" )
562 if mode not in { "minimal" , "standard" }:
563 raise HelperFailure(
564 "desired-definition-invalid" ,
565 "File contentExtractionMode must be minimal or standard." ,
566 blocked_at = "input-resolution" ,
567 )
568 credential_environment = plan.get( "ai_services_api_key_environment" )
569 if mode == "standard" :
570 if (
571 (acquisition is None and mi is not True and (
572 not isinstance (credential_environment, str )
573 or ENVIRONMENT_NAME .fullmatch(credential_environment) is None
574 ))
575 or not isinstance (ingestion.get( "aiServices" ), dict )
576 or not ingestion[ "aiServices" ].get( "uri" )
577 ):
578 raise HelperFailure(
579 "credential-channel-invalid" ,
580 "Standard extraction requires aiServices.uri and an approved private ARM acquisition or existing API-key environment variable name." ,
581 blocked_at = "input-resolution" ,
582 )
583 if acquisition is not None :
584 file_cu_auth.validate_acquisition(acquisition, ingestion[ "aiServices" ][ "uri" ])
585 elif credential_environment is not None or acquisition is not None or mi is not None :
586 raise HelperFailure(
587 "credential-channel-invalid" ,
588 "Minimal extraction must not declare an AI Services credential channel." ,
589 blocked_at = "input-resolution" ,
590 )
591 else :
592 if plan.get( "ai_services_api_key_environment" ) is not None :
593 raise HelperFailure(
594 "credential-channel-invalid" , "Blob ingestion uses keyless dependencies, not File credential channels." ,
595 blocked_at = "input-resolution" ,
596 )
597 parameters = desired.get( "azureBlobParameters" )
598 evidence = plan.get( "source_evidence" )
599 if isinstance (evidence, dict ):
600 require_allowed_fields(
601 evidence,
602 {
603 "verified" ,
604 "inventory_digest" ,
605 "path_verified" ,
606 "acl_verified" ,
607 },
608 label = "Source evidence" ,
609 )
610 if (
611 not isinstance (parameters, dict )
612 or not isinstance (parameters.get( "connectionString" ), str )
613 or RESOURCE_ID_CONNECTION .fullmatch(
614 parameters[ "connectionString" ]
615 )
616 is None
617 or not isinstance (parameters.get( "containerName" ), str )
618 or not parameters[ "containerName" ]
619 or "folderPath" not in parameters
620 or not isinstance (parameters.get( "isADLSGen2" ), bool )
621 or not isinstance (evidence, dict )
622 or evidence.get( "verified" ) is not True
623 or SHA256 .fullmatch( str (evidence.get( "inventory_digest" ))) is None
624 ):
625 raise HelperFailure(
626 "source-evidence-invalid" ,
627 "Blob and ADLS plans require an exact ResourceId boundary and verified immutable inventory evidence." ,
628 blocked_at = "input-resolution" ,
629 )
630 if parameters[ "isADLSGen2" ] and (
631 evidence.get( "path_verified" ) is not True
632 or evidence.get( "acl_verified" ) is not True
633 ):
634 raise HelperFailure(
635 "adls-evidence-unverified" ,
636 "ADLS reconciliation requires exact path and ACL readback evidence." ,
637 blocked_at = "reconciliation" ,
638 )
639 if plan.get( "resource_type" ) == "knowledge-base" :
640 sources = desired.get( "knowledgeSources" )
641 if not isinstance (sources, list ) or len (sources) != 1 :
642 raise HelperFailure(
643 "source-count-invalid" ,
644 "A knowledge base must name exactly one knowledge source." ,
645 blocked_at = "input-resolution" ,
646 )
647 verified_source = plan.get( "verified_source" )
648 if isinstance (verified_source, dict ):
649 require_allowed_fields(
650 verified_source,
651 { "name" , "verified" , "definition_digest" },
652 label = "Verified source" ,
653 )
654 if (
655 not isinstance (sources[ 0 ], dict )
656 or not isinstance (verified_source, dict )
657 or verified_source.get( "name" ) != sources[ 0 ].get( "name" )
658 or verified_source.get( "verified" ) is not True
659 or SHA256 .fullmatch( str (verified_source.get( "definition_digest" )))
660 is None
661 ):
662 raise HelperFailure(
663 "source-drift" ,
664 "Knowledge-base reconciliation requires exact verified source readback." ,
665 blocked_at = "reconciliation" ,
666 )
667 _odata_name(verified_source.get( "name" ))
668 if plan.get( "api_version" ) == "2026-04-01" :
669 preview_fields = {
670 "outputMode" ,
671 "retrievalReasoningEffort" ,
672 "retrievalInstructions" ,
673 "answerInstructions" ,
674 }
675 if preview_fields.intersection(desired) or desired.get( "models" ):
676 raise HelperFailure(
677 "mode-api-mismatch" ,
678 "GA knowledge bases must omit preview output, reasoning, instruction, and model fields." ,
679 blocked_at = "input-resolution" ,
680 )
681 else :
682 mode = desired.get( "outputMode" )
683 effort = desired.get( "retrievalReasoningEffort" )
684 if (
685 mode not in { "extractiveData" , "answerSynthesis" }
686 or not isinstance (effort, dict )
687 or effort.get( "kind" ) not in { "minimal" , "low" , "medium" }
688 or (
689 (mode == "answerSynthesis" or effort.get( "kind" ) in { "low" , "medium" })
690 and (
691 not isinstance (desired.get( "models" ), list )
692 or len (desired[ "models" ]) != 1
693 )
694 )
695 ):
696 raise HelperFailure(
697 "mode-api-mismatch" ,
698 "Preview knowledge-base output mode, reasoning effort, and model selection are inconsistent." ,
699 blocked_at = "input-resolution" ,
700 )
701 else :
702 if plan.get( "plan_kind" ) != "cleanup" or plan.get( "cleanup_approved" ) is not True :
703 raise HelperFailure(
704 "cleanup-approval-mismatch" ,
705 "Delete requires a separate cleanup plan with cleanup_approved true." ,
706 blocked_at = "confirmation" ,
707 )
708 if not isinstance (plan.get( "owned_definition_digest" ), str ):
709 raise HelperFailure(
710 "ownership-unproven" ,
711 "Delete requires the approved owned definition digest." ,
712 blocked_at = "reconciliation" ,
713 )
714
715
716 def resource_url (plan: dict[ str , Any]) -> str :
717 """Validate and address the exact selected Search resource."""
718 return _resource_url(plan)
719
720
721 def read_resource (
722 url: str , token: str , * , transport: Transport
723 ) -> tuple[dict[ str , Any] | None , str | None ]:
724 """Read one identity; only a definitive 404 means absent."""
725 return _get(url, token, transport = transport)
726
727
728 def definitions_match (desired: dict[ str , Any], current: dict[ str , Any]) -> bool :
729 return _definition(desired) == _definition(current)
730
731
732 def execute (
733 document: dict[ str , Any],
734 * ,
735 token_provider: TokenProvider = azure_cli_token,
736 transport: Transport = http_request,
737 credential_provider = None ,
738 managed_identity_verified: bool = False ,
739 on_created: Callable[ ... , None ] | None = None ,
740 cleanup_capture: cleanup_receipts.Capture | None = None ,
741 on_file_acknowledged: Callable[ ... , None ] | None = None ,
742 ) -> dict[ str , Any]:
743 """Full callbacks are keyless Blob/File MI only; File ACK callbacks receive provenance only."""
744 guard = document.get( "plan" , {}).get( "dependency_guard" )
745 if guard is None :
746 return _execute(document, token_provider = token_provider, transport = transport, on_created = on_created,
747 cleanup_capture = cleanup_capture, credential_provider = credential_provider,
748 on_file_acknowledged = on_file_acknowledged,
749 managed_identity_verified = managed_identity_verified)
750 tokens: dict[ str , str ] = {}
751
752 def cached_token (audience: str ) -> str :
753 if audience not in tokens:
754 tokens[audience] = token_provider(audience)
755 return tokens[audience]
756
757 try :
758 result = _execute(document, token_provider = cached_token, transport = transport, on_created = on_created,
759 cleanup_capture = cleanup_capture, credential_provider = credential_provider,
760 on_file_acknowledged = on_file_acknowledged,
761 managed_identity_verified = managed_identity_verified)
762 except HelperFailure as failure:
763 if isinstance (guard, dict ) and guard.get( "kind" ) == "search-source" and (failure.partial or failure.writes):
764 failure.resources_remaining.extend(
765 { "type" : item[ "type" ], "name" : item[ "name" ], "absence" : "unverified" }
766 for item in guard[ "generated" ]
767 )
768 raise
769 if guard is None or guard.get( "kind" ) != "search-source" :
770 return result
771 plan = document[ "plan" ]
772 written = bool (result[ "resources" ].get( "deleted" ))
773 writes = [{ "action" : "deleted" , "type" : "knowledge-source" , "name" : plan[ "name" ]}] if written else []
774 token = cached_token( SEARCH_AUDIENCE )
775 verified = []
776 for position, item in enumerate (guard[ "generated" ]):
777 pending = [
778 f "generated-absence-unverified: { later[ 'type' ] } : { later[ 'name' ] } "
779 for later in guard[ "generated" ][position + 1 :]
780 ]
781 url = (
782 f " { plan[ 'endpoint' ].rstrip( '/' ) } / { cleanup_dependencies. COLLECTIONS [item[ 'type' ]] } "
783 f "(' { odata_name(item[ 'name' ]) } ')?api-version= { plan[ 'api_version' ] } "
784 )
785 identity = { "type" : item[ "type" ], "name" : item[ "name" ], "service_managed" : True }
786 try :
787 child, _ = read_resource(url, token, transport = transport)
788 except HelperFailure as failure:
789 raise HelperFailure(
790 failure.code, "Source is absent but generated-child absence readback failed." ,
791 blocked_at = "verification" , writes = writes, partial = written,
792 resources_remaining = [{ ** identity, "absence" : "unverified" }],
793 warnings = pending,
794 request_id = failure.request_id, status = failure.http_status,
795 ) from failure
796 if child is not None :
797 unchanged = child.get( "@odata.etag" ) == item[ "etag" ] and digest(child) == item[ "definition_digest" ]
798 raise HelperFailure(
799 "generated-absence-unverified" , "Source is absent but a generated identity remains; never delete it independently." ,
800 blocked_at = "verification" , writes = writes, partial = written,
801 resources_remaining = [identity] if unchanged else [],
802 warnings = pending + ([] if unchanged else [ f "generated-incarnation-unowned: { item[ 'type' ] } : { item[ 'name' ] } " ]),
803 )
804 verified.append(identity)
805 result[ "verification" ][ "generated_absence" ] = True
806 result[ "resources" ][ "deleted" if written else "skipped" ].extend(verified)
807 return result
808
809
810 def _execute (
811 document: dict[ str , Any],
812 * ,
813 token_provider: TokenProvider = azure_cli_token,
814 transport: Transport = http_request,
815 credential_provider = None ,
816 managed_identity_verified: bool = False ,
817 on_created: Callable[ ... , None ] | None = None ,
818 cleanup_capture: cleanup_receipts.Capture | None = None ,
819 on_file_acknowledged: Callable[ ... , None ] | None = None ,
820 ) -> dict[ str , Any]:
821 plan = document[ "plan" ]
822 fingerprint = document[ "_computed_fingerprint" ]
823 _validate_plan(plan)
824 if on_file_acknowledged is not None and (
825 not callable (on_file_acknowledged) or plan.get( "operation" ) != "reconcile"
826 or plan.get( "action" ) != "create" or plan.get( "resource_type" ) != "knowledge-source"
827 or plan.get( "desired" , {}).get( "kind" ) != "file"
828 ):
829 raise HelperFailure( "creation-callback-unsupported" , "Filtered File ACK capture requires an approved File create." ,
830 blocked_at = "input-resolution" )
831 acquisition = plan.get( "ai_services_key_acquisition" )
832 if (plan.get( "ai_services_managed_identity" ) is True ) != (managed_identity_verified is True ):
833 raise file_cu_auth.failure( "cu-mi-executor-required" , "MI requires the approved versioned File executor and fresh Search identity/CU role readbacks." )
834 if cleanup_capture is not None and (
835 not isinstance (cleanup_capture, cleanup_receipts.Capture)
836 or cleanup_capture.plan_digest != fingerprint or cleanup_capture.owner != plan.get( "owner" )
837 or plan.get( "operation" ) != "reconcile"
838 ):
839 raise HelperFailure( "cleanup-receipt-input-invalid" , "Capture must bind this approved creation operation." , blocked_at = "confirmation" )
840 if on_created is not None and (
841 not callable (on_created) or plan.get( "operation" ) != "reconcile"
842 or plan.get( "resource_type" ) != "knowledge-source"
843 or (plan.get( "desired" , {}).get( "kind" ) != "azureBlob"
844 and not (plan.get( "ai_services_managed_identity" ) is True and managed_identity_verified is True ))
845 or plan.get( "ai_services_api_key_environment" ) is not None
846 or acquisition is not None
847 ):
848 raise HelperFailure(
849 "creation-callback-unsupported" , "Creation callbacks support only keyless Blob or verified File MI plans, never credentialized File/CU wire." ,
850 blocked_at = "input-resolution" ,
851 )
852 if (acquisition is not None ) != (credential_provider is not None ):
853 raise file_cu_auth.failure( "cu-executor-required" , "Private acquisition requires the approved versioned file_source executor; no standalone credential reads." )
854 url = _resource_url(plan)
855 if plan.get( "kb_plan_version" ) == "1.0" :
856 transport = _kb_guard(plan, transport)
857 token = token_provider( SEARCH_AUDIENCE )
858 current, initial_request_id = _get(url, token, transport = transport)
859 operation = plan[ "operation" ]
860 outcome = str (plan.get( "outcome" ) or f "search- { operation } " )
861 owner = plan.get( "owner" )
862
863 if operation == "delete" :
864 if current is None :
865 return _completed(
866 outcome,
867 fingerprint,
868 plan,
869 action = "skipped" ,
870 readback = None ,
871 request_ids = [initial_request_id],
872 absence = True ,
873 )
874 current_definition = _definition(current)
875 if digest(current_definition) != plan[ "owned_definition_digest" ]:
876 raise HelperFailure(
877 "ownership-unproven" ,
878 "Current definition no longer matches the approved owned definition." ,
879 blocked_at = "reconciliation" ,
880 )
881 expected_etag = plan.get( "expected_etag" )
882 if not expected_etag or current.get( "@odata.etag" ) != expected_etag:
883 raise HelperFailure(
884 "definition-drift" ,
885 "Current ETag does not match the approved cleanup plan." ,
886 blocked_at = "reconciliation" ,
887 )
888 if plan.get( "dependency_guard" ) is not None :
889 cleanup_dependencies.verify_search(plan, current, token, transport = transport)
890 try :
891 result = transport(
892 "DELETE" ,
893 url,
894 token,
895 headers = { "If-Match" : expected_etag},
896 )
897 except HelperFailure as failure:
898 if failure.http_status == 404 :
899 return _completed(
900 outcome,
901 fingerprint,
902 plan,
903 action = "skipped" ,
904 readback = None ,
905 request_ids = [initial_request_id, failure.request_id],
906 absence = True ,
907 )
908 if not is_ambiguous_mutation_failure(failure):
909 raise
910 return _recover_ambiguous_delete(
911 outcome,
912 fingerprint,
913 plan,
914 url,
915 token,
916 initial_request_id,
917 failure,
918 transport = transport,
919 )
920 if result.status not in { 200 , 204 }:
921 if result.status == 404 :
922 return _completed(
923 outcome,
924 fingerprint,
925 plan,
926 action = "skipped" ,
927 readback = None ,
928 request_ids = [initial_request_id, result.request_id],
929 absence = True ,
930 )
931 if result.status in { 408 , 429 } or result.status >= 500 :
932 return _recover_ambiguous_delete(
933 outcome,
934 fingerprint,
935 plan,
936 url,
937 token,
938 initial_request_id,
939 HelperFailure(
940 "delete-outcome-ambiguous" ,
941 f "Delete returned ambiguous HTTP { result.status } ." ,
942 blocked_at = "execution" ,
943 request_id = result.request_id,
944 status = result.status,
945 partial = True ,
946 ),
947 transport = transport,
948 )
949 raise HelperFailure(
950 "delete-failed" ,
951 f "Delete returned unexpected HTTP { result.status } ." ,
952 blocked_at = "execution" ,
953 request_id = result.request_id,
954 status = result.status,
955 )
956 write = { "action" : "deleted" , "name" : plan[ "name" ]}
957 try :
958 after, verify_request_id = _get(url, token, transport = transport)
959 except HelperFailure as failure:
960 raise HelperFailure(
961 failure.code,
962 failure.message,
963 blocked_at = failure.blocked_at,
964 writes = [write, * failure.writes],
965 resources_remaining = [_target_identity(plan)],
966 request_id = failure.request_id,
967 status = failure.http_status,
968 partial = True ,
969 ) from failure
970 if after is not None :
971 raise HelperFailure(
972 "absence-unverified" ,
973 "The exact resource still exists after delete." ,
974 blocked_at = "verification" ,
975 writes = [write],
976 resources_remaining = [_target_identity(plan)],
977 request_id = verify_request_id,
978 partial = True ,
979 )
980 return _completed(
981 outcome,
982 fingerprint,
983 plan,
984 action = "deleted" ,
985 readback = None ,
986 request_ids = [initial_request_id, result.request_id, verify_request_id],
987 absence = True ,
988 )
989
990 source_request_id = None
991 if plan[ "resource_type" ] == "knowledge-base" :
992 verified_source = plan[ "verified_source" ]
993 source_url = _resource_url({
994 ** plan, "resource_type" : "knowledge-source" , "name" : verified_source[ "name" ],
995 })
996 source, source_request_id = _get(source_url, token, transport = transport)
997 if source is None or digest(_definition(source)) != verified_source[ "definition_digest" ]:
998 raise HelperFailure(
999 "source-drift" ,
1000 "The source is absent or its current definition differs from the approved source." ,
1001 blocked_at = "reconciliation" ,
1002 request_id = source_request_id,
1003 )
1004
1005 desired = plan[ "desired" ]
1006 if current is not None and _definition(desired) == _definition(current):
1007 if (
1008 plan.get( "action" ) == "reuse"
1009 and plan.get( "expected_etag" ) is not None
1010 and current.get( "@odata.etag" ) != plan[ "expected_etag" ]
1011 ):
1012 raise HelperFailure(
1013 "definition-drift" ,
1014 "Current ETag does not match the approved reuse plan." ,
1015 blocked_at = "reconciliation" ,
1016 )
1017 return _completed(
1018 outcome,
1019 fingerprint,
1020 plan,
1021 action = "reused" ,
1022 readback = current,
1023 request_ids = [initial_request_id, source_request_id],
1024 absence = False ,
1025 )
1026
1027 action = plan.get( "action" )
1028 headers = {
1029 "Content-Type" : "application/json" ,
1030 "Prefer" : "return=representation" ,
1031 }
1032 if current is None :
1033 if action != "create" :
1034 raise HelperFailure(
1035 "target-absent" ,
1036 "The approved update target does not exist." ,
1037 blocked_at = "reconciliation" ,
1038 )
1039 headers[ "If-None-Match" ] = "*"
1040 completed_action = "created"
1041 else :
1042 if action != "update" :
1043 raise HelperFailure(
1044 "definition-conflict" ,
1045 "An existing non-equivalent resource cannot be overwritten by this plan." ,
1046 blocked_at = "reconciliation" ,
1047 )
1048 expected_etag = plan.get( "expected_etag" )
1049 if not expected_etag or current.get( "@odata.etag" ) != expected_etag:
1050 raise HelperFailure(
1051 "definition-drift" ,
1052 "Current ETag does not match the approved update plan." ,
1053 blocked_at = "reconciliation" ,
1054 )
1055 headers[ "If-Match" ] = expected_etag
1056 completed_action = "updated"
1057
1058 request_desired = copy.deepcopy(desired)
1059 credential_environment = plan.get( "ai_services_api_key_environment" )
1060 if credential_environment is not None or acquisition is not None :
1061 secret = credential_provider() if acquisition is not None else os.environ.get(credential_environment)
1062 if not secret:
1063 raise HelperFailure(
1064 "credential-unavailable" ,
1065 "The approved AI Services credential environment variable is unset." ,
1066 blocked_at = "execution" ,
1067 )
1068 request_desired[ "fileParameters" ][ "ingestionParameters" ][ "aiServices" ][
1069 "apiKey"
1070 ] = secret
1071 try :
1072 result = transport(
1073 "PUT" ,
1074 url,
1075 token,
1076 body = canonical_bytes(request_desired),
1077 headers = headers,
1078 )
1079 except HelperFailure as failure:
1080 if not is_ambiguous_mutation_failure(failure):
1081 raise
1082 return _recover_ambiguous_put(
1083 outcome,
1084 fingerprint,
1085 plan,
1086 url,
1087 token,
1088 initial_request_id,
1089 completed_action,
1090 failure,
1091 transport = transport,
1092 source_request_id = source_request_id,
1093 )
1094 finally :
1095 if credential_environment is not None or acquisition is not None :
1096 request_desired[ "fileParameters" ][ "ingestionParameters" ][ "aiServices" ].pop( "apiKey" , None )
1097 secret = None
1098 if result.status not in { 200 , 201 }:
1099 if result.status in { 408 , 429 } or result.status >= 500 :
1100 return _recover_ambiguous_put(
1101 outcome,
1102 fingerprint,
1103 plan,
1104 url,
1105 token,
1106 initial_request_id,
1107 completed_action,
1108 HelperFailure(
1109 "mutation-outcome-ambiguous" ,
1110 f "Create or update returned ambiguous HTTP { result.status } ." ,
1111 blocked_at = "execution" ,
1112 request_id = result.request_id,
1113 status = result.status,
1114 partial = True ,
1115 retry_after = result.retry_after,
1116 recovery_deadline = result.recovery_deadline,
1117 ),
1118 transport = transport,
1119 source_request_id = source_request_id,
1120 )
1121 raise HelperFailure(
1122 "mutation-failed" ,
1123 f "Create or update returned unexpected HTTP { result.status } ." ,
1124 blocked_at = "execution" ,
1125 request_id = result.request_id,
1126 status = result.status,
1127 )
1128 write = { "action" : completed_action, "name" : plan[ "name" ]}
1129 recovery = ReadRecovery( deadline = result.recovery_deadline)
1130 try :
1131 if completed_action == "created" and cleanup_capture is not None :
1132 cleanup_receipts.search_ack(cleanup_capture, plan, result)
1133 # Outside the ambiguous-PUT handler: callback/IO failure cannot replay or undo this write.
1134 if completed_action == "created" and (on_created is not None or on_file_acknowledged is not None ):
1135 try :
1136 if on_created is not None :
1137 on_created(
1138 response = _creation_callback_response(result), url = url, body = canonical_bytes(request_desired),
1139 headers = {key: value for key, value in headers.items()
1140 if key in { "Content-Type" , "Prefer" , "If-None-Match" , "If-Match" }},
1141 )
1142 if on_file_acknowledged is not None :
1143 on_file_acknowledged({
1144 "status" : result.status, "request_id" : ReadRecovery.safe_id(result.request_id) if result.request_id else None ,
1145 "etag_evidence" : response_etags(_creation_callback_response(result)),
1146 })
1147 except OSError as failure:
1148 raise HelperFailure(
1149 "creation-receipt-persistence-failed" ,
1150 f "Acknowledged creation receipt persistence failed ( { type (failure). __name__ } ); private details withheld." ,
1151 blocked_at = "local-persistence" , request_id = result.request_id, status = result.status,
1152 ) from failure
1153 after, verify_request_id = _get(url, token, transport = transport, recovery = recovery)
1154 if completed_action == "created" and cleanup_capture is not None :
1155 cleanup_receipts.search_finish(cleanup_capture, plan, after, token, transport)
1156 except HelperFailure as failure:
1157 recovery.annotate(failure)
1158 if result.request_id:
1159 failure.warnings.append( "Acknowledged write request ID: " + recovery.safe_id(result.request_id))
1160 raise HelperFailure(
1161 failure.code,
1162 failure.message,
1163 blocked_at = failure.blocked_at,
1164 writes = [write, * failure.writes],
1165 resources_remaining = ([_target_identity(plan)] if completed_action == "created" else []) + failure.resources_remaining,
1166 resources_reused = ([_target_identity(plan)] if completed_action == "updated" else []) + failure.resources_reused,
1167 resources_unverified = failure.resources_unverified,
1168 warnings = failure.warnings,
1169 request_id = failure.request_id,
1170 status = failure.http_status,
1171 partial = True ,
1172 ) from failure
1173 if after is None or _definition(desired) != _definition(after):
1174 raise HelperFailure(
1175 "readback-mismatch" ,
1176 "Readback does not contain the approved definition." ,
1177 blocked_at = "verification" ,
1178 writes = [write],
1179 resources_remaining = [_target_identity(plan)] if completed_action == "created" else [],
1180 resources_reused = [_target_identity(plan)] if completed_action == "updated" else [],
1181 request_id = verify_request_id,
1182 partial = True ,
1183 warnings = [ * recovery.diagnostics(), * (
1184 [ "Acknowledged write request ID: " + recovery.safe_id(result.request_id)]
1185 if result.request_id else []
1186 )],
1187 )
1188 completed = _completed(
1189 outcome,
1190 fingerprint,
1191 plan,
1192 action = completed_action,
1193 readback = after,
1194 request_ids = [initial_request_id, source_request_id, result.request_id, * recovery.request_ids],
1195 absence = False ,
1196 )
1197 completed[ "warnings" ].extend(recovery.warnings)
1198 return completed
1199
1200
1201 def _target_identity (plan: dict[ str , Any]) -> dict[ str , Any]:
1202 return { "type" : plan[ "resource_type" ], "name" : plan[ "name" ]}
1203
1204
1205 def _recover_ambiguous_put (
1206 outcome: str ,
1207 fingerprint: str ,
1208 plan: dict[ str , Any],
1209 url: str ,
1210 token: str ,
1211 initial_request_id: str | None ,
1212 completed_action: str ,
1213 failure: HelperFailure,
1214 * ,
1215 transport: Transport,
1216 source_request_id: str | None = None ,
1217 ) -> dict[ str , Any]:
1218 identity = _target_identity(plan)
1219 recovery = ReadRecovery()
1220
1221 def readback ():
1222 recovery.delay(failure)
1223 return _get(url, token, transport = transport, recovery = recovery)
1224
1225 if completed_action == "created" :
1226 observed = []
1227 warnings = list (failure.warnings)
1228 read_ids = [value for value in (initial_request_id, source_request_id) if value]
1229 try :
1230 after, read_id = readback()
1231 except HelperFailure as readback_failure:
1232 if readback_failure.request_id:
1233 read_ids.append(readback_failure.request_id)
1234 warnings.append( f "Ambiguous-create readback also failed ( { readback_failure.code } ); original mutation error retained." )
1235 else :
1236 if read_id:
1237 read_ids.append(read_id)
1238 if after is not None :
1239 observed.append(identity)
1240 warnings.append( "Readback resources are observations only, not creation ownership or authorized reuse." )
1241 if read_ids:
1242 warnings.append( "Read-only request IDs: " + ", " .join(recovery.safe_id(value) for value in read_ids))
1243 warnings.extend(recovery.diagnostics())
1244 raise HelperFailure(
1245 failure.code, failure.message, blocked_at = failure.blocked_at,
1246 resources_reused = observed, request_id = failure.request_id, status = failure.http_status,
1247 partial = True , warnings = warnings,
1248 ) from failure
1249 message = "Same-identity readback did not prove the approved create or update."
1250 try :
1251 after, verify_request_id = readback()
1252 except HelperFailure as readback_failure:
1253 verify_request_id = readback_failure.request_id
1254 message = "Create or update outcome and same-identity readback are ambiguous."
1255 detail = (
1256 f " { message } Readback failure: { readback_failure.code } ; "
1257 f "HTTP { readback_failure.http_status } ; request ID { recovery.safe_id(verify_request_id) } ."
1258 )
1259 else :
1260 if after is not None and _definition(plan[ "desired" ]) == _definition(after):
1261 completed = _completed(
1262 outcome,
1263 fingerprint,
1264 plan,
1265 action = completed_action,
1266 readback = after,
1267 request_ids = [
1268 initial_request_id,
1269 source_request_id,
1270 failure.request_id,
1271 * recovery.request_ids,
1272 ],
1273 absence = False ,
1274 )
1275 completed[ "warnings" ].extend(recovery.warnings)
1276 return completed
1277 detail = f " { message } Readback request ID: { verify_request_id } ."
1278 if completed_action == "updated" :
1279 raise HelperFailure(
1280 failure.code,
1281 failure.message,
1282 blocked_at = failure.blocked_at,
1283 writes = failure.writes,
1284 resources_reused = [identity],
1285 request_id = failure.request_id,
1286 status = failure.http_status,
1287 partial = True ,
1288 warnings = [ * failure.warnings, detail, * recovery.diagnostics()],
1289 )
1290 raise HelperFailure(
1291 "mutation-outcome-ambiguous" ,
1292 message,
1293 blocked_at = "verification" ,
1294 resources_remaining = [identity],
1295 request_id = verify_request_id or failure.request_id,
1296 status = failure.http_status,
1297 partial = True ,
1298 )
1299
1300
1301 def _recover_ambiguous_delete (
1302 outcome: str ,
1303 fingerprint: str ,
1304 plan: dict[ str , Any],
1305 url: str ,
1306 token: str ,
1307 initial_request_id: str | None ,
1308 failure: HelperFailure,
1309 * ,
1310 transport: Transport,
1311 ) -> dict[ str , Any]:
1312 identity = _target_identity(plan)
1313 try :
1314 after, verify_request_id = _get(url, token, transport = transport)
1315 except HelperFailure as readback_failure:
1316 raise HelperFailure(
1317 "delete-outcome-ambiguous" ,
1318 "Delete outcome and same-identity readback are ambiguous." ,
1319 blocked_at = "verification" ,
1320 resources_remaining = [identity],
1321 request_id = readback_failure.request_id or failure.request_id,
1322 status = failure.http_status,
1323 partial = True ,
1324 ) from readback_failure
1325 if after is not None :
1326 raise HelperFailure(
1327 "delete-outcome-ambiguous" ,
1328 "Same-identity readback still found the resource after an ambiguous delete." ,
1329 blocked_at = "verification" ,
1330 resources_remaining = [identity],
1331 request_id = verify_request_id or failure.request_id,
1332 status = failure.http_status,
1333 partial = True ,
1334 )
1335 return _completed(
1336 outcome,
1337 fingerprint,
1338 plan,
1339 action = "deleted" ,
1340 readback = None ,
1341 request_ids = [
1342 initial_request_id,
1343 failure.request_id,
1344 verify_request_id,
1345 ],
1346 absence = True ,
1347 )
1348
1349
1350 def _completed (
1351 outcome: str ,
1352 fingerprint: str ,
1353 plan: dict[ str , Any],
1354 * ,
1355 action: str ,
1356 readback: dict[ str , Any] | None ,
1357 request_ids: list[ str | None ],
1358 absence: bool ,
1359 ) -> dict[ str , Any]:
1360 resource = {
1361 "type" : plan[ "resource_type" ],
1362 "name" : plan[ "name" ],
1363 "etag" : readback.get( "@odata.etag" ) if readback else None ,
1364 "definition_digest" : digest(_definition(readback)) if readback else None ,
1365 }
1366 resources = { "created" : [], "reused" : [], "updated" : [], "skipped" : []}
1367 if action in resources:
1368 resources[action].append(resource)
1369 elif action == "deleted" :
1370 resources[ "deleted" ] = [resource]
1371 return {
1372 "status" : "completed" ,
1373 "outcome" : outcome,
1374 "approved_plan" : { "fingerprint" : fingerprint, "confirmed" : True },
1375 "resources" : resources,
1376 "api_contracts" : [
1377 {
1378 "operation" : plan[ "operation" ],
1379 "version" : plan[ "api_version" ],
1380 "preview" : plan[ "api_version" ].endswith( "-preview" ),
1381 }
1382 ],
1383 "data_movement" : plan.get( "data_movement" , { "boundary" : None , "result" : "none" }),
1384 "auth" : { "mode" : "entra-user" , "principals" : []},
1385 "rbac" : plan.get( "rbac" , { "assignments" : []}),
1386 "network" : plan.get( "network" , { "posture" : "preserved" , "evidence" : None }),
1387 "verification" : {
1388 "readback" : resource,
1389 "absence" : absence,
1390 "request_ids" : [item for item in request_ids if item],
1391 "idempotency" : "exact readback is zero-write" ,
1392 },
1393 "warnings" : [],
1394 "ownership" : {
1395 "run_owned" : [resource] if action in { "created" , "deleted" } else [],
1396 "reused_not_owned" : [resource] if action in { "reused" , "updated" } else [],
1397 "owner" : plan.get( "owner" ),
1398 },
1399 "cleanup" : {
1400 "status" : (
1401 "completed"
1402 if plan[ "operation" ] == "delete" and absence
1403 else "not-requested"
1404 ),
1405 "separate_confirmation_required" : True ,
1406 },
1407 }
1408
1409
1410 def main (argv: list[ str ] | None = None ) -> int :
1411 try :
1412 from .private_artifacts import add_execution_output_argument, emit_plan_result, validate_execution_output_mode
1413 except ImportError :
1414 from private_artifacts import add_execution_output_argument, emit_plan_result, validate_execution_output_mode
1415 parser = argparse.ArgumentParser()
1416 modes = parser.add_mutually_exclusive_group( required = True )
1417 modes.add_argument( "--input" , type = Path)
1418 modes.add_argument( "--plan" , type = Path, help = "Build an unapproved KB plan using exact read-only discovery." )
1419 cleanup_receipts.add_argument(parser)
1420 add_execution_output_argument(parser)
1421 args = parser.parse_args(argv)
1422 fingerprint: str | None = None
1423 owner: Any = None
1424 outcome = "search-resource-reconciliation"
1425 capture = None
1426 try :
1427 validate_execution_output_mode(args)
1428 if args.plan and args.cleanup_receipt_dir:
1429 raise HelperFailure( "input-schema-invalid" , "Creation capture is only available with approved --input." , blocked_at = "confirmation" )
1430 if args.plan:
1431 try :
1432 from ._bootstrap_io import read_json
1433 except ImportError :
1434 from _bootstrap_io import read_json
1435 outcome = "create-knowledge-base"
1436 request = read_json(args.plan)
1437 owner = request.get( "owner" ) if isinstance (request, dict ) else None
1438 result = plan_knowledge_base(request)
1439 emit_plan_result(result, args.execution_output, preserve_unapproved_input = True )
1440 return 0
1441 document, plan, fingerprint = load_approved_input(args.input)
1442 document[ "_computed_fingerprint" ] = fingerprint
1443 owner = plan.get( "owner" )
1444 outcome = str (plan.get( "outcome" ) or outcome)
1445 capture = cleanup_receipts.Capture(args.cleanup_receipt_dir, document) if args.cleanup_receipt_dir else None
1446 result = execute(document, ** ({ "cleanup_capture" : capture} if capture else {}))
1447 except HelperFailure as failure:
1448 result = blocked_result(
1449 failure,
1450 outcome = outcome,
1451 fingerprint = fingerprint,
1452 owner = owner,
1453 )
1454 if capture is not None :
1455 result[ "cleanup_receipts" ] = capture.summaries
1456 emit_result(result)
1457 return 3 if result[ "status" ] == "partial" else 2
1458 if capture is not None :
1459 result[ "cleanup_receipts" ] = capture.summaries
1460 emit_result(result)
1461 return 0
1462
1463
1464 if __name__ == "__main__" :
1465 sys.exit(main())