refactor: clarify publication recovery count

This commit is contained in:
xavierk
2026-09-28 09:06:09 +05:30
parent 847bde1be0
commit 8a5ef05188
+4 -8
View File
@@ -422,20 +422,16 @@ def run_collection(
if history_path.exists(): if history_path.exists():
import_legacy_history(conn, history_path, clock=clock) 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. # Hold the writer reservation across the capacity check and acquisition.
# A concurrent collector will recheck pending work before it acquires. # A concurrent collector will recheck pending work before it acquires.
while True: while True:
conn.execute("BEGIN IMMEDIATE") conn.execute("BEGIN IMMEDIATE")
waiting = _pending_count(conn) waiting = _pending_count(conn)
if waiting >= PENDING_PUBLICATION_LIMIT:
conn.rollback()
recovered_count += _recover_pending_for_collection(conn)
continue
if waiting: if waiting:
conn.rollback() conn.rollback()
recovered_count += _recover_pending_for_collection(conn) published_count += _recover_pending_for_collection(conn)
continue continue
from .tz_util import detect_system_tz from .tz_util import detect_system_tz
@@ -458,10 +454,10 @@ def run_collection(
conn.commit() conn.commit()
break break
recovered_count += _recover_pending_for_collection(conn) published_count += _recover_pending_for_collection(conn)
return { return {
"ok": True, "ok": True,
"sample_count": recovered_count, "sample_count": published_count,
"store_path": str(store_path), "store_path": str(store_path),
} }