refactor: 移动pool到pipeline包

This commit is contained in:
RockChinQ
2024-02-23 17:46:22 +08:00
committed by Junyan Qin
parent cacd21bde7
commit 1f07a8a9e3
3 changed files with 3 additions and 3 deletions

View File

@@ -13,7 +13,7 @@ from ..config import manager as config_mgr
from ..audit.center import v2 as center_mgr
from ..command import cmdmgr
from ..plugin import manager as plugin_mgr
from . import pool
from ..pipeline import pool
from ..pipeline import controller, stagemgr
from ..utils import version as version_mgr, proxy as proxy_mgr

View File

@@ -9,7 +9,7 @@ from .bootutils import log
from .bootutils import config
from . import app
from . import pool
from ..pipeline import pool
from ..pipeline import controller
from ..pipeline import stagemgr
from ..audit import identifier

View File

@@ -1,57 +0,0 @@
from __future__ import annotations
import asyncio
import mirai
from . import entities
from ..platform import adapter as msadapter
class QueryPool:
query_id_counter: int = 0
pool_lock: asyncio.Lock
queries: list[entities.Query]
condition: asyncio.Condition
def __init__(self):
self.query_id_counter = 0
self.pool_lock = asyncio.Lock()
self.queries = []
self.condition = asyncio.Condition(self.pool_lock)
async def add_query(
self,
launcher_type: entities.LauncherTypes,
launcher_id: int,
sender_id: int,
message_event: mirai.MessageEvent,
message_chain: mirai.MessageChain,
adapter: msadapter.MessageSourceAdapter
) -> entities.Query:
async with self.condition:
query = entities.Query(
query_id=self.query_id_counter,
launcher_type=launcher_type,
launcher_id=launcher_id,
sender_id=sender_id,
message_event=message_event,
message_chain=message_chain,
resp_messages=[],
resp_message_chain=None,
adapter=adapter
)
self.queries.append(query)
self.query_id_counter += 1
self.condition.notify_all()
async def __aenter__(self):
await self.pool_lock.acquire()
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
self.pool_lock.release()