Files
Bithumb/scripts/3_run_vol_monitor.py
dsyoon 48dcdd0ae8 feat(rsi_dca): 7종 RSI 정액 매수 라이브 전략·모니터 고도화·수집 안정화
RSI DCA 전략 (신규)
- rsi_dca_engine/runner/control: 1분봉 RSI(14) 종목별 기준선 상향 돌파 시 정액 매수,
  일 상한, 신호 45분 유효, 킬스위치, 인터벌 변경 시 커서 재초기화, 체결가 거래소 보정
- scripts: 3_run_rsi_dca(.py/_cron.sh), 백테스트·인터벌 비교, go-live 스위치
- 설정: RSI_DCA_* (모드·종목·기준선·종목별 오버라이드·일 상한 등)

모니터 (vol_live_monitor / vol_monitor_chart)
- 분봉 탭(/api/candles), RSI(14) 패널·종목별 기준선, 3패널 시간축 정렬, KST 표기
- 자동매수 ON/OFF 패널(/api/rsi/status·toggle), 빌드 해시 기반 자동 새로고침, 지연 경고
- 요약표: 거래소 평균매입가 기준 보유원금·수익률, 총평가 손익, 수익률순 동적 정렬
- 잔고 스냅샷을 계좌 전체 조회 1회로 통합, 실시간 시세 반영

데이터·수집
- candle_store/loader: SQL 범위·LIMIT 조회로 대형 테이블 전량 스캔 제거
- 절전·재부팅 후 공백 자동 백필(gap_backfill, 00_backfill_gaps) 및 cron 연동
- 수집 cron 분할(매분 핵심 분봉·5분 전체), 프로젝트 한정 lock 패턴, exec 제거로 lock 정리 복구
- 모니터 종목(VOL_MONITOR_SYMBOLS)·수집 종목 7종 분리, 한글 코인명 추가

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-24 18:30:31 +09:00

386 lines
13 KiB
Python
Executable File

