feat(lifecycle): alive doar la deschidere/închidere, /window persistent, /status cu fereastra
Heartbeat-ul periodic (activ/IDLE la 30 min) e oprit implicit (heartbeat_min = 0). Semnalul „alive" devin alertele zilnice „🟢 Piața deschisă" / „🔴 Piața închisă", emise pe tranziția ferestrei (operating_hours AND /window) independent de pauză — body-ul spune dacă monitorizarea e pauzată manual sau pe drift. /window era doar în memorie și se pierdea la repornire (09-09), așa că programul revenea la 16:30–23:00 și trimitea mesaje înainte/după fereastra 19:30–22:00. Acum e persistat în logs/session_window.json și restaurat la pornire; /window off șterge fișierul. /status afișează blocul de fereastră: deschisă/închisă + următoarea deschidere/închidere (ora locală), /window, orele bursei cu echivalentul local, și config | heartbeat. Config activ: heartbeat_min = 0; include și baseline_phash rescris de auto-rebase-ul canary din 09-09. DOX pass: notifier + configs. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Em4jtmh4pRpYMdyd9ueHYF
This commit is contained in:
1
.gitignore
vendored
1
.gitignore
vendored
@@ -49,6 +49,7 @@ logs/dead_letter.jsonl
|
|||||||
logs/detections/
|
logs/detections/
|
||||||
logs/fires
|
logs/fires
|
||||||
logs/pause.flag
|
logs/pause.flag
|
||||||
|
logs/session_window.json
|
||||||
samples/*.png
|
samples/*.png
|
||||||
samples/*.jpg
|
samples/*.jpg
|
||||||
samples/labels.json
|
samples/labels.json
|
||||||
|
|||||||
@@ -53,7 +53,7 @@ p2_y = 664
|
|||||||
p2_price = 483.2
|
p2_price = 483.2
|
||||||
|
|
||||||
[canary]
|
[canary]
|
||||||
baseline_phash = "c11f4a852ec09f3a8de4e4cf4ad76d84f10b19d3e708663c38f5b538877c6624"
|
baseline_phash = "fbe145390c1abec23204017757a326b8e37077288ef79947310a89c70e07ffff"
|
||||||
drift_threshold = 8
|
drift_threshold = 8
|
||||||
|
|
||||||
[canary.roi]
|
[canary.roi]
|
||||||
@@ -73,7 +73,7 @@ h = 1029
|
|||||||
[options]
|
[options]
|
||||||
debounce_depth = 1
|
debounce_depth = 1
|
||||||
loop_interval_s = 5.0
|
loop_interval_s = 5.0
|
||||||
heartbeat_min = 30
|
heartbeat_min = 0 # 0 = fără heartbeat periodic; alive = mesajele de deschidere/închidere
|
||||||
lockout_s = 240
|
lockout_s = 240
|
||||||
low_conf_threshold = 0.2
|
low_conf_threshold = 0.2
|
||||||
low_conf_run = 3
|
low_conf_run = 3
|
||||||
|
|||||||
@@ -23,8 +23,14 @@ alerts). Schema e definită și validată în `src/atm/config.py` (fail fast la
|
|||||||
|
|
||||||
`enabled`, `timezone` (NYSE local, ex. `America/New_York`), `weekdays`,
|
`enabled`, `timezone` (NYSE local, ex. `America/New_York`), `weekdays`,
|
||||||
`start_hhmm`, `stop_hhmm`. Timezone validat la load; `_tz_cache` reutilizat per
|
`start_hhmm`, `stop_hhmm`. Timezone validat la load; `_tz_cache` reutilizat per
|
||||||
tick. Boundary crossings logează `market_open`/`market_closed` și notifică o dată.
|
tick. Boundary crossings logează `market_open`/`market_closed` și notifică o dată
|
||||||
Startup in-window e silent.
|
(și când monitorizarea e pauzată — sunt semnalul zilnic „alive"). Startup
|
||||||
|
in-window e silent.
|
||||||
|
|
||||||
|
### `[options] heartbeat_min`
|
||||||
|
|
||||||
|
Default `0` = fără heartbeat periodic (alive = alertele de deschidere/închidere).
|
||||||
|
`N > 0` repornește heartbeat-ul la N minute, doar în fereastra de tranzacționare.
|
||||||
|
|
||||||
### `[options.alerts]`
|
### `[options.alerts]`
|
||||||
|
|
||||||
|
|||||||
@@ -70,7 +70,9 @@ h = 50
|
|||||||
[options]
|
[options]
|
||||||
debounce_depth = 1
|
debounce_depth = 1
|
||||||
loop_interval_s = 5.0
|
loop_interval_s = 5.0
|
||||||
heartbeat_min = 30
|
# Periodic heartbeat (minutes). 0 = off: the daily "Piața deschisă/închisă"
|
||||||
|
# alerts (sent even while paused) are the alive signal.
|
||||||
|
heartbeat_min = 0
|
||||||
lockout_s = 240
|
lockout_s = 240
|
||||||
low_conf_threshold = 0.2
|
low_conf_threshold = 0.2
|
||||||
low_conf_run = 3
|
low_conf_run = 3
|
||||||
|
|||||||
@@ -30,9 +30,9 @@ _BASE = "https://api.telegram.org/bot{token}/{method}"
|
|||||||
# command must match [a-z0-9_]{1,32}; description is 1-256 chars. Telegram has
|
# command must match [a-z0-9_]{1,32}; description is 1-256 chars. Telegram has
|
||||||
# no typed args, so arg formats live in the description text.
|
# no typed args, so arg formats live in the description text.
|
||||||
TELEGRAM_COMMANDS: list[tuple[str, str]] = [
|
TELEGRAM_COMMANDS: list[tuple[str, str]] = [
|
||||||
("status", "Stare FSM, uptime, ultima detecție, fereastră open/closed"),
|
("status", "Stare FSM, uptime, fereastră (/window, ore bursă, următoarea deschidere/închidere)"),
|
||||||
("ss", "Screenshot acum (top-3 buline din ROI)"),
|
("ss", "Screenshot acum (top-3 buline din ROI)"),
|
||||||
("pause", "Suspendă detecția (heartbeat-urile continuă)"),
|
("pause", "Suspendă detecția (mesajele de deschidere/închidere continuă)"),
|
||||||
("resume", "Reia detecția (șterge user-pause + drift-pause)"),
|
("resume", "Reia detecția (șterge user-pause + drift-pause)"),
|
||||||
("rebase", "Propune phash nou pentru canary — aplici cu: /rebase confirm"),
|
("rebase", "Propune phash nou pentru canary — aplici cu: /rebase confirm"),
|
||||||
("interval", "Auto-screenshot la N minute, ex: /interval 3"),
|
("interval", "Auto-screenshot la N minute, ex: /interval 3"),
|
||||||
|
|||||||
@@ -248,7 +248,7 @@ class Config:
|
|||||||
chart_window_region: ROI | None = None # virtual-desktop absolute region; when set, runtime uses full-desktop capture + crop
|
chart_window_region: ROI | None = None # virtual-desktop absolute region; when set, runtime uses full-desktop capture + crop
|
||||||
debounce_depth: int = 1
|
debounce_depth: int = 1
|
||||||
loop_interval_s: float = 5.0
|
loop_interval_s: float = 5.0
|
||||||
heartbeat_min: int = 30
|
heartbeat_min: int = 0 # 0 = off; alive signal = daily open/close alerts
|
||||||
lockout_s: int = 240
|
lockout_s: int = 240
|
||||||
low_conf_threshold: float = 0.2
|
low_conf_threshold: float = 0.2
|
||||||
low_conf_run: int = 3
|
low_conf_run: int = 3
|
||||||
@@ -369,7 +369,7 @@ class Config:
|
|||||||
chart_window_region=region,
|
chart_window_region=region,
|
||||||
debounce_depth=int(opts.get("debounce_depth", 1)),
|
debounce_depth=int(opts.get("debounce_depth", 1)),
|
||||||
loop_interval_s=float(opts.get("loop_interval_s", 5.0)),
|
loop_interval_s=float(opts.get("loop_interval_s", 5.0)),
|
||||||
heartbeat_min=int(opts.get("heartbeat_min", 30)),
|
heartbeat_min=int(opts.get("heartbeat_min", 0)),
|
||||||
lockout_s=int(opts.get("lockout_s", 240)),
|
lockout_s=int(opts.get("lockout_s", 240)),
|
||||||
low_conf_threshold=float(opts.get("low_conf_threshold", 0.2)),
|
low_conf_threshold=float(opts.get("low_conf_threshold", 0.2)),
|
||||||
low_conf_run=int(opts.get("low_conf_run", 3)),
|
low_conf_run=int(opts.get("low_conf_run", 3)),
|
||||||
|
|||||||
215
src/atm/main.py
215
src/atm/main.py
@@ -4,6 +4,7 @@ from __future__ import annotations
|
|||||||
import argparse
|
import argparse
|
||||||
import asyncio
|
import asyncio
|
||||||
import contextlib
|
import contextlib
|
||||||
|
import json
|
||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
import time
|
import time
|
||||||
@@ -1176,6 +1177,30 @@ class LifecycleState:
|
|||||||
# Local wall-clock (start_hhmm, stop_hhmm) set via /window; None = no window.
|
# Local wall-clock (start_hhmm, stop_hhmm) set via /window; None = no window.
|
||||||
# ANDed with operating_hours: detection runs only when inside BOTH.
|
# ANDed with operating_hours: detection runs only when inside BOTH.
|
||||||
session_window: tuple[str, str] | None = None
|
session_window: tuple[str, str] | None = None
|
||||||
|
# Where /window is persisted so it survives restarts. None = in-memory only
|
||||||
|
# (unit tests); run_live_async points it at logs/session_window.json.
|
||||||
|
session_window_path: Path | None = None
|
||||||
|
|
||||||
|
|
||||||
|
def _load_session_window(path: Path) -> tuple[str, str] | None:
|
||||||
|
"""Read a persisted /window; None when missing or unreadable."""
|
||||||
|
try:
|
||||||
|
data = json.loads(path.read_text(encoding="utf-8"))
|
||||||
|
start, stop = str(data["start"]), str(data["stop"])
|
||||||
|
except (OSError, ValueError, KeyError, TypeError):
|
||||||
|
return None
|
||||||
|
return (start, stop)
|
||||||
|
|
||||||
|
|
||||||
|
def _save_session_window(path: Path | None, window: tuple[str, str] | None) -> None:
|
||||||
|
"""Persist /window (or remove the file on /window off). No-op when path is None."""
|
||||||
|
if path is None:
|
||||||
|
return
|
||||||
|
if window is None:
|
||||||
|
path.unlink(missing_ok=True)
|
||||||
|
return
|
||||||
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
path.write_text(json.dumps({"start": window[0], "stop": window[1]}), encoding="utf-8")
|
||||||
|
|
||||||
|
|
||||||
# Locale-independent weekday names; index matches datetime.weekday() (MON=0).
|
# Locale-independent weekday names; index matches datetime.weekday() (MON=0).
|
||||||
@@ -1210,6 +1235,15 @@ def _should_skip(now_ts: float, state: LifecycleState, cfg, canary) -> str | Non
|
|||||||
return "user_paused"
|
return "user_paused"
|
||||||
if getattr(canary, "is_paused", False):
|
if getattr(canary, "is_paused", False):
|
||||||
return "drift_paused"
|
return "drift_paused"
|
||||||
|
return _window_skip_reason(now_ts, state, cfg)
|
||||||
|
|
||||||
|
|
||||||
|
def _window_skip_reason(now_ts: float, state: LifecycleState, cfg) -> str | None:
|
||||||
|
"""Market-window part of _should_skip: operating_hours AND session window.
|
||||||
|
|
||||||
|
Ignores user/drift pauses on purpose — the open/close ("alive") alerts
|
||||||
|
follow the trading window even while detection is paused.
|
||||||
|
"""
|
||||||
oh = getattr(cfg, "operating_hours", None)
|
oh = getattr(cfg, "operating_hours", None)
|
||||||
if oh is not None and oh.enabled:
|
if oh is not None and oh.enabled:
|
||||||
tz = getattr(oh, "_tz_cache", None)
|
tz = getattr(oh, "_tz_cache", None)
|
||||||
@@ -1232,6 +1266,71 @@ def _should_skip(now_ts: float, state: LifecycleState, cfg, canary) -> str | Non
|
|||||||
return None
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
_WEEKDAY_RO: tuple[str, ...] = ("Lu", "Ma", "Mi", "Jo", "Vi", "Sâ", "Du")
|
||||||
|
|
||||||
|
|
||||||
|
def _next_window_change(
|
||||||
|
now_ts: float, state: LifecycleState, cfg, horizon_days: int = 8,
|
||||||
|
) -> float | None:
|
||||||
|
"""First minute boundary after now_ts where the window flips open↔closed.
|
||||||
|
|
||||||
|
Brute-force minute scan (boundaries are HH:MM, so minute resolution is
|
||||||
|
exact); ~11k cheap evaluations, only run on /status. None = no flip within
|
||||||
|
the horizon (e.g. no window configured → always open).
|
||||||
|
"""
|
||||||
|
is_open = _window_skip_reason(now_ts, state, cfg) is None
|
||||||
|
t = (int(now_ts) // 60 + 1) * 60
|
||||||
|
end = now_ts + horizon_days * 86400
|
||||||
|
while t <= end:
|
||||||
|
if (_window_skip_reason(t, state, cfg) is None) != is_open:
|
||||||
|
return float(t)
|
||||||
|
t += 60
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _fmt_local_when(ts: float, now_ts: float) -> str:
|
||||||
|
"""'azi 22:00' / 'mâine 19:30' / 'Lu 19:30' in local wall-clock."""
|
||||||
|
d = datetime.fromtimestamp(ts)
|
||||||
|
days = (d.date() - datetime.fromtimestamp(now_ts).date()).days
|
||||||
|
day = "azi" if days == 0 else "mâine" if days == 1 else _WEEKDAY_RO[d.weekday()]
|
||||||
|
return f"{day} {d:%H:%M}"
|
||||||
|
|
||||||
|
|
||||||
|
def _window_status_lines(now_ts: float, state: LifecycleState, cfg) -> list[str]:
|
||||||
|
"""/status block: window open/closed + next change, /window, exchange hours.
|
||||||
|
|
||||||
|
Empty when neither operating_hours nor /window is configured (always open).
|
||||||
|
"""
|
||||||
|
oh = getattr(cfg, "operating_hours", None)
|
||||||
|
oh_on = oh is not None and oh.enabled and isinstance(getattr(oh, "_tz_cache", None), tzinfo)
|
||||||
|
sw = state.session_window
|
||||||
|
if not oh_on and sw is None:
|
||||||
|
return []
|
||||||
|
|
||||||
|
is_open = _window_skip_reason(now_ts, state, cfg) is None
|
||||||
|
head = f"fereastră: {'deschisă' if is_open else 'închisă'}"
|
||||||
|
nxt = _next_window_change(now_ts, state, cfg)
|
||||||
|
if nxt is not None:
|
||||||
|
head += f" → {'închidere' if is_open else 'deschidere'} {_fmt_local_when(nxt, now_ts)}"
|
||||||
|
lines = [head, f"/window: {sw[0]}–{sw[1]} (local)" if sw else "/window: —"]
|
||||||
|
|
||||||
|
if oh_on:
|
||||||
|
tz = oh._tz_cache
|
||||||
|
now_ex = datetime.fromtimestamp(now_ts, tz=tz)
|
||||||
|
|
||||||
|
def _to_local(hhmm: str) -> str:
|
||||||
|
h, m = (int(p) for p in hhmm.split(":"))
|
||||||
|
return now_ex.replace(hour=h, minute=m, second=0, microsecond=0).astimezone().strftime("%H:%M")
|
||||||
|
|
||||||
|
days = tuple(oh.weekdays)
|
||||||
|
days_txt = "MON–FRI" if days == _WEEKDAY_NAMES[:5] else ",".join(days)
|
||||||
|
lines.append(
|
||||||
|
f"bursă: {oh.start_hhmm}–{oh.stop_hhmm} {oh.timezone} {days_txt} "
|
||||||
|
f"(local {_to_local(oh.start_hhmm)}–{_to_local(oh.stop_hhmm)})"
|
||||||
|
)
|
||||||
|
return lines
|
||||||
|
|
||||||
|
|
||||||
def _maybe_log_transition(
|
def _maybe_log_transition(
|
||||||
reason: str | None,
|
reason: str | None,
|
||||||
state: LifecycleState,
|
state: LifecycleState,
|
||||||
@@ -1267,15 +1366,17 @@ def _maybe_log_transition(
|
|||||||
|
|
||||||
event_name = "market_open" if window_reason == "open" else "market_closed"
|
event_name = "market_open" if window_reason == "open" else "market_closed"
|
||||||
audit.log({"ts": now, "event": event_name, "reason": reason})
|
audit.log({"ts": now, "event": event_name, "reason": reason})
|
||||||
|
# These two alerts double as the daily "alive" signal (the periodic
|
||||||
|
# heartbeat is off by default), so they fire even while paused.
|
||||||
if event_name == "market_closed":
|
if event_name == "market_closed":
|
||||||
body = "Piața închisă — monitorizare pauzată până la următoarea deschidere."
|
body = "Sfârșit sesiune — ATM activ, monitorizare pauzată până la următoarea deschidere."
|
||||||
else:
|
else:
|
||||||
body = "Piața deschisă — monitorizare reluată."
|
body = "Început sesiune — ATM activ."
|
||||||
if status_body:
|
if status_body:
|
||||||
body = f"{body}\n{status_body}"
|
body = f"{body}\n{status_body}"
|
||||||
notifier.send(Alert(
|
notifier.send(Alert(
|
||||||
kind="status",
|
kind="status",
|
||||||
title="Piața deschisă" if event_name == "market_open" else "Piața închisă",
|
title="🟢 Piața deschisă" if event_name == "market_open" else "🔴 Piața închisă",
|
||||||
body=body,
|
body=body,
|
||||||
))
|
))
|
||||||
state.last_window_state = window_reason
|
state.last_window_state = window_reason
|
||||||
@@ -1403,6 +1504,40 @@ def _brief_status(ctx) -> str:
|
|||||||
return f"{fsm_state} | semnale: {ctx.state.fire_count} | {h:.1f}h"
|
return f"{fsm_state} | semnale: {ctx.state.fire_count} | {h:.1f}h"
|
||||||
|
|
||||||
|
|
||||||
|
def _session_status_body(ctx, skip: str | None) -> str:
|
||||||
|
"""Body tail for the open/close alerts: FSM brief + pause + /window info."""
|
||||||
|
lines = [_brief_status(ctx)]
|
||||||
|
if skip == "user_paused":
|
||||||
|
lines.append("⏸ monitorizare oprită manual — /resume pentru a relua")
|
||||||
|
elif skip == "drift_paused":
|
||||||
|
lines.append("⚠️ detecție pauzată (drift) — /resume pentru a relua")
|
||||||
|
sw = ctx.lifecycle.session_window if ctx.lifecycle is not None else None
|
||||||
|
if sw is not None:
|
||||||
|
lines.append(f"fereastră: {sw[0]}–{sw[1]} (ora locală)")
|
||||||
|
return "\n".join(lines)
|
||||||
|
|
||||||
|
|
||||||
|
async def _lifecycle_gate(ctx: RunContext, now: float) -> str | None:
|
||||||
|
"""Emit window open/close transitions, return the detection skip reason.
|
||||||
|
|
||||||
|
The transition follows the market window only (not user/drift pause), so
|
||||||
|
the daily open/close alerts arrive even while monitoring is paused.
|
||||||
|
"""
|
||||||
|
skip = _should_skip(now, ctx.lifecycle, ctx.cfg, ctx.canary)
|
||||||
|
transition = _maybe_log_transition(
|
||||||
|
_window_skip_reason(now, ctx.lifecycle, ctx.cfg),
|
||||||
|
ctx.lifecycle, now, ctx.audit, ctx.notifier,
|
||||||
|
status_body=_session_status_body(ctx, skip),
|
||||||
|
)
|
||||||
|
if transition == "market_open" and skip is None and ctx.cfg.window_title:
|
||||||
|
title = await asyncio.to_thread(_focus_window_by_title, ctx.cfg.window_title)
|
||||||
|
ctx.audit.log({"ts": now, "event": "window_focused", "command": "market_open", "title": title})
|
||||||
|
await asyncio.sleep(0.15)
|
||||||
|
elif transition == "market_closed":
|
||||||
|
_handle_market_closed(ctx, now)
|
||||||
|
return skip
|
||||||
|
|
||||||
|
|
||||||
async def _run_tick(ctx: RunContext) -> _TickSyncResult:
|
async def _run_tick(ctx: RunContext) -> _TickSyncResult:
|
||||||
"""Execute one `_sync_detection_tick` in a thread; returns result or empty.
|
"""Execute one `_sync_detection_tick` in a thread; returns result or empty.
|
||||||
|
|
||||||
@@ -1413,18 +1548,7 @@ async def _run_tick(ctx: RunContext) -> _TickSyncResult:
|
|||||||
"""
|
"""
|
||||||
now = time.time()
|
now = time.time()
|
||||||
if ctx.lifecycle is not None:
|
if ctx.lifecycle is not None:
|
||||||
skip = _should_skip(now, ctx.lifecycle, ctx.cfg, ctx.canary)
|
if await _lifecycle_gate(ctx, now) is not None:
|
||||||
sb = _brief_status(ctx)
|
|
||||||
transition = _maybe_log_transition(
|
|
||||||
skip, ctx.lifecycle, now, ctx.audit, ctx.notifier, status_body=sb,
|
|
||||||
)
|
|
||||||
if transition == "market_open" and ctx.cfg.window_title:
|
|
||||||
title = await asyncio.to_thread(_focus_window_by_title, ctx.cfg.window_title)
|
|
||||||
ctx.audit.log({"ts": now, "event": "window_focused", "command": "market_open", "title": title})
|
|
||||||
await asyncio.sleep(0.15)
|
|
||||||
elif transition == "market_closed":
|
|
||||||
_handle_market_closed(ctx, now)
|
|
||||||
if skip is not None:
|
|
||||||
# No detection this tick. Empty result → _handle_fsm_result no-op.
|
# No detection this tick. Empty result → _handle_fsm_result no-op.
|
||||||
return _TickSyncResult()
|
return _TickSyncResult()
|
||||||
return await asyncio.to_thread(
|
return await asyncio.to_thread(
|
||||||
@@ -1570,18 +1694,7 @@ async def _run_multi_tick(ctx: RunContext) -> "list[_TickSyncResult]":
|
|||||||
"""
|
"""
|
||||||
now = time.time()
|
now = time.time()
|
||||||
if ctx.lifecycle is not None:
|
if ctx.lifecycle is not None:
|
||||||
skip = _should_skip(now, ctx.lifecycle, ctx.cfg, ctx.canary)
|
if await _lifecycle_gate(ctx, now) is not None:
|
||||||
sb = _brief_status(ctx)
|
|
||||||
transition = _maybe_log_transition(
|
|
||||||
skip, ctx.lifecycle, now, ctx.audit, ctx.notifier, status_body=sb,
|
|
||||||
)
|
|
||||||
if transition == "market_open" and ctx.cfg.window_title:
|
|
||||||
title = await asyncio.to_thread(_focus_window_by_title, ctx.cfg.window_title)
|
|
||||||
ctx.audit.log({"ts": now, "event": "window_focused", "command": "market_open", "title": title})
|
|
||||||
await asyncio.sleep(0.15)
|
|
||||||
elif transition == "market_closed":
|
|
||||||
_handle_market_closed(ctx, now)
|
|
||||||
if skip is not None:
|
|
||||||
return []
|
return []
|
||||||
|
|
||||||
frame = await asyncio.to_thread(ctx.capture)
|
frame = await asyncio.to_thread(ctx.capture)
|
||||||
@@ -1827,15 +1940,14 @@ async def _dispatch_command(ctx: RunContext, cmd) -> None:
|
|||||||
f"{line1_state} | semnale: {ctx.state.fire_count} | {uptime_h:.1f}h",
|
f"{line1_state} | semnale: {ctx.state.fire_count} | {uptime_h:.1f}h",
|
||||||
f"{last_color} ({last_conf}) | poller: {sched_info}",
|
f"{last_color} ({last_conf}) | poller: {sched_info}",
|
||||||
]
|
]
|
||||||
oh = getattr(ctx.cfg, "operating_hours", None)
|
if ctx.lifecycle is not None:
|
||||||
if oh is not None and oh.enabled:
|
lines.extend(_window_status_lines(time.time(), ctx.lifecycle, ctx.cfg))
|
||||||
window_val = ctx.lifecycle.last_window_state if ctx.lifecycle else "—"
|
|
||||||
window_ro = {"open": "deschisă", "closed": "închisă"}.get(window_val or "", window_val or "—")
|
|
||||||
lines.append(f"fereastră: {window_ro}")
|
|
||||||
|
|
||||||
cfg_name = getattr(ctx.cfg, "config_version", None) or getattr(ctx.cfg, "version", None)
|
cfg_name = getattr(ctx.cfg, "config_version", None) or getattr(ctx.cfg, "version", None)
|
||||||
if isinstance(cfg_name, str) and cfg_name not in ("", "unknown"):
|
if isinstance(cfg_name, str) and cfg_name not in ("", "unknown"):
|
||||||
lines.append(f"config: {cfg_name}")
|
hb = getattr(ctx.cfg, "heartbeat_min", 0)
|
||||||
|
hb_txt = f"{hb} min" if isinstance(hb, int) and hb > 0 else "oprit"
|
||||||
|
lines.append(f"config: {cfg_name} | heartbeat: {hb_txt}")
|
||||||
|
|
||||||
ctx.notifier.send(Alert(kind="status", title="ATM Status", body="\n".join(lines)))
|
ctx.notifier.send(Alert(kind="status", title="ATM Status", body="\n".join(lines)))
|
||||||
elif cmd.action == "ss":
|
elif cmd.action == "ss":
|
||||||
@@ -1959,6 +2071,7 @@ async def _dispatch_command(ctx: RunContext, cmd) -> None:
|
|||||||
elif cmd.action == "window":
|
elif cmd.action == "window":
|
||||||
if ctx.lifecycle is None:
|
if ctx.lifecycle is None:
|
||||||
return
|
return
|
||||||
|
_save_session_window(ctx.lifecycle.session_window_path, cmd.window)
|
||||||
if cmd.window is None:
|
if cmd.window is None:
|
||||||
ctx.lifecycle.session_window = None
|
ctx.lifecycle.session_window = None
|
||||||
ctx.audit.log({"ts": time.time(), "event": "session_window_cleared"})
|
ctx.audit.log({"ts": time.time(), "event": "session_window_cleared"})
|
||||||
@@ -1977,15 +2090,18 @@ async def _dispatch_command(ctx: RunContext, cmd) -> None:
|
|||||||
ctx.notifier.send(Alert(
|
ctx.notifier.send(Alert(
|
||||||
kind="status",
|
kind="status",
|
||||||
title=f"Fereastră monitorizare: {start_w}–{stop_w} (ora locală, zilnic)",
|
title=f"Fereastră monitorizare: {start_w}–{stop_w} (ora locală, zilnic)",
|
||||||
body="În afara intervalului monitorizarea se pauzează automat.",
|
body=(
|
||||||
|
"În afara intervalului monitorizarea se pauzează automat.\n"
|
||||||
|
"Salvată — rămâne activă și după repornire."
|
||||||
|
),
|
||||||
))
|
))
|
||||||
elif cmd.action == "rebase":
|
elif cmd.action == "rebase":
|
||||||
await _dispatch_rebase(ctx, cmd)
|
await _dispatch_rebase(ctx, cmd)
|
||||||
elif cmd.action == "help":
|
elif cmd.action == "help":
|
||||||
body = (
|
body = (
|
||||||
"/status — stare FSM, uptime, ultima detecție\n"
|
"/status — stare FSM, uptime, fereastră + ore bursă\n"
|
||||||
"/ss — screenshot acum\n"
|
"/ss — screenshot acum\n"
|
||||||
"/pause — oprește detecția (heartbeat continuă)\n"
|
"/pause — oprește detecția (mesajele de deschidere/închidere continuă)\n"
|
||||||
"/resume — reia detecția (șterge user-pause și drift-pause)\n"
|
"/resume — reia detecția (șterge user-pause și drift-pause)\n"
|
||||||
"/rebase — propune phash nou pentru canary (confirm cu /rebase confirm)\n"
|
"/rebase — propune phash nou pentru canary (confirm cu /rebase confirm)\n"
|
||||||
"/3 — screenshot automat la fiecare 3 min (sau orice număr)\n"
|
"/3 — screenshot automat la fiecare 3 min (sau orice număr)\n"
|
||||||
@@ -2286,11 +2402,24 @@ async def run_live_async(cfg, duration_s=None, capture_stub: bool = False) -> No
|
|||||||
)
|
)
|
||||||
poller = TelegramPoller(cfg.telegram, cmd_queue, audit)
|
poller = TelegramPoller(cfg.telegram, cmd_queue, audit)
|
||||||
|
|
||||||
lifecycle = LifecycleState()
|
sw_path = Path("logs/session_window.json")
|
||||||
|
lifecycle = LifecycleState(
|
||||||
|
session_window=_load_session_window(sw_path), session_window_path=sw_path,
|
||||||
|
)
|
||||||
|
if lifecycle.session_window is not None:
|
||||||
|
sw_start, sw_stop = lifecycle.session_window
|
||||||
|
audit.log({"ts": time.time(), "event": "session_window_restored",
|
||||||
|
"start": sw_start, "stop": sw_stop})
|
||||||
|
notifier.send(Alert(
|
||||||
|
kind="status",
|
||||||
|
title=f"Fereastră monitorizare restaurată: {sw_start}–{sw_stop} (ora locală)",
|
||||||
|
body="/window off pentru a o dezactiva.",
|
||||||
|
silent=True,
|
||||||
|
))
|
||||||
# Seed lifecycle.last_window_state with the current status so we don't emit
|
# Seed lifecycle.last_window_state with the current status so we don't emit
|
||||||
# a spurious market_open alert on the very first tick (R2).
|
# a spurious market_open alert on the very first tick (R2).
|
||||||
_pre_skip = _should_skip(time.time(), lifecycle, cfg, canary)
|
_now = time.time()
|
||||||
_maybe_log_transition(_pre_skip, lifecycle, time.time(), audit, notifier)
|
_maybe_log_transition(_window_skip_reason(_now, lifecycle, cfg), lifecycle, _now, audit, notifier)
|
||||||
|
|
||||||
ctx = RunContext(
|
ctx = RunContext(
|
||||||
cfg=cfg, capture=capture, canary=canary, detector=detector, fsm=fsm,
|
cfg=cfg, capture=capture, canary=canary, detector=detector, fsm=fsm,
|
||||||
@@ -2350,7 +2479,12 @@ async def run_live_async(cfg, duration_s=None, capture_stub: bool = False) -> No
|
|||||||
# Launch background tasks
|
# Launch background tasks
|
||||||
t_scheduler = asyncio.create_task(scheduler.run(), name="scheduler")
|
t_scheduler = asyncio.create_task(scheduler.run(), name="scheduler")
|
||||||
t_poller = asyncio.create_task(poller.run(), name="poller")
|
t_poller = asyncio.create_task(poller.run(), name="poller")
|
||||||
t_heartbeat = asyncio.create_task(_heartbeat_loop(), name="heartbeat")
|
# heartbeat_min = 0 (default): no periodic heartbeat — the daily
|
||||||
|
# "Piața deschisă/închisă" alerts are the alive signal instead.
|
||||||
|
t_heartbeat = (
|
||||||
|
asyncio.create_task(_heartbeat_loop(), name="heartbeat")
|
||||||
|
if cfg.heartbeat_min > 0 else None
|
||||||
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await _detection_loop()
|
await _detection_loop()
|
||||||
@@ -2365,6 +2499,7 @@ async def run_live_async(cfg, duration_s=None, capture_stub: bool = False) -> No
|
|||||||
with contextlib.suppress(asyncio.CancelledError, Exception):
|
with contextlib.suppress(asyncio.CancelledError, Exception):
|
||||||
await t_poller
|
await t_poller
|
||||||
# 3. cancel heartbeat
|
# 3. cancel heartbeat
|
||||||
|
if t_heartbeat is not None:
|
||||||
t_heartbeat.cancel()
|
t_heartbeat.cancel()
|
||||||
with contextlib.suppress(asyncio.CancelledError, Exception):
|
with contextlib.suppress(asyncio.CancelledError, Exception):
|
||||||
await t_heartbeat
|
await t_heartbeat
|
||||||
|
|||||||
@@ -35,6 +35,10 @@ comenzi live Telegram. Fan-out trimite același eveniment pe mai multe canale.
|
|||||||
`/rebase confirm` în ≤180s rescrie `baseline_phash` în TOML-ul activ (păstrează
|
`/rebase confirm` în ≤180s rescrie `baseline_phash` în TOML-ul activ (păstrează
|
||||||
comentariile), mirror în `cfg` la runtime, clear `user_paused` + `drift_paused`.
|
comentariile), mirror în `cfg` la runtime, clear `user_paused` + `drift_paused`.
|
||||||
Fără confirm, nimic nu se schimbă.
|
Fără confirm, nimic nu se schimbă.
|
||||||
|
- **`/status`** — stare/semnale/uptime, ultima culoare, poller; blocul de
|
||||||
|
fereastră (`_window_status_lines`, doar dacă e configurat `operating_hours`
|
||||||
|
sau `/window`): `deschisă/închisă → următoarea deschidere/închidere` (ora
|
||||||
|
locală), `/window`, orele bursei + echivalentul local; `config | heartbeat`.
|
||||||
- **`/ss`** — top-3 buline din `dot_roi`: cerc roșu gros pe pick-ul FSM, cercuri
|
- **`/ss`** — top-3 buline din `dot_roi`: cerc roșu gros pe pick-ul FSM, cercuri
|
||||||
colorate subțiri pe vecini; caption cu nume/RGB/distanță/confidence +
|
colorate subțiri pe vecini; caption cu nume/RGB/distanță/confidence +
|
||||||
`config: {version}`. Culoarea cercului = `cfg.colors[name].rgb` (DRY cu paleta).
|
`config: {version}`. Culoarea cercului = `cfg.colors[name].rgb` (DRY cu paleta).
|
||||||
@@ -43,13 +47,20 @@ comenzi live Telegram. Fan-out trimite același eveniment pe mai multe canale.
|
|||||||
Dacă capture pică, title conține `⚠️ captură eșuată` și resume se execută oricum.
|
Dacă capture pică, title conține `⚠️ captură eșuată` și resume se execută oricum.
|
||||||
- **Drift-pause** — un singur alert Telegram pe tranziție. Cât e pauzat,
|
- **Drift-pause** — un singur alert Telegram pe tranziție. Cât e pauzat,
|
||||||
`/set_interval` e refuzat, caption-ul `/ss` avertizează că detecția e oprită,
|
`/set_interval` e refuzat, caption-ul `/ss` avertizează că detecția e oprită,
|
||||||
heartbeat arată `⚠️ pauzat (drift)` în loc de `activ`.
|
heartbeat-ul (dacă e pornit) arată `⚠️ pauzat (drift)` în loc de `activ`.
|
||||||
|
- **Alive signal** — heartbeat-ul periodic e **oprit implicit**
|
||||||
|
(`heartbeat_min = 0`). Semnalul „sunt viu" = alertele zilnice
|
||||||
|
„🟢 Piața deschisă" / „🔴 Piața închisă", emise pe tranziția ferestrei
|
||||||
|
(`_window_skip_reason` = `operating_hours` AND `/window`), **independent de
|
||||||
|
pauză** (user/drift) — body-ul spune dacă monitorizarea e pauzată.
|
||||||
- **`/window HH:MM-HH:MM`** (sau `HH:MM HH:MM`) — fereastră de monitorizare în
|
- **`/window HH:MM-HH:MM`** (sau `HH:MM HH:MM`) — fereastră de monitorizare în
|
||||||
**ora locală**, **recurentă zilnic**, stocată în `LifecycleState.session_window`.
|
**ora locală**, **recurentă zilnic**, stocată în `LifecycleState.session_window`
|
||||||
Se combină prin **AND** cu `operating_hours` (vezi `../AGENTS.md` → scheduler):
|
și **persistată** în `logs/session_window.json` (restaurată la pornire, cu o
|
||||||
în afara intervalului `_should_skip` întoarce `out_of_window_hours` ⇒ pauză
|
alertă silent). Se combină prin **AND** cu `operating_hours` (vezi
|
||||||
automată (alertă „Piața închisă" o dată + scheduler oprit + FSM reset).
|
`../AGENTS.md` → scheduler): în afara intervalului `_should_skip` întoarce
|
||||||
`/window off` (sau `clear`) șterge fereastra. Format invalid e ignorat.
|
`out_of_window_hours` ⇒ pauză automată (alertă „Piața închisă" o dată +
|
||||||
|
scheduler oprit + FSM reset). `/window off` (sau `clear`) șterge fereastra și
|
||||||
|
fișierul. Format invalid e ignorat.
|
||||||
|
|
||||||
## Work Guidance
|
## Work Guidance
|
||||||
|
|
||||||
|
|||||||
@@ -870,6 +870,63 @@ async def test_window_command_sets_and_clears_session_window():
|
|||||||
assert any(e.get("event") == "session_window_cleared" for e in ctx.audit.events)
|
assert any(e.get("event") == "session_window_cleared" for e in ctx.audit.events)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_window_command_persists_across_restart(tmp_path):
|
||||||
|
"""/window is saved to disk and restored on the next start; /window off deletes it."""
|
||||||
|
import atm.main as _main
|
||||||
|
from atm.commands import Command
|
||||||
|
|
||||||
|
sw_path = tmp_path / "logs" / "session_window.json"
|
||||||
|
ctx = _dispatch_ctx(lifecycle=_main.LifecycleState(session_window_path=sw_path))
|
||||||
|
await _main._dispatch_command(ctx, Command(action="window", window=("19:30", "22:00")))
|
||||||
|
assert _main._load_session_window(sw_path) == ("19:30", "22:00")
|
||||||
|
|
||||||
|
await _main._dispatch_command(ctx, Command(action="window", window=None))
|
||||||
|
assert not sw_path.exists()
|
||||||
|
assert _main._load_session_window(sw_path) is None
|
||||||
|
|
||||||
|
|
||||||
|
def test_load_session_window_corrupt_file_is_none(tmp_path):
|
||||||
|
import atm.main as _main
|
||||||
|
|
||||||
|
p = tmp_path / "session_window.json"
|
||||||
|
p.write_text("{not json", encoding="utf-8")
|
||||||
|
assert _main._load_session_window(p) is None
|
||||||
|
p.write_text('{"start": "19:30"}', encoding="utf-8")
|
||||||
|
assert _main._load_session_window(p) is None
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_open_close_alerts_sent_while_user_paused():
|
||||||
|
"""Daily open/close alerts are the alive signal → they fire even when paused."""
|
||||||
|
import atm.main as _main
|
||||||
|
|
||||||
|
cfg = _oh_cfg()
|
||||||
|
cfg.lockout_s = 240
|
||||||
|
tz = cfg.operating_hours._tz_cache
|
||||||
|
lifecycle = _main.LifecycleState(user_paused=True, last_window_state="closed")
|
||||||
|
ctx = _dispatch_ctx(lifecycle=lifecycle, cfg=cfg)
|
||||||
|
ctx.charts = []
|
||||||
|
|
||||||
|
mid = _dt.datetime(2026, 4, 20, 12, 0, tzinfo=tz).timestamp()
|
||||||
|
assert await _main._lifecycle_gate(ctx, mid) == "user_paused"
|
||||||
|
assert lifecycle.last_window_state == "open"
|
||||||
|
opened = [a for a in ctx.notifier.alerts if "deschisă" in a.title]
|
||||||
|
assert len(opened) == 1 and "oprită manual" in opened[0].body
|
||||||
|
# Paused → TradeStation isn't pulled to the front on open.
|
||||||
|
assert not any(e.get("event") == "window_focused" for e in ctx.audit.events)
|
||||||
|
|
||||||
|
close = _dt.datetime(2026, 4, 20, 16, 5, tzinfo=tz).timestamp()
|
||||||
|
assert await _main._lifecycle_gate(ctx, close) == "user_paused"
|
||||||
|
assert lifecycle.last_window_state == "closed"
|
||||||
|
assert sum("închisă" in a.title for a in ctx.notifier.alerts) == 1
|
||||||
|
|
||||||
|
# Same side of the boundary again → no duplicate.
|
||||||
|
n = len(ctx.notifier.alerts)
|
||||||
|
await _main._lifecycle_gate(ctx, close + 60)
|
||||||
|
assert len(ctx.notifier.alerts) == n
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
# Commit 5: /pause /resume dispatch (plan tests #11-15, #16, R2 #21)
|
# Commit 5: /pause /resume dispatch (plan tests #11-15, #16, R2 #21)
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
@@ -1741,7 +1798,60 @@ async def test_status_window_line_when_oh_enabled():
|
|||||||
|
|
||||||
status = [a for a in ctx.notifier.alerts if a.kind == "status"]
|
status = [a for a in ctx.notifier.alerts if a.kind == "status"]
|
||||||
body = status[0].body
|
body = status[0].body
|
||||||
assert "fereastră: deschisă" in body
|
# Open/closed depends on the real clock here; exact content is covered by
|
||||||
|
# the _window_status_lines tests below with fixed timestamps.
|
||||||
|
assert "fereastră: " in body
|
||||||
|
assert "/window: —" in body
|
||||||
|
assert "bursă: 09:30–16:00 America/New_York MON–FRI" in body
|
||||||
|
|
||||||
|
|
||||||
|
def test_window_status_lines_session_window_only():
|
||||||
|
"""Local /window only (oh off) → open/closed + next change in local time."""
|
||||||
|
import atm.main as _main
|
||||||
|
|
||||||
|
cfg = _oh_cfg(enabled=False)
|
||||||
|
lifecycle = _main.LifecycleState(session_window=("19:30", "22:00"))
|
||||||
|
inside = _dt.datetime(2026, 4, 20, 20, 0).timestamp()
|
||||||
|
assert _main._window_status_lines(inside, lifecycle, cfg) == [
|
||||||
|
"fereastră: deschisă → închidere azi 22:00",
|
||||||
|
"/window: 19:30–22:00 (local)",
|
||||||
|
]
|
||||||
|
after = _dt.datetime(2026, 4, 20, 22, 30).timestamp()
|
||||||
|
assert _main._window_status_lines(after, lifecycle, cfg)[0] == (
|
||||||
|
"fereastră: închisă → deschidere mâine 19:30"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_window_status_lines_nothing_configured_is_empty():
|
||||||
|
import atm.main as _main
|
||||||
|
lines = _main._window_status_lines(0.0, _main.LifecycleState(), _oh_cfg(enabled=False))
|
||||||
|
assert lines == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_window_status_lines_exchange_hours_local_conversion():
|
||||||
|
import atm.main as _main
|
||||||
|
|
||||||
|
cfg = _oh_cfg()
|
||||||
|
tz = cfg.operating_hours._tz_cache
|
||||||
|
mid = _dt.datetime(2026, 4, 20, 12, 0, tzinfo=tz).timestamp()
|
||||||
|
loc = lambda h, m: _dt.datetime(2026, 4, 20, h, m, tzinfo=tz).astimezone().strftime("%H:%M") # noqa: E731
|
||||||
|
lines = _main._window_status_lines(mid, _main.LifecycleState(), cfg)
|
||||||
|
assert lines[0].startswith("fereastră: deschisă → închidere ")
|
||||||
|
assert lines[1] == "/window: —"
|
||||||
|
assert lines[2] == (
|
||||||
|
f"bursă: 09:30–16:00 America/New_York MON–FRI (local {loc(9, 30)}–{loc(16, 0)})"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_next_window_change_skips_weekend():
|
||||||
|
"""Friday after close → next flip is Monday's open, not Saturday."""
|
||||||
|
import atm.main as _main
|
||||||
|
|
||||||
|
cfg = _oh_cfg()
|
||||||
|
tz = cfg.operating_hours._tz_cache
|
||||||
|
fri_evening = _dt.datetime(2026, 4, 24, 17, 0, tzinfo=tz).timestamp()
|
||||||
|
mon_open = _dt.datetime(2026, 4, 27, 9, 30, tzinfo=tz).timestamp()
|
||||||
|
assert _main._next_window_change(fri_evening, _main.LifecycleState(), cfg) == mon_open
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
|
|||||||
Reference in New Issue
Block a user