dev-bucket-sync / tests /test_notify_unit.py
cmpatino's picture
cmpatino HF Staff
Upload folder using huggingface_hub
04fcb03 verified
Raw History Blame Contribute Delete
4.57 kB
"""Direct unit coverage of app/notify.py's Notifier internals.
This file is deliberately narrow. Most of the Notifier's behaviour is already
proven end-to-end through the real app by test_longpoll_api.py and
test_updates_api.py. Three behaviours, though, are subtle enough — and internal
enough — that the API-level tests only exercise them indirectly, as a side
effect of some other assertion, rather than pinning them directly:
- waking a parked waiter from a foreign OS thread, which is the exact bridge
(``loop.call_soon_threadsafe`` on the future captured at ``register`` time)
production relies on when the verifier's worker thread announces a verdict
from outside the event loop;
- a wake's latch absorbing a signal that arrives before ``wait`` is ever
called, so a wake racing a client's reconnect is never lost;
- double-wake idempotency, so two writers touching the same key in quick
succession don't double-count or raise.
Kept here rather than folded into an API test because sinking a bug in any of
these three would surface only as an intermittent, hard-to-reproduce timing
flake three layers up (a long-poll that occasionally holds for its full
timeout instead of waking promptly) — a direct unit test turns that into a
deterministic, immediate failure at the primitive itself.
"""
from __future__ import annotations
import asyncio
import threading
import time
from app.notify import Notifier
def _notifier(per_owner: int = 4, total: int = 256) -> Notifier:
# Mirrors production Settings defaults (config.py) for the two spread
# knobs: a threshold of 20 keeps every wake in these tests (at most a
# couple of subscriptions) on the instant path, never the spread-out-over-
# wake_spread_s path meant for large broadcasts.
return Notifier(
max_waiters_per_owner=per_owner,
max_waiters_total=total,
wake_spread_s=8.0,
wake_spread_threshold=20,
)
def test_latch_absorbs_wake_before_wait():
n = _notifier()
async def scenario():
sub = n.register("a", {"k"})
n.wake({"k"}) # arrives before the first wait()
first = await sub.wait(0.01) # consumes the latch immediately
second = await sub.wait(0.05) # latch cleared -> times out
return first, second
first, second = asyncio.run(scenario())
assert first is True # not lost despite arriving between/around waits
assert second is False
def test_double_wake_is_idempotent():
n = _notifier()
async def scenario():
sub = n.register("a", {"k"})
n.wake({"k"})
n.wake({"k"}) # second wake must not error or double-count
first = await sub.wait(0.5)
second = await sub.wait(0.05) # only one latch was pending
return first, second
first, second = asyncio.run(scenario())
assert first is True
assert second is False
def test_wake_from_foreign_thread():
n = _notifier()
async def scenario():
sub = n.register("a", {"k"})
# Fire the wake from a plain OS thread while the coroutine is parked;
# the bridge is loop.call_soon_threadsafe on the captured future.
def waker():
time.sleep(0.05)
n.wake({"k"})
t = threading.Thread(target=waker)
t.start()
result = await sub.wait(2.0)
t.join(1.0)
return result
assert asyncio.run(scenario()) is True
def test_mode_is_parked_within_the_window_then_poll():
"""A plain read after a park keeps reporting parked until the parked stamp
is older than parked_window_s; stream is the parked poll's until then,
the latest read's after."""
now = [0.0]
n = Notifier(
max_waiters_per_owner=4,
max_waiters_total=256,
wake_spread_s=8.0,
wake_spread_threshold=20,
parked_window_s=110.0,
clock=lambda: now[0],
)
n.note_poll("a", "updates", parked=True)
now[0] = 50.0
n.note_poll("a", "digest", parked=False)
assert n.last_poll("a")[1:3] == ("parked", "updates")
now[0] = 111.0
assert n.last_poll("a")[1:3] == ("poll", "digest")
def test_last_cursor_keeps_the_newest_of_sent_and_handed_out():
n = _notifier()
n.note_poll("a", "updates", parked=False, after="20260101-000000-000_x.md")
n.note_cursor("a", "20260102-000000-000_y.md")
n.note_poll("a", "updates", parked=False, after="20260101-000000-000_x.md")
n.note_cursor("a", None)
assert n.last_poll("a").last_cursor == "20260102-000000-000_y.md"