fix: bound pending publication admission (#101)

This commit is contained in:
xavierk
2026-09-28 15:02:58 +05:30
parent 5aeee7b6f1
commit fd4db45ae1
3 changed files with 384 additions and 59 deletions
+127 -59
View File
@@ -32,6 +32,26 @@ class InvariantViolationError(Exception):
pass
class _PendingDerivationError(Exception):
"""A valid staged observation failed while deriving dependent evidence."""
def __init__(self, cause: Exception):
super().__init__(str(cause))
self.cause = cause
class _PendingRecoveryFailure(Exception):
"""A publication attempt failed after earlier pending rows committed."""
def __init__(self, error: Exception, published_count: int):
derivation_error = isinstance(error, _PendingDerivationError)
cause = error.cause if derivation_error else error
super().__init__(str(cause))
self.cause = cause
self.published_count = published_count
self.derivation_error = derivation_error
PENDING_PUBLICATION_LIMIT = 6720
@@ -300,47 +320,50 @@ def _publish_observation(
current_id = seg_info["sample_id"]
prev = find_previous_sample(conn, seg_info.get("segment_id"), current_id)
if prev is not None:
current = {
"id": current_id,
"ts": sample["ts"],
"bytes_written": sample["bytes_written"],
"bytes_read": sample["bytes_read"],
"power_on_hours": sample["power_on_hours"],
"temperature_c": sample["temperature_c"],
"data_units_written": sample["data_units_written"],
"data_units_read": sample["data_units_read"],
"local_tz": tz_name,
}
derive_hours_from_interval(conn, prev, current)
from .local_day import record_local_activity_interval
record_local_activity_interval(
conn,
prev,
current,
start_sample_id=prev["id"],
end_sample_id=current_id,
segment_id=seg_info.get("segment_id"),
)
else:
previous_any_segment = find_previous_sample(conn, None, current_id)
if previous_any_segment is not None:
from .local_day import mark_local_activity_gap
mark_local_activity_gap(
try:
if prev is not None:
current = {
"id": current_id,
"ts": sample["ts"],
"bytes_written": sample["bytes_written"],
"bytes_read": sample["bytes_read"],
"power_on_hours": sample["power_on_hours"],
"temperature_c": sample["temperature_c"],
"data_units_written": sample["data_units_written"],
"data_units_read": sample["data_units_read"],
"local_tz": tz_name,
}
derive_hours_from_interval(conn, prev, current)
from .local_day import record_local_activity_interval
record_local_activity_interval(
conn,
sample["ts"],
tz_name,
previous_any_segment,
prev,
current,
start_sample_id=prev["id"],
end_sample_id=current_id,
segment_id=seg_info.get("segment_id"),
)
else:
previous_any_segment = find_previous_sample(conn, None, current_id)
if previous_any_segment is not None:
from .local_day import mark_local_activity_gap
mark_local_activity_gap(
conn,
sample["ts"],
tz_name,
previous_any_segment,
)
from .day_aggregate import derive_all_days, persist_day_aggregate
for aggregate in derive_all_days(conn):
persist_day_aggregate(conn, aggregate)
from .day_aggregate import derive_all_days, persist_day_aggregate
for aggregate in derive_all_days(conn):
persist_day_aggregate(conn, aggregate)
from .local_day import derive_local_day_summary, persist_local_day
local_summary = derive_local_day_summary(conn, tz_name, observed_at)
if local_summary is not None:
persist_local_day(conn, local_summary)
from .local_day import derive_local_day_summary, persist_local_day
local_summary = derive_local_day_summary(conn, tz_name, observed_at)
if local_summary is not None:
persist_local_day(conn, local_summary)
except Exception as exc: # noqa: BLE001 - preserve valid staged evidence for retry.
raise _PendingDerivationError(exc) from exc
def _pending_count(conn: sqlite3.Connection) -> int:
@@ -381,12 +404,15 @@ def _recover_pending(conn: sqlite3.Connection) -> int:
pending_id, payload = row
observation = json.loads(payload)
_publish_observation(
conn,
observation["sample"],
observation["identity"],
observation["tz_name"],
)
try:
_publish_observation(
conn,
observation["sample"],
observation["identity"],
observation["tz_name"],
)
except Exception as exc:
raise _PendingRecoveryFailure(exc, recovered) from exc
conn.execute("DELETE FROM pending_publications WHERE id = ?", (pending_id,))
conn.commit()
recovered += 1
@@ -395,26 +421,61 @@ def _recover_pending(conn: sqlite3.Connection) -> int:
raise
def _is_store_or_invariant_failure(exc: Exception) -> bool:
return (
isinstance(exc, sqlite3.Error)
and not isinstance(exc, sqlite3.IntegrityError)
) or isinstance(exc, InvariantViolationError)
def _pending_capacity_error(reason: str) -> RuntimeError:
return RuntimeError(
"pending publication capacity full "
f"({PENDING_PUBLICATION_LIMIT} observations); no new observation acquired; "
f"{reason}"
)
def _recover_pending_for_collection(conn: sqlite3.Connection) -> int:
"""Retry queued work and report exhausted capacity without masking store faults."""
try:
return _recover_pending(conn)
except Exception as exc:
if isinstance(exc, sqlite3.Error) and not isinstance(exc, sqlite3.IntegrityError):
raise
if isinstance(exc, InvariantViolationError):
raise
except Exception as exc: # noqa: BLE001 - preserve pending evidence on recovery errors.
cause = exc.cause if isinstance(exc, _PendingRecoveryFailure) else exc
if _is_store_or_invariant_failure(cause):
raise cause
try:
capacity_full = _pending_count(conn) >= PENDING_PUBLICATION_LIMIT
except sqlite3.Error:
raise exc
raise cause
if capacity_full:
raise RuntimeError(
"pending publication capacity full "
f"({PENDING_PUBLICATION_LIMIT} observations); no new observation acquired; "
f"recovery failed: {exc}"
) from exc
raise
raise _pending_capacity_error(f"recovery failed: {cause}") from cause
raise cause
def _recover_pending_for_admission(conn: sqlite3.Connection) -> int:
"""Recover in order, then admit one observation if bounded space remains."""
try:
recovered = _recover_pending(conn)
except _PendingRecoveryFailure as failure:
if (
not failure.derivation_error
or _is_store_or_invariant_failure(failure.cause)
):
raise failure.cause
try:
pending_count = _pending_count(conn)
except sqlite3.Error:
raise failure.cause
if pending_count >= PENDING_PUBLICATION_LIMIT:
raise _pending_capacity_error(
f"recovery failed: {failure.cause}"
) from failure.cause
return failure.published_count
if _pending_count(conn) >= PENDING_PUBLICATION_LIMIT:
raise _pending_capacity_error("recovery left the queue full")
return recovered
def run_collection(
@@ -444,16 +505,23 @@ def run_collection(
if history_path.exists():
import_legacy_history(conn, history_path, clock=clock)
published_count = _recover_pending_for_collection(conn)
published_count = _recover_pending_for_admission(conn)
checked_pending_count = _pending_count(conn)
# Hold the writer reservation across the capacity check and acquisition.
# A concurrent collector will recheck pending work before it acquires.
# Recheck recovery if another collector appended work during preflight.
while True:
conn.execute("BEGIN IMMEDIATE")
waiting = _pending_count(conn)
if waiting:
if waiting >= PENDING_PUBLICATION_LIMIT:
conn.rollback()
published_count += _recover_pending_for_collection(conn)
published_count += _recover_pending_for_admission(conn)
checked_pending_count = _pending_count(conn)
continue
if waiting > checked_pending_count:
conn.rollback()
published_count += _recover_pending_for_admission(conn)
checked_pending_count = _pending_count(conn)
continue
from .tz_util import detect_system_tz