#!/usr/bin/env python3
"""vol_live 모니터 — JSON/HTML 갱신 + HTTP 서버 (통합).
기본 (인자 없음): 전체 갱신(--full) 후 서버 기동
python scripts/3_run_vol_monitor.py
갱신만:
python scripts/3_run_vol_monitor.py --refresh-only
서버만:
python scripts/3_run_vol_monitor.py --serve-only
"""
from __future__ import annotations
import argparse
import json
import logging
import os
import sys
import threading
import time
from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from urllib.parse import parse_qs, urlparse
_ROOT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(_ROOT / "src"))
from bithumb.config import load_settings # noqa: E402
from bithumb.operations.rsi_dca_control import rsi_status, set_rsi_enabled # noqa: E402
from bithumb.operations.vol_breakout_engine import load_vol_state # noqa: E402
from bithumb.operations.vol_live_monitor import ( # noqa: E402
build_candles_api_payload,
fetch_live_balance_snapshot,
patch_vol_monitor_balance,
write_vol_monitor,
)
logger = logging.getLogger("vol_monitor")
_refresh_lock = threading.Lock()
_balance_lock = threading.Lock()
_CLIENT_GONE = (BrokenPipeError, ConnectionResetError)
def _client_gone(exc: BaseException) -> bool:
"""브라우저가 응답 전 연결을 끊은 경우."""
return isinstance(exc, _CLIENT_GONE)
def refresh_vol_live_monitor(*, write_html: bool = True) -> dict:
"""state + DB 캔들 기준 전체 JSON/HTML 갱신."""
settings = load_settings()
state = load_vol_state(settings.vol_state_json)
if settings.ops_mode == "live":
try:
bal = fetch_live_balance_snapshot(settings)
snap = state.setdefault("portfolio_snapshot", {})
snap["cash_krw"] = bal.get("cash_krw", snap.get("cash_krw"))
positions = snap.setdefault("positions", {})
for sym, qty in (bal.get("positions") or {}).items():
positions[sym] = qty
snap["avg_prices"] = dict(bal.get("avg_prices") or {})
except Exception as exc: # noqa: BLE001
logger.warning("live balance sync skipped: %s", exc)
json_path, html_path = write_vol_monitor(settings, state)
logger.debug("monitor written: %s", json_path)
if write_html:
logger.debug("html: %s", html_path)
return {"ok": True, "json": str(json_path), "html": str(html_path)}
def refresh_vol_live_balance() -> dict:
"""거래소 잔고만 JSON summary 패치."""
settings = load_settings()
bal = fetch_live_balance_snapshot(settings)
return patch_vol_monitor_balance(settings.vol_monitor_json, bal)
def fetch_live_balance() -> dict:
"""서버 /api/balance용."""
settings = load_settings()
if settings.ops_mode != "live":
state = load_vol_state(settings.vol_state_json)
snap = state.get("portfolio_snapshot") or {}
return {
"ok": True,
"cash_krw": snap.get("cash_krw", 0),
"positions": snap.get("positions") or {},
"mode": settings.ops_mode,
}
return fetch_live_balance_snapshot(settings)
def _out_dir() -> Path:
return load_settings().vol_monitor_html.parent
def _api_refresh() -> dict:
out_dir = _out_dir()
json_path = out_dir / "vol_live_chart.json"
with _refresh_lock:
if not json_path.is_file():
return refresh_vol_live_monitor(write_html=False)
return refresh_vol_live_balance()
def _api_balance() -> dict:
with _balance_lock:
return fetch_live_balance()
_rsi_lock = threading.Lock()
def _api_rsi_status() -> dict:
return rsi_status(load_settings())
def _api_rsi_toggle(query: str) -> dict:
"""/api/rsi/toggle[?enable=1|0] — 파라미터 없으면 반전."""
qs = parse_qs(query or "")
raw = (qs.get("enable") or [""])[0].strip().lower()
settings = load_settings()
with _rsi_lock:
if raw in ("1", "true", "on", "yes"):
target = True
elif raw in ("0", "false", "off", "no"):
target = False
else:
target = not rsi_status(settings)["enabled"]
out = set_rsi_enabled(settings, target)
logger.warning("RSI 자동매수 %s (mode=%s)", "ON" if target else "OFF", out.get("mode"))
return out
def _api_candles(query: str) -> dict:
"""/api/candles?symbol=TRX&interval=5[&bars=800] — DB 직접 조회."""
qs = parse_qs(query or "")
symbol = (qs.get("symbol") or [""])[0]
interval = (qs.get("interval") or ["15"])[0]
bars_raw = (qs.get("bars") or [""])[0]
settings = load_settings()
max_bars = None
if bars_raw:
try:
max_bars = max(10, min(int(bars_raw), 20000))
except ValueError:
max_bars = None
return build_candles_api_payload(settings, symbol, interval, max_bars=max_bars)
class MonitorHandler(SimpleHTTPRequestHandler):
"""vol_live 정적 파일 + /api/chart · /api/candles · /api/balance · /api/refresh · /api/rsi/status · /api/rsi/toggle."""
_static_dir: str | None = None
_access_log: bool = False
def __init__(self, *args, **kwargs) -> None:
if MonitorHandler._static_dir is None:
MonitorHandler._static_dir = str(_out_dir())
super().__init__(*args, directory=MonitorHandler._static_dir, **kwargs)
def _chart_json_path(self) -> Path:
return Path(self.directory) / "vol_live_chart.json"
def _serve_chart_json(self) -> None:
path = self._chart_json_path()
if not path.is_file():
self.send_error(404, "chart json not found")
return
body: bytes | None = None
for attempt in range(3):
try:
body = path.read_bytes()
json.loads(body.decode("utf-8"))
break
except (json.JSONDecodeError, OSError):
if attempt >= 2:
self.send_error(503, "chart json temporarily unavailable")
return
time.sleep(0.05)
if body is None:
self.send_error(503, "chart json unavailable")
return
self.send_response(200)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Cache-Control", "no-store, must-revalidate")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
try:
self.wfile.write(body)
except _CLIENT_GONE:
logger.debug("client disconnected during chart json")
def log_message(self, fmt: str, *args) -> None:
"""HTTP 접근 로그 — 기본 off (--verbose 시에만 출력)."""
if not MonitorHandler._access_log:
return
logger.info("%s - %s", self.address_string(), fmt % args)
def log_error(self, fmt: str, *args) -> None:
"""5xx 등 서버 오류만 기록 (favicon 404 제외)."""
msg = fmt % args
if "404" in msg and "File not found" in msg:
return
logger.warning("%s - %s", self.address_string(), msg)
def _send_json(self, payload: dict, *, status: int = 200) -> None:
body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
try:
self.send_response(status)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Cache-Control", "no-store, must-revalidate")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
except _CLIENT_GONE:
logger.debug("client disconnected before response sent")
def _handle_refresh(self) -> None:
try:
self._send_json(_api_refresh())
except _CLIENT_GONE:
pass
except Exception as exc: # noqa: BLE001
if not _client_gone(exc):
self._send_json({"ok": False, "error": str(exc)}, status=500)
def _handle_balance(self) -> None:
try:
self._send_json(_api_balance())
except _CLIENT_GONE:
pass
except Exception as exc: # noqa: BLE001
if not _client_gone(exc):
self._send_json({"ok": False, "error": str(exc)}, status=500)
def end_headers(self) -> None:
p = urlparse(self.path).path
if p.endswith((".json", ".html")) or p in ("", "/"):
self.send_header("Cache-Control", "no-store, must-revalidate")
super().end_headers()
def _request_path(self) -> str:
return urlparse(self.path).path.rstrip("/")
def _handle_json_call(self, fn, *args) -> None:
try:
out = fn(*args)
self._send_json(out, status=200 if out.get("ok") else 400)
except _CLIENT_GONE:
pass
except Exception as exc: # noqa: BLE001
if not _client_gone(exc):
self._send_json({"ok": False, "error": str(exc)}, status=500)
def do_POST(self) -> None:
path = self._request_path()
if path == "/api/refresh":
self._handle_refresh()
return
if path == "/api/rsi/toggle":
self._handle_json_call(_api_rsi_toggle, urlparse(self.path).query)
return
self.send_error(404, "not found")
def do_GET(self) -> None:
path = self._request_path()
if path == "/api/refresh":
self._handle_refresh()
return
if path == "/api/rsi/status":
self._handle_json_call(_api_rsi_status)
return
if path == "/api/candles":
try:
out = _api_candles(urlparse(self.path).query)
self._send_json(out, status=200 if out.get("ok") else 400)
except _CLIENT_GONE:
pass
except Exception as exc: # noqa: BLE001
if not _client_gone(exc):
self._send_json({"ok": False, "error": str(exc)}, status=500)
return
if path == "/api/balance":
self._handle_balance()
return
if path == "/api/chart":
self._serve_chart_json()
return
super().do_GET()
def run_serve(*, access_log: bool = False, quiet: bool = True) -> int:
"""HTTP 서버 기동 (블로킹)."""
from dotenv import load_dotenv
MonitorHandler._access_log = access_log
load_dotenv(_ROOT / ".env", override=False)
port = int(os.environ.get("VOL_MONITOR_PORT", "8766"))
out = _out_dir()
out.mkdir(parents=True, exist_ok=True)
url = f"http://127.0.0.1:{port}/vol_live_monitor.html"
if quiet and not access_log:
print(f"vol monitor {url} (Ctrl+C 종료)", flush=True)
else:
logger.info("모니터: %s", url)
logger.info("출력 디렉터리: %s", out)
try:
server = ThreadingHTTPServer(("127.0.0.1", port), MonitorHandler)
except OSError as exc:
if exc.errno == 48:
logger.error(
"포트 %s 이미 사용 중 — lsof -iTCP:%s -sTCP:LISTEN 후 종료",
port,
port,
)
else:
logger.error("서버 bind 실패: %s", exc)
return 1
try:
server.serve_forever()
except KeyboardInterrupt:
logger.info("종료")
return 0
def main(argv: list[str] | None = None) -> int:
"""CLI — 기본: 갱신 + 서버."""
parser = argparse.ArgumentParser(
description="Bithumb vol_live 모니터 (갱신 + HTTP 서버)",
)
mode = parser.add_mutually_exclusive_group()
mode.add_argument(
"--refresh-only",
action="store_true",
help="JSON/HTML 갱신만 (서버 미기동)",
)
mode.add_argument(
"--serve-only",
action="store_true",
help="HTTP 서버만 (갱신 생략)",
)
parser.add_argument(
"--balance-only",
action="store_true",
help="--refresh-only 와 함께: 잔고 summary만 패치",
)
parser.add_argument(
"-v",
"--verbose",
action="store_true",
help="HTTP 접근·갱신 상세 로그 출력",
)
args = parser.parse_args(argv)
serve_mode = not args.refresh_only
quiet_serve = serve_mode and not args.verbose
log_level = logging.INFO if (args.verbose or args.refresh_only) else logging.WARNING
logging.basicConfig(
level=log_level,
format="%(asctime)s [%(levelname)s] %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
if not args.serve_only:
if args.balance_only:
out = refresh_vol_live_balance()
else:
out = refresh_vol_live_monitor()
if args.verbose or args.refresh_only:
logger.info("refresh done: %s", out)
if not out.get("ok"):
return 1
if args.refresh_only:
return 0
return run_serve(access_log=args.verbose, quiet=quiet_serve)
if __name__ == "__main__":
raise SystemExit(main())