fix(bot-relay): claim the outbox oldest first, so two DMs to one agent keep their order
The Desktop delivers each target's claimed envelopes in the order the drain returns them, one
turn at a time, and says so at the call site; the user guide states the same guarantee
("messages to the same Bot are delivered in order, one turn at a time"). The claim sorted the
outbox directory by filename, which is uuid4().hex — so the order handed to the Desktop was
random, and a sender's second instruction to one agent could be delivered before its first.
Sort by the mtime _sweep_stale already treats as an envelope's age, with the name only breaking
ties. The whole-second created_at field cannot separate two DMs sent in the same second.
The test forces the filenames into the reverse of the send order, which random ids reproduce
half the time: it fails on the filename sort, on newest-first, and when the mtime is dropped.
This commit is contained in:
committed by
Teknium
parent
b5f3ed77a4
commit
135bf07a75
@@ -528,6 +528,33 @@ def test_drain_expires_old_envelope_with_queued_expired_reply(root):
|
||||
assert not reply["reply"]
|
||||
|
||||
|
||||
def test_the_outbox_is_claimed_oldest_first(root):
|
||||
"""Two DMs from one sender to one agent must arrive in the order they were sent. The Desktop
|
||||
delivers each target's claimed envelopes in the order this list gives them, one turn at a
|
||||
time, so the claim IS the delivery order — and sorting by filename ordered them by
|
||||
``uuid4().hex``. The names here are forced into the reverse of the send order to pin that
|
||||
deterministically, which random ids reproduce half the time."""
|
||||
first = bot_relay.enqueue_envelope(
|
||||
root, target=_target(), message="do this first",
|
||||
sender_profile="default", sender_handle="hermes",
|
||||
)
|
||||
second = bot_relay.enqueue_envelope(
|
||||
root, target=_target(), message="then this",
|
||||
sender_profile="default", sender_handle="hermes",
|
||||
)
|
||||
outbox = bot_relay.relay_root(root) / bot_relay.OUTBOX_DIR
|
||||
now = _time2.time()
|
||||
for env, name, sent_at in ((first, "f" * 32, now - 2), (second, "0" * 32, now - 1)):
|
||||
path = outbox / f"{name}.json"
|
||||
(outbox / f"{env['id']}.json").rename(path)
|
||||
_os2.utime(path, (sent_at, sent_at))
|
||||
|
||||
claimed = bot_relay.claim_pending_envelopes(root)
|
||||
|
||||
assert [e["id"] for e in claimed] == [first["id"], second["id"]]
|
||||
assert [e["message"] for e in claimed] == ["do this first", "then this"]
|
||||
|
||||
|
||||
def test_drain_delivers_fresh_envelope_under_ttl(root):
|
||||
env = bot_relay.enqueue_envelope(
|
||||
root, target=_target(), message="on time",
|
||||
|
||||
@@ -242,6 +242,15 @@ def _expire_if_stale(root: Path | str, path: Path, ttl: float, now: float) -> bo
|
||||
return True
|
||||
|
||||
|
||||
def _queued_at(path: Path) -> tuple[float, str]:
|
||||
"""Claim order for one outbox entry: oldest first. ``mtime`` is what ``_sweep_stale`` already
|
||||
treats as an envelope's age, and unlike the whole-second ``created_at`` field it separates two
|
||||
DMs sent in the same second. The name only breaks ties."""
|
||||
with contextlib.suppress(OSError):
|
||||
return (path.stat().st_mtime, path.name)
|
||||
return (0.0, path.name)
|
||||
|
||||
|
||||
def claim_pending_envelopes(root: Path | str) -> list[dict]:
|
||||
"""Drain the outbox (rename → claimed/ so a second drain can't double-deliver).
|
||||
TTL-expired envelopes get a 'queued_expired' reply and are removed instead.
|
||||
@@ -255,7 +264,10 @@ def claim_pending_envelopes(root: Path | str) -> list[dict]:
|
||||
ttl = _envelope_ttl_seconds()
|
||||
now = time.time()
|
||||
out: list[dict] = []
|
||||
for path in sorted((base / OUTBOX_DIR).glob("*.json")):
|
||||
# Oldest first: the Desktop delivers each target's claimed envelopes in the order this list
|
||||
# gives them, so a sender's two DMs to one agent arrive in the order they were sent. Sorting
|
||||
# by filename ordered them by ``uuid4().hex`` — at random.
|
||||
for path in sorted((base / OUTBOX_DIR).glob("*.json"), key=_queued_at):
|
||||
if ttl > 0 and _expire_if_stale(root, path, ttl, now):
|
||||
with contextlib.suppress(OSError):
|
||||
path.unlink()
|
||||
|
||||
Reference in New Issue
Block a user