@@ -133,6 +133,13 @@ async def create_run(
133133 for s in staged :
134134 await self .db .link_run_document (run_id , s .document_id , status = "pending" )
135135
136+ logger .info (
137+ "run=%s queued docs=%d pages=%d name=%s" ,
138+ run_id ,
139+ len (staged ),
140+ pages_total ,
141+ name or "-" ,
142+ )
136143 observability .job_created ()
137144 return CreateRunResult (
138145 run_id = run_id , documents = staged , pages_total_estimate = pages_total
@@ -167,6 +174,7 @@ async def retry_incomplete_run(self, run_id: str) -> CreateRunResult:
167174 raise ValueError ("No incomplete documents to retry" )
168175
169176 name = run .get ("name" ) or run_id
177+ logger .info ("run=%s retry queued incomplete_docs=%d" , run_id , len (retry_paths ))
170178 result = await self .create_run (
171179 retry_paths ,
172180 name = f"{ name } retry" ,
@@ -189,15 +197,20 @@ async def _run(
189197 ) -> None :
190198 run_id = result .run_id
191199 started_at = _now ()
200+ pages_total = result .pages_total_estimate
201+ pages_completed = 0
202+ documents_meta : list = []
192203 await self .db .update_run (
193204 run_id , status = "processing" , stage = "ocr" , started_at = started_at
194205 )
206+ logger .info (
207+ "run=%s started docs=%d pages=%d" ,
208+ run_id ,
209+ len (result .documents ),
210+ pages_total ,
211+ )
195212 await self ._emit (run_id , "run_started" , {"started_at" : started_at })
196213
197- pages_total = result .pages_total_estimate
198- pages_completed = 0
199- documents_meta : list = []
200-
201214 async def page_event (event : dict ) -> None :
202215 nonlocal pages_completed
203216 etype = event .get ("type" )
@@ -214,15 +227,39 @@ async def page_event(event: dict) -> None:
214227 token_count = event .get ("token_count" , 0 ),
215228 validation_status = event .get ("validation_status" , "pass" ),
216229 )
230+ logger .info (
231+ "run=%s page=%d/%d doc=%s status=%s time=%.1fms" ,
232+ run_id ,
233+ pages_completed ,
234+ pages_total ,
235+ event .get ("document" ),
236+ event .get ("validation_status" ),
237+ event .get ("processing_time_ms" , 0.0 ),
238+ )
217239 elif etype == "page_retry" :
218240 observability .page_retry ()
241+ logger .info (
242+ "run=%s page=%s retry attempt=%s strategy=%s reason=%s" ,
243+ run_id ,
244+ event .get ("page" ),
245+ event .get ("attempt" ),
246+ event .get ("new_strategy" ),
247+ event .get ("reason" ),
248+ )
219249 await self ._emit (run_id , etype , event )
220250
221251 try :
222- for staged in result .documents :
252+ for index , staged in enumerate ( result .documents , start = 1 ) :
223253 paths = self .storage .artifact_paths (
224254 run_id , staged .document_id , staged .filename
225255 )
256+ logger .info (
257+ "run=%s doc=%d/%d started %s" ,
258+ run_id ,
259+ index ,
260+ len (result .documents ),
261+ staged .filename ,
262+ )
226263 processor = BatchProcessor (
227264 self .db , event_callback = page_event , strip_refs = strip_refs
228265 )
@@ -238,6 +275,16 @@ async def page_event(event: dict) -> None:
238275 await self .db .update_run (
239276 run_id , documents_completed = len (documents_meta )
240277 )
278+ logger .info (
279+ "run=%s doc=%d/%d completed %s pass=%d warn=%d fail=%d" ,
280+ run_id ,
281+ index ,
282+ len (result .documents ),
283+ staged .filename ,
284+ doc_meta .pages_pass ,
285+ doc_meta .pages_warn ,
286+ doc_meta .pages_fail ,
287+ )
241288
242289 dataset_bundle = await self ._maybe_export (
243290 run_id , documents_meta , export_parquet
@@ -253,6 +300,7 @@ async def page_event(event: dict) -> None:
253300 dataset_bundle = dataset_bundle ,
254301 completed_at = completed_at ,
255302 )
303+ logger .info ("run=%s completed bundle=%s" , run_id , dataset_bundle or "-" )
256304 observability .job_completed ()
257305 await self ._emit (
258306 run_id ,
@@ -285,6 +333,7 @@ async def _maybe_export(
285333 if not (export_parquet and documents_meta ):
286334 return None
287335 await self .db .update_run (run_id , stage = "exporting" )
336+ logger .info ("run=%s exporting dataset docs=%d" , run_id , len (documents_meta ))
288337 await self ._emit (run_id , "dataset_export_started" , {})
289338 exports = []
290339 for did , paths , meta in documents_meta :
0 commit comments