51 lines
2.3 KiB
Python
51 lines
2.3 KiB
Python
|
|
"""Single-instance leader election for in-process background jobs.
|
||
|
|
|
||
|
|
With multiple uvicorn workers each worker process runs its own app lifespan, so a
|
||
|
|
background thread started there runs once *per worker*. That is correct for the
|
||
|
|
neighbor warmer (each worker drains its own request-fed queue) but wrong for the
|
||
|
|
subscription notifier: its periodic sweep scans every subscription and issues
|
||
|
|
upstream forecast/archive fetches on a timer, independent of any request. Running
|
||
|
|
it in N workers multiplies the Open-Meteo quota use N-fold (the source of the
|
||
|
|
overnight 429/503 rate-limit errors after prod went to 3 workers) and does the
|
||
|
|
same DB scan N times.
|
||
|
|
|
||
|
|
``claim`` elects one worker to own such jobs via a non-blocking exclusive advisory
|
||
|
|
lock (``flock``) on a lockfile:
|
||
|
|
|
||
|
|
* the first worker to call it acquires the lock and holds the fd for the process
|
||
|
|
lifetime, so the lock is never dropped while it lives -> returns ``True``;
|
||
|
|
* other workers fail the non-blocking lock and stand down -> return ``False``;
|
||
|
|
* the OS releases the lock when the leader exits, so a restart (or a crash) lets a
|
||
|
|
surviving/new worker win on its next call -> clean re-election, no stale state.
|
||
|
|
|
||
|
|
No lockfile configured (single-worker deploy, dev, tests) => always the leader.
|
||
|
|
This mirrors the metrics store's env-selected coordination: unset keeps the
|
||
|
|
zero-dependency single-worker behavior; a path opts into cross-worker arbitration.
|
||
|
|
"""
|
||
|
|
import fcntl
|
||
|
|
import os
|
||
|
|
|
||
|
|
# lock_path -> held fd. Kept open for the whole process lifetime: closing the fd (or
|
||
|
|
# letting it be GC'd) would drop the flock and let a second worker also become leader.
|
||
|
|
_held: dict[str, int] = {}
|
||
|
|
|
||
|
|
|
||
|
|
def claim(lock_path: "str | None") -> bool:
|
||
|
|
"""Return True if this process should own the single-instance job guarded by
|
||
|
|
``lock_path``. ``None``/empty (no multi-worker coordination configured) is always
|
||
|
|
the leader. Idempotent: re-claiming a path this process already holds returns True
|
||
|
|
without touching the lock."""
|
||
|
|
if not lock_path:
|
||
|
|
return True
|
||
|
|
if lock_path in _held:
|
||
|
|
return True
|
||
|
|
fd = os.open(lock_path, os.O_RDWR | os.O_CREAT, 0o644)
|
||
|
|
try:
|
||
|
|
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||
|
|
except OSError:
|
||
|
|
# Another process already holds it — we're not the leader.
|
||
|
|
os.close(fd)
|
||
|
|
return False
|
||
|
|
_held[lock_path] = fd
|
||
|
|
return True
|