commit 9117c3ac1ec19e4adab65f1ecaa786ee58e3750e Author: Cxx0822 <1556464090@qq.com> Date: Thu Oct 23 23:09:48 2025 +0800 feat:初始化工程 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..fe8e8f5 --- /dev/null +++ b/config/setting.py @@ -0,0 +1,16 @@ +from pydantic_settings import BaseSettings + + +class Settings(BaseSettings): + host: str + port: int + max_workers: int + timeout: int + + class Config: + env_file = ".env" # 指定.env文件路径 + case_sensitive = False # 忽略变量名大小写 + + +# 全局配置实例 +settings = Settings() diff --git a/endpoint/tcp_server.py b/endpoint/tcp_server.py new file mode 100644 index 0000000..62c1828 --- /dev/null +++ b/endpoint/tcp_server.py @@ -0,0 +1,161 @@ +import socket +import signal +import threading +from concurrent.futures import ThreadPoolExecutor + +from config.setting import settings +from model.index import UNKNOWN_MESSAGE +from processor.tcp_processor import TcpMessageProcessor +from config.logger_config import logger + + +class TCPServer: + def __init__(self): + self.host = settings.host + self.port = settings.port + + self.max_workers = settings.max_workers + self.timeout = settings.timeout + + self.server_socket = None + self.running = False + self.thread_pool = ThreadPoolExecutor(max_workers=self.max_workers) + self.message_processor = TcpMessageProcessor() + + # 已经连接的客户端 + self.clients = {} + self.client_lock = threading.Lock() + + self.server_thread = None + + signal.signal(signal.SIGTERM, self._handle_signal) + signal.signal(signal.SIGINT, self._handle_signal) + + def _handle_signal(self, signum, frame): + """处理终止信号,触发优雅关闭""" + logger.info(f"收到信号 {signum},准备关闭服务器...") + self.running = False + + def _start_loop(self): + """启动服务器""" + try: + self.server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + self.server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + self.server_socket.bind((self.host, self.port)) + self.server_socket.listen(5) + self.server_socket.settimeout(1.0) + self.running = True + + logger.info(f"TCP服务器启动,监听 {self.host}:{self.port} " + f"(最大线程: {self.max_workers}, 超时: {self.timeout}s)") + + while self.running: + try: + client_socket, client_address = self.server_socket.accept() + client_ip = client_address[0] + client_socket.settimeout(self.timeout) + logger.info(f"新连接: {client_address}") + + # 存储客户端 + with self.client_lock: + self.clients[client_ip] = client_socket + + # 提交到线程池处理 + self.thread_pool.submit(self.handle_client, client_socket, client_ip) + except socket.timeout: + continue + except Exception as e: + if self.running: + logger.error(f"接受连接失败: {str(e)}") + except Exception as e: + logger.error(f"TCP服务器启动失败: {str(e)}") + finally: + self.stop() + + def start(self): + self.server_thread = threading.Thread(target=self._start_loop, daemon=True) + self.server_thread.start() + + def stop(self): + if not self.running: + return + + self.running = False + logger.info("开始关闭服务器...") + + # 移除客户端 + with self.client_lock: + for client_socket in self.clients.values(): + try: + client_socket.close() + except Exception as e: + logger.warning(f"关闭客户端连接失败: {str(e)}") + self.clients.clear() + + # 关闭线程池 + self.thread_pool.shutdown(wait=True) + logger.info("所有客户端处理线程已结束") + + # 关闭连接 + if self.server_socket: + self.server_socket.close() + logger.info(f"服务器已关闭({self.host}:{self.port})") + + def handle_client(self, client_socket, client_ip): + """处理客户端连接""" + try: + while True: + data = client_socket.recv(1024) + if not data: + logger.info(f"客户端 {client_ip} 主动断开连接") + break + + message = data.decode('utf-8').strip() + logger.info(f"收到 {client_ip} 的消息: {message}") + + # 放到消息处理器里面处理 + response = self.message_processor.process(message) + + # 回复消息 + if response != UNKNOWN_MESSAGE: + client_socket.sendall(response.encode('utf-8')) + logger.info(f"回复 {client_ip}: {response}") + except socket.timeout: + logger.warning(f"客户端 {client_ip} 超时未活动") + except Exception as e: + logger.error(f"处理 {client_ip} 出错: {str(e)}") + finally: + # 异常情况下关闭连接 + with self.client_lock: + if client_ip in self.clients: + del self.clients[client_ip] + + try: + client_socket.close() + logger.info(f"客户端 {client_ip} 连接已关闭") + except Exception as e: + logger.warning(f"关闭 {client_ip} 连接失败: {str(e)}") + + def send_to_client(self, client_ip, message): + # 先获取客户端连接(加锁保护) + with self.client_lock: + client_socket = self.clients.get(client_ip) + if not client_socket: + logger.warning(f"客户端 {client_ip} 不存在或已断开连接") + return False + + # 发送消息 + try: + client_socket.sendall(message.encode('utf-8')) + logger.info(f"主动发送消息给 {client_ip}: {message}") + return True + except Exception as e: + logger.error(f"向 {client_ip} 发送消息失败: {str(e)}") + # 发送失败时移除无效连接 + with self.client_lock: + if client_ip in self.clients: + del self.clients[client_ip] + return False + + +tcp_server = TCPServer() diff --git a/main.py b/main.py new file mode 100644 index 0000000..bf7a5c1 --- /dev/null +++ b/main.py @@ -0,0 +1,19 @@ +from endpoint.tcp_server import tcp_server +from config.logger_config import logger +import time + +if __name__ == "__main__": + try: + tcp_server.start() + + try: + while True: + time.sleep(1) + except KeyboardInterrupt: + logger.info("收到终止信号,开始关闭程序...") + finally: + tcp_server.stop() + logger.info("程序已退出") + except Exception as e: + logger.critical(f"程序启动失败: {str(e)}", exc_info=True) + exit(1) 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..2302bd8 --- /dev/null +++ b/processor/tcp_processor.py @@ -0,0 +1,26 @@ +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 + \ No newline at end of file diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..91405a6 --- /dev/null +++ b/requirements.txt @@ -0,0 +1 @@ +loguru~=0.7.3 \ No newline at end of file