Persist alert cooldowns across restarts so deploys stop re-firing every zone
This commit is contained in:
parent
e02f919b27
commit
617d7c75fb
7 changed files with 147 additions and 11 deletions
1
.gitignore
vendored
1
.gitignore
vendored
|
|
@ -4,4 +4,5 @@ __pycache__/
|
|||
.env
|
||||
.schwab_token.json
|
||||
data/manual_lines.json
|
||||
data/alert_state.json
|
||||
artifacts/playwright/
|
||||
|
|
|
|||
24
README.md
24
README.md
|
|
@ -347,15 +347,23 @@ whole point of the phone push — and opening two tabs does not double-notify.
|
|||
`NTFY_TOPIC` must be set or nothing sends; `send_ntfy` returns immediately on a blank
|
||||
topic. Set it in **Coolify's environment variables** for production, not in this repo.
|
||||
|
||||
**Use a different topic locally — or better, none.** Cooldown state is in memory, so
|
||||
every restart starts with empty cooldowns and the first closed bar re-alerts whatever
|
||||
zone price is sitting on. Locally that means every `--reload` save. Leaving
|
||||
`NTFY_TOPIC` blank keeps the in-browser sound and banner while suppressing the push;
|
||||
set a `-dev` topic only while testing the push path itself.
|
||||
**Cooldown state survives a restart.** The fired-zone table is written to
|
||||
`ALERT_STATE_PATH` (`./data/alert_state.json`), which in production is the same
|
||||
persistent volume as the trendlines. Before that, every deploy started with empty
|
||||
cooldowns and the next closed bar re-alerted whatever zone price was sitting on —
|
||||
with a four-hour cooldown, each push produced a burst of notifications for zones
|
||||
that had already had their say.
|
||||
|
||||
The same applies to production, more slowly: **a deploy resets the cooldowns**, so a
|
||||
zone that alerted an hour ago can alert again right after a redeploy. Persisting the
|
||||
fired-zone table would fix it.
|
||||
Two consequences worth knowing. The file has to be on the volume, or the problem
|
||||
comes straight back on the next deploy. And a corrupt or unreadable state file is
|
||||
deliberately non-fatal: it logs and starts empty, costing one burst of duplicate
|
||||
alerts rather than refusing to start the stream.
|
||||
|
||||
**Still prefer a blank topic locally.** Persistence removes the restart bursts, but
|
||||
an in-memory engine is only half the story — a dev instance watching the same
|
||||
symbol will happily push real alerts to your phone. Leaving `NTFY_TOPIC` blank keeps
|
||||
the in-browser sound and banner while suppressing the push; set a `-dev` topic only
|
||||
while testing the push path itself.
|
||||
|
||||
Note that ntfy topics are public by default: anyone who knows the name can both read
|
||||
your alerts and publish to it. Treat the topic name as a secret.
|
||||
|
|
|
|||
|
|
@ -1,8 +1,13 @@
|
|||
import json
|
||||
import logging
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
|
||||
from app.analysis.confluence import Cluster
|
||||
from app.analysis.levels import LevelKind
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class Alert:
|
||||
|
|
@ -27,12 +32,55 @@ class AlertEngine:
|
|||
cluster's identity while a human still sees one zone sitting at the prior
|
||||
day's close. Keying on identity let every reshuffle through as a fresh
|
||||
alert; keying on where the zone *is* does not.
|
||||
|
||||
Suppression is persisted when given a ``state_path``. Without it the list
|
||||
lives only in memory, so every restart re-fires every zone that currently
|
||||
qualifies — with a four-hour cooldown that turned each deploy into a burst
|
||||
of pushes for zones that had already had their say.
|
||||
"""
|
||||
|
||||
def __init__(self, min_score: float, cooldown_seconds: int = 900):
|
||||
def __init__(
|
||||
self,
|
||||
min_score: float,
|
||||
cooldown_seconds: int = 900,
|
||||
state_path: Path | None = None,
|
||||
):
|
||||
self.min_score = min_score
|
||||
self.cooldown_seconds = cooldown_seconds
|
||||
self._fired: list[_Fired] = []
|
||||
self.state_path = Path(state_path) if state_path else None
|
||||
self._fired: list[_Fired] = self._load()
|
||||
|
||||
def _load(self) -> list[_Fired]:
|
||||
if not self.state_path or not self.state_path.exists():
|
||||
return []
|
||||
try:
|
||||
payload = json.loads(self.state_path.read_text(encoding="utf-8"))
|
||||
return [_Fired(float(item["center"]), int(item["at"])) for item in payload]
|
||||
except Exception:
|
||||
# Corrupt state costs one burst of duplicate alerts, which is a far
|
||||
# better failure than refusing to start the stream.
|
||||
logger.warning("Could not read alert state; starting empty", exc_info=True)
|
||||
return []
|
||||
|
||||
def _save(self) -> None:
|
||||
if not self.state_path:
|
||||
return
|
||||
try:
|
||||
self.state_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
temporary = self.state_path.with_suffix(self.state_path.suffix + ".tmp")
|
||||
temporary.write_text(
|
||||
json.dumps(
|
||||
[{"center": entry.center, "at": entry.at} for entry in self._fired],
|
||||
indent=2,
|
||||
sort_keys=True,
|
||||
)
|
||||
+ "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
temporary.replace(self.state_path)
|
||||
except Exception:
|
||||
# Losing a write means duplicate alerts later, never a missed one.
|
||||
logger.warning("Could not persist alert state", exc_info=True)
|
||||
|
||||
def evaluate(
|
||||
self,
|
||||
|
|
@ -51,6 +99,7 @@ class AlertEngine:
|
|||
|
||||
# Re-arming needs both elapsed time and real separation. Time alone lets
|
||||
# price oscillating on a level alert forever.
|
||||
before = len(self._fired)
|
||||
self._fired = [
|
||||
entry
|
||||
for entry in self._fired
|
||||
|
|
@ -59,6 +108,7 @@ class AlertEngine:
|
|||
and abs(entry.center - current_price) > merge_distance
|
||||
)
|
||||
]
|
||||
changed = len(self._fired) != before
|
||||
|
||||
alerts: list[Alert] = []
|
||||
# Strongest first, so when several overlapping zones qualify at once the
|
||||
|
|
@ -87,6 +137,7 @@ class AlertEngine:
|
|||
):
|
||||
continue
|
||||
self._fired.append(_Fired(cluster.center, now))
|
||||
changed = True
|
||||
direction = "BEARISH" if cluster.side.value == "resistance" else "BULLISH"
|
||||
timeframes = ", ".join(dict.fromkeys(member.tf.value for member in cluster.members))
|
||||
# Naming the line matters: "your line" is actionable in a way that
|
||||
|
|
@ -102,4 +153,6 @@ class AlertEngine:
|
|||
f"{direction} {headline} {symbol} {current_price:.2f}\n{detail}\n{timeframes}"
|
||||
)
|
||||
alerts.append(Alert(cluster, message, tuple(member.id for member in drawn)))
|
||||
if changed:
|
||||
self._save()
|
||||
return alerts
|
||||
|
|
|
|||
|
|
@ -54,6 +54,9 @@ class Settings(BaseSettings):
|
|||
# per price zone, so an unrelated zone still alerts immediately; this only
|
||||
# governs how often the *same* area repeats itself.
|
||||
alert_cooldown_seconds: int = 14400
|
||||
# On the persistent volume in production: suppression has to outlive a
|
||||
# deploy or every push re-fires every zone that currently qualifies.
|
||||
alert_state_path: Path = Path("./data/alert_state.json")
|
||||
ntfy_topic: str = ""
|
||||
ntfy_server: str = "https://ntfy.sh"
|
||||
chart_auth_token: str = ""
|
||||
|
|
|
|||
|
|
@ -48,7 +48,9 @@ class Runtime:
|
|||
# are only meaningful if they outlive a page reload, and a phone push
|
||||
# must not depend on a tab being open to produce it.
|
||||
self.alert_engine = AlertEngine(
|
||||
self.settings.confluence_min_score, self.settings.alert_cooldown_seconds
|
||||
self.settings.confluence_min_score,
|
||||
self.settings.alert_cooldown_seconds,
|
||||
self.settings.alert_state_path,
|
||||
)
|
||||
self.levels = self.manual_lines.levels()
|
||||
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
|
||||
|
|
|
|||
65
tests/test_alert_state.py
Normal file
65
tests/test_alert_state.py
Normal file
|
|
@ -0,0 +1,65 @@
|
|||
"""Suppression has to survive a restart, or every deploy re-alerts."""
|
||||
import json
|
||||
|
||||
from app.analysis.alerts import AlertEngine
|
||||
from app.analysis.confluence import cluster_levels
|
||||
from app.analysis.levels import Level, LevelKind, Side
|
||||
from app.bars.models import Timeframe
|
||||
|
||||
|
||||
def level(id_: str, price: float, weight: float):
|
||||
return Level(
|
||||
id_, LevelKind.MA, Timeframe.D1, Side.RESISTANCE, weight, 1, id_,
|
||||
100, price, 0, None, 0, 100, 100, False, False,
|
||||
)
|
||||
|
||||
|
||||
def zone(price: float = 100.0):
|
||||
return cluster_levels([level("a", price, 3), level("b", price + 0.1, 4)], 100, price, 1)
|
||||
|
||||
|
||||
def engine(tmp_path, cooldown=14400):
|
||||
return AlertEngine(6, cooldown, tmp_path / "alert_state.json")
|
||||
|
||||
|
||||
def test_fires_once_then_suppresses_within_the_process(tmp_path):
|
||||
one = engine(tmp_path)
|
||||
assert len(one.evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
||||
assert one.evaluate(zone(), 100, 1, 60, "/ES") == []
|
||||
|
||||
|
||||
def test_suppression_survives_a_restart(tmp_path):
|
||||
one = engine(tmp_path)
|
||||
assert len(one.evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
||||
|
||||
# A second engine over the same state file stands in for a redeploy.
|
||||
two = engine(tmp_path)
|
||||
assert two.evaluate(zone(), 100, 1, 60, "/ES") == []
|
||||
|
||||
|
||||
def test_without_a_state_path_a_restart_still_refires(tmp_path):
|
||||
"""Unchanged behaviour for local runs, which should not write files."""
|
||||
assert len(AlertEngine(6, 14400).evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
||||
assert len(AlertEngine(6, 14400).evaluate(zone(), 100, 1, 60, "/ES")) == 1
|
||||
|
||||
|
||||
def test_rearms_across_a_restart_after_cooldown_and_separation(tmp_path):
|
||||
one = engine(tmp_path, cooldown=900)
|
||||
assert len(one.evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
||||
|
||||
two = engine(tmp_path, cooldown=900)
|
||||
# Price genuinely left the zone, and the cooldown has elapsed.
|
||||
assert two.evaluate(cluster_levels([level("a", 100, 3)], 100, 103, 1), 103, 1, 902, "/ES") == []
|
||||
assert len(two.evaluate(zone(), 100, 1, 903, "/ES")) == 1
|
||||
|
||||
|
||||
def test_corrupt_state_does_not_prevent_alerting(tmp_path):
|
||||
(tmp_path / "alert_state.json").write_text("{not json", encoding="utf-8")
|
||||
assert len(engine(tmp_path).evaluate(zone(), 100, 1, 0, "/ES")) == 1
|
||||
|
||||
|
||||
def test_state_file_records_centre_and_time(tmp_path):
|
||||
engine(tmp_path).evaluate(zone(), 100, 1, 42, "/ES")
|
||||
payload = json.loads((tmp_path / "alert_state.json").read_text(encoding="utf-8"))
|
||||
assert len(payload) == 1
|
||||
assert payload[0]["at"] == 42
|
||||
|
|
@ -14,6 +14,10 @@ from app.runtime import Runtime
|
|||
def runtime(tmp_path, **overrides) -> Runtime:
|
||||
settings = Settings(
|
||||
manual_lines_path=tmp_path / "manual_lines.json",
|
||||
# Isolated per test: the default is relative to the working directory,
|
||||
# so without this every test shares one alert-suppression file and they
|
||||
# silence each other.
|
||||
alert_state_path=tmp_path / "alert_state.json",
|
||||
ntfy_topic=overrides.pop("ntfy_topic", ""),
|
||||
**overrides,
|
||||
)
|
||||
|
|
|
|||
Loading…
Reference in a new issue