feat(proxy): integrate rectifier config into request forwarding

- Load rectifier config from database in RequestContext
- Pass config to RequestForwarder for runtime checking
- Use should_rectify_thinking_signature with config parameter
This commit is contained in:
YoVinchen
2026-01-12 14:42:03 +08:00
parent 8506522e26
commit 0f8533ea98
2 changed files with 153 additions and 139 deletions
+142 -138
View File
@@ -7,9 +7,9 @@ use super::{
error::*, error::*,
failover_switch::FailoverSwitchManager, failover_switch::FailoverSwitchManager,
provider_router::ProviderRouter, provider_router::ProviderRouter,
providers::{get_adapter, ProviderAdapter}, providers::{get_adapter, ProviderAdapter, ProviderType},
thinking_rectifier::{detect_rectifier_trigger, rectify_anthropic_request}, thinking_rectifier::{rectify_anthropic_request, should_rectify_thinking_signature},
types::ProxyStatus, types::{ProxyStatus, RectifierConfig},
ProxyError, ProxyError,
}; };
use crate::{app_config::AppType, provider::Provider}; use crate::{app_config::AppType, provider::Provider};
@@ -94,6 +94,8 @@ pub struct RequestForwarder {
app_handle: Option<tauri::AppHandle>, app_handle: Option<tauri::AppHandle>,
/// 请求开始时的"当前供应商 ID"(用于判断是否需要同步 UI/托盘) /// 请求开始时的"当前供应商 ID"(用于判断是否需要同步 UI/托盘)
current_provider_id_at_start: String, current_provider_id_at_start: String,
/// 整流器配置
rectifier_config: RectifierConfig,
} }
impl RequestForwarder { impl RequestForwarder {
@@ -108,6 +110,7 @@ impl RequestForwarder {
current_provider_id_at_start: String, current_provider_id_at_start: String,
_streaming_first_byte_timeout: u64, _streaming_first_byte_timeout: u64,
_streaming_idle_timeout: u64, _streaming_idle_timeout: u64,
rectifier_config: RectifierConfig,
) -> Self { ) -> Self {
// 全局超时设置为 1800 秒(30 分钟),确保业务层超时配置能正常工作 // 全局超时设置为 1800 秒(30 分钟),确保业务层超时配置能正常工作
// 参考 Claude Code Hub 的 undici 全局超时设计 // 参考 Claude Code Hub 的 undici 全局超时设计
@@ -154,6 +157,7 @@ impl RequestForwarder {
failover_manager, failover_manager,
app_handle, app_handle,
current_provider_id_at_start, current_provider_id_at_start,
rectifier_config,
} }
} }
@@ -188,6 +192,9 @@ impl RequestForwarder {
let mut last_provider = None; let mut last_provider = None;
let mut attempted_providers = 0usize; let mut attempted_providers = 0usize;
// 整流器重试标记:确保整流最多触发一次
let mut rectifier_retried = false;
// 单 Provider 场景下跳过熔断器检查(故障转移关闭时) // 单 Provider 场景下跳过熔断器检查(故障转移关闭时)
let bypass_circuit_breaker = providers.len() == 1; let bypass_circuit_breaker = providers.len() == 1;
@@ -282,155 +289,152 @@ impl RequestForwarder {
}); });
} }
Err(e) => { Err(e) => {
// 检测是否需要触发整流器(仅 Claude 供应商) // 检测是否需要触发整流器(仅 Claude/ClaudeAuth 供应商)
let is_claude_provider = adapter.name() == "Claude"; let provider_type = ProviderType::from_app_type_and_config(app_type, provider);
let is_anthropic_provider = matches!(
provider_type,
ProviderType::Claude | ProviderType::ClaudeAuth
);
if is_claude_provider { if is_anthropic_provider {
let error_message = extract_error_message(&e); let error_message = extract_error_message(&e);
if let Some(trigger) = detect_rectifier_trigger(error_message.as_deref()) { if should_rectify_thinking_signature(
error_message.as_deref(),
&self.rectifier_config,
) {
// 已经重试过:直接返回错误(不可重试客户端错误)
if rectifier_retried {
log::warn!("[{app_type_str}] [RECT-005] 整流器已触发过,不再重试");
let mut status = self.status.write().await;
status.failed_requests += 1;
status.last_error = Some(e.to_string());
if status.total_requests > 0 {
status.success_rate = (status.success_requests as f32
/ status.total_requests as f32)
* 100.0;
}
return Err(ForwardError {
error: e,
provider: Some(provider.clone()),
});
}
// 首次触发:整流请求体
let rectified = rectify_anthropic_request(&mut body); let rectified = rectify_anthropic_request(&mut body);
if rectified.applied { // 整流未生效:直接返回错误(不可重试客户端错误)
log::info!( if !rectified.applied {
"[{}] [RECT-001] 整流器触发: {:?}, 移除 {} thinking blocks, {} redacted_thinking blocks, {} signature fields", log::warn!(
app_type_str, "[{app_type_str}] [RECT-006] 整流器触发但无可整流内容,不做无意义重试"
trigger,
rectified.removed_thinking_blocks,
rectified.removed_redacted_thinking_blocks,
rectified.removed_signature_fields
); );
let mut status = self.status.write().await;
status.failed_requests += 1;
status.last_error = Some(e.to_string());
if status.total_requests > 0 {
status.success_rate = (status.success_requests as f32
/ status.total_requests as f32)
* 100.0;
}
return Err(ForwardError {
error: e,
provider: Some(provider.clone()),
});
}
// 使用同一供应商重试(不计入熔断器) log::info!(
match self "[{}] [RECT-001] thinking 签名整流器触发, 移除 {} thinking blocks, {} redacted_thinking blocks, {} signature fields",
.forward(provider, endpoint, &body, &headers, adapter.as_ref()) app_type_str,
.await rectified.removed_thinking_blocks,
{ rectified.removed_redacted_thinking_blocks,
Ok(response) => { rectified.removed_signature_fields
log::info!("[{app_type_str}] [RECT-002] 整流重试成功"); );
// 记录成功
let _ = self
.router
.record_result(
&provider.id,
app_type_str,
used_half_open_permit,
true,
None,
)
.await;
// 更新当前应用类型使用的 provider // 标记已重试(当前逻辑下重试后必定 return,保留标记以备将来扩展)
{ let _ = std::mem::replace(&mut rectifier_retried, true);
let mut current_providers =
self.current_providers.write().await;
current_providers.insert(
app_type_str.to_string(),
(provider.id.clone(), provider.name.clone()),
);
}
// 更新成功统计 // 使用同一供应商重试(不计入熔断器)
{ match self
let mut status = self.status.write().await; .forward(provider, endpoint, &body, &headers, adapter.as_ref())
status.success_requests += 1; .await
status.last_error = None; {
let should_switch = self Ok(response) => {
.current_provider_id_at_start log::info!("[{app_type_str}] [RECT-002] 整流重试成功");
.as_str() // 记录成功
!= provider.id.as_str(); let _ = self
if should_switch { .router
status.failover_count += 1; .record_result(
&provider.id,
app_type_str,
used_half_open_permit,
true,
None,
)
.await;
// 异步触发供应商切换,更新 UI/托盘 // 更新当前应用类型使用的 provider
let fm = self.failover_manager.clone(); {
let ah = self.app_handle.clone(); let mut current_providers =
let pid = provider.id.clone(); self.current_providers.write().await;
let pname = provider.name.clone(); current_providers.insert(
let at = app_type_str.to_string(); app_type_str.to_string(),
(provider.id.clone(), provider.name.clone()),
tokio::spawn(async move {
let _ =
fm.try_switch(ah.as_ref(), &at, &pid, &pname)
.await;
});
}
if status.total_requests > 0 {
status.success_rate = (status.success_requests
as f32
/ status.total_requests as f32)
* 100.0;
}
}
return Ok(ForwardResult {
response,
provider: provider.clone(),
});
}
Err(retry_err) => {
// 整流重试仍失败:走正常错误分类逻辑
log::warn!(
"[{app_type_str}] [RECT-003] 整流重试仍失败: {retry_err}"
); );
}
// 记录失败并更新熔断器 // 更新成功统计
let _ = self {
.router let mut status = self.status.write().await;
.record_result( status.success_requests += 1;
&provider.id, status.last_error = None;
app_type_str, let should_switch =
used_half_open_permit, self.current_provider_id_at_start.as_str()
false, != provider.id.as_str();
Some(retry_err.to_string()), if should_switch {
) status.failover_count += 1;
.await;
// 分类错误,决定是否继续 failover // 异步触发供应商切换,更新 UI/托盘
let category = self.categorize_proxy_error(&retry_err); let fm = self.failover_manager.clone();
match category { let ah = self.app_handle.clone();
ErrorCategory::Retryable => { let pid = provider.id.clone();
// 可重试:继续尝试下一个供应商 let pname = provider.name.clone();
{ let at = app_type_str.to_string();
let mut status = self.status.write().await;
status.last_error = Some(format!( tokio::spawn(async move {
"Provider {} 整流重试失败: {}", let _ = fm
provider.name, retry_err .try_switch(ah.as_ref(), &at, &pid, &pname)
)); .await;
} });
log::warn!( }
"[{}] [RECT-004] Provider {} 整流重试失败,切换下一个 ({}/{})", if status.total_requests > 0 {
app_type_str, status.success_rate = (status.success_requests as f32
provider.name, / status.total_requests as f32)
attempted_providers, * 100.0;
providers.len()
);
last_error = Some(retry_err);
last_provider = Some(provider.clone());
// 继续尝试下一个供应商
continue;
}
ErrorCategory::NonRetryable
| ErrorCategory::ClientAbort => {
// 不可重试:直接返回错误
{
let mut status = self.status.write().await;
status.failed_requests += 1;
status.last_error =
Some(retry_err.to_string());
if status.total_requests > 0 {
status.success_rate =
(status.success_requests as f32
/ status.total_requests as f32)
* 100.0;
}
}
return Err(ForwardError {
error: retry_err,
provider: Some(provider.clone()),
});
}
} }
} }
return Ok(ForwardResult {
response,
provider: provider.clone(),
});
}
Err(retry_err) => {
// 整流重试仍失败:直接返回错误(不可重试客户端错误)
// 不记录熔断器、不继续 failover
log::warn!(
"[{app_type_str}] [RECT-003] 整流重试仍失败: {retry_err}"
);
let mut status = self.status.write().await;
status.failed_requests += 1;
status.last_error = Some(retry_err.to_string());
if status.total_requests > 0 {
status.success_rate = (status.success_requests as f32
/ status.total_requests as f32)
* 100.0;
}
return Err(ForwardError {
error: retry_err,
provider: Some(provider.clone()),
});
} }
} }
} }
+11 -1
View File
@@ -5,7 +5,10 @@
use crate::app_config::AppType; use crate::app_config::AppType;
use crate::provider::Provider; use crate::provider::Provider;
use crate::proxy::{ use crate::proxy::{
extract_session_id, forwarder::RequestForwarder, server::ProxyState, types::AppProxyConfig, extract_session_id,
forwarder::RequestForwarder,
server::ProxyState,
types::{AppProxyConfig, RectifierConfig},
ProxyError, ProxyError,
}; };
use axum::http::HeaderMap; use axum::http::HeaderMap;
@@ -54,6 +57,8 @@ pub struct RequestContext {
pub app_type: AppType, pub app_type: AppType,
/// Session ID(从客户端请求提取或新生成) /// Session ID(从客户端请求提取或新生成)
pub session_id: String, pub session_id: String,
/// 整流器配置
pub rectifier_config: RectifierConfig,
} }
impl RequestContext { impl RequestContext {
@@ -86,6 +91,9 @@ impl RequestContext {
.await .await
.map_err(|e| ProxyError::DatabaseError(e.to_string()))?; .map_err(|e| ProxyError::DatabaseError(e.to_string()))?;
// 从数据库读取整流器配置
let rectifier_config = state.db.get_rectifier_config().unwrap_or_default();
let current_provider_id = let current_provider_id =
crate::settings::get_current_provider(&app_type).unwrap_or_default(); crate::settings::get_current_provider(&app_type).unwrap_or_default();
@@ -147,6 +155,7 @@ impl RequestContext {
app_type_str, app_type_str,
app_type, app_type,
session_id, session_id,
rectifier_config,
}) })
} }
@@ -206,6 +215,7 @@ impl RequestContext {
self.current_provider_id.clone(), self.current_provider_id.clone(),
first_byte_timeout, first_byte_timeout,
idle_timeout, idle_timeout,
self.rectifier_config.clone(),
) )
} }