From 59d806f9fd609d28c2db862e6b9d9826899c7c45 Mon Sep 17 00:00:00 2001 From: ysz Date: Tue, 29 Sep 2026 11:38:14 +0200 Subject: [PATCH 1/2] fix(ingestion): read generation id before rolling back in the failure handler AsyncSession.rollback() expires every loaded ORM object even with expire_on_commit=False. The failure handler then read generation.id, which needs a lazy refresh and fails in async code ("greenlet_spawn has not been called"), so the error state was again never recorded and the document kept a stale processing_attempt_id. Capture the generation id before the rollback and skip marking the generation failed when none was created. --- .../application/ingestion_service.py | 12 +++-- .../test_ingestion_service_failed_cleanup.py | 45 +++++++++++++++++++ 2 files changed, 53 insertions(+), 4 deletions(-) diff --git a/src/core/ingestion/application/ingestion_service.py b/src/core/ingestion/application/ingestion_service.py index 07d91583..acff8da3 100644 --- a/src/core/ingestion/application/ingestion_service.py +++ b/src/core/ingestion/application/ingestion_service.py @@ -1215,6 +1215,9 @@ async def _on_graph_progress(completed: int, total: int): except Exception as e: logger.exception(f"Failed to process document {document_id}") + # rollback() expires every loaded ORM object; read ids while they are + # still loaded, lazy refreshes are not possible in async code. + generation_id = generation.id if generation is not None else None try: # A failed flush poisons the session; without a rollback every # write below fails and the doc keeps a stale processing_attempt_id. @@ -1233,11 +1236,12 @@ async def _on_graph_progress(completed: int, total: int): logger.error(f"Failed to map error for {document_id}: {map_err}") error_message = f"{type(e).__name__}: {str(e)}" - await self.document_repository.mark_generation_failed( - generation.id, error_message - ) + if generation_id: + await self.document_repository.mark_generation_failed( + generation_id, error_message + ) if preserve_published: - if document.pending_generation_id == generation.id: + if document.pending_generation_id == generation_id: document.pending_generation_id = None await self.document_repository.save(document) else: diff --git a/tests/unit/test_ingestion_service_failed_cleanup.py b/tests/unit/test_ingestion_service_failed_cleanup.py index f08d2170..fba5606d 100644 --- a/tests/unit/test_ingestion_service_failed_cleanup.py +++ b/tests/unit/test_ingestion_service_failed_cleanup.py @@ -319,3 +319,48 @@ async def capture(*args, **kwargs): await service.process_document("doc_10") assert seen["result"].content == "ab" + + +class ExpiringUnitOfWork(PoisonedSessionUnitOfWork): + """Like AsyncSession.rollback(): every loaded ORM object is expired, so + reading any attribute afterwards needs IO (MissingGreenlet in async code).""" + + def __init__(self, repository: FakeDocumentRepositoryForFailure) -> None: + super().__init__() + self.repository = repository + + async def rollback(self) -> None: + from sqlalchemy import inspect + + await super().rollback() + generation = self.repository.generation + if generation is not None: + state = inspect(generation) + state._expire(state.dict, set()) + + +@pytest.mark.asyncio +async def test_process_document_failure_handler_does_not_read_expired_generation(): + document = StubDocument( + id="doc_11", + tenant_id="tenant-1", + status=DocumentStatus.INGESTED, + storage_path="tenant-1/doc_11/file.txt", + filename="file.txt", + content_hash="hash-11", + metadata_={}, + ) + repository = FakeDocumentRepositoryForFailure(document) + uow = ExpiringUnitOfWork(repository) + service = make_service(vector_store=FakeVectorStore(), neo4j_client=FakeNeo4jClient()) + service.document_repository = repository + service.unit_of_work = uow + service.storage = PoisoningStorage(uow) + + with pytest.raises(ValueError, match="storage is down"): + await service.process_document("doc_11") + + assert uow.rollbacks >= 1 + assert document.status == DocumentStatus.FAILED + assert document.error_message + assert document.processing_attempt_id is None From e92c36e6d74f31d6d6e947778ecadccfffed9ce6 Mon Sep 17 00:00:00 2001 From: ysz Date: Fri, 2 Oct 2026 10:28:07 +0200 Subject: [PATCH 2/2] fix(ingestion): capture generation identity before ORM expiration Preserve the scalar identity before writes can expire ORM state. Cover fresh and existing generations and retain published content on reprocess failure. --- .../application/ingestion_service.py | 7 +- .../test_ingestion_service_failed_cleanup.py | 83 +++++++++++++++++++ 2 files changed, 87 insertions(+), 3 deletions(-) diff --git a/src/core/ingestion/application/ingestion_service.py b/src/core/ingestion/application/ingestion_service.py index f3dfe6dd..670e146a 100644 --- a/src/core/ingestion/application/ingestion_service.py +++ b/src/core/ingestion/application/ingestion_service.py @@ -538,10 +538,13 @@ async def process_document(self, document_id: str, force: bool = False) -> None: ) return generation = None + # A failed flush may expire ORM attributes before explicit rollback. + generation_id: str | None = None try: if pending_generation_id: generation = await self.document_repository.get_generation(pending_generation_id) if generation is not None: + generation_id = generation.id # A retry (re-upload with the same content-hash, or a # stale-lock sweep) reaches here with pending_generation_id # still pointing at the PREVIOUS attempt's generation row. @@ -577,6 +580,7 @@ async def process_document(self, document_id: str, force: bool = False) -> None: keywords=list(getattr(document, "keywords", None) or []), hashtags=list(getattr(document, "hashtags", None) or []), ) + generation_id = generation.id await self.document_repository.save_generation(generation) document.pending_generation_id = generation.id await self.document_repository.save(document) @@ -1278,9 +1282,6 @@ async def _on_graph_progress(completed: int, total: int): except Exception as e: logger.exception(f"Failed to process document {document_id}") - # rollback() expires every loaded ORM object; read ids while they are - # still loaded, lazy refreshes are not possible in async code. - generation_id = generation.id if generation is not None else None try: # A failed flush poisons the session; without a rollback every # write below fails and the doc keeps a stale processing_attempt_id. diff --git a/tests/unit/test_ingestion_service_failed_cleanup.py b/tests/unit/test_ingestion_service_failed_cleanup.py index a604249c..e8f46971 100644 --- a/tests/unit/test_ingestion_service_failed_cleanup.py +++ b/tests/unit/test_ingestion_service_failed_cleanup.py @@ -363,3 +363,86 @@ async def test_process_document_failure_handler_does_not_read_expired_generation assert document.status == DocumentStatus.FAILED assert document.error_message assert document.processing_attempt_id is None + + +class ExpiringPoisoningStorage(PoisoningStorage): + """A failed flush can expire ORM state before the handler starts.""" + + def __init__(self, uow, repository): + super().__init__(uow) + self.repository = repository + + def get_file(self, storage_path): + from sqlalchemy import inspect + + generation = self.repository.generation + self.generation_id = generation.id + state = inspect(generation) + state._expire(state.dict, set()) + return super().get_file(storage_path) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("existing_generation", [False, True]) +@pytest.mark.parametrize("preserve_published", [False, True]) +async def test_failure_handler_handles_generation_already_expired_before_rollback( + existing_generation, preserve_published, +): + document = StubDocument( + id="doc_expired_flush", + tenant_id="tenant-1", + status=DocumentStatus.READY if preserve_published else DocumentStatus.INGESTED, + storage_path="tenant-1/doc_expired_flush/file.txt", + filename="file.txt", + content_hash="hash_expired_flush", + metadata_={}, + active_generation_id="published_gen" if preserve_published else None, + ) + repository = FakeDocumentRepositoryForFailure(document) + if existing_generation: + repository.generation = service_module.DocumentGeneration( + id="gen_existing", + document_id=document.id, + tenant_id=document.tenant_id, + filename=document.filename, + content_hash=document.content_hash, + storage_path=document.storage_path, + metadata_={}, + ) + document.pending_generation_id = "gen_existing" + + async def delete_chunks(generation_id): + return 0 + + repository.delete_chunks_by_generation = delete_chunks + + failed_ids = [] + mark_failed = repository.mark_generation_failed + + async def record_failed(generation_id, error_message): + failed_ids.append(generation_id) + await mark_failed(generation_id, error_message) + + repository.mark_generation_failed = record_failed + uow = ExpiringUnitOfWork(repository) + service = make_service(vector_store=FakeVectorStore(), neo4j_client=FakeNeo4jClient()) + service.document_repository = repository + service.unit_of_work = uow + service.storage = ExpiringPoisoningStorage(uow, repository) + + with pytest.raises(ValueError, match="storage is down"): + await service.process_document(document.id, force=preserve_published) + + assert uow.rollbacks >= 1 + assert repository.generation.status == "failed" + assert repository.generation.error_message + assert document.processing_attempt_id is None + if preserve_published: + assert document.status == DocumentStatus.READY + assert document.active_generation_id == "published_gen" + assert document.pending_generation_id is None + assert service.vector_store.delete_calls == [] + else: + assert document.status == DocumentStatus.FAILED + assert document.error_message + assert failed_ids == [service.storage.generation_id]