377 lines
16 KiB
Python
377 lines
16 KiB
Python
"""
|
|
텔레그램 봇 최적화 버전
|
|
- Polling 최적화 (CPU 사용률 감소)
|
|
- 별도 프로세스로 분리
|
|
- 봇 재시작 명령어
|
|
- 원격 명령어 실행
|
|
"""
|
|
import os
|
|
import asyncio
|
|
import logging
|
|
import subprocess
|
|
import sys
|
|
from telegram import Update
|
|
from telegram.ext import Application, CommandHandler, ContextTypes
|
|
from dotenv import load_dotenv
|
|
|
|
# 로깅 설정
|
|
logging.basicConfig(
|
|
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
|
|
level=logging.INFO
|
|
)
|
|
logging.getLogger("httpx").setLevel(logging.WARNING)
|
|
|
|
class TelegramBotServer:
|
|
def __init__(self, bot_token):
|
|
# [최적화] 연결 풀 설정 추가
|
|
self.application = Application.builder()\
|
|
.token(bot_token)\
|
|
.concurrent_updates(True)\
|
|
.build()
|
|
|
|
self.bot_instance = None
|
|
self.is_shutting_down = False
|
|
self.should_restart = False
|
|
|
|
def set_bot_instance(self, bot):
|
|
"""AutoTradingBot 인스턴스를 주입받음"""
|
|
self.bot_instance = bot
|
|
|
|
def refresh_bot_instance(self):
|
|
"""IPC에서 최신 봇 인스턴스 데이터 읽기"""
|
|
from modules.utils.ipc import BotIPC
|
|
ipc = BotIPC()
|
|
self.bot_instance = ipc.get_bot_instance_data()
|
|
return self.bot_instance is not None
|
|
|
|
async def start_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/start 명령어 핸들러"""
|
|
print(f"📨 [Telegram] /start received from user {update.effective_user.id}")
|
|
await update.message.reply_text(
|
|
"🤖 <b>AI Trading Bot Command Center</b>\n"
|
|
"명령어 목록:\n"
|
|
"/status - 현재 봇 및 시장 상태 조회\n"
|
|
"/portfolio - 현재 보유 종목 및 평가액\n"
|
|
"/watchlist - 현재 감시 중인 종목 리스트\n"
|
|
"/update_watchlist - Watchlist 즉시 업데이트\n"
|
|
"/macro - 거시경제 지표 및 시장 위험도\n"
|
|
"/system - PC 리소스(CPU/GPU) 상태\n"
|
|
"/ai - AI 모델 학습 상태 조회\n\n"
|
|
"<b>[관리 명령어]</b>\n"
|
|
"/restart - 봇 재시작\n"
|
|
"/exec <code>명령어</code> - 원격 명령어 실행\n"
|
|
"/stop - 봇 종료",
|
|
parse_mode="HTML"
|
|
)
|
|
|
|
async def status_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/status: 종합 상태 브리핑"""
|
|
if not self.refresh_bot_instance():
|
|
await update.message.reply_text("⚠️ 메인 봇이 실행 중이 아닙니다.")
|
|
return
|
|
|
|
from datetime import datetime
|
|
now = datetime.now()
|
|
is_market_open = (9 <= now.hour < 15) or (now.hour == 15 and now.minute < 30)
|
|
|
|
status_msg = "✅ <b>System Status: ONLINE</b>\n"
|
|
status_msg += f"🕒 <b>Market:</b> {'OPEN 🟢' if is_market_open else 'CLOSED 🔴'}\n"
|
|
|
|
macro_warn = self.bot_instance.is_macro_warning_sent
|
|
status_msg += f"🌍 <b>Macro Filter:</b> {'DANGER 🚨 (Trading Halted)' if macro_warn else 'SAFE 🟢'}\n"
|
|
|
|
await update.message.reply_text(status_msg, parse_mode="HTML")
|
|
|
|
async def portfolio_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/portfolio: 잔고 조회"""
|
|
if not self.refresh_bot_instance():
|
|
await update.message.reply_text("⚠️ 봇 인스턴스가 연결되지 않았습니다.")
|
|
return
|
|
|
|
await update.message.reply_text("⏳ 잔고를 조회 중입니다...")
|
|
|
|
try:
|
|
balance = self.bot_instance.kis.get_balance()
|
|
if "error" in balance:
|
|
await update.message.reply_text(f"❌ 잔고 조회 실패: {balance['error']}")
|
|
return
|
|
|
|
msg = f"💰 <b>Total Asset:</b> <code>{int(balance['total_eval']):,} KRW</code>\n" \
|
|
f"💵 <b>Deposit:</b> <code>{int(balance['deposit']):,} KRW</code>\n\n"
|
|
|
|
if balance['holdings']:
|
|
msg += "<b>[Holdings]</b>\n"
|
|
for stock in balance['holdings']:
|
|
icon = "🔴" if stock['yield'] > 0 else "🔵"
|
|
msg += f"{icon} <b>{stock['name']}</b> <code>{stock['yield']}%</code>\n" \
|
|
f" (수량: {stock['qty']} / 평가손익: {stock['profit_loss']:,})\n"
|
|
else:
|
|
msg += "보유 중인 종목이 없습니다."
|
|
|
|
await update.message.reply_text(msg, parse_mode="HTML")
|
|
|
|
except Exception as e:
|
|
await update.message.reply_text(f"❌ Error: {str(e)}")
|
|
|
|
async def watchlist_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/watchlist: 감시 대상 종목"""
|
|
if not self.refresh_bot_instance():
|
|
await update.message.reply_text("⚠️ 봇 인스턴스가 연결되지 않았습니다.")
|
|
return
|
|
|
|
target_dict = self.bot_instance.load_watchlist()
|
|
discovered = list(self.bot_instance.discovered_stocks)
|
|
|
|
msg = f"👀 <b>Watchlist: {len(target_dict)} items</b>\n"
|
|
for code, name in target_dict.items():
|
|
themes = self.bot_instance.theme_manager.get_themes(code)
|
|
theme_str = f" ({', '.join(themes)})" if themes else ""
|
|
msg += f"- {name}{theme_str}\n"
|
|
|
|
if discovered:
|
|
msg += f"\n✨ <b>Discovered Today ({len(discovered)}):</b>\n"
|
|
for code in discovered:
|
|
msg += f"- {code}\n"
|
|
|
|
await update.message.reply_text(msg, parse_mode="HTML")
|
|
|
|
async def update_watchlist_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/update_watchlist: Watchlist 즉시 업데이트"""
|
|
await update.message.reply_text("🔄 Watchlist를 업데이트하고 있습니다... (30초 소요)")
|
|
|
|
try:
|
|
from modules.services.kis import KISClient
|
|
from watchlist_manager import WatchlistManager
|
|
from modules.config import Config
|
|
|
|
temp_kis = KISClient()
|
|
mgr = WatchlistManager(temp_kis, watchlist_file=Config.WATCHLIST_FILE)
|
|
|
|
summary = mgr.update_watchlist_daily()
|
|
# HTML 특수문자 이스케이프
|
|
summary = summary.replace("&", "&").replace("<", "<").replace(">", ">")
|
|
await update.message.reply_text(summary)
|
|
|
|
except Exception as e:
|
|
await update.message.reply_text(f"❌ 업데이트 실패: {e}")
|
|
|
|
async def macro_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/macro: 거시경제 지표 조회 (IPC 데이터 사용)"""
|
|
if not self.refresh_bot_instance():
|
|
await update.message.reply_text("⚠️ 메인 봇 연결 대기 중...")
|
|
return
|
|
|
|
await update.message.reply_text("⏳ 거시경제 데이터를 불러옵니다...")
|
|
|
|
try:
|
|
indices = getattr(self.bot_instance.kis, '_macro_indices', {})
|
|
|
|
if not indices:
|
|
await update.message.reply_text("⚠️ 데이터가 아직 수집되지 않았습니다. 잠시 후 다시 시도하세요.")
|
|
return
|
|
|
|
status = "SAFE"
|
|
msi = indices.get('MSI', 0)
|
|
if msi >= 50: status = "DANGER"
|
|
elif msi >= 30: status = "CAUTION"
|
|
|
|
color = "🟢" if status == "SAFE" else "🔴" if status == "DANGER" else "🟡"
|
|
msg = f"{color} <b>Market Risk: {status}</b>\n\n"
|
|
|
|
if 'MSI' in indices:
|
|
msg += f"🌡️ <b>Stress Index:</b> <code>{indices['MSI']}</code>\n"
|
|
|
|
for k, v in indices.items():
|
|
if k != "MSI":
|
|
icon = "🔺" if v.get('change', 0) > 0 else "🔻"
|
|
msg += f"{icon} <b>{k}</b>: {v.get('price', 0)} ({v.get('change', 0)}%)\n"
|
|
|
|
await update.message.reply_text(msg, parse_mode="HTML")
|
|
|
|
except Exception as e:
|
|
await update.message.reply_text(f"❌ Error: {e}")
|
|
|
|
async def system_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/system: 시스템 리소스 상태"""
|
|
if not self.refresh_bot_instance():
|
|
await update.message.reply_text("⚠️ 메인 봇이 실행 중이 아닙니다.")
|
|
return
|
|
|
|
import psutil
|
|
|
|
cpu = psutil.cpu_percent(interval=1)
|
|
ram = psutil.virtual_memory().percent
|
|
|
|
# CPU 점유율 상위 3개 프로세스 수집
|
|
top_processes = []
|
|
for proc in psutil.process_iter(['pid', 'name', 'cpu_percent']):
|
|
try:
|
|
proc_info = proc.info
|
|
if proc_info['name'] == 'System Idle Process':
|
|
continue
|
|
top_processes.append(proc_info)
|
|
except (psutil.NoSuchProcess, psutil.AccessDenied, psutil.ZombieProcess):
|
|
pass
|
|
|
|
top_processes.sort(key=lambda x: x.get('cpu_percent', 0), reverse=True)
|
|
top_3 = top_processes[:3]
|
|
|
|
gpu_status = self.bot_instance.ollama_monitor.get_gpu_status()
|
|
gpu_msg = "N/A"
|
|
if gpu_status and gpu_status.get('name') != 'N/A':
|
|
gpu_name = gpu_status.get('name', 'GPU')
|
|
gpu_msg = f"{gpu_name}\n Temp: {gpu_status.get('temp', 0)}°C / VRAM: {gpu_status.get('vram_used', 0)}GB / {gpu_status.get('vram_total', 0)}GB"
|
|
|
|
msg = "🖥️ <b>PC System Status</b>\n" \
|
|
f"🧠 <b>CPU:</b> <code>{cpu}%</code>\n" \
|
|
f"💾 <b>RAM:</b> <code>{ram}%</code>\n" \
|
|
f"🎮 <b>GPU:</b> {gpu_msg}\n\n"
|
|
|
|
if top_3:
|
|
msg += "⚙️ <b>Top CPU Processes:</b>\n"
|
|
for i, proc in enumerate(top_3, 1):
|
|
proc_name = proc.get('name', 'Unknown')
|
|
proc_cpu = proc.get('cpu_percent', 0)
|
|
msg += f" {i}. <code>{proc_name}</code> - {proc_cpu:.1f}%\n"
|
|
|
|
await update.message.reply_text(msg, parse_mode="HTML")
|
|
|
|
async def ai_status_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/ai: AI 모델 학습 상태 조회"""
|
|
if not self.refresh_bot_instance():
|
|
await update.message.reply_text("⚠️ 메인 봇이 실행 중이 아닙니다.")
|
|
return
|
|
|
|
gpu = self.bot_instance.ollama_monitor.get_gpu_status()
|
|
|
|
msg = "🧠 <b>AI Model Status</b>\n"
|
|
msg += "• <b>LLM Engine:</b> Ollama (Llama 3.1)\n"
|
|
|
|
gpu_name = gpu.get('name', 'NVIDIA RTX 5070 Ti')
|
|
msg += f"• <b>Device:</b> {gpu_name}\n"
|
|
|
|
if gpu:
|
|
msg += f"• <b>GPU Load:</b> <code>{gpu.get('load', 0)}%</code>\n"
|
|
msg += f"• <b>VRAM Usage:</b> <code>{gpu.get('vram_used', 0)}GB</code> / {gpu.get('vram_total', 0)}GB"
|
|
|
|
await update.message.reply_text(msg, parse_mode="HTML")
|
|
|
|
async def restart_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/restart: 텔레그램 봇 모듈만 재시작"""
|
|
await update.message.reply_text("🔄 <b>텔레그램 인터페이스를 재시작합니다...</b>", parse_mode="HTML")
|
|
|
|
self.should_restart = True
|
|
self.application.stop_running()
|
|
|
|
async def stop_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/stop: 봇 종료"""
|
|
await update.message.reply_text("🛑 <b>텔레그램 봇을 종료합니다.</b>", parse_mode="HTML")
|
|
|
|
self.should_restart = False
|
|
self.application.stop_running()
|
|
|
|
async def exec_command(self, update: Update, context: ContextTypes.DEFAULT_TYPE):
|
|
"""/exec: 원격 명령어 실행 (Non-blocking)"""
|
|
text = update.message.text.strip()
|
|
parts = text.split(maxsplit=1)
|
|
|
|
if len(parts) < 2:
|
|
await update.message.reply_text("❌ 사용법: /exec 명령어")
|
|
return
|
|
|
|
command = parts[1]
|
|
await update.message.reply_text(f"⚙️ 실행 중: <code>{command}</code>", parse_mode="HTML")
|
|
|
|
try:
|
|
# 보안: 위험한 명령어 차단
|
|
dangerous_keywords = ['rm', 'del', 'format', 'shutdown', 'reboot']
|
|
if any(keyword in command.lower() for keyword in dangerous_keywords):
|
|
await update.message.reply_text("⛔ 위험한 명령어는 실행할 수 없습니다.")
|
|
return
|
|
|
|
import platform
|
|
is_windows = platform.system() == 'Windows'
|
|
|
|
if is_windows:
|
|
exec_cmd = ['powershell', '-Command', command]
|
|
else:
|
|
exec_cmd = command
|
|
|
|
def run_subprocess():
|
|
return subprocess.run(
|
|
exec_cmd,
|
|
shell=not is_windows,
|
|
capture_output=True,
|
|
text=True,
|
|
encoding='utf-8',
|
|
errors='replace',
|
|
timeout=30,
|
|
cwd=os.getcwd()
|
|
)
|
|
|
|
loop = asyncio.get_running_loop()
|
|
result = await loop.run_in_executor(None, run_subprocess)
|
|
|
|
output = result.stdout.strip() if result.stdout else ""
|
|
error_output = result.stderr.strip() if result.stderr else ""
|
|
|
|
if output and error_output:
|
|
combined = f"[STDOUT]\n{output}\n\n[STDERR]\n{error_output}"
|
|
elif output:
|
|
combined = output
|
|
elif error_output:
|
|
combined = f"[ERROR]\n{error_output}"
|
|
else:
|
|
combined = "명령어 실행 완료 (출력 없음)"
|
|
|
|
if len(combined) > 3000:
|
|
combined = combined[:3000] + "\n... (Truncated)"
|
|
|
|
# HTML 특수문자 이스케이프
|
|
combined = combined.replace("&", "&").replace("<", "<").replace(">", ">")
|
|
await update.message.reply_text(f"<pre>{combined}</pre>", parse_mode="HTML")
|
|
|
|
except asyncio.TimeoutError:
|
|
await update.message.reply_text("⏱️ 명령어 실행 시간 초과 (30초)")
|
|
except Exception as e:
|
|
await update.message.reply_text(f"❌ 실행 오류: {e}")
|
|
|
|
def run(self):
|
|
"""봇 실행 (Handler 등록 및 Polling)"""
|
|
handlers = [
|
|
("start", self.start_command),
|
|
("status", self.status_command),
|
|
("portfolio", self.portfolio_command),
|
|
("watchlist", self.watchlist_command),
|
|
("update_watchlist", self.update_watchlist_command),
|
|
("macro", self.macro_command),
|
|
("system", self.system_command),
|
|
("ai", self.ai_status_command),
|
|
("restart", self.restart_command),
|
|
("stop", self.stop_command),
|
|
("exec", self.exec_command)
|
|
]
|
|
|
|
for cmd, func in handlers:
|
|
self.application.add_handler(CommandHandler(cmd, func))
|
|
|
|
async def error_handler(update: object, context: ContextTypes.DEFAULT_TYPE) -> None:
|
|
if "Conflict" in str(context.error):
|
|
print(f"⚠️ [Telegram] Conflict detected. Stopping...")
|
|
if self.application.running:
|
|
await self.application.stop()
|
|
return
|
|
print(f"❌ [Telegram Error] {context.error}")
|
|
|
|
self.application.add_error_handler(error_handler)
|
|
|
|
print("🤖 [Telegram] Command Server Started (Standard Polling Mode).")
|
|
|
|
try:
|
|
self.application.run_polling(
|
|
allowed_updates=Update.ALL_TYPES,
|
|
drop_pending_updates=True
|
|
)
|
|
except Exception as e:
|
|
print(f"❌ [Telegram] Polling Error: {e}")
|