重构: F-09 batch5 llm_concurrency会话级改造(决策c)
- state.rs: per_conv HashMap<conv_id,Semaphore>(每对话permits=2) + per_conv_permits热改 + acquire_per_conv(conv_id) + release_conv - agentic.rs: loop入口acquire_global+acquire_per_conv(conv_id)整loop持有(会话级并发上限3) + 删每轮acquire/drop - compress/title/knowledge_inject: acquire_per_conv(conv_id)+签名conv_id - project.rs: 扫描合成key + ai_conversation_delete release_conv global: LLM调用并发→并发会话数上限;F-260616-12 retry permit自洽 主代兜底: cargo check --workspace 0 + test 96 + grep印证
This commit is contained in:
@@ -509,6 +509,17 @@ pub(crate) async fn run_agentic_loop(
|
||||
// 其 token 估算在 loop 外算一次缓存复用,避免每轮/每次重试重复 estimate_text(低收益优化,行为不变)。
|
||||
let sys_tokens = TokenEstimator::default().estimate_text(&system_prompt);
|
||||
|
||||
// F-260616-09 B 批5 / 决策 c-1:global 改「并发会话数上限」,每 loop 入口拿 1 permit 持整个 loop
|
||||
// 生命周期(含工具执行/审批等待/重试)。同一时刻最多 N 个对话并发跑 loop(N=global permits,默认 3),
|
||||
// 第 N+1 个对话的 acquire_global await 阻塞排队。permit 绑 guard(函数返回)Drop 自动释放——
|
||||
// 各 return 点(save_conversation/stop/conv 删除/收敛/达 MAX)退出即释放槽位。
|
||||
// per_conv 同样整 loop 持有:同一对话内主循环 stream_llm + 后台标题/压缩/提炼共享该 conv 的 permits=2,
|
||||
// 防单对话内并发 LLM 调用失控。**F-260616-12 核验**:retry 同 loop 内,持 per_conv 合理;
|
||||
// global 是会话级,retry 不再阻塞他对话(原每轮 acquire/drop 语义下 global 短暂释放,
|
||||
// 现整 loop 持有更贴合"会话级并发上限"语义)。
|
||||
let _conv_global_permit = llm_concurrency.acquire_global().await;
|
||||
let _conv_per_conv_permit = llm_concurrency.acquire_per_conv(&conv_id).await;
|
||||
|
||||
for iteration in start_iteration..max_iterations {
|
||||
// 用户请求停止 → 收尾退出(已生成文本已在上一轮入库)
|
||||
if stop_flag.load(Ordering::SeqCst) {
|
||||
@@ -640,6 +651,7 @@ pub(crate) async fn run_agentic_loop(
|
||||
&provider_config,
|
||||
active_msgs,
|
||||
&lang,
|
||||
&conv_id,
|
||||
&llm_concurrency,
|
||||
).await.map(Some)
|
||||
};
|
||||
@@ -720,20 +732,19 @@ pub(crate) async fn run_agentic_loop(
|
||||
messages.iter().map(|m| est.estimate_message(m)).sum()
|
||||
};
|
||||
|
||||
// LLM 并发限流(全局 + 单对话双层),仅覆盖 stream_llm 调用本身;
|
||||
// 工具执行(process_tool_calls)是本地操作无 RPM 成本,permit 在 stream 后立即释放避免占槽
|
||||
// LLM 并发限流(global + per_conv 双层)— F-260616-09 B 批5 已上移到 loop 入口
|
||||
// (L520-521 _conv_global_permit/_conv_per_conv_permit),整 loop 持有(含工具执行/审批等待/重试)。
|
||||
// 原每轮 acquire + stream 后 drop(L747-748 acquire / L914-915 drop)已移除——会话级并发语义下,
|
||||
// permit 应绑 loop 生命周期而非单次 stream。工具执行期间占槽是决策 c 的有意行为
|
||||
// (会话级并发上限含全部 LLM 相关工作,非仅 stream 调用)。
|
||||
//
|
||||
// CR-30-1 / F-260616-07 / 决策 a1: 流前失败(Init Err)重试,流中途失败(MidStream
|
||||
// Partial)不重试保文。重试退避复用 retry::backoff_delay(1s→2s→4s±20% jitter) +
|
||||
// retry::is_status_retryable Fatal 分类(stream_recv classify_status_or_class 镜像,
|
||||
// 4xx 非429 立即放弃) + 30s 总挂钟预算。重试期间持有 permit 不释放(防新请求挤占)。
|
||||
//
|
||||
// F-260614-04 / F-260614-04b: global/per_conv permit 整轮迭代共享(不随 candidate 切换
|
||||
// 释放——并发语义是「全局/单对话」级,与具体 provider 无关);per-provider permit 改为
|
||||
// 在 candidate 循环内取(切换 candidate 时释放旧取新,避免占用未用 provider 的槽)。
|
||||
// 两 permit 均 Drop 释放(L803-804 显式 drop _global/_per_conv)。
|
||||
let _global_permit = llm_concurrency.acquire_global().await;
|
||||
let _per_conv_permit = llm_concurrency.acquire_per_conv().await;
|
||||
// F-260614-04 / F-260614-04b: per-provider permit 仍在 candidate 循环内取
|
||||
// (切换 candidate 时释放旧取新,避免占用未用 provider 的槽);global/per_conv 由 loop 入口持有。
|
||||
|
||||
// 重试总预算(挂钟,含 sleep + 各次请求耗时),对齐 retry::MAX_TOTAL_BUDGET 30s。
|
||||
// F-260614-04b:本轮各 candidate 共享一个 30s 预算(切换 provider 不重置预算,
|
||||
@@ -896,11 +907,9 @@ pub(crate) async fn run_agentic_loop(
|
||||
});
|
||||
return;
|
||||
}
|
||||
// stream 结束立即释放 permit,后续工具执行不受限流(本地操作无 RPM 成本)
|
||||
// per-provider permit(_provider_permit)在 candidate 循环内随作用域 Drop 自动释放,
|
||||
// 此处仅显式释放 global/per_conv(F-04b:per-provider 改为循环内取)。
|
||||
drop(_global_permit);
|
||||
drop(_per_conv_permit);
|
||||
// F-260616-09 B 批5: global/per_conv permit 已上移 loop 入口(L520-521)整 loop 持有,
|
||||
// stream 后不再立即释放(会话级并发语义:工具执行期间也占槽)。per-provider permit
|
||||
// (_provider_permit)仍在 candidate 循环内随作用域 Drop 自动释放(F-04b 语义不变)。
|
||||
|
||||
// 累加本轮 token:provider 流式 usage 的 prompt_tokens 为 0 时(GLM 等),用预估输入兜底
|
||||
let round_prompt = if round_usage.prompt_tokens == 0 { estimated_prompt } else { round_usage.prompt_tokens };
|
||||
|
||||
@@ -692,6 +692,7 @@ pub async fn ai_chat_compress_context(
|
||||
&provider_config,
|
||||
active_msgs,
|
||||
&lang,
|
||||
&conv_id,
|
||||
&llm_concurrency,
|
||||
).await {
|
||||
Ok(s) => s,
|
||||
@@ -1727,6 +1728,13 @@ pub async fn ai_conversation_delete(
|
||||
session.active_conversation_id = None;
|
||||
session.messages.clear();
|
||||
}
|
||||
// F-260616-09 B 批5:同时清理 LlmConcurrency 的 per_conv Semaphore 条目(防 HashMap 无限增长)。
|
||||
// conv 已删=LlmConcurrency 该条目不再被 acquire(无 conv 则无 loop/标题/提炼/压缩针对它)。
|
||||
// 已持 permit 不受影响(permit 绑旧 Arc,随 Drop 释放),仅阻止新条目累积。
|
||||
// 时机:conv 删除即清理(比"loop 结束 + 无 pending"更确定——conv 删了必无 pending,
|
||||
// 上述 retain 已清)。loop 正常收敛/达 MAX/stop 但 conv 未删时不清理(下次发消息复用,限流计数连续)。
|
||||
drop(session); // 释放 AiSession 锁再取 LlmConcurrency 锁(避免潜在锁序问题)
|
||||
state.llm_concurrency.release_conv(&conversation_id).await;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -48,6 +48,7 @@ pub(crate) async fn compress_via_llm(
|
||||
provider_config: &AiProviderRecord,
|
||||
active_msgs: Vec<ChatMessage>,
|
||||
lang: &str,
|
||||
conv_id: &str,
|
||||
llm_concurrency: &LlmConcurrency,
|
||||
) -> Result<String, String> {
|
||||
if active_msgs.is_empty() {
|
||||
@@ -79,8 +80,9 @@ pub(crate) async fn compress_via_llm(
|
||||
};
|
||||
|
||||
// LLM 并发限流(压缩属独立调用,纳入双层 Semaphore)
|
||||
// F-09 B 批5: per_conv 改 HashMap<conv_id>,压缩在 loop 内针对本对话,用 conv_id 共享限流槽。
|
||||
let _global_permit = llm_concurrency.acquire_global().await;
|
||||
let _per_conv_permit = llm_concurrency.acquire_per_conv().await;
|
||||
let _per_conv_permit = llm_concurrency.acquire_per_conv(conv_id).await;
|
||||
|
||||
let resp = provider
|
||||
.complete(request)
|
||||
|
||||
@@ -388,8 +388,9 @@ async fn extract_knowledge_from_conversation(
|
||||
Err(e) => return Err(anyhow::anyhow!("provider 密钥不可用: {}", e)),
|
||||
};
|
||||
// LLM 并发限流(知识提炼属独立调用,纳入双层 Semaphore)
|
||||
// F-09 B 批5: per_conv 改 HashMap<conv_id>,知识提炼针对本对话,用 conv_id 共享该 conv 限流槽。
|
||||
let _global_permit = llm_concurrency.acquire_global().await;
|
||||
let _per_conv_permit = llm_concurrency.acquire_per_conv().await;
|
||||
let _per_conv_permit = llm_concurrency.acquire_per_conv(conv_id).await;
|
||||
let resp = provider.complete(request).await?;
|
||||
let raw = resp.text.trim();
|
||||
|
||||
|
||||
@@ -122,7 +122,7 @@ pub(crate) async fn ensure_conversation_title(
|
||||
// 超时或返回 None 均保留上方已落的 extract 兜底,不覆盖。
|
||||
let title = match tokio::time::timeout(
|
||||
std::time::Duration::from_secs(20),
|
||||
generate_title_via_llm(&*provider, &title_model, summary_msgs, &llm_concurrency),
|
||||
generate_title_via_llm(&*provider, &title_model, summary_msgs, conv_id, &llm_concurrency),
|
||||
).await {
|
||||
Ok(Some(t)) => t,
|
||||
Ok(None) => {
|
||||
@@ -167,6 +167,7 @@ async fn generate_title_via_llm(
|
||||
provider: &dyn LlmProvider,
|
||||
model: &str,
|
||||
msgs: Vec<ChatMessage>,
|
||||
conv_id: &str,
|
||||
llm_concurrency: &LlmConcurrency,
|
||||
) -> Option<String> {
|
||||
let mut prompt = vec![ChatMessage::system(
|
||||
@@ -184,8 +185,9 @@ async fn generate_title_via_llm(
|
||||
reasoning_content: None,
|
||||
};
|
||||
// LLM 并发限流(标题生成属独立调用,纳入双层 Semaphore)
|
||||
// F-09 B 批5: per_conv 改 HashMap<conv_id>,标题针对本对话,用 conv_id 共享限流槽。
|
||||
let _global_permit = llm_concurrency.acquire_global().await;
|
||||
let _per_conv_permit = llm_concurrency.acquire_per_conv().await;
|
||||
let _per_conv_permit = llm_concurrency.acquire_per_conv(conv_id).await;
|
||||
let resp = provider.complete(request).await.ok()?;
|
||||
Some(clean_title(&resp.text))
|
||||
}
|
||||
|
||||
@@ -549,7 +549,8 @@ async fn extract_description_via_llm(
|
||||
};
|
||||
|
||||
let _g = state.llm_concurrency.acquire_global().await;
|
||||
let _c = state.llm_concurrency.acquire_per_conv().await;
|
||||
// F-09 B 批5: 项目扫描无 conv_id(非会话内 LLM 调用),用合成 key 共享一槽(扫描低频,同类聚合限流)。
|
||||
let _c = state.llm_concurrency.acquire_per_conv("__project_scan__").await;
|
||||
let resp = provider.complete(request).await.map_err(err_str)?;
|
||||
// 只取 description,其它字段丢弃(批量场景不需要 project_type/stack 细化)
|
||||
let desc = parse_scan_result(&resp.text)
|
||||
@@ -644,7 +645,8 @@ pub async fn scan_project_with_ai(
|
||||
|
||||
// 4. 双层限流 + complete
|
||||
let _global_permit = state.llm_concurrency.acquire_global().await;
|
||||
let _per_conv_permit = state.llm_concurrency.acquire_per_conv().await;
|
||||
// F-09 B 批5: 项目扫描无 conv_id,用合成 key 共享一槽(扫描低频,同类聚合限流)。
|
||||
let _per_conv_permit = state.llm_concurrency.acquire_per_conv("__project_scan__").await;
|
||||
let llm_result = provider.complete(request).await;
|
||||
|
||||
// 5. 解析 + 合并 stack(LLM 失败降级纯规则)
|
||||
|
||||
@@ -81,20 +81,28 @@ impl Default for KnowledgeConfig {
|
||||
// LLM 调用并发控制(双层 Semaphore + 可选 per-provider 层)
|
||||
// ============================================================
|
||||
|
||||
/// LLM 调用并发控制 — 全局 + 单对话双层 Semaphore + 可选 per-provider 层
|
||||
/// LLM 调用并发控制 — 会话级 global + 单对话内 per_conv 双层 Semaphore + 可选 per-provider 层
|
||||
///
|
||||
/// 限流对象:所有真实 LLM 调用(主循环 stream_llm / 标题生成 / 知识提炼)。
|
||||
/// 不限流本地工具执行(tools.execute)——本地操作无外部成本、不受 RPM 约束。
|
||||
///
|
||||
/// ## F-260616-09 B 批5 / 决策 c-1:global 改「并发会话数上限」
|
||||
/// `global` permits 默认 3(`new(3, 2)`),语义从「全局 LLM 调用并发上限」改为
|
||||
/// 「并发会话数上限」:`run_agentic_loop` 入口 `acquire_global()` 拿 1 permit,持有整个 loop
|
||||
/// 生命周期(含工具执行/审批等待/重试)。同一时刻最多 3 个对话并发跑 loop,第 4 个排队。
|
||||
/// 收益:token 暴增护栏(决策 c 原意)。
|
||||
///
|
||||
/// ## per_conv 改 HashMap<conv_id, Semaphore>
|
||||
/// 原应用级单信号量(AiSession 单例 + generating 互斥下退化)改为
|
||||
/// `HashMap<String, Arc<Semaphore>>`:每对话一份(permits=2,主循环 + 标题 + 提炼各自限流)。
|
||||
/// `acquire_per_conv(conv_id)`:lock HashMap → 无则建(permits=2)→ clone Arc → 释放 lock → acquire_owned。
|
||||
/// conv 退出清理(`release_conv(conv_id)`):loop 结束 + 无 pending 审批时 remove 条目(防 HashMap 无限增长)。
|
||||
///
|
||||
/// 运行时调整:tokio Semaphore 的 permits 数构造时固定、不可增减,
|
||||
/// 故用 `Arc<Mutex<Arc<Semaphore>>>` 双层包装——替换内层 Arc 即重建 Semaphore。
|
||||
/// 故 `global` 用 `Arc<Mutex<Arc<Semaphore>>>` 双层包装——替换内层 Arc 即重建 Semaphore。
|
||||
/// 已持有旧 permit 的任务不受影响(permit 绑定旧 Semaphore,软收敛),
|
||||
/// 新请求 lock 后克隆到最新 Arc、自动走新限制。旧 Semaphore 随最后 permit 释放而 drop。
|
||||
///
|
||||
/// 注意:per_conv 当前是应用级单一信号量(非 per-conv map)。因 AiSession 为单例 +
|
||||
/// generating 互斥,同一时刻仅一个对话的 loop 在跑,per_conv 退化为"单对话内并发"
|
||||
/// (主循环 stream_llm + 标题生成 + 知识提炼)。未来若支持多对话并发,
|
||||
/// 需改为 HashMap<conv_id, Semaphore>。
|
||||
/// per_conv HashMap 内每条 Arc<Semaphore> 不需运行时改 permits(对话内并发上限 2 固定),故无重建需求。
|
||||
///
|
||||
/// ## F-260614-04: per-provider 层(可选)
|
||||
/// `per_provider` 为 HashMap<provider_id, Arc<Semaphore>>。调用方(agentic loop)经
|
||||
@@ -105,8 +113,16 @@ impl Default for KnowledgeConfig {
|
||||
/// (set_provider_caps 传 min(sum, global_cap)),非运行时强约束。
|
||||
#[derive(Clone)]
|
||||
pub struct LlmConcurrency {
|
||||
/// 会话级并发上限(F--09 B 批5/决策 c-1):permits=默认 3,run_agentic_loop 入口拿 1 持整 loop。
|
||||
global: Arc<Mutex<Arc<Semaphore>>>,
|
||||
per_conv: Arc<Mutex<Arc<Semaphore>>>,
|
||||
/// 单对话内并发上限(F-09 B 批5):HashMap<conv_id, Semaphore>,每对话 permits=2。
|
||||
/// acquire 时按 conv_id 取/建;release_conv 在 loop 结束 + 无 pending 时 remove。
|
||||
per_conv: Arc<Mutex<HashMap<String, Arc<Semaphore>>>>,
|
||||
/// 每对话内并发上限(permits 默认 2,构造时传入 new(3, 2) 的第二参)。
|
||||
/// F-09 B 批5:AtomicUsize 支持运行时热改(ai_set_concurrency_config 的 per_conv_limit)。
|
||||
/// 热改后**已建对话**的旧 Semaphore 不变(permits 构造时固定),**新建对话**用新值;
|
||||
/// 为使热改立即全量生效,set_per_conv 同时清空 HashMap 强制重建(软收敛:旧 permit 随 Drop 释放)。
|
||||
per_conv_permits: Arc<AtomicUsize>,
|
||||
/// F-260614-04: per-provider 信号量表。空 = 无 per-provider 限流(单 provider 路径零变化)。
|
||||
/// Arc<Mutex<HashMap>>:运行时增删 provider 配置时替换/插入,acquire 时 clone Arc。
|
||||
per_provider: Arc<Mutex<HashMap<String, Arc<Semaphore>>>>,
|
||||
@@ -116,31 +132,53 @@ impl LlmConcurrency {
|
||||
pub fn new(global: usize, per_conv: usize) -> Self {
|
||||
Self {
|
||||
global: Arc::new(Mutex::new(Arc::new(Semaphore::new(global)))),
|
||||
per_conv: Arc::new(Mutex::new(Arc::new(Semaphore::new(per_conv)))),
|
||||
per_conv: Arc::new(Mutex::new(HashMap::new())),
|
||||
per_conv_permits: Arc::new(AtomicUsize::new(per_conv)),
|
||||
per_provider: Arc::new(Mutex::new(HashMap::new())),
|
||||
}
|
||||
}
|
||||
|
||||
/// 取全局并发 permit(重建后新请求自动走最新 Semaphore)
|
||||
/// 取会话级并发 permit(F-09 B 批5/决策 c-1):run_agentic_loop 入口拿 1 持整个 loop 生命周期。
|
||||
/// 重建后(set_global)新请求自动走最新 Semaphore。
|
||||
pub async fn acquire_global(&self) -> tokio::sync::OwnedSemaphorePermit {
|
||||
let sema = self.global.lock().await.clone();
|
||||
sema.acquire_owned().await.expect("llm global semaphore closed")
|
||||
}
|
||||
|
||||
/// 取单对话并发 permit
|
||||
pub async fn acquire_per_conv(&self) -> tokio::sync::OwnedSemaphorePermit {
|
||||
let sema = self.per_conv.lock().await.clone();
|
||||
/// 取单对话内并发 permit(F-09 B 批5):按 conv_id 取/建 Semaphore,permits=当前 per_conv_permits(默认 2)。
|
||||
/// lock HashMap → 无则建 → clone Arc → 释放 lock → acquire_owned。锁持有短(不含 await acquire)。
|
||||
pub async fn acquire_per_conv(&self, conv_id: &str) -> tokio::sync::OwnedSemaphorePermit {
|
||||
let sema = {
|
||||
let mut map = self.per_conv.lock().await;
|
||||
let permits = self.per_conv_permits.load(std::sync::atomic::Ordering::SeqCst);
|
||||
map.entry(conv_id.to_string())
|
||||
.or_insert_with(|| Arc::new(Semaphore::new(permits)))
|
||||
.clone()
|
||||
};
|
||||
sema.acquire_owned().await.expect("llm per_conv semaphore closed")
|
||||
}
|
||||
|
||||
/// F-09 B 批5:conv 退出清理。loop 结束 + 无 pending 审批时 remove 该 conv 的 Semaphore 条目。
|
||||
/// **时机由调用方判断**(agentic loop 退出点):仅在确信无后续 acquire 时调用,否则误删会致
|
||||
/// 该 conv 下次 acquire 重建 Semaphore(限流计数清零,非致命,但语义偏离)。
|
||||
/// remove 后已持 permit 不受影响(permit 绑旧 Arc,随 Drop 释放),仅阻止新条目累积。
|
||||
pub async fn release_conv(&self, conv_id: &str) {
|
||||
self.per_conv.lock().await.remove(conv_id);
|
||||
}
|
||||
|
||||
/// 重建全局 Semaphore(软收敛:旧 permit 不回收,待其释放后新限制完全生效)
|
||||
pub async fn set_global(&self, permits: usize) {
|
||||
*self.global.lock().await = Arc::new(Semaphore::new(permits));
|
||||
}
|
||||
|
||||
/// 重建单对话 Semaphore
|
||||
/// 重建单对话 Semaphore(F-09 B 批5 后 per_conv 为 HashMap)。
|
||||
/// 行为:更新 per_conv_permits(AtomicUsize) + 清空 HashMap(软收敛:旧 permit 随 Drop 释放,
|
||||
/// 新对话 acquire 用新 permits 值重建 Semaphore)。已建对话若仍在跑,旧 Semaphore 不变;
|
||||
/// 下次该 conv 新 acquire 时因 HashMap 已清空会重建为新 permits。
|
||||
pub async fn set_per_conv(&self, permits: usize) {
|
||||
*self.per_conv.lock().await = Arc::new(Semaphore::new(permits));
|
||||
self.per_conv_permits
|
||||
.store(permits, std::sync::atomic::Ordering::SeqCst);
|
||||
self.per_conv.lock().await.clear();
|
||||
}
|
||||
|
||||
/// F-260614-04: 取 per-provider 并发 permit(可选)。
|
||||
@@ -151,7 +189,7 @@ impl LlmConcurrency {
|
||||
/// 调用方(agentic loop)用法:
|
||||
/// ```ignore
|
||||
/// let _global_permit = llm_concurrency.acquire_global().await;
|
||||
/// let _per_conv_permit = llm_concurrency.acquire_per_conv().await;
|
||||
/// let _per_conv_permit = llm_concurrency.acquire_per_conv(&conv_id).await;
|
||||
/// let _provider_permit = llm_concurrency.acquire_for_provider(&provider_id).await;
|
||||
/// ```
|
||||
/// 三 permit 均绑 guard Drop 自动释放。None 时无 permit 需释放(行为同 F-01 前)。
|
||||
|
||||
Reference in New Issue
Block a user