diff --git a/backend/app/task/__init__.py b/backend/app/task/__init__.py index f0416040..8a280cdc 100644 --- a/backend/app/task/__init__.py +++ b/backend/app/task/__init__.py @@ -2,7 +2,9 @@ # -*- coding: utf-8 -*- import sys -from pathlib import Path +from backend.core.path_conf import BASE_PATH + +from .actions import * # noqa: F403 # 导入项目根目录 -sys.path.append(str(Path(__file__).resolve().parent.parent.parent.parent)) +sys.path.append(str(BASE_PATH.parent)) diff --git a/backend/app/task/actions.py b/backend/app/task/actions.py new file mode 100644 index 00000000..f6b849cd --- /dev/null +++ b/backend/app/task/actions.py @@ -0,0 +1,13 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +from starlette.concurrency import run_in_threadpool + +from backend.app.task.celery import celery_app +from backend.common.socketio.server import sio + + +@sio.event +async def task_worker_status(sid, data): + """任务 Worker 状态事件""" + worker = await run_in_threadpool(celery_app.control.ping) + await sio.emit('task_worker_status', worker, sid) diff --git a/backend/app/task/celery.py b/backend/app/task/celery.py index d543317e..24076a11 100644 --- a/backend/app/task/celery.py +++ b/backend/app/task/celery.py @@ -35,7 +35,7 @@ def init_celery() -> celery.Celery: if settings.CELERY_BROKER == 'redis' else f'amqp://{settings.CELERY_RABBITMQ_USERNAME}:{settings.CELERY_RABBITMQ_PASSWORD}@{settings.CELERY_RABBITMQ_HOST}:{settings.CELERY_RABBITMQ_PORT}', broker_connection_retry_on_startup=True, - backend=f'db+{settings.DATABASE_TYPE + "+pymysql" if settings.DATABASE_TYPE == "mysql" else settings.DATABASE_TYPE}' # noqa: E501 + backend=f'db+{settings.DATABASE_TYPE}+{"pymysql" if settings.DATABASE_TYPE == "mysql" else "psycopg"}' f'://{settings.DATABASE_USER}:{settings.DATABASE_PASSWORD}@{settings.DATABASE_HOST}:{settings.DATABASE_PORT}/{settings.DATABASE_SCHEMA}', database_engine_options={'echo': settings.DATABASE_ECHO}, database_table_names={ diff --git a/backend/common/socketio/__init__.py b/backend/common/socketio/__init__.py index 56fafa58..f6eb45c3 100644 --- a/backend/common/socketio/__init__.py +++ b/backend/common/socketio/__init__.py @@ -1,2 +1,3 @@ #!/usr/bin/env python3 # -*- coding: utf-8 -*- +from .actions import * # noqa: F403 diff --git a/backend/common/socketio/server.py b/backend/common/socketio/server.py index a2d0a331..6edda817 100644 --- a/backend/common/socketio/server.py +++ b/backend/common/socketio/server.py @@ -9,17 +9,8 @@ from backend.database.redis import redis_client # 创建 Socket.IO 服务器实例 sio = socketio.AsyncServer( - # 集成 Celery 实现消息订阅 client_manager=socketio.AsyncRedisManager( - f'redis://:{settings.REDIS_PASSWORD}@{settings.REDIS_HOST}:' - f'{settings.REDIS_PORT}/{settings.CELERY_BROKER_REDIS_DATABASE}' - ) - if settings.CELERY_BROKER == 'redis' - else socketio.AsyncAioPikaManager( - ( - f'amqp://{settings.CELERY_RABBITMQ_USERNAME}:{settings.CELERY_RABBITMQ_PASSWORD}@' - f'{settings.CELERY_RABBITMQ_HOST}:{settings.CELERY_RABBITMQ_PORT}' - ) + f'redis://:{settings.REDIS_PASSWORD}@{settings.REDIS_HOST}:{settings.REDIS_PORT}/{settings.REDIS_DATABASE}' ), async_mode='asgi', cors_allowed_origins=settings.CORS_ALLOWED_ORIGINS, @@ -30,7 +21,7 @@ sio = socketio.AsyncServer( @sio.event async def connect(sid, environ, auth): - """处理 WebSocket 连接事件""" + """Socket 连接事件""" if not auth: log.error('WebSocket 连接失败:无授权') return False @@ -57,6 +48,6 @@ async def connect(sid, environ, auth): @sio.event -async def disconnect(sid: str) -> None: - """处理 WebSocket 断开连接事件""" +async def disconnect(sid) -> None: + """Socket 断开连接事件""" await redis_client.spop(settings.TOKEN_ONLINE_REDIS_PREFIX) diff --git a/pyproject.toml b/pyproject.toml index 6d098fce..eda19eb8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -40,6 +40,7 @@ dependencies = [ "msgspec>=0.19.0", "path>=17.0.0", "psutil>=7.0.0", + "psycopg>=3.2.9", "pwdlib>=0.2.1", "pydantic>=2.11.0", "pydantic-settings>=2.10.0", diff --git a/requirements.txt b/requirements.txt index 6da6c82b..75d75a51 100644 --- a/requirements.txt +++ b/requirements.txt @@ -185,6 +185,8 @@ prompt-toolkit==3.0.51 # via click-repl psutil==7.0.0 # via fastapi-best-architecture +psycopg==3.2.9 + # via fastapi-best-architecture pwdlib==0.2.1 # via fastapi-best-architecture pyasn1==0.6.1 @@ -296,6 +298,7 @@ typing-extensions==4.14.1 # exceptiongroup # fastapi # fastapi-pagination + # psycopg # pydantic # pydantic-core # rich @@ -310,7 +313,9 @@ typing-inspection==0.4.1 # pydantic # pydantic-settings tzdata==2025.2 - # via kombu + # via + # kombu + # psycopg ua-parser==1.0.1 # via user-agents ua-parser-builtins==0.18.0.post1 diff --git a/uv.lock b/uv.lock index 9cf78aea..3eca8d5b 100644 --- a/uv.lock +++ b/uv.lock @@ -661,6 +661,7 @@ dependencies = [ { name = "msgspec" }, { name = "path" }, { name = "psutil" }, + { name = "psycopg" }, { name = "pwdlib" }, { name = "pydantic" }, { name = "pydantic-settings" }, @@ -716,6 +717,7 @@ requires-dist = [ { name = "msgspec", specifier = ">=0.19.0" }, { name = "path", specifier = ">=17.0.0" }, { name = "psutil", specifier = ">=7.0.0" }, + { name = "psycopg", specifier = ">=3.2.9" }, { name = "pwdlib", specifier = ">=0.2.1" }, { name = "pydantic", specifier = ">=2.11.0" }, { name = "pydantic-settings", specifier = ">=2.10.0" }, @@ -1771,6 +1773,19 @@ wheels = [ { url = "https://mirrors.aliyun.com/pypi/packages/50/1b/6921afe68c74868b4c9fa424dad3be35b095e16687989ebbb50ce4fceb7c/psutil-7.0.0-cp37-abi3-win_amd64.whl", hash = "sha256:4cf3d4eb1aa9b348dec30105c55cd9b7d4629285735a102beb4441e38db90553" }, ] +[[package]] +name = "psycopg" +version = "3.2.9" +source = { registry = "https://mirrors.aliyun.com/pypi/simple" } +dependencies = [ + { name = "typing-extensions", marker = "python_full_version < '3.13'" }, + { name = "tzdata", marker = "sys_platform == 'win32'" }, +] +sdist = { url = "https://mirrors.aliyun.com/pypi/packages/27/4a/93a6ab570a8d1a4ad171a1f4256e205ce48d828781312c0bbaff36380ecb/psycopg-3.2.9.tar.gz", hash = "sha256:2fbb46fcd17bc81f993f28c47f1ebea38d66ae97cc2dbc3cad73b37cefbff700" } +wheels = [ + { url = "https://mirrors.aliyun.com/pypi/packages/44/b0/a73c195a56eb6b92e937a5ca58521a5c3346fb233345adc80fd3e2f542e2/psycopg-3.2.9-py3-none-any.whl", hash = "sha256:01a8dadccdaac2123c916208c96e06631641c0566b22005493f09663c7a8d3b6" }, +] + [[package]] name = "pwdlib" version = "0.2.1"