diff --git a/cron/bot_chat_delivery.py b/cron/bot_chat_delivery.py index 15f692c666..90d0803784 100644 --- a/cron/bot_chat_delivery.py +++ b/cron/bot_chat_delivery.py @@ -16,6 +16,7 @@ from hermes_constants import get_hermes_home from utils import atomic_json_write logger = logging.getLogger(__name__) +_warned_unreadable: set[Path] = set() _running: set[Path] = set() _running_lock = threading.Lock() @@ -36,11 +37,15 @@ def _records(root: Path) -> list[tuple[Path, dict]]: for path in root.glob("*.json"): try: record = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError, UnicodeDecodeError) as exc: + except (OSError, ValueError) as exc: # ValueError: corrupt JSON and invalid UTF-8 alike # Keep damaged or unreadable receipts as evidence; never replay them or block peers # (same rule as tools/bot_live_delivery.py::_scan_read — one bad file must not wedge the dir). - logger.error("Unreadable deferred Bot Chat receipt %s: %s", path, exc) + # The scheduler drains every tick: ERROR once per receipt per process, DEBUG after. + level = logging.DEBUG if path in _warned_unreadable else logging.ERROR + _warned_unreadable.add(path) + logger.log(level, "Unreadable deferred Bot Chat receipt %s: %s", path, exc) continue + _warned_unreadable.discard(path) records.append((path, record)) return records diff --git a/tests/cron/test_bot_chat_pending.py b/tests/cron/test_bot_chat_pending.py index a1801135fb..3f1aefca7c 100644 --- a/tests/cron/test_bot_chat_pending.py +++ b/tests/cron/test_bot_chat_pending.py @@ -120,6 +120,8 @@ def test_unreadable_deferred_receipt_does_not_block_siblings(tmp_path, monkeypat monkeypatch.setattr(delivery, "_deliver_to_bot_chat", lambda j, c, p, **kw: seen.append(c)) with caplog.at_level("ERROR", logger=queue.logger.name): queue.drain() + queue.drain() # every scheduler tick drains; the same bad receipt must not re-log assert seen == ["healthy"] - assert any("Unreadable deferred Bot Chat receipt" in r.message and "Permission denied" in r.message - for r in caplog.records) + assert [r for r in caplog.records + if "Unreadable deferred Bot Chat receipt" in r.message and "Permission denied" in r.message] and \ + sum("Unreadable deferred Bot Chat receipt" in r.message for r in caplog.records) == 1 diff --git a/tests/tools/test_bot_live_owner_delivery.py b/tests/tools/test_bot_live_owner_delivery.py index 21f268c5ec..ee56c31dcd 100644 --- a/tests/tools/test_bot_live_owner_delivery.py +++ b/tests/tools/test_bot_live_owner_delivery.py @@ -119,22 +119,26 @@ def test_unreadable_ticket_does_not_wedge_bulk_scans(tmp_path, caplog): lease_id="lease", live_session_id="live") queued = mailbox.deliver_to_live_owner(tmp_path, owner, "readable", delivery_id="d" * 32) root = tmp_path / "runtime" / mailbox.DELIVERY_DIR_NAME - unreadable = root / f"{'e' * 32}.json" - unreadable.write_text('{"status": "queued"}', encoding="utf-8") - unreadable.chmod(0) + # A real admission that later turns unreadable: its sequence must survive the skip. + hidden = mailbox.deliver_to_live_owner(tmp_path, owner, "hidden", delivery_id="e" * 32) + (root / f"{'e' * 32}.json").chmod(0) corrupt = root / f"{'1' * 32}.json" corrupt.write_text("{not json", encoding="utf-8") + (root / f"{'2' * 32}.json").write_bytes(b"\xff\xfe\x00garbage") # invalid UTF-8, not just bad JSON with caplog.at_level(logging.WARNING, logger="tools.bot_live_delivery"): # Sender side: admission of a fresh id must survive the sequence sweep. admitted = mailbox.deliver_to_live_owner(tmp_path, owner, "second", delivery_id="f" * 32) # Receiver side: every readable queued ticket must still be claimed, in order. assert mailbox.claim_pending_delivery(tmp_path, owner)["delivery_id"] == queued["delivery_id"] assert mailbox.claim_pending_delivery(tmp_path, owner)["delivery_id"] == admitted["delivery_id"] - assert mailbox.claim_pending_delivery(tmp_path, owner) is None + for _ in range(10): # the idle poller rescans twice a second + assert mailbox.claim_pending_delivery(tmp_path, owner) is None assert admitted["status"] == "queued" - assert any(record.message.startswith(f"bot_live_delivery: skipping unreadable ticket {'e' * 32}.json") - and "Permission denied" in record.message - for record in caplog.records) + assert admitted["sequence"] > hidden["sequence"] > queued["sequence"] + denied = [record for record in caplog.records + if record.message.startswith(f"bot_live_delivery: skipping unreadable ticket {'e' * 32}.json") + and "Permission denied" in record.message] + assert len(denied) == 1, "one persistent bad ticket must warn once per process, not per scan" @pytest.mark.skipif(os.name == "nt" or getattr(os, "geteuid", lambda: 1)() == 0, diff --git a/tools/bot_live_delivery.py b/tools/bot_live_delivery.py index f31b244ddb..eecb329dc4 100644 --- a/tools/bot_live_delivery.py +++ b/tools/bot_live_delivery.py @@ -15,7 +15,7 @@ import time import uuid from contextlib import contextmanager -from utils import atomic_json_write, fsync_directory +from utils import atomic_json_write, atomic_write_text, fsync_directory from pathlib import Path from typing import Any @@ -24,6 +24,7 @@ from hermes_cli.active_sessions import _FileLock log = logging.getLogger(__name__) DELIVERY_DIR_NAME = "bot_live_delivery" +_SEQUENCE_FILE = ".sequence" _OWNER_KEYS = ("profile_home", "session_id", "lease_id", "live_session_id") _TERMINAL = frozenset({"settled", "failed", "cancelled", "ambiguous"}) @@ -108,21 +109,51 @@ def _read(path: Path) -> dict[str, Any] | None: return None +# Tickets already reported unreadable by this process. The live poller rescans the +# dir twice a second, so a persistent bad ticket is WARNING once and DEBUG after. +_warned_unreadable: set[Path] = set() + + def _scan_read(path: Path) -> dict[str, Any] | None: """Bulk-scan variant: one unreadable ticket must not wedge the whole dir. Directory scans (sequence high-water mark, queued-claim sweep) may only treat a file as absent when it is provably absent; an unreadable ticket - degrades to "that one delivery is uninspectable" with a loud warning. + degrades to "that one delivery is uninspectable" with a warning. Exact-id reads (admission idempotency, completion, result lookup) keep using _read so a permission error still fails closed instead of licensing an overwrite of a possibly-live receipt. """ try: - return _read(path) - except (OSError, json.JSONDecodeError) as exc: - log.warning("bot_live_delivery: skipping unreadable ticket %s (%s)", path.name, exc) + record = _read(path) + except (OSError, ValueError) as exc: # ValueError: corrupt JSON and invalid UTF-8 alike + level = logging.DEBUG if path in _warned_unreadable else logging.WARNING + _warned_unreadable.add(path) + log.log(level, "bot_live_delivery: skipping unreadable ticket %s (%s)", path.name, exc) return None + _warned_unreadable.discard(path) + return record + + +def _next_sequence(root: Path) -> int: + """Allocate the next admission sequence under the dir lock. + + The high-water mark lives in a counter file beside the tickets, so a ticket + the scan cannot read does not drop its sequence and hand a later admission + a duplicate or lower one. Readable tickets still bootstrap dirs written + before the counter existed. Wall time can roll back; sequences never do. + """ + counter = root / _SEQUENCE_FILE + try: + persisted = int(counter.read_text(encoding="utf-8")) + except (OSError, ValueError): + persisted = 0 + scanned = max((record.get("sequence", record["created_at"]) + for candidate in root.glob("*.json") + if (record := _scan_read(candidate)) is not None), default=0) + sequence = max(persisted, scanned) + 1 + atomic_write_text(counter, str(sequence), mode=0o600, fsync_dir=True) + return sequence def _write(path: Path, record: dict[str, Any]) -> None: @@ -149,14 +180,9 @@ def deliver_to_live_owner( if existing["owner"] != pinned or existing["message"] != message or existing.get("author") != author: raise ValueError("delivery id already belongs to a different payload") return existing - # Wall time can roll back. Permanent receipts retain the admission - # high-water mark, allocated while holding the cross-process lock. - sequence = max((record.get("sequence", record["created_at"]) - for candidate in root.glob("*.json") - if (record := _scan_read(candidate)) is not None), default=0) + 1 record = dict(delivery_id=key, id=key, owner=pinned, **pinned, message=message, status="queued", created_at=time.time_ns(), - sequence=sequence, **({"author": dict(author)} if author else {})) + sequence=_next_sequence(root), **({"author": dict(author)} if author else {})) _write(path, record) return record