#!/usr/bin/env python3 """토글 켜진 종목의 급등락 감시 → 레이 텔레그램 알림. LLM을 깨우지 않음. 세 기준을 병행한다: 1. 거래소 상·하한가 도달 — 그날 갈 수 있는 끝까지 간 것. 가장 강한 신호라 무엇에도 가리지 않는다. 2. 거래소 VI 발동 — "이건 급등락이다"를 거래소가 공식 판정. 문턱을 우리가 정할 필요가 없다. 3. 자기 이력 분위수 — 장중 이탈폭(전일종가 대비)이 그 종목 과거 이탈폭의 p90 을 넘으면 알림. 고정 %를 쓰지 않는 이유 (실측 관심·감시 65종목 × 280거래일, 2026-08-04): - 장중 이탈폭 중간값이 3.8% → ±3%는 급등락이 아니라 평범한 날 - 종목별 변동성이 22배 차(ATR14 1.0%~22.3%). ±5% 문턱이면 KODEX 미국S&P500 은 발생 0회, SK이터닉스는 75회. 같은 5%가 한쪽엔 도달 불가, 한쪽엔 노이즈다. ⚠️ ATR 배수(ATR% × K)를 쓰지 않는 이유 — 2026-08-04 관리자님 지적으로 교체: 국내 주식 하루 가격제한폭이 ±30% 인데 ATR 배수는 상한이 없어 변동성 큰 종목의 문턱이 제한폭 밖으로 밀려난다. 실측: 1차(ATR×1.5)가 5/65 종목, 확대(ATR×3.0)가 31/65 종목(48%)에서 30% 초과 = **영원히 발동 불가**였다(ATR% 중간값이 10%라 확대는 절반이 죽는다). 분위수는 실제로 관측된 이탈폭이라 구조적으로 제한폭을 넘을 수 없다(p90~p99 전부 30% 초과 0종목). 덤으로 알림량이 정의상 (1-p) 비율로 확정된다 — p90 = 종목당 연 25회. ⚠️ 신규상장 첫날·정리매매는 제한폭 예외지만, 그런 종목은 이력이 없어 애초에 VI만 감시한다. 데이터원은 behive_web 의 /api/realtime/quotes — 현재가와 VI를 한 번에 주고 키움 호출이 0이다. VI(1h)는 시장 전역 broadcast라 구독이 필요 없고, 현재가(0B)는 구독이 필요한데 behive_web 의 _rt_gather_codes() 가 이미 보유+관심+감시를 구독하므로 토글 종목은 그 부분집합이다. 구독 상한에 걸려 빠진 종목만 ka10095 배치 1콜로 폴백한다. ⚠️ 상·하한가·VI 가 걸려 있는 동안 분위수 트리거는 **보류**했다가 풀린 뒤에 내보낸다 (2026-08-06 관리자님 요청). 그전엔 조용히 버려져서, VI 중에 문턱을 넘었다가 해제 시점에 되밀린 움직임은 알림이 아예 없었다. 같은 사건을 두 번 알리지 않으면서 사실은 잃지 않기 위함. state: behive_surge_toggles.json — 감시 대상 (behive_web 의 🔔 토글이 씀) surge_thresholds.json — {date, by_code: {code: {thr, thr_big, prev_close, n}}}. 하루 1회만 계산 surge_limits.json — {date, by_code: {code: {upl, lst}}}. 거래소 상·하한가, 하루 1콜 surge_alerts.json — {date: {code: [트리거키], __pending__: {code: {키: 보류정보}}}} 방향별 1회 + 확대단계 1회 + 상·하한가 1회, VI 는 발동 건마다 1회 Usage: python3 surge_monitor.py check # 1회 감시 (launchd 용) python3 surge_monitor.py check --force # 장외에도 실행 (테스트용) python3 surge_monitor.py dry-run # 판정만 출력, 텔레그램 발송 없음 python3 surge_monitor.py list # 토글 종목의 문턱(%)·발동가 출력 """ from __future__ import annotations import fcntl import json import sys import time import urllib.parse import urllib.request from contextlib import contextmanager from datetime import datetime, timezone, timedelta from pathlib import Path KST = timezone(timedelta(hours=9)) WORKSPACE = Path('/Users/snowoyh/.openclaw/agents/stock/workspace') sys.path.insert(0, str(WORKSPACE / 'scripts')) sys.path.insert(0, str(WORKSPACE)) import kiwoom_client as kc # noqa: E402 STATE_DIR = WORKSPACE / 'state' TOGGLES = STATE_DIR / 'behive_surge_toggles.json' THR_CACHE = STATE_DIR / 'surge_thresholds.json' LIMITS_CACHE = STATE_DIR / 'surge_limits.json' ALERTS_STATE = STATE_DIR / 'surge_alerts.json' # 보류 트리거를 담는 예약 키. 종목코드는 6자리라 이 이름과 충돌하지 않는다. PENDING_KEY = '__pending__' CONFIG_PATH = Path('/Users/snowoyh/.openclaw/openclaw.json') TELEGRAM_ACCOUNT = 'stock' # 레이 봇 # 1차 문턱 = 그 종목 과거 장중 이탈폭의 이 분위수. 통수를 정하는 유일한 손잡이 — # 분위수라 알림량이 정의상 (1-p) 비율로 확정된다. 0.90 = 종목당 연 25회(≈2주 1번). SURGE_PCTL = 0.90 # 확대 단계 — 1차 알림 후 여기까지 더 벌어지면 한 번 더. 종목당 하루 최대 2통. SURGE_PCTL_BIG = 0.98 # 분위수 산출에 필요한 최소 관측일. 미달(신규상장 등)이면 VI만 감시한다. MIN_HISTORY = 60 # 이력 조회 봉 수 — sqlite 캐시에서 읽는 양만 늘린다(추가 API 콜 없음). HISTORY_COUNT = 400 # 문턱 산출 페이싱 — ka10081 유량이 초당 5건이라 캐시가 stale 한 종목이 많으면 429가 쏟아진다. MAX_COMPUTE_PER_CYCLE = 12 COMPUTE_PACE_SEC = 0.35 # behive_web 실시간 엔드포인트. BIND_HOST='' 라 localhost 로 도달한다. RT_URL = 'http://127.0.0.1:18790/api/realtime/quotes' RT_TIMEOUT = 5 MARKET_OPEN = (9, 0) MARKET_CLOSE = (15, 35) # 15:30 마감 + 최종 체결 버퍼 def load_json(path: Path, default): if path.exists(): try: return json.loads(path.read_text()) except Exception: return default return default def save_json(path: Path, data): path.parent.mkdir(parents=True, exist_ok=True) tmp = path.with_suffix(path.suffix + '.tmp') tmp.write_text(json.dumps(data, ensure_ascii=False, indent=2)) tmp.replace(path) @contextmanager def alerts_lock(): """surge_alerts.json 동시 쓰기 직렬화. watchlist_monitor 와 동일 패턴, 별도 lock 파일.""" lock_path = ALERTS_STATE.with_suffix(ALERTS_STATE.suffix + '.lock') lock_path.parent.mkdir(parents=True, exist_ok=True) f = open(lock_path, 'a') try: fcntl.flock(f.fileno(), fcntl.LOCK_EX) yield finally: try: fcntl.flock(f.fileno(), fcntl.LOCK_UN) finally: f.close() def is_market_hours() -> bool: now = datetime.now(KST) if now.weekday() >= 5: return False # KRX 휴장일도 거래 없음. 데이터 파일 누락이나 import 실패 시엔 평소대로 진행. try: from holiday_sync import is_holiday_today if is_holiday_today(): return False except Exception: pass mins = now.hour * 60 + now.minute return (MARKET_OPEN[0] * 60 + MARKET_OPEN[1]) <= mins <= (MARKET_CLOSE[0] * 60 + MARKET_CLOSE[1]) def send_telegram(text: str) -> bool: cfg = json.loads(CONFIG_PATH.read_text()) acct = cfg['channels']['telegram']['accounts'][TELEGRAM_ACCOUNT] token = acct['botToken'] chat_ids = acct.get('allowFrom') or [] if not chat_ids: print('no telegram chat_ids', file=sys.stderr) return False url = f'https://api.telegram.org/bot{token}/sendMessage' ok = True for chat_id in chat_ids: data = urllib.parse.urlencode({ 'chat_id': chat_id, 'text': text[:4000], 'disable_web_page_preview': 'true', }).encode() try: req = urllib.request.Request(url, data=data, method='POST') with urllib.request.urlopen(req, timeout=15) as r: if r.status != 200: ok = False print(f'telegram HTTP {r.status}', file=sys.stderr) except Exception as e: print(f'telegram error: {e}', file=sys.stderr) ok = False return ok # ---------------- 감시 대상 ---------------- def load_toggles() -> dict[str, str]: """{code: name} — behive_web 의 🔔 토글이 켠 종목만.""" raw = load_json(TOGGLES, {}) by_code = raw.get('by_code') if isinstance(raw, dict) else None if not isinstance(by_code, dict): return {} out: dict[str, str] = {} for code, v in by_code.items(): cc = kc._clean_code(code) if not cc or not v: continue # v 는 True(레거시) 또는 {'name': ...} out[cc] = (v.get('name') or '') if isinstance(v, dict) else '' return out # ---------------- 데이터원 ---------------- def fetch_realtime(codes: list[str]) -> tuple[dict, dict, bool]: """behive_web 실시간 허브에서 (quotes, vi, ok). 키움 호출 0. 허브가 콜드하거나 behive_web 이 죽어 있으면 ok=False — 호출측이 폴백/경고를 결정한다. """ if not codes: return {}, {}, False url = f'{RT_URL}?codes={",".join(codes)}' try: with urllib.request.urlopen(url, timeout=RT_TIMEOUT) as r: d = json.loads(r.read().decode()) except Exception as e: print(f'[warn] behive_web 실시간 조회 실패 ({e}) — VI 감시 불가, ATR만 진행', file=sys.stderr) return {}, {}, False return (d.get('quotes') or {}), (d.get('vi') or {}), bool(d.get('connected')) def fetch_quotes_fallback(codes: list[str]) -> dict: """구독 상한에 걸려 허브에 없는 종목만 ka10095 배치 1콜로 보충.""" if not codes: return {} try: return kc.get_watchlist_quotes(codes) or {} except Exception as e: print(f'[warn] ka10095 배치 실패: {e}', file=sys.stderr) return {} # ---------------- 문턱 (자기 이력 분위수) ---------------- def _quantile(sorted_vals: list[float], p: float) -> float: """오름차순 리스트의 p 분위수. 실측 스크립트와 같은 식(nearest-rank)을 써야 값이 일치한다.""" return sorted_vals[min(len(sorted_vals) - 1, int(len(sorted_vals) * p))] def _compute_thresholds(code: str) -> tuple[dict | None, bool]: """(문턱 dict 또는 None, 재시도해야 하는가). 장중 이탈폭 = max(|고가−전일종가|, |저가−전일종가|) / 전일종가 × 100. 종가 대비가 아니라 **전일종가 대비 장중 최대 이탈폭**을 쓰는 이유 — 실시간 감시가 보는 값이 `현재가 vs 전일종가`(키움 pct)라서 같은 축으로 비교해야 한다. 종가기준은 장중 움직임을 과소평가해(중간값 3.8% vs 종가기준 1.x%) 문턱이 너무 낮게 잡힌다. ⚠️ 두 번째 반환값이 필요한 이유 — 조회 실패(429·네트워크)를 이력부족과 같이 취급해 빈 dict 로 캐시하면 그 종목이 **하루 내내 조용히 VI만 감시**하게 된다. 실패는 재시도 대상. """ try: import daily_candles_cache as dcc candles = dcc.get_candles(code, HISTORY_COUNT) except Exception as e: print(f'[{code}] 일봉 조회 실패 (다음 사이클 재시도): {e}', file=sys.stderr) return None, True if not candles or len(candles) < MIN_HISTORY: return None, False exc: list[float] = [] for i in range(1, len(candles)): pc = candles[i - 1]['close'] if not pc: continue exc.append(max(abs(candles[i]['high'] - pc), abs(candles[i]['low'] - pc)) / pc * 100) if len(exc) < MIN_HISTORY: return None, False exc.sort() prev_close = candles[-1]['close'] if not prev_close: return None, False return { 'thr': round(_quantile(exc, SURGE_PCTL), 2), 'thr_big': round(_quantile(exc, SURGE_PCTL_BIG), 2), 'prev_close': prev_close, 'n': len(exc), }, False def get_threshold_map(codes: list[str]) -> dict[str, dict]: """종목별 문턱 맵. 하루 1회만 계산하고 캐시한다. ⚠️ daily_candles_cache.get_candles 는 캐시 최신봉이 어제보다 오래되면 ka10081 을 때린다. 매 사이클(1분) × 종목수만큼 호출하면 폭주하므로 날짜가 바뀔 때와 새 토글이 생길 때만 계산. ⚠️ 분위수 파라미터를 바꿨으면 캐시가 옛 값을 들고 있으니 `pctl` 서명이 다르면 재계산한다. ⚠️ **ka10081 유량 제한이 초당 5건**이라 종목을 한꺼번에 돌리면 429가 쏟아진다(2026-08-04 실측: 70종목 일괄 산출 시 26종목 실패). 그래서 사이클당 `MAX_COMPUTE_PER_CYCLE` 개까지만, `COMPUTE_PACE_SEC` 간격으로 계산한다. 남은 종목은 다음 사이클이 이어받는다(1분 간격 × 396회라 장 시작 몇 분 안에 전부 채워진다). 조회 실패는 캐시하지 않아 다음 사이클에 재시도된다. """ today = datetime.now(KST).strftime('%Y-%m-%d') sig = f'{SURGE_PCTL}/{SURGE_PCTL_BIG}' cache = load_json(THR_CACHE, {}) fresh = cache.get('date') == today and cache.get('pctl') == sig by_code = cache.get('by_code') if fresh else {} if not isinstance(by_code, dict): by_code = {} missing = [c for c in codes if c not in by_code][:MAX_COMPUTE_PER_CYCLE] if missing: for i, c in enumerate(missing): if i: time.sleep(COMPUTE_PACE_SEC) r, retry = _compute_thresholds(c) if retry: continue # 실패는 기록하지 않는다 — 다음 사이클에 다시 시도 by_code[c] = r if r else {} save_json(THR_CACHE, {'date': today, 'pctl': sig, 'by_code': by_code}) return {c: v for c, v in by_code.items() if v.get('thr')} # ---------------- 상·하한가 (거래소 계산값) ---------------- def get_limit_map(codes: list[str]) -> dict[str, dict]: """{code: {'upl', 'lst'}} — 거래소가 계산한 가격제한폭. 하루 1콜(ka10095 배치). 기준가(전일종가)로 정해져 장중 불변이라 하루 1회면 충분하다. 분위수 문턱과 달리 일봉 이력이 필요 없어서 신규상장 등 이력 부족 종목도 이 기준으론 감시된다. ⚠️ 배치 호출 자체가 실패하면 캐시하지 않는다 — 문턱 캐시와 같은 이유로, 실패를 '제한가 없음'으로 굳히면 그 종목이 하루 내내 조용히 빠진다. 응답에 값이 없는 종목만 빈 dict 로 굳혀 매 사이클 재조회를 막는다. """ if not codes: return {} today = datetime.now(KST).strftime('%Y-%m-%d') cache = load_json(LIMITS_CACHE, {}) by_code = cache.get('by_code') if cache.get('date') == today else {} if not isinstance(by_code, dict): by_code = {} missing = [c for c in codes if c not in by_code] if missing: try: rows = kc.get_watchlist_quotes(missing) or {} except Exception as e: print(f'[warn] 상하한가 조회 실패 (다음 사이클 재시도): {e}', file=sys.stderr) return {c: v for c, v in by_code.items() if v.get('upl')} for c in missing: r = rows.get(c) or {} upl, lst = r.get('upper_limit') or 0, r.get('lower_limit') or 0 by_code[c] = {'upl': upl, 'lst': lst} if (upl and lst) else {} save_json(LIMITS_CACHE, {'date': today, 'by_code': by_code}) return {c: v for c, v in by_code.items() if v.get('upl')} # ---------------- 판정 ---------------- def _limit_hit(price: int, limits: dict | None) -> str | None: """'up'(상한가) | 'down'(하한가) | None. 거래소 제한가와 현재가 비교.""" if not limits or not price: return None upl, lst = limits.get('upl') or 0, limits.get('lst') or 0 if upl and price >= upl: return 'up' if lst and price <= lst: return 'down' return None def _quantile_triggers(pct: float, thr: dict | None) -> list[dict]: """분위수 트리거만 — [{key, reason, word}].""" if not thr or not thr.get('thr'): return [] direction = 'up' if pct > 0 else 'down' word = '급등' if pct > 0 else '급락' t1, t2 = thr['thr'], thr.get('thr_big') or 0 top1 = round((1 - SURGE_PCTL) * 100) top2 = round((1 - SURGE_PCTL_BIG) * 100) out: list[dict] = [] if t2 and abs(pct) >= t2: out.append({'key': f'{direction}2', 'word': word, 'reason': f'{thr["n"]}일 중 상위 {top2}% 움직임 (문턱 {t2:.1f}%)'}) if abs(pct) >= t1: out.append({'key': direction, 'word': word, 'reason': f'{thr["n"]}일 중 상위 {top1}% 움직임 (문턱 {t1:.1f}%)'}) return out def evaluate(pct: float, price: int, thr: dict | None, vi: dict | None, limits: dict | None) -> tuple[list[dict], list[dict], str]: """(즉시 알릴 트리거, 보류할 트리거, 가림사유). 거래소 판정(상·하한가·VI)이 걸려 있는 동안 분위수 트리거는 **보류**한다. 같은 사건을 두 번 알리지 않으면서도, 그 사이 문턱을 넘은 사실은 버리지 않고 풀린 뒤에 내보내기 위해서다. 가림사유가 빈 문자열이면 가린 것이 없다는 뜻 — 호출측은 이때 보류분을 방출한다. ⚠️ 가림 여부는 `defer` 가 비었는지로 판단하면 안 된다. VI 중에 주가가 문턱 아래로 되밀리면 defer 도 비는데, 그걸 '가림 해제'로 읽으면 VI 도중에 보류분이 새어나간다. """ now: list[dict] = [] hit = _limit_hit(price, limits) if hit == 'up': now.append({'key': 'limit_up', 'word': '상한가', 'reason': f'거래소 상한가 도달 ({(limits or {}).get("upl", 0):,}원)'}) elif hit == 'down': now.append({'key': 'limit_down', 'word': '하한가', 'reason': f'거래소 하한가 도달 ({(limits or {}).get("lst", 0):,}원)'}) vi_on = bool(vi and vi.get('active')) if vi_on: t = (vi.get('trigger_time') or '').strip() or 'na' kind = (vi.get('apply_kind') or '').strip() or 'VI' now.append({'key': f'vi:{t}', 'word': None, 'reason': f'거래소 VI 발동 ({kind})'}) quantile = _quantile_triggers(pct, thr) if hit or vi_on: return now, quantile, ({'up': '상한가', 'down': '하한가'}.get(hit) or 'VI') return now + quantile, [], '' def build_message(records: list[dict]) -> str: ts = datetime.now(KST).strftime('%m/%d %H:%M') lines = [f'[급등락] {ts}', f'{len(records)}건 감지'] for r in records: pct = r['pct'] # ⚠️ 아이콘은 pct 가 아니라 방향어에서 파생한다. 보류 방출 건은 돌파 시점 방향(급락)과 # 방출 시점 현재가 부호(+)가 어긋날 수 있어, pct 로 고르면 '🚀 급락'이 찍힌다. word = r.get('word') or ('급등' if pct > 0 else '급락') icon = {'상한가': '⏫', '하한가': '⏬', '급등': '🚀', '급락': '🔻'}.get(word, '🚀') label = r['name'] or r['code'] lines.append('') lines.append(f'{icon} #{label} {word} ({pct:+.2f}%)') lines.append(f'• 현재가: {r["price"]:,}원') lines.append(f'• 사유: {r["reason"]}') if r.get('vi_price'): lines.append(f'• VI 발동가: {r["vi_price"]:,}원') return '\n'.join(lines) # ---------------- 실행 ---------------- def _prune(state: dict, today: str) -> dict: """오늘·어제만 남긴다 (무한 증식 방지).""" keep = {today, (datetime.now(KST) - timedelta(days=1)).strftime('%Y-%m-%d')} return {k: v for k, v in state.items() if k in keep} def run(dry: bool = False, force: bool = False) -> int: if not force and not is_market_hours(): print(f'장외 시간 — skip ({datetime.now(KST).strftime("%Y-%m-%d %H:%M")})') return 0 toggles = load_toggles() if not toggles: print('감시 대상 없음 — 자산웹에서 🔔 토글을 켜주세요') return 0 codes = sorted(toggles) quotes, vi_map, connected = fetch_realtime(codes) missing = [c for c in codes if not (quotes.get(c) or {}).get('price')] if missing: for c, q in fetch_quotes_fallback(missing).items(): quotes[kc._clean_code(c)] = {'price': q.get('price'), 'pct': q.get('change_pct'), 'change': q.get('change')} thr_map = get_threshold_map(codes) limit_map = get_limit_map(codes) today = datetime.now(KST).strftime('%Y-%m-%d') day_state = load_json(ALERTS_STATE, {}).get(today) or {} already = {f'{c}:{k}' for c, ks in day_state.items() if isinstance(ks, list) for k in ks} # 보류분 — 감시에서 빠진 종목 것은 버린다(토글이 꺼졌으면 방출할 이유가 없다). held = day_state.get(PENDING_KEY) or {} pending = {c: v for c, v in held.items() if c in toggles and isinstance(v, dict)} if isinstance(held, dict) else {} pending_sig = json.dumps(pending, sort_keys=True) to_send: list[dict] = [] flushed: list[tuple[str, str]] = [] # (code, key) — 발송 성공 시에만 보류에서 지운다 def _rec(code: str, price: int, pct: float, trig: dict, vi: dict | None) -> dict: r = { 'code': code, 'name': toggles.get(code) or '', 'price': price, 'pct': pct, 'key': trig['key'], 'word': trig.get('word'), 'reason': trig['reason'], } if vi and vi.get('active'): r['vi_price'] = vi.get('trigger_price') or 0 return r for code in codes: q = quotes.get(code) or {} if not q.get('price') or q.get('pct') is None: continue price, pct = int(q['price']), float(q['pct']) vi = vi_map.get(code) now_trigs, defer_trigs, blocked_by = evaluate( pct, price, thr_map.get(code), vi, limit_map.get(code)) for trig in now_trigs: if f'{code}:{trig["key"]}' in already: continue to_send.append(_rec(code, price, pct, trig, vi)) already.add(f'{code}:{trig["key"]}') if blocked_by: # 가림 중 — 문턱 돌파 사실만 적어둔다. 이탈폭이 가장 컸던 시점을 남긴다. for trig in defer_trigs: if f'{code}:{trig["key"]}' in already: continue cur = (pending.get(code) or {}).get(trig['key']) or {} if abs(pct) > abs(cur.get('pct') or 0): pending.setdefault(code, {})[trig['key']] = { 'pct': pct, 'word': trig.get('word'), 'reason': trig['reason'], 'via': blocked_by, } continue # 가림 해제 — 보류분 방출. 같은 키를 위에서 이미 즉시 알렸으면 already 가 걸러낸다. for key, d in sorted((pending.get(code) or {}).items()): if f'{code}:{key}' in already: flushed.append((code, key)) continue trig = { 'key': key, 'word': d.get('word'), 'reason': f'{d.get("reason") or "문턱 돌파"} — {d.get("via") or "VI"} 중 최대 {d.get("pct", 0):+.2f}% 도달', } to_send.append(_rec(code, price, pct, trig, vi)) already.add(f'{code}:{key}') flushed.append((code, key)) if dry: print(f'감시 {len(codes)}종목 / 허브연결 {connected} / 문턱 {len(thr_map)}종목 / 제한가 {len(limit_map)}종목') for code in codes: q = quotes.get(code) or {} t = thr_map.get(code) or {} lm = limit_map.get(code) or {} thr = f'{t["thr"]:.1f}% / 확대 {t["thr_big"]:.1f}%' if t.get('thr') else '— (이력부족, VI·제한가만)' lim = f'{lm["upl"]:,}/{lm["lst"]:,}' if lm.get('upl') else '—' print(f' {code} {toggles.get(code) or "":<12} 등락 {q.get("pct")}% / 문턱 {thr} / 상하한 {lim}' f'{" / VI" if (vi_map.get(code) or {}).get("active") else ""}' f'{" / 보류 " + ",".join(sorted(pending[code])) if pending.get(code) else ""}') if to_send: print('--- DRY ---') print(build_message(to_send)) print(f'done. watched={len(codes)} triggered={len(to_send)} pending={sum(len(v) for v in pending.values())}') return 0 ok = send_telegram(build_message(to_send)) if to_send else True if ok: for code, key in flushed: bucket = pending.get(code) or {} bucket.pop(key, None) if not bucket: pending.pop(code, None) if not to_send and json.dumps(pending, sort_keys=True) == pending_sig: print(f'done. watched={len(codes)} triggered=0') return 0 with alerts_lock(): latest = _prune(load_json(ALERTS_STATE, {}), today) day = latest.setdefault(today, {}) if ok: for r in to_send: bucket = day.setdefault(r['code'], []) if r['key'] not in bucket: bucket.append(r['key']) print(f'alerted {r["code"]} {r["name"]}:{r["key"]} {r["pct"]:+.2f}% @ {r["price"]:,}원') if pending: day[PENDING_KEY] = pending else: day.pop(PENDING_KEY, None) save_json(ALERTS_STATE, latest) print(f'done. watched={len(codes)} triggered={len(to_send)} pending={sum(len(v) for v in pending.values())}') return 0 def cmd_list() -> int: toggles = load_toggles() if not toggles: print('감시 대상 없음 — 자산웹에서 🔔 토글을 켜주세요') return 0 codes = sorted(toggles) thr_map = get_threshold_map(codes) limit_map = get_limit_map(codes) top1 = round((1 - SURGE_PCTL) * 100) top2 = round((1 - SURGE_PCTL_BIG) * 100) print(f'감시 {len(codes)}종목 — 1차=자기이력 상위 {top1}% / 확대=상위 {top2}%') print(f'{"종목":<16}{"전일종가":>10}{"1차":>7}{"급등가":>10}{"급락가":>10}{"확대":>7}' f'{"상한가":>10}{"하한가":>10}') for c in codes: t = thr_map.get(c) or {} lm = limit_map.get(c) or {} name = (toggles[c] or c)[:15] lim = f'{lm["upl"]:>10,}{lm["lst"]:>10,}' if lm.get('upl') else f'{"—":>10}{"—":>10}' if not t.get('thr'): print(f'{name:<16}{"이력 부족 — VI·제한가만":>30}{lim}') continue pv, t1, t2 = t['prev_close'], t['thr'], t['thr_big'] print(f'{name:<16}{pv:>10,}{t1:>6.1f}%{round(pv * (1 + t1 / 100)):>10,}' f'{round(pv * (1 - t1 / 100)):>10,}{t2:>6.1f}%{lim}') print(f'\n※ 상·하한가 도달은 문턱과 별개로 항상 알림 (거래소 계산값, 하루 1콜 캐시)') print(f'※ VI 발동은 이 문턱과 별개로 먼저 알림 (정적VI 전일종가 ±10% 부근)') print(f'※ 상·하한가·VI 중 문턱을 넘으면 보류했다가 풀린 뒤 알림 (사실을 버리지 않음)') print(f'※ 문턱은 하루 가격제한폭 ±30% 안에 있음이 보장됨 (실제 관측된 이탈폭의 분위수)') return 0 def main(): cmd = sys.argv[1] if len(sys.argv) > 1 else 'help' force = '--force' in sys.argv[2:] try: if cmd == 'check': return run(dry=False, force=force) if cmd == 'dry-run': return run(dry=True, force=True) if cmd == 'list': return cmd_list() print(__doc__, file=sys.stderr) return 2 except Exception as e: import traceback tb = traceback.format_exc() print(f'[fatal] {e}\n{tb}', file=sys.stderr) try: send_telegram(f'⚠️ [surge_monitor] 실행 실패\n{type(e).__name__}: {e}') except Exception: pass return 1 if __name__ == '__main__': sys.exit(main())