1717上限(CodeArts 每账号 3),无上限放行会把第 4 个起全部打成 400 并发超限。
18180/None 表示不限(保持旧行为)。
1919
20+ `window_seconds > 0` 时上限改用**滑动窗口**口径:同一个桶在最近
21+ `window_seconds` 秒内最多放行 `max_concurrency` 次请求启动。这是对齐
22+ CodeArts 实测行为的修正——它限制的不是「同时在途」,而是「每账号每约 60s
23+ 最多 3 个会话」:会话在 HTTP 流结束后仍滞留数十秒(130s 打满 3 并发后,
24+ 单请求直到约 68s 才恢复)。纯在途上限挡不住「3 个并发刚结束就立刻再发 3 个」
25+ 这类突发。窗口模式下 `release` 不再让出窗口配额(配额按启动时刻自然滑出),
26+ 但仍用于计数在途、维持最小间隔。
27+
2028**`max_concurrency > 0` 时 `release` 必须与 `wait_turn` 严格配对。**
2129provider 在 `stream_chat` 的 `finally` 里归还名额,而 `async for ... break`
2230不会关闭 async generator(CPython 只在耗尽 / 显式 `aclose()` / GC 的
@@ -55,6 +63,7 @@ class Pacer:
5563 def __init__ (self , min_seconds : float , max_seconds : float , * ,
5664 allow_concurrent : bool = False ,
5765 max_concurrency : int | Callable [[], int ] = 0 ,
66+ window_seconds : float | Callable [[], float ] = 0 ,
5867 sleep : Callable [[float ], Awaitable [None ]] | None = None ,
5968 now : Callable [[], float ] | None = None ) -> None :
6069 # 上下限可热更(B3.2):存取值器,每次计间隔读当前值。校验必须放在
@@ -73,12 +82,19 @@ def __init__(self, min_seconds: float, max_seconds: float, *,
7382 if not callable (max_concurrency ) and max_concurrency < 0 :
7483 raise ValueError ("pacer max_concurrency must be non-negative" )
7584 self ._max_concurrency = live (max_concurrency )
85+ # 滑动窗口长度(秒,0 = 关闭):与 max_concurrency 组合成「窗口内启动
86+ # 次数上限」。同样可热更。
87+ if not callable (window_seconds ) and window_seconds < 0 :
88+ raise ValueError ("pacer window_seconds must be non-negative" )
89+ self ._window = live (window_seconds )
7690 self ._lock = asyncio .Lock ()
7791 # 严格模式:全局单桶
7892 self ._last_started : float | None = None
7993 # 并发模式:桶 → 上次(预留)开始时刻 / 在途计数
8094 self ._bucket_started : dict [str , float ] = {}
8195 self ._inflight : dict [str , int ] = {}
96+ # 窗口模式:桶 → 落在当前窗口内的启动时刻(升序,用于滑动窗口计数)
97+ self ._starts : dict [str , list [float ]] = {}
8298 # 并发模式 + 上限:名额让出时置位,挂起的请求醒来重查是否有空位
8399 self ._wake = asyncio .Event ()
84100
@@ -112,6 +128,12 @@ def max_concurrency(self) -> int:
112128 limit = int (self ._max_concurrency ())
113129 return limit if limit > 0 else 0
114130
131+ @property
132+ def window_seconds (self ) -> float :
133+ """当前滑动窗口长度;0/负 = 关闭窗口口径(每次现读,支持热更)。"""
134+ value = float (self ._window ())
135+ return value if value > 0 else 0.0
136+
115137 def next_interval (self ) -> float :
116138 low , high = self ._validate ()
117139 if low == 0 and high == 0 :
@@ -124,6 +146,22 @@ def _capped(self) -> bool:
124146 """并发模式且配了在途上限——只有这种组合需要排队等名额。"""
125147 return self ._allow_concurrent and self .max_concurrency > 0
126148
149+ def _windowed (self ) -> bool :
150+ """并发 + 在途上限 + 窗口三者齐备时才走滑动窗口口径。"""
151+ return (self ._allow_concurrent and self .max_concurrency > 0
152+ and self .window_seconds > 0 )
153+
154+ def _prune_starts (self , bucket : str , now : float ) -> None :
155+ """**调用方须持 `_lock`**:丢掉已经滑出窗口的启动时刻。"""
156+ starts = self ._starts .get (bucket )
157+ if not starts :
158+ return
159+ cutoff = now - self .window_seconds
160+ while starts and starts [0 ] <= cutoff :
161+ starts .pop (0 )
162+ if not starts :
163+ del self ._starts [bucket ]
164+
127165 def _reserve (self , bucket : str ) -> float :
128166 """**调用方须持 `_lock`**:登记一次在途并返回本次需补足的间隔秒数。
129167
@@ -151,7 +189,8 @@ async def wait_turn(self, key: str | None = None) -> None:
151189 并发模式下返回即视为「已占用一个在途名额」,调用方必须在请求结束后
152190 用同一 key 调 `release` 归还,否则该桶会被当成永远有请求在途而失去
153191 节流(见模块 docstring)。配了 `max_concurrency` 时,名额满会让新
154- 请求在此挂起,直到有 `release` 让位。
192+ 请求在此挂起,直到有 `release` 让位;再配了 `window_seconds` 时按
193+ 滑动窗口计数,窗口内启动次数满则挂起到最早一次启动滑出窗口。
155194 """
156195 if self .disabled and not self ._capped ():
157196 return
@@ -168,28 +207,57 @@ async def wait_turn(self, key: str | None = None) -> None:
168207 bucket = key or ""
169208 remaining = 0.0
170209 reserved = False
210+ windowed_start : float | None = None
171211 try :
172212 while True :
213+ window_wait = 0.0
173214 async with self ._lock :
215+ now = self ._now ()
174216 limit = self .max_concurrency
175- if not limit or self ._inflight .get (bucket , 0 ) < limit :
217+ if self ._windowed ():
218+ self ._prune_starts (bucket , now )
219+ starts = self ._starts .get (bucket , [])
220+ if limit and len (starts ) >= limit :
221+ # 窗口未过期:等到最早一次启动滑出窗口再重查。
222+ window_wait = starts [0 ] + self .window_seconds - now
223+ self ._wake .clear ()
224+ else :
225+ self ._starts .setdefault (bucket , []).append (now )
226+ windowed_start = now
227+ remaining = self ._reserve (bucket )
228+ reserved = True
229+ elif not limit or self ._inflight .get (bucket , 0 ) < limit :
176230 remaining = self ._reserve (bucket )
177231 reserved = True
232+ else :
233+ # 名额已满:清事件后到锁外等待,release 会 set 唤醒。
234+ self ._wake .clear ()
235+ if reserved :
178236 break
179- # 名额已满:清事件后到锁外等待,release 会 set 唤醒。
180- self ._wake .clear ()
181- await self ._wake .wait ()
237+ if window_wait > 0 :
238+ await self ._sleep (window_wait )
239+ else :
240+ await self ._wake .wait ()
182241 if remaining > 0 :
183242 await self ._sleep (remaining )
184243 except BaseException :
185244 # 未成功「占用」就退出(取消 / 间隔校验异常):把名额还回去,
186- # 否则该桶的名额会永久少一个(最终把渠道卡死)。
245+ # 否则该桶的名额会永久少一个(最终把渠道卡死)。窗口模式下还要
246+ # 抹掉刚登记、并未真正发起的启动时刻,否则白占一个窗口配额。
187247 if reserved :
188248 self .release (bucket )
249+ if windowed_start is not None :
250+ starts = self ._starts .get (bucket )
251+ if starts and windowed_start in starts :
252+ starts .remove (windowed_start )
189253 raise
190254
191255 def release (self , key : str | None = None ) -> None :
192- """并发模式归还一个在途名额;严格模式与多余的 release 都是空操作。"""
256+ """并发模式归还一个在途名额;严格模式与多余的 release 都是空操作。
257+
258+ 窗口模式下只减少在途计数、维持最小间隔;窗口配额由启动时刻自然滑出,
259+ 不因 `release` 提前让出(否则又退回纯在途口径,挡不住突发)。
260+ """
193261 if not self ._allow_concurrent :
194262 return
195263 bucket = key or ""
0 commit comments