message_handler.py

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
消息处理模块
负责发送选股结果和实时监控结果
"""

import os
import logging
import requests
from datetime import datetime, timedelta
from dotenv import load_dotenv

# 加载环境变量
load_dotenv()

# 配置日志
logger = logging.getLogger(__name__)

class MessageHandler:
    """消息处理类"""
    
    def __init__(self):
        """初始化消息处理对象"""
        # 配置Feishu消息发送
        self.feishu_webhook_url = os.getenv("FEISHU_WEBHOOK_URL")
        self.feishu_token = os.getenv("FEISHU_TOKEN")
    
    def send_message(self, content):
        """发送消息"""
        try:
            if self.feishu_webhook_url:
                self._send_feishu_message(content)
            else:
                logger.warning("未配置Feishu Webhook URL,消息发送失败")
        except Exception as e:
            logger.error(f"消息发送失败:{e}")
    
    def _send_feishu_message(self, content):
        """发送Feishu消息"""
        try:
            headers = {
                "Content-Type": "application/json"
            }
            
            data = {
                "msg_type": "post",
                "content": {
                    "zh_cn": {
                        "content": [
                            [
                                {
                                    "tag": "text",
                                    "text": content
                                }
                            ]
                        ]
                    }
                }
            }
            
            response = requests.post(self.feishu_webhook_url, headers=headers, json=data, timeout=10)
            
            if response.status_code == 200:
                logger.info("Feishu消息发送成功")
            else:
                logger.error(f"Feishu消息发送失败,状态码:{response.status_code}")
        
        except Exception as e:
            logger.error(f"Feishu消息发送失败:{e}")
    
    def send_stock_selection_result(self, selected_stocks):
        """发送选股结果消息"""
        try:
            content = "【选股结果】\n"
            for i, (ts_code, name) in enumerate(selected_stocks.items(), 1):
                content += f"{i}. {name} ({ts_code})\n"
            
            self.send_message(content)
            logger.info("选股结果消息发送成功")
        
        except Exception as e:
            logger.error(f"选股结果消息发送失败:{e}")
    
    def send_monitoring_result(self, ts_code, name, current_price, optimal_price, action):
        """发送实时监控结果消息"""
        try:
            content = f"【实时监控】\n股票:{name} ({ts_code})\n当前价格:{current_price:.2f}\n最佳卖出价格:{optimal_price:.2f}\n建议动作:{action}"
            
            self.send_message(content)
            logger.info("实时监控结果消息发送成功")
        
        except Exception as e:
            logger.error(f"实时监控结果消息发送失败:{e}")
    
    def test_message_sending(self):
        """测试消息发送功能"""
        try:
            logger.info("开始测试消息发送功能...")
            
            # 测试发送选股结果消息
            test_stocks = {"000001.SZ": "平安银行", "000002.SZ": "万科A"}
            self.send_stock_selection_result(test_stocks)
            
            # 测试发送实时监控结果消息
            self.send_monitoring_result("000001.SZ", "平安银行", 10.5, 11.0, "卖出")
            
            logger.info("消息发送功能测试完成")
            
        except Exception as e:
            logger.error(f"消息发送功能测试失败:{e}")

if __name__ == "__main__":
    # 配置日志
    logging.basicConfig(
        level=logging.INFO,
        format="%(asctime)s - %(name)s - %(levelname)s - %(message)s"
    )
    
    # 测试消息处理模块
    message_handler = MessageHandler()
    message_handler.test_message_sending()