diff --git a/management/openjiuwen_runtime/management/session/GLOSSARY.md b/management/openjiuwen_runtime/management/session/GLOSSARY.md new file mode 100644 index 0000000..55c8257 --- /dev/null +++ b/management/openjiuwen_runtime/management/session/GLOSSARY.md @@ -0,0 +1,178 @@ +# 术语表(Session SDK) + +本文件是 Session SDK 各名词的**权威定义**,消除"session"一词的多义混淆。 + +--- + +## 0. 核心澄清:三种 "session" + +历史上 "session" 在本项目中指代三种**完全不同**的东西,是所有命名混淆的根源: + +| 含义 | 实际是什么 | 标识 | 归属 | +|------|-----------|------|------| +| **Page session** | 用户的一次浏览器会话(对话) | `sess_*`(如 `sess_19f02e1...`) | **gateway** 侧,SDK 不感知 | +| **ServiceScope** | (group, bot) 资源分配域——同 group 同 bot 的所有用户共享 | `service_id` = md5(group+bot) | **SDK** 核心概念 | +| **会话状态** | 对话上下文(记忆/历史) | AgentServer 进程内或外置存储 | **AgentServer**,SDK 不持有 | + +> **关键**:SDK 层**没有"用户会话"概念**。SDK 管理的是 ServiceScope(Pod 资源域),不是用户对话。对话连续性由 AgentServer 外置状态保证。 + +--- + +## 1. ServiceScope 体系(SDK 核心) + +### ServiceScope +以 **(group_id, bot_id)** 为键的**资源分配与路由域**。同一 group 内所有用户与同一 bot 的交互共享一个 ServiceScope。用于:查找 Pod 模版(`service_config_template`)、分配/管理一组 Pod、在这一组 Pod 间路由请求。 + +- **不是**:某个用户的某次对话 +- **是**:一个 bot 在一个 group 里的 Pod 资源单元 + +### `service_id` +ServiceScope 的唯一标识 = `md5(group_id + bot_id)`。在整个 SDK 中用作 registry 键、计时器键、路由标识。 + +> **历史**:曾叫 `session_id`,因与 page session 同名造成混淆,已改名。gateway 的 `_SessionRequest` 内部早已用 `self._service_id` 存储。 + +### ServiceScopeHandler +一个 ServiceScope 的运行时处理器(`service_scope_handler.py`)。持有多个 Pod endpoint,做: +- 并发限流(semaphore,上限 = `scope_concurrency`) +- 容量感知路由(用户亲和 → 最少负载) +- 弹性扩缩 endpoint(multi-pod) + +> **历史**:曾叫 `SessionHandler`。 + +### ServiceScopeRegistry +`service_id → ServiceScopeHandler` 的归属表(在 `SessionRuntimeManager` 内)。管理 ServiceScopeHandler 的创建、查找、移除。 + +> **历史**:曾叫 `SessionRegistry`。 + +--- + +## 2. 三层架构 + +``` +Access (入口) + ↓ +SessionRuntimeManager (编排层) ← 1 per 进程 + ├── ServiceScopeRegistry ← service_id → ServiceScopeHandler + ├── handle_user_request ← 消息入口、TTL、pending 过期 + └── 借/还 Pod ←→ ServiceManager + ↓ +ServiceScopeHandler (域层) ← 1 per (group,bot) + ├── semaphore(scope_concurrency) + ├── _pick_endpoint (用户亲和→最少负载) + └── List[ISendEndpoint] + ↓ send_message +ServiceHandler (Pod 层) ← 1 per Pod + ├── quota (_session_reserved) + ├── inflight / WebSocket channel + └── deploy / delete / evict_session +``` + +### ServiceManager +**Pod 池管理者**(`service_manager.py`)。管理 Pod 生命周期:deploy/delete、autoscale(维护 min_idle 热备)、idle/reclaim(老化回收)、监控失效 Pod。向编排层暴露 `pick_or_create_pod` / `find_service_handler` / `reconsider_idle_transition`。 + +### SessionRuntimeManager +**Session 编排层**(`session_runtime_manager.py`)。职责:消息入口(handle_user_request)、ServiceScope TTL/pending 过期、向 ServiceManager 借还 Pod。内含 `ServiceScopeRegistry`。 + +> **命名说明**:叫 "SessionRuntime" 而非 "ServiceScopeRuntime" 是因为它编排的是**运行时生命周期**(TTL/路由/pending),与 gateway 侧已有的 `SessionManager`(任务队列)同名但不同层。 + +### ServiceHandler +**单个 Pod 的管理者**(`service_handler.py`)。同时实现 `IServiceHandler`(Pod 池视角:quota/lifecycle)和 `ISendEndpoint`(发送视角:send_message/inflight/endpoint_id)。每个 ServiceHandler 对应一个 AgentServer Pod,持有一条 WebSocket 连接。 + +### ISendEndpoint +**发送端点接口**(`interfaces.py`)。ServiceScopeHandler 通过此接口向 Pod 发请求,不绑定 ServiceHandler 类型。属性:`endpoint_id`、`inflight`、`send_message`。由 ServiceHandler 实现。 + +--- + +## 3. 并发与 TTL(DB 列名 ↔ 概念名对照) + +> **重要**:以下四个值同时是 Python 字段名和 `service_config_template` **DB 列名**。为兼容现网 DB,DB 列名保留旧称,概念名对称改名仅在文档/讨论中使用。 + +| 概念名(本文档) | DB 列名(代码实际字段) | 含义 | 默认 | +|----------------|---------------------|------|------| +| **pod_concurrency** | `service_concurrency` | 单 Pod 最大并发容量 | 10 | +| **scope_concurrency** | `session_concurrency` | 单 ServiceScope 总并发预算(跨所有 Pod) | 100 | +| **pod_ttl** | `service_ttl` | Pod 无业务后转 idle 池的等待秒数 | 180 | +| **scope_ttl** | `session_ttl` | ServiceScope 保活窗口(idle 后多久回收 Pod) | 60 | + +**代码中统一用 DB 列名**(`service_concurrency` / `session_concurrency` / `service_ttl` / `session_ttl`),本文档用概念名便于理解层级关系。 + +### min_idle / max_services +- `min_idle`(`min_idle_services`):全局 idle 池至少保持的**热备** Pod 数(autoscale 自动补位)。新请求到达时唤醒热备,避免冷启动。 +- `max_services`:Pod 总数上限。 + +### reserve_per_pod +每个业务 Pod 从 ServiceScope 抽取的容量 = `min(scope_concurrency, pod_concurrency)`。multi-pod 弹性扩缩的计量单位。 + +### max_business +一个 ServiceScope 的最大业务 Pod 数 = `ceil(scope_concurrency / reserve_per_pod)`。 + +--- + +## 4. 请求与消息类型 + +### IRequest +**Gateway 入口请求接口**(`interfaces.py`)。携带 page session 信息:`request_id`、`chat_id`、`bot_id`、`user_id`、`session_id`(**page session**)。 + +### ISessionRequest +**SDK 内部的 scope 请求接口**(`interfaces.py`)。由策略从 `IRequest` 聚合而成,携带:`service_id`(scope 键)、`session_concurrency`、`session_ttl`、`request_id`、`raw_msg`(原始 IRequest)、`service_template`。 + +> **关键区分**:`IRequest.session_id`(page session)≠ `ISessionRequest.service_id`(scope 键)。 + +### ScopeRequestWrapper +封装 `ISessionRequest` + `response_queue` + `cancel future` 的传输包装,贯穿 SDK 管道。 + +> **历史**:曾叫 `SessionRequestWrapper`。 + +### SessionConfig +Gateway 传入的 scope 配置(`models.py`):`concurrency`(→ scope_concurrency)和 `ttl`(→ scope_ttl)。`Access.init` 的参数,属 gateway↔SDK 契约。 + +--- + +## 5. Pod 生命周期 + +| 术语 | 含义 | +|------|------| +| **in_use 池** | 正在服务请求的 Pod(有 session 预留额度) | +| **idle 池** | 空闲热备 Pod(无 session,可被唤醒) | +| **deploy** | 创建新 Pod(K8s 创建 Pod + WebSocket 建链) | +| **pick_or_create_pod** | 编排层借 Pod:找有容量的 in_use → 唤醒 idle → 新 deploy | +| **try_reserve_session_quota** | 为 ServiceScope 在某 Pod 上预留额度(隔离) | +| **evict_session** | 释放某 Pod 上该 ServiceScope 的额度 + 取消 pod-local 在途请求 | +| **reconsider_idle_transition** | evict 后推动 Pod 重新评估 in_use→idle(老化推进) | +| **autoscale** | 周期维护 min_idle 热备 + 回收多余 idle | +| **scope TTL** | ServiceScope idle 后回收 Pod 的计时器 | + +--- + +## 6. Gateway 侧概念(不属于 SDK,但易混淆) + +| 术语 | 含义 | +|------|------| +| **page session** (`sess_*`) | 用户的一次浏览器会话,gateway 的 SessionMap 管理 | +| **SessionMapScope** | 会话映射策略(`per_chat_bot` / `per_chat_bot_user`) | +| **service_config_template** | DB 表,按 (group, bot) 配置 Pod 模版(镜像、资源、并发、TTL) | +| **AgentServer** | AI Agent 运行时进程(每个 Pod 一个),执行 LLM/Skill/Tool | +| **Gateway** | 消息路由层,连接 IM 平台与 AgentServer | + +--- + +## 7. 名称变更历史 + +| 旧名 | 新名 | 改名原因 | +|------|------|---------| +| `SessionHandler` | `ServiceScopeHandler` | 它管的是 (group,bot) 资源域,不是用户会话 | +| `SessionRegistry` | `ServiceScopeRegistry` | 同上 | +| `ISessionHandler` | `IServiceScopeHandler` | 接口对齐 | +| `SessionRequestWrapper` | `ScopeRequestWrapper` | 同上 | +| `session_id`(SDK 内) | `service_id` | 与 page session 同名混淆;值本就是 service_id | +| `_arm_session_timer` | `_arm_scope_timer` | 方法名对齐 | +| `_on_session_expired` | `_on_scope_expired` | 同上 | + +**未改名(保留)**: +- `ISessionRequest` / `SessionConfig` / `SessionRequest`:gateway 面向契约,改名需跨模块同步。 +- `service_concurrency` / `session_concurrency` / `service_ttl` / `session_ttl`:`service_config_template` DB 列名 + 配置协议,改名需 DB 迁移。 +- `IRequest.session_id`:page session,gateway 概念,不属 SDK。 + +--- + +*本文档随代码演进同步更新。如有疑问,以代码实际命名为准。* diff --git a/management/openjiuwen_runtime/management/session/__init__.py b/management/openjiuwen_runtime/management/session/__init__.py index 76e44c7..4d40d7a 100644 --- a/management/openjiuwen_runtime/management/session/__init__.py +++ b/management/openjiuwen_runtime/management/session/__init__.py @@ -12,10 +12,10 @@ IServiceManager, IServiceMessageChannel, OnRequestCompleteCallback, - ISessionHandler, + IServiceScopeHandler, ISessionRequest, ISessionStrategy, - SessionRequestWrapper, + ScopeRequestWrapper, ) from .models import AccessConfig, MessagePriority, MessageType, SessionConfig from .runtime import IDeployController, NoOpDeployController @@ -36,7 +36,7 @@ "IServiceInstanceFactory", "IServiceManager", "IServiceMessageChannel", - "ISessionHandler", + "IServiceScopeHandler", "ISessionRequest", "ISessionStrategy", "MessagePriority", @@ -48,7 +48,7 @@ "ServiceManager", "SessionConfig", "SessionRequest", - "SessionRequestWrapper", + "ScopeRequestWrapper", "serialize_request_payload", "Timer", "WSServiceMessageChannel", diff --git a/management/openjiuwen_runtime/management/session/access.py b/management/openjiuwen_runtime/management/session/access.py index c8c85c3..c846adb 100644 --- a/management/openjiuwen_runtime/management/session/access.py +++ b/management/openjiuwen_runtime/management/session/access.py @@ -15,7 +15,7 @@ IResponseParser, ISessionStrategy, IServiceManager, - SessionRequestWrapper, ISessionRequest, + ScopeRequestWrapper, ISessionRequest, ) from .models import AccessConfig, SessionConfig @@ -278,7 +278,7 @@ async def send_message(self, msg: IRequest | ISessionRequest) -> AsyncIterator[A rid = session_request.request_id logger.debug( "Access receive session: session_id=%s session_conc=%s session_ttl=%s request_id=%s", - session_request.session_id, + session_request.service_id, session_request.session_concurrency, session_request.session_ttl, rid, @@ -303,7 +303,7 @@ async def send_message(self, msg: IRequest | ISessionRequest) -> AsyncIterator[A session_request = self._strategy.handle_session(msg) logger.debug( "Access 策略生成 session: session_id=%s session_conc=%s session_ttl=%s", - session_request.session_id, + session_request.service_id, session_request.session_concurrency, session_request.session_ttl, ) @@ -311,7 +311,7 @@ async def send_message(self, msg: IRequest | ISessionRequest) -> AsyncIterator[A # 3) 每个入口请求独占一条响应队列 + cancel,用于多路复用/取消 response_queue: asyncio.Queue[Any] = asyncio.Queue() cancel: asyncio.Future = asyncio.get_running_loop().create_future() - wrapper = SessionRequestWrapper(session_request, response_queue, cancel) + wrapper = ScopeRequestWrapper(session_request, response_queue, cancel) # 4) 入用户队列,由 ServiceManager 异步消费并路由到具体服务实例 await self._service_manager.handle_message(wrapper) diff --git a/management/openjiuwen_runtime/management/session/interfaces.py b/management/openjiuwen_runtime/management/session/interfaces.py index 61ae7a3..ccae02a 100644 --- a/management/openjiuwen_runtime/management/session/interfaces.py +++ b/management/openjiuwen_runtime/management/session/interfaces.py @@ -78,7 +78,8 @@ def session_id(self) -> Optional[str]: class ISessionRequest(PriorityMessage): @property @abstractmethod - def session_id(self) -> str: + def service_id(self) -> str: + """ServiceScope 标识(group+bot 派生);注意与 :class:`IRequest.session_id`(page session)区分。""" pass @property @@ -114,7 +115,7 @@ def __init__(self, code: int, message: str, data: dict = None) -> None: self.data = data -class SessionRequestWrapper: +class ScopeRequestWrapper: def __init__( self, request: ISessionRequest, @@ -231,13 +232,13 @@ class IServiceMessageChannel(Protocol): 约定: * ``send`` 中 **上行** 发送一帧业务负载(通常 JSON 序列化自 ``ISessionRequest.raw_msg``); - * **下行** 由实现类在独接收循环里按 ``IResponseParser.request_id`` 分片写入对应 ``SessionRequestWrapper.response_queue``; + * **下行** 由实现类在独接收循环里按 ``IResponseParser.request_id`` 分片写入对应 ``ScopeRequestWrapper.response_queue``; * 当某 ``request_id`` 的响应用 ``IResponseParser.is_completed`` 判定结束时,**必须** ``await on_request_complete(request_id)`` 归还本实例并发。 可选实现(鸭子类型): ``bind_handler(handler, parser)``、``on_pod_ready(service_id, pod_info)``、``close()``. """ - async def send(self, service_id: str, wrapper: SessionRequestWrapper, *, response_parser: IResponseParser, + async def send(self, service_id: str, wrapper: ScopeRequestWrapper, *, response_parser: IResponseParser, on_request_complete: OnRequestCompleteCallback, ) -> None: pass @@ -245,15 +246,15 @@ async def send(self, service_id: str, wrapper: SessionRequestWrapper, *, respons @runtime_checkable class ISendEndpoint(Protocol): - """SessionHandler 通过此接口向某个 Pod 发送请求,不绑定具体 ServiceHandler 类型。 + """ServiceScopeHandler 通过此接口向某个 Pod 发送请求,不绑定具体 ServiceHandler 类型。 约定: * 由 ``ServiceHandler`` 实现(结构子类型,不要让具体子类再继承本 Protocol); * ``send_message`` 完成 inflight 计数、request_id → wrapper 映射(供下行分片)、 实际的 ``IServiceMessageChannel.send`` 以及完成回调(inflight 归零通知 idle_pool_hook); - * ``inflight`` / ``endpoint_id`` 用于 SessionHandler 的最少负载路由与用户亲和记录。 + * ``inflight`` / ``endpoint_id`` 用于 ServiceScopeHandler 的最少负载路由与用户亲和记录。 - 单 Pod 场景:SessionHandler 持有 1 个端点;多 Pod 场景:持有 N 个端点。 + 单 Pod 场景:ServiceScopeHandler 持有 1 个端点;多 Pod 场景:持有 N 个端点。 """ @property @@ -266,8 +267,8 @@ def inflight(self) -> int: """当前端点在途请求数,用于最少负载路由。""" ... - async def send_message(self, wrapper: SessionRequestWrapper) -> None: - """发送请求并追踪完成。SessionHandler 在获取信号量后调用此方法。""" + async def send_message(self, wrapper: ScopeRequestWrapper) -> None: + """发送请求并追踪完成。ServiceScopeHandler 在获取信号量后调用此方法。""" ... @@ -293,7 +294,7 @@ async def stop(self) -> None: pass @abstractmethod - async def handle_message(self, msg: "SessionRequestWrapper") -> None: + async def handle_message(self, msg: "ScopeRequestWrapper") -> None: pass @abstractmethod @@ -347,7 +348,7 @@ def inflight_requests(self) -> int: def active_session_count(self) -> int: """当前实例上仍预留额度的 session 数(基于 ``_session_reserved``)。 - 解耦后语义为 quota 计数,用于 idle/TTL 判定,而非 SessionHandler 实例数。 + 解耦后语义为 quota 计数,用于 idle/TTL 判定,而非 ServiceScopeHandler 实例数。 """ pass @@ -360,7 +361,7 @@ async def evict_session(self, session_id: str) -> int: """释放该 session 在本实例的预留额度,并取消本实例上该 session 的在途请求(pod-local)。 解耦后取代原 ``remove_session``:只处理 pod-local 状态(quota + 取消), - 不销毁 SessionHandler 实例(后者由 SessionRegistry 负责)。 + 不销毁 ServiceScopeHandler 实例(后者由 ServiceScopeRegistry 负责)。 """ pass @@ -385,7 +386,7 @@ def set_idle_pool_transition_hook( return -class ISessionHandler(ABC): +class IServiceScopeHandler(ABC): @abstractmethod - async def handle_message(self, msg: "SessionRequestWrapper") -> None: + async def handle_message(self, msg: "ScopeRequestWrapper") -> None: pass diff --git a/management/openjiuwen_runtime/management/session/service_handler.py b/management/openjiuwen_runtime/management/session/service_handler.py index f030da5..3d88b64 100644 --- a/management/openjiuwen_runtime/management/session/service_handler.py +++ b/management/openjiuwen_runtime/management/session/service_handler.py @@ -3,7 +3,7 @@ """单实例(Pod):服务级并发按 session 预留额度、消息级在途计数、发送端点与下行通道。 -解耦后 ServiceHandler 不再创建/持有 SessionHandler 实例(归属权在 SessionRegistry), +解耦后 ServiceHandler 不再创建/持有 ServiceScopeHandler 实例(归属权在 ServiceScopeRegistry), 仅负责:quota 预留、inflight 计数、``ISendEndpoint.send_message``、下行分片、Pod 生命周期。 """ @@ -21,7 +21,7 @@ ISendEndpoint, IServiceHandler, IServiceMessageChannel, - SessionRequestWrapper, + ScopeRequestWrapper, ) from .runtime import IDeployController, NoOpDeployController @@ -34,10 +34,10 @@ class ServiceHandler(IServiceHandler, ISendEndpoint): """单个 Pod 的管理者 + 发送端点。 * 实现 :class:`ISendEndpoint`(``send_message`` / ``endpoint_id`` / ``inflight``), - 供 SessionHandler 跨多端点路由调用; + 供 ServiceScopeHandler 跨多端点路由调用; * 实现 :class:`IServiceHandler`(quota / 生命周期 / idle 钩子),供 ServiceManager 管理 Pod 池; - * **不再** 创建或持有 SessionHandler 实例,**不再** 实现 ``handle_message`` - (消息入口由 SessionRuntimeManager → SessionHandler 承担)。 + * **不再** 创建或持有 ServiceScopeHandler 实例,**不再** 实现 ``handle_message`` + (消息入口由 SessionRuntimeManager → ServiceScopeHandler 承担)。 """ def __init__( @@ -63,7 +63,7 @@ def __init__( self._channel = message_channel self._parser = response_parser self._deploy: IDeployController = deploy_controller or NoOpDeployController() - self._by_request: Dict[str, SessionRequestWrapper] = {} + self._by_request: Dict[str, ScopeRequestWrapper] = {} self._pod_info: Any = None self._closed = False # ServiceManager 注入: 每次 inflight 归零后被回调一次, Manager 据此推动到期 session 清理与 service_ttl 计时 @@ -103,10 +103,10 @@ def inflight(self) -> int: """ISendEndpoint.inflight:当前端点在途请求数(同 inflight_requests)。""" return self._inflight - async def send_message(self, wrapper: SessionRequestWrapper) -> None: + async def send_message(self, wrapper: ScopeRequestWrapper) -> None: """ISendEndpoint.send_message:登记在途请求并发通道(不占服务额度)。 - 由 SessionHandler 在取得 session 信号量后调用。此处 + 由 ServiceScopeHandler 在取得 session 信号量后调用。此处 ``await self._channel.send(...)`` 为 ``WSServiceMessageChannel`` 等 ``IServiceMessageChannel`` 实现的业务上行入口。 """ @@ -156,7 +156,7 @@ async def _complete(r: Optional[str]) -> None: raise # 兼容别名:原 invoke_channel 改名为 send_message,保留旧名一个版本便于过渡 - async def invoke_channel(self, wrapper: SessionRequestWrapper) -> None: + async def invoke_channel(self, wrapper: ScopeRequestWrapper) -> None: """deprecated 别名,等价于 :meth:`send_message`。""" await self.send_message(wrapper) @@ -217,7 +217,7 @@ async def evict_session(self, session_id: str) -> int: """释放该 session 在本实例的预留额度,并取消本实例上该 session 的在途请求。 解耦后取代原 ``remove_session``:只处理 pod-local 状态(quota + 取消), - 不销毁 SessionHandler 实例(后者由 SessionRegistry 负责)。 + 不销毁 ServiceScopeHandler 实例(后者由 ServiceScopeRegistry 负责)。 Returns: 1 表示该 session 曾在本实例预留额度(已驱逐);0 表示未找到(无操作)。 @@ -229,7 +229,7 @@ async def evict_session(self, session_id: str) -> int: w = self._by_request.get(rid) if w is None: continue - if w.session_request.session_id == session_id and not w.cancel.done(): + if w.session_request.service_id == session_id and not w.cancel.done(): w.cancel.set_result(None) cancelled += 1 logger.info( diff --git a/management/openjiuwen_runtime/management/session/service_manager.py b/management/openjiuwen_runtime/management/session/service_manager.py index bf94a07..46da151 100644 --- a/management/openjiuwen_runtime/management/session/service_manager.py +++ b/management/openjiuwen_runtime/management/session/service_manager.py @@ -25,7 +25,7 @@ IServiceManager, ITimer, RawMessage, - SessionRequestWrapper, + ScopeRequestWrapper, ) from .k8s_service_handler import K8sServiceHandler, POD_LABEL_SELECTOR from .models import MessageType @@ -335,11 +335,11 @@ async def stop(self) -> None: for h in all_handlers: self._deleting_services.discard(h.id) - async def handle_message(self, msg: SessionRequestWrapper) -> None: + async def handle_message(self, msg: ScopeRequestWrapper) -> None: if self._deprecated: logger.warning( "ServiceManager 已标记为待老化,但仍收到新消息: session_id=%s request_id=%s", - msg.session_request.session_id, + msg.session_request.service_id, msg.session_request.request_id, ) sreq = msg.session_request @@ -350,7 +350,7 @@ async def handle_message(self, msg: SessionRequestWrapper) -> None: await self._q.put_user(raw) logger.debug( "ServiceManager 用户消息已入队: session_id=%s request_id=%s user_q~=%s", - sreq.session_id, + sreq.service_id, sreq.request_id, self._q.user_qsize(), ) @@ -699,7 +699,7 @@ async def _cleanup_displaced_handler(self, old: IServiceHandler) -> None: """清理被同 id 挤出的旧 handler,并 delete 其底层资源。""" await self._cancel_in_use_to_idle_timer(old.id) await self._cancel_excess_idle_timer(old.id) - # session 侧:从引用该 Pod 的 SessionHandler 摘除 endpoint、清 pending + # session 侧:从引用该 Pod 的 ServiceScopeHandler 摘除 endpoint、清 pending if self._session_runtime is not None: try: await self._session_runtime.on_pod_removed(old.id) @@ -709,7 +709,7 @@ async def _cleanup_displaced_handler(self, old: IServiceHandler) -> None: old.id, e, exc_info=True, ) for session_id in list(old.open_session_ids()): - await self._timer.cancel_timer(f"sess:{session_id}") + await self._timer.cancel_timer(f"scope:{session_id}") try: # 不用 _safe_delete_handler:勿把 service_id 放进 _deleting_services, # 否则会与新实例的后续清理路径互相干扰。 @@ -1104,7 +1104,7 @@ def _pick_existing_locked( ) -> Optional[IServiceHandler]: """在池中按容量选实例。调用方须持有 self._lock。不 deploy。 - 注:session 亲和由 SessionHandler / SessionRuntimeManager 处理,此处不做亲和查找。 + 注:session 亲和由 ServiceScopeHandler / SessionRuntimeManager 处理,此处不做亲和查找。 """ in_use_pool = self._in_use.get(template_id, {}) for h in in_use_pool.values(): @@ -1158,7 +1158,7 @@ def _try_begin_deploy_locked( logger.warning( "预检拦截(会话并发超过单实例总并发): session_id=%s template_id=%s " "session_concurrency=%s service_concurrency=%s, 拒绝扩容", - sreq.session_id, template_id, need, service_concurrency_for_tpl, + sreq.service_id, template_id, need, service_concurrency_for_tpl, ) return None @@ -1195,13 +1195,13 @@ async def _pick_or_create(self, sreq) -> Optional[IServiceHandler]: # noqa: ANN 避免惊群重复占满 max_services。 缩容 delete 未完成前仍占用 max 名额(reclaim_occupancy)。 - 注:session 亲和(同 session 复用同一 Pod)由 SessionHandler(持有 endpoints) + 注:session 亲和(同 session 复用同一 Pod)由 ServiceScopeHandler(持有 endpoints) 在 SessionRuntimeManager 层处理,本方法不再做亲和查找 / 额度预留。 注意:本方法自行获取/释放 ``self._lock``,调用方勿再持锁调用。 """ need = max(1, int(sreq.session_concurrency)) - session_id = sreq.session_id + session_id = sreq.service_id template_id: Optional[str] = None if sreq.service_template: template_id = sreq.service_template.get("template_id") @@ -1590,8 +1590,8 @@ async def _cleanup_dead_pods(self, dead_pods: Dict[str, str]) -> None: self._excess_idle_timer_armed.discard(service_id) # session 侧清理(委托 session 编排层):从引用该 Pod 的 - # SessionHandler 摘除 endpoint、清 _pending_expired 记录。 - # 注:ServiceManager 不再持有 _service_router / _pending_expired_sessions + # ServiceScopeHandler 摘除 endpoint、清 _pending_scope_expiry 记录。 + # 注:ServiceManager 不再持有 _service_router / _pending_scope_expiry_sessions # (已迁至 SessionRuntimeManager)。 if self._session_runtime is not None: try: @@ -1608,7 +1608,7 @@ async def _cleanup_dead_pods(self, dead_pods: Dict[str, str]) -> None: for service_id, h, pod_name, reason, session_ids in handlers_to_delete: try: for session_id in session_ids: - await self._timer.cancel_timer(f"sess:{session_id}") + await self._timer.cancel_timer(f"scope:{session_id}") await self._cancel_in_use_to_idle_timer(service_id) await self._cancel_excess_idle_timer(service_id) try: diff --git a/management/openjiuwen_runtime/management/session/session_handler.py b/management/openjiuwen_runtime/management/session/service_scope_handler.py similarity index 73% rename from management/openjiuwen_runtime/management/session/session_handler.py rename to management/openjiuwen_runtime/management/session/service_scope_handler.py index 33570cc..4b4f5d8 100644 --- a/management/openjiuwen_runtime/management/session/session_handler.py +++ b/management/openjiuwen_runtime/management/session/service_scope_handler.py @@ -3,7 +3,7 @@ """单 session 限流 + 多发送端点路由(与 ServiceHandler 类型解耦)。 -SessionHandler 不再引用 ServiceHandler 类型,仅依赖 :class:`ISendEndpoint` 接口。 +ServiceScopeHandler 不再引用 ServiceHandler 类型,仅依赖 :class:`ISendEndpoint` 接口。 单 Pod 场景持有 1 个端点;多 Pod 场景持有 N 个端点,按「用户亲和 → 最少负载」路由。 """ @@ -14,13 +14,13 @@ from openjiuwen_runtime.foundation.log import get_logger -from .interfaces import IResponseParser, ISendEndpoint, ISessionHandler, SessionRequestWrapper +from .interfaces import IResponseParser, ISendEndpoint, IServiceScopeHandler, ScopeRequestWrapper from .router import SessionRouter logger = get_logger(__name__) -class SessionHandler(ISessionHandler): +class ServiceScopeHandler(IServiceScopeHandler): """同 session 内限流 + 多端点路由。 最多 ``session_concurrency`` 路并行,路由策略:用户亲和 → 最少负载。 @@ -29,12 +29,12 @@ class SessionHandler(ISessionHandler): def __init__( self, - session_id: str, + service_id: str, max_parallel: int, endpoints: List[ISendEndpoint], session_router: SessionRouter, ) -> None: - self._session_id = session_id + self._service_id = service_id m = max(1, int(max_parallel)) self._sem = asyncio.BoundedSemaphore(m) self._endpoints: List[ISendEndpoint] = list(endpoints) @@ -43,13 +43,13 @@ def __init__( # 用户亲和: user_id -> endpoint_id self._user_affinity: Dict[str, str] = {} logger.debug( - "SessionHandler 构造: session_id=%s max_parallel=%s endpoints=%s", - session_id, m, [ep.endpoint_id for ep in self._endpoints], + "ServiceScopeHandler 构造: session_id=%s max_parallel=%s endpoints=%s", + self._service_id, m, [ep.endpoint_id for ep in self._endpoints], ) @property - def session_id(self) -> str: - return self._session_id + def service_id(self) -> str: + return self._service_id @property def active_rids(self) -> Set[str]: @@ -69,14 +69,14 @@ def add_endpoint(self, endpoint: ISendEndpoint) -> None: """弹性扩容:追加一个发送端点。""" if any(ep.endpoint_id == endpoint.endpoint_id for ep in self._endpoints): logger.debug( - "SessionHandler 忽略重复 endpoint: session_id=%s endpoint_id=%s", - self._session_id, endpoint.endpoint_id, + "ServiceScopeHandler 忽略重复 endpoint: session_id=%s endpoint_id=%s", + self._service_id, endpoint.endpoint_id, ) return self._endpoints.append(endpoint) logger.info( - "SessionHandler 添加 endpoint: session_id=%s endpoint_id=%s endpoint_count=%s", - self._session_id, endpoint.endpoint_id, len(self._endpoints), + "ServiceScopeHandler 添加 endpoint: session_id=%s endpoint_id=%s endpoint_count=%s", + self._service_id, endpoint.endpoint_id, len(self._endpoints), ) def remove_endpoint(self, endpoint_id: str) -> bool: @@ -87,8 +87,8 @@ def remove_endpoint(self, endpoint_id: str) -> bool: """ if len(self._endpoints) <= 1: logger.debug( - "SessionHandler 拒绝移除最后一个 endpoint: session_id=%s endpoint_id=%s", - self._session_id, endpoint_id, + "ServiceScopeHandler 拒绝移除最后一个 endpoint: session_id=%s endpoint_id=%s", + self._service_id, endpoint_id, ) return False for i, ep in enumerate(self._endpoints): @@ -100,8 +100,8 @@ def remove_endpoint(self, endpoint_id: str) -> bool: if eid != endpoint_id } logger.info( - "SessionHandler 移除 endpoint: session_id=%s endpoint_id=%s remaining=%s", - self._session_id, endpoint_id, len(self._endpoints), + "ServiceScopeHandler 移除 endpoint: session_id=%s endpoint_id=%s remaining=%s", + self._service_id, endpoint_id, len(self._endpoints), ) return True return False @@ -131,7 +131,7 @@ def _pick_endpoint(self, user_id: Optional[str]) -> Optional[ISendEndpoint]: return best - async def handle_message(self, msg: SessionRequestWrapper) -> None: + async def handle_message(self, msg: ScopeRequestWrapper) -> None: sreq = msg.session_request # 会话内并发位:满则在此等待,不阻塞其它 session await self._sem.acquire() @@ -139,7 +139,7 @@ async def handle_message(self, msg: SessionRequestWrapper) -> None: if sreq.request_id: rid = sreq.request_id self._active_rids.add(rid) - await self._session_router.set_request_session(rid, self._session_id) + await self._session_router.set_request_session(rid, self._service_id) user_id = None raw = sreq.raw_msg @@ -154,15 +154,15 @@ async def handle_message(self, msg: SessionRequestWrapper) -> None: await self._session_router.delete_request_session(rid) self._sem.release() logger.error( - "SessionHandler 无可用 endpoint, 丢弃消息: session_id=%s request_id=%s", - self._session_id, rid, + "ServiceScopeHandler 无可用 endpoint, 丢弃消息: session_id=%s request_id=%s", + self._service_id, rid, ) - raise RuntimeError(f"SessionHandler {self._session_id} 无可用 endpoint") + raise RuntimeError(f"ServiceScopeHandler {self._service_id} 无可用 endpoint") logger.debug( - "SessionHandler 已获会话并发: session_id=%s request_id=%s " + "ServiceScopeHandler 已获会话并发: session_id=%s request_id=%s " "endpoint=%s 活跃rid数=%s", - self._session_id, rid, endpoint.endpoint_id, len(self._active_rids), + self._service_id, rid, endpoint.endpoint_id, len(self._active_rids), ) try: await endpoint.send_message(msg) @@ -172,6 +172,6 @@ async def handle_message(self, msg: SessionRequestWrapper) -> None: await self._session_router.delete_request_session(rid) self._sem.release() logger.debug( - "SessionHandler 释放会话并发: session_id=%s request_id=%s endpoint=%s", - self._session_id, rid, endpoint.endpoint_id, + "ServiceScopeHandler 释放会话并发: session_id=%s request_id=%s endpoint=%s", + self._service_id, rid, endpoint.endpoint_id, ) diff --git a/management/openjiuwen_runtime/management/session/session_request.py b/management/openjiuwen_runtime/management/session/session_request.py index 9b0ac75..32e50f6 100644 --- a/management/openjiuwen_runtime/management/session/session_request.py +++ b/management/openjiuwen_runtime/management/session/session_request.py @@ -14,14 +14,14 @@ class SessionRequest(ISessionRequest): def __init__( self, - session_id: str, + service_id: str, concurrency: int, ttl: int, request_id: Optional[str], raw: IRequest, service_template: Optional[dict] = None, ): - self._session_id = session_id + self._service_id = service_id self._concurrency = concurrency self._ttl = ttl self._request_id = request_id @@ -29,8 +29,8 @@ def __init__( self._service_template = service_template @property - def session_id(self) -> str: - return self._session_id + def service_id(self) -> str: + return self._service_id @property def session_concurrency(self) -> int: diff --git a/management/openjiuwen_runtime/management/session/session_runtime_manager.py b/management/openjiuwen_runtime/management/session/session_runtime_manager.py index e98724c..de98f40 100644 --- a/management/openjiuwen_runtime/management/session/session_runtime_manager.py +++ b/management/openjiuwen_runtime/management/session/session_runtime_manager.py @@ -1,13 +1,13 @@ # coding: utf-8 # Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved -"""Session 运行时管理:SessionHandler 归属、消息入口、session TTL/亲和/pending 过期。 +"""Session 运行时管理:ServiceScopeHandler 归属、消息入口、session TTL/亲和/pending 过期。 第二层解耦:从 ServiceManager 抽取的 session 编排职责。与 ServiceManager(Pod 池)协作: - 向 ServiceManager 借/还 endpoint(Pod); - ServiceManager 在 Pod inflight 归零时回调 :meth:`flush_pending_for_service`。 -内含 :class:`SessionRegistry`(SessionHandler 归属,替代原 ServiceHandler._sessions)。 +内含 :class:`ServiceScopeRegistry`(ServiceScopeHandler 归属,替代原 ServiceHandler._sessions)。 """ from __future__ import annotations @@ -22,10 +22,10 @@ IServiceHandler, ITimer, RawMessage, - SessionRequestWrapper, + ScopeRequestWrapper, ) from .router import SessionRouter -from .session_handler import SessionHandler +from .service_scope_handler import ServiceScopeHandler if TYPE_CHECKING: # 仅用于类型提示,避免与 ServiceManager 的运行时循环导入 @@ -34,15 +34,15 @@ logger = get_logger(__name__) -class SessionRegistry: - """SessionHandler 的归属与生命周期管理(替代原 ServiceHandler._sessions)。 +class ServiceScopeRegistry: + """ServiceScopeHandler 的归属与生命周期管理(替代原 ServiceHandler._sessions)。 - 每个 session_id 对应一个 SessionHandler;endpoints 由 SessionRuntimeManager 装填。 + 每个 session_id 对应一个 ServiceScopeHandler;endpoints 由 SessionRuntimeManager 装填。 """ def __init__(self, session_router: SessionRouter) -> None: self._lock = asyncio.Lock() - self._handlers: Dict[str, SessionHandler] = {} + self._handlers: Dict[str, ServiceScopeHandler] = {} self._session_router = session_router @property @@ -51,28 +51,28 @@ def session_router(self) -> SessionRouter: async def get_or_create( self, session_id: str, max_parallel: int, - ) -> SessionHandler: - """获取或创建 SessionHandler(endpoints 为空,由 SessionRuntimeManager 装填)。""" + ) -> ServiceScopeHandler: + """获取或创建 ServiceScopeHandler(endpoints 为空,由 SessionRuntimeManager 装填)。""" async with self._lock: sh = self._handlers.get(session_id) if sh is None: - sh = SessionHandler( + sh = ServiceScopeHandler( session_id, max_parallel, [], self._session_router ) self._handlers[session_id] = sh logger.info( - "SessionRegistry 新建 SessionHandler: session_id=%s max_parallel=%s", + "ServiceScopeRegistry 新建 ServiceScopeHandler: session_id=%s max_parallel=%s", session_id, max_parallel, ) return sh - def get(self, session_id: str) -> Optional[SessionHandler]: + def get(self, session_id: str) -> Optional[ServiceScopeHandler]: return self._handlers.get(session_id) def has_session(self, session_id: str) -> bool: return session_id in self._handlers - async def remove(self, session_id: str) -> Optional[SessionHandler]: + async def remove(self, session_id: str) -> Optional[ServiceScopeHandler]: async with self._lock: return self._handlers.pop(session_id, None) @@ -96,24 +96,24 @@ def __init__( ) -> None: self._timer = timer self._sm = service_manager - self._registry = SessionRegistry(SessionRouter()) + self._registry = ServiceScopeRegistry(SessionRouter()) self._lock = asyncio.Lock() # session_id -> 代表性 service_id(TTL 到期时仍有 inflight 的延期清理队列) - self._pending_expired: Dict[str, str] = {} + self._pending_scope_expiry: Dict[str, str] = {} @property - def registry(self) -> SessionRegistry: + def registry(self) -> ServiceScopeRegistry: return self._registry # ==================== 消息入口 ==================== async def handle_user_request(self, raw: RawMessage) -> None: w = raw.message - if not isinstance(w, SessionRequestWrapper): - logger.error("用户消息体类型错误, 预期 SessionRequestWrapper, 实际 %s", type(w)) + if not isinstance(w, ScopeRequestWrapper): + logger.error("用户消息体类型错误, 预期 ScopeRequestWrapper, 实际 %s", type(w)) return sreq = w.session_request - session_id = sreq.session_id + session_id = sreq.service_id if self._sm.is_deprecated(): logger.warning( @@ -123,10 +123,10 @@ async def handle_user_request(self, raw: RawMessage) -> None: ) # 新消息进入:先取消该 session 待生效的 session_ttl 计时与延期清理标记 - await self._timer.cancel_timer(f"sess:{session_id}") - self._pending_expired.pop(session_id, None) + await self._timer.cancel_timer(f"scope:{session_id}") + self._pending_scope_expiry.pop(session_id, None) - sh: Optional[SessionHandler] = None + sh: Optional[ServiceScopeHandler] = None failed_to_reserve: Optional[IServiceHandler] = None try: sh = await self._registry.get_or_create(session_id, sreq.session_concurrency) @@ -167,7 +167,7 @@ async def handle_user_request(self, raw: RawMessage) -> None: await self._fail(w, 100001, session_id=session_id) return - # 转发到 SessionHandler(不再经过 ServiceHandler.handle_message) + # 转发到 ServiceScopeHandler(不再经过 ServiceHandler.handle_message) endpoint_id = self._peek_endpoint_id(sh) pod_name = self._pod_name_for(endpoint_id) logger.info( @@ -190,7 +190,7 @@ async def handle_user_request(self, raw: RawMessage) -> None: if sh is not None: if sreq.session_ttl > 0: try: - await self._arm_session_timer(session_id, sreq.session_ttl) + await self._arm_scope_timer(session_id, sreq.session_ttl) except Exception as e2: # noqa: BLE001 logger.error("arm session 计时器失败: %s", e2, exc_info=True) else: @@ -201,7 +201,7 @@ async def handle_user_request(self, raw: RawMessage) -> None: else: rep = self._peek_endpoint_id(sh) if rep is not None: - self._pending_expired[session_id] = rep + self._pending_scope_expiry[session_id] = rep logger.debug( "session ttl<=0 pending expired: session_id=%s service_id=%s", session_id, rep, @@ -209,7 +209,8 @@ async def handle_user_request(self, raw: RawMessage) -> None: # 由 inflight 归零钩子触发 flush;此处主动尝试一次 await self.flush_pending_for_service(rep) - def _peek_endpoint_id(self, sh: SessionHandler) -> Optional[str]: + @staticmethod + def _peek_endpoint_id(sh: ServiceScopeHandler) -> Optional[str]: ids = sh.endpoint_ids return ids[0] if ids else None @@ -223,21 +224,21 @@ def _pod_name_for(self, service_id: Optional[str]) -> str: # ==================== TTL / 过期 ==================== - async def _arm_session_timer(self, session_id: str, ttl: int) -> None: + async def _arm_scope_timer(self, session_id: str, ttl: int) -> None: if ttl <= 0: return - key = f"sess:{session_id}" + key = f"scope:{session_id}" await self._timer.cancel_timer(key) # 重新 arm 等同于「session 又活跃」, 清掉旧的延期清理标记 - self._pending_expired.pop(session_id, None) + self._pending_scope_expiry.pop(session_id, None) async def _expired() -> None: - await self._on_session_expired(session_id) + await self._on_scope_expired(session_id) await self._timer.start_timer(key, ttl, _expired) logger.info("已 arm session TTL 计时: session_id=%s ttl=%s", session_id, ttl) - async def _on_session_expired(self, session_id: str) -> None: + async def _on_scope_expired(self, session_id: str) -> None: """session_ttl 到期: 若该 session 已无 inflight 则立即移除并触发 flush;否则入 pending 等待。""" logger.info("session TTL 到期, 准备回收: session_id=%s", session_id) target_svc: Optional[str] = None @@ -251,14 +252,14 @@ async def _on_session_expired(self, session_id: str) -> None: if len(sh.active_rids) > 0: rep = self._peek_endpoint_id(sh) if rep is not None: - self._pending_expired[session_id] = rep + self._pending_scope_expiry[session_id] = rep logger.info( "session 到期但仍有 inflight, 入延期清理队列: session_id=%s rids=%s", session_id, len(sh.active_rids), ) return - # 在途已空:逐 Pod 释放 quota + 取消,再销毁 SessionHandler + # 在途已空:逐 Pod 释放 quota + 取消,再销毁 ServiceScopeHandler evicted_endpoints: list[str] = [] for endpoint_id in sh.endpoint_ids: h = self._sm.find_service_handler(endpoint_id) @@ -268,7 +269,7 @@ async def _on_session_expired(self, session_id: str) -> None: if target_svc is None: target_svc = endpoint_id await self._registry.remove(session_id) - self._pending_expired.pop(session_id, None) + self._pending_scope_expiry.pop(session_id, None) removed = True logger.info( "session 已移除并归还 service 并发: session_id=%s", session_id, @@ -295,13 +296,13 @@ async def flush_pending_for_service(self, service_id: str) -> None: async with self._lock: sids = [ sid - for sid, svc in list(self._pending_expired.items()) + for sid, svc in list(self._pending_scope_expiry.items()) if svc == service_id ] for sid in sids: sh = self._registry.get(sid) if sh is None: - self._pending_expired.pop(sid, None) + self._pending_scope_expiry.pop(sid, None) continue if len(sh.active_rids) == 0: for endpoint_id in sh.endpoint_ids: @@ -309,7 +310,7 @@ async def flush_pending_for_service(self, service_id: str) -> None: if h is not None: await h.evict_session(sid) await self._registry.remove(sid) - self._pending_expired.pop(sid, None) + self._pending_scope_expiry.pop(sid, None) logger.info( "延期清理已移除 session: session_id=%s service_id=%s", sid, service_id, @@ -318,7 +319,7 @@ async def flush_pending_for_service(self, service_id: str) -> None: # ==================== 失败响应 ==================== async def _fail( - self, w: SessionRequestWrapper, code: int, *, session_id: str | None = None + self, w: ScopeRequestWrapper, code: int, *, session_id: str | None = None ) -> None: em = exception_message(code) if code == 100001: @@ -360,7 +361,7 @@ async def _fail( # ==================== 生命周期(由 ServiceManager.stop 调用)==================== async def on_pod_removed(self, service_id: str) -> None: - """某 Pod 已失效被移除:从所有引用它的 SessionHandler 中摘除该 endpoint。 + """某 Pod 已失效被移除:从所有引用它的 ServiceScopeHandler 中摘除该 endpoint。 - 单 endpoint 的 session:摘除后 endpoint_count==0,下条消息会触发 pick_or_create 重新装填; - 多 endpoint 的 session:摘除该 endpoint,其余继续服务; @@ -377,12 +378,12 @@ async def on_pod_removed(self, service_id: str) -> None: sid, service_id, sh.endpoint_count, ) # 清理指向该 Pod 的 pending 记录 - stale = [sid for sid, svc in self._pending_expired.items() if svc == service_id] + stale = [sid for sid, svc in self._pending_scope_expiry.items() if svc == service_id] for sid in stale: - self._pending_expired.pop(sid, None) + self._pending_scope_expiry.pop(sid, None) async def shutdown(self) -> None: - """清理所有 session TTL 计时与 pending 标记(SessionHandler 由 registry 丢弃)。""" + """清理所有 session TTL 计时与 pending 标记(ServiceScopeHandler 由 registry 丢弃)。""" for sid in list(self._registry.all_session_ids()): - await self._timer.cancel_timer(f"sess:{sid}") - self._pending_expired.clear() + await self._timer.cancel_timer(f"scope:{sid}") + self._pending_scope_expiry.clear() diff --git a/management/openjiuwen_runtime/management/session/strategies/_base.py b/management/openjiuwen_runtime/management/session/strategies/_base.py index 7dafc9f..2f41c77 100644 --- a/management/openjiuwen_runtime/management/session/strategies/_base.py +++ b/management/openjiuwen_runtime/management/session/strategies/_base.py @@ -30,7 +30,7 @@ def handle_session(self, msg: IRequest) -> ISessionRequest: logger.warning("策略生成的 session 键为空或全空白, request_id=%s", msg.request_id) logger.debug("策略 handle_session: session_id=%s conc=%s ttl=%s", sid, self._concurrency, self._ttl) return SessionRequest( - session_id=sid, + service_id=sid, concurrency=self._concurrency, ttl=self._ttl, request_id=msg.request_id, diff --git a/management/openjiuwen_runtime/management/session/ws_client_channel.py b/management/openjiuwen_runtime/management/session/ws_client_channel.py index 68d836a..b33db03 100644 --- a/management/openjiuwen_runtime/management/session/ws_client_channel.py +++ b/management/openjiuwen_runtime/management/session/ws_client_channel.py @@ -20,7 +20,7 @@ IResponseParser, IServiceHandler, OnRequestCompleteCallback, - SessionRequestWrapper, + ScopeRequestWrapper, ) logger = get_logger(__name__) @@ -295,7 +295,7 @@ async def _recv_loop(self) -> None: async def send( self, service_id: str, - wrapper: SessionRequestWrapper, + wrapper: ScopeRequestWrapper, *, response_parser: IResponseParser, on_request_complete: OnRequestCompleteCallback, diff --git a/management/tests/system_tests/management_session/main_k8s_access.py b/management/tests/system_tests/management_session/main_k8s_access.py index 0b014ee..6c7d668 100644 --- a/management/tests/system_tests/management_session/main_k8s_access.py +++ b/management/tests/system_tests/management_session/main_k8s_access.py @@ -246,8 +246,8 @@ def _load_message_payload(ns: argparse.Namespace) -> dict[str, Any]: data["bot_id"] = ns.bot_id if ns.user_id is not None: data["user_id"] = ns.user_id - if ns.session_id is not None: - data["session_id"] = ns.session_id + if ns.service_id is not None: + data["session_id"] = ns.service_id return data diff --git a/management/tests/system_tests/management_session/main_k8s_multi_session.py b/management/tests/system_tests/management_session/main_k8s_multi_session.py index ca88279..b0f2779 100644 --- a/management/tests/system_tests/management_session/main_k8s_multi_session.py +++ b/management/tests/system_tests/management_session/main_k8s_multi_session.py @@ -332,7 +332,7 @@ async def _run_one( rid = f"req_{spec.sid}_{seq}_{uuid.uuid4().hex[:6]}" msg = _render_message(template, spec.sid, seq, rid) sreq = SessionRequest( - session_id=spec.sid, + service_id=spec.sid, concurrency=spec.conc, ttl=spec.ttl, request_id=rid, diff --git a/management/tests/system_tests/management_session/test_session_sdk.py b/management/tests/system_tests/management_session/test_session_sdk.py index e438834..8e1c238 100644 --- a/management/tests/system_tests/management_session/test_session_sdk.py +++ b/management/tests/system_tests/management_session/test_session_sdk.py @@ -18,7 +18,7 @@ IServiceInstanceFactory, IServiceHandler, IRequest, - SessionRequestWrapper, + ScopeRequestWrapper, ) from openjiuwen_runtime.management.session.models import AccessConfig, SessionConfig from openjiuwen_runtime.management.session.runtime import NoOpDeployController @@ -67,13 +67,13 @@ def __init__(self) -> None: async def send( self, service_id: str, - wrapper: SessionRequestWrapper, + wrapper: ScopeRequestWrapper, *, response_parser: IResponseParser, on_request_complete: Callable[[Optional[str]], Awaitable[None]], ) -> None: sreq = wrapper.session_request - self.send_calls.append((service_id, sreq.session_id, sreq.request_id)) + self.send_calls.append((service_id, sreq.service_id, sreq.request_id)) if wrapper.cancel.done(): await on_request_complete(sreq.request_id) return diff --git a/management/tests/unit_tests/management_session/test_access_session_st.py b/management/tests/unit_tests/management_session/test_access_session_st.py index 7010e2c..af23e68 100644 --- a/management/tests/unit_tests/management_session/test_access_session_st.py +++ b/management/tests/unit_tests/management_session/test_access_session_st.py @@ -20,7 +20,7 @@ IServiceInstanceFactory, IServiceHandler, IRequest, - SessionRequestWrapper, + ScopeRequestWrapper, ) from openjiuwen_runtime.management.session.models import AccessConfig, SessionConfig from openjiuwen_runtime.management.session.runtime import NoOpDeployController @@ -59,12 +59,12 @@ def __init__(self) -> None: async def send( self, service_id: str, - wrapper: SessionRequestWrapper, + wrapper: ScopeRequestWrapper, *, response_parser: IResponseParser, on_request_complete: Callable[[Optional[str]], Awaitable[None]], ) -> None: - self.send_log.append((service_id, wrapper.session_request.session_id)) + self.send_log.append((service_id, wrapper.session_request.service_id)) if wrapper.cancel.done(): await on_request_complete(wrapper.session_request.request_id) return @@ -130,12 +130,12 @@ def __init__(self) -> None: async def send( self, service_id: str, - wrapper: SessionRequestWrapper, + wrapper: ScopeRequestWrapper, *, response_parser: IResponseParser, on_request_complete: Callable[[Optional[str]], Awaitable[None]], ) -> None: - sid = wrapper.session_request.session_id + sid = wrapper.session_request.service_id rid = wrapper.session_request.request_id async with self._lock: self.in_send += 1 diff --git a/management/tests/unit_tests/management_session/test_service_handler_unit.py b/management/tests/unit_tests/management_session/test_service_handler_unit.py index 8a536b0..c0bff4d 100644 --- a/management/tests/unit_tests/management_session/test_service_handler_unit.py +++ b/management/tests/unit_tests/management_session/test_service_handler_unit.py @@ -11,12 +11,12 @@ IResponseParser, IRequest, ISessionRequest, - SessionRequestWrapper, + ScopeRequestWrapper, ) from openjiuwen_runtime.management.session.router import SessionRouter from openjiuwen_runtime.management.session.runtime import NoOpDeployController from openjiuwen_runtime.management.session.service_handler import ServiceHandler -from openjiuwen_runtime.management.session.session_handler import SessionHandler +from openjiuwen_runtime.management.session.service_scope_handler import ServiceScopeHandler from openjiuwen_runtime.management.session.session_request import SessionRequest @@ -39,7 +39,7 @@ def __init__(self) -> None: async def send( self, service_id: str, - wrapper: SessionRequestWrapper, + wrapper: ScopeRequestWrapper, *, response_parser: IResponseParser, on_request_complete: Callable[[Optional[str]], Awaitable[None]], @@ -57,7 +57,7 @@ async def send( def _sreq() -> ISessionRequest: return SessionRequest( - session_id="s1", + service_id="s1", concurrency=1, ttl=0, request_id="r1", @@ -80,9 +80,9 @@ async def test_one_inflight_decrements() -> None: assert h.available_concurrency == 1 assert h.try_reserve_session_quota("s1", 1) assert h.available_concurrency == 0 - # 消息经 SessionHandler(持有 [h] 作为 endpoint)路由到 ServiceHandler.send_message - sh = SessionHandler("s1", 1, [h], SessionRouter()) - w = SessionRequestWrapper( + # 消息经 ServiceScopeHandler(持有 [h] 作为 endpoint)路由到 ServiceHandler.send_message + sh = ServiceScopeHandler("s1", 1, [h], SessionRouter()) + w = ScopeRequestWrapper( _sreq(), asyncio.Queue(), asyncio.get_running_loop().create_future() ) await sh.handle_message(w) diff --git a/management/tests/unit_tests/management_session/test_service_manager_deploy_lock.py b/management/tests/unit_tests/management_session/test_service_manager_deploy_lock.py index 7de5d72..aebcfbe 100644 --- a/management/tests/unit_tests/management_session/test_service_manager_deploy_lock.py +++ b/management/tests/unit_tests/management_session/test_service_manager_deploy_lock.py @@ -50,7 +50,7 @@ async def send(self, *args: Any, **kwargs: Any) -> None: def _sreq(session_id: str, request_id: str = "r") -> SessionRequest: return SessionRequest( - session_id=session_id, + service_id=session_id, concurrency=1, ttl=0, request_id=request_id, diff --git a/management/tests/unit_tests/management_session/test_session_concurrency_interleave.py b/management/tests/unit_tests/management_session/test_session_concurrency_interleave.py index c454744..f813ed7 100644 --- a/management/tests/unit_tests/management_session/test_session_concurrency_interleave.py +++ b/management/tests/unit_tests/management_session/test_session_concurrency_interleave.py @@ -15,12 +15,12 @@ IResponseParser, IRequest, ISessionRequest, - SessionRequestWrapper, + ScopeRequestWrapper, ) from openjiuwen_runtime.management.session.router import SessionRouter from openjiuwen_runtime.management.session.runtime import NoOpDeployController from openjiuwen_runtime.management.session.service_handler import ServiceHandler -from openjiuwen_runtime.management.session.session_handler import SessionHandler +from openjiuwen_runtime.management.session.service_scope_handler import ServiceScopeHandler from openjiuwen_runtime.management.session.session_request import SessionRequest @@ -49,12 +49,12 @@ def __init__(self) -> None: async def send( self, service_id: str, - wrapper: SessionRequestWrapper, + wrapper: ScopeRequestWrapper, *, response_parser: IResponseParser, on_request_complete: Callable[[Optional[str]], Awaitable[None]], ) -> None: - sid = wrapper.session_request.session_id + sid = wrapper.session_request.service_id rid = wrapper.session_request.request_id async with self._lock: self.in_send += 1 @@ -76,15 +76,15 @@ async def send( def _wrap( session_id: str, rid: str, cap: int, loop: asyncio.AbstractEventLoop -) -> SessionRequestWrapper: +) -> ScopeRequestWrapper: sreq: ISessionRequest = SessionRequest( - session_id=session_id, + service_id=session_id, concurrency=cap, ttl=0, request_id=rid, raw=cast(IRequest, object()), ) - return SessionRequestWrapper(sreq, asyncio.Queue(), loop.create_future()) + return ScopeRequestWrapper(sreq, asyncio.Queue(), loop.create_future()) @pytest.mark.asyncio @@ -104,12 +104,12 @@ async def test_session_cap_interleaves_other_sessions() -> None: ) loop = asyncio.get_running_loop() cap = 10 - # 解耦后:每 session 一个 SessionHandler(持有 [h] 作为 endpoint),各自 semaphore(cap) + # 解耦后:每 session 一个 ServiceScopeHandler(持有 [h] 作为 endpoint),各自 semaphore(cap) router = SessionRouter() assert h.try_reserve_session_quota("sess1", cap) assert h.try_reserve_session_quota("sess2", cap) - sh1 = SessionHandler("sess1", cap, [h], router) - sh2 = SessionHandler("sess2", cap, [h], router) + sh1 = ServiceScopeHandler("sess1", cap, [h], router) + sh2 = ServiceScopeHandler("sess2", cap, [h], router) tasks = [] for i in range(11): w = _wrap("sess1", f"s1-r{i}", cap, loop)