mirror of
https://github.com/aiclientproxy/proxycast.git
synced 2026-09-24 23:10:56 +08:00
refactor: 迁移 config/observer 到 proxycast-config crate
- 创建 proxycast-config crate,包含配置观察者系统 - ConfigEventEmit trait 抽象 Tauri AppHandle 的 emit 功能 - TauriConfigEmitter 和 TauriObserver 保留在主 crate - 11 个 crate 测试 + 58 个主 crate config 测试全部通过 - 主 crate config/observer 重写为重新导出层
This commit is contained in:
Generated
+17
@@ -6669,6 +6669,7 @@ dependencies = [
|
||||
"parking_lot",
|
||||
"portable-pty",
|
||||
"proptest",
|
||||
"proxycast-config",
|
||||
"proxycast-core",
|
||||
"proxycast-infra",
|
||||
"proxycast-providers",
|
||||
@@ -6720,6 +6721,22 @@ dependencies = [
|
||||
"zip",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "proxycast-config"
|
||||
version = "0.60.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"parking_lot",
|
||||
"proptest",
|
||||
"proxycast-core",
|
||||
"proxycast-infra",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"tempfile",
|
||||
"tokio",
|
||||
"tracing",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "proxycast-core"
|
||||
version = "0.60.0"
|
||||
|
||||
@@ -12,6 +12,7 @@ homepage = "https://github.com/aiclientproxy/proxycast"
|
||||
[workspace.dependencies]
|
||||
# 项目内 crate 依赖
|
||||
proxycast-core = { path = "crates/core" }
|
||||
proxycast-config = { path = "crates/config" }
|
||||
proxycast-infra = { path = "crates/infra" }
|
||||
proxycast-providers = { path = "crates/providers" }
|
||||
proxycast-services = { path = "crates/services" }
|
||||
@@ -188,6 +189,7 @@ tauri-build.workspace = true
|
||||
[dependencies]
|
||||
# 项目内 crate
|
||||
proxycast-core.workspace = true
|
||||
proxycast-config.workspace = true
|
||||
proxycast-infra.workspace = true
|
||||
proxycast-providers.workspace = true
|
||||
proxycast-services.workspace = true
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
[package]
|
||||
name = "proxycast-config"
|
||||
version.workspace = true
|
||||
edition.workspace = true
|
||||
authors.workspace = true
|
||||
|
||||
[dependencies]
|
||||
proxycast-core.workspace = true
|
||||
proxycast-infra.workspace = true
|
||||
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
async-trait.workspace = true
|
||||
tokio.workspace = true
|
||||
tracing.workspace = true
|
||||
parking_lot.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
proptest.workspace = true
|
||||
tempfile.workspace = true
|
||||
@@ -0,0 +1,38 @@
|
||||
//! ProxyCast 配置观察者模块
|
||||
//!
|
||||
//! 提供基于观察者模式的配置管理系统。
|
||||
//! Tauri 相关的具体实现(TauriObserver、AppHandle 集成)
|
||||
//! 保留在主 crate 中。
|
||||
|
||||
pub mod observer;
|
||||
|
||||
// 重新导出观察者模块的核心类型
|
||||
pub use observer::emitter::ConfigEventEmit;
|
||||
pub use observer::events::{
|
||||
AmpConfigChangeEvent, ConfigChangeEvent, ConfigChangeSource, CredentialPoolChangeEvent,
|
||||
EndpointProvidersChangeEvent, FullReloadEvent, InjectionChangeEvent, LoggingChangeEvent,
|
||||
NativeAgentChangeEvent, RetryChangeEvent, RoutingChangeEvent, ServerChangeEvent,
|
||||
};
|
||||
pub use observer::manager::GlobalConfigManager;
|
||||
pub use observer::observers::{
|
||||
DefaultProviderRefObserver, EndpointObserver, InjectorObserver, LoggingObserver, RouterObserver,
|
||||
};
|
||||
pub use observer::subject::{ConfigSubject, CONFIG_CHANGED_EVENT, CONFIG_RELOAD_EVENT};
|
||||
pub use observer::traits::{ConfigObserver, FnObserver, SyncConfigObserver, SyncObserverWrapper};
|
||||
|
||||
/// 全局配置管理器状态(用于 Tauri 状态管理)
|
||||
pub struct GlobalConfigManagerState(pub std::sync::Arc<GlobalConfigManager>);
|
||||
|
||||
impl GlobalConfigManagerState {
|
||||
pub fn new(manager: GlobalConfigManager) -> Self {
|
||||
Self(std::sync::Arc::new(manager))
|
||||
}
|
||||
}
|
||||
|
||||
impl std::ops::Deref for GlobalConfigManagerState {
|
||||
type Target = std::sync::Arc<GlobalConfigManager>;
|
||||
|
||||
fn deref(&self) -> &Self::Target {
|
||||
&self.0
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,39 @@
|
||||
//! 配置事件发射器 Trait
|
||||
//!
|
||||
//! 抽象 Tauri AppHandle 的 emit 功能,
|
||||
//! 使 observer 模块不直接依赖 Tauri。
|
||||
|
||||
use super::events::ConfigChangeEvent;
|
||||
|
||||
/// 配置事件发射器 Trait
|
||||
///
|
||||
/// 用于向前端发送配置变更事件。
|
||||
/// 主 crate 中通过 Tauri AppHandle 实现此 trait。
|
||||
pub trait ConfigEventEmit: Send + Sync {
|
||||
/// 发送配置变更事件
|
||||
fn emit_config_event(
|
||||
&self,
|
||||
event_name: &str,
|
||||
payload: &ConfigChangeEvent,
|
||||
) -> Result<(), String>;
|
||||
|
||||
/// 发送无负载事件
|
||||
fn emit_empty_event(&self, event_name: &str) -> Result<(), String>;
|
||||
}
|
||||
|
||||
/// 空操作发射器(用于测试)
|
||||
pub struct NoOpEmitter;
|
||||
|
||||
impl ConfigEventEmit for NoOpEmitter {
|
||||
fn emit_config_event(
|
||||
&self,
|
||||
_event_name: &str,
|
||||
_payload: &ConfigChangeEvent,
|
||||
) -> Result<(), String> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn emit_empty_event(&self, _event_name: &str) -> Result<(), String> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
-9
@@ -11,31 +11,22 @@ use std::collections::HashMap;
|
||||
pub enum ConfigChangeEvent {
|
||||
/// 完整配置重载
|
||||
FullReload(FullReloadEvent),
|
||||
|
||||
/// 路由配置变更
|
||||
RoutingChanged(RoutingChangeEvent),
|
||||
|
||||
/// 注入配置变更
|
||||
InjectionChanged(InjectionChangeEvent),
|
||||
|
||||
/// 端点 Provider 配置变更
|
||||
EndpointProvidersChanged(EndpointProvidersChangeEvent),
|
||||
|
||||
/// 服务器配置变更
|
||||
ServerChanged(ServerChangeEvent),
|
||||
|
||||
/// 日志配置变更
|
||||
LoggingChanged(LoggingChangeEvent),
|
||||
|
||||
/// 重试配置变更
|
||||
RetryChanged(RetryChangeEvent),
|
||||
|
||||
/// Amp CLI 配置变更
|
||||
AmpConfigChanged(AmpConfigChangeEvent),
|
||||
|
||||
/// 凭证池配置变更
|
||||
CredentialPoolChanged(CredentialPoolChangeEvent),
|
||||
|
||||
/// Native Agent 配置变更
|
||||
NativeAgentChanged(NativeAgentChangeEvent),
|
||||
}
|
||||
+22
-53
@@ -1,24 +1,24 @@
|
||||
//! 全局配置管理器
|
||||
//!
|
||||
//! 整合配置主题、热重载和观察者管理
|
||||
//! 整合配置主题、热重载和观察者管理。
|
||||
//! register_processor_observers 和 register_tauri_observer
|
||||
//! 保留在主 crate(依赖 Tauri / RequestProcessor)。
|
||||
|
||||
use super::emitter::ConfigEventEmit;
|
||||
use super::events::ConfigChangeSource;
|
||||
use super::observers::{
|
||||
DefaultProviderRefObserver, EndpointObserver, InjectorObserver, LoggingObserver,
|
||||
RouterObserver, TauriObserver,
|
||||
DefaultProviderRefObserver, EndpointObserver, InjectorObserver, LoggingObserver, RouterObserver,
|
||||
};
|
||||
use super::subject::ConfigSubject;
|
||||
use super::traits::ConfigObserver;
|
||||
use crate::config::{Config, EndpointProvidersConfig, HotReloadManager, ReloadResult};
|
||||
use crate::processor::RequestProcessor;
|
||||
use proxycast_core::config::{Config, EndpointProvidersConfig, HotReloadManager, ReloadResult};
|
||||
use proxycast_core::router::{ModelMapper, Router};
|
||||
use proxycast_infra::Injector;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use tauri::AppHandle;
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
/// 全局配置管理器
|
||||
///
|
||||
/// 整合配置主题、热重载和观察者管理
|
||||
pub struct GlobalConfigManager {
|
||||
/// 配置主题
|
||||
subject: Arc<ConfigSubject>,
|
||||
@@ -54,9 +54,9 @@ impl GlobalConfigManager {
|
||||
self.subject.config()
|
||||
}
|
||||
|
||||
/// 设置 Tauri AppHandle
|
||||
pub fn set_app_handle(&self, handle: AppHandle) {
|
||||
self.subject.set_app_handle(handle);
|
||||
/// 设置事件发射器(替代原来的 set_app_handle)
|
||||
pub fn set_emitter(&self, emitter: Arc<dyn ConfigEventEmit>) {
|
||||
self.subject.set_emitter(emitter);
|
||||
}
|
||||
|
||||
/// 注册观察者
|
||||
@@ -69,24 +69,23 @@ impl GlobalConfigManager {
|
||||
self.subject.unregister(name);
|
||||
}
|
||||
|
||||
/// 注册 RequestProcessor 相关的观察者
|
||||
pub fn register_processor_observers(&self, processor: &RequestProcessor) {
|
||||
// 路由器观察者
|
||||
let router_observer = Arc::new(RouterObserver::new(
|
||||
processor.router.clone(),
|
||||
processor.mapper.clone(),
|
||||
));
|
||||
/// 注册路由器相关观察者
|
||||
pub fn register_router_observers(
|
||||
&self,
|
||||
router: Arc<RwLock<Router>>,
|
||||
mapper: Arc<RwLock<ModelMapper>>,
|
||||
injector: Arc<RwLock<Injector>>,
|
||||
) {
|
||||
let router_observer = Arc::new(RouterObserver::new(router, mapper));
|
||||
self.subject.register(router_observer);
|
||||
|
||||
// 注入器观察者
|
||||
let injector_observer = Arc::new(InjectorObserver::new(processor.injector.clone()));
|
||||
let injector_observer = Arc::new(InjectorObserver::new(injector));
|
||||
self.subject.register(injector_observer);
|
||||
|
||||
// 日志观察者
|
||||
let logging_observer = Arc::new(LoggingObserver);
|
||||
self.subject.register(logging_observer);
|
||||
|
||||
tracing::info!("[GlobalConfigManager] 已注册 RequestProcessor 观察者");
|
||||
tracing::info!("[GlobalConfigManager] 已注册路由器相关观察者");
|
||||
}
|
||||
|
||||
/// 注册端点 Provider 观察者
|
||||
@@ -107,21 +106,12 @@ impl GlobalConfigManager {
|
||||
self.subject.register(observer);
|
||||
}
|
||||
|
||||
/// 注册 Tauri 前端观察者
|
||||
pub fn register_tauri_observer(&self, app_handle: AppHandle) {
|
||||
let observer = Arc::new(TauriObserver::new(app_handle));
|
||||
self.subject.register(observer);
|
||||
}
|
||||
|
||||
/// 更新配置并通知观察者
|
||||
pub async fn update_config(&self, new_config: Config, source: ConfigChangeSource) {
|
||||
// 更新热重载管理器
|
||||
{
|
||||
let hot_reload = self.hot_reload.read();
|
||||
hot_reload.update_config(new_config.clone());
|
||||
}
|
||||
|
||||
// 通知观察者
|
||||
self.subject.update_config(new_config, source).await;
|
||||
}
|
||||
|
||||
@@ -142,7 +132,6 @@ impl GlobalConfigManager {
|
||||
self.subject
|
||||
.update_config(new_config, ConfigChangeSource::HotReload)
|
||||
.await;
|
||||
|
||||
tracing::info!("[GlobalConfigManager] 热重载成功");
|
||||
}
|
||||
ReloadResult::RolledBack { error, .. } => {
|
||||
@@ -158,12 +147,9 @@ impl GlobalConfigManager {
|
||||
|
||||
/// 保存配置到文件并通知观察者
|
||||
pub async fn save_config(&self, config: &Config) -> Result<(), String> {
|
||||
crate::config::save_config(config).map_err(|e| e.to_string())?;
|
||||
|
||||
// 更新内部状态并通知观察者
|
||||
proxycast_core::config::save_config(config).map_err(|e| e.to_string())?;
|
||||
self.update_config(config.clone(), ConfigChangeSource::ApiCall)
|
||||
.await;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -188,23 +174,6 @@ impl GlobalConfigManager {
|
||||
}
|
||||
}
|
||||
|
||||
/// 全局配置管理器状态(用于 Tauri 状态管理)
|
||||
pub struct GlobalConfigManagerState(pub Arc<GlobalConfigManager>);
|
||||
|
||||
impl GlobalConfigManagerState {
|
||||
pub fn new(manager: GlobalConfigManager) -> Self {
|
||||
Self(Arc::new(manager))
|
||||
}
|
||||
}
|
||||
|
||||
impl std::ops::Deref for GlobalConfigManagerState {
|
||||
type Target = Arc<GlobalConfigManager>;
|
||||
|
||||
fn deref(&self) -> &Self::Target {
|
||||
&self.0
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -0,0 +1,10 @@
|
||||
//! 配置观察者模块
|
||||
//!
|
||||
//! 提供基于观察者模式的全局配置管理系统
|
||||
|
||||
pub mod emitter;
|
||||
pub mod events;
|
||||
pub mod manager;
|
||||
pub mod observers;
|
||||
pub mod subject;
|
||||
pub mod traits;
|
||||
+9
-66
@@ -1,15 +1,15 @@
|
||||
//! 内置配置观察者实现
|
||||
//!
|
||||
//! 提供常用组件的配置观察者
|
||||
//! 提供常用组件的配置观察者。
|
||||
//! TauriObserver 保留在主 crate(依赖 Tauri)。
|
||||
|
||||
use super::events::ConfigChangeEvent;
|
||||
use super::traits::ConfigObserver;
|
||||
use crate::config::{Config, EndpointProvidersConfig};
|
||||
use crate::injection::Injector;
|
||||
use crate::router::{ModelMapper, Router};
|
||||
use async_trait::async_trait;
|
||||
use proxycast_core::config::{Config, EndpointProvidersConfig};
|
||||
use proxycast_core::router::{ModelMapper, Router};
|
||||
use proxycast_infra::Injector;
|
||||
use std::sync::Arc;
|
||||
use tauri::{AppHandle, Emitter};
|
||||
use tokio::sync::RwLock;
|
||||
|
||||
/// 路由器观察者
|
||||
@@ -33,7 +33,7 @@ impl ConfigObserver for RouterObserver {
|
||||
}
|
||||
|
||||
fn priority(&self) -> i32 {
|
||||
10 // 高优先级,路由器需要先更新
|
||||
10
|
||||
}
|
||||
|
||||
fn is_interested_in(&self, event: &ConfigChangeEvent) -> bool {
|
||||
@@ -52,7 +52,7 @@ impl ConfigObserver for RouterObserver {
|
||||
if let Ok(provider_type) = config
|
||||
.routing
|
||||
.default_provider
|
||||
.parse::<crate::ProviderType>()
|
||||
.parse::<proxycast_core::ProviderType>()
|
||||
{
|
||||
let mut router = self.router.write().await;
|
||||
router.set_default_provider(provider_type);
|
||||
@@ -80,8 +80,6 @@ impl ConfigObserver for RouterObserver {
|
||||
}
|
||||
|
||||
/// 注入器观察者
|
||||
///
|
||||
/// 监听注入配置变更,更新 Injector
|
||||
pub struct InjectorObserver {
|
||||
injector: Arc<RwLock<Injector>>,
|
||||
}
|
||||
@@ -131,8 +129,6 @@ impl ConfigObserver for InjectorObserver {
|
||||
}
|
||||
|
||||
/// 端点 Provider 观察者
|
||||
///
|
||||
/// 监听端点 Provider 配置变更
|
||||
pub struct EndpointObserver {
|
||||
endpoint_providers: Arc<RwLock<EndpointProvidersConfig>>,
|
||||
}
|
||||
@@ -167,16 +163,12 @@ impl ConfigObserver for EndpointObserver {
|
||||
) -> Result<(), String> {
|
||||
let mut ep = self.endpoint_providers.write().await;
|
||||
*ep = config.endpoint_providers.clone();
|
||||
|
||||
tracing::info!("[EndpointObserver] 更新端点 Provider 配置");
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// 日志观察者
|
||||
///
|
||||
/// 监听日志配置变更
|
||||
pub struct LoggingObserver;
|
||||
|
||||
#[async_trait]
|
||||
@@ -206,58 +198,11 @@ impl ConfigObserver for LoggingObserver {
|
||||
config.logging.enabled,
|
||||
config.logging.level
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// Tauri 前端通知观察者
|
||||
///
|
||||
/// 将配置变更事件转发到前端
|
||||
pub struct TauriObserver {
|
||||
app_handle: AppHandle,
|
||||
}
|
||||
|
||||
impl TauriObserver {
|
||||
pub fn new(app_handle: AppHandle) -> Self {
|
||||
Self { app_handle }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ConfigObserver for TauriObserver {
|
||||
fn name(&self) -> &str {
|
||||
"TauriObserver"
|
||||
}
|
||||
|
||||
fn priority(&self) -> i32 {
|
||||
1000 // 最低优先级,确保其他观察者先处理
|
||||
}
|
||||
|
||||
async fn on_config_changed(
|
||||
&self,
|
||||
event: &ConfigChangeEvent,
|
||||
_config: &Config,
|
||||
) -> Result<(), String> {
|
||||
// 发送详细事件
|
||||
self.app_handle
|
||||
.emit("config-changed-detail", event)
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
// 发送简化的刷新通知
|
||||
self.app_handle
|
||||
.emit("config-refresh-needed", ())
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
tracing::debug!("[TauriObserver] 已通知前端配置变更: {}", event.event_type());
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
/// 默认 Provider 引用观察者
|
||||
///
|
||||
/// 更新 default_provider_ref(用于向后兼容)
|
||||
pub struct DefaultProviderRefObserver {
|
||||
default_provider_ref: Arc<RwLock<String>>,
|
||||
}
|
||||
@@ -277,7 +222,7 @@ impl ConfigObserver for DefaultProviderRefObserver {
|
||||
}
|
||||
|
||||
fn priority(&self) -> i32 {
|
||||
5 // 最高优先级,确保引用先更新
|
||||
5
|
||||
}
|
||||
|
||||
fn is_interested_in(&self, event: &ConfigChangeEvent) -> bool {
|
||||
@@ -294,12 +239,10 @@ impl ConfigObserver for DefaultProviderRefObserver {
|
||||
) -> Result<(), String> {
|
||||
let mut dp = self.default_provider_ref.write().await;
|
||||
*dp = config.routing.default_provider.clone();
|
||||
|
||||
tracing::debug!(
|
||||
"[DefaultProviderRefObserver] 更新 default_provider_ref: {}",
|
||||
config.routing.default_provider
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -307,7 +250,7 @@ impl ConfigObserver for DefaultProviderRefObserver {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::config::observer::events::{ConfigChangeSource, FullReloadEvent};
|
||||
use crate::observer::events::{ConfigChangeSource, FullReloadEvent};
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_router_observer_priority() {
|
||||
+27
-40
@@ -2,14 +2,14 @@
|
||||
//!
|
||||
//! 管理配置观察者的注册、注销和通知
|
||||
|
||||
use super::emitter::ConfigEventEmit;
|
||||
use super::events::{ConfigChangeEvent, ConfigChangeSource, FullReloadEvent};
|
||||
use super::traits::ConfigObserver;
|
||||
use crate::config::Config;
|
||||
use parking_lot::RwLock;
|
||||
use proxycast_core::config::Config;
|
||||
use std::collections::BTreeMap;
|
||||
use std::sync::Arc;
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use tauri::{AppHandle, Emitter};
|
||||
use tokio::sync::broadcast;
|
||||
|
||||
/// Tauri 事件名称常量
|
||||
@@ -33,10 +33,10 @@ pub struct ConfigSubject {
|
||||
current_config: RwLock<Config>,
|
||||
/// 事件广播通道
|
||||
event_tx: broadcast::Sender<ConfigChangeEvent>,
|
||||
/// Tauri AppHandle(用于向前端发送事件)
|
||||
app_handle: RwLock<Option<AppHandle>>,
|
||||
/// 是否启用 Tauri 事件
|
||||
tauri_events_enabled: RwLock<bool>,
|
||||
/// 事件发射器(抽象 Tauri AppHandle)
|
||||
emitter: RwLock<Option<Arc<dyn ConfigEventEmit>>>,
|
||||
/// 是否启用事件发射
|
||||
events_enabled: RwLock<bool>,
|
||||
}
|
||||
|
||||
impl ConfigSubject {
|
||||
@@ -48,21 +48,21 @@ impl ConfigSubject {
|
||||
observers: RwLock::new(BTreeMap::new()),
|
||||
current_config: RwLock::new(initial_config),
|
||||
event_tx,
|
||||
app_handle: RwLock::new(None),
|
||||
tauri_events_enabled: RwLock::new(true),
|
||||
emitter: RwLock::new(None),
|
||||
events_enabled: RwLock::new(true),
|
||||
}
|
||||
}
|
||||
|
||||
/// 设置 Tauri AppHandle
|
||||
pub fn set_app_handle(&self, handle: AppHandle) {
|
||||
let mut app_handle = self.app_handle.write();
|
||||
*app_handle = Some(handle);
|
||||
tracing::debug!("[ConfigSubject] AppHandle 已设置");
|
||||
/// 设置事件发射器
|
||||
pub fn set_emitter(&self, emitter: Arc<dyn ConfigEventEmit>) {
|
||||
let mut e = self.emitter.write();
|
||||
*e = Some(emitter);
|
||||
tracing::debug!("[ConfigSubject] 事件发射器已设置");
|
||||
}
|
||||
|
||||
/// 启用/禁用 Tauri 事件
|
||||
pub fn set_tauri_events_enabled(&self, enabled: bool) {
|
||||
let mut flag = self.tauri_events_enabled.write();
|
||||
/// 启用/禁用事件发射
|
||||
pub fn set_events_enabled(&self, enabled: bool) {
|
||||
let mut flag = self.events_enabled.write();
|
||||
*flag = enabled;
|
||||
}
|
||||
|
||||
@@ -88,9 +88,7 @@ impl ConfigSubject {
|
||||
for entries in observers.values_mut() {
|
||||
entries.retain(|e| e.observer.name() != name);
|
||||
}
|
||||
// 清理空的优先级组
|
||||
observers.retain(|_, v| !v.is_empty());
|
||||
|
||||
tracing::info!("[ConfigSubject] 注销观察者: {}", name);
|
||||
}
|
||||
|
||||
@@ -112,25 +110,18 @@ impl ConfigSubject {
|
||||
|
||||
/// 更新配置并通知观察者
|
||||
pub async fn update_config(&self, new_config: Config, source: ConfigChangeSource) {
|
||||
// 创建完整重载事件
|
||||
let event = ConfigChangeEvent::FullReload(FullReloadEvent {
|
||||
timestamp_ms: Self::current_timestamp_ms(),
|
||||
source,
|
||||
});
|
||||
|
||||
// 更新配置
|
||||
{
|
||||
let mut config = self.current_config.write();
|
||||
*config = new_config.clone();
|
||||
}
|
||||
|
||||
// 通知观察者
|
||||
self.notify_observers(&event, &new_config).await;
|
||||
|
||||
// 发送 Tauri 事件
|
||||
self.emit_tauri_event(&event);
|
||||
|
||||
// 广播事件
|
||||
self.emit_event(&event);
|
||||
let _ = self.event_tx.send(event);
|
||||
}
|
||||
|
||||
@@ -138,7 +129,7 @@ impl ConfigSubject {
|
||||
pub async fn notify_event(&self, event: ConfigChangeEvent) {
|
||||
let config = self.config();
|
||||
self.notify_observers(&event, &config).await;
|
||||
self.emit_tauri_event(&event);
|
||||
self.emit_event(&event);
|
||||
let _ = self.event_tx.send(event);
|
||||
}
|
||||
|
||||
@@ -149,7 +140,6 @@ impl ConfigSubject {
|
||||
|
||||
/// 通知所有观察者
|
||||
async fn notify_observers(&self, event: &ConfigChangeEvent, config: &Config) {
|
||||
// 收集所有感兴趣的观察者
|
||||
let observers: Vec<Arc<dyn ConfigObserver>> = {
|
||||
let observers = self.observers.read();
|
||||
observers
|
||||
@@ -166,7 +156,6 @@ impl ConfigSubject {
|
||||
event.event_type()
|
||||
);
|
||||
|
||||
// 按优先级顺序通知(BTreeMap 已排序)
|
||||
for observer in observers {
|
||||
let name = observer.name().to_string();
|
||||
match observer.on_config_changed(event, config).await {
|
||||
@@ -180,19 +169,19 @@ impl ConfigSubject {
|
||||
}
|
||||
}
|
||||
|
||||
/// 发送 Tauri 事件到前端
|
||||
fn emit_tauri_event(&self, event: &ConfigChangeEvent) {
|
||||
let enabled = *self.tauri_events_enabled.read();
|
||||
/// 发送事件(通过抽象发射器)
|
||||
fn emit_event(&self, event: &ConfigChangeEvent) {
|
||||
let enabled = *self.events_enabled.read();
|
||||
if !enabled {
|
||||
return;
|
||||
}
|
||||
|
||||
let app_handle = self.app_handle.read();
|
||||
if let Some(handle) = app_handle.as_ref() {
|
||||
if let Err(e) = handle.emit(CONFIG_CHANGED_EVENT, event) {
|
||||
tracing::error!("[ConfigSubject] 发送 Tauri 事件失败: {}", e);
|
||||
let emitter = self.emitter.read();
|
||||
if let Some(emitter) = emitter.as_ref() {
|
||||
if let Err(e) = emitter.emit_config_event(CONFIG_CHANGED_EVENT, event) {
|
||||
tracing::error!("[ConfigSubject] 发送事件失败: {}", e);
|
||||
} else {
|
||||
tracing::debug!("[ConfigSubject] 已发送 Tauri 事件: {}", event.event_type());
|
||||
tracing::debug!("[ConfigSubject] 已发送事件: {}", event.event_type());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -231,7 +220,7 @@ impl Default for ConfigSubject {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::config::observer::events::ConfigChangeSource;
|
||||
use crate::observer::events::ConfigChangeSource;
|
||||
use async_trait::async_trait;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
@@ -306,7 +295,6 @@ mod tests {
|
||||
async fn test_priority_order() {
|
||||
let subject = ConfigSubject::new(Config::default());
|
||||
|
||||
// 注册不同优先级的观察者
|
||||
for (name, priority) in [("low", 100), ("high", 10), ("medium", 50)] {
|
||||
let observer = Arc::new(CountingObserver {
|
||||
name: name.to_string(),
|
||||
@@ -316,7 +304,6 @@ mod tests {
|
||||
subject.register(observer);
|
||||
}
|
||||
|
||||
// 验证观察者名称列表
|
||||
let names = subject.observer_names();
|
||||
assert_eq!(names.len(), 3);
|
||||
}
|
||||
+4
-17
@@ -3,8 +3,8 @@
|
||||
//! 定义观察者接口,支持异步和同步两种模式
|
||||
|
||||
use super::events::ConfigChangeEvent;
|
||||
use crate::config::Config;
|
||||
use async_trait::async_trait;
|
||||
use proxycast_core::config::Config;
|
||||
use std::sync::Arc;
|
||||
|
||||
/// 配置观察者 Trait
|
||||
@@ -16,30 +16,18 @@ pub trait ConfigObserver: Send + Sync {
|
||||
fn name(&self) -> &str;
|
||||
|
||||
/// 处理配置变更事件
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `event` - 配置变更事件
|
||||
/// * `config` - 变更后的完整配置
|
||||
///
|
||||
/// # Returns
|
||||
/// * `Ok(())` - 处理成功
|
||||
/// * `Err(String)` - 处理失败,包含错误信息
|
||||
async fn on_config_changed(
|
||||
&self,
|
||||
event: &ConfigChangeEvent,
|
||||
config: &Config,
|
||||
) -> Result<(), String>;
|
||||
|
||||
/// 是否对特定事件类型感兴趣
|
||||
///
|
||||
/// 默认实现对所有事件感兴趣
|
||||
/// 是否对特定事件类型感兴趣(默认对所有事件感兴趣)
|
||||
fn is_interested_in(&self, _event: &ConfigChangeEvent) -> bool {
|
||||
true
|
||||
}
|
||||
|
||||
/// 观察者优先级(数字越小优先级越高)
|
||||
///
|
||||
/// 默认优先级为 100
|
||||
/// 观察者优先级(数字越小优先级越高,默认 100)
|
||||
fn priority(&self) -> i32 {
|
||||
100
|
||||
}
|
||||
@@ -97,7 +85,6 @@ impl<T: SyncConfigObserver + 'static> ConfigObserver for SyncObserverWrapper<T>
|
||||
}
|
||||
|
||||
/// 函数式观察者(用于简单的回调场景)
|
||||
/// 目前主要用于测试,将来可用于动态注册观察者
|
||||
#[allow(dead_code)]
|
||||
pub struct FnObserver<F>
|
||||
where
|
||||
@@ -153,7 +140,7 @@ where
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::config::observer::events::{ConfigChangeSource, FullReloadEvent};
|
||||
use crate::observer::events::{ConfigChangeSource, FullReloadEvent};
|
||||
|
||||
struct TestObserver {
|
||||
name: String,
|
||||
@@ -199,12 +199,15 @@ pub fn run() {
|
||||
}
|
||||
}
|
||||
|
||||
// 设置 GlobalConfigManager 的 AppHandle(用于向前端发送事件)
|
||||
// 设置 GlobalConfigManager 的事件发射器(用于向前端发送事件)
|
||||
if let Some(config_manager) =
|
||||
app.try_state::<crate::config::GlobalConfigManagerState>()
|
||||
{
|
||||
config_manager.0.set_app_handle(app.handle().clone());
|
||||
tracing::info!("[启动] GlobalConfigManager AppHandle 已设置");
|
||||
let emitter = std::sync::Arc::new(
|
||||
crate::config::observer::TauriConfigEmitter::new(app.handle().clone()),
|
||||
);
|
||||
config_manager.0.set_emitter(emitter);
|
||||
tracing::info!("[启动] GlobalConfigManager 事件发射器已设置");
|
||||
}
|
||||
|
||||
// 设置 MCP Manager 的 AppHandle(用于发送 mcp:* 事件)
|
||||
|
||||
@@ -1,22 +1,30 @@
|
||||
//! 配置观察者模块
|
||||
//!
|
||||
//! 提供基于观察者模式的全局配置管理系统
|
||||
//! 核心逻辑已迁移到 proxycast-config crate。
|
||||
//! 本模块保留 Tauri 相关实现和重新导出。
|
||||
|
||||
mod events;
|
||||
mod manager;
|
||||
mod observers;
|
||||
mod subject;
|
||||
mod traits;
|
||||
mod tauri_emitter;
|
||||
mod tauri_observer;
|
||||
|
||||
pub use events::{
|
||||
// 从 proxycast-config crate 重新导出所有类型
|
||||
pub use proxycast_config::observer::emitter::{ConfigEventEmit, NoOpEmitter};
|
||||
pub use proxycast_config::observer::events::{
|
||||
AmpConfigChangeEvent, ConfigChangeEvent, ConfigChangeSource, CredentialPoolChangeEvent,
|
||||
EndpointProvidersChangeEvent, FullReloadEvent, InjectionChangeEvent, LoggingChangeEvent,
|
||||
NativeAgentChangeEvent, RetryChangeEvent, RoutingChangeEvent, ServerChangeEvent,
|
||||
};
|
||||
pub use manager::{GlobalConfigManager, GlobalConfigManagerState};
|
||||
pub use observers::{
|
||||
DefaultProviderRefObserver, EndpointObserver, InjectorObserver, LoggingObserver,
|
||||
RouterObserver, TauriObserver,
|
||||
pub use proxycast_config::observer::manager::GlobalConfigManager;
|
||||
pub use proxycast_config::observer::observers::{
|
||||
DefaultProviderRefObserver, EndpointObserver, InjectorObserver, LoggingObserver, RouterObserver,
|
||||
};
|
||||
pub use subject::{ConfigSubject, CONFIG_CHANGED_EVENT, CONFIG_RELOAD_EVENT};
|
||||
pub use traits::{ConfigObserver, FnObserver, SyncConfigObserver, SyncObserverWrapper};
|
||||
pub use proxycast_config::observer::subject::{
|
||||
ConfigSubject, CONFIG_CHANGED_EVENT, CONFIG_RELOAD_EVENT,
|
||||
};
|
||||
pub use proxycast_config::observer::traits::{
|
||||
ConfigObserver, FnObserver, SyncConfigObserver, SyncObserverWrapper,
|
||||
};
|
||||
pub use proxycast_config::GlobalConfigManagerState;
|
||||
|
||||
// Tauri 相关实现
|
||||
pub use tauri_emitter::TauriConfigEmitter;
|
||||
pub use tauri_observer::TauriObserver;
|
||||
|
||||
@@ -0,0 +1,36 @@
|
||||
//! Tauri 配置事件发射器
|
||||
//!
|
||||
//! 实现 ConfigEventEmit trait,通过 Tauri AppHandle 发送事件
|
||||
|
||||
use proxycast_config::ConfigChangeEvent;
|
||||
use proxycast_config::ConfigEventEmit;
|
||||
use tauri::{AppHandle, Emitter};
|
||||
|
||||
/// Tauri 配置事件发射器
|
||||
pub struct TauriConfigEmitter {
|
||||
app_handle: AppHandle,
|
||||
}
|
||||
|
||||
impl TauriConfigEmitter {
|
||||
pub fn new(app_handle: AppHandle) -> Self {
|
||||
Self { app_handle }
|
||||
}
|
||||
}
|
||||
|
||||
impl ConfigEventEmit for TauriConfigEmitter {
|
||||
fn emit_config_event(
|
||||
&self,
|
||||
event_name: &str,
|
||||
payload: &ConfigChangeEvent,
|
||||
) -> Result<(), String> {
|
||||
self.app_handle
|
||||
.emit(event_name, payload)
|
||||
.map_err(|e| e.to_string())
|
||||
}
|
||||
|
||||
fn emit_empty_event(&self, event_name: &str) -> Result<(), String> {
|
||||
self.app_handle
|
||||
.emit(event_name, ())
|
||||
.map_err(|e| e.to_string())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
//! Tauri 前端通知观察者
|
||||
//!
|
||||
//! 将配置变更事件转发到前端
|
||||
|
||||
use async_trait::async_trait;
|
||||
use proxycast_config::observer::events::ConfigChangeEvent;
|
||||
use proxycast_config::observer::traits::ConfigObserver;
|
||||
use proxycast_core::config::Config;
|
||||
use tauri::{AppHandle, Emitter};
|
||||
|
||||
/// Tauri 前端通知观察者
|
||||
pub struct TauriObserver {
|
||||
app_handle: AppHandle,
|
||||
}
|
||||
|
||||
impl TauriObserver {
|
||||
pub fn new(app_handle: AppHandle) -> Self {
|
||||
Self { app_handle }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ConfigObserver for TauriObserver {
|
||||
fn name(&self) -> &str {
|
||||
"TauriObserver"
|
||||
}
|
||||
|
||||
fn priority(&self) -> i32 {
|
||||
1000
|
||||
}
|
||||
|
||||
async fn on_config_changed(
|
||||
&self,
|
||||
event: &ConfigChangeEvent,
|
||||
_config: &Config,
|
||||
) -> Result<(), String> {
|
||||
self.app_handle
|
||||
.emit("config-changed-detail", event)
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
self.app_handle
|
||||
.emit("config-refresh-needed", ())
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
tracing::debug!("[TauriObserver] 已通知前端配置变更: {}", event.event_type());
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user