Run push notification channels on independent fixed cadences #714
|
|
@ -210,6 +210,8 @@ export STACKCHAIN_VAPID_SUBJECT='mailto:ops@example.com'
|
||||||
# and STACKCHAIN_STATE_DIR/push-subscriptions.sqlite3. Threads stay ordered
|
# and STACKCHAIN_STATE_DIR/push-subscriptions.sqlite3. Threads stay ordered
|
||||||
# within each device; one slow device does not delay healthy devices behind it.
|
# within each device; one slow device does not delay healthy devices behind it.
|
||||||
export STACKCHAIN_PUSH_POLL_SECONDS=30
|
export STACKCHAIN_PUSH_POLL_SECONDS=30
|
||||||
|
# Unread updates and deadline reminders run on independent, fixed-cadence
|
||||||
|
# workers, so a slow channel cannot delay the other or add drift to its ticks.
|
||||||
export STACKCHAIN_PUSH_SEND_TIMEOUT_SECONDS=10
|
export STACKCHAIN_PUSH_SEND_TIMEOUT_SECONDS=10
|
||||||
export STACKCHAIN_PUSH_MAX_CONCURRENCY=8
|
export STACKCHAIN_PUSH_MAX_CONCURRENCY=8
|
||||||
# Maximum individual alerts per device and poll before one digest covers the rest.
|
# Maximum individual alerts per device and poll before one digest covers the rest.
|
||||||
|
|
|
||||||
52
src/main.py
52
src/main.py
|
|
@ -87,6 +87,23 @@ async def _drain_authored_action_operations() -> None:
|
||||||
_authored_action_operations.pop(key, None)
|
_authored_action_operations.pop(key, None)
|
||||||
|
|
||||||
|
|
||||||
|
async def _push_channel_loop(dispatch, *, interval: float) -> None:
|
||||||
|
"""Run one push channel on fixed ticks without overlapping or catch-up bursts."""
|
||||||
|
loop = asyncio.get_running_loop()
|
||||||
|
next_tick = loop.time() + interval
|
||||||
|
while True:
|
||||||
|
await asyncio.sleep(max(0.0, next_tick - loop.time()))
|
||||||
|
try:
|
||||||
|
await dispatch()
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
raise
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
now = loop.time()
|
||||||
|
elapsed_intervals = max(1, int((now - next_tick) // interval) + 1)
|
||||||
|
next_tick += elapsed_intervals * interval
|
||||||
|
|
||||||
|
|
||||||
async def _push_poll_loop() -> None:
|
async def _push_poll_loop() -> None:
|
||||||
interval = max(5.0, float(os.getenv("STACKCHAIN_PUSH_POLL_SECONDS", "30")))
|
interval = max(5.0, float(os.getenv("STACKCHAIN_PUSH_POLL_SECONDS", "30")))
|
||||||
send_timeout = max(
|
send_timeout = max(
|
||||||
|
|
@ -103,10 +120,8 @@ async def _push_poll_loop() -> None:
|
||||||
0,
|
0,
|
||||||
int(os.getenv("STACKCHAIN_PUSH_MAX_INDIVIDUAL_NOTIFICATIONS", "3")),
|
int(os.getenv("STACKCHAIN_PUSH_MAX_INDIVIDUAL_NOTIFICATIONS", "3")),
|
||||||
)
|
)
|
||||||
while True:
|
async def dispatch_unread() -> None:
|
||||||
await asyncio.sleep(interval)
|
await dispatch_unread_updates(
|
||||||
try:
|
|
||||||
await dispatch_unread_updates(
|
|
||||||
_push_subscription_store,
|
_push_subscription_store,
|
||||||
_push_configuration(),
|
_push_configuration(),
|
||||||
gitea_proxy.unread_notification_snapshot,
|
gitea_proxy.unread_notification_snapshot,
|
||||||
|
|
@ -116,12 +131,9 @@ async def _push_poll_loop() -> None:
|
||||||
max_concurrency=max_concurrency,
|
max_concurrency=max_concurrency,
|
||||||
max_individual_notifications=max_individual_notifications,
|
max_individual_notifications=max_individual_notifications,
|
||||||
)
|
)
|
||||||
except asyncio.CancelledError:
|
|
||||||
raise
|
async def dispatch_deadlines() -> None:
|
||||||
except Exception:
|
await dispatch_deadline_reminders(
|
||||||
pass
|
|
||||||
try:
|
|
||||||
await dispatch_deadline_reminders(
|
|
||||||
_push_subscription_store,
|
_push_subscription_store,
|
||||||
_push_configuration(),
|
_push_configuration(),
|
||||||
gitea_proxy.assigned_issue_snapshot,
|
gitea_proxy.assigned_issue_snapshot,
|
||||||
|
|
@ -130,11 +142,21 @@ async def _push_poll_loop() -> None:
|
||||||
lease_seconds=lease_seconds,
|
lease_seconds=lease_seconds,
|
||||||
max_concurrency=max_concurrency,
|
max_concurrency=max_concurrency,
|
||||||
)
|
)
|
||||||
except asyncio.CancelledError:
|
|
||||||
raise
|
channel_tasks = (
|
||||||
except Exception:
|
asyncio.create_task(_push_channel_loop(dispatch_unread, interval=interval)),
|
||||||
# Gitea and push endpoints retry independently on the next poll.
|
asyncio.create_task(_push_channel_loop(dispatch_deadlines, interval=interval)),
|
||||||
continue
|
)
|
||||||
|
try:
|
||||||
|
done, _pending = await asyncio.wait(
|
||||||
|
channel_tasks, return_when=asyncio.FIRST_COMPLETED
|
||||||
|
)
|
||||||
|
await next(iter(done))
|
||||||
|
finally:
|
||||||
|
for task in channel_tasks:
|
||||||
|
if not task.done():
|
||||||
|
task.cancel()
|
||||||
|
await asyncio.gather(*channel_tasks, return_exceptions=True)
|
||||||
|
|
||||||
|
|
||||||
async def _readiness_monitor() -> None:
|
async def _readiness_monitor() -> None:
|
||||||
|
|
|
||||||
|
|
@ -235,11 +235,16 @@ async def test_push_poll_still_dispatches_deadlines_when_unread_dispatch_fails(m
|
||||||
calls = []
|
calls = []
|
||||||
|
|
||||||
sleeps = 0
|
sleeps = 0
|
||||||
|
first_tick = asyncio.Event()
|
||||||
|
|
||||||
async def no_wait(_seconds):
|
async def no_wait(_seconds):
|
||||||
nonlocal sleeps
|
nonlocal sleeps
|
||||||
sleeps += 1
|
sleeps += 1
|
||||||
if sleeps > 1:
|
if sleeps <= 2:
|
||||||
|
if sleeps == 2:
|
||||||
|
first_tick.set()
|
||||||
|
await first_tick.wait()
|
||||||
|
else:
|
||||||
raise asyncio.CancelledError
|
raise asyncio.CancelledError
|
||||||
|
|
||||||
async def fail_unread(*_args, **_kwargs):
|
async def fail_unread(*_args, **_kwargs):
|
||||||
|
|
@ -256,7 +261,71 @@ async def test_push_poll_still_dispatches_deadlines_when_unread_dispatch_fails(m
|
||||||
with pytest.raises(asyncio.CancelledError):
|
with pytest.raises(asyncio.CancelledError):
|
||||||
await main._push_poll_loop()
|
await main._push_poll_loop()
|
||||||
|
|
||||||
assert calls == ["unread", "deadline"]
|
assert sorted(calls) == ["deadline", "unread"]
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_push_poll_deadlines_continue_while_unread_dispatch_is_blocked(monkeypatch):
|
||||||
|
unread_started = asyncio.Event()
|
||||||
|
release_unread = asyncio.Event()
|
||||||
|
deadline_started = asyncio.Event()
|
||||||
|
|
||||||
|
async def no_wait(_seconds):
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def blocked_unread(*_args, **_kwargs):
|
||||||
|
unread_started.set()
|
||||||
|
await release_unread.wait()
|
||||||
|
|
||||||
|
async def dispatch_deadlines(*_args, **_kwargs):
|
||||||
|
deadline_started.set()
|
||||||
|
await release_unread.wait()
|
||||||
|
|
||||||
|
monkeypatch.setattr(main.asyncio, "sleep", no_wait)
|
||||||
|
monkeypatch.setattr(main, "dispatch_unread_updates", blocked_unread)
|
||||||
|
monkeypatch.setattr(main, "dispatch_deadline_reminders", dispatch_deadlines)
|
||||||
|
|
||||||
|
poll = asyncio.create_task(main._push_poll_loop())
|
||||||
|
await asyncio.wait_for(unread_started.wait(), timeout=0.5)
|
||||||
|
try:
|
||||||
|
await asyncio.wait_for(deadline_started.wait(), timeout=0.5)
|
||||||
|
finally:
|
||||||
|
poll.cancel()
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await poll
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.anyio
|
||||||
|
async def test_push_channel_loop_skips_missed_ticks_without_overlapping(monkeypatch):
|
||||||
|
starts = []
|
||||||
|
now = 0.0
|
||||||
|
|
||||||
|
class Loop:
|
||||||
|
def time(self):
|
||||||
|
return now
|
||||||
|
|
||||||
|
async def record_dispatch():
|
||||||
|
nonlocal now
|
||||||
|
starts.append("dispatch")
|
||||||
|
now = 16.0 if len(starts) == 1 else now
|
||||||
|
if len(starts) == 2:
|
||||||
|
raise asyncio.CancelledError
|
||||||
|
|
||||||
|
sleeps = []
|
||||||
|
|
||||||
|
async def capture_sleep(seconds):
|
||||||
|
nonlocal now
|
||||||
|
sleeps.append(seconds)
|
||||||
|
now += seconds
|
||||||
|
|
||||||
|
monkeypatch.setattr(main.asyncio, "get_running_loop", lambda: Loop())
|
||||||
|
monkeypatch.setattr(main.asyncio, "sleep", capture_sleep)
|
||||||
|
|
||||||
|
with pytest.raises(asyncio.CancelledError):
|
||||||
|
await main._push_channel_loop(record_dispatch, interval=5.0)
|
||||||
|
|
||||||
|
assert starts == ["dispatch", "dispatch"]
|
||||||
|
assert sleeps == [5.0, 4.0]
|
||||||
|
|
||||||
|
|
||||||
def test_dispatch_leases_are_isolated_by_channel(tmp_path):
|
def test_dispatch_leases_are_isolated_by_channel(tmp_path):
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue
Block a user