Skip to main content

flare_core/server/events/
factory.rs

1//! 服务端消息观察者工厂
2//!
3//! 提供观察者创建接口,支持用户自定义观察者实现
4
5use crate::common::MessageParser;
6use crate::server::connection::ConnectionManager;
7use crate::server::events::{ServerMessageWrapper, observer::ConnectionHandlerObserverAdapter};
8use crate::transport::events::ConnectionObserver;
9use std::sync::Arc;
10use tracing::error;
11
12/// ServerCore 的轻量级引用
13///
14/// 用于在工厂中访问 ServerCore 的组件,避免循环引用
15pub struct ServerCoreRef {
16    pub device_manager: Option<Arc<crate::server::device::DeviceManager>>,
17    pub event_handler: Option<Arc<dyn crate::server::events::handler::ServerEventHandler>>,
18}
19
20/// 服务端消息观察者工厂
21///
22/// 用于创建连接观察者,支持用户自定义实现
23///
24/// # 设计原则
25/// - **依赖倒置**:transports 依赖工厂接口,不依赖具体实现
26/// - **开闭原则**:对扩展开放,对修改关闭
27/// - **单一职责**:只负责创建观察者
28///
29/// # 使用示例
30/// ```rust,no_run
31/// use flare_core::server::events::factory::ServerMessageObserverFactory;
32///
33/// struct MyCustomObserverFactory;
34///
35/// impl ServerMessageObserverFactory for MyCustomObserverFactory {
36///     fn create_observer(
37///         &self,
38///         manager: Arc<ConnectionManager>,
39///         parser: MessageParser,
40///         event_handler: Arc<dyn crate::server::events::handler::ServerEventHandler>,
41///         connection_id: String,
42///         core_ref: Arc<ServerCoreRef>,
43///         core: Arc<crate::server::transports::server_core::ServerCore>,
44///     ) -> Arc<dyn ConnectionObserver> {
45///         // 创建自定义观察者
46///         Arc::new(MyCustomObserver::new(...))
47///     }
48/// }
49/// ```
50pub trait ServerMessageObserverFactory: Send + Sync {
51    /// 为指定连接创建观察者
52    ///
53    /// # 参数
54    /// - `manager`: 连接管理器
55    /// - `parser`: 消息解析器
56    /// - `event_handler`: 事件处理器(必需)
57    /// - `connection_id`: 连接 ID
58    /// - `core_ref`: 服务器核心的轻量级引用
59    /// - `core`: 服务器核心的完整引用(用于需要完整功能的情况)
60    ///
61    /// # 返回
62    /// 创建的观察者实例
63    fn create_observer(
64        &self,
65        manager: Arc<ConnectionManager>,
66        parser: MessageParser,
67        event_handler: Arc<dyn crate::server::events::handler::ServerEventHandler>,
68        connection_id: String,
69        core_ref: Arc<ServerCoreRef>,
70        core: Arc<crate::server::transports::server_core::ServerCore>,
71    ) -> Arc<dyn ConnectionObserver>;
72}
73
74/// 默认观察者工厂
75///
76/// 使用 `DefaultServerMessageObserver` 创建观察者
77///
78/// # 使用示例
79/// ```rust,no_run
80/// use flare_core::server::events::factory::DefaultServerMessageObserverFactory;
81///
82/// // 使用默认工厂
83/// let factory = DefaultServerMessageObserverFactory::new();
84///
85/// // 或配置设备管理器和事件处理器
86/// let factory = DefaultServerMessageObserverFactory::new()
87///     .with_device_manager(Some(device_manager))
88///     .with_event_handler(Some(event_handler));
89/// ```
90pub struct DefaultServerMessageObserverFactory {
91    /// 设备管理器(可选)
92    device_manager: Option<Arc<crate::server::device::DeviceManager>>,
93    /// 事件处理器(可选)
94    event_handler: Option<Arc<dyn crate::server::events::handler::ServerEventHandler>>,
95}
96
97impl DefaultServerMessageObserverFactory {
98    /// 创建新的默认工厂
99    pub fn new() -> Self {
100        Self {
101            device_manager: None,
102            event_handler: None,
103        }
104    }
105
106    /// 设置设备管理器
107    pub fn with_device_manager(
108        mut self,
109        device_manager: Option<Arc<crate::server::device::DeviceManager>>,
110    ) -> Self {
111        self.device_manager = device_manager;
112        self
113    }
114
115    /// 设置事件处理器
116    pub fn with_event_handler(
117        mut self,
118        event_handler: Option<Arc<dyn crate::server::events::handler::ServerEventHandler>>,
119    ) -> Self {
120        self.event_handler = event_handler;
121        self
122    }
123}
124
125impl Default for DefaultServerMessageObserverFactory {
126    fn default() -> Self {
127        Self::new()
128    }
129}
130
131/// 服务端观察者工厂实现
132impl ServerMessageObserverFactory for DefaultServerMessageObserverFactory {
133    fn create_observer(
134        &self,
135        manager: Arc<ConnectionManager>,
136        parser: MessageParser,
137        event_handler: Arc<dyn crate::server::events::handler::ServerEventHandler>,
138        connection_id: String,
139        core_ref: Arc<ServerCoreRef>,
140        _core: Arc<crate::server::transports::server_core::ServerCore>,
141    ) -> Arc<dyn ConnectionObserver> {
142        // 优先使用工厂配置,如果没有则使用 core_ref 中的配置
143        let device_manager = self
144            .device_manager
145            .clone()
146            .or_else(|| core_ref.device_manager.clone());
147
148        // 优先使用传入的 event_handler,如果没有则使用工厂配置,最后使用 core_ref 中的配置
149        let event_handler = Some(event_handler)
150            .or_else(|| self.event_handler.clone())
151            .or_else(|| core_ref.event_handler.clone())
152            .ok_or_else(|| {
153                error!("[DefaultServerMessageObserverFactory] ServerEventHandler is required but not provided");
154                "ServerEventHandler is required"
155            })
156            .expect("ServerEventHandler is required");
157
158        // 创建 ServerMessageWrapper(实现 ConnectionHandler)
159        let wrapper = Arc::new(ServerMessageWrapper::new(
160            event_handler,
161            Some(Arc::clone(&manager)),
162            device_manager,
163            parser.clone(), // 克隆 parser 用于适配器
164        ));
165
166        // 将 ConnectionHandler 适配为 ConnectionObserver
167        // 不传递 parser,而是从连接信息中动态获取协商结果来创建 parser
168        Arc::new(ConnectionHandlerObserverAdapter::new(
169            wrapper,
170            connection_id,
171            manager,
172            Some(_core),
173        ))
174    }
175}
176
177/// 观察者链工厂
178///
179/// 支持创建多个观察者,按顺序处理事件
180///
181/// # 使用场景
182/// - 需要多个观察者协同工作
183/// - 需要按优先级处理事件
184/// - 需要组合不同的观察者功能
185///
186/// # 使用示例
187/// ```rust,no_run
188/// use flare_core::server::events::factory::{ChainedObserverFactory, ServerMessageObserverFactory};
189///
190/// let factory1 = Arc::new(MyObserverFactory1::new());
191/// let factory2 = Arc::new(MyObserverFactory2::new());
192///
193/// let chained = ChainedObserverFactory::new()
194///     .add_factory(factory1)
195///     .add_factory(factory2);
196/// ```
197pub struct ChainedObserverFactory {
198    factories: Vec<Arc<dyn ServerMessageObserverFactory>>,
199}
200
201impl ChainedObserverFactory {
202    /// 创建新的链式工厂
203    pub fn new() -> Self {
204        Self {
205            factories: Vec::new(),
206        }
207    }
208
209    /// 添加工厂到链中
210    pub fn add_factory(mut self, factory: Arc<dyn ServerMessageObserverFactory>) -> Self {
211        self.factories.push(factory);
212        self
213    }
214}
215
216impl Default for ChainedObserverFactory {
217    fn default() -> Self {
218        Self::new()
219    }
220}
221
222impl ServerMessageObserverFactory for ChainedObserverFactory {
223    fn create_observer(
224        &self,
225        manager: Arc<ConnectionManager>,
226        parser: MessageParser,
227        event_handler: Arc<dyn crate::server::events::handler::ServerEventHandler>,
228        connection_id: String,
229        core_ref: Arc<ServerCoreRef>,
230        core: Arc<crate::server::transports::server_core::ServerCore>,
231    ) -> Arc<dyn ConnectionObserver> {
232        // 如果只有一个工厂,直接返回
233        if self.factories.len() == 1 {
234            return self.factories[0].create_observer(
235                manager,
236                parser,
237                event_handler,
238                connection_id,
239                core_ref,
240                core,
241            );
242        }
243
244        // 创建多个观察者并组合
245        let observers: Vec<Arc<dyn ConnectionObserver>> = self
246            .factories
247            .iter()
248            .map(|factory| {
249                factory.create_observer(
250                    Arc::clone(&manager),
251                    parser.clone(),
252                    Arc::clone(&event_handler),
253                    connection_id.clone(),
254                    Arc::clone(&core_ref),
255                    Arc::clone(&core),
256                )
257            })
258            .collect();
259
260        // 创建链式观察者包装器
261        Arc::new(ChainedObserver { observers })
262    }
263}
264
265/// 链式观察者包装器
266///
267/// 按顺序调用所有观察者
268struct ChainedObserver {
269    observers: Vec<Arc<dyn ConnectionObserver>>,
270}
271
272impl crate::transport::events::ConnectionObserver for ChainedObserver {
273    fn on_event(&self, event: &crate::transport::events::ConnectionEvent) {
274        for observer in &self.observers {
275            observer.on_event(event);
276        }
277    }
278}