Setting the file. One moment.
File Upload · Foundry Iq · microsoft/azure-skills · Skills Docs
ContentsBack to the top of the page File Cu Canary
— line 220
This file
Number 10.50
Position 50 of 77
Type Python
Size 55 KB
Lines 853 helpers/ file_upload.py
Python · 853 lines · 55 KB
pathlib
import
Path
13
14 try :
15 from . import _bootstrap_io as private_io, cu_ingestion_auth, file_ingest, search_reconcile
16 from ._common import (HelperFailure, ReadRecovery, RetryAfter, RetryAfterTiming, SEARCH_AUDIENCE ,
17 retry_after_timing, retry_after_not_before, valid_utc_timestamp, azure_cli_token, blocked_result,
18 digest, emit_result, http_request, load_approved_input, reject_secrets, require_allowed_fields)
19 from ._progress import Progress, add_progress_argument, reporting
20 except ImportError :
21 import _bootstrap_io as private_io, cu_ingestion_auth, file_ingest, search_reconcile
22 from _common import (HelperFailure, ReadRecovery, RetryAfter, RetryAfterTiming, SEARCH_AUDIENCE ,
23 retry_after_timing, retry_after_not_before, valid_utc_timestamp, azure_cli_token, blocked_result,
24 digest, emit_result, http_request, load_approved_input, reject_secrets, require_allowed_fields)
25 from _progress import Progress, add_progress_argument, reporting
26
27
28 RETAIN = " Retain the source and evidence; no creation replay, reset, deletion or new roles."
29 UPLOAD_TIMEOUT = 180
30 MAX_BACKOFF_RECORDS = 200
31
32
33 def _rejected (status):
34 return type (status) is int and 400 <= status < 500 and status not in ( 408 , 409 )
35
36
37 def fail (code, message):
38 return HelperFailure(code, message + RETAIN , blocked_at = "file-upload-resume" )
39
40
41 class Session :
42 def __init__ (self, directory, document = None , * , context_provider = cu_ingestion_auth.account_context):
43 self .directory = private_io.validate_private_artifact_directory( str (directory))
44 self .context_provider = context_provider
45 self .records = {}
46 self .network_failure = None
47 self .active_attempt = None
48 if document is not None :
49 self ._validate_document(document)
50 if any ( self .directory.iterdir()):
51 raise fail( "file-upload-journal-exists" , "Use a new empty private receipt directory for original creation." )
52 context = context_provider()
53 cu_ingestion_auth.validate_context(context)
54 self .seed = { "document" : {key: copy.deepcopy(document[key]) for key in ( "schema_version" , "plan" , "approval" )},
55 "context" : context, "backoff_version" : "1.0" }
56 self ._write( "run.json" , self .seed)
57 else :
58 self ._load()
59 self .seed = self .records[ "run.json" ]
60 self ._validate_document( self .seed[ "document" ])
61 cu_ingestion_auth.validate_context( self .seed[ "context" ])
62 self .document = self .seed[ "document" ]
63 self .plan = self .document[ "plan" ]
64 self .ingestion = self .plan[ "ingestion" ]
65
66 @ staticmethod
67 def _validate_document (document):
68 try :
69 from . import file_source
70 except ImportError :
71 import file_source
72 if ( not isinstance (document, dict ) or document.get( "schema_version" ) != "1.0"
73 or not isinstance (document.get( "plan" ), dict )
74 or "_computed_fingerprint" in document and document[ "_computed_fingerprint" ] != digest(document.get( "plan" ))
75 or not isinstance (document.get( "approval" ), dict ) or document[ "approval" ].get( "confirmed" ) is not True
76 or document.get( "approval" ) != { "confirmed" : True , "fingerprint" : digest(document.get( "plan" ))}):
77 raise fail( "file-upload-approval-missing" , "Original unchanged File creation approval is required." )
78 require_allowed_fields(document, { "schema_version" , "plan" , "approval" , "_computed_fingerprint" }, label = "original File envelope" )
79 file_source._validate_plan(document[ "plan" ])
80 if document[ "plan" ][ "source" ][ "action" ] != "create" :
81 raise fail( "file-upload-provenance-missing" , "This continuation is only for original acknowledged creation, not generic reuse." )
82
83 def _write (self, name, payload):
84 value = { "schema_version" : "1.0" , "kind" : "file-upload-journal" , "payload" : payload, "integrity" : digest(payload)}
85 reject_secrets(value)
86 try :
87 private_io.atomic_private_file( self .directory, name, value, max_bytes = private_io. MAX_BYTES )
88 except HelperFailure as error:
89 raise HelperFailure( "file-upload-receipt-failed" , "Private upload evidence persistence failed." + RETAIN ,
90 blocked_at = "local-persistence" , warnings = error.warnings) from error
91 self .records[name] = copy.deepcopy(payload)
92
93 def _load (self):
94 # Reuse the existing bounded, private, no-link evidence reader.
95 try :
96 from .blob_recheck import read_private
97 except ImportError :
98 from blob_recheck import read_private
99 try :
100 paths = list ( self .directory.iterdir())
101 except OSError as error:
102 raise fail( "file-upload-evidence-unreadable" , "Original upload evidence is inaccessible." ) from error
103 if len (paths) > 803 + MAX_BACKOFF_RECORDS or not { "run.json" , "source-ack.json" }.issubset({p.name for p in paths}):
104 raise fail( "file-upload-provenance-missing" , "Necessary original ACK and journal were not retained." )
105 for path in paths:
106 if (path.name not in { "run.json" , "source-ack.json" , "pending-request.json" }
107 and not re.fullmatch( r " (?:[ 0-9 ] {4} - ( retry- ) ? ( attempt | result ) | backoff- [ 0-9 ] {4} ) \. json" , path.name)):
108 raise fail( "file-upload-evidence-invalid" , "Unexpected or incomplete journal entry; do not discard it to resume." )
109 value = read_private(path)
110 if ( set (value) != { "schema_version" , "kind" , "payload" , "integrity" }
111 or value[ "schema_version" ] != "1.0" or value[ "kind" ] != "file-upload-journal"
112 or value[ "integrity" ] != digest(value[ "payload" ])):
113 raise fail( "file-upload-evidence-invalid" , "Original private upload evidence changed." )
114 self .records[path.name] = value[ "payload" ]
115 seed = self .records[ "run.json" ]
116 if ( not isinstance (seed, dict ) or set (seed) not in ({ "document" , "context" }, { "document" , "context" , "backoff_version" })
117 or "backoff_version" in seed and seed[ "backoff_version" ] != "1.0" ):
118 raise fail( "file-upload-evidence-invalid" , "Original approval/context evidence is incomplete." )
119
120 def acknowledge (self, metadata):
121 self ._write( "source-ack.json" , { "original_plan_digest" : digest( self .plan),
122 "source_url" : search_reconcile.resource_url( self .plan[ "source" ]),
123 "acknowledgement" : metadata})
124
125 def require_ack (self):
126 value = self .records.get( "source-ack.json" )
127 if ( not isinstance (value, dict ) or set (value) != { "original_plan_digest" , "source_url" , "acknowledgement" }
128 or value[ "original_plan_digest" ] != digest( self .plan)
129 or value[ "source_url" ] != search_reconcile.resource_url( self .plan[ "source" ])):
130 raise fail( "file-upload-provenance-missing" , "An original acknowledged conditional source creation must be retained." )
131 ack = value[ "acknowledgement" ]
132 if ( not isinstance (ack, dict ) or set (ack) != { "status" , "request_id" , "etag_evidence" }
133 or ack[ "status" ] not in ( 200 , 201 ) or not ack[ "request_id" ]
134 or not isinstance (ack[ "request_id" ], str ) or not re.fullmatch( r " [ A-Za-z0-9 ][ A-Za-z0-9._:- ] {0,127} " , ack[ "request_id" ])
135 or not isinstance (ack[ "etag_evidence" ], dict ) or set (ack[ "etag_evidence" ]) != { "body" , "headers" }
136 or not isinstance (ack[ "etag_evidence" ][ "headers" ], list )):
137 raise fail( "file-upload-provenance-missing" , "Original successful creation ACK/request ID is required, not GET ownership." )
138 etag = search_reconcile.resolve_etag(ack[ "etag_evidence" ], ack[ "request_id" ])
139 if not etag:
140 raise fail( "file-upload-provenance-missing" , "The original creation ACK lacks version proof; GET cannot supply it." )
141 return etag
142
143 def states (self):
144 states = []
145 events = self .backoff_events()
146 allowed = { "run.json" , "source-ack.json" , "pending-request.json" } | set (events)
147 pending = self .records.get( "pending-request.json" )
148 if "pending-request.json" in self .records and (
149 not isinstance (pending, dict ) or set (pending) != { "version" , "plan_digest" , "nonce" , "method" , "url_digest" }
150 or pending[ "version" ] != "1.0" or pending[ "plan_digest" ] != digest( self .plan)
151 or not isinstance (pending[ "nonce" ], str ) or not re.fullmatch( r " [ 0-9a-f ] {32} " , pending[ "nonce" ])
152 or pending[ "method" ] not in ( "GET" , "POST" , "PUT" , "HEAD" )
153 or not isinstance (pending[ "url_digest" ], str ) or not re.fullmatch( r "sha256: [ 0-9a-f ] {64} " , pending[ "url_digest" ])
154 ):
155 raise fail( "file-upload-evidence-invalid" , "Pending request evidence is invalid." )
156 for index, record in enumerate ( self .ingestion[ "files" ]):
157 state = { "upload" : "not_attempted" , "request_ids" : [], "ingested" : False }
158 for retry in ( False , True ):
159 prefix = f " { index :04} -" + ( "retry-" if retry else "" )
160 attempt_name, result_name = prefix + "attempt.json" , prefix + "result.json"
161 allowed.update((attempt_name, result_name))
162 attempt, result = self .records.get(attempt_name), self .records.get(result_name)
163 expected = self ._attempt_record(index, retry = retry) if attempt is not None else None
164 if attempt is not None and attempt != expected or result is not None and attempt is None :
165 raise fail( "file-upload-evidence-invalid" , "Upload attempt evidence does not bind the original source/corpus." )
166 if attempt is not None :
167 state.update( upload = "unverified" , ingested = False )
168 state.pop( "file_proof" , None )
169 state.pop( "ack_file_id" , None )
170 if result is not None :
171 if ( not isinstance (result, dict ) or not { "upload" , "status" , "request_id" } <= set (result)
172 or set (result) - { "upload" , "status" , "request_id" , "file_proof" , "file_id" , "throttle_event" }
173 or not isinstance (result[ "upload" ], str )
174 or result[ "upload" ] not in { "accepted" , "rejected" , "unverified" }
175 or result[ "upload" ] == "accepted" and result[ "status" ] not in ( 200 , 201 )
176 or result[ "upload" ] == "rejected" and not _rejected(result[ "status" ])
177 or result[ "status" ] is not None and ( type (result[ "status" ]) is not int or not 100 <= result[ "status" ] <= 599 )
178 or result[ "request_id" ] is not None and ReadRecovery.safe_id(result[ "request_id" ]) != result[ "request_id" ]):
179 raise fail( "file-upload-evidence-invalid" , "Upload result is not an original supported ACK or failure." )
180 linked = [event for event in events.values() if event[ "attempt" ] == attempt_name]
181 if "throttle_event" in result:
182 if (result[ "status" ] != 429 or len (linked) != 1
183 or result[ "throttle_event" ] != digest(linked[ 0 ])
184 or result[ "request_id" ] != linked[ 0 ][ "request_id" ]):
185 raise fail( "file-upload-evidence-invalid" , "Upload 429 timing does not bind its original attempt/result." )
186 elif result[ "status" ] == 429 and self .seed.get( "backoff_version" ) is not None :
187 raise fail( "file-upload-evidence-invalid" , "Required original upload 429 timing evidence is missing." )
188 if "file_id" in result:
189 if result[ "upload" ] != "accepted" or not _file_id(result[ "file_id" ]):
190 raise fail( "file-upload-evidence-invalid" , "Original upload ACK file identity is invalid." )
191 state[ "ack_file_id" ] = result[ "file_id" ]
192 proof = result.get( "file_proof" )
193 if proof is not None :
194 if (result[ "upload" ] != "accepted" or not isinstance (proof, dict )
195 or set (proof) != { "file_id" , "record_digest" } or not _file_id(proof[ "file_id" ])
196 or proof[ "record_digest" ] != digest(record)
197 or result.get( "file_id" , proof[ "file_id" ]) != proof[ "file_id" ]):
198 raise fail( "file-upload-evidence-invalid" , "Persisted File completion proof does not match the approved record." )
199 state.update( file_proof = proof, ingested = True )
200 state.update( upload = result[ "upload" ], status = result[ "status" ])
201 if result[ "request_id" ]:
202 state[ "request_ids" ].append(result[ "request_id" ])
203 states.append(state)
204 if set ( self .records) - allowed:
205 raise fail( "file-upload-evidence-invalid" , "Upload journal contains records outside the approved corpus." )
206 return states
207
208 def _attempt_record (self, index, * , retry = False ):
209 record = {
210 "original_plan_digest" : digest( self .plan), "record_digest" : digest( self .ingestion[ "files" ][index]),
211 "creation_ack_digest" : digest( self .records[ "source-ack.json" ]),
212 }
213 if retry:
214 rejection = self .records.get( f " { index :04} -result.json" )
215 if not isinstance (rejection, dict ) or rejection.get( "status" ) != 429 or rejection.get( "upload" ) != "rejected" :
216 raise fail( "file-upload-replay-forbidden" , "A retry attempt requires the same operation's retained File upload 429." )
217 record[ "rejected_result_digest" ] = digest(rejection)
218 return record
219
220 def attempt (self, index, * , retry = False ):
221 prefix = f " { index :04} -" + ( "retry-" if retry else "" )
222 self ._write(prefix + "attempt.json" , self ._attempt_record(index, retry = retry))
223 self .active_attempt = prefix + "attempt.json"
224
225 def result (self, index, upload, status, request_id, * , retry = False , proof = None , file_id = None ):
226 prefix = f " { index :04} -" + ( "retry-" if retry else "" )
227 event = {}
228 if status == 429 :
229 linked = [value for value in self .backoff_events().values() if value[ "attempt" ] == prefix + "attempt.json" ]
230 if len (linked) != 1 :
231 raise fail( "file-upload-evidence-invalid" , "Original upload 429 timing was not retained." )
232 event[ "throttle_event" ] = digest(linked[ 0 ])
233 self ._write(prefix + "result.json" , { "upload" : upload, "status" : status,
234 "request_id" : ReadRecovery.safe_id(request_id) if request_id else None ,
235 ** ({ "file_id" : file_id} if file_id is not None else {}),
236 ** ({ "file_proof" : proof} if proof is not None else {}), ** event})
237 self .active_attempt = None
238
239 def backoff_events (self):
240 events = {key: self .records[key] for key in sorted ( self .records) if key.startswith( "backoff-" )}
241 if len (events) > MAX_BACKOFF_RECORDS :
242 raise fail( "file-upload-evidence-invalid" , "Too many retained backoff records." )
243 plan_digest = digest( self .plan)
244 previous = None
245 for index, (name, event) in enumerate (events.items()):
246 if (name != f "backoff- { index :04} .json" or not isinstance (event, dict )
247 or set (event) != { "version" , "plan_digest" , "previous" , "method" , "url_digest" ,
248 "attempt" , "attempt_digest" , "request_id" , "metadata" , "timing" , "origin" }
249 or event[ "version" ] != "1.0" or event[ "plan_digest" ] != plan_digest
250 or event[ "previous" ] != previous or event[ "method" ] not in ( "GET" , "POST" , "PUT" , "HEAD" )
251 or not isinstance (event[ "url_digest" ], str ) or not re.fullmatch( r "sha256: [ 0-9a-f ] {64} " , event[ "url_digest" ])
252 or event[ "request_id" ] is not None and ReadRecovery.safe_id(event[ "request_id" ]) != event[ "request_id" ]
253 or event[ "origin" ] not in ( "response-headers" , "native-http-error" , "typed-metadata" , "unavailable" )):
254 raise fail( "file-upload-evidence-invalid" , "Original backoff provenance is invalid or incomplete." )
255 if event[ "attempt" ] is not None :
256 attempt = self .records.get(event[ "attempt" ]) if isinstance (event[ "attempt" ], str ) else None
257 if (attempt is None or not re.fullmatch( r " [ 0-9 ] {4} - ( retry- ) ? attempt \. json" , event[ "attempt" ])
258 or event[ "attempt_digest" ] != digest(attempt) or event[ "method" ] != "POST"
259 or event[ "url_digest" ] != digest(file_ingest._list_url( self .ingestion))):
260 raise fail( "file-upload-evidence-invalid" , "Backoff does not bind the original upload attempt." )
261 elif event[ "attempt_digest" ] is not None :
262 raise fail( "file-upload-evidence-invalid" , "Backoff attempt provenance is inconsistent." )
263 metadata, timing = event[ "metadata" ], event[ "timing" ]
264 if ( not isinstance (metadata, dict ) or set (metadata) != { "kind" , "value" }
265 or metadata[ "kind" ] not in ( "seconds" , "date" , "date-rfc850" , "missing" , "invalid" , "overlong" )
266 or not valid_utc_timestamp(metadata[ "value" ])
267 or metadata[ "kind" ] == "seconds" and ( type (metadata[ "value" ]) is not int or not 0 <= metadata[ "value" ] <= 30 )
268 or metadata[ "kind" ] in ( "missing" , "invalid" , "overlong" ) and metadata[ "value" ] != 0
269 or not isinstance (timing, dict ) or set (timing) != { "received_at_utc" , "not_before_utc" , "server_delay_seconds" }):
270 raise fail( "file-upload-evidence-invalid" , "Retained Retry-After metadata is invalid." )
271 received, deadline = timing[ "received_at_utc" ], timing[ "not_before_utc" ]
272 seconds = timing[ "server_delay_seconds" ]
273 if (received is not None and not valid_utc_timestamp(received)
274 or deadline is not None and (received is None or not valid_utc_timestamp(deadline))
275 or seconds is not None and ( type (seconds) is not int or not 0 <= seconds < 315537897600
276 or metadata[ "kind" ] not in ( "seconds" , "overlong" ))
277 or metadata[ "kind" ] == "seconds" and seconds not in ( None , metadata[ "value" ])):
278 raise fail( "file-upload-evidence-invalid" , "Retained backoff UTC timing is invalid." )
279 if deadline is not None :
280 if (event[ "origin" ] == "unavailable"
281 or event[ "origin" ] == "typed-metadata" and metadata[ "kind" ] not in ( "seconds" , "date" , "date-rfc850" )
282 or metadata[ "kind" ] == "overlong" and (seconds is None or seconds <= 30 or deadline != received + seconds)
283 or metadata[ "kind" ] != "overlong"
284 and deadline != retry_after_not_before(RetryAfter( ** metadata), received)):
285 raise fail( "file-upload-evidence-invalid" , "Retained backoff deadline contradicts original metadata." )
286 previous = digest(event)
287 return events
288
289 def backoff_status (self):
290 events = self .backoff_events()
291 attempts = {name: value for name, value in self .records.items()
292 if re.fullmatch( r " [ 0-9 ] {4} - ( retry- ) ? ( attempt | result ) \. json" , name)}
293 if any ( not isinstance (value, dict ) for value in attempts.values()):
294 raise fail( "file-upload-evidence-invalid" , "Attempt/result evidence is malformed." )
295 # Legacy 429 or an interrupted receipt cannot prove what restriction was received.
296 unknown = "pending-request.json" in self .records or any (
297 name.endswith( "result.json" ) and value.get( "status" ) == 429 and "throttle_event" not in value
298 or name.endswith( "attempt.json" ) and name.replace( "attempt.json" , "result.json" ) not in self .records
299 and name != self .active_attempt
300 for name, value in attempts.items()
301 )
302 now = time.time()
303 if unknown or any (event[ "timing" ][ "not_before_utc" ] is None for event in events.values()):
304 return { "status" : "unresolved" , "reason" : "Original response timing is unavailable; no deadline is inferred." }
305 if not events:
306 return { "status" : "clear" }
307 if ( not valid_utc_timestamp(now)
308 or any (now < event[ "timing" ][ "received_at_utc" ] for event in events.values())):
309 return { "status" : "unresolved" , "reason" : "UTC clock is invalid or precedes retained response receipt." }
310 deadline = max (event[ "timing" ][ "not_before_utc" ] for event in events.values())
311 return { "status" : "waiting" if now < deadline else "elapsed" , "not_before_utc" : deadline,
312 "seconds_remaining" : max ( 0 , deadline - now)}
313
314 def require_backoff (self):
315 status = self .backoff_status()
316 if status[ "status" ] in ( "waiting" , "unresolved" ):
317 error = fail( "file-upload-backoff-" + status[ "status" ],
318 "Server backoff still applies; fresh approval does not waive it. "
319 + ( "Wait until the retained UTC not-before time." if status[ "status" ] == "waiting"
320 else "Original timing evidence or scoped operator recovery is required." ))
321 error.partial = "source-ack.json" in self .records
322 error.file_batch = _summary( self .ingestion, self .states(), self )
323 raise error
324
325 def transport (self, raw):
326 if getattr (raw, "_file_journal_owner" , None ) is self :
327 return raw
328
329 def invoke (method, url, token, ** kwargs):
330 if self .network_failure is not None :
331 raise self .network_failure
332 self .require_backoff()
333 if len ( self .backoff_events()) >= MAX_BACKOFF_RECORDS :
334 raise fail( "file-upload-backoff-limit" , "Backoff journal limit reached; no more requests." )
335 pending = { "version" : "1.0" , "plan_digest" : digest( self .plan), "nonce" : uuid.uuid4().hex,
336 "method" : method, "url_digest" : digest(url)}
337 self ._write( "pending-request.json" , pending)
338 try :
339 response = raw(method, url, token, ** kwargs)
340 except HelperFailure as error:
341 if error.http_status == 429 :
342 self .record_backoff(method, url, error, None )
343 self .finish_request(error)
344 raise
345 if response.status == 429 :
346 error = HelperFailure( "azure-http-error" , "Azure request failed with HTTP 429." ,
347 blocked_at = "execution" , status = 429 , request_id = response.request_id,
348 retry_after = response.retry_after, recovery_deadline = response.recovery_deadline)
349 self .record_backoff(method, url, error, response.headers)
350 try :
351 self .finish_request()
352 except HelperFailure as persistence:
353 if method in ( "PUT" , "POST" ) and response.status in ( 200 , 201 ):
354 return replace(response, ack_failure = persistence)
355 raise
356 return response
357
358 invoke._file_journal_owner = self
359 return invoke
360
361 def finish_request (self, original = None ):
362 try :
363 path = self .directory / "pending-request.json"
364 expected = self .records[ "pending-request.json" ]
365 retained = private_io.read_json(path)
366 if not isinstance (retained, dict ) or retained.get( "payload" ) != expected or retained.get( "integrity" ) != digest(expected):
367 raise fail( "file-upload-evidence-invalid" , "Pending request receipt changed." )
368 path.unlink()
369 self .records.pop( "pending-request.json" )
370 except ( OSError , HelperFailure) as cleanup:
371 error = original or HelperFailure( "file-upload-receipt-failed" , "Pending request evidence could not be finalized." ,
372 blocked_at = "local-persistence" )
373 error.blocked_at = "local-persistence"
374 error.warnings.append( "Pending request evidence remains unresolved; no further requests." )
375 self .network_failure = error
376 raise error from cleanup
377
378 def record_backoff (self, method, url, error, headers):
379 metadata = error.retry_after
380 origin = "response-headers" if headers is not None else "native-http-error"
381 timing = retry_after_timing(headers) if headers is not None else error.retry_after_timing
382 if timing is None :
383 received = time.time()
384 received = received if valid_utc_timestamp(received) else None
385 known = isinstance (metadata, RetryAfter) and metadata.kind in ( "seconds" , "date" , "date-rfc850" )
386 timing = RetryAfterTiming(received, retry_after_not_before(metadata, received) if known else None ,
387 metadata.value if known and metadata.kind == "seconds" else None )
388 origin = "typed-metadata" if known else "unavailable"
389 if not isinstance (metadata, RetryAfter):
390 metadata, origin = RetryAfter( "missing" ), "unavailable"
391 timing = RetryAfterTiming(timing.received_at_utc, None )
392 events = self .backoff_events()
393 attempt = self .active_attempt if method == "POST" and url == file_ingest._list_url( self .ingestion) else None
394 event = { "version" : "1.0" , "plan_digest" : digest( self .plan),
395 "previous" : digest( next ( reversed (events.values()))) if events else None ,
396 "method" : method, "url_digest" : digest(url), "attempt" : attempt,
397 "attempt_digest" : digest( self .records[attempt]) if attempt else None ,
398 "request_id" : ReadRecovery.safe_id(error.request_id) if error.request_id else None ,
399 "metadata" : metadata._asdict(), "timing" : timing._asdict(), "origin" : origin}
400 try :
401 self ._write( f "backoff- { len (events) :04} .json" , event)
402 except HelperFailure as persistence:
403 error.blocked_at = "local-persistence"
404 error.warnings.extend(persistence.warnings)
405 error.warnings.append( "Original HTTP 429 retained; backoff receipt persistence failed. No further requests." )
406 self .network_failure = error
407 raise error from persistence
408
409
410 def _file_id (value):
411 return isinstance (value, str ) and re.fullmatch( r " [ A-Za-z0-9 ][ A-Za-z0-9._:- ] {0,127} " , value) is not None
412
413
414 def _proof (plan, record, item):
415 if isinstance (item, dict ) and _file_id(item.get( "fileId" )) and file_ingest._matches(item, plan, record):
416 return { "file_id" : item[ "fileId" ], "record_digest" : digest(record)}
417 return None
418
419
420 def _observe (plan, states, files):
421 names = {record[ "path" ] for record in plan[ "files" ]}
422 if any (item.get( "fileName" ) not in names for item in files):
423 raise fail( "server-inventory-conflict" , "File inventory contains an unapproved identity." )
424 observed = {}
425 for index, record in enumerate (plan[ "files" ]):
426 matches = [item for item in files if item.get( "fileName" ) == record[ "path" ]]
427 if len (matches) > 1 or matches and not file_ingest._matches({ ** matches[ 0 ], "errorMessage" : None }, plan, record):
428 raise fail( "file-record-conflict" , "File inventory has duplicate or conflicting approved markers." )
429 state = states[index]
430 state[ "verification" ] = { "accepted" : "pending" , "rejected" : "failed" }.get(state[ "upload" ], state[ "upload" ])
431 if matches:
432 item = matches[ 0 ]
433 proof = _proof(plan, record, item)
434 if (proof and state.get( "file_proof" ) and proof != state[ "file_proof" ]
435 or state.get( "ack_file_id" ) and item.get( "fileId" ) is not None
436 and item[ "fileId" ] != state[ "ack_file_id" ]):
437 raise fail( "file-record-conflict" , "Current File identity differs from the persisted upload ACK." )
438 state[ "verification" ] = ( "failed" if state[ "upload" ] == "rejected" and state.get( "status" ) != 429
439 or item.get( "errorMessage" ) is not None else
440 "confirmed" if proof else "unverified" )
441 state[ "ingested" ] = state[ "verification" ] == "confirmed"
442 observed[index] = item
443 return observed
444
445
446 def _summary (plan, states, session):
447 counts = {name: 0 for name in ( "accepted" , "confirmed" , "failed" , "pending" , "unverified" , "not_attempted" )}
448 files = []
449 for record, state in zip (plan[ "files" ], states):
450 counts[ "accepted" ] += state[ "upload" ] == "accepted"
451 verification = state.get( "verification" , { "accepted" : "pending" , "rejected" : "failed" }.get(state[ "upload" ], state[ "upload" ]))
452 counts[verification] += 1
453 files.append({ "fileName" : record[ "path" ], "sha256" : record[ "sha256" ], "upload" : state[ "upload" ],
454 "verification" : verification, "http_status" : state.get( "status" ),
455 "request_ids" : state[ "request_ids" ][: 3 ]})
456 ingested = sum ( bool (state.get( "ingested" )) for state in states)
457 return { "counts" : counts, "files" : files, "ingested" : ingested,
458 "readiness" : "ingested" if ingested == len (states) and counts[ "confirmed" ] == len (states) else "partial" ,
459 "retrieval" : "unverified" ,
460 ** ({ "backoff" : session.backoff_status()} if session else {}),
461 ** ({ "upload_retry" : {
462 "status" : "blocked" , "code" : "file-upload-retry-safety-unproven" ,
463 "reason" : "Unknown upload outcomes cannot be replayed; only a received File upload 429 permits one bounded in-operation retry." ,
464 }} if counts[ "unverified" ] else {}),
465 "resume" : ( "Plan newly approved upload-only continuation for never-attempted files; uncertain attempts are never replayed."
466 if session else "Original durable ACK/attempt evidence was not retained; uploads cannot safely resume from GET alone." ),
467 "subset_handoff" : "Report the proved ingested subset and remaining outcomes; subset isolation and KB retrieval are separate. KB changes need separate approval." }
468
469
470 def _progress (progress, plan, states):
471 summary = _summary(plan, states, None )
472 counts = summary[ "counts" ]
473 progress.update( "file-upload" , uploads_acknowledged = counts[ "accepted" ], files_verified = counts[ "confirmed" ],
474 files_ingested = summary[ "ingested" ],
475 files_reused = sum (state[ "upload" ] == "not_attempted" and state.get( "verification" ) == "confirmed" for state in states),
476 files_failed = counts[ "failed" ], files_unverified = counts[ "unverified" ],
477 files_not_attempted = counts[ "not_attempted" ], files_pending = counts[ "pending" ])
478
479
480 def run_batch (document, * , token_provider, transport, progress, allow_new_uploads = False , session = None ,
481 eligible = None , source_check = None , allow_upload_retry = True ):
482 plan, fingerprint = document[ "plan" ], document[ "_computed_fingerprint" ]
483 root, records = file_ingest._validate_plan(plan)
484 states = session.states() if session else [{ "upload" : "not_attempted" , "request_ids" : []} for _ in records]
485 if session:
486 session.require_backoff()
487 transport = session.transport(transport)
488 created, reused, request_ids, warnings = [], [], [], []
489 primary = None
490 stopped = False
491 retry_used = False
492 readback_failed = False
493 token = token_provider( SEARCH_AUDIENCE )
494 url = file_ingest._list_url(plan)
495 recovery = ReadRecovery( on_wait = progress.waiting)
496 progress.update( "file-inventory" )
497 try :
498 before, ids = file_ingest._list_files(url, token, transport = transport, recovery = recovery)
499 request_ids.extend(ids)
500 if session is None and file_ingest.inventory_digest(before) != plan[ "expected_server_inventory_digest" ]:
501 raise fail( "server-inventory-drift" , "Approved initial file inventory changed." )
502 observed = _observe(plan, states, before)
503 if not allow_new_uploads and len (observed) != len (records):
504 raise fail( "reused-source-upload-forbidden" , "Generic reuse cannot authorize new uploads." )
505 except HelperFailure as failure:
506 failure.file_batch = _summary(plan, states, session)
507 failure.file_batch[ "request_ids" ] = recovery.request_ids
508 raise
509 warnings.extend(recovery.warnings)
510 eligible = set ( range ( len (records))) if eligible is None else set (eligible)
511 after = None
512 for index, record in enumerate (records):
513 _progress(progress, plan, states)
514 state = states[index]
515 if index in observed:
516 reused.append({ "fileId" : observed[index].get( "fileId" ), "fileName" : record[ "path" ], "sha256" : record[ "sha256" ]})
517 continue
518 if index not in eligible or state[ "upload" ] != "not_attempted" :
519 continue
520 try :
521 path = file_ingest._resolve_file(root, record)
522 content = path.read_bytes()
523 if ( len (content) != record[ "size" ] or "sha256:" + hashlib.sha256(content).hexdigest() != record[ "sha256" ]
524 or path.stat().st_mtime_ns != record[ "mtime_ns" ]):
525 raise fail( "inventory-drift" , "Approved corpus changed before upload." )
526 except OSError :
527 primary = primary or fail( "inventory-unreadable" , "Approved corpus became inaccessible." )
528 warnings.append( "Batch stopped: approved corpus became inaccessible." )
529 stopped = True
530 break
531 except HelperFailure as failure:
532 primary = primary or failure
533 warnings.append( f "Batch stopped ( { failure.code } ); no further uploads were attempted." )
534 stopped = True
535 break
536 body, boundary = file_ingest._multipart(plan, record, content, fingerprint)
537 for attempt in range ( 2 ):
538 try :
539 if attempt:
540 file_ingest._resolve_file(root, record)
541 if session:
542 session.attempt(index, retry = bool (attempt))
543 except HelperFailure as failure:
544 primary = primary or failure
545 warnings.append( f "Upload attempt evidence/drift check failed ( { failure.code } ); no further requests." )
546 stopped = True
547 break
548 state.update( upload = "unverified" , verification = "unverified" , status = None )
549 _progress(progress, plan, states)
550 response = None
551 transport_failed = False
552 try :
553 with progress.processing_file(index + 1 , len (records), attempt + 1 ):
554 response = transport( "POST" , url, token, body = body,
555 headers = { "Content-Type" : f "multipart/form-data; boundary= { boundary } " },
556 follow_redirects = False , timeout = UPLOAD_TIMEOUT ,
557 response_deadline = time.monotonic() + UPLOAD_TIMEOUT ,
558 max_response_bytes = 1024 * 1024 )
559 except HelperFailure as caught:
560 transport_failed = True
561 failure = caught
562 failure_status, failure_id = failure.http_status, failure.request_id
563 else :
564 failure_status, failure_id = response.status, response.request_id
565 failure = HelperFailure( "upload-failed" , f "Upload returned HTTP { response.status } ." ,
566 blocked_at = "execution" , status = response.status, request_id = response.request_id,
567 retry_after = response.retry_after, recovery_deadline = response.recovery_deadline)
568 if failure_id:
569 safe_id = ReadRecovery.safe_id(failure_id)
570 failure.request_id = safe_id
571 state[ "request_ids" ].append(safe_id)
572 request_ids.append(safe_id)
573 state[ "status" ] = failure_status
574 proof = None
575 if failure_status in ( 200 , 201 ):
576 state.update( upload = "accepted" , verification = "pending" )
577 created.append({ "fileName" : record[ "path" ], "sha256" : record[ "sha256" ]})
578 if response is not None :
579 proof = _proof(plan, record, response.body)
580 if isinstance (response.body, dict ) and _file_id(response.body.get( "fileId" )):
581 state[ "ack_file_id" ] = response.body[ "fileId" ]
582 elif _rejected(failure_status):
583 state.update( upload = "rejected" , verification = "failed" )
584 try :
585 if response is not None and response.ack_failure is not None :
586 raise response.ack_failure
587 if session:
588 session.result(index, state[ "upload" ], failure_status, failure_id, retry = bool (attempt),
589 proof = proof, file_id = state.get( "ack_file_id" ))
590 except HelperFailure as persistence:
591 primary = primary or (failure if failure_status not in ( 200 , 201 ) else persistence)
592 primary.warnings.extend(persistence.warnings)
593 warnings.append( "Upload journal persistence failed; no further requests were issued." )
594 stopped = True
595 break
596 if failure_status in ( 200 , 201 ) and not transport_failed:
597 if proof:
598 state.update( file_proof = proof, ingested = True , verification = "confirmed" )
599 elif isinstance (response.body, dict ) and (
600 response.body.get( "errorMessage" ) is not None
601 or "fileName" in response.body and not file_ingest._matches(response.body, plan, record)
602 ):
603 primary = primary or fail( "upload-metadata-unverified" , "Upload ACK metadata conflicts with the approved file." )
604 stopped = True
605 _progress(progress, plan, states)
606 break
607 primary = primary or failure
608 if failure_status == 415 :
609 break
610 _progress(progress, plan, states)
611 if failure_status == 429 and not retry_used and allow_upload_retry and failure.blocked_at != "local-persistence" :
612 retry_used = True
613 recovery = ReadRecovery( on_wait = progress.waiting)
614 try :
615 recovery.delay(failure)
616 if source_check is not None :
617 source_check(recovery, token)
618 inventory, _ = file_ingest._list_files(url, token, transport = transport, recovery = recovery)
619 seen = _observe(plan, states, inventory)
620 if index in seen:
621 if state[ "verification" ] != "confirmed" :
622 raise fail( "file-upload-retry-conflict" , "Existing file does not prove the approved completed upload." )
623 created.append({ "fileName" : record[ "path" ], "sha256" : record[ "sha256" ]})
624 break
625 file_ingest._resolve_file(root, record)
626 if time.monotonic() >= recovery.deadline:
627 raise fail( "read-recovery-budget-exhausted" , "Retry preflight exceeded its complete read budget." )
628 warnings.append( "File upload HTTP 429: one same-operation retry after bounded backoff; no source recreation." )
629 except HelperFailure as read_failure:
630 warnings.append( f "File retry stopped ( { read_failure.code } ; HTTP { read_failure.http_status } )." )
631 warnings.extend(read_failure.warnings)
632 stopped = True
633 finally :
634 request_ids.extend(recovery.request_ids)
635 warnings.extend(recovery.diagnostics())
636 if not stopped:
637 continue
638 else :
639 stopped = True
640 warnings.append( f "Batch stopped ( { failure.code } ; HTTP { failure_status } ); no further uploads." )
641 if failure.blocked_at != "local-persistence" and (
642 failure_status in ( 408 , 409 ) or failure_status is None or failure_status >= 500
643 ):
644 recovery = ReadRecovery( on_wait = progress.waiting)
645 try :
646 after, _ = file_ingest._list_files(url, token, transport = transport, recovery = recovery)
647 _observe(plan, states, after)
648 if state[ "verification" ] == "confirmed" :
649 created.append({ "fileName" : record[ "path" ], "sha256" : record[ "sha256" ]})
650 except HelperFailure as read_failure:
651 warnings.append( f "File readback stopped ( { read_failure.code } ; HTTP { read_failure.http_status } )." )
652 warnings.extend(read_failure.warnings)
653 request_ids.extend(recovery.request_ids)
654 warnings.extend(recovery.diagnostics())
655 break
656 if stopped:
657 break
658 if session:
659 session.active_attempt = None
660 _progress(progress, plan, states)
661 progress.update( "file-readback" )
662 if all (state[ "upload" ] == "rejected" and state[ "verification" ] != "confirmed" for state in states):
663 stopped = True
664 if after is None and not stopped:
665 recovery = ReadRecovery( on_wait = progress.waiting)
666 try :
667 after, ids = file_ingest._list_files(url, token, transport = transport, recovery = recovery)
668 _observe(plan, states, after)
669 except HelperFailure as failure:
670 primary = primary or failure
671 readback_failed = True
672 warnings.append( f "File readback stopped ( { failure.code } ; HTTP { failure.http_status } )." )
673 request_ids.extend(recovery.request_ids)
674 warnings.extend(recovery.diagnostics())
675 batch = _summary(plan, states, session)
676 batch[ "request_ids" ] = request_ids[: 804 ]
677 batch[ "upload_retry_used" ] = retry_used
678 progress.update( "file-readback" , files_verified = batch[ "counts" ][ "confirmed" ], files_failed = batch[ "counts" ][ "failed" ],
679 files_pending = batch[ "counts" ][ "pending" ], files_unverified = batch[ "counts" ][ "unverified" ],
680 files_not_attempted = batch[ "counts" ][ "not_attempted" ], files_ingested = batch[ "ingested" ])
681 if batch[ "counts" ][ "confirmed" ] != len (records) or readback_failed:
682 failure = primary or fail( "readback-mismatch" , "Some approved files remain pending, failed or unverified." )
683 uncertain = [{ "action" : "upload-unverified" , "type" : "knowledge-source-file" ,
684 "fileName" : record[ "path" ], "sha256" : record[ "sha256" ]}
685 for record, state in zip (records, states)
686 if state[ "upload" ] == "unverified" and state[ "verification" ] != "confirmed" ]
687 failure.writes = created + uncertain + failure.writes
688 failure.resources_remaining = file_ingest._remaining_files(created) + uncertain + failure.resources_remaining
689 failure.partial = bool (created) or any (s[ "upload" ] == "unverified" for s in states) or (
690 failure.partial and not _rejected(failure.http_status))
691 failure.warnings.extend(warnings)
692 failure.file_batch = batch
693 raise failure
694 verified = [{ "fileId" : item.get( "fileId" ), "fileName" : record[ "path" ], "sha256" : record[ "sha256" ], "size" : record[ "size" ]}
695 for record in records for item in after or before if item.get( "fileName" ) == record[ "path" ]]
696 return { "status" : "completed" , "outcome" : plan.get( "outcome" , "file-knowledge-source-ingestion" ),
697 "approved_plan" : { "fingerprint" : fingerprint, "confirmed" : True },
698 "resources" : { "created" : created, "reused" : reused, "updated" : [], "skipped" : []},
699 "api_contracts" : [{ "operation" : "upload-file" , "version" : file_ingest. API_VERSION , "preview" : True }],
700 "data_movement" : { "boundary" : { "local_root_digest" : digest( str (root))}, "result" : "Approved direct File uploads" },
701 "auth" : { "mode" : "entra-user" , "principals" : []}, "rbac" : plan.get( "rbac" , { "assignments" : []}),
702 "network" : plan.get( "network" , { "posture" : "preserved" , "evidence" : None }),
703 "verification" : { "readback" : verified, "request_ids" : request_ids, "server_inventory_digest" : file_ingest.inventory_digest(after or before),
704 "idempotency" : "Only a received File upload 429 permits one bounded same-operation retry; unknown outcomes are never replayed." },
705 "warnings" : warnings, "file_batch" : batch,
706 "ownership" : { "run_owned" : created, "reused_not_owned" : reused, "owner" : plan[ "owner" ]},
707 "cleanup" : { "status" : "not-requested" , "separate_confirmation_required" : True }}
708
709
710 def _source_check (session, token, transport, recovery):
711 try :
712 from . import source_vector
713 except ImportError :
714 import source_vector
715 plan = session.plan
716 file_ingest._validate_plan(session.ingestion)
717 etag = session.require_ack()
718 current, _ = search_reconcile._get(search_reconcile.resource_url(plan[ "source" ]), token,
719 transport = source_vector.guard_readback_transport(plan, transport), recovery = recovery)
720 if (current is None or current.get( "@odata.etag" ) != etag
721 or not search_reconcile.definitions_match(plan[ "source" ][ "desired" ], current)):
722 raise fail( "file-upload-source-drift" , "Current source does not match original acknowledged identity/version/definition." )
723 return current
724
725
726 def _verify (session, token_provider, transport, context_provider):
727 try :
728 from . import file_source, file_cu_mi
729 except ImportError :
730 import file_source, file_cu_mi
731 session.require_ack()
732 session.states()
733 session.require_backoff()
734 transport = session.transport(transport)
735 context = context_provider()
736 cu_ingestion_auth.validate_context(context)
737 if context != session.seed[ "context" ]:
738 raise fail( "file-upload-context-drift" , "Original tenant/subscription/principal changed." )
739 plan = session.plan
740 if "content_understanding" in plan:
741 state, _ = file_source.read_content_understanding(plan[ "content_understanding" ], token_provider = token_provider, transport = transport)
742 if not file_source._cu_states_match(state, plan[ "cu_resource_state" ]):
743 raise fail( "file-upload-auth-drift" , "Original CU account/auth/network state changed." )
744 if plan[ "source" ].get( "ai_services_managed_identity" ):
745 binding, _ = file_cu_mi.read_binding(plan[ "content_understanding" ], plan[ "source" ][ "endpoint" ],
746 token_provider = token_provider, transport = transport)
747 if binding != plan[ "cu_identity_state" ]:
748 raise fail( "file-upload-auth-drift" , "Original Search MI/CU role/network binding changed." )
749 return _source_check(session, token_provider( SEARCH_AUDIENCE ), transport, ReadRecovery())
750
751
752 def plan_resume (request, * , token_provider = azure_cli_token, transport = http_request,
753 context_provider = cu_ingestion_auth.account_context):
754 if not isinstance (request, dict ) or not isinstance (request.get( "receipt_directory" ), str ):
755 raise fail( "file-upload-input-invalid" , "Use an object with the original private receipt directory." )
756 require_allowed_fields(request, { "schema_version" , "receipt_directory" }, label = "upload resume request" )
757 if request.get( "schema_version" ) != "1.0" :
758 raise fail( "file-upload-input-invalid" , "Use the supported resume request schema." )
759 session = Session(request.get( "receipt_directory" ))
760 transport = session.transport(transport)
761 _verify(session, token_provider, transport, context_provider)
762 states = session.states()
763 inventory, _ = file_ingest._list_files(file_ingest._list_url(session.ingestion), token_provider( SEARCH_AUDIENCE ),
764 transport = transport, recovery = ReadRecovery())
765 observed = _observe(session.ingestion, states, inventory)
766 eligible = [i for i, state in enumerate (states) if state[ "upload" ] == "not_attempted" and i not in observed]
767 plan = { "operation" : "resume-file-uploads" , "version" : "1.0" , "receipt_directory" : str (session.directory),
768 "original_plan_digest" : digest(session.plan), "creation_ack_digest" : digest(session.records[ "source-ack.json" ]),
769 "journal_digest" : digest(session.records), "eligible" : eligible, "owner" : session.plan[ "owner" ], "cleanup_approved" : False }
770 return { "status" : "planned" , "execution_input" : { "schema_version" : "1.0" , "plan" : plan,
771 "approval" : { "confirmed" : False , "fingerprint" : digest(plan)}},
772 "approval_summary" : { "execution_required" : bool (eligible), "mutation_approval_required" : bool (eligible), "uploads" : len (eligible)},
773 "file_batch" : _summary(session.ingestion, states, session), "writes_performed" : [],
774 "next_step" : "Approve only never-attempted uploads. No eligible files: retain this observation; do not replay uncertain files." }
775
776
777 @reporting ( "file-source" )
778 def execute (document, * , token_provider = azure_cli_token, transport = http_request,
779 context_provider = cu_ingestion_auth.account_context, progress = None ):
780 if not isinstance (document, dict ) or not isinstance (document.get( "plan" ), dict ):
781 raise fail( "file-upload-input-invalid" , "Use the newly approved upload-only envelope." )
782 require_allowed_fields(document, { "schema_version" , "plan" , "approval" , "_computed_fingerprint" }, label = "upload resume envelope" )
783 plan = document.get( "plan" )
784 reject_secrets(document)
785 require_allowed_fields(plan, { "operation" , "version" , "receipt_directory" , "original_plan_digest" , "creation_ack_digest" ,
786 "journal_digest" , "eligible" , "owner" , "cleanup_approved" }, label = "upload resume plan" )
787 if (document.get( "schema_version" ) != "1.0" or plan.get( "operation" ) != "resume-file-uploads" or plan.get( "version" ) != "1.0"
788 or set (plan) != { "operation" , "version" , "receipt_directory" , "original_plan_digest" , "creation_ack_digest" ,
789 "journal_digest" , "eligible" , "owner" , "cleanup_approved" }
790 or not isinstance (document.get( "approval" ), dict ) or document[ "approval" ].get( "confirmed" ) is not True
791 or plan.get( "cleanup_approved" ) is not False or document.get( "approval" ) != { "confirmed" : True , "fingerprint" : digest(plan)}
792 or document.get( "_computed_fingerprint" ) != digest(plan)):
793 raise fail( "file-upload-approval-missing" , "A newly approved unchanged upload-only plan is required." )
794 session = Session(plan.get( "receipt_directory" ))
795 transport = session.transport(transport)
796 if (plan[ "original_plan_digest" ] != digest(session.plan) or plan[ "journal_digest" ] != digest(session.records)
797 or plan[ "creation_ack_digest" ] != digest(session.records[ "source-ack.json" ]) or plan[ "owner" ] != session.plan[ "owner" ]
798 or not isinstance (plan[ "eligible" ], list ) or any ( type (i) is not int or not 0 <= i < len (session.ingestion[ "files" ]) for i in plan[ "eligible" ])
799 or len ( set (plan[ "eligible" ])) != len (plan[ "eligible" ])):
800 raise fail( "file-upload-plan-drift" , "Original approval/ACK/journal changed; refresh the upload-only plan." )
801 states = session.states()
802 if any (states[i][ "upload" ] != "not_attempted" for i in plan[ "eligible" ]):
803 raise fail( "file-upload-replay-forbidden" , "Attempted files cannot be submitted again, even after an empty inventory." )
804 progress.update( "source-reconciliation" )
805 try :
806 _verify(session, token_provider, transport, context_provider)
807 result = run_batch({ "plan" : session.ingestion, "_computed_fingerprint" : digest(session.plan)}, token_provider = token_provider,
808 transport = transport, progress = progress, allow_new_uploads = True , session = session, eligible = plan[ "eligible" ],
809 source_check =lambda recovery, token: _source_check(session, token, transport, recovery))
810 except HelperFailure as failure:
811 if failure.file_batch is None :
812 failure.file_batch = _summary(session.ingestion, states, session)
813 failure.partial = True
814 failure.resources_remaining.insert( 0 , { "type" : "knowledge-source" , "name" : session.plan[ "source" ][ "name" ]})
815 raise
816 result[ "approved_plan" ] = document[ "approval" ]
817 result[ "original_run" ] = { "plan_digest" : digest(session.plan), "creation_ack_digest" : plan[ "creation_ack_digest" ],
818 "source_retained" : session.plan[ "source" ][ "name" ]}
819 return result
820
821
822 def main (argv = None ):
823 try :
824 from .private_artifacts import add_execution_output_argument, emit_plan_result, validate_execution_output_mode
825 except ImportError :
826 from private_artifacts import add_execution_output_argument, emit_plan_result, validate_execution_output_mode
827 parser = argparse.ArgumentParser()
828 modes = parser.add_mutually_exclusive_group( required = True )
829 modes.add_argument( "--plan" , type = Path)
830 modes.add_argument( "--input" , type = Path)
831 add_execution_output_argument(parser)
832 add_progress_argument(parser)
833 args = parser.parse_args(argv)
834 fingerprint = None
835 try :
836 validate_execution_output_mode(args)
837 if args.plan:
838 emit_plan_result(plan_resume(private_io.read_json(args.plan)), args.execution_output)
839 return 0
840 document, _, fingerprint = load_approved_input(args.input)
841 document[ "_computed_fingerprint" ] = fingerprint
842 result = execute(document, progress = Progress( "file-source" , enabled = args.progress))
843 except HelperFailure as failure:
844 result = blocked_result(failure, outcome = "resume-file-uploads" , fingerprint = fingerprint)
845 result[ "safe_next_decision" ] = "Retain the source. Missing/changed original evidence blocks upload continuation, not retention; no creation replay or cleanup workaround."
846 emit_result(result)
847 return 3 if result[ "status" ] == "partial" else 2
848 emit_result(result)
849 return 0
850
851
852 if __name__ == "__main__" :
853 sys.exit(main())