diff --git a/.gitignore b/.gitignore index 23240c8..7d2ef95 100644 --- a/.gitignore +++ b/.gitignore @@ -4,4 +4,5 @@ __pycache__/ .env .schwab_token.json data/manual_lines.json +data/alert_state.json artifacts/playwright/ diff --git a/README.md b/README.md index 72a29e4..fcb94f6 100644 --- a/README.md +++ b/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. diff --git a/app/analysis/alerts.py b/app/analysis/alerts.py index dee87f1..ee872c4 100644 --- a/app/analysis/alerts.py +++ b/app/analysis/alerts.py @@ -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 diff --git a/app/config.py b/app/config.py index c418604..3a0402c 100644 --- a/app/config.py +++ b/app/config.py @@ -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 = "" diff --git a/app/runtime.py b/app/runtime.py index e88096d..582b520 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -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) diff --git a/tests/test_alert_state.py b/tests/test_alert_state.py new file mode 100644 index 0000000..af0b99f --- /dev/null +++ b/tests/test_alert_state.py @@ -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 diff --git a/tests/test_runtime_alerts.py b/tests/test_runtime_alerts.py index 7f32a5b..cff6eca 100644 --- a/tests/test_runtime_alerts.py +++ b/tests/test_runtime_alerts.py @@ -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, )