diff --git a/.env.example b/.env.example index 25667e3..f6cbf75 100644 --- a/.env.example +++ b/.env.example @@ -136,6 +136,8 @@ OPS_SIZING_RULES_JSON=data/spot/operations/sizing_rules.json # fractal_swing(3분)과 별도 프로세스 — 동시 live 시 KRW 경합 주의 OPS_SYMBOLS=TRX,NEAR,WLD VOL_STATE_JSON=data/spot/operations/vol_breakout_state.json +# 모니터 차트 표시 종목 (매매 OPS_SYMBOLS와 분리). 비우면 OPS_SYMBOLS +VOL_MONITOR_SYMBOLS=XRP,TRX,WLD,SOL,ETH,ADA,SUI VOL_LOOKBACK=14 VOL_ATR_MULT=2.0 VOL_LOOKBACK_DAYS=60 @@ -162,6 +164,28 @@ VOL_MONITOR_DAYS=14 # cron (vol_breakout): bash scripts/install_crontab.sh --apply # BITHUMB_PYTHON=/Users/dsyoon/opt/anaconda3/envs/coin/bin/python3 +# --- RSI DCA 정액 매수 (15m RSI 30/35 상향 돌파 · 매도 없음 · 2026-09-07) --- +# 모드는 OPS_MODE와 분리. live 전환: RSI_DCA_MODE=live (실제 주문 발생) +RSI_DCA_MODE=paper +RSI_DCA_SYMBOLS=XRP,TRX,WLD,SOL,ETH,ADA,SUI +RSI_DCA_INTERVAL_MIN=15 +RSI_DCA_PERIOD=14 +# level:원 — 30 상향 돌파 2만원, 35 상향 돌파 1만원 (둘 다 발생 가능) +RSI_DCA_LEVELS=30:20000,35:10000 +# 종목별 기준선 오버라이드 (없는 종목은 RSI_DCA_LEVELS). 예: XRP=19:10000;TRX=32:10000 +RSI_DCA_LEVELS_BY_SYMBOL= +# 일(KST) 총 매수 상한 — 초과하는 주문은 스킵 +RSI_DCA_DAILY_MAX_KRW=60000 +RSI_DCA_LOOKBACK_DAYS=20 +RSI_DCA_MAX_BARS_PER_TICK=8 +# 봉 마감 후 N분 지난 신호는 매수하지 않음 (장애 복구 시 몰아 매수 방지) +RSI_DCA_MAX_SIGNAL_AGE_MIN=45 +RSI_DCA_STATE_JSON=data/spot/operations/rsi_dca_state.json +RSI_DCA_REPORT_JSON=docs/spot/3_operations/rsi_dca_report.json +RSI_DCA_TICK_LOCK_PATH=data/spot/operations/rsi.tick.lock +# 파일이 존재하면 신규 매수 차단: touch data/spot/operations/rsi.kill +RSI_DCA_KILL_SWITCH_PATH=data/spot/operations/rsi.kill + # 폴더 구조: data|docs / {common, spot} # common — coins.db 등 공유 리소스 # spot — 현물 GT·기법·분석·운영 diff --git a/.gitignore b/.gitignore index 91ec39c..a124413 100644 --- a/.gitignore +++ b/.gitignore @@ -114,3 +114,7 @@ ENV/ # Rope project settings .ropeproject + +# 로컬 설정 백업 (API 키 포함) — 커밋 금지 +.env.bak* +.env.local.backup diff --git a/README.md b/README.md index 1d594ca..2ce0014 100644 --- a/README.md +++ b/README.md @@ -302,6 +302,7 @@ crontab -l # 확인 | 모니터 JSON | 5분 | `3_run_vol_monitor_cron.sh` | `data/spot/operations/vol_monitor_cron.log` | - hung 프로세스: 다운로드 20분·vol tick 10분 초과 시 자동 종료 후 lock 정리 (`scripts/_cron_env.sh`) +- **공백 자동 백필**: 절전·재부팅 복귀 시 `00_run_download_cron.sh`가 `00_backfill_gaps.py`로 증분 범위(200봉)를 넘는 (심볼, 분봉)만 감지해 `--full --days N`으로 채운 뒤 증분 수집을 진행 (`--dry-run --assume-now`로 계획 확인) - Python: `coin` / `ncue` conda 우선. 다른 환경이면 `.env` 또는 crontab에 `BITHUMB_PYTHON=...` 설정 - 모니터 UI (8766): ```bash @@ -325,6 +326,30 @@ bash scripts/00_run_download_cron.sh # 수동 1회 bash scripts/3_run_vol_breakout_cron.sh ``` +### RSI DCA 정액 매수 (15m RSI 30/35 상향 돌파 · 매도 없음) + +7종(XRP, TRX, WLD, SOL, ETH, ADA, SUI) 15분봉 RSI(14)가 **종가 확정 기준**으로 기준선을 상향 돌파하면 정액 매수한다. 자동 매도는 없다(보유분 수동 관리). + +| 규칙 | 값 | 설정 | +|------|----|------| +| 30 상향 돌파 | 20,000원 매수 | `RSI_DCA_LEVELS=30:20000,35:10000` | +| 35 상향 돌파 | 10,000원 매수 (30 돌파와 독립, 둘 다 발생 가능) | 〃 | +| 일(KST) 매수 상한 | 60,000원 — 초과 주문은 스킵 | `RSI_DCA_DAILY_MAX_KRW` | +| 신호 유효 시간 | 봉 마감 후 45분 (장애 복구 시 몰아 매수 방지) | `RSI_DCA_MAX_SIGNAL_AGE_MIN` | +| 모드 | `RSI_DCA_MODE=paper` 기본. **live는 OPS_MODE와 별개**로 명시 전환 | `RSI_DCA_MODE` | +| 킬스위치 | `data/spot/operations/rsi.kill` 존재 시 신규 매수 차단 | `RSI_DCA_KILL_SWITCH_PATH` | + +```bash +python scripts/3_run_rsi_dca_backtest.py --days 90 # 신호·투입·평가 (docs/spot/3_operations/rsi_dca_backtest.json) +python scripts/3_run_rsi_dca.py # paper 1회 tick (첫 실행은 커서 초기화만, 매수 없음) +python scripts/3_run_rsi_dca.py --status # 상태 확인 +touch data/spot/operations/rsi.kill # 긴급 차단 +``` + +cron(1분): `scripts/crontab.bithumb.example`의 `3_run_rsi_dca_cron.sh` 줄 주석 해제 후 `bash scripts/install_crontab.sh --apply`. +live 전환은 `.env`에서 `RSI_DCA_MODE=live` 로 바꾼 뒤 cron 또는 `python scripts/3_run_rsi_dca.py --mode live` 를 **사용자가 직접** 실행한다. +체결·상태: `data/spot/operations/rsi_dca_state.json`, `docs/spot/3_operations/rsi_dca_report.json`, 텔레그램 알림. 모니터 차트(8766)에 매수 마커 표시. + ### fractal watch 점검 ```bash @@ -502,6 +527,7 @@ OPS_DAILY_MAX_TRADES=20 ## 변경 이력 +- **2026-09-07:** 수집 대상 7종(XRP,TRX,WLD,SOL,ETH,ADA,SUI)·`VOL_MONITOR_SYMBOLS` 분리, 모니터 분봉 탭(`/api/candles`)·RSI(14) 패널, RSI DCA 정액 매수 전략(`rsi_dca_engine/runner`, 백테스트, cron 래퍼) 추가 - **2026-06-14:** ledger pending, exchange reconcile, max_age backlog, watch 5분 감시·조치, ops.tick.lock, README 전면 갱신 - **2026-06-13:** fractal_swing live — 슬리피지·sync·tail·텔레그램; ops_default sim **+1,873,140%** - **2026-06-13:** 프로젝트명 Bithumb, 선물 파이프라인 제거 diff --git a/scripts/00_backfill_gaps.py b/scripts/00_backfill_gaps.py new file mode 100755 index 0000000..5a956a4 --- /dev/null +++ b/scripts/00_backfill_gaps.py @@ -0,0 +1,100 @@ +#!/usr/bin/env python3 +"""절전·재부팅 후 캔들 공백 자동 백필 (cron 수집 전 단계). + +증분 수집(최신 200봉)으로 못 채우는 공백이 있는 (심볼, 인터벌)만 골라 +00_download_candles.py --full --days N 으로 채운다. 공백이 없으면 즉시 종료(약 0.1초). + + python scripts/00_backfill_gaps.py # 점검 후 필요 시 백필 + python scripts/00_backfill_gaps.py --dry-run # 계획만 출력 + python scripts/00_backfill_gaps.py --dry-run --assume-now "2026-09-12 09:00:00" +""" + +from __future__ import annotations + +import argparse +import logging +import subprocess +import sys +from datetime import datetime +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +SRC = ROOT / "src" +if str(SRC) not in sys.path: + sys.path.insert(0, str(SRC)) + +from bithumb.config import load_settings # noqa: E402 +from bithumb.data.candle_store import CandleStore # noqa: E402 +from bithumb.data.gap_backfill import plan_backfill # noqa: E402 + +logger = logging.getLogger("backfill") + + +def main() -> int: + parser = argparse.ArgumentParser(description="캔들 공백 자동 백필") + parser.add_argument("--dry-run", action="store_true", help="계획만 출력") + parser.add_argument("--assume-now", default=None, help="테스트용 기준 시각 'YYYY-MM-DD HH:MM:SS'") + parser.add_argument("--max-days", type=int, default=30, help="백필 상한 일수 (기본 30)") + parser.add_argument("--safety", type=float, default=0.9, help="증분 커버리지 안전계수 (기본 0.9)") + parser.add_argument("-v", "--verbose", action="store_true") + args = parser.parse_args() + logging.basicConfig( + level=logging.DEBUG if args.verbose else logging.INFO, + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", + ) + + settings = load_settings() + now = datetime.strptime(args.assume_now, "%Y-%m-%d %H:%M:%S") if args.assume_now else datetime.now() + symbols = list(settings.download_symbols) + intervals = list(settings.download_intervals) + + store = CandleStore(settings.db_path) + try: + db_max = { + sym: {iv: store.get_max_datetime(sym, iv) for iv in intervals} + for sym in symbols + } + finally: + store.close() + + plan = plan_backfill( + db_max, now=now, intervals=intervals, + batch_size=settings.candle_count, safety=args.safety, max_days=args.max_days, + ) + if plan.empty: + logger.info("공백 없음 (증분 수집 범위 내) — 백필 생략") + return 0 + + logger.warning("캔들 공백 감지 → 백필 계획: %s", plan.describe()) + if args.dry_run: + return 0 + + # 인터벌별 필요 일수가 다르므로 일수별로 묶어 실행 (요청 수 최소화) + by_days: dict[int, dict[str, set[int]]] = {} + for sym, m in plan.needs.items(): + for iv, days in m.items(): + by_days.setdefault(days, {}).setdefault(sym, set()).add(iv) + + python = sys.executable + rc_all = 0 + for days in sorted(by_days): + sym_ivs = by_days[days] + syms = ",".join(sorted(sym_ivs)) + ivs = ",".join(str(i) for i in sorted({iv for s in sym_ivs.values() for iv in s})) + cmd = [ + python, str(ROOT / "scripts" / "00_download_candles.py"), + "--full", "--days", str(days), "--symbols", syms, "--intervals", ivs, + ] + logger.warning("백필 실행: --full --days %s --symbols %s --intervals %s", days, syms, ivs) + rc = subprocess.call(cmd, cwd=str(ROOT)) + if rc != 0: + logger.error("백필 실패 rc=%s (days=%s)", rc, days) + rc_all = rc + if rc_all == 0: + logger.warning("백필 완료: %s", plan.describe()) + return rc_all + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/00_run_download_cron.sh b/scripts/00_run_download_cron.sh index 3f4c9ce..d113b80 100755 --- a/scripts/00_run_download_cron.sh +++ b/scripts/00_run_download_cron.sh @@ -6,10 +6,27 @@ source "$(dirname "$0")/_cron_env.sh" ensure_cron_log_dir "data/common" LOCKDIR="data/common/download.lock.d" -# 3종목×11 TF 증분 — 20분 초과 시 hung 으로 간주 -if ! acquire_cron_lock "$LOCKDIR" "scripts/00_download.py" 1200; then +DL_SCRIPT="${CRON_PROJECT_ROOT}/scripts/00_download.py" +# 7종목 증분 — 20분 초과 시 hung 으로 간주. 패턴은 이 프로젝트의 절대 경로로 한정 +# (다른 프로젝트(Binance 등)의 scripts/00_download.py 와 pgrep/pkill 이 섞이지 않도록) +if ! acquire_cron_lock "$LOCKDIR" "$DL_SCRIPT" 1200; then exit 0 fi PYTHON="$(resolve_bithumb_python)" || exit 1 -"$PYTHON" scripts/00_download.py "$@" + +# 절전·재부팅 복귀 시: 증분 범위(200봉)를 넘는 공백이 있으면 필요한 분봉·일수만 먼저 백필 (공백 없으면 ~0.1초) +"$PYTHON" "${CRON_PROJECT_ROOT}/scripts/00_backfill_gaps.py" || echo "$(date '+%Y-%m-%d %H:%M:%S') [WARN] 공백 백필 실패 — 증분 수집은 계속 진행" >&2 + +# 인자가 없으면: 매분 핵심 분봉(1,3,5,15)만, 5분 배수 분에는 DOWNLOAD_INTERVALS 전체. +# 1분봉 전략의 캔들 지연을 줄이기 위한 분할 (환경변수 DOWNLOAD_CORE_INTERVALS 로 조정) +if [ "$#" -eq 0 ]; then + minute="$(date +%M)" + if [ $((10#$minute % 5)) -ne 0 ]; then + CORE="${DOWNLOAD_CORE_INTERVALS:-1,3,5,15}" + echo "$(date '+%Y-%m-%d %H:%M:%S') [INFO] core intervals only: ${CORE}" + "$PYTHON" "$DL_SCRIPT" --intervals "$CORE" + exit 0 # 핵심 분봉 회차 종료 (EXIT 트랩이 lock 정리) + fi +fi +"$PYTHON" "$DL_SCRIPT" "$@" diff --git a/scripts/3_go_live_rsi_dca.sh b/scripts/3_go_live_rsi_dca.sh new file mode 100755 index 0000000..1e9504e --- /dev/null +++ b/scripts/3_go_live_rsi_dca.sh @@ -0,0 +1,64 @@ +#!/usr/bin/env bash +# RSI DCA 정액 매수 live 전환/복귀 — 사용자가 직접 실행하는 스위치. +# bash scripts/3_go_live_rsi_dca.sh # live 전환: .env RSI_DCA_MODE=live + cron 줄 활성화 + crontab 적용 +# bash scripts/3_go_live_rsi_dca.sh --paper # paper 복귀 (cron 유지, 주문 없음) +# bash scripts/3_go_live_rsi_dca.sh --dry-run # 변경 없이 수행 내용만 출력 +set -euo pipefail +ROOT="$(cd "$(dirname "$0")/.." && pwd)" +ENV="${ROOT}/.env" +EXAMPLE="${ROOT}/scripts/crontab.bithumb.example" +RSI_LINE_RE='3_run_rsi_dca_cron\.sh' + +MODE="live" +DRY=0 +for a in "$@"; do + case "$a" in + --paper) MODE="paper" ;; + --live) MODE="live" ;; + --dry-run) DRY=1 ;; + -h|--help) sed -n 2,5p "$0"; exit 0 ;; + *) echo "unknown arg: $a" >&2; exit 1 ;; + esac +done + +[ -f "$ENV" ] || { echo ".env 없음: $ENV" >&2; exit 1; } +cur="$(grep -E '^RSI_DCA_MODE=' "$ENV" | tail -1 | cut -d= -f2 | tr -d '[:space:]' || true)" +echo "현재 RSI_DCA_MODE=${cur:-<미설정>} → 목표 ${MODE}" +echo "대상: $(grep -E '^RSI_DCA_SYMBOLS=' "$ENV" | cut -d= -f2) | 규칙: $(grep -E '^RSI_DCA_LEVELS=' "$ENV" | cut -d= -f2) | 일 상한: $(grep -E '^RSI_DCA_DAILY_MAX_KRW=' "$ENV" | cut -d= -f2)원" +if [ "$MODE" = "live" ] && [ -f "${ROOT}/data/spot/operations/rsi.kill" ]; then + echo "주의: 킬스위치 파일(data/spot/operations/rsi.kill)이 있어 매수가 차단됩니다. 해제: rm data/spot/operations/rsi.kill" +fi + +if [ "$DRY" = "1" ]; then + echo "[dry-run] 1) .env RSI_DCA_MODE=${MODE} 로 변경" + echo "[dry-run] 2) ${EXAMPLE} 의 RSI cron 줄 주석 해제" + echo "[dry-run] 3) bash scripts/install_crontab.sh --apply" + exit 0 +fi + +# 1) .env 모드 +if grep -qE '^RSI_DCA_MODE=' "$ENV"; then + sed -i '' -E "s/^RSI_DCA_MODE=.*/RSI_DCA_MODE=${MODE}/" "$ENV" +else + printf '\nRSI_DCA_MODE=%s\n' "$MODE" >> "$ENV" +fi +echo "1) .env: $(grep -E '^RSI_DCA_MODE=' "$ENV")" + +# 2) cron 예시에서 RSI 줄 활성화 (이미 활성화면 그대로) +if grep -qE "^# \* \* \* \* \* .*${RSI_LINE_RE}" "$EXAMPLE"; then + sed -i '' -E "s|^# (\* \* \* \* \* .*${RSI_LINE_RE}.*)$|\1|" "$EXAMPLE" +fi +echo "2) cron 예시: $(grep -E "${RSI_LINE_RE}" "$EXAMPLE" | grep -vE '^#' | head -1 | cut -c1-70)…" + +# 3) crontab 적용 (BITHUMB 블록 갱신 — 캔들 수집·모니터·RSI tick) +bash "${ROOT}/scripts/install_crontab.sh" --apply +echo "3) crontab RSI 줄:"; crontab -l | grep -E "${RSI_LINE_RE}" || echo " (없음)" + +echo +if [ "$MODE" = "live" ]; then + echo "완료: 다음 1분 tick부터 실거래 매수. 첫 tick은 미초기화 종목 커서 설정만 하고, 이후 봉부터 매수합니다." + echo "긴급 차단: touch ${ROOT}/data/spot/operations/rsi.kill | paper 복귀: bash scripts/3_go_live_rsi_dca.sh --paper" + echo "확인: python scripts/3_run_rsi_dca.py --status / tail -f data/spot/operations/rsi_dca_cron.log" +else + echo "완료: paper 모드. cron은 유지되며 주문은 나가지 않습니다." +fi diff --git a/scripts/3_run_rsi_dca.py b/scripts/3_run_rsi_dca.py new file mode 100755 index 0000000..9bcc248 --- /dev/null +++ b/scripts/3_run_rsi_dca.py @@ -0,0 +1,106 @@ +#!/usr/bin/env python3 +"""RSI DCA 현물 정액 매수 tick (15m RSI 30/35 상향 돌파, 매도 없음). + + python scripts/3_run_rsi_dca.py # RSI_DCA_MODE(.env, 기본 paper) 1회 tick + python scripts/3_run_rsi_dca.py --mode paper --loop 60 + python scripts/3_run_rsi_dca.py --mode live # 실제 주문 — 사용자 직접 실행 + python scripts/3_run_rsi_dca.py --status # 상태만 출력 +""" + +from __future__ import annotations + +import argparse +import json +import logging +import sys +import time +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +SRC = ROOT / "src" +if str(SRC) not in sys.path: + sys.path.insert(0, str(SRC)) + +from bithumb.config import load_settings # noqa: E402 +from bithumb.operations.rsi_dca_engine import load_state # noqa: E402 +from bithumb.operations.rsi_dca_runner import RsiDcaRunner # noqa: E402 + + +def _configure_logging(verbose: bool) -> None: + logging.basicConfig( + level=logging.DEBUG if verbose else logging.INFO, + format="%(asctime)s [%(levelname)s] %(message)s", + datefmt="%Y-%m-%d %H:%M:%S", + ) + + +def _print_status(settings, mode: str) -> None: + st = load_state(settings.rsi_dca_state_json, mode) + d = st.get("daily") or {} + t = st.get("totals") or {} + print(f"strategy={st.get('strategy')} mode={st.get('mode')} last_run={st.get('last_run_at')}") + print(f"daily {d.get('date')} spent={float(d.get('spent_krw') or 0):,.0f}/{settings.rsi_dca_daily_max_krw:,.0f} count={d.get('count')}") + print(f"totals spent={float(t.get('spent_krw') or 0):,.0f} count={t.get('count')} trades={len(st.get('trades') or [])}") + for sym, s in (st.get("symbols") or {}).items(): + rsi = s.get("last_rsi") + print(f" {sym:4s} init={s.get('initialized')} cursor={s.get('last_confirm_time')} " + f"rsi={rsi if rsi is None else round(rsi, 1)} buys={s.get('buy_count')} spent={float(s.get('spent_krw') or 0):,.0f}") + + +def main() -> int: + parser = argparse.ArgumentParser(description="Bithumb RSI DCA 정액 매수 tick") + parser.add_argument("--mode", choices=("paper", "live"), default=None, + help="기본: .env RSI_DCA_MODE (paper)") + parser.add_argument("--loop", type=int, default=0, metavar="SEC") + parser.add_argument("--status", action="store_true", help="상태 출력만") + parser.add_argument("-v", "--verbose", action="store_true") + args = parser.parse_args() + _configure_logging(args.verbose) + + settings = load_settings() + mode = (args.mode or settings.rsi_dca_mode or "paper").lower() + if args.status: + _print_status(settings, mode) + return 0 + + if mode == "live": + print("경고: live — 실제 주문 발생 (RSI 정액 매수, 매도 없음)") + levels = ", ".join(f"RSI {lv:g}↑ {krw:,.0f}원" for lv, krw in settings.rsi_dca_levels) + print( + f"rsi_dca {mode} | symbols={settings.rsi_dca_symbols} | {settings.rsi_dca_interval_min}m " + f"RSI({settings.rsi_dca_period}) | {levels} | 일 상한 {settings.rsi_dca_daily_max_krw:,.0f}원" + ) + + def _once() -> dict: + runner = RsiDcaRunner(settings, mode=mode) + report = runner.tick() + if not report.get("ok"): + print(f" skip: {report.get('note')}") + return report + d = report.get("daily") or {} + for row in report.get("results") or []: + if "error" in row: + print(f" {row['symbol']}: ERROR {row['error']}") + continue + rsi = row.get("rsi") + print(f" {row['symbol']:4s} rsi={'-' if rsi is None else f'{rsi:.1f}':>5s} fills={row['fills']} {row['note']}") + print(f" daily {d.get('date')} spent={float(d.get('spent_krw') or 0):,.0f} " + f"remaining={float(d.get('remaining_krw') or 0):,.0f} kill_switch={report.get('kill_switch')}") + return report + + if args.loop <= 0: + _once() + return 0 + while True: + try: + _once() + except KeyboardInterrupt: + print("\n종료") + return 0 + except Exception: # noqa: BLE001 + logging.exception("rsi_dca loop tick failed") + time.sleep(args.loop) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/3_run_rsi_dca_backtest.py b/scripts/3_run_rsi_dca_backtest.py new file mode 100755 index 0000000..f50c04f --- /dev/null +++ b/scripts/3_run_rsi_dca_backtest.py @@ -0,0 +1,86 @@ +#!/usr/bin/env python3 +"""RSI DCA 백테스트 — DB 15m 캔들로 최근 N일 신호·매수·평가 산출. + + python scripts/3_run_rsi_dca_backtest.py --days 90 + python scripts/3_run_rsi_dca_backtest.py --days 90 --levels 30:20000,35:10000 --daily-max 60000 +""" + +from __future__ import annotations + +import argparse +import json +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +SRC = ROOT / "src" +if str(SRC) not in sys.path: + sys.path.insert(0, str(SRC)) + +from bithumb.config import load_settings # noqa: E402 +from bithumb.data.candle_loader import load_candles # noqa: E402 +from bithumb.operations.rsi_dca_engine import ( # noqa: E402 + RsiDcaConfig, + backtest_rsi_dca, + parse_levels, +) +from bithumb.operations.rsi_dca_runner import config_from_settings # noqa: E402 + + +def main() -> int: + parser = argparse.ArgumentParser(description="RSI DCA 백테스트") + parser.add_argument("--days", type=int, default=90) + parser.add_argument("--symbols", default=None, help="쉼표 구분 (기본 RSI_DCA_SYMBOLS)") + parser.add_argument("--levels", default=None, help="예 30:20000,35:10000") + parser.add_argument("--daily-max", type=float, default=None) + parser.add_argument("--out", default="docs/spot/3_operations/rsi_dca_backtest.json") + args = parser.parse_args() + + settings = load_settings() + base = config_from_settings(settings) + cfg = RsiDcaConfig( + symbols=[s.strip().upper() for s in args.symbols.split(",")] if args.symbols else base.symbols, + interval_min=base.interval_min, + period=base.period, + levels=parse_levels(args.levels) if args.levels else base.levels, + daily_max_krw=args.daily_max if args.daily_max is not None else base.daily_max_krw, + lookback_days=base.lookback_days, + max_bars_per_tick=base.max_bars_per_tick, + max_signal_age_min=base.max_signal_age_min, + min_order_krw=base.min_order_krw, + fee_rate=base.fee_rate, + slippage_rate=base.slippage_rate, + fee_lock_rate=base.fee_lock_rate, + ) + + candles = {} + for sym in cfg.symbols: + # RSI 워밍업을 위해 여유 있게 로드 + df = load_candles(settings.db_path, sym, cfg.interval_min, lookback_days=args.days + 10) + candles[sym] = df + rep = backtest_rsi_dca(candles, cfg, days=args.days) + + levels_txt = ", ".join(f"RSI {lv:g}↑ {krw:,.0f}원" for lv, krw in cfg.levels) + print(f"RSI DCA 백테스트 · 최근 {args.days}일 · {cfg.interval_min}m RSI({cfg.period}) · {levels_txt} · 일 상한 {cfg.daily_max_krw:,.0f}원") + print(f"수수료 {cfg.fee_rate*100:.3f}% · 슬리피지 {cfg.slippage_rate*100:.3f}% · 매도 없음(마지막 종가 평가)") + print() + print(f"{'종목':5s} {'신호':>5s} {'매수':>5s} {'투입(원)':>12s} {'평가(원)':>12s} {'손익%':>7s}") + for sym in cfg.symbols: + v = rep["per_symbol"].get(sym) + if not v: + print(f"{sym:5s} {'데이터없음':>5s}") + continue + print(f"{sym:5s} {int(v['signals']):5d} {int(v['buys']):5d} {v['spent_krw']:12,.0f} {v['value_krw']:12,.0f} {v['pnl_pct']:7.2f}") + print("-" * 52) + print(f"{'합계':5s} {rep['signals']:5d} {rep['buys']:5d} {rep['total_spent_krw']:12,.0f} {rep['total_value_krw']:12,.0f} {rep['total_pnl_pct']:7.2f}") + print(f"일 상한으로 스킵된 신호: {rep['skipped_daily_cap']}건 · 일평균 투입 {rep['avg_daily_spent_krw']:,.0f}원 (신호 구간 {rep['span_days']}일)") + + out = ROOT / args.out + out.parent.mkdir(parents=True, exist_ok=True) + out.write_text(json.dumps(rep, ensure_ascii=False, indent=2, default=str), encoding="utf-8") + print(f"\n저장: {out}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/3_run_rsi_dca_cron.sh b/scripts/3_run_rsi_dca_cron.sh new file mode 100755 index 0000000..30bbf44 --- /dev/null +++ b/scripts/3_run_rsi_dca_cron.sh @@ -0,0 +1,18 @@ +#!/usr/bin/env bash +# RSI DCA 정액 매수 tick (cron 1분). 모드는 .env RSI_DCA_MODE (기본 paper). +set -euo pipefail +# shellcheck source=scripts/_cron_env.sh +source "$(dirname "$0")/_cron_env.sh" + +ensure_cron_log_dir "data/spot/operations" +LOCKDIR="data/spot/operations/rsi.tick.lock.d" +RSI_SCRIPT="${CRON_PROJECT_ROOT}/scripts/3_run_rsi_dca.py" +if ! acquire_cron_lock "$LOCKDIR" "$RSI_SCRIPT" 600; then + exit 0 +fi + +PYTHON="$(resolve_bithumb_python)" || exit 1 +# 같은 분에 시작하는 캔들 수집(핵심 분봉 12~21초)이 끝난 뒤 판정하도록 지연 → 봉 마감 후 약 1.5분 내 매수 +sleep "${RSI_DCA_TICK_DELAY_SEC:-25}" +echo "$(date '+%Y-%m-%d %H:%M:%S') tick" +"$PYTHON" "$RSI_SCRIPT" "$@" diff --git a/scripts/3_run_rsi_dca_interval_sweep.py b/scripts/3_run_rsi_dca_interval_sweep.py new file mode 100755 index 0000000..a5c5aac --- /dev/null +++ b/scripts/3_run_rsi_dca_interval_sweep.py @@ -0,0 +1,154 @@ +#!/usr/bin/env python3 +"""RSI DCA 인터벌 비교 실험 — 동일 규칙(30↑ 2만원, 35↑ 1만원, 일 상한)으로 분봉별 수익률 비교. + + python scripts/3_run_rsi_dca_interval_sweep.py --days 90 + python scripts/3_run_rsi_dca_interval_sweep.py --days 90 --intervals 1,3,5,10,15,30,60,240,1440 +""" + +from __future__ import annotations + +import argparse +import json +import math +import sys +from pathlib import Path + +import pandas as pd + +ROOT = Path(__file__).resolve().parents[1] +SRC = ROOT / "src" +if str(SRC) not in sys.path: + sys.path.insert(0, str(SRC)) + +from bithumb.config import load_settings # noqa: E402 +from bithumb.data.candle_loader import load_candles # noqa: E402 +from bithumb.operations.rsi_dca_engine import ( # noqa: E402 + RsiDcaConfig, + backtest_rsi_dca, + parse_levels, +) +from bithumb.operations.rsi_dca_runner import config_from_settings # noqa: E402 + +LABEL = {1: "1분", 3: "3분", 5: "5분", 10: "10분", 15: "15분", 30: "30분", 60: "1시간", 240: "4시간", 1440: "1일"} + + +def _with(cfg: RsiDcaConfig, **kw) -> RsiDcaConfig: + d = {k: getattr(cfg, k) for k in cfg.__dataclass_fields__} + d.update(kw) + return RsiDcaConfig(**d) + + +def daily_dca_benchmark( + daily_by_symbol: dict[str, pd.DataFrame], + *, + days: int, + daily_krw: float, + fee_rate: float, + slippage_rate: float, +) -> dict: + """기준선: 매일 종가에 일 예산을 종목 수로 균등 분할 매수 (매도 없음).""" + syms = [s for s, df in daily_by_symbol.items() if df is not None and not df.empty] + if not syms: + return {} + per = daily_krw / len(syms) + spent = 0.0 + value = 0.0 + buys = 0 + for sym in syms: + d = daily_by_symbol[sym].copy() + d["datetime"] = pd.to_datetime(d["datetime"]) + d = d.sort_values("datetime") + start = d["datetime"].max() - pd.Timedelta(days=days) + w = d[d["datetime"] >= start] + last = float(d["close"].iloc[-1]) + coin = 0.0 + for px in w["close"].astype(float): + fill = px * (1.0 + slippage_rate) + coin += per * (1.0 - fee_rate) / fill + spent += per + buys += 1 + value += coin * last + return { + "label": "매일 정액 분할(기준선)", + "buys": buys, + "total_spent_krw": spent, + "total_value_krw": value, + "total_pnl_pct": (value / spent - 1.0) * 100.0 if spent else 0.0, + "avg_daily_spent_krw": daily_krw, + } + + +def main() -> int: + parser = argparse.ArgumentParser(description="RSI DCA 인터벌 비교") + parser.add_argument("--days", type=int, default=90) + parser.add_argument("--intervals", default="1,3,5,10,15,30,60,240,1440") + parser.add_argument("--levels", default=None) + parser.add_argument("--daily-max", type=float, default=None) + parser.add_argument("--out", default="docs/spot/3_operations/rsi_dca_interval_sweep.json") + args = parser.parse_args() + + settings = load_settings() + base = config_from_settings(settings) + if args.levels: + base = _with(base, levels=parse_levels(args.levels)) + if args.daily_max is not None: + base = _with(base, daily_max_krw=args.daily_max) + intervals = [int(x) for x in args.intervals.split(",") if x.strip()] + + results = [] + for iv in intervals: + cfg = _with(base, interval_min=iv) + warm_days = math.ceil(cfg.period * 3 * iv / 1440) + 2 + candles = {} + for sym in cfg.symbols: + candles[sym] = load_candles(settings.db_path, sym, iv, lookback_days=args.days + warm_days) + rep = backtest_rsi_dca(candles, cfg, days=args.days) + bars = sum(len(df) for df in candles.values() if df is not None) + rep["interval_min"] = iv + rep["label"] = LABEL.get(iv, f"{iv}분") + rep["bars_loaded"] = bars + rep.pop("trades", None) + results.append(rep) + print(f" {rep['label']:>4s} 완료 · 신호 {rep['signals']:,} · 매수 {rep['buys']} · 손익 {rep['total_pnl_pct']:+.2f}%", flush=True) + + daily = {s: load_candles(settings.db_path, s, 1440, lookback_days=args.days + 5) for s in base.symbols} + bench = daily_dca_benchmark( + daily, days=args.days, daily_krw=base.daily_max_krw, + fee_rate=base.fee_rate, slippage_rate=base.slippage_rate, + ) + + levels_txt = ", ".join(f"RSI {lv:g}↑ {krw:,.0f}원" for lv, krw in base.levels) + print() + print(f"RSI DCA 인터벌 비교 · 최근 {args.days}일 · {len(base.symbols)}종 · RSI({base.period}) · {levels_txt} · 일 상한 {base.daily_max_krw:,.0f}원") + print(f"수수료 {base.fee_rate*100:.3f}% · 슬리피지 {base.slippage_rate*100:.3f}% · 매도 없음 · 마지막 종가 평가") + print() + hdr = f"{'인터벌':>6s} {'신호':>7s} {'매수':>5s} {'상한스킵':>7s} {'투입(원)':>11s} {'평가(원)':>11s} {'손익%':>7s} {'손익(원)':>10s} {'일평균투입':>9s}" + print(hdr) + for r in sorted(results, key=lambda x: -x["total_pnl_pct"]): + pnl = r["total_value_krw"] - r["total_spent_krw"] + print(f"{r['label']:>6s} {r['signals']:7,d} {r['buys']:5d} {r['skipped_daily_cap']:7,d} {r['total_spent_krw']:11,.0f} {r['total_value_krw']:11,.0f} {r['total_pnl_pct']:7.2f} {pnl:10,.0f} {r['avg_daily_spent_krw']:9,.0f}") + if bench: + pnl = bench["total_value_krw"] - bench["total_spent_krw"] + print("-" * len(hdr)) + print(f"{'기준선':>6s} {'-':>7s} {bench['buys']:5d} {'-':>7s} {bench['total_spent_krw']:11,.0f} {bench['total_value_krw']:11,.0f} {bench['total_pnl_pct']:7.2f} {pnl:10,.0f} {bench['avg_daily_spent_krw']:9,.0f} (매일 종가 6만원 균등 분할)") + + best = max(results, key=lambda x: x["total_pnl_pct"]) + print(f"\n최고 수익률: {best['label']} ({best['total_pnl_pct']:+.2f}%)") + print("종목별 (최고 인터벌):") + for sym in base.symbols: + v = best["per_symbol"].get(sym) + if v: + print(f" {sym:4s} 신호 {int(v['signals']):5d} 매수 {int(v['buys']):4d} 투입 {v['spent_krw']:10,.0f} 평가 {v['value_krw']:10,.0f} {v['pnl_pct']:+7.2f}%") + + out = ROOT / args.out + out.parent.mkdir(parents=True, exist_ok=True) + out.write_text(json.dumps({ + "days": args.days, "levels": base.levels, "daily_max_krw": base.daily_max_krw, + "symbols": base.symbols, "results": results, "benchmark_daily_dca": bench, + }, ensure_ascii=False, indent=2, default=str), encoding="utf-8") + print(f"\n저장: {out}") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/3_run_vol_monitor.py b/scripts/3_run_vol_monitor.py index 6a9c984..78ceaca 100755 --- a/scripts/3_run_vol_monitor.py +++ b/scripts/3_run_vol_monitor.py @@ -22,14 +22,16 @@ import threading import time from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path -from urllib.parse import urlparse +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, @@ -58,6 +60,7 @@ def refresh_vol_live_monitor(*, write_html: bool = True) -> dict: 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) @@ -108,8 +111,48 @@ def _api_balance() -> dict: 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/balance · /api/refresh.""" + """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 @@ -195,18 +238,32 @@ class MonitorHandler(SimpleHTTPRequestHandler): self._send_json({"ok": False, "error": str(exc)}, status=500) def end_headers(self) -> None: - if self.path.endswith(".json"): + 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: @@ -214,6 +271,19 @@ class MonitorHandler(SimpleHTTPRequestHandler): 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 diff --git a/scripts/3_run_vol_monitor_cron.sh b/scripts/3_run_vol_monitor_cron.sh index 0d1aa25..cc0304a 100755 --- a/scripts/3_run_vol_monitor_cron.sh +++ b/scripts/3_run_vol_monitor_cron.sh @@ -6,9 +6,10 @@ source "$(dirname "$0")/_cron_env.sh" ensure_cron_log_dir "docs/spot/3_operations" LOCKDIR="data/spot/operations/vol.monitor.lock.d" -if ! acquire_cron_lock "$LOCKDIR" "scripts/3_run_vol_monitor.py" 300; then +MON_SCRIPT="${CRON_PROJECT_ROOT}/scripts/3_run_vol_monitor.py" +if ! acquire_cron_lock "$LOCKDIR" "$MON_SCRIPT --refresh-only" 300; then exit 0 fi PYTHON="$(resolve_bithumb_python)" || exit 1 -"$PYTHON" scripts/3_run_vol_monitor.py --refresh-only "$@" +"$PYTHON" "$MON_SCRIPT" --refresh-only "$@" diff --git a/scripts/crontab.bithumb.example b/scripts/crontab.bithumb.example index 29654de..eb678ac 100644 --- a/scripts/crontab.bithumb.example +++ b/scripts/crontab.bithumb.example @@ -7,17 +7,20 @@ # 비활성화(주석): bash scripts/install_crontab.sh --disable # 제거: bash scripts/install_crontab.sh --remove -# 캔들 증분 (TRX,NEAR,WLD × DOWNLOAD_INTERVALS) — 매 1분 -# * * * * * /Users/dsyoon/workspace/bithumb/scripts/00_run_download_cron.sh >> /Users/dsyoon/workspace/bithumb/data/common/download_cron.log 2>&1 +# 캔들 증분 (DOWNLOAD_SYMBOLS × DOWNLOAD_INTERVALS) — 매 1분 +* * * * * /Users/dsyoon/workspace/bithumb/scripts/00_run_download_cron.sh >> /Users/dsyoon/workspace/bithumb/data/common/download_cron.log 2>&1 # vol_breakout 15m flip tick — 매 1분 # * * * * * /Users/dsyoon/workspace/bithumb/scripts/3_run_vol_breakout_cron.sh >> /Users/dsyoon/workspace/bithumb/data/spot/operations/vol_breakout_cron.log 2>&1 +# RSI DCA 정액 매수 tick — 매 1분 (모드 .env RSI_DCA_MODE, 기본 paper). 활성화 시 # 제거 +* * * * * /Users/dsyoon/workspace/bithumb/scripts/3_run_rsi_dca_cron.sh >> /Users/dsyoon/workspace/bithumb/data/spot/operations/rsi_dca_cron.log 2>&1 + # vol_live 모니터 JSON/HTML 백업 갱신 — 5분마다 -# */5 * * * * /Users/dsyoon/workspace/bithumb/scripts/3_run_vol_monitor_cron.sh >> /Users/dsyoon/workspace/bithumb/data/spot/operations/vol_monitor_cron.log 2>&1 +*/5 * * * * /Users/dsyoon/workspace/bithumb/scripts/3_run_vol_monitor_cron.sh >> /Users/dsyoon/workspace/bithumb/data/spot/operations/vol_monitor_cron.log 2>&1 # vol_live 모니터 HTTP 서버(8766) — 2분마다 미기동 시 nohup 기동 -# */2 * * * * /Users/dsyoon/workspace/bithumb/scripts/3_ensure_vol_monitor_serve.sh >> /Users/dsyoon/workspace/bithumb/data/spot/operations/vol_monitor_serve.log 2>&1 +*/2 * * * * /Users/dsyoon/workspace/bithumb/scripts/3_ensure_vol_monitor_serve.sh >> /Users/dsyoon/workspace/bithumb/data/spot/operations/vol_monitor_serve.log 2>&1 # (선택) fractal 운영 감시 — vol 전용이면 주석 유지 # */5 * * * * /Users/dsyoon/workspace/bithumb/scripts/3_run_watch_cron.sh >> /Users/dsyoon/workspace/bithumb/data/spot/operations/watch_cron.log 2>&1 diff --git a/scripts/install_crontab.sh b/scripts/install_crontab.sh index 4518b36..ba9a261 100755 --- a/scripts/install_crontab.sh +++ b/scripts/install_crontab.sh @@ -8,7 +8,7 @@ MARKER_END="# BITHUMB vol_breakout cron (end)" MARKER_BEGIN_ALT="# BITHUMB (begin)" MARKER_END_ALT="# BITHUMB (end)" EXAMPLE="${ROOT}/scripts/crontab.bithumb.example" -BITHUMB_CRON_RE='bithumb/scripts/(00_run_download_cron|3_run_vol_breakout_cron|3_run_vol_monitor_cron|3_ensure_vol_monitor_serve|3_run_watch_cron)\.sh' +BITHUMB_CRON_RE='bithumb/scripts/(00_run_download_cron|3_run_vol_breakout_cron|3_run_rsi_dca_cron|3_run_vol_monitor_cron|3_ensure_vol_monitor_serve|3_run_watch_cron)\.sh' usage() { cat < dict[str, Any]: + """단일 주문 상세 (GET /v1/order) — executed_volume, executed_funds, paid_fee, trades.""" + payload = self._request("GET", "/v1/order", params={"uuid": str(uuid)}) + if isinstance(payload, dict) and isinstance(payload.get("data"), dict): + return payload["data"] + return payload if isinstance(payload, dict) else {} + def get_orders( self, market: str, diff --git a/src/bithumb/config.py b/src/bithumb/config.py index 2ad2cb9..c8c6d4b 100644 --- a/src/bithumb/config.py +++ b/src/bithumb/config.py @@ -42,6 +42,11 @@ def resolve_coin_name(symbol: str) -> str: "TRX": "트론", "NEAR": "니어프로토콜", "WLD": "월드코인", + "XRP": "리플", + "SOL": "솔라나", + "ETH": "이더리움", + "ADA": "에이다", + "SUI": "수이", } return names.get(symbol.upper(), symbol.upper()) @@ -175,6 +180,25 @@ class Settings: vol_monitor_json: Path vol_monitor_html: Path vol_monitor_days: int + vol_monitor_symbols: list[str] + vol_monitor_intervals: list[int] + vol_monitor_max_bars: int + vol_monitor_avg_price: dict[str, float] + # RSI DCA 정액 매수 (15m RSI 30/35 상향 돌파, 매도 없음) + rsi_dca_mode: str + rsi_dca_symbols: list[str] + rsi_dca_interval_min: int + rsi_dca_period: int + rsi_dca_levels: list[tuple[float, float]] + rsi_dca_levels_by_symbol: dict[str, list[tuple[float, float]]] + rsi_dca_daily_max_krw: float + rsi_dca_lookback_days: int + rsi_dca_max_bars_per_tick: int + rsi_dca_max_signal_age_min: int + rsi_dca_state_json: Path + rsi_dca_report_json: Path + rsi_dca_tick_lock_path: Path | None + rsi_dca_kill_switch_path: Path | None telegram_bot_token: str telegram_chat_id: str ops_telegram_enabled: bool @@ -444,6 +468,40 @@ def load_settings(env_path: Path | None = None) -> Settings: ) ), vol_monitor_days=int(os.getenv("VOL_MONITOR_DAYS", "14")), + # 모니터 표시 종목 — 비우면 ops_symbols(매매 종목) + vol_monitor_symbols=_parse_symbol_list(os.getenv("VOL_MONITOR_SYMBOLS", "")), + # 모니터 분봉 탭 (분 단위) · 분봉당 최대 표시 봉수 + vol_monitor_intervals=_parse_int_list( + os.getenv("VOL_MONITOR_INTERVALS", "1,3,5,10,15,30,60,240,1440") + ), + vol_monitor_max_bars=int(os.getenv("VOL_MONITOR_MAX_BARS", "1500")), + # 수동 보유분 평균 매입가 (수익률 표시용). 예: "TRX:470,XRP:1900" + vol_monitor_avg_price=_parse_symbol_price_map(os.getenv("VOL_MONITOR_AVG_PRICE", "")), + # RSI DCA — 모드는 OPS_MODE와 분리 (기본 paper). live 전환은 RSI_DCA_MODE=live + rsi_dca_mode=os.getenv("RSI_DCA_MODE", "paper").strip().lower(), + rsi_dca_symbols=_parse_symbol_list(os.getenv("RSI_DCA_SYMBOLS", "")) + or [s for s in download_symbols if s != "BTC"], + rsi_dca_interval_min=int(os.getenv("RSI_DCA_INTERVAL_MIN", "15")), + rsi_dca_period=int(os.getenv("RSI_DCA_PERIOD", "14")), + rsi_dca_levels=_parse_rsi_levels(os.getenv("RSI_DCA_LEVELS", "30:20000,35:10000")), + # 종목별 오버라이드: "XRP=19:10000;TRX=32:10000" (없는 종목은 RSI_DCA_LEVELS) + rsi_dca_levels_by_symbol=_parse_rsi_levels_by_symbol(os.getenv("RSI_DCA_LEVELS_BY_SYMBOL", "")), + rsi_dca_daily_max_krw=float(os.getenv("RSI_DCA_DAILY_MAX_KRW", "60000")), + rsi_dca_lookback_days=int(os.getenv("RSI_DCA_LOOKBACK_DAYS", "20")), + rsi_dca_max_bars_per_tick=int(os.getenv("RSI_DCA_MAX_BARS_PER_TICK", "8")), + rsi_dca_max_signal_age_min=int(os.getenv("RSI_DCA_MAX_SIGNAL_AGE_MIN", "45")), + rsi_dca_state_json=_resolve_project_path( + os.getenv("RSI_DCA_STATE_JSON", "data/spot/operations/rsi_dca_state.json") + ), + rsi_dca_report_json=_resolve_project_path( + os.getenv("RSI_DCA_REPORT_JSON", "docs/spot/3_operations/rsi_dca_report.json") + ), + rsi_dca_tick_lock_path=_resolve_project_path_optional( + os.getenv("RSI_DCA_TICK_LOCK_PATH", "data/spot/operations/rsi.tick.lock") + ), + rsi_dca_kill_switch_path=_resolve_project_path_optional( + os.getenv("RSI_DCA_KILL_SWITCH_PATH", "data/spot/operations/rsi.kill") + ), telegram_bot_token=os.getenv("COIN_TELEGRAM_BOT_TOKEN", "").strip(), telegram_chat_id=os.getenv("COIN_TELEGRAM_CHAT_ID", "").strip(), ops_telegram_enabled=_parse_ops_telegram_enabled( @@ -454,6 +512,47 @@ def load_settings(env_path: Path | None = None) -> Settings: ) +def _parse_rsi_levels(raw: str) -> list[tuple[float, float]]: + """'30:20000,35:10000' → [(30.0, 20000.0), (35.0, 10000.0)] (level 오름차순).""" + out: list[tuple[float, float]] = [] + for part in (raw or "").split(","): + part = part.strip() + if not part: + continue + lv, _, krw = part.partition(":") + out.append((float(lv), float(krw))) + return sorted(out, key=lambda x: x[0]) + + +def _parse_symbol_price_map(raw: str) -> dict[str, float]: + """'TRX:470,XRP:1900' → {"TRX": 470.0, "XRP": 1900.0}.""" + out: dict[str, float] = {} + for part in (raw or "").split(","): + part = part.strip() + if not part or ":" not in part: + continue + sym, _, px = part.partition(":") + try: + out[sym.strip().upper()] = float(px) + except ValueError: + continue + return out + + +def _parse_rsi_levels_by_symbol(raw: str) -> dict[str, list[tuple[float, float]]]: + """'XRP=19:10000;TRX=32:10000,35:5000' → 종목별 레벨 dict.""" + out: dict[str, list[tuple[float, float]]] = {} + for part in (raw or "").split(";"): + part = part.strip() + if not part or "=" not in part: + continue + sym, _, lv = part.partition("=") + levels = _parse_rsi_levels(lv) + if levels: + out[sym.strip().upper()] = levels + return out + + def _parse_ops_telegram_enabled(raw: str, *, bot_token: str, chat_id: str) -> bool: """운영 텔레그램 알림 on/off. diff --git a/src/bithumb/data/candle_loader.py b/src/bithumb/data/candle_loader.py index ff3ec41..c5522b5 100644 --- a/src/bithumb/data/candle_loader.py +++ b/src/bithumb/data/candle_loader.py @@ -30,7 +30,12 @@ def load_candles( """ store = CandleStore(db_path) try: - df = store.read_dataframe(symbol, interval_min) + since = None + if lookback_days is not None and lookback_days > 0: + _, _, db_max = store.get_range(symbol, interval_min) + if db_max is not None: + since = db_max - timedelta(days=lookback_days) + df = store.read_dataframe(symbol, interval_min, since=since) finally: store.close() @@ -42,3 +47,17 @@ def load_candles( df = df[df["datetime"] >= cutoff].reset_index(drop=True) return df + + +def load_last_candles( + db_path: Path | str, + symbol: str, + interval_min: int, + last_n: int, +) -> pd.DataFrame: + """최신 N봉만 로드 (차트 API용, SQL LIMIT).""" + store = CandleStore(db_path) + try: + return store.read_dataframe(symbol, interval_min, last_n=int(last_n)) + finally: + store.close() diff --git a/src/bithumb/data/candle_store.py b/src/bithumb/data/candle_store.py index 30a329e..7eb35ba 100644 --- a/src/bithumb/data/candle_store.py +++ b/src/bithumb/data/candle_store.py @@ -78,18 +78,24 @@ class CandleStore: ``(row_count, min_dt, max_dt)``. 데이터 없으면 ``(0, None, None)``. """ table = self.table_name(symbol, interval_min) + code = symbol.upper() try: - row = self._conn.execute( - f"SELECT COUNT(*), MIN(ymdhms), MAX(ymdhms) FROM {table} WHERE CODE = ?", - (symbol.upper(),), + # (CODE, ymdhms) 인덱스를 타는 ORDER BY … LIMIT 1 — 대형 테이블(수백만 행)에서도 즉시 응답 + mn = self._conn.execute( + f"SELECT ymdhms FROM {table} WHERE CODE = ? ORDER BY ymdhms ASC LIMIT 1", (code,) + ).fetchone() + if mn is None or mn[0] is None: + return 0, None, None + mx = self._conn.execute( + f"SELECT ymdhms FROM {table} WHERE CODE = ? ORDER BY ymdhms DESC LIMIT 1", (code,) + ).fetchone() + cnt = self._conn.execute( + f"SELECT COUNT(*) FROM {table} WHERE CODE = ?", (code,) ).fetchone() except sqlite3.OperationalError: return 0, None, None - if row is None or row[0] == 0 or row[1] is None: - return 0, None, None - - return int(row[0]), parse_kst_datetime(str(row[1])), parse_kst_datetime(str(row[2])) + return int(cnt[0] if cnt else 0), parse_kst_datetime(str(mn[0])), parse_kst_datetime(str(mx[0])) def delete_incomplete_tail( self, @@ -116,28 +122,61 @@ class CandleStore: self._conn.commit() return cur.rowcount - def read_dataframe(self, symbol: str, interval_min: int) -> pd.DataFrame: + def get_max_datetime(self, symbol: str, interval_min: int) -> datetime | None: + """마지막 봉 시각만 (인덱스 ORDER BY … LIMIT 1, COUNT 없이 즉시).""" + table = self.table_name(symbol, interval_min) + try: + row = self._conn.execute( + f"SELECT ymdhms FROM {table} WHERE CODE = ? ORDER BY ymdhms DESC LIMIT 1", + (symbol.upper(),), + ).fetchone() + except sqlite3.OperationalError: + return None + return parse_kst_datetime(str(row[0])) if row and row[0] else None + + def read_dataframe( + self, + symbol: str, + interval_min: int, + *, + since: datetime | str | None = None, + last_n: int | None = None, + ) -> pd.DataFrame: """캔들을 pandas DataFrame으로 읽는다. Args: symbol: 코인 심볼. interval_min: 분 단위 인터벌. + since: 이 시각 이상만 (SQL WHERE, 대형 테이블 전체 스캔 방지). + last_n: 최신 N행만 (SQL ORDER BY DESC LIMIT). Returns: - 소문자 OHLCV 컬럼 DataFrame. 테이블 없으면 빈 DataFrame. + 소문자 OHLCV 컬럼 DataFrame(시간 오름차순). 테이블 없으면 빈 DataFrame. """ table = self.table_name(symbol, interval_min) - try: - raw = pd.read_sql_query( - f""" + params: list = [symbol.upper()] + where = "WHERE CODE = ?" + if since is not None: + since_s = since.strftime("%Y-%m-%d %H:%M:%S") if isinstance(since, datetime) else str(since) + where += " AND ymdhms >= ?" + params.append(since_s) + if last_n is not None and last_n > 0: + sql = f""" + SELECT ymdhms, Open, High, Low, Close, Volume FROM ( + SELECT ymdhms, Open, High, Low, Close, Volume + FROM {table} {where} + ORDER BY ymdhms DESC LIMIT ? + ) ORDER BY ymdhms ASC + """ + params.append(int(last_n)) + else: + sql = f""" SELECT ymdhms, Open, High, Low, Close, Volume - FROM {table} - WHERE CODE = ? + FROM {table} {where} ORDER BY ymdhms ASC - """, - self._conn, - params=(symbol.upper(),), - ) + """ + try: + raw = pd.read_sql_query(sql, self._conn, params=tuple(params)) except Exception: return pd.DataFrame( columns=["datetime", "open", "high", "low", "close", "volume"] diff --git a/src/bithumb/data/gap_backfill.py b/src/bithumb/data/gap_backfill.py new file mode 100644 index 0000000..a50bc68 --- /dev/null +++ b/src/bithumb/data/gap_backfill.py @@ -0,0 +1,90 @@ +"""캔들 공백 자동 백필 계획 — 절전·재부팅 후 증분 수집 범위(200봉)를 넘는 공백을 감지한다. + +증분 수집은 인터벌당 최신 200봉만 받으므로, DB 마지막 봉과 현재 시각의 차이가 +``200 × interval × 안전계수`` 를 넘으면 그 인터벌은 풀 다운(--full --days N)으로 채워야 한다. +""" + +from __future__ import annotations + +import math +from dataclasses import dataclass, field +from datetime import datetime +from typing import Iterable + + +@dataclass(frozen=True) +class BackfillPlan: + """백필 계획 — 심볼별 (인터벌 → 필요 일수).""" + + needs: dict[str, dict[int, int]] = field(default_factory=dict) # {symbol: {interval_min: days}} + + @property + def empty(self) -> bool: + return not self.needs + + def symbols(self) -> list[str]: + return sorted(self.needs) + + def intervals(self) -> list[int]: + out: set[int] = set() + for m in self.needs.values(): + out.update(m) + return sorted(out) + + def max_days(self) -> int: + return max((d for m in self.needs.values() for d in m.values()), default=0) + + def describe(self) -> str: + parts = [] + for sym in self.symbols(): + iv_txt = ", ".join(f"{iv}m→{days}d" for iv, days in sorted(self.needs[sym].items())) + parts.append(f"{sym}[{iv_txt}]") + return " ".join(parts) if parts else "no gap" + + +def coverage_minutes(interval_min: int, *, batch_size: int = 200, safety: float = 0.9) -> float: + """증분 1회(batch_size봉)가 덮는 시간(분) × 안전계수.""" + return batch_size * int(interval_min) * safety + + +def days_needed(lag_minutes: float, *, extra_days: int = 1) -> int: + """공백(분)을 채우기 위한 풀 다운 일수 (여유 extra_days 포함, 최소 1).""" + return max(1, math.ceil(lag_minutes / 1440.0) + int(extra_days)) + + +def plan_backfill( + db_max_by_symbol_interval: dict[str, dict[int, datetime | None]], + *, + now: datetime, + intervals: Iterable[int], + batch_size: int = 200, + safety: float = 0.9, + max_days: int = 30, +) -> BackfillPlan: + """DB 마지막 봉 시각으로 백필 필요 여부·일수를 계산한다. + + Args: + db_max_by_symbol_interval: {symbol: {interval_min: 마지막 봉 시각 or None}}. + now: 기준 시각. + intervals: 점검할 인터벌 목록. + batch_size: 증분 1회 봉 수(빗썸 200). + safety: 커버리지 안전계수(0.9 → 200봉의 90% 넘으면 백필). + max_days: 백필 상한 일수(그 이상 공백은 상한만큼만). + + Returns: + BackfillPlan. 데이터가 전혀 없는 (symbol, interval)은 최초 적재 대상이 아니므로 제외한다. + """ + needs: dict[str, dict[int, int]] = {} + for sym, by_iv in db_max_by_symbol_interval.items(): + for iv in intervals: + db_max = by_iv.get(iv) + if db_max is None: + continue # 미적재 테이블은 별도 초기 적재(--full) 대상 + lag_min = (now - db_max).total_seconds() / 60.0 + # 마지막 봉 이후 아직 마감되지 않은 1봉은 공백이 아님 + lag_min -= iv + if lag_min <= coverage_minutes(iv, batch_size=batch_size, safety=safety): + continue + days = min(days_needed(lag_min), max_days) + needs.setdefault(sym.upper(), {})[int(iv)] = days + return BackfillPlan(needs=needs) diff --git a/src/bithumb/operations/rsi_dca_control.py b/src/bithumb/operations/rsi_dca_control.py new file mode 100644 index 0000000..fa2c538 --- /dev/null +++ b/src/bithumb/operations/rsi_dca_control.py @@ -0,0 +1,142 @@ +"""RSI DCA 온/오프 제어 — 킬스위치 파일 기반 (모니터 UI·API용).""" + +from __future__ import annotations + +import json +from datetime import datetime +from pathlib import Path +from typing import Any + +TICK_ALIVE_SEC = 180 # 마지막 tick 이후 이 시간 내면 '가동 중' + + +def _load_json(path: Path | None) -> dict[str, Any]: + if not path: + return {} + p = Path(path) + if not p.exists(): + return {} + try: + return json.loads(p.read_text(encoding="utf-8")) + except Exception: # noqa: BLE001 + return {} + + +def rsi_enabled(settings: Any) -> bool: + """킬스위치 파일이 없으면 ON.""" + p = getattr(settings, "rsi_dca_kill_switch_path", None) + return not (p is not None and Path(p).exists()) + + +def set_rsi_enabled(settings: Any, enabled: bool) -> dict[str, Any]: + """ON → 킬스위치 삭제, OFF → 킬스위치 생성. 결과 status 반환.""" + p = getattr(settings, "rsi_dca_kill_switch_path", None) + if p is None: + return {"ok": False, "error": "RSI_DCA_KILL_SWITCH_PATH 미설정"} + path = Path(p) + path.parent.mkdir(parents=True, exist_ok=True) + if enabled: + if path.exists(): + path.unlink() + else: + path.write_text( + f"off by monitor {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}\n", + encoding="utf-8", + ) + st = rsi_status(settings) + st["changed"] = True + return st + + +def _env_file_mode(settings: Any) -> str | None: + """.env 파일의 RSI_DCA_MODE를 매번 직접 읽는다. + + 장기 실행 서버는 기동 시 os.environ에 올라간 값이 고정되어(load_dotenv override=False) + .env 변경이 반영되지 않으므로, 파일을 직접 파싱해 현재 값을 보여준다. + """ + env_path = getattr(settings, "env_path", None) + if not env_path: + env_path = Path(__file__).resolve().parents[3] / ".env" + try: + for line in Path(env_path).read_text(encoding="utf-8").splitlines(): + line = line.strip() + if line.startswith("RSI_DCA_MODE="): + val = line.split("=", 1)[1].split("#", 1)[0].strip().strip('"').strip("'").lower() + if val in ("paper", "live"): + return val + except Exception: # noqa: BLE001 + return None + return None + + +def html_build_stamp(settings: Any) -> str | None: + """서버가 제공하는 모니터 HTML의 PAGE_BUILD 값 (탭 자동 새로고침 판단용).""" + path = getattr(settings, "vol_monitor_html", None) + if not path: + return None + try: + import re + + text = Path(path).read_text(encoding="utf-8") + m = re.search(r'const PAGE_BUILD = "([^"]+)"', text) + return m.group(1) if m else None + except Exception: # noqa: BLE001 + return None + + +def rsi_status(settings: Any, *, now: datetime | None = None) -> dict[str, Any]: + """모니터 표시용 상태 요약.""" + now = now or datetime.now() + state = _load_json(getattr(settings, "rsi_dca_state_json", None)) + report = _load_json(getattr(settings, "rsi_dca_report_json", None)) + daily = state.get("daily") or {} + totals = state.get("totals") or {} + last_run = state.get("last_run_at") or report.get("last_run_at") + tick_age = None + if last_run: + try: + tick_age = (now - datetime.strptime(str(last_run)[:19], "%Y-%m-%d %H:%M:%S")).total_seconds() + except ValueError: + tick_age = None + # 실행 중 규칙(러너가 기록) 우선 — 서버 프로세스의 env 캐시 회피 + daily_max = float(state.get("daily_max_krw") or getattr(settings, "rsi_dca_daily_max_krw", 0) or 0) + today = now.strftime("%Y-%m-%d") + spent_today = float(daily.get("spent_krw") or 0.0) if daily.get("date") == today else 0.0 + symbols = {} + for sym, s in (state.get("symbols") or {}).items(): + symbols[sym] = { + "rsi": s.get("last_rsi"), + "cursor": s.get("last_confirm_time"), + "buys": s.get("buy_count", 0), + "spent_krw": s.get("spent_krw", 0.0), + "initialized": s.get("initialized", False), + } + trades = list(state.get("trades") or [])[-5:] + return { + "ok": True, + "enabled": rsi_enabled(settings), + "mode": _env_file_mode(settings) or str(getattr(settings, "rsi_dca_mode", "paper") or "paper"), + "state_mode": state.get("mode"), + "last_run_at": last_run, + "tick_age_sec": None if tick_age is None else int(tick_age), + "tick_alive": tick_age is not None and 0 <= tick_age <= TICK_ALIVE_SEC, + "daily": { + "date": today, + "spent_krw": spent_today, + "max_krw": daily_max, + "remaining_krw": max(daily_max - spent_today, 0.0), + "count": int(daily.get("count") or 0) if daily.get("date") == today else 0, + }, + "totals": { + "spent_krw": float(totals.get("spent_krw") or 0.0), + "count": int(totals.get("count") or 0), + }, + # 실제 실행 중인 규칙은 러너가 state에 기록한 값을 우선 (서버 프로세스의 env 캐시 회피) + "levels": state.get("levels") or [list(x) for x in (getattr(settings, "rsi_dca_levels", None) or [])], + "levels_by_symbol": state.get("levels_by_symbol") or {}, + "interval_min": state.get("interval_min") or getattr(settings, "rsi_dca_interval_min", None), + "symbols": symbols, + "recent_trades": trades, + "kill_switch_path": str(getattr(settings, "rsi_dca_kill_switch_path", "") or ""), + "html_build": html_build_stamp(settings), + } diff --git a/src/bithumb/operations/rsi_dca_engine.py b/src/bithumb/operations/rsi_dca_engine.py new file mode 100644 index 0000000..0031882 --- /dev/null +++ b/src/bithumb/operations/rsi_dca_engine.py @@ -0,0 +1,518 @@ +"""RSI DCA 현물 정액 매수 — 15분봉 RSI가 기준선을 상향 돌파하는 봉 마감에 고정 금액 매수. + +규칙(2026-09-07 확정): +- RSI(14) 종가 확정 기준. prev <= level < now 이면 해당 level 매수. +- 30 상향 돌파 20,000원, 35 상향 돌파 10,000원. 두 조건은 독립적으로 모두 발생 가능. +- 일(KST) 총 매수 상한 60,000원. 상한을 넘기는 주문은 스킵. +- 자동 매도 없음 (보유분은 수동 관리). +""" + +from __future__ import annotations + +import json +import logging +import math +from dataclasses import dataclass, field +from datetime import datetime, timedelta +from pathlib import Path +from typing import Any, Callable + +import numpy as np +import pandas as pd + +logger = logging.getLogger(__name__) + +STRATEGY = "rsi_dca_15m_spot_buy" + +# 매수 실행 콜백: (symbol, krw, ref_price) -> {"ok", "order_krw", "order_coin", "price", "error", "api_response"} +BuyFn = Callable[[str, float, float], dict[str, Any]] + + +# --------------------------------------------------------------------------- +# 지표 +# --------------------------------------------------------------------------- +def _rsi_value(avg_gain: float, avg_loss: float) -> float: + if avg_loss == 0: + return 100.0 + if avg_gain == 0: + return 0.0 + return 100.0 - 100.0 / (1.0 + avg_gain / avg_loss) + + +def wilder_rsi(close: pd.Series | list[float], period: int = 14) -> pd.Series: + """Wilder RSI — 첫 평균은 단순평균, 이후 (prev*(n-1)+cur)/n. + + 모니터 차트(JS computeRSI)와 동일한 수식·시드를 사용한다. + index period 이전 값은 NaN. + """ + series = pd.Series(close, dtype="float64") + values = series.to_numpy() + n = len(values) + out = np.full(n, np.nan) + if n <= period or period <= 0: + return pd.Series(out, index=series.index) + deltas = np.diff(values) + gains = np.where(deltas > 0, deltas, 0.0) + losses = np.where(deltas < 0, -deltas, 0.0) + avg_g = float(gains[:period].mean()) + avg_l = float(losses[:period].mean()) + out[period] = _rsi_value(avg_g, avg_l) + for i in range(period + 1, n): + avg_g = (avg_g * (period - 1) + gains[i - 1]) / period + avg_l = (avg_l * (period - 1) + losses[i - 1]) / period + out[i] = _rsi_value(avg_g, avg_l) + return pd.Series(out, index=series.index) + + +def closed_candles(df: pd.DataFrame, interval_min: int, now: datetime) -> pd.DataFrame: + """마감된 봉만 남긴다 (봉 시작 + interval <= now).""" + if df.empty: + return df + out = df.copy() + out["datetime"] = pd.to_datetime(out["datetime"]) + cutoff = pd.Timestamp(now) - pd.Timedelta(minutes=interval_min) + out = out[out["datetime"] <= cutoff] + return out.reset_index(drop=True) + + +def cross_up_levels( + prev_rsi: float, + cur_rsi: float, + levels: list[tuple[float, float]], +) -> list[tuple[float, float]]: + """prev <= level < cur 를 만족하는 (level, krw) 목록 (level 오름차순).""" + if prev_rsi is None or cur_rsi is None: + return [] + if math.isnan(prev_rsi) or math.isnan(cur_rsi): + return [] + hits = [(lv, krw) for lv, krw in levels if prev_rsi <= lv < cur_rsi] + return sorted(hits, key=lambda x: x[0]) + + +def parse_levels_by_symbol(raw: str) -> dict[str, list[tuple[float, float]]]: + """'XRP=19:10000;TRX=32:10000,35:5000' → {"XRP": [(19,10000)], "TRX": [(32,10000),(35,5000)]}.""" + out: dict[str, list[tuple[float, float]]] = {} + for part in (raw or "").split(";"): + part = part.strip() + if not part or "=" not in part: + continue + sym, _, lv = part.partition("=") + levels = parse_levels(lv) + if levels: + out[sym.strip().upper()] = levels + return out + + +def parse_levels(raw: str) -> list[tuple[float, float]]: + """'30:20000,35:10000' → [(30.0, 20000.0), (35.0, 10000.0)].""" + out: list[tuple[float, float]] = [] + for part in (raw or "").split(","): + part = part.strip() + if not part: + continue + lv, _, krw = part.partition(":") + out.append((float(lv), float(krw))) + return sorted(out, key=lambda x: x[0]) + + +# --------------------------------------------------------------------------- +# 설정·상태 +# --------------------------------------------------------------------------- +@dataclass(frozen=True) +class RsiDcaConfig: + """전략 파라미터.""" + + symbols: list[str] + interval_min: int = 15 + period: int = 14 + levels: list[tuple[float, float]] = field(default_factory=lambda: [(30.0, 20000.0), (35.0, 10000.0)]) + daily_max_krw: float = 60000.0 + lookback_days: int = 20 + max_bars_per_tick: int = 8 + max_signal_age_min: int = 45 + min_order_krw: float = 5000.0 + fee_rate: float = 0.0005 + slippage_rate: float = 0.0005 + fee_lock_rate: float = 0.0025 + # 종목별 레벨 오버라이드 {"TRX": [(32, 10000)], ...} — 없으면 levels 사용 + levels_by_symbol: dict[str, list[tuple[float, float]]] = field(default_factory=dict) + + def levels_for(self, symbol: str) -> list[tuple[float, float]]: + """종목별 레벨 (없으면 공통 levels).""" + ov = (self.levels_by_symbol or {}).get(str(symbol).upper()) + return list(ov) if ov else list(self.levels) + + +def _default_sym_state() -> dict[str, Any]: + return { + "initialized": False, + "last_confirm_time": None, + "last_rsi": None, + "last_price": 0.0, + "buy_count": 0, + "spent_krw": 0.0, + "coin_qty_est": 0.0, + } + + +def empty_state(mode: str) -> dict[str, Any]: + return { + "strategy": STRATEGY, + "mode": mode, + "symbols": {}, + "daily": {"date": None, "spent_krw": 0.0, "count": 0}, + "totals": {"spent_krw": 0.0, "count": 0}, + "trades": [], + "events": [], + "last_run_at": None, + } + + +def load_state(path: Path, mode: str) -> dict[str, Any]: + if not path.exists(): + return empty_state(mode) + with path.open(encoding="utf-8") as f: + state = json.load(f) + base = empty_state(mode) + for k, v in base.items(): + state.setdefault(k, v) + return state + + +def save_state(path: Path, state: dict[str, Any]) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + tmp = path.with_suffix(path.suffix + ".tmp") + with tmp.open("w", encoding="utf-8") as f: + json.dump(state, f, ensure_ascii=False, indent=2) + tmp.replace(path) + + +def paper_buy_fn(cfg: RsiDcaConfig) -> BuyFn: + """paper 체결 — 슬리피지·수수료 반영 추정 수량.""" + + def _buy(symbol: str, krw: float, ref_price: float) -> dict[str, Any]: + px = float(ref_price) * (1.0 + cfg.slippage_rate) + if px <= 0: + return {"ok": False, "error": "price<=0"} + order_krw = float(math.floor(krw)) + coin = order_krw * (1.0 - cfg.fee_rate) / px + return {"ok": True, "order_krw": order_krw, "order_coin": coin, "price": px, "api_response": None} + + return _buy + + +def apply_fill_to_trade(rec: dict[str, Any], order: dict[str, Any]) -> bool: + """거래소 주문 상세(executed_funds/volume/paid_fee)로 기록의 체결가·수량·수수료를 실제값으로 보정. + + Returns: + 보정 완료(주문 done·체결량>0) 여부. + """ + try: + vol = float(order.get("executed_volume") or 0.0) + funds = float(order.get("executed_funds") or 0.0) + fee = float(order.get("paid_fee") or 0.0) + except (TypeError, ValueError): + return False + state = str(order.get("state") or "") + if vol <= 0 or funds <= 0: + return False + rec.setdefault("price_ref", rec.get("price")) + rec["order_coin"] = vol + rec["order_krw"] = funds + rec["price"] = funds / vol # 실제 평균 체결가 + rec["fee_krw"] = fee + rec["fill_reconciled"] = state == "done" + rec["fill_state"] = state + return state == "done" + + +# --------------------------------------------------------------------------- +# 엔진 +# --------------------------------------------------------------------------- +@dataclass +class SymbolTickResult: + symbol: str + note: str + processed_bars: int = 0 + fills: int = 0 + last_rsi: float | None = None + last_price: float = 0.0 + trade_records: list[dict[str, Any]] = field(default_factory=list) + + +class RsiDcaEngine: + """상태(state dict)를 갱신하며 종목별 tick을 처리한다.""" + + def __init__( + self, + cfg: RsiDcaConfig, + state: dict[str, Any], + *, + mode: str, + buy_fn: BuyFn, + available_cash_fn: Callable[[], float | None] | None = None, + ) -> None: + self.cfg = cfg + self.state = state + self.mode = mode + self._buy = buy_fn + self._avail_cash = available_cash_fn + + # -- 상태 헬퍼 ----------------------------------------------------------- + def sym_state(self, symbol: str) -> dict[str, Any]: + root = self.state.setdefault("symbols", {}) + st = root.setdefault(symbol.upper(), _default_sym_state()) + for k, v in _default_sym_state().items(): + st.setdefault(k, v) + return st + + def _daily(self, now: datetime) -> dict[str, Any]: + d = self.state.setdefault("daily", {"date": None, "spent_krw": 0.0, "count": 0}) + today = now.strftime("%Y-%m-%d") + if d.get("date") != today: + d["date"] = today + d["spent_krw"] = 0.0 + d["count"] = 0 + return d + + def daily_remaining_krw(self, now: datetime) -> float: + d = self._daily(now) + return max(self.cfg.daily_max_krw - float(d.get("spent_krw") or 0.0), 0.0) + + def _push_event(self, rec: dict[str, Any], *, keep: int = 300) -> None: + events = self.state.setdefault("events", []) + events.append(rec) + if len(events) > keep: + del events[: len(events) - keep] + + def _push_trade(self, rec: dict[str, Any], *, keep: int = 2000) -> None: + trades = self.state.setdefault("trades", []) + trades.append(rec) + if len(trades) > keep: + del trades[: len(trades) - keep] + + # -- 핵심 ------------------------------------------------------------------ + def process_symbol( + self, + symbol: str, + df_closed: pd.DataFrame, + *, + now: datetime, + block_entry: bool = False, + ) -> SymbolTickResult: + """마감 봉 DataFrame(datetime, close …)으로 신규 봉을 판정·매수한다.""" + sym = symbol.upper() + st = self.sym_state(sym) + cfg = self.cfg + + if df_closed.empty or len(df_closed) <= cfg.period + 1: + return SymbolTickResult(symbol=sym, note="no_data") + + df = df_closed.copy() + df["datetime"] = pd.to_datetime(df["datetime"]) + rsi = wilder_rsi(df["close"].astype(float), cfg.period) + last_idx = len(df) - 1 + last_time = df["datetime"].iloc[last_idx] + last_price = float(df["close"].iloc[last_idx]) + last_rsi = float(rsi.iloc[last_idx]) if not math.isnan(rsi.iloc[last_idx]) else None + + if not st.get("initialized"): + st["initialized"] = True + st["last_confirm_time"] = str(last_time)[:19] + st["last_rsi"] = last_rsi + st["last_price"] = last_price + return SymbolTickResult( + symbol=sym, note=f"initialized {str(last_time)[:19]}", + last_rsi=last_rsi, last_price=last_price, + ) + + cursor = pd.Timestamp(st["last_confirm_time"]) if st.get("last_confirm_time") else None + new_idx = [i for i in range(len(df)) if cursor is None or df["datetime"].iloc[i] > cursor] + if not new_idx: + st["last_rsi"] = last_rsi + st["last_price"] = last_price + return SymbolTickResult( + symbol=sym, note=f"no_new_bar {str(last_time)[:19]}", + last_rsi=last_rsi, last_price=last_price, + ) + if len(new_idx) > cfg.max_bars_per_tick: + skipped = len(new_idx) - cfg.max_bars_per_tick + self._push_event({ + "ts": now.strftime("%Y-%m-%d %H:%M:%S"), "symbol": sym, + "type": "catchup_truncated", "skipped_bars": skipped, + }) + new_idx = new_idx[-cfg.max_bars_per_tick:] + + result = SymbolTickResult(symbol=sym, note="", last_rsi=last_rsi, last_price=last_price) + notes: list[str] = [] + for i in new_idx: + if i == 0: + st["last_confirm_time"] = str(df["datetime"].iloc[i])[:19] + continue + bar_open = df["datetime"].iloc[i] + bar_close = bar_open + pd.Timedelta(minutes=cfg.interval_min) + prev_rsi = float(rsi.iloc[i - 1]) if not math.isnan(rsi.iloc[i - 1]) else None + cur_rsi = float(rsi.iloc[i]) if not math.isnan(rsi.iloc[i]) else None + price = float(df["close"].iloc[i]) + result.processed_bars += 1 + + hits = cross_up_levels(prev_rsi, cur_rsi, cfg.levels_for(sym)) if prev_rsi is not None and cur_rsi is not None else [] + for level, krw in hits: + bar_key = str(bar_open)[:19] + base = { + "ts": now.strftime("%Y-%m-%d %H:%M:%S"), "symbol": sym, "bar_time": bar_key, + "level": level, "krw": krw, "rsi_prev": round(prev_rsi, 2), "rsi": round(cur_rsi, 2), + } + age_min = (pd.Timestamp(now) - bar_close).total_seconds() / 60.0 + if age_min > cfg.max_signal_age_min: + self._push_event({**base, "type": "expired", "age_min": round(age_min, 1)}) + notes.append(f"expired L{level:g} {bar_key}") + continue + if block_entry: + self._push_event({**base, "type": "kill_switch"}) + notes.append(f"kill_switch L{level:g} {bar_key}") + continue + remaining = self.daily_remaining_krw(now) + if krw > remaining + 1e-9: + self._push_event({**base, "type": "daily_cap", "remaining_krw": remaining}) + notes.append(f"daily_cap L{level:g} {bar_key}") + continue + if krw < cfg.min_order_krw: + self._push_event({**base, "type": "below_min_order"}) + notes.append(f"below_min L{level:g}") + continue + if self._avail_cash is not None: + avail = self._avail_cash() + need = krw * (1.0 + cfg.fee_lock_rate) + if avail is not None and avail < need: + self._push_event({**base, "type": "insufficient_cash", "available_krw": avail}) + notes.append(f"no_cash L{level:g} {bar_key}") + continue + + try: + fill = self._buy(sym, krw, price) + except Exception as exc: # noqa: BLE001 + logger.exception("rsi_dca buy failed %s", sym) + fill = {"ok": False, "error": str(exc)} + if not fill.get("ok"): + self._push_event({**base, "type": "buy_failed", "error": str(fill.get("error"))}) + notes.append(f"buy_fail L{level:g} {bar_key}") + continue + + order_krw = float(fill.get("order_krw") or krw) + order_coin = float(fill.get("order_coin") or 0.0) + fill_px = float(fill.get("price") or price) + rec = { + "symbol": sym, "side": "buy", "ts": str(bar_close)[:19], "bar_time": bar_key, + "price": fill_px, "price_ref": price, "order_krw": order_krw, "order_coin": order_coin, + "fill_reconciled": bool(fill.get("fill_reconciled")) or self.mode != "live", + "fee_krw": fill.get("fee_krw"), + "level": level, "rsi_prev": round(prev_rsi, 2), "rsi": round(cur_rsi, 2), + "reason": f"rsi_cross_up_{level:g}", "mode": self.mode, + "executed_at": now.strftime("%Y-%m-%d %H:%M:%S"), + "api_response": fill.get("api_response"), + } + self._push_trade(rec) + result.trade_records.append(rec) + result.fills += 1 + d = self._daily(now) + d["spent_krw"] = float(d.get("spent_krw") or 0.0) + order_krw + d["count"] = int(d.get("count") or 0) + 1 + t = self.state.setdefault("totals", {"spent_krw": 0.0, "count": 0}) + t["spent_krw"] = float(t.get("spent_krw") or 0.0) + order_krw + t["count"] = int(t.get("count") or 0) + 1 + st["buy_count"] = int(st.get("buy_count") or 0) + 1 + st["spent_krw"] = float(st.get("spent_krw") or 0.0) + order_krw + st["coin_qty_est"] = float(st.get("coin_qty_est") or 0.0) + order_coin + notes.append(f"buy L{level:g} {order_krw:,.0f}원 {bar_key}") + + st["last_confirm_time"] = str(bar_open)[:19] + st["last_rsi"] = cur_rsi + st["last_price"] = price + + result.note = "; ".join(notes) if notes else f"no_signal {str(last_time)[:19]} rsi={last_rsi:.1f}" if last_rsi is not None else "no_signal" + return result + + +# --------------------------------------------------------------------------- +# 백테스트 (종목 병합 · 일 상한 공유) +# --------------------------------------------------------------------------- +def backtest_rsi_dca( + candles_by_symbol: dict[str, pd.DataFrame], + cfg: RsiDcaConfig, + *, + days: int | None = None, +) -> dict[str, Any]: + """15m 종가 기준 RSI 교차 정액 매수 재생. 일 상한은 전 종목 공유(시간순).""" + rows: list[dict[str, Any]] = [] + last_price: dict[str, float] = {} + for sym, df in candles_by_symbol.items(): + if df is None or df.empty: + continue + d = df.copy() + d["datetime"] = pd.to_datetime(d["datetime"]) + d = d.sort_values("datetime").reset_index(drop=True) + rsi = wilder_rsi(d["close"].astype(float), cfg.period) + start = d["datetime"].max() - pd.Timedelta(days=days) if days else None + last_price[sym.upper()] = float(d["close"].iloc[-1]) + for i in range(1, len(d)): + if start is not None and d["datetime"].iloc[i] < start: + continue + p, c = rsi.iloc[i - 1], rsi.iloc[i] + if math.isnan(p) or math.isnan(c): + continue + for level, krw in cross_up_levels(float(p), float(c), cfg.levels_for(sym)): + rows.append({ + "symbol": sym.upper(), "bar_time": d["datetime"].iloc[i], + "level": level, "krw": krw, "rsi_prev": float(p), "rsi": float(c), + "price": float(d["close"].iloc[i]), + }) + rows.sort(key=lambda r: (r["bar_time"], r["symbol"], r["level"])) + + daily_date = None + daily_spent = 0.0 + trades: list[dict[str, Any]] = [] + skipped_cap = 0 + per_sym: dict[str, dict[str, float]] = {} + for r in rows: + day = r["bar_time"].strftime("%Y-%m-%d") + if day != daily_date: + daily_date, daily_spent = day, 0.0 + if r["krw"] > cfg.daily_max_krw - daily_spent + 1e-9: + skipped_cap += 1 + continue + px = r["price"] * (1.0 + cfg.slippage_rate) + coin = r["krw"] * (1.0 - cfg.fee_rate) / px + daily_spent += r["krw"] + trades.append({**r, "bar_time": str(r["bar_time"])[:19], "fill_price": px, "order_coin": coin}) + ps = per_sym.setdefault(r["symbol"], {"signals": 0, "buys": 0, "spent_krw": 0.0, "coin": 0.0}) + ps["buys"] += 1 + ps["spent_krw"] += r["krw"] + ps["coin"] += coin + for r in rows: + per_sym.setdefault(r["symbol"], {"signals": 0, "buys": 0, "spent_krw": 0.0, "coin": 0.0})["signals"] += 1 + + total_spent = sum(v["spent_krw"] for v in per_sym.values()) + total_value = 0.0 + for sym, v in per_sym.items(): + v["value_krw"] = v["coin"] * last_price.get(sym, 0.0) + v["pnl_pct"] = (v["value_krw"] / v["spent_krw"] - 1.0) * 100.0 if v["spent_krw"] > 0 else 0.0 + total_value += v["value_krw"] + span_days = None + if rows: + span_days = max((rows[-1]["bar_time"] - rows[0]["bar_time"]).days, 1) + return { + "strategy": STRATEGY, + "config": { + "levels": cfg.levels, "levels_by_symbol": cfg.levels_by_symbol, + "daily_max_krw": cfg.daily_max_krw, "period": cfg.period, + "interval_min": cfg.interval_min, "fee_rate": cfg.fee_rate, "slippage_rate": cfg.slippage_rate, + }, + "days": days, "span_days": span_days, + "signals": len(rows), "buys": len(trades), "skipped_daily_cap": skipped_cap, + "total_spent_krw": total_spent, "total_value_krw": total_value, + "total_pnl_pct": (total_value / total_spent - 1.0) * 100.0 if total_spent > 0 else 0.0, + "avg_daily_spent_krw": total_spent / span_days if span_days else 0.0, + "per_symbol": per_sym, + "trades": trades, + } diff --git a/src/bithumb/operations/rsi_dca_runner.py b/src/bithumb/operations/rsi_dca_runner.py new file mode 100644 index 0000000..d1e74e5 --- /dev/null +++ b/src/bithumb/operations/rsi_dca_runner.py @@ -0,0 +1,240 @@ +"""RSI DCA 러너 — 설정 로드, 거래소 클라이언트, lock, 텔레그램, 상태 저장.""" + +from __future__ import annotations + +import json +import logging +import math +from datetime import datetime +from pathlib import Path +from typing import Any + +from bithumb.api.bithumb_private import BithumbPrivateClient +from bithumb.config import Settings +from bithumb.data.candle_loader import load_candles +from bithumb.notifications.telegram import create_telegram_notifier +from bithumb.operations.ops_lock import ops_tick_lock +from bithumb.operations.rsi_dca_engine import ( + STRATEGY, + BuyFn, + RsiDcaConfig, + RsiDcaEngine, + apply_fill_to_trade, + closed_candles, + load_state, + paper_buy_fn, + save_state, +) + +logger = logging.getLogger(__name__) + + +def config_from_settings(settings: Settings) -> RsiDcaConfig: + return RsiDcaConfig( + symbols=list(settings.rsi_dca_symbols), + interval_min=settings.rsi_dca_interval_min, + period=settings.rsi_dca_period, + levels=list(settings.rsi_dca_levels), + daily_max_krw=settings.rsi_dca_daily_max_krw, + lookback_days=settings.rsi_dca_lookback_days, + max_bars_per_tick=settings.rsi_dca_max_bars_per_tick, + max_signal_age_min=settings.rsi_dca_max_signal_age_min, + min_order_krw=settings.ops_min_order_krw, + fee_rate=settings.gt_trading_fee_rate, + slippage_rate=settings.ops_slippage_rate, + fee_lock_rate=settings.ops_exchange_fee_lock_rate, + levels_by_symbol=dict(getattr(settings, "rsi_dca_levels_by_symbol", None) or {}), + ) + + +def live_buy_fn(client: BithumbPrivateClient, cfg: RsiDcaConfig) -> BuyFn: + """빗썸 시장가 매수(원화 금액). 수량은 참조가 기준 추정치.""" + + def _buy(symbol: str, krw: float, ref_price: float) -> dict[str, Any]: + market = f"KRW-{symbol.upper()}" + order_krw = float(math.floor(krw)) + resp = client.market_buy_krw(market, order_krw) + px = float(ref_price) if ref_price > 0 else 0.0 + coin_est = order_krw * (1.0 - cfg.fee_rate) / px if px > 0 else 0.0 + out = {"ok": True, "order_krw": order_krw, "order_coin": coin_est, "price": px, "api_response": resp} + # 시장가는 즉시 체결되므로 상세 조회로 실제 체결가·수량 반영 (실패 시 추정값 유지, 다음 tick에서 보정) + uuid = (resp or {}).get("uuid") if isinstance(resp, dict) else None + if uuid: + try: + order = client.get_order(str(uuid)) + tmp: dict[str, Any] = {} + if apply_fill_to_trade(tmp, order): + out.update({"order_krw": tmp["order_krw"], "order_coin": tmp["order_coin"], + "price": tmp["price"], "fee_krw": tmp["fee_krw"], "fill_reconciled": True}) + except Exception as exc: # noqa: BLE001 + logger.warning("체결 상세 조회 실패 %s: %s", uuid, exc) + return out + + return _buy + + +class RsiDcaRunner: + """7종 RSI 정액 매수 tick.""" + + def __init__(self, settings: Settings, *, mode: str | None = None) -> None: + self.settings = settings + self.mode = (mode or settings.rsi_dca_mode or "paper").lower() + self.cfg = config_from_settings(settings) + self.state = load_state(settings.rsi_dca_state_json, self.mode) + self.state["strategy"] = STRATEGY + self.state["mode"] = self.mode + # 인터벌이 바뀌면 이전 인터벌 커서로 소급 판정하지 않도록 종목 커서 재초기화 + prev_iv = self.state.get("interval_min") + if prev_iv is not None and int(prev_iv) != int(self.cfg.interval_min): + for sym_st in (self.state.get("symbols") or {}).values(): + sym_st["initialized"] = False + sym_st["last_confirm_time"] = None + logger.warning("rsi_dca interval %s→%s: 종목 커서 재초기화", prev_iv, self.cfg.interval_min) + self.state.setdefault("events", []).append({ + "ts": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "type": "interval_changed", + "from": prev_iv, "to": self.cfg.interval_min, + }) + self.state["interval_min"] = int(self.cfg.interval_min) + self.state["daily_max_krw"] = float(self.cfg.daily_max_krw) + self.state["levels"] = [list(x) for x in self.cfg.levels] + self.state["levels_by_symbol"] = {k: [list(x) for x in v] for k, v in self.cfg.levels_by_symbol.items()} + self._client: BithumbPrivateClient | None = None + if self.mode == "live": + if not settings.bithumb_access_key or not settings.bithumb_secret_key: + raise RuntimeError("live: BITHUMB_ACCESS_KEY / BITHUMB_SECRET_KEY 필요") + self._client = BithumbPrivateClient( + access_key=settings.bithumb_access_key, + secret_key=settings.bithumb_secret_key, + base_url=settings.api_url, + sleep_sec=settings.request_sleep_sec, + retries=settings.request_retries, + ) + buy_fn = live_buy_fn(self._client, self.cfg) + avail_fn = self._available_krw + else: + buy_fn = paper_buy_fn(self.cfg) + avail_fn = None + self.engine = RsiDcaEngine( + self.cfg, self.state, mode=self.mode, buy_fn=buy_fn, available_cash_fn=avail_fn, + ) + self.telegram = create_telegram_notifier( + settings.telegram_bot_token, + settings.telegram_chat_id, + enabled=settings.ops_telegram_enabled, + ) + + # -- 헬퍼 ---------------------------------------------------------------- + def _available_krw(self) -> float | None: + if self._client is None: + return None + try: + avail, _ = self._client.get_balance("KRW") + return float(avail) + except Exception as exc: # noqa: BLE001 + logger.warning("KRW 잔고 조회 실패: %s", exc) + return None + + def _reconcile_fills(self, *, max_orders: int = 10) -> int: + """live 매수 기록 중 실제 체결가 미반영 건을 거래소 주문 상세로 보정.""" + if self._client is None: + return 0 + fixed = 0 + for rec in reversed(self.state.get("trades") or []): + if rec.get("fill_reconciled") or rec.get("mode") != "live": + continue + resp = rec.get("api_response") + uuid = resp.get("uuid") if isinstance(resp, dict) else None + if not uuid: + continue + try: + order = self._client.get_order(str(uuid)) + except Exception as exc: # noqa: BLE001 + logger.warning("체결 보정 조회 실패 %s: %s", uuid, exc) + continue + if apply_fill_to_trade(rec, order): + fixed += 1 + logger.info("체결 보정 %s %s: price=%.4f coin=%.6f fee=%s", rec.get("symbol"), uuid, rec["price"], rec["order_coin"], rec.get("fee_krw")) + max_orders -= 1 + if max_orders <= 0: + break + return fixed + + def _kill_switch_active(self) -> bool: + p = self.settings.rsi_dca_kill_switch_path + return p is not None and Path(p).exists() + + def _notify_trade(self, rec: dict[str, Any], now: datetime) -> None: + if not self.telegram.is_active: + return + d = self.state.get("daily") or {} + mode_txt = "실거래" if self.mode == "live" else "페이퍼" + text = ( + f"[{mode_txt}] RSI 정액 매수\n" + f"{rec['symbol']}KRW · {self.cfg.interval_min}분 RSI {rec['rsi_prev']}→{rec['rsi']} " + f"({rec['level']:g} 상향 돌파)\n" + f"금액 {rec['order_krw']:,.0f}원 · 가격 {rec['price']:,.2f} · 수량 {rec['order_coin']:.4f}\n" + f"봉 {rec['bar_time']} · 오늘 누적 {float(d.get('spent_krw') or 0):,.0f}/{self.cfg.daily_max_krw:,.0f}원" + ) + try: + self.telegram.send_message(text) + except Exception: # noqa: BLE001 + logger.exception("텔레그램 알림 실패") + + # -- tick ----------------------------------------------------------------- + def tick(self, *, skip_lock: bool = False) -> dict[str, Any]: + lock_path = self.settings.rsi_dca_tick_lock_path + if lock_path and not skip_lock: + with ops_tick_lock(Path(lock_path), blocking=False) as acquired: + if not acquired: + return {"ok": False, "note": "lock_busy"} + return self._tick_impl() + return self._tick_impl() + + def _tick_impl(self) -> dict[str, Any]: + now = datetime.now() + now_s = now.strftime("%Y-%m-%d %H:%M:%S") + block = self._kill_switch_active() + results: list[dict[str, Any]] = [] + fills = 0 + try: + self._reconcile_fills() + except Exception: # noqa: BLE001 + logger.exception("fill reconcile failed") + for sym in self.cfg.symbols: + try: + df = load_candles( + self.settings.db_path, sym, self.cfg.interval_min, + lookback_days=self.cfg.lookback_days, + ) + df = closed_candles(df, self.cfg.interval_min, now) + res = self.engine.process_symbol(sym, df, now=now, block_entry=block) + except Exception as exc: # noqa: BLE001 + logger.exception("rsi_dca tick failed %s", sym) + results.append({"symbol": sym, "error": str(exc)}) + continue + for rec in res.trade_records: + self._notify_trade(rec, now) + fills += res.fills + results.append({ + "symbol": sym, "note": res.note, "fills": res.fills, + "bars": res.processed_bars, "rsi": res.last_rsi, "price": res.last_price, + }) + + self.state["last_run_at"] = now_s + save_state(self.settings.rsi_dca_state_json, self.state) + + d = self.state.get("daily") or {} + report = { + "ok": True, "strategy": STRATEGY, "mode": self.mode, "symbols": self.cfg.symbols, + "kill_switch": block, "fills": fills, "results": results, "last_run_at": now_s, + "daily": {**d, "max_krw": self.cfg.daily_max_krw, + "remaining_krw": self.engine.daily_remaining_krw(now)}, + "totals": self.state.get("totals"), + } + try: + p = Path(self.settings.rsi_dca_report_json) + p.parent.mkdir(parents=True, exist_ok=True) + p.write_text(json.dumps(report, ensure_ascii=False, indent=2, default=str), encoding="utf-8") + except Exception: # noqa: BLE001 + logger.exception("rsi_dca report write failed") + return report diff --git a/src/bithumb/operations/vol_live_monitor.py b/src/bithumb/operations/vol_live_monitor.py index 379fd0e..66a3f5c 100644 --- a/src/bithumb/operations/vol_live_monitor.py +++ b/src/bithumb/operations/vol_live_monitor.py @@ -10,7 +10,7 @@ from typing import Any import pandas as pd from bithumb.config import Settings, resolve_coin_name -from bithumb.data.candle_loader import load_candles +from bithumb.data.candle_loader import load_candles, load_last_candles from bithumb.operations.multi_portfolio import in_long_position from bithumb.operations.vol_monitor_chart import write_vol_monitor_html from bithumb.simulation.vol_breakout import drop_incomplete_base_bar @@ -250,6 +250,45 @@ def _summary_html(summary: dict[str, Any]) -> str: return " · ".join(lines) +def fetch_ticker_prices(settings: Settings, symbols: list[str]) -> dict[str, float]: + """공개 ticker로 현재가 일괄 조회. 실패 시 빈 dict (호출측은 캔들 종가로 대체).""" + if not symbols: + return {} + try: + import requests + + markets = ",".join(f"KRW-{s.upper()}" for s in symbols) + url = f"{settings.api_url.rstrip('/')}/v1/ticker" + resp = requests.get(url, params={"markets": markets}, timeout=5) + resp.raise_for_status() + out: dict[str, float] = {} + for row in resp.json() or []: + m = str(row.get("market", "")) + if m.startswith("KRW-"): + out[m[4:].upper()] = float(row.get("trade_price") or 0.0) + return out + except Exception: # noqa: BLE001 + return {} + + +def _load_rsi_dca_trades(settings: Settings) -> list[dict[str, Any]]: + """RSI DCA 상태 파일의 체결 기록 (차트 마커용). 없으면 빈 목록.""" + path = getattr(settings, "rsi_dca_state_json", None) + if not path: + return [] + try: + p = Path(path) + if not p.exists(): + return [] + data = json.loads(p.read_text(encoding="utf-8")) + out = [] + for tr in data.get("trades") or []: + out.append({**tr, "reason": tr.get("reason") or "rsi_dca"}) + return out + except Exception: # noqa: BLE001 + return [] + + def build_vol_monitor_payload( settings: Settings, state: dict[str, Any], @@ -257,10 +296,22 @@ def build_vol_monitor_payload( tick_report: dict[str, Any] | None = None, ) -> dict[str, Any]: """모니터 JSON 페이로드.""" - symbols = list(settings.ops_symbols) + symbols = list(getattr(settings, "vol_monitor_symbols", None) or settings.ops_symbols) days = float(settings.vol_monitor_days or 14) sym_state = state.get("symbols") or {} trades = list(state.get("trades") or []) + rsi_trades = _load_rsi_dca_trades(settings) + # 전략 매수 체결(자동)의 종목별 누적 투입·수량 → 평균 매입가. 수동 보유분은 VOL_MONITOR_AVG_PRICE 로 보완 + cost_krw: dict[str, float] = {} + cost_coin: dict[str, float] = {} + for tr in rsi_trades: # RSI 자동 매수만 (매도 없는 전략이므로 누적 = 보유 원가) + if str(tr.get("side", "")) != "buy": + continue + s_ = str(tr.get("symbol", "")).upper() + cost_krw[s_] = cost_krw.get(s_, 0.0) + float(tr.get("order_krw") or 0.0) + cost_coin[s_] = cost_coin.get(s_, 0.0) + float(tr.get("order_coin") or 0.0) + manual_avg = dict(getattr(settings, "vol_monitor_avg_price", None) or {}) + live_px = fetch_ticker_prices(settings, symbols) snap = state.get("portfolio_snapshot") or {} cash = float(snap.get("cash_krw") or 0.0) @@ -276,7 +327,7 @@ def build_vol_monitor_payload( next_15m, sec_until = _next_15m_close(df_closed) st = sym_state.get(sym.upper()) or sym_state.get(sym) or {} qty = float((snap.get("positions") or {}).get(sym, 0) or 0.0) - price = float(df_closed["close"].iloc[-1]) if not df_closed.empty else 0.0 + price = float(live_px.get(sym.upper()) or (df_closed["close"].iloc[-1] if not df_closed.empty else 0.0)) holding = in_long_position( {"positions": {sym: {"coin_qty": qty}}}, sym, @@ -285,7 +336,23 @@ def build_vol_monitor_payload( ) if price > 0: total_equity += qty * price + # 평균매입가 우선순위: 거래소 계좌(avg_buy_price) > 자동매수 체결 누적 > 수동 설정(VOL_MONITOR_AVG_PRICE) + exch_avg = float((snap.get("avg_prices") or {}).get(sym.upper()) or 0.0) + avg_price = 0.0 + if exch_avg > 0: + avg_price = exch_avg + elif cost_coin.get(sym.upper(), 0.0) > 0: + avg_price = cost_krw[sym.upper()] / cost_coin[sym.upper()] + elif sym.upper() in manual_avg and manual_avg[sym.upper()] > 0: + avg_price = float(manual_avg[sym.upper()]) + pnl_pct = (price / avg_price - 1.0) * 100.0 if (avg_price > 0 and price > 0 and qty > 0) else None symbol_summary[sym] = { + "avg_price": round(avg_price, 6) if avg_price else None, + "avg_price_source": "exchange" if exch_avg > 0 else ("auto" if cost_coin.get(sym.upper(), 0.0) > 0 else ("manual" if avg_price else None)), + "holding_cost_krw": round(qty * avg_price, 0) if avg_price else None, # 현재 보유분 원금 (매도 시 즉시 감소) + "auto_cost_krw": round(cost_krw.get(sym.upper(), 0.0), 0), + "auto_buys": sum(1 for t in rsi_trades if str(t.get("symbol", "")).upper() == sym.upper()), + "pnl_pct": None if pnl_pct is None else round(pnl_pct, 2), "name": resolve_coin_name(sym), "in_position": holding, "coin_qty": qty, @@ -297,24 +364,26 @@ def build_vol_monitor_payload( } symbol_blocks[sym] = { "candles_15m": _candles_payload(df_closed, days=days), - "markers": _trade_markers(trades, sym), + "markers": _trade_markers(trades + rsi_trades, sym), } panel = _merged_close_panel(symbol_dfs, symbols, days=days) seed_krw = max(total_equity, 1.0) equity_strategy: list[dict[str, float | int]] = [] equity_buyhold: list[dict[str, float | int]] = [] - if not panel.empty: + # 캔들이 아직 없는 종목(신규 수집 중)은 수익률 곡선에서 제외 + panel_syms = [s for s in symbols if s.upper() in panel.columns] + if not panel.empty and panel_syms: window_start = pd.Timestamp(panel["datetime"].iloc[0]) equity_strategy = build_spot_strategy_equity_series( panel, - symbols, + panel_syms, trades, seed_krw=seed_krw, current_equity=total_equity, window_start=window_start, ) - equity_buyhold = build_multi_buyhold_series(panel, symbols, seed_krw) + equity_buyhold = build_multi_buyhold_series(panel, panel_syms, seed_krw) summary = { "mode": settings.ops_mode, @@ -336,16 +405,72 @@ def build_vol_monitor_payload( "trades": trades[-100:], "last_tick": tick_report or {}, "ops_symbols": symbols, + "intervals": list(getattr(settings, "vol_monitor_intervals", None) or [INTERVAL_MIN]), "equity": { "strategy": equity_strategy, "buyhold": equity_buyhold, "seed_krw": round(seed_krw, 0), "label_strategy": "vol_breakout", - "label_buyhold": "B&H 1/3×3", + "label_buyhold": f"B&H 1/{len(panel_syms)}×{len(panel_syms)}", }, } +_INTERVAL_LABELS = { + 1: "1분", 3: "3분", 5: "5분", 10: "10분", 15: "15분", 30: "30분", + 60: "1시간", 240: "4시간", 1440: "1일", 10080: "1주", 43200: "1월", +} + + +def interval_label(interval_min: int) -> str: + """분봉 코드 → 표시 라벨.""" + return _INTERVAL_LABELS.get(int(interval_min), f"{int(interval_min)}분") + + +def build_candles_api_payload( + settings: Settings, + symbol: str, + interval_min: int, + *, + max_bars: int | None = None, +) -> dict[str, Any]: + """`/api/candles` — DB 기준 특정 종목·분봉 최근 N봉 OHLC (lightweight-charts용).""" + sym = str(symbol or "").strip().upper() + allowed_syms = { + s.upper() for s in (getattr(settings, "vol_monitor_symbols", None) or settings.ops_symbols) + } + allowed_iv = set(getattr(settings, "vol_monitor_intervals", None) or [INTERVAL_MIN]) + if sym not in allowed_syms: + return {"ok": False, "error": f"unknown symbol: {sym}"} + try: + iv = int(interval_min) + except (TypeError, ValueError): + return {"ok": False, "error": f"bad interval: {interval_min}"} + if iv not in allowed_iv: + return {"ok": False, "error": f"interval not allowed: {iv}"} + limit = int(max_bars or getattr(settings, "vol_monitor_max_bars", 0) or 1500) + + df = load_last_candles(settings.db_path, sym, iv, limit) + rows: list[dict[str, float | int]] = [] + for _, row in df.iterrows(): + rows.append({ + "time": _epoch_kst(row["datetime"]), + "open": float(row["open"]), + "high": float(row["high"]), + "low": float(row["low"]), + "close": float(row["close"]), + }) + return { + "ok": True, + "symbol": sym, + "interval": iv, + "label": interval_label(iv), + "count": len(rows), + "last": str(df["datetime"].iloc[-1])[:19] if not df.empty else None, + "candles": rows, + } + + def write_vol_monitor( settings: Settings, state: dict[str, Any], @@ -412,17 +537,35 @@ def fetch_live_balance_snapshot(settings: Settings) -> dict[str, Any]: sleep_sec=settings.request_sleep_sec, retries=settings.request_retries, ) - krw, _ = client.get_balance("KRW") - positions: dict[str, float] = {} - total = float(krw) - for sym in settings.ops_symbols: - qty, _ = client.get_balance(sym) - positions[sym] = float(qty) - if qty > 0: - pass # price optional for total + # 모니터 표시 종목 + 매매 종목 + RSI 종목 전체를 한 번의 계정 조회로 채운다 + symbols: list[str] = [] + for group in ( + getattr(settings, "vol_monitor_symbols", None) or [], + settings.ops_symbols or [], + getattr(settings, "rsi_dca_symbols", None) or [], + ): + for sym in group: + if sym.upper() not in symbols: + symbols.append(sym.upper()) + accounts = client.get_accounts() + by_cur: dict[str, tuple[float, float]] = {} + avg_prices: dict[str, float] = {} + for acc in accounts or []: + cur = str(acc.get("currency", "")).upper() + try: + by_cur[cur] = (float(acc.get("balance") or 0.0), float(acc.get("locked") or 0.0)) + avg = float(acc.get("avg_buy_price") or 0.0) # 빗썸 계좌 평균 매입가 (수동·자동 매수 모두 반영) + if avg > 0: + avg_prices[cur] = avg + except (TypeError, ValueError): + continue + krw = by_cur.get("KRW", (0.0, 0.0))[0] + positions: dict[str, float] = {sym: by_cur.get(sym, (0.0, 0.0))[0] for sym in symbols} return { "ok": True, "cash_krw": round(float(krw), 0), "positions": positions, + "avg_prices": {sym: avg_prices[sym] for sym in symbols if sym in avg_prices}, + "prices": fetch_ticker_prices(settings, symbols), "updated_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), } diff --git a/src/bithumb/operations/vol_monitor_chart.py b/src/bithumb/operations/vol_monitor_chart.py index a4f8700..d454d50 100644 --- a/src/bithumb/operations/vol_monitor_chart.py +++ b/src/bithumb/operations/vol_monitor_chart.py @@ -24,6 +24,20 @@ _MONITOR_HTML = """ #btnUpdate:hover {{ background:#eee; }} #btnUpdate:disabled {{ opacity:0.55; cursor:wait; }} #meta {{ font-size:12px; color:#555; margin-bottom:6px; }} + .rsiBox {{ display:flex; align-items:center; gap:10px; flex-wrap:wrap; margin:6px 0 2px; padding:6px 10px; + border:1px solid #e5e5e5; border-radius:6px; background:#fafafa; font-size:12px; }} + .rsiBox .ttl {{ font-weight:bold; color:#333; }} + .badge {{ display:inline-block; padding:1px 8px; border-radius:10px; font-size:11px; color:#fff; background:#888; }} + .badge.live {{ background:#c62828; }} + .badge.paper {{ background:#607d8b; }} + .badge.on {{ background:#2e7d32; }} + .badge.off {{ background:#9e9e9e; }} + .badge.dead {{ background:#ef6c00; }} + #btnRsiToggle {{ font-size:12px; padding:3px 14px; border-radius:4px; border:1px solid #999; cursor:pointer; background:#fff; }} + #btnRsiToggle.on {{ background:#2e7d32; color:#fff; border-color:#2e7d32; }} + #btnRsiToggle.off {{ background:#eee; color:#333; }} + #btnRsiToggle:disabled {{ opacity:0.55; cursor:wait; }} + .rsiBox .dim {{ color:#777; }} table.summary {{ border-collapse:collapse; font-size:12px; width:100%; margin-top:6px; }} table.summary th, table.summary td {{ border:1px solid #ddd; padding:5px 10px; text-align:center; white-space:nowrap; }} table.summary th {{ background:#f5f5f5; color:#666; font-weight:normal; font-size:11px; }} @@ -33,35 +47,66 @@ _MONITOR_HTML = """ cursor:pointer; font-size:13px; }} .tab.active {{ background:#333; color:#fff; border-color:#333; }} - #priceChart {{ width:100%; height:44vh; border-bottom:1px solid #eee; }} - #equityChart {{ width:100%; height:32vh; }} + .tabs.intervals {{ padding-top:0; border-bottom:1px solid #eee; }} + .tabs.intervals .tab {{ padding:3px 10px; font-size:12px; }} + .tabs.intervals .tab.active {{ background:#555; border-color:#555; }} + .tabs.intervals .label {{ font-size:12px; color:#777; align-self:center; margin-right:4px; }} + #priceChart {{ width:100%; height:40vh; border-bottom:1px solid #eee; }} + #rsiChart {{ width:100%; height:18vh; border-bottom:1px solid #eee; }} + #equityChart {{ width:100%; height:24vh; }} + .panelLabel {{ font-size:11px; color:#888; padding:2px 14px 0; }} #err {{ display:none; padding:10px 14px; color:#b91c1c; font-size:13px; }} + #stale {{ display:none; padding:8px 14px; background:#fff7e6; color:#8a5a00; border-bottom:1px solid #f3d9a4; font-size:13px; }} #reloadHint {{ font-size:12px; color:#888; padding:4px 14px; }}
-

Bithumb vol_breakout (15m spot long)

+

Bithumb 라이브 모니터

로딩 중…
+
+ RSI 자동매수 + - + - + + tick - + 오늘 - + +
+
+
+
RSI(14)
+
+
수익률(%) · 전략 vs B&H