diff --git a/src/fenris/collector.py b/src/fenris/collector.py index 21f829f..a09e317 100644 --- a/src/fenris/collector.py +++ b/src/fenris/collector.py @@ -422,20 +422,16 @@ def run_collection( if history_path.exists(): import_legacy_history(conn, history_path, clock=clock) - recovered_count = _recover_pending_for_collection(conn) + published_count = _recover_pending_for_collection(conn) # Hold the writer reservation across the capacity check and acquisition. # A concurrent collector will recheck pending work before it acquires. while True: conn.execute("BEGIN IMMEDIATE") waiting = _pending_count(conn) - if waiting >= PENDING_PUBLICATION_LIMIT: - conn.rollback() - recovered_count += _recover_pending_for_collection(conn) - continue if waiting: conn.rollback() - recovered_count += _recover_pending_for_collection(conn) + published_count += _recover_pending_for_collection(conn) continue from .tz_util import detect_system_tz @@ -458,10 +454,10 @@ def run_collection( conn.commit() break - recovered_count += _recover_pending_for_collection(conn) + published_count += _recover_pending_for_collection(conn) return { "ok": True, - "sample_count": recovered_count, + "sample_count": published_count, "store_path": str(store_path), }