已合并
feat(rm): deploy 锁输家改为 follower 等待室——跨副本冷竞争零多余 Pod(M8) #438
feat(rm): deploy 锁输家改为 follower 等待室——跨副本冷竞争零多余 Pod(M8) #438
已合并
王明琦创建于 8月24日
共 13 个文件变更+336-43
@@ -32,7 +32,7 @@ cp agent_runtime.server.env.example .env.production.local
32```bash32```bash
33cd applications/agent_runtime33cd applications/agent_runtime
34uv sync --extra local # 含 dev 组(fakeredis[lua]/pytest)34uv sync --extra local # 含 dev 组(fakeredis[lua]/pytest)
35-uv run pytest # 106 个用例:状态层 Lua / config 层 / 组件全链路 / HTTP 冒烟 / 双实例多副本35+uv run pytest # 114 个用例:状态层 Lua / config 层 / 组件全链路 / HTTP 冒烟 / 双实例多副本
36```36```
37 37 
38### 集成冒烟测试(真环境,M6 用例固化)38### 集成冒烟测试(真环境,M6 用例固化)
@@ -1,12 +1,16 @@
1# coding: utf-81# coding: utf-8
2-"""Resource Manager 的 4 个 Lua 脚本(所有编排态变更,原子)。2+"""Resource Manager 的 6 个 Lua 脚本(所有编排态变更,原子)。
3 3 
4约定同 SM:``ARGV[1]`` 恒为键前缀(``resource_manager:``);返回扁平字符串数组。4约定同 SM:``ARGV[1]`` 恒为键前缀(``resource_manager:``);返回扁平字符串数组。
5脚本清单(语义见 RM 设计 §5.1):5脚本清单(语义见 RM 设计 §5.1):
6- LUA_ACQUIRE 取暖 Pod 复用(deploy_ver 过滤)/ 判 max_pods(含 deploying 占位)/ 占位6- LUA_ACQUIRE 取暖 Pod 复用(deploy_ver 过滤)/ 判 max_pods(含 deploying 占位)/ 占位
7+- LUA_PLACEHOLDER autoscale 专用占位(判 max_pods + SADD,不碰 idle 池)
7- LUA_REGISTER deploy 成功登记(info / scope:pods / pods:all,清占位;热备入 idle)8- LUA_REGISTER deploy 成功登记(info / scope:pods / pods:all,清占位;热备入 idle)
8- LUA_RELEASE idle_consider:转 idle 暖池 + 起 pod_ttl 计时(幂等)9- LUA_RELEASE idle_consider:转 idle 暖池 + 起 pod_ttl 计时(幂等)
9- LUA_PURGE Pod 死亡 / reclaim 后清全部 RM key(幂等)10- LUA_PURGE Pod 死亡 / reclaim 后清全部 RM key(幂等)
11+- LUA_DEPLOY_FOLLOWER_GATE deploy 锁输家的等待室原子准入(ZSET+deadline,
12+ 上限 pod_concurrency-1;先清过期成员再 ZADD 先行+超限自退——同
13+ LUA_WAITER_GATE 纪律,禁止先查后加)
10"""14"""
11 15 
12from __future__ import annotations16from __future__ import annotations
@@ -128,3 +132,25 @@ end
128redis.call('SREM', pfx .. 'resource:pods:all', pod)132redis.call('SREM', pfx .. 'resource:pods:all', pod)
129return {'ok', scope or ''}133return {'ok', scope or ''}
130"""134"""
135+ 
136+# Argv: prefix, scope_id, follower_id, max_followers, deadline, now
137+# deploy 锁输家(follower)等待室原子准入。ZSET 以 deadline(秒级时间戳)为
138+# score:先 ZREMRANGEBYSCORE 清过期成员(等待进程崩溃的兜底,不泄漏),
139+# 再 ZADD 先行 + ZCARD 超限自退(并发同时到达也不会超收)。
140+LUA_DEPLOY_FOLLOWER_GATE = r"""
141+local pfx = ARGV[1]
142+local scope = ARGV[2]
143+local follower = ARGV[3]
144+local max = tonumber(ARGV[4])
145+local deadline = tonumber(ARGV[5])
146+local now = tonumber(ARGV[6])
147+ 
148+local key = pfx .. 'resource:scope:' .. scope .. ':deploy_followers'
149+redis.call('ZREMRANGEBYSCORE', key, '-inf', now)
150+redis.call('ZADD', key, deadline, follower)
151+if redis.call('ZCARD', key) > max then
152+ redis.call('ZREM', key, follower)
153+ return {'false'}
154+end
155+return {'true'}
156+"""
@@ -20,14 +20,15 @@ from uuid import uuid4
20from ..errors import DeployFailed, MaxPodsReached20from ..errors import DeployFailed, MaxPodsReached
21from ..spec_fields import DEPLOY_VER_FIELDS21from ..spec_fields import DEPLOY_VER_FIELDS
22from ..util import fingerprint, now_ts22from ..util import fingerprint, now_ts
23-from .k8s import K8sPodClient23+from .k8s import DEFAULT_READY_TIMEOUT, K8sPodClient
24from .state import ResourceState24from .state import ResourceState
25 25 
26logger = logging.getLogger("agent_runtime.resource_manager")26logger = logging.getLogger("agent_runtime.resource_manager")
27 27 
28ACQUIRE_IDEM_TTL = 60 # acquire 结果幂等缓存窗口28ACQUIRE_IDEM_TTL = 60 # acquire 结果幂等缓存窗口
29DEPLOY_LOCK_TTL = 360 # per-scope deploy 锁(盖住 ready_timeout 300s + 余量)29DEPLOY_LOCK_TTL = 360 # per-scope deploy 锁(盖住 ready_timeout 300s + 余量)
30-DEPLOY_WAIT_ON_BUSY = 0.3 # 他副本在 deploy 时的重试间隔30+DEPLOY_WAIT_ON_BUSY = 0.3 # follower 轮询间隔(原输家自旋间隔沿用)
31+FOLLOWER_WAIT_MARGIN = 10 # follower 等待上界 = ready_timeout + 此余量(注册开销)
31 32 
32 33 
33def _deploy_ver(pod_spec: dict[str, Any]) -> str:34def _deploy_ver(pod_spec: dict[str, Any]) -> str:
@@ -70,6 +71,7 @@ class ResourceOrchestrator:
70 "min_idle_pods": int(pool_config.get("min_idle_pods", 0)),71 "min_idle_pods": int(pool_config.get("min_idle_pods", 0)),
71 "max_pods": int(pool_config.get("max_pods", 1)),72 "max_pods": int(pool_config.get("max_pods", 1)),
72 "pod_ttl": int(pool_config.get("pod_ttl", 300)),73 "pod_ttl": int(pool_config.get("pod_ttl", 300)),
74+ "pod_concurrency": int(pool_config.get("pod_concurrency", 1)),
73 "deploy_ver": deploy_ver,75 "deploy_ver": deploy_ver,
74 "pod_spec_json": json.dumps(pod_spec),76 "pod_spec_json": json.dumps(pod_spec),
75 },77 },
@@ -97,9 +99,16 @@ class ResourceOrchestrator:
97 lock_key = self.state.k.lock_deploy(scope_id)99 lock_key = self.state.k.lock_deploy(scope_id)
98 lock_token = f"deploy-{token}"100 lock_token = f"deploy-{token}"
99 if not await self.state.try_lock(lock_key, DEPLOY_LOCK_TTL, lock_token):101 if not await self.state.try_lock(lock_key, DEPLOY_LOCK_TTL, lock_token):
102+ # 输家:清占位 → 进 follower 等待室(上限 pc-1,overflow 严格
103+ # 快失败),等 leader 的 Pod 注册后**直接复用**——RM 全程不读
104+ # SM 的容量键(闸门归属不变),也修掉「输家自建第 2 个空 Pod」
105+ # 的跨副本冷竞争浪费(见 handoff §十一.1 开放问题)
100 await self.state.clear_deploy_token(scope_id, token)106 await self.state.clear_deploy_token(scope_id, token)
101- await asyncio.sleep(DEPLOY_WAIT_ON_BUSY) # 他副本在 deploy,稍后复用其成果107+ result = await self._follow_leader(
102- continue108+ scope_id, pod_spec, pool_config, request_id)
109+ if request_id:
110+ await self._idem_put(request_id, result)
111+ return result
103 try:112 try:
104 pod_id, sse_url = await self._deploy_and_register(113 pod_id, sse_url = await self._deploy_and_register(
105 scope_id, pod_spec, deploy_ver, token, idle_flag=False114 scope_id, pod_spec, deploy_ver, token, idle_flag=False
@@ -112,6 +121,68 @@ class ResourceOrchestrator:
112 await self._idem_put(request_id, result)121 await self._idem_put(request_id, result)
113 return result122 return result
114 123 
124+ async def _follow_leader(
125+ self,
126+ scope_id: str,
127+ pod_spec: dict[str, Any],
128+ pool_config: dict[str, Any],
129+ request_id: str,
130+ ) -> dict[str, str]:
131+ """deploy 锁输家的 follower 等待室:等 leader 的 Pod 注册后直接复用。
132+ 
133+ 设计定案(M8,讨论见 e2e-test-cases §8.2 / handoff §十一.1):
134+ - 准入 = 原子闸门(ZSET+deadline),上限 ``pod_concurrency - 1``——
135+ leader 会话之外新 Pod 恰剩这些槽;**overflow 严格快失败**;
136+ - 等待有界(ready_timeout + 余量);leader 失败(锁空闲且无进展)
137+ → follower **不接管**直接失败(同镜像同环境大概率也失败);
138+ - 检测到新 Pod 注册(进展)→ 返回该 Pod,与 reuse 分支同构——
139+ SM 侧重跑仲裁即可,RM 全程不读 SM 容量键(红线不破)。
140+ """
141+ pc = max(int(pool_config.get("pod_concurrency", 1)), 1)
142+ max_followers = pc - 1
143+ ready_timeout = int(pod_spec.get("ready_timeout") or DEFAULT_READY_TIMEOUT)
144+ now = now_ts()
145+ admitted = await self.state.try_add_deploy_follower(
146+ scope_id, request_id, max_followers,
147+ now + ready_timeout + FOLLOWER_WAIT_MARGIN, now,
148+ )
149+ if not admitted:
150+ raise MaxPodsReached(
151+ f"scope {scope_id} deploy followers full ({max_followers}); "
152+ f"retry after leader registers")
153+ 
154+ pods_before = set(await self.state.pod_ids(scope_id))
155+ deadline = time.monotonic() + ready_timeout + FOLLOWER_WAIT_MARGIN
156+ lock_key = self.state.k.lock_deploy(scope_id)
157+ try:
158+ while True:
159+ await asyncio.sleep(DEPLOY_WAIT_ON_BUSY)
160+ new_pods = [p for p in await self.state.pod_ids(scope_id)
161+ if p not in pods_before]
162+ if new_pods:
163+ for pod in new_pods:
164+ info = await self.state.pod_info(pod)
165+ if info.get("pod_sse_url"):
166+ logger.info(
167+ "acquire follower reuses leader pod: scope=%s "
168+ "pod=%s follower=%s", scope_id, pod, request_id)
169+ return {"pod_id": pod,
170+ "pod_sse_url": info["pod_sse_url"]}
171+ # 新 Pod 已入池但 info 尚未可见(原子脚本内不存在此窗口,
172+ # 防御性兜底):视为有进展,跳过失败判定等下一轮
173+ continue
174+ if not await self.state.lock_held(lock_key):
175+ raise DeployFailed(
176+ f"scope {scope_id} leader deploy aborted; follower aborts")
177+ if time.monotonic() >= deadline:
178+ raise MaxPodsReached(
179+ f"scope {scope_id} follower wait timeout "
180+ f"({ready_timeout + FOLLOWER_WAIT_MARGIN}s)")
181+ finally:
182+ # 错误路径双清纪律:follower 成员必须退出(防虚占 pc-1 名额);
183+ # 崩溃遗留由闸门的 ZREMRANGEBYSCORE(deadline) 兜底
184+ await self.state.remove_deploy_follower(scope_id, request_id)
185+ 
115 async def _deploy_and_register(186 async def _deploy_and_register(
116 self,187 self,
117 scope_id: str,188 scope_id: str,
Mapplications/agent_runtime/src/agent_runtime/resource_manager/state.py+37-0文件内容审核中,请稍后刷新重试
Mapplications/agent_runtime/src/agent_runtime/session_manager/models.py+6-1文件内容审核中,请稍后刷新重试
@@ -90,6 +90,10 @@ class SlowFakeK8sPodClient:
90 def __getattr__(self, name):90 def __getattr__(self, name):
91 return getattr(self._inner, name)91 return getattr(self._inner, name)
92 92 
93+ def set_deploy_failures(self, count: int) -> None:
94+ """写透内层 FakeK8s 的失败旋钮(__getattr__ 只代理读,不代理写)。"""
95+ self._inner.deploy_failures = count
96+ 
93 async def deploy(self, pod_spec: dict):97 async def deploy(self, pod_spec: dict):
94 t0 = time.monotonic()98 t0 = time.monotonic()
95 if self.deploy_delay:99 if self.deploy_delay:
Mapplications/agent_runtime/tests/integration/test_multi_replica.py+81-18文件内容审核中,请稍后刷新重试
Mapplications/agent_runtime/tests/resource_manager/test_rm_state.py+61-0文件内容审核中,请稍后刷新重试
@@ -45,7 +45,7 @@ src/agent_runtime/
45 ├── orchestrator.py acquire(取暖/选主 deploy/封顶)+ idle_consider45 ├── orchestrator.py acquire(取暖/选主 deploy/封顶)+ idle_consider
46 │ + update_pool_config + cleanup + 结果幂等缓存46 │ + update_pool_config + cleanup + 结果幂等缓存
47 ├── state.py RM Redis 键 schema 唯一出口47 ├── state.py RM Redis 键 schema 唯一出口
48- ├── lua_scripts.py 5 个 Lua 全文(见 §4)48+ ├── lua_scripts.py 6 个 Lua 全文(见 §4)
49 ├── k8s.py RealK8sPodClient(kubernetes_asyncio)/ FakeK8sPodClient49 ├── k8s.py RealK8sPodClient(kubernetes_asyncio)/ FakeK8sPodClient
50 │ (deploy 等 Ready、409 重命名重试、判死归一化、/health 探测)50 │ (deploy 等 Ready、409 重命名重试、判死归一化、/health 探测)
51 ├── sweeper.py autoscale / reclaim / watch(死 Pod+健康探测)/ reconcile51 ├── sweeper.py autoscale / reclaim / watch(死 Pod+健康探测)/ reconcile
@@ -55,7 +55,7 @@ src/agent_runtime/
55scripts/deploy.sh local|server 启动55scripts/deploy.sh local|server 启动
56scripts/e2e_hld_acceptance.py 集成冒烟(场景 A–L,真环境)56scripts/e2e_hld_acceptance.py 集成冒烟(场景 A–L,真环境)
57scripts/integration_smoke.sh 冒烟入口包装57scripts/integration_smoke.sh 冒烟入口包装
58-tests/ 106 个单测(见 §7)58+tests/ 114 个单测(见 §7)
59```59```
60 60 
61## 3. 关键流程(读代码的切入点)61## 3. 关键流程(读代码的切入点)
@@ -79,7 +79,8 @@ tests/ 106 个单测(见 §7)
79 reuse(暖 Pod,deploy_ver 过滤)→ 返回79 reuse(暖 Pod,deploy_ver 过滤)→ 返回
80 max_reached → 抛 MaxPodsReached(SM 映射 503 NO_POD_AVAILABLE)80 max_reached → 抛 MaxPodsReached(SM 映射 503 NO_POD_AVAILABLE)
81 need_deploy → 抢 lock:rm:deploy:{scope} 选主串行 deploy81 need_deploy → 抢 lock:rm:deploy:{scope} 选主串行 deploy
82- (他副本持锁 → 0.3s 等待重试,复用其成果)82+ (输家 → follower 等待室:准入≤pc-1/overflow 快失败/等 leader
83+ Pod 注册即复用/leader 失败不接管/等待有界)
83k8s.deploy:create + wait Ready(409 重命名重试;超时/镜像失败 → DeployFailed)84k8s.deploy:create + wait Ready(409 重命名重试;超时/镜像失败 → DeployFailed)
84错误路径必须清 deploying 占位(红线)85错误路径必须清 deploying 占位(红线)
85```86```
@@ -128,6 +129,7 @@ lock:config_sync 串行化(忙 → 409 CONFIG_SYNC_BUSY)
128| RM | `LUA_RELEASE` | idle_consider 转 idle 暖池(起 pod_ttl 计时) |129| RM | `LUA_RELEASE` | idle_consider 转 idle 暖池(起 pod_ttl 计时) |
129| RM | `LUA_PURGE` | 清该 Pod 全部 RM key(返回其 scope) |130| RM | `LUA_PURGE` | 清该 Pod 全部 RM key(返回其 scope) |
130| RM | `LUA_PLACEHOLDER` | autoscale 专用占位(计入 max_pods,不碰 idle 池) |131| RM | `LUA_PLACEHOLDER` | autoscale 专用占位(计入 max_pods,不碰 idle 池) |
132+| RM | `LUA_DEPLOY_FOLLOWER_GATE` | deploy 锁输家等待室原子准入(ZSET+deadline,≤pc-1;先清过期再 ZADD 先行+超限自退) |
131 133 
132约定:脚本不传 KEYS(键由 `ARGV[1]` 前缀在脚本内拼);调用统一经各自 `state.py` 的 `eval()`。134约定:脚本不传 KEYS(键由 `ARGV[1]` 前缀在脚本内拼);调用统一经各自 `state.py` 的 `eval()`。
133 135 
@@ -143,7 +145,8 @@ session_manager:scope:{sid}:waiters SET 等待队列(LUA_WAITER_
143session_manager:scope:{sid}:free PubSub 额度释放信号145session_manager:scope:{sid}:free PubSub 额度释放信号
144session_manager:pod:{scope}:{pod}:sessions|info SET|HASH per-Pod 会话 / sse_url+deploy_ver146session_manager:pod:{scope}:{pod}:sessions|info SET|HASH per-Pod 会话 / sse_url+deploy_ver
145session_manager:pods:registered SET "{scope}:{pod}"(不变量 5)147session_manager:pods:registered SET "{scope}:{pod}"(不变量 5)
146-resource_manager:resource:scope:{sid}:pods|idle|config|deploying ZSET|SET|HASH|SET148+resource_manager:resource:scope:{sid}:pods|idle|config|deploying|deploy_followers
149+ ZSET|SET|HASH|SET|ZSET(follower 等待室,≤pc-1)
147resource_manager:resource:pod:{pod}:info|idle_since|health_fails HASH|STR|STR150resource_manager:resource:pod:{pod}:info|idle_since|health_fails HASH|STR|STR
148resource_manager:resource:pods:all SET 全部 pod_id(watch/reconcile 枚举)151resource_manager:resource:pods:all SET 全部 pod_id(watch/reconcile 枚举)
149resource_manager:lock:rm:deploy:{sid}|autoscale|reclaim|watch|reconcile 选主/串行化锁152resource_manager:lock:rm:deploy:{sid}|autoscale|reclaim|watch|reconcile 选主/串行化锁
@@ -168,7 +171,7 @@ Facade 间以 Python 异常传播,handler 捕获后映射为错误信封。
168 171 
169```bash172```bash
170cd applications/agent_runtime173cd applications/agent_runtime
171-uv sync --extra local && uv run pytest # 106 用例,fakeredis+SQLite+FakeK8s174+uv sync --extra local && uv run pytest # 114 用例,fakeredis+SQLite+FakeK8s
172./scripts/integration_smoke.sh # 真环境冒烟(场景 A–L;FLUSHDB 目标库,有防误刷)175./scripts/integration_smoke.sh # 真环境冒烟(场景 A–L;FLUSHDB 目标库,有防误刷)
173```176```
174 177 
@@ -276,7 +276,7 @@ flowchart TD
276 A2 -- 已达 --> MX276 A2 -- 已达 --> MX
277```277```
278 278 
279-**读图**:优先取该 scope idle 池里的暖 Pod(`SREM idle` 后返回)→ 无暖 Pod 且该 scope 未达 `max_pods` 则 per-`scope_id` 选主 deploy +1 → `LUA_REGISTER` 登记新 Pod 即返回 → 达 `max_pods`(含 `deploying` 占位)则 `MAX_PODS_REACHED`。deploy 走 **per-`scope_id` 选主串行**(防并发超配;他副本在 deploy 时本请求稍后重试即可复用其成果)。**config-agnostic**:RM 不解析 `pod_spec` 语义;一个 Pod 只服务一个 scope,容量由 SM 的 `SCARD < pod_concurrency` 闸门保证,RM 不做容量叠加判定。279+**读图**:优先取该 scope idle 池里的暖 Pod(`SREM idle` 后返回)→ 无暖 Pod 且该 scope 未达 `max_pods` 则 per-`scope_id` 选主 deploy +1 → `LUA_REGISTER` 登记新 Pod 即返回 → 达 `max_pods`(含 `deploying` 占位)则 `MAX_PODS_REACHED`。deploy 走 **per-`scope_id` 选主串行**(防并发超配;锁输家进 **follower 等待室**——原子准入上限 `pod_concurrency-1`、overflow 快失败、等待有界、leader 的 Pod 注册即直接复用、leader 失败则不接管直接失败)。**config-agnostic**:RM 不解析 `pod_spec` 语义;一个 Pod 只服务一个 scope,容量由 SM 的 `SCARD < pod_concurrency` 闸门保证,RM 不做容量叠加判定。
280 280 
281#### 4.1.4 Resource Manager —— Pod 生命周期与自治回收(后台)281#### 4.1.4 Resource Manager —— Pod 生命周期与自治回收(后台)
282 282 
@@ -377,6 +377,7 @@ flowchart TB
377 SI[("resource:scope:{scope}:idle<br/>SET: idle pod_id<br/>SCARD = min_idle 计数")]:::set377 SI[("resource:scope:{scope}:idle<br/>SET: idle pod_id<br/>SCARD = min_idle 计数")]:::set
378 SCFG[("resource:scope:{scope}:config<br/>HASH: min_idle_pods / max_pods / pod_ttl")]:::hash378 SCFG[("resource:scope:{scope}:config<br/>HASH: min_idle_pods / max_pods / pod_ttl")]:::hash
379 SD[("resource:scope:{scope}:deploying<br/>SET: deploy 占位 token<br/>(计入 max_pods)")]:::str379 SD[("resource:scope:{scope}:deploying<br/>SET: deploy 占位 token<br/>(计入 max_pods)")]:::str
380+ SDF[("resource:scope:{scope}:deploy_followers<br/>ZSET: request_id → deadline<br/>follower 等待室(≤ pc-1)")]:::zset
380 end381 end
381 382 
382 subgraph POD["pod_X"]383 subgraph POD["pod_X"]
@@ -409,6 +410,7 @@ flowchart TB
409| `resource:scope:{scope_id}:idle` | SET | idle 的 pod_id | **SCARD = idle Pod 数**(autoscale / reclaim 闸门;acquire 从此取暖 Pod) |410| `resource:scope:{scope_id}:idle` | SET | idle 的 pod_id | **SCARD = idle Pod 数**(autoscale / reclaim 闸门;acquire 从此取暖 Pod) |
410| `resource:scope:{scope_id}:config` | HASH | min_idle_pods / max_pods / pod_ttl | 首 acquire 存,后续读;config_sync 经 `update_pool_config` **主动刷新**(`autoscale_interval` 全局默认,不入) |411| `resource:scope:{scope_id}:config` | HASH | min_idle_pods / max_pods / pod_ttl | 首 acquire 存,后续读;config_sync 经 `update_pool_config` **主动刷新**(`autoscale_interval` 全局默认,不入) |
411| `resource:scope:{scope_id}:deploying` | SET | deploy 占位 token(uuid) | 计入 `max_pods` 判定(防并发 deploy 超配);register / 失败时清 |412| `resource:scope:{scope_id}:deploying` | SET | deploy 占位 token(uuid) | 计入 `max_pods` 判定(防并发 deploy 超配);register / 失败时清 |
413+| `resource:scope:{scope_id}:deploy_followers` | ZSET | deploy 锁输家的 request_id → deadline | follower 等待室(M8):`ZCARD ≤ pod_concurrency-1`;检测 leader Pod 注册即复用;score=deadline 供闸门清崩溃遗留 |
412| `resource:pod:{pod_id}:info` | HASH | scope_id / pod_sse_url / pod_ip / namespace / phase / created_ts / deploy_ver | Pod 元信息;`scope_id` 标识所属池;`deploy_ver` 供 acquire 版本过滤(只发当前版本暖 Pod) |414| `resource:pod:{pod_id}:info` | HASH | scope_id / pod_sse_url / pod_ip / namespace / phase / created_ts / deploy_ver | Pod 元信息;`scope_id` 标识所属池;`deploy_ver` 供 acquire 版本过滤(只发当前版本暖 Pod) |
413| `resource:pod:{pod_id}:idle_since` | STRING | idle 起始时间戳 | reclaim 计时(aged ≥ `pod_ttl` 回收) |415| `resource:pod:{pod_id}:idle_since` | STRING | idle 起始时间戳 | reclaim 计时(aged ≥ `pod_ttl` 回收) |
414| `resource:pods:all` | SET | 全部 pod_id | 孤儿对账 / 枚举 |416| `resource:pods:all` | SET | 全部 pod_id | 孤儿对账 / 枚举 |
@@ -123,10 +123,10 @@
123### 2.1 `ResourceManagerFacade.acquire(...)` —— 请求 Pod(取 scope 暖 Pod 或 deploy +1)123### 2.1 `ResourceManagerFacade.acquire(...)` —— 请求 Pod(取 scope 暖 Pod 或 deploy +1)
124- **in**:`{ scope_id:str, pod_spec:dict, pool_config:dict, request_id:str }`124- **in**:`{ scope_id:str, pod_spec:dict, pool_config:dict, request_id:str }`
125 - `pod_spec` = deploy 字段子集(`agent_image`/`namespace`/`container_name`/`kubeconfig`/`readiness_*`/`nfs_*`/资源限额)+ **SSE 端口/路径**(⚠️ §13.1)。125 - `pod_spec` = deploy 字段子集(`agent_image`/`namespace`/`container_name`/`kubeconfig`/`readiness_*`/`nfs_*`/资源限额)+ **SSE 端口/路径**(⚠️ §13.1)。
126- - `pool_config` = per-scope 池参数(`min_idle_pods`/`max_pods`(SM 派生)/`pod_ttl`),RM 首 acquire 缓存为 `scope:config`(`pod_concurrency` 不入 RM;SM 自用作 per-Pod 容量闸门)。126+ - `pool_config` = per-scope 池参数(`min_idle_pods`/`max_pods`(SM 派生)/`pod_ttl`/`pod_concurrency`),RM 首 acquire 缓存为 `scope:config`。`pod_concurrency` **仅用于 deploy follower 等待室推导上限(pc-1)**——per-Pod 容量闸门仍在 SM 侧(`SCARD < pod_concurrency`),RM 不做容量叠加判定(红线不变)。
127 - **out**:`{ pod_id:str, pod_sse_url:str }`127 - **out**:`{ pod_id:str, pod_sse_url:str }`
128 - **错**(Facade 抛异常):`MaxPodsReached`(对应原 `MAX_PODS_REACHED`)、`DeployFailed`(对应原 `DEPLOY_FAILED`)、`ValidationError`(对应原 `VALIDATION`)。错误码语义不变,见 §7;SM 的 `route` handler 捕获后映射为自身对外 HTTP 响应。128 - **错**(Facade 抛异常):`MaxPodsReached`(对应原 `MAX_PODS_REACHED`)、`DeployFailed`(对应原 `DEPLOY_FAILED`)、`ValidationError`(对应原 `VALIDATION`)。错误码语义不变,见 §7;SM 的 `route` handler 捕获后映射为自身对外 HTTP 响应。
129- - **语义**:按 `scope_id` 取 Pod——若 `scope:idle` 有暖 Pod 则复用,否则未达 `max_pods` → 选主 deploy +1,达 `max_pods` → `MaxPodsReached`。从 `scope:idle` 取暖 Pod 时**跳过 `deploy_ver` 不匹配当前 deploy 字段的**(A 类配置变更后的老版本暖 Pod 留在 idle 池按 `pod_ttl` 回收,不外发给新流量;见 HLD 场景 M / §2.2.1)。**config-agnostic**:不解析 `pod_spec` 语义,池参数经 `acquire` 传入并缓存于 `scope:config`。129+ - **语义**:按 `scope_id` 取 Pod——若 `scope:idle` 有暖 Pod 则复用,否则未达 `max_pods` → 选主 deploy +1,达 `max_pods` → `MaxPodsReached`。deploy 锁的**输家进 follower 等待室**(M8):原子闸门准入上限 `pod_concurrency-1`(leader 会话之外新 Pod 恰剩这些槽),overflow 严格快失败;等待有界(`ready_timeout`+余量),leader 失败(锁空闲且无进展)则 follower 直接失败**不接管**;检测到 leader 的 Pod 注册即**直接复用返回**(与 reuse 分支同构,SM 侧重跑仲裁)——RM 全程不读 SM 容量键。从 `scope:idle` 取暖 Pod 时**跳过 `deploy_ver` 不匹配当前 deploy 字段的**(A 类配置变更后的老版本暖 Pod 留在 idle 池按 `pod_ttl` 回收,不外发给新流量;见 HLD 场景 M / §2.2.1)。**config-agnostic**:不解析 `pod_spec` 语义,池参数经 `acquire` 传入并缓存于 `scope:config`。
130 130 
131### 2.2 `ResourceManagerFacade.idle_consider(...)` —— 该 scope 在该 Pod 上已无会话(⚠️ scope 级,修订)131### 2.2 `ResourceManagerFacade.idle_consider(...)` —— 该 scope 在该 Pod 上已无会话(⚠️ scope 级,修订)
132- **in**:`{ pod_id:str, scope_id:str }`132- **in**:`{ pod_id:str, scope_id:str }`
@@ -189,6 +189,7 @@
189| `resource:scope:{scope_id}:idle` | LUA_RELEASE `SADD` / LUA_REGISTER(idle_flag=true)`SADD` | autoscale / reclaim `SCARD` | reuse `SREM`(LUA_ACQUIRE)/ LUA_PURGE | 无 |189| `resource:scope:{scope_id}:idle` | LUA_RELEASE `SADD` / LUA_REGISTER(idle_flag=true)`SADD` | autoscale / reclaim `SCARD` | reuse `SREM`(LUA_ACQUIRE)/ LUA_PURGE | 无 |
190| `resource:scope:{scope_id}:config` | 首 acquire `HSET` | LUA_ACQUIRE / autoscale / reclaim | 无(随 scope 生命周期,长期) | 无 |190| `resource:scope:{scope_id}:config` | 首 acquire `HSET` | LUA_ACQUIRE / autoscale / reclaim | 无(随 scope 生命周期,长期) | 无 |
191| `resource:scope:{scope_id}:deploying` | LUA_ACQUIRE(need_deploy)`SADD token` | LUA_ACQUIRE max_pods 判定 `SCARD` | LUA_REGISTER / deploy 失败 `SREM token` | 无(token 短命) |191| `resource:scope:{scope_id}:deploying` | LUA_ACQUIRE(need_deploy)`SADD token` | LUA_ACQUIRE max_pods 判定 `SCARD` | LUA_REGISTER / deploy 失败 `SREM token` | 无(token 短命) |
192+| `resource:scope:{scope_id}:deploy_followers` | LUA_DEPLOY_FOLLOWER_GATE `ZADD request_id→deadline` | follower 闸门 `ZCARD ≤ pc-1` | follower 退出 `ZREM` / 闸门 `ZREMRANGEBYSCORE`(过期兜底) | 无(成员短命) |
192| `resource:pods:all` | LUA_REGISTER `SADD` | 对账 / RM 冷启动 | LUA_PURGE `SREM` | 无 |193| `resource:pods:all` | LUA_REGISTER `SADD` | 对账 / RM 冷启动 | LUA_PURGE `SREM` | 无 |
193| `lock:rm:*` | sweeper `SET NX EX` | sweeper 判定 | 自动过期 / 下 tick 覆盖 | 见 §5.4 |194| `lock:rm:*` | sweeper `SET NX EX` | sweeper 判定 | 自动过期 / 下 tick 覆盖 | 见 §5.4 |
194 195 
@@ -212,7 +213,18 @@ Resource Manager **不查 config DB、不解析 template 语义**。所有部署
212 213 
213### 5.1 Lua 脚本(承担所有编排态变更,原子)214### 5.1 Lua 脚本(承担所有编排态变更,原子)
214 215 
215-> 脚本全集(5 个,实现与本文对齐):`LUA_ACQUIRE` / `LUA_REGISTER` / `LUA_RELEASE` / `LUA_PURGE`(核心 4 个,下文全文)+ `LUA_PLACEHOLDER`(实现期从 ACQUIRE 的占位逻辑抽出,供 autoscale 热备 deploy 专用——同样计入 max_pods 但**不碰 idle 池**,补位不该消耗既有暖 Pod)。216+> 脚本全集(6 个,实现与本文对齐):`LUA_ACQUIRE` / `LUA_REGISTER` / `LUA_RELEASE` / `LUA_PURGE`(核心 4 个,下文全文)+ `LUA_PLACEHOLDER`(实现期从 ACQUIRE 的占位逻辑抽出,供 autoscale 热备 deploy 专用——同样计入 max_pods 但**不碰 idle 池**,补位不该消耗既有暖 Pod)+ `LUA_DEPLOY_FOLLOWER_GATE`(M8:deploy 锁输家的等待室原子准入,ZSET+deadline,上限 `pod_concurrency-1`;先清过期成员再 ZADD 先行+超限自退——纪律同 SM 的 `LUA_WAITER_GATE`)。
217+>
218+> `LUA_DEPLOY_FOLLOWER_GATE(scope_id, follower_id, max_followers, deadline, now)` 全文:
219+> ```
220+> key = resource:scope:{scope_id}:deploy_followers
221+> ZREMRANGEBYSCORE(key, -inf, now) # 清过期成员(等待进程崩溃兜底,不泄漏)
222+> ZADD(key, deadline, follower_id) # 先行
223+> if ZCARD(key) > max_followers:
224+> ZREM(key, follower_id) # 超限自退
225+> return {false}
226+> return {true}
227+> ```
216 228 
217**`LUA_ACQUIRE(scope_id, deploy_token, now)`** —— acquire 的原子核心:取 scope 暖 Pod 复用,或判定 need_deploy/max_reached(占位)。中途无并发插入(无 race)。229**`LUA_ACQUIRE(scope_id, deploy_token, now)`** —— acquire 的原子核心:取 scope 暖 Pod 复用,或判定 need_deploy/max_reached(占位)。中途无并发插入(无 race)。
218```230```
@@ -297,7 +309,13 @@ loop:
297 # 选主串行 deploy:抢 lock:rm:deploy:{scope_id}(或经 autoscale 同一选主通道),309 # 选主串行 deploy:抢 lock:rm:deploy:{scope_id}(或经 autoscale 同一选主通道),
298 # 保证同 scope 不会并发 deploy 超配;deploy 是重操作(K8s create+wait Ready),串行可接受。310 # 保证同 scope 不会并发 deploy 超配;deploy 是重操作(K8s create+wait Ready),串行可接受。
299 if not acquire_lock("lock:rm:deploy:" + scope_id, ex=deploy_timeout):311 if not acquire_lock("lock:rm:deploy:" + scope_id, ex=deploy_timeout):
300- SREM deploying token; await short_sleep; continue # 他副本在 deploy,稍后重跑 ACQUIRE 复用其成果312+ # 输家 → follower 等待室(M8):复用 leader 的 Pod,不再自建第 2 个空 Pod
313+ SREM deploying token
314+ out = await follow_leader(scope_id, pod_spec, pool_config, request_id)
315+ idempotency.put(request_id, out); return out
316+ # follow_leader 内部(见下):闸门准入(≤pc-1,overflow 快失败)→
317+ # 轮询:新 Pod 注册 → 返回该 Pod;锁空闲无进展(leader 失败)→ 失败不接管;
318+ # deadline(ready_timeout+余量)→ 失败;finally 必 ZREM 出室
301 try:319 try:
302 info = await k8s.deploy(pod_spec) # create + _wait_running_ready → pod_ip320 info = await k8s.deploy(pod_spec) # create + _wait_running_ready → pod_ip
303 pod_id = info.pod_name321 pod_id = info.pod_name
@@ -311,7 +329,7 @@ loop:
311 release_lock("lock:rm:deploy:" + scope_id)329 release_lock("lock:rm:deploy:" + scope_id)
312 continue # 重跑 ACQUIRE:新 Pod 必被取作暖 Pod 复用330 continue # 重跑 ACQUIRE:新 Pod 必被取作暖 Pod 复用
313```331```
314-**并发安全**:`LUA_ACQUIRE` 单脚本原子(取暖 Pod + 移出 idle 不可分);`need_deploy` 走 `lock:rm:deploy:{scope_id}` 选主串行,避免并发 deploy 超配(SM spec §13.6 的"并发 acquire 超配"在此靠选主根治)。多个 `scope_id` 的 deploy 互不阻塞(按 scope 分锁)。332+**并发安全**:`LUA_ACQUIRE` 单脚本原子(取暖 Pod + 移出 idle 不可分);`need_deploy` 走 `lock:rm:deploy:{scope_id}` 选主串行,避免并发 deploy 超配(SM spec §13.6 的"并发 acquire 超配"在此靠选主根治)。多个 `scope_id` 的 deploy 互不阻塞(按 scope 分锁)。锁输家经 follower 等待室复用 leader 成果——准入原子(`LUA_DEPLOY_FOLLOWER_GATE`,上限 pc-1)、等待有界、退出必清成员;跨副本冷竞争不再产生自建空 Pod。
315 333 
316### 5.3 idle_consider(ResourceManagerFacade.idle_consider 方法)334### 5.3 idle_consider(ResourceManagerFacade.idle_consider 方法)
317```335```
@@ -11,7 +11,7 @@
11 11 
12| 层 | 入口 | 规模 | 依赖环境 | 退出码 |12| 层 | 入口 | 规模 | 依赖环境 | 退出码 |
13|---|---|---|---|---|13|---|---|---|---|---|
14-| 进程内双实例 | `uv run pytest tests/integration/test_multi_replica.py` | 12 用例 | 无(离线,fakeredis) | pytest 标准 |14+| 进程内双实例 | `uv run pytest tests/integration/test_multi_replica.py` | 14 用例 | 无(离线,fakeredis) | pytest 标准 |
15| 集成冒烟(M6) | `./scripts/integration_smoke.sh` | 65 项断言 | 单实例 server 模式 + 真 Redis/MySQL/K8s | 0/1/2 |15| 集成冒烟(M6) | `./scripts/integration_smoke.sh` | 65 项断言 | 单实例 server 模式 + 真 Redis/MySQL/K8s | 0/1/2 |
16| 多副本 e2e(M7) | `uv run --no-sync python scripts/e2e_multi_replica.py` | 35 项断言 | K8s 多副本 + Service LB + 真 Redis | 0/1/2 |16| 多副本 e2e(M7) | `uv run --no-sync python scripts/e2e_multi_replica.py` | 35 项断言 | K8s 多副本 + Service LB + 真 Redis | 0/1/2 |
17| 压测/浸泡 | `uv run --no-sync python scripts/load_test.py` | 3 场景 | 任意入口(建议 LB) | 0/1 |17| 压测/浸泡 | `uv run --no-sync python scripts/load_test.py` | 3 场景 | 任意入口(建议 LB) | 0/1 |
@@ -288,7 +288,7 @@ config_sync **串行**;起点 FLUSHDB + TRUNCATE + 删残留 Pod。
288 288 
289---289---
290 290 
291-## 5. 进程内双实例:`tests/integration/test_multi_replica.py`(12 用例)291+## 5. 进程内双实例:`tests/integration/test_multi_replica.py`(14 用例)
292 292 
293同进程两个完整 App(各自 SystemContext + 5 个后台 Job)共享一组293同进程两个完整 App(各自 SystemContext + 5 个后台 Job)共享一组
294fakeredis/SQLite/FakeK8s,`instance_id` 显式 `replica-a`/`replica-b`,294fakeredis/SQLite/FakeK8s,`instance_id` 显式 `replica-a`/`replica-b`,
@@ -300,7 +300,7 @@ httpx ASGITransport 单事件循环并发驱动。**输入全部走完整 HTTP**
300| 1 | 身份与共享态 | route 经 A,touch 经 B | instance_id 互异且 RM 镜像;B touch 到 A 建的会话 `touched=true` |300| 1 | 身份与共享态 | route 经 A,touch 经 B | instance_id 互异且 RM 镜像;B touch 到 A 建的会话 `touched=true` |
301| 2 | 交替亲和 | 同 session A→B→A→B route | 恒同 Pod;SCARD=1 |301| 2 | 交替亲和 | 同 session A→B→A→B route | 恒同 Pod;SCARD=1 |
302| 3 | 跨副本突发不超收 | cc=2/pc=1 占满后 8 并发交替 A/B | 0×200;4×503 队列满 + 4×504 超时;终态 SCARD=2、waiters=0、deploying=0 |302| 3 | 跨副本突发不超收 | cc=2/pc=1 占满后 8 并发交替 A/B | 0×200;4×503 队列满 + 4×504 超时;终态 SCARD=2、waiters=0、deploying=0 |
303-| 4 | deploy 锁串行化 | SlowFakeK8s(deploy 0.4s),A/B 并发冷启动 + 追加 s3 | 部署窗口**零重叠**;Pod ≤ max_pods=2;占位清空;s3 first-fit 复用 |303+| 4 | deploy 锁串行化 + follower 复用 | SlowFakeK8s(deploy 0.4s),A/B 并发冷启动 + 追加 s3 | 并发对**恰好 1 次部署**(输家进等待室复用同 Pod);pod1 满后 s3 才第 2 次部署;窗口零重叠;占位/等待室清空 |
304| 5 | 输家复用暖 Pod | 手持 deploy 锁 + 后台注册 idle Pod 后释放;A route | 返回他副本 Pod;本侧零部署;占位清空 |304| 5 | 输家复用暖 Pod | 手持 deploy 锁 + 后台注册 idle Pod 后释放;A route | 返回他副本 Pod;本侧零部署;占位清空 |
305| 6 | 跨副本唤醒 | A 占满→A 排队→回拨过期→**B** touch | B 的 touch 返回 `touched=false`(惰性驱逐);A 的等待者 <2s 被唤醒并占释放额度 |305| 6 | 跨副本唤醒 | A 占满→A 排队→回拨过期→**B** touch | B 的 touch 返回 `touched=false`(惰性驱逐);A 的等待者 <2s 被唤醒并占释放额度 |
306| 7 | 幂等跨副本 | 同 request_id A 首发、B 重放 | 响应一致;仅一会话 |306| 7 | 幂等跨副本 | 同 request_id A 首发、B 重放 | 响应一致;仅一会话 |
@@ -309,6 +309,8 @@ httpx ASGITransport 单事件循环并发驱动。**输入全部走完整 HTTP**
309| 10 | sweeper 互斥 | 手持 lock:sweep 后 A sweep_once | 直退不误扫;锁释放后补扫完成 |309| 10 | sweeper 互斥 | 手持 lock:sweep 后 A sweep_once | 直退不误扫;锁释放后补扫完成 |
310| 11 | 并发收敛 | A/B sweep_once 并发 gather | 无异常;全部老化;锁正常释放;`pods:registered` 不变 |310| 11 | 并发收敛 | A/B sweep_once 并发 gather | 无异常;全部老化;锁正常释放;`pods:registered` 不变 |
311| 12 | /healthz | 分别 GET 两 App | 200 + 各自 instance_id |311| 12 | /healthz | 分别 GET 两 App | 200 + 各自 instance_id |
312+| 13 | follower 上限严格快失败 | cc=8/pc=2,4 并发冷启动(deploy 0.4s) | 2×200(同 Pod)+ 2×503 NO_POD_AVAILABLE(闸门拒);恰好 1 次部署;占位/等待室清空 |
313+| 14 | leader 失败 follower 不接管 | deploy 慢速失败(0.5s 后抛),2 并发 | 双 503 NO_POD_AVAILABLE;占位/等待室全清 |
312 314 
313---315---
314 316 
@@ -405,10 +407,11 @@ uv run --no-sync python scripts/load_test.py \
4051. **多副本冷突发**:并发冷启动时占位先封顶 `max_pods`,多余请求立即 5034071. **多副本冷突发**:并发冷启动时占位先封顶 `max_pods`,多余请求立即 503
406 `NO_POD_AVAILABLE`(retry_after=1)而非排队——多副本 e2e 的 S4 已按此语义预填;408 `NO_POD_AVAILABLE`(retry_after=1)而非排队——多副本 e2e 的 S4 已按此语义预填;
407 压测 queued 场景冷启动期同样可见。不阻塞 M6 冒烟(部署参数对齐后经 LB 65/65)。409 压测 queued 场景冷启动期同样可见。不阻塞 M6 冒烟(部署参数对齐后经 LB 65/65)。
408-2. **跨副本冷竞争双 Pod**:deploy 锁输家 retry 时赢家 Pod 为 in-use 不在 idle 池,410+2. ~~跨副本冷竞争双 Pod~~(**M8 已解决**):deploy 锁输家原会自建第 2 个空 Pod;
409- 输家自建第 2 个 Pod(`max_pods` 内,空 Pod 经 empty-pod pass→idle_consider→reclaim411+ 现进 **follower 等待室**——原子准入上限 `pod_concurrency-1`(overflow 严格快失败)、
410- 自愈)。测试断言「窗口零重叠 + Pod ≤ max_pods」,不断言「恰好 1 个 Pod」。412+ 等待有界(ready_timeout+余量)、leader 的 Pod 注册即直接复用、leader 失败不接管
411- 收紧需把 SM 容量感知下沉 RM(开放问题,见 handoff §十一.1)。413+ 直接失败(`LUA_DEPLOY_FOLLOWER_GATE`,ZSET+deadline 防崩溃泄漏)。
414+ 双实例用例 4 已收紧断言「冷竞争恰好 1 次部署」;实测冷启动尾延迟 30.5s→10.2s。
4123. **场景 N(半死探测)**:待 AgentServer 原生支持 `GET /health` 后补端到端4153. **场景 N(半死探测)**:待 AgentServer 原生支持 `GET /health` 后补端到端
413 (机制已有单测)。416 (机制已有单测)。
4144. **config_sync 串行**:全局锁,任何脚本/客户端并发下发即 409 `CONFIG_SYNC_BUSY`。4174. **config_sync 串行**:全局锁,任何脚本/客户端并发下发即 409 `CONFIG_SYNC_BUSY`。