diff --git a/src/fenris/collector.py b/src/fenris/collector.py index 99ec28e..40faac7 100644 --- a/src/fenris/collector.py +++ b/src/fenris/collector.py @@ -428,10 +428,15 @@ def _is_store_or_invariant_failure(exc: Exception) -> bool: ) or isinstance(exc, InvariantViolationError) -def _pending_capacity_error(reason: str) -> RuntimeError: +def _pending_capacity_error( + reason: str, + *, + acquisition_skipped: bool = False, +) -> RuntimeError: + outcome = "; no new observation acquired" if acquisition_skipped else "" return RuntimeError( "pending publication capacity full " - f"({PENDING_PUBLICATION_LIMIT} observations); no new observation acquired; " + f"({PENDING_PUBLICATION_LIMIT} observations){outcome}; " f"{reason}" ) @@ -469,12 +474,16 @@ def _recover_pending_for_admission(conn: sqlite3.Connection) -> int: raise failure.cause if pending_count >= PENDING_PUBLICATION_LIMIT: raise _pending_capacity_error( - f"recovery failed: {failure.cause}" + f"recovery failed: {failure.cause}", + acquisition_skipped=True, ) from failure.cause return failure.published_count if _pending_count(conn) >= PENDING_PUBLICATION_LIMIT: - raise _pending_capacity_error("recovery left the queue full") + raise _pending_capacity_error( + "recovery left the queue full", + acquisition_skipped=True, + ) return recovered diff --git a/tests/test_issue_101.py b/tests/test_issue_101.py index 9078792..fd7d2df 100644 --- a/tests/test_issue_101.py +++ b/tests/test_issue_101.py @@ -139,6 +139,7 @@ def test_partial_recovery_frees_admission_and_keeps_pending_order( ) assert partial["ok"] is False assert acquired is True, partial + assert "no new observation acquired" not in partial["error"] with sqlite3.connect(store_path) as conn: assert conn.execute(