From c85f481a184fa2279f27c1f3bd7c3b13c134cf84 Mon Sep 17 00:00:00 2001 From: Cxx0822 <1556464090@qq.com> Date: Mon, 27 Oct 2025 19:05:54 +0800 Subject: [PATCH] =?UTF-8?q?feat:=E5=88=9D=E5=A7=8B=E5=8C=96=E5=B7=A5?= =?UTF-8?q?=E7=A8=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 67 ++++++++++++++++++ config/logger_config.py | 33 +++++++++ config/setting.py | 16 +++++ endpoint/mqtt_client.py | 136 +++++++++++++++++++++++++++++++++++++ endpoint/tcp_server.py | 84 +++++++++++++++++++++++ main.py | 22 ++++++ model/index.py | 1 + parser/abstract_parsers.py | 26 +++++++ parser/test_parser.py | 16 +++++ processor/tcp_processor.py | 25 +++++++ requirements.txt | 3 + 11 files changed, 429 insertions(+) create mode 100644 .gitignore create mode 100644 config/logger_config.py create mode 100644 config/setting.py create mode 100644 endpoint/mqtt_client.py create mode 100644 endpoint/tcp_server.py create mode 100644 main.py create mode 100644 model/index.py create mode 100644 parser/abstract_parsers.py create mode 100644 parser/test_parser.py create mode 100644 processor/tcp_processor.py create mode 100644 requirements.txt diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..97f9175 --- /dev/null +++ b/.gitignore @@ -0,0 +1,67 @@ +# Python 字节码文件 +__pycache__/ +*.py[cod] +*$py.class + +# C 扩展 +*.so + +# 分发/打包 +.Python +build/ +develop-eggs/ +dist/ +downloads/ +eggs/ +.eggs/ +lib/ +lib64/ +parts/ +sdist/ +var/ +wheels/ +*.egg-info/ +.installed.cfg +*.egg + +# 虚拟环境 +venv/ +env/ +ENV/ +.env +.venv + +# 测试 +htmlcov/ +.tox/ +.nox/ +.coverage +.coverage.* +.cache +nosetests.xml +coverage.xml +*.cover +.hypothesis/ + +# Django 相关 +*.log +local_settings.py +db.sqlite3 +db.sqlite3-journal +media/ + +# PyCharm IDE +.idea/ +*.iml +*.iws +*.ipr + +# VS Code +.vscode/ +*.code-workspace +.history/ + +# 其他 +.DS_Store + +tcp-logs/ \ No newline at end of file diff --git a/config/logger_config.py b/config/logger_config.py new file mode 100644 index 0000000..e278841 --- /dev/null +++ b/config/logger_config.py @@ -0,0 +1,33 @@ +import sys + +from loguru import logger + +# 日志格式 +STDOUT_FORMAT = ( + "{time:YYYY-MM-DD HH:mm:ss.SSS} | " + "{level: <8} | " + "{name}:{function}:{line} - " + "{message}" +) + +FILE_FORMAT = ( + "{time:YYYY-MM-DD HH:mm:ss.SSS} | " + "{level: <8} | " + "{name}:{function}:{line} - {message}" +) + +# 移除默认处理器 +logger.remove() + +# 添加控制台处理器 +logger.add( + sink=sys.stdout, + level="INFO", + format=STDOUT_FORMAT, + colorize=True, + backtrace=True, # 显示完整异常堆栈 + diagnose=True, # 显示详细异常信息 +) + +# 导出配置好的logger +__all__ = ["logger"] diff --git a/config/setting.py b/config/setting.py new file mode 100644 index 0000000..8ccc0bc --- /dev/null +++ b/config/setting.py @@ -0,0 +1,16 @@ +from pydantic_settings import BaseSettings + + +class Settings(BaseSettings): + tcp_host: str + tcp_port: int + mqtt_host: str + mqtt_port: int + + class Config: + env_file = ".env" # 指定.env文件路径 + case_sensitive = False # 忽略变量名大小写 + + +# 全局配置实例 +settings = Settings() diff --git a/endpoint/mqtt_client.py b/endpoint/mqtt_client.py new file mode 100644 index 0000000..b30180a --- /dev/null +++ b/endpoint/mqtt_client.py @@ -0,0 +1,136 @@ +import paho.mqtt.client as mqtt +import uuid +from config.logger_config import logger +from config.setting import settings + + +class MQTTClient: + def __init__(self): + self.broker_host = settings.mqtt_host + self.broker_port = settings.mqtt_port + self.client_id = self.generate_client_id() + + self.client = mqtt.Client(client_id=self.client_id) + + # 注册回调函数 + self.client.on_connect = self._on_connect + self.client.on_disconnect = self._on_disconnect + self.client.on_message = self._on_message + + # 存储消息处理器的字典,键为主题,值为处理函数 + self.message_handlers = {} + + self.connected = False + + def _on_connect(self, client, userdata, flags, rc): + """连接回调函数""" + self.connected = rc == 0 + if self.connected: + logger.info(f"成功连接到MQTT broker {self.broker_host}:{self.broker_port}") + else: + logger.error(f"连接MQTT broker失败,错误代码: {rc}") + + def _on_disconnect(self, client, userdata, rc): + """断开连接回调函数""" + self.connected = False + if rc != 0: + logger.warning(f"意外断开与MQTT broker的连接,错误代码: {rc}") + else: + logger.info("已与MQTT broker断开连接") + + def _on_message(self, client, userdata, msg): + """消息接收回调函数""" + logger.info(f"收到消息: {msg.topic}") + try: + if msg.topic in self.message_handlers: + self.message_handlers[msg.topic](client, msg) + else: + logger.warning(f"未找到 {msg.topic} 的消息处理器") + + except Exception as e: + logger.error(f"处理消息时出错: {str(e)}", exc_info=True) + + @staticmethod + def generate_client_id(prefix: str = "client") -> str: + """生成唯一客户端ID""" + uuid_str = str(uuid.uuid4()).split('-')[0] + return f"{prefix}-{uuid_str}" + + def add_message_handler(self, topic, handler): + """添加消息处理器""" + if topic in self.message_handlers: + logger.warning(f"主题 {topic} 已存在处理器,将被覆盖") + self.message_handlers[topic] = handler + logger.info(f"为主题 {topic} 注册了消息处理器") + logger.debug(f"当前消息处理器: {self.message_handlers.keys()}") + + def connect(self): + """连接到MQTT服务器""" + try: + self.client.connect(self.broker_host, self.broker_port) + self.start_loop() + return True + except Exception as e: + logger.error(f"连接MQTT broker {self.broker_host}:{self.broker_port}时发生错误: {str(e)}", exc_info=True) + return False + + def disconnect(self): + """断开与MQTT服务器的连接""" + try: + self.stop_loop() + self.client.disconnect() + logger.info("正在断开与MQTT broker的连接") + except Exception as e: + logger.error(f"断开连接时发生错误: {str(e)}", exc_info=True) + + def subscribe(self, topic, qos=0): + """订阅主题""" + try: + result, mid = self.client.subscribe(topic, qos) + if result == mqtt.MQTT_ERR_SUCCESS: + logger.info(f"已订阅主题: {topic} (QoS: {qos})") + return True + else: + logger.error(f"订阅主题 {topic} 失败,错误代码: {result}") + return False + except Exception as e: + logger.error(f"订阅主题时发生错误: {str(e)}", exc_info=True) + return False + + def publish(self, topic, payload, qos=0, retain=False): + """发布消息到指定主题""" + if not self.connected: + logger.warning("未连接到MQTT broker,无法发布消息") + return False + + try: + result = self.client.publish(topic, payload, qos, retain) + result.wait_for_publish() + if result.rc == mqtt.MQTT_ERR_SUCCESS: + logger.debug(f"已发布消息到主题 {topic}: {payload}") + return True + else: + logger.error(f"发布消息到主题 {topic} 失败,错误代码: {result.rc}") + return False + except Exception as e: + logger.error(f"发布消息时发生错误: {str(e)}", exc_info=True) + return False + + def start_loop(self): + """启动MQTT网络循环""" + try: + self.client.loop_start() + logger.info("已启动MQTT网络循环") + except Exception as e: + logger.error(f"启动网络循环时发生错误: {str(e)}", exc_info=True) + + def stop_loop(self): + """停止MQTT网络循环""" + try: + self.client.loop_stop() + logger.info("已停止MQTT网络循环") + except Exception as e: + logger.error(f"停止网络循环时发生错误: {str(e)}", exc_info=True) + + +mqtt_client = MQTTClient() diff --git a/endpoint/tcp_server.py b/endpoint/tcp_server.py new file mode 100644 index 0000000..5cf4b41 --- /dev/null +++ b/endpoint/tcp_server.py @@ -0,0 +1,84 @@ +import asyncio +import signal +import platform +from config.logger_config import logger +from config.setting import settings +from model.index import UNKNOWN_MESSAGE +from processor.tcp_processor import TcpMessageProcessor + + +class TCPServer: + def __init__(self): + self.host = settings.tcp_host + self.port = settings.tcp_port + self.server = None + self.clients = set() + self.running = False + self.message_processor = TcpMessageProcessor() + + async def handle_client(self, reader, writer): + """处理客户端连接""" + addr = writer.get_extra_info('peername') + logger.info(f"新客户端连接: {addr}") + self.clients.add(writer) + + try: + while self.running: + data = await reader.read(1024) + if not data: + break # 连接断开 + + message = data.decode('utf-8').strip() + logger.info(f"收到 {addr} 的消息: {message}") + + # 放到消息处理器里面处理 + response = self.message_processor.process(message) + + # 回复消息 + if response != UNKNOWN_MESSAGE: + writer.write(response.encode()) + await writer.drain() + logger.info(f"回复 {addr}: {response}") + except Exception as e: + logger.error(f"客户端 {addr} 处理错误: {e}") + finally: + if writer in self.clients: + self.clients.remove(writer) + writer.close() + await writer.wait_closed() + logger.warning(f"客户端 {addr} 断开") + + async def start(self): + """启动服务器""" + self.running = True + self.server = await asyncio.start_server(self.handle_client, self.host, self.port) + + # 仅在非Windows平台设置信号处理 + if platform.system() != 'Windows': + loop = asyncio.get_running_loop() + for sig in (signal.SIGINT, signal.SIGTERM): + loop.add_signal_handler(sig, self.stop) + + logger.info(f"TCP服务器已启动 {self.host}:{self.port}") + async with self.server: + await self.server.serve_forever() + + async def _safe_stop(self): + """安全的停止服务器""" + self.running = False + if self.server: + self.server.close() + await self.server.wait_closed() + + # 关闭所有客户端连接 + for writer in list(self.clients): + writer.close() + await writer.wait_closed() + + def stop(self): + """停止服务器""" + logger.info("正在关闭TCP服务器...") + asyncio.create_task(self._safe_stop()) + + +tcp_server = TCPServer() diff --git a/main.py b/main.py new file mode 100644 index 0000000..418d3d5 --- /dev/null +++ b/main.py @@ -0,0 +1,22 @@ +import asyncio + +from endpoint.tcp_server import TCPServer +from config.logger_config import logger + + +def start_tcp_server(): + server = TCPServer() + try: + asyncio.run(server.start()) + except KeyboardInterrupt: + logger.info("收到中断信号,正在关闭服务器...") + server.stop() + except Exception as e: + logger.error(f"服务器运行错误: {e}") + server.stop() + finally: + logger.info("服务器进程结束") + + +if __name__ == "__main__": + start_tcp_server() diff --git a/model/index.py b/model/index.py new file mode 100644 index 0000000..e236af2 --- /dev/null +++ b/model/index.py @@ -0,0 +1 @@ +UNKNOWN_MESSAGE = "unknown" diff --git a/parser/abstract_parsers.py b/parser/abstract_parsers.py new file mode 100644 index 0000000..7efe825 --- /dev/null +++ b/parser/abstract_parsers.py @@ -0,0 +1,26 @@ +from abc import ABC, abstractmethod + + +class AbstractMessageParser(ABC): + """设备消息解析器基类""" + + @property + @abstractmethod + def name(self): + """消息名称,用于标识该解析器处理的消息类型""" + pass + + @abstractmethod + def check(self, message) -> bool: + """判断当前解析器是否能处理该消息""" + pass + + @abstractmethod + def parse(self, message): + """解析消息内容""" + pass + + @abstractmethod + def response(self, parsed_data): + """根据解析后的数据生成响应""" + pass diff --git a/parser/test_parser.py b/parser/test_parser.py new file mode 100644 index 0000000..7dbb36c --- /dev/null +++ b/parser/test_parser.py @@ -0,0 +1,16 @@ +from parser.abstract_parsers import AbstractMessageParser + + +class TestParser(AbstractMessageParser): + @property + def name(self): + return "name" + + def check(self, message) -> bool: + return True + + def parse(self, message): + return message + + def response(self, parsed_data): + return parsed_data diff --git a/processor/tcp_processor.py b/processor/tcp_processor.py new file mode 100644 index 0000000..6f6ac58 --- /dev/null +++ b/processor/tcp_processor.py @@ -0,0 +1,25 @@ +from config.logger_config import logger +from model.index import UNKNOWN_MESSAGE +from parser.test_parser import TestParser + + +class TcpMessageProcessor: + """消息处理器,负责将消息路由到正确的设备解析器""" + + def __init__(self): + # 注册所有支持的设备解析器 + self.parsers = [TestParser()] + + def process(self, message): + """处理消息,返回响应""" + # 尝试找到能处理该消息的解析器 + for parser in self.parsers: + if parser.check(message): + parsed_data = parser.parse(message) + if parsed_data: + logger.info(f"处理{parser.name}消息: {message}") + return parser.response(parsed_data) + + # 没有找到合适的解析器 + logger.warning(f"未识别的消息格式: {message}") + return UNKNOWN_MESSAGE diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..c9091c3 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,3 @@ +loguru~=0.7.3 +pydantic-settings~=2.11.0 +paho-mqtt~=2.1.0 \ No newline at end of file