mirror of
https://github.com/aiclientproxy/proxycast.git
synced 2026-09-24 23:10:56 +08:00
feat: add request dedup, response cache, and capability routing metrics to API server
- Add response cache middleware (backed by aster-rust implementation) - Add request deduplication middleware to prevent duplicate in-flight requests - Add capability routing metrics middleware for model/provider fallback tracking - Add idempotency stats (atomic counters) to IdempotencyStore - Expose all four stores in ServerState and surface stats in ServerStatus - Add ResponseCacheSettings to ServerConfig with defaults (600s TTL, 200 max entries) - Add API server observability infrastructure (IdempotencyGuard, RequestDedupGuard, ResponseCacheGuard) - Add inline capability detection for vision/tools/context when routing requests - Register workspace_ensure_ready and workspace_ensure_default_ready commands in runner - Update docs with response_cache configuration example Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Sonnet 4.6
parent
2d00911c28
commit
0b2d3caa6c
@@ -72,6 +72,25 @@ routing:
|
||||
|
||||
适用:有脚本联动、自动化流程需求的用户。
|
||||
|
||||
## 示例 4:API 缓存策略(高级)
|
||||
|
||||
目标:精确控制哪些响应状态码参与短时缓存(非流式)。
|
||||
|
||||
```yaml
|
||||
server:
|
||||
host: "127.0.0.1"
|
||||
port: 8999
|
||||
api_key: "your-api-key"
|
||||
response_cache:
|
||||
enabled: true
|
||||
ttl_secs: 600
|
||||
max_entries: 200
|
||||
max_body_bytes: 1048576
|
||||
cacheable_status_codes: [200] # 默认仅缓存 200;可按需扩展 [200, 201]
|
||||
```
|
||||
|
||||
适用:对缓存命中与语义一致性有要求的自动化/API 调用场景。
|
||||
|
||||
## 调整顺序建议
|
||||
|
||||
1. 先确认导航与主题
|
||||
|
||||
Generated
+1
-2
@@ -370,7 +370,6 @@ checksum = "7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50"
|
||||
[[package]]
|
||||
name = "aster-core"
|
||||
version = "0.15.0"
|
||||
source = "git+https://github.com/astercloud/aster-rust?tag=v0.15.0#eb0377b71bcc056589cd802c46c98fec8880544f"
|
||||
dependencies = [
|
||||
"ahash",
|
||||
"anyhow",
|
||||
@@ -463,7 +462,6 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "aster-models"
|
||||
version = "0.15.0"
|
||||
source = "git+https://github.com/astercloud/aster-rust?tag=v0.15.0#eb0377b71bcc056589cd802c46c98fec8880544f"
|
||||
dependencies = [
|
||||
"serde",
|
||||
"serde_json",
|
||||
@@ -7338,6 +7336,7 @@ dependencies = [
|
||||
name = "proxycast-server"
|
||||
version = "0.76.0"
|
||||
dependencies = [
|
||||
"aster-core",
|
||||
"async-stream",
|
||||
"axum 0.7.9",
|
||||
"base64 0.22.1",
|
||||
|
||||
@@ -30,10 +30,10 @@ pub use types::{
|
||||
MemoryConfig, MemoryProfileConfig, MemoryResolveConfig, MemorySourcesConfig, ModelInfo,
|
||||
ModelsConfig, NativeAgentConfig, NavigationConfig, OpenAIAsrConfig, PairingSettings,
|
||||
ProviderConfig, ProviderModelsConfig, ProvidersConfig, QuotaExceededConfig, RateLimitSettings,
|
||||
RemoteManagementConfig, RetrySettings, RoutingConfig, ScreenshotChatConfig, SearchEngine,
|
||||
ServerConfig, TaskSchedule, TlsConfig, UpdateCheckConfig, UserProfile, VertexApiKeyEntry,
|
||||
VertexModelAlias, VoiceConfig, VoiceInputConfig, VoiceInstruction, VoiceOutputConfig,
|
||||
VoiceOutputMode, VoiceProcessorConfig, WebSearchConfig, WhisperLocalConfig, WhisperModelSize,
|
||||
WorkspaceSandboxConfig, XunfeiConfig, DEFAULT_API_KEY,
|
||||
RemoteManagementConfig, ResponseCacheSettings, RetrySettings, RoutingConfig,
|
||||
ScreenshotChatConfig, SearchEngine, ServerConfig, TaskSchedule, TlsConfig, UpdateCheckConfig,
|
||||
UserProfile, VertexApiKeyEntry, VertexModelAlias, VoiceConfig, VoiceInputConfig,
|
||||
VoiceInstruction, VoiceOutputConfig, VoiceOutputMode, VoiceProcessorConfig, WebSearchConfig,
|
||||
WhisperLocalConfig, WhisperModelSize, WorkspaceSandboxConfig, XunfeiConfig, DEFAULT_API_KEY,
|
||||
};
|
||||
pub use yaml::{load_config, save_config, ConfigError, ConfigManager, YamlService};
|
||||
|
||||
@@ -39,6 +39,7 @@ fn arb_server_config() -> impl Strategy<Value = ServerConfig> {
|
||||
port,
|
||||
api_key,
|
||||
tls: crate::config::TlsConfig::default(),
|
||||
response_cache: crate::config::ResponseCacheSettings::default(),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -357,6 +358,7 @@ fn arb_valid_server_config() -> impl Strategy<Value = ServerConfig> {
|
||||
port,
|
||||
api_key,
|
||||
tls: crate::config::TlsConfig::default(),
|
||||
response_cache: crate::config::ResponseCacheSettings::default(),
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -977,6 +977,61 @@ pub struct ServerConfig {
|
||||
/// TLS 配置
|
||||
#[serde(default)]
|
||||
pub tls: TlsConfig,
|
||||
/// 响应缓存配置(仅影响非流式请求)
|
||||
#[serde(default)]
|
||||
pub response_cache: ResponseCacheSettings,
|
||||
}
|
||||
|
||||
/// 响应缓存配置
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
|
||||
pub struct ResponseCacheSettings {
|
||||
/// 是否启用响应缓存
|
||||
#[serde(default = "default_response_cache_enabled")]
|
||||
pub enabled: bool,
|
||||
/// 缓存 TTL(秒)
|
||||
#[serde(default = "default_response_cache_ttl_secs")]
|
||||
pub ttl_secs: u64,
|
||||
/// 最大缓存条目数
|
||||
#[serde(default = "default_response_cache_max_entries")]
|
||||
pub max_entries: usize,
|
||||
/// 单响应最大缓存字节数
|
||||
#[serde(default = "default_response_cache_max_body_bytes")]
|
||||
pub max_body_bytes: usize,
|
||||
/// 可缓存的 HTTP 状态码列表(默认仅 200)
|
||||
#[serde(default = "default_response_cache_cacheable_status_codes")]
|
||||
pub cacheable_status_codes: Vec<u16>,
|
||||
}
|
||||
|
||||
fn default_response_cache_enabled() -> bool {
|
||||
true
|
||||
}
|
||||
|
||||
fn default_response_cache_ttl_secs() -> u64 {
|
||||
600
|
||||
}
|
||||
|
||||
fn default_response_cache_max_entries() -> usize {
|
||||
200
|
||||
}
|
||||
|
||||
fn default_response_cache_max_body_bytes() -> usize {
|
||||
1_048_576
|
||||
}
|
||||
|
||||
fn default_response_cache_cacheable_status_codes() -> Vec<u16> {
|
||||
vec![200]
|
||||
}
|
||||
|
||||
impl Default for ResponseCacheSettings {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
enabled: default_response_cache_enabled(),
|
||||
ttl_secs: default_response_cache_ttl_secs(),
|
||||
max_entries: default_response_cache_max_entries(),
|
||||
max_body_bytes: default_response_cache_max_body_bytes(),
|
||||
cacheable_status_codes: default_response_cache_cacheable_status_codes(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// TLS 配置
|
||||
@@ -1113,6 +1168,7 @@ impl Default for ServerConfig {
|
||||
port: default_port(),
|
||||
api_key: default_api_key(),
|
||||
tls: TlsConfig::default(),
|
||||
response_cache: ResponseCacheSettings::default(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2181,6 +2237,11 @@ mod unit_tests {
|
||||
assert_eq!(config.server.host, "127.0.0.1");
|
||||
assert_eq!(config.server.port, 8999);
|
||||
assert_eq!(config.server.api_key, "proxy_cast");
|
||||
assert!(config.server.response_cache.enabled);
|
||||
assert_eq!(
|
||||
config.server.response_cache.cacheable_status_codes,
|
||||
vec![200]
|
||||
);
|
||||
assert!(config.providers.kiro.enabled);
|
||||
assert!(!config.providers.gemini.enabled);
|
||||
assert_eq!(config.default_provider, "kiro");
|
||||
@@ -2304,6 +2365,21 @@ mod unit_tests {
|
||||
assert_eq!(config.host, "127.0.0.1");
|
||||
assert_eq!(config.port, 8999);
|
||||
assert_eq!(config.api_key, "proxy_cast");
|
||||
assert!(config.response_cache.enabled);
|
||||
assert_eq!(config.response_cache.ttl_secs, 600);
|
||||
assert_eq!(config.response_cache.max_entries, 200);
|
||||
assert_eq!(config.response_cache.max_body_bytes, 1_048_576);
|
||||
assert_eq!(config.response_cache.cacheable_status_codes, vec![200]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_response_cache_settings_default() {
|
||||
let config = ResponseCacheSettings::default();
|
||||
assert!(config.enabled);
|
||||
assert_eq!(config.ttl_secs, 600);
|
||||
assert_eq!(config.max_entries, 200);
|
||||
assert_eq!(config.max_body_bytes, 1_048_576);
|
||||
assert_eq!(config.cacheable_status_codes, vec![200]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -821,6 +821,32 @@ server:
|
||||
// 其他字段应使用默认值
|
||||
assert_eq!(config.server.host, "127.0.0.1");
|
||||
assert_eq!(config.retry.max_retries, 3);
|
||||
assert_eq!(
|
||||
config.server.response_cache.cacheable_status_codes,
|
||||
vec![200]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_parse_yaml_with_response_cache_settings() {
|
||||
let yaml = r#"
|
||||
server:
|
||||
response_cache:
|
||||
enabled: true
|
||||
ttl_secs: 120
|
||||
max_entries: 64
|
||||
max_body_bytes: 262144
|
||||
cacheable_status_codes: [200, 201]
|
||||
"#;
|
||||
let config = ConfigManager::parse_yaml(yaml).unwrap();
|
||||
assert!(config.server.response_cache.enabled);
|
||||
assert_eq!(config.server.response_cache.ttl_secs, 120);
|
||||
assert_eq!(config.server.response_cache.max_entries, 64);
|
||||
assert_eq!(config.server.response_cache.max_body_bytes, 262144);
|
||||
assert_eq!(
|
||||
config.server.response_cache.cacheable_status_codes,
|
||||
vec![200, 201]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -15,6 +15,7 @@ proxycast-processor.workspace = true
|
||||
proxycast-server-utils.workspace = true
|
||||
proxycast-scheduler.workspace = true
|
||||
proxycast-agent.workspace = true
|
||||
aster.workspace = true
|
||||
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -171,6 +171,15 @@ pub struct ServerStatus {
|
||||
pub open_circuit_count: u32,
|
||||
/// 当前活跃请求数(近似值,当前版本默认 0)
|
||||
pub active_requests: u64,
|
||||
/// 能力过滤与跨 Provider 回退指标
|
||||
pub capability_routing:
|
||||
middleware::capability_routing_metrics::CapabilityRoutingMetricsSnapshot,
|
||||
/// 响应缓存运行时统计
|
||||
pub response_cache: middleware::response_cache::ResponseCacheStats,
|
||||
/// 请求去重运行时统计
|
||||
pub request_dedup: middleware::request_dedup::RequestDedupStats,
|
||||
/// 幂等运行时统计
|
||||
pub idempotency: middleware::idempotency::IdempotencyStats,
|
||||
}
|
||||
|
||||
pub struct ServerState {
|
||||
@@ -191,6 +200,15 @@ pub struct ServerState {
|
||||
pub running_api_key: Option<String>,
|
||||
/// 服务器实际监听的 host(可能与配置不同,因为会自动切换到有效的 IP)
|
||||
pub running_host: Option<String>,
|
||||
/// 能力路由指标(能力过滤/模型回退/Provider 回退)
|
||||
pub capability_routing_metrics_store:
|
||||
Arc<middleware::capability_routing_metrics::CapabilityRoutingMetricsStore>,
|
||||
/// 响应缓存存储(用于状态统计与运行时共享)
|
||||
pub response_cache_store: Arc<middleware::response_cache::ResponseCacheStore>,
|
||||
/// 请求去重存储(用于状态统计与运行时共享)
|
||||
pub request_dedup_store: Arc<middleware::request_dedup::RequestDedupStore>,
|
||||
/// 幂等性存储(用于状态统计与运行时共享)
|
||||
pub idempotency_store: Arc<middleware::idempotency::IdempotencyStore>,
|
||||
}
|
||||
|
||||
impl ServerState {
|
||||
@@ -200,6 +218,21 @@ impl ServerState {
|
||||
let openai_custom = OpenAICustomProvider::new();
|
||||
let claude_custom = ClaudeCustomProvider::new();
|
||||
let default_provider_ref = Arc::new(RwLock::new(config.default_provider.clone()));
|
||||
let idempotency_store = Arc::new(middleware::idempotency::IdempotencyStore::new(
|
||||
middleware::idempotency::IdempotencyConfig::default(),
|
||||
));
|
||||
let request_dedup_store = Arc::new(middleware::request_dedup::RequestDedupStore::new(
|
||||
middleware::request_dedup::RequestDedupConfig::default(),
|
||||
));
|
||||
let response_cache_store = Arc::new(middleware::response_cache::ResponseCacheStore::new(
|
||||
middleware::response_cache::ResponseCacheConfig {
|
||||
enabled: config.server.response_cache.enabled,
|
||||
ttl_secs: config.server.response_cache.ttl_secs,
|
||||
max_entries: config.server.response_cache.max_entries,
|
||||
max_body_bytes: config.server.response_cache.max_body_bytes,
|
||||
cacheable_status_codes: config.server.response_cache.cacheable_status_codes.clone(),
|
||||
},
|
||||
));
|
||||
|
||||
Self {
|
||||
config,
|
||||
@@ -215,6 +248,12 @@ impl ServerState {
|
||||
shutdown_tx: None,
|
||||
running_api_key: None,
|
||||
running_host: None,
|
||||
capability_routing_metrics_store: Arc::new(
|
||||
middleware::capability_routing_metrics::CapabilityRoutingMetricsStore::new(),
|
||||
),
|
||||
response_cache_store,
|
||||
request_dedup_store,
|
||||
idempotency_store,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -233,6 +272,10 @@ impl ServerState {
|
||||
p95_latency_ms_1m: None,
|
||||
open_circuit_count: 0,
|
||||
active_requests: 0,
|
||||
capability_routing: self.capability_routing_metrics_store.snapshot(),
|
||||
response_cache: self.response_cache_store.stats(),
|
||||
request_dedup: self.request_dedup_store.stats(),
|
||||
idempotency: self.idempotency_store.stats(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -392,6 +435,27 @@ impl ServerState {
|
||||
|
||||
// 保存实际使用的 host(在移动到 spawn 之前克隆)
|
||||
let running_host = host.clone();
|
||||
let idempotency_store = Arc::new(middleware::idempotency::IdempotencyStore::new(
|
||||
middleware::idempotency::IdempotencyConfig::default(),
|
||||
));
|
||||
self.idempotency_store = idempotency_store.clone();
|
||||
let request_dedup_store = Arc::new(middleware::request_dedup::RequestDedupStore::new(
|
||||
middleware::request_dedup::RequestDedupConfig::default(),
|
||||
));
|
||||
self.request_dedup_store = request_dedup_store.clone();
|
||||
let capability_routing_metrics_store =
|
||||
Arc::new(middleware::capability_routing_metrics::CapabilityRoutingMetricsStore::new());
|
||||
self.capability_routing_metrics_store = capability_routing_metrics_store.clone();
|
||||
let response_cache_store = Arc::new(middleware::response_cache::ResponseCacheStore::new(
|
||||
middleware::response_cache::ResponseCacheConfig {
|
||||
enabled: config.server.response_cache.enabled,
|
||||
ttl_secs: config.server.response_cache.ttl_secs,
|
||||
max_entries: config.server.response_cache.max_entries,
|
||||
max_body_bytes: config.server.response_cache.max_body_bytes,
|
||||
cacheable_status_codes: config.server.response_cache.cacheable_status_codes.clone(),
|
||||
},
|
||||
));
|
||||
self.response_cache_store = response_cache_store.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = run_server(
|
||||
@@ -413,6 +477,10 @@ impl ServerState {
|
||||
Some(config),
|
||||
Some(config_path),
|
||||
Some(processor),
|
||||
capability_routing_metrics_store,
|
||||
response_cache_store,
|
||||
request_dedup_store,
|
||||
idempotency_store,
|
||||
None, // dev_bridge_callback: 由主 crate 在重新导出层注入
|
||||
)
|
||||
.await
|
||||
@@ -477,6 +545,9 @@ pub struct AppState {
|
||||
pub amp_router: Arc<proxycast_core::router::AmpRouter>,
|
||||
/// 端点 Provider 配置
|
||||
pub endpoint_providers: Arc<RwLock<EndpointProvidersConfig>>,
|
||||
/// Provider 维度模型配置(用于能力感知回退)
|
||||
pub provider_models:
|
||||
Arc<std::collections::HashMap<String, proxycast_core::config::ProviderModelsConfig>>,
|
||||
/// Kiro 事件服务
|
||||
pub kiro_event_service: Arc<KiroEventService>,
|
||||
/// API Key Provider 服务(用于智能降级)
|
||||
@@ -488,6 +559,13 @@ pub struct AppState {
|
||||
pub rate_limiter: Option<Arc<middleware::rate_limit::SlidingWindowRateLimiter>>,
|
||||
/// 幂等性存储
|
||||
pub idempotency_store: Arc<middleware::idempotency::IdempotencyStore>,
|
||||
/// 请求去重存储(请求指纹 in-flight + 短 TTL 回放)
|
||||
pub request_dedup_store: Arc<middleware::request_dedup::RequestDedupStore>,
|
||||
/// 响应缓存存储(非流式短时缓存)
|
||||
pub response_cache_store: Arc<middleware::response_cache::ResponseCacheStore>,
|
||||
/// 能力路由指标(能力过滤/模型回退/Provider 回退)
|
||||
pub capability_routing_metrics_store:
|
||||
Arc<middleware::capability_routing_metrics::CapabilityRoutingMetricsStore>,
|
||||
/// 凭证清理器
|
||||
pub sanitizer: Arc<proxycast_core::sanitizer::CredentialSanitizer>,
|
||||
}
|
||||
@@ -775,6 +853,12 @@ async fn run_server(
|
||||
config: Option<Config>,
|
||||
config_path: Option<PathBuf>,
|
||||
processor: Option<Arc<RequestProcessor>>,
|
||||
capability_routing_metrics_store: Arc<
|
||||
middleware::capability_routing_metrics::CapabilityRoutingMetricsStore,
|
||||
>,
|
||||
response_cache_store: Arc<middleware::response_cache::ResponseCacheStore>,
|
||||
request_dedup_store: Arc<middleware::request_dedup::RequestDedupStore>,
|
||||
idempotency_store: Arc<middleware::idempotency::IdempotencyStore>,
|
||||
dev_bridge_callback: Option<DevBridgeCallback>,
|
||||
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
let base_url = format!("http://{host}:{port}");
|
||||
@@ -865,6 +949,12 @@ async fn run_server(
|
||||
.map(|c| c.endpoint_providers.clone())
|
||||
.unwrap_or_default(),
|
||||
));
|
||||
let provider_models = Arc::new(
|
||||
config
|
||||
.as_ref()
|
||||
.map(|c| c.models.providers.clone())
|
||||
.unwrap_or_default(),
|
||||
);
|
||||
|
||||
// 创建 Kiro 事件服务
|
||||
let kiro_event_service = Arc::new(KiroEventService::new());
|
||||
@@ -878,7 +968,6 @@ async fn run_server(
|
||||
.as_ref()
|
||||
.map(|c| c.retry.auto_switch_provider)
|
||||
.unwrap_or(true);
|
||||
|
||||
let state = AppState {
|
||||
api_key: api_key.to_string(),
|
||||
base_url,
|
||||
@@ -900,6 +989,7 @@ async fn run_server(
|
||||
request_logger: shared_logger,
|
||||
amp_router,
|
||||
endpoint_providers,
|
||||
provider_models,
|
||||
kiro_event_service,
|
||||
api_key_service,
|
||||
batch_executor: Arc::new(tokio::sync::RwLock::new(None)),
|
||||
@@ -908,9 +998,10 @@ async fn run_server(
|
||||
middleware::rate_limit::RateLimitConfig::default(),
|
||||
),
|
||||
)),
|
||||
idempotency_store: Arc::new(middleware::idempotency::IdempotencyStore::new(
|
||||
middleware::idempotency::IdempotencyConfig::default(),
|
||||
)),
|
||||
idempotency_store,
|
||||
request_dedup_store,
|
||||
response_cache_store,
|
||||
capability_routing_metrics_store,
|
||||
sanitizer: Arc::new(proxycast_core::sanitizer::CredentialSanitizer::with_defaults()),
|
||||
};
|
||||
|
||||
|
||||
@@ -0,0 +1,7 @@
|
||||
//! 能力路由指标适配层
|
||||
//!
|
||||
//! 复用 aster-rust 中的通用实现,避免本地重复维护。
|
||||
|
||||
pub use aster::network::{
|
||||
CapabilityFilterExcludedReason, CapabilityRoutingMetricsSnapshot, CapabilityRoutingMetricsStore,
|
||||
};
|
||||
@@ -5,6 +5,7 @@
|
||||
use parking_lot::Mutex;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashMap;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
/// 幂等性配置
|
||||
@@ -62,10 +63,28 @@ enum RequestState {
|
||||
},
|
||||
}
|
||||
|
||||
/// 幂等性运行时统计
|
||||
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
|
||||
pub struct IdempotencyStats {
|
||||
pub entries_size: u64,
|
||||
pub in_progress_size: u64,
|
||||
pub completed_size: u64,
|
||||
pub check_new_total: u64,
|
||||
pub check_in_progress_total: u64,
|
||||
pub check_completed_total: u64,
|
||||
pub complete_total: u64,
|
||||
pub remove_total: u64,
|
||||
}
|
||||
|
||||
/// 幂等性存储
|
||||
pub struct IdempotencyStore {
|
||||
config: IdempotencyConfig,
|
||||
entries: Mutex<HashMap<String, RequestState>>,
|
||||
check_new_total: AtomicU64,
|
||||
check_in_progress_total: AtomicU64,
|
||||
check_completed_total: AtomicU64,
|
||||
complete_total: AtomicU64,
|
||||
remove_total: AtomicU64,
|
||||
}
|
||||
|
||||
impl IdempotencyStore {
|
||||
@@ -73,6 +92,11 @@ impl IdempotencyStore {
|
||||
Self {
|
||||
config,
|
||||
entries: Mutex::new(HashMap::new()),
|
||||
check_new_total: AtomicU64::new(0),
|
||||
check_in_progress_total: AtomicU64::new(0),
|
||||
check_completed_total: AtomicU64::new(0),
|
||||
complete_total: AtomicU64::new(0),
|
||||
remove_total: AtomicU64::new(0),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -94,8 +118,10 @@ impl IdempotencyStore {
|
||||
key.to_string(),
|
||||
RequestState::InProgress { started_at: now },
|
||||
);
|
||||
self.check_new_total.fetch_add(1, Ordering::Relaxed);
|
||||
IdempotencyCheck::New
|
||||
} else {
|
||||
self.check_in_progress_total.fetch_add(1, Ordering::Relaxed);
|
||||
IdempotencyCheck::InProgress
|
||||
}
|
||||
}
|
||||
@@ -109,8 +135,10 @@ impl IdempotencyStore {
|
||||
key.to_string(),
|
||||
RequestState::InProgress { started_at: now },
|
||||
);
|
||||
self.check_new_total.fetch_add(1, Ordering::Relaxed);
|
||||
IdempotencyCheck::New
|
||||
} else {
|
||||
self.check_completed_total.fetch_add(1, Ordering::Relaxed);
|
||||
IdempotencyCheck::Completed {
|
||||
status: *status,
|
||||
body: body.clone(),
|
||||
@@ -122,6 +150,7 @@ impl IdempotencyStore {
|
||||
key.to_string(),
|
||||
RequestState::InProgress { started_at: now },
|
||||
);
|
||||
self.check_new_total.fetch_add(1, Ordering::Relaxed);
|
||||
IdempotencyCheck::New
|
||||
}
|
||||
}
|
||||
@@ -141,12 +170,16 @@ impl IdempotencyStore {
|
||||
completed_at: Instant::now(),
|
||||
},
|
||||
);
|
||||
self.complete_total.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
/// 移除键(请求失败时调用,允许重试)
|
||||
pub fn remove(&self, key: &str) {
|
||||
let mut entries = self.entries.lock();
|
||||
entries.remove(key);
|
||||
let removed = entries.remove(key);
|
||||
if removed.is_some() {
|
||||
self.remove_total.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
/// 清理过期条目
|
||||
@@ -168,6 +201,31 @@ impl IdempotencyStore {
|
||||
pub fn is_empty(&self) -> bool {
|
||||
self.entries.lock().is_empty()
|
||||
}
|
||||
|
||||
pub fn stats(&self) -> IdempotencyStats {
|
||||
let entries = self.entries.lock();
|
||||
let entries_size = entries.len() as u64;
|
||||
let mut in_progress_size = 0u64;
|
||||
let mut completed_size = 0u64;
|
||||
for state in entries.values() {
|
||||
match state {
|
||||
RequestState::InProgress { .. } => in_progress_size += 1,
|
||||
RequestState::Completed { .. } => completed_size += 1,
|
||||
}
|
||||
}
|
||||
drop(entries);
|
||||
|
||||
IdempotencyStats {
|
||||
entries_size,
|
||||
in_progress_size,
|
||||
completed_size,
|
||||
check_new_total: self.check_new_total.load(Ordering::Relaxed),
|
||||
check_in_progress_total: self.check_in_progress_total.load(Ordering::Relaxed),
|
||||
check_completed_total: self.check_completed_total.load(Ordering::Relaxed),
|
||||
complete_total: self.complete_total.load(Ordering::Relaxed),
|
||||
remove_total: self.remove_total.load(Ordering::Relaxed),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -267,4 +325,29 @@ mod tests {
|
||||
assert_eq!(config.ttl_secs, 86400);
|
||||
assert_eq!(config.header_name, "Idempotency-Key");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_stats_tracking() {
|
||||
let store = IdempotencyStore::new(enabled_config(60));
|
||||
|
||||
assert_eq!(store.check("key1"), IdempotencyCheck::New);
|
||||
assert_eq!(store.check("key1"), IdempotencyCheck::InProgress);
|
||||
store.complete("key1", 200, "ok".to_string());
|
||||
assert_eq!(
|
||||
store.check("key1"),
|
||||
IdempotencyCheck::Completed {
|
||||
status: 200,
|
||||
body: "ok".to_string(),
|
||||
}
|
||||
);
|
||||
store.remove("key1");
|
||||
|
||||
let stats = store.stats();
|
||||
assert_eq!(stats.entries_size, 0);
|
||||
assert_eq!(stats.check_new_total, 1);
|
||||
assert_eq!(stats.check_in_progress_total, 1);
|
||||
assert_eq!(stats.check_completed_total, 1);
|
||||
assert_eq!(stats.complete_total, 1);
|
||||
assert_eq!(stats.remove_total, 1);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,4 +1,7 @@
|
||||
//! 服务器中间件模块
|
||||
|
||||
pub mod capability_routing_metrics;
|
||||
pub mod idempotency;
|
||||
pub mod rate_limit;
|
||||
pub mod request_dedup;
|
||||
pub mod response_cache;
|
||||
|
||||
@@ -0,0 +1,8 @@
|
||||
//! 请求去重能力适配层
|
||||
//!
|
||||
//! 复用 aster-rust 中的通用实现,避免本地重复维护。
|
||||
|
||||
pub use aster::network::{
|
||||
build_request_fingerprint, CompletedReplay, RequestDedupCheck, RequestDedupConfig,
|
||||
RequestDedupStats, RequestDedupStore,
|
||||
};
|
||||
@@ -0,0 +1,7 @@
|
||||
//! 响应缓存能力适配层
|
||||
//!
|
||||
//! 复用 aster-rust 中的通用实现,避免本地重复维护。
|
||||
|
||||
pub use aster::network::{
|
||||
CachedHttpResponse, ResponseCacheConfig, ResponseCacheStats, ResponseCacheStore,
|
||||
};
|
||||
@@ -119,6 +119,10 @@ pub async fn get_server_status(
|
||||
.iter()
|
||||
.filter(|log| matches!(log.status, RequestStatus::Retrying))
|
||||
.count() as u64;
|
||||
status.capability_routing = s.capability_routing_metrics_store.snapshot();
|
||||
status.response_cache = s.response_cache_store.stats();
|
||||
status.request_dedup = s.request_dedup_store.stats();
|
||||
status.idempotency = s.idempotency_store.stats();
|
||||
|
||||
Ok(status)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user