1use std::any::{Any, TypeId};
7use std::collections::HashMap;
8use std::path::PathBuf;
9use std::sync::{Arc, LazyLock};
10
11use dashmap::DashMap;
12use tokio::signal;
13use tokio::sync::RwLock;
14use tokio::task::JoinHandle;
15use tokio::time::Instant;
16use tokio_util::sync::CancellationToken;
17use tracing::{debug, info};
18
19use crate::component::Component;
20use crate::config::AppAllConfig;
21use crate::registry::{ComponentMeta, COMPONENT_REGISTRY};
22use crate::scope::Scope;
23use crate::store::{CompRef, Store, TraitImplEntry};
24use crate::topology::{all_metas, topo_sort};
25use crate::RIE;
26
27pub type InnerContext = DashMap<TypeId, CompRef>;
29
30static SYS_CONFIG: LazyLock<DashMap<String, String>> = LazyLock::new(DashMap::new);
32
33pub fn get_sys_config(key: &str) -> Option<String> {
35 SYS_CONFIG.get(key).map(|v| v.value().clone())
36}
37
38pub fn set_sys_config(key: &str, value: String) {
40 SYS_CONFIG.insert(key.to_string(), value);
41}
42
43pub const CONFIG_PATH: &str = "config_path";
45
46pub struct BuildContext {
50 store: Store,
51 metas: Vec<&'static ComponentMeta>,
52}
53
54impl BuildContext {
55 pub fn inner_new(ctx: InnerContext) -> Self {
57 BuildContext {
58 store: Store::from_dashmap(ctx),
59 metas: vec![],
60 }
61 }
62
63 #[inline]
69 pub fn new<P: Into<PathBuf>>(config_path: Option<P>) -> Self {
70 let mut ctx = Self {
71 store: Store::new(),
72 metas: vec![],
73 };
74
75 let app_configs = AppAllConfig::new(config_path);
77 ctx.store.insert_cached(app_configs);
78
79 ctx.auto_register_all();
81
82 ctx
83 }
84
85 fn auto_register_all(&mut self) {
87 for meta in COMPONENT_REGISTRY.iter() {
89 if !meta.trait_impls.is_empty() {
90 for trait_fn in meta.impl_traits {
91 let trait_tid = trait_fn();
92 self.store
93 .trait_impls
94 .entry(trait_tid)
95 .or_default()
96 .extend(meta.trait_impls.to_vec());
97 debug!("组件 '{}' 实现了 trait {:?}", meta.name, trait_tid);
98 }
99 }
100 }
101
102 let metas: Vec<&'static ComponentMeta> = COMPONENT_REGISTRY.iter().collect();
104 let sorted_ids = topo_sort(&metas, &self.store.trait_impls).unwrap_or_else(|e| {
105 panic!("{}", e);
106 });
107
108 for tid in &sorted_ids {
110 if let Some(meta) = metas.iter().find(|m| (m.type_id)() == *tid) {
111 self.register_factory(meta);
112 self.metas.push(meta);
113 }
114 }
115 }
116
117 fn register_factory(&mut self, meta: &ComponentMeta) {
122 let type_id = (meta.type_id)();
123 let scope = meta.scope;
124 let factory = meta.factory;
125
126 match scope {
127 Scope::Singleton => {
128 let instance = factory(&self.store);
129 let arc: Arc<dyn Any + Send + Sync> = Arc::from(instance);
130 self.store.inner().insert(type_id, CompRef::Cached(arc));
131 }
132 Scope::Prototype => {
133 let closure =
134 move |store: &Store| -> Arc<dyn Any + Send + Sync> {
135 let boxed = factory(store);
136 Arc::from(boxed)
137 };
138 self.store
139 .inner()
140 .insert(type_id, CompRef::Factory(Arc::new(closure)));
141 }
142 }
143 }
144
145 pub fn inject<T: Component>(&self) -> Arc<T> {
149 self.store.inject_or_panic::<T>()
150 }
151
152 pub fn try_inject<T: Component>(&self) -> Option<Arc<T>> {
154 self.store.try_inject::<T>()
155 }
156
157 pub fn store(&self) -> &Store {
159 &self.store
160 }
161
162 #[inline]
166 pub fn len(&self) -> usize {
167 self.store.len()
168 }
169
170 #[inline]
172 pub fn is_empty(&self) -> bool {
173 self.store.is_empty()
174 }
175
176 pub fn debug_registry() -> RIE<()> {
178 let metas = all_metas();
179 let id_to_idx: HashMap<TypeId, (usize, &str)> = metas
180 .iter()
181 .enumerate()
182 .map(|(i, m)| ((m.type_id)(), (i, m.name)))
183 .collect();
184
185 let temp_trait_impls: DashMap<TypeId, Vec<TraitImplEntry>> = DashMap::new();
187 for meta in COMPONENT_REGISTRY.iter() {
188 if !meta.trait_impls.is_empty() {
189 for trait_fn in meta.impl_traits {
190 let trait_tid = trait_fn();
191 temp_trait_impls
192 .entry(trait_tid)
193 .or_default()
194 .extend(meta.trait_impls.to_vec());
195 }
196 }
197 }
198
199 let ans = topo_sort(&metas, &temp_trait_impls).map_err(|e| {
200 crate::AppError::Internal(anyhow::anyhow!("{}", e))
201 })?;
202
203 debug!("组件注册表(拓扑排序后):");
204 debug!("{:20} {:10} deps", "name", "scope");
205 for tid in ans.iter() {
206 let meta = metas[id_to_idx
207 .get(tid)
208 .ok_or_else(|| crate::AppError::Internal(anyhow::anyhow!("RegistryError")))?
209 .0];
210 let dep_names: Vec<&str> = meta
211 .dep_type_ids
212 .iter()
213 .map(|dep_fn| {
214 COMPONENT_REGISTRY
215 .iter()
216 .find(|m| (m.type_id)() == dep_fn())
217 .map(|m| m.name)
218 .unwrap_or("unknown")
219 })
220 .collect();
221 debug!(
222 "{:20} {:10} [{}]",
223 meta.name,
224 format!("{:?}", meta.scope),
225 dep_names.join(", ")
226 )
227 }
228 Ok(())
229 }
230
231 pub fn build(mut self) -> RIE<App> {
235 let shutdown_token = CancellationToken::new();
236 let store = std::mem::replace(&mut self.store, Store::new());
237 let metas = std::mem::take(&mut self.metas);
238 Ok(App {
239 store,
240 metas,
241 shutdown_token,
242 task_handle: RwLock::new(None),
243 })
244 }
245
246 pub async fn build_and_run(self) -> RIE<()> {
248 let app = self.build()?;
249 let arc_app = Arc::new(app);
250 App::run(arc_app.clone(), arc_app.shutdown_token.clone()).await
251 }
252}
253
254impl Default for BuildContext {
255 fn default() -> Self {
256 Self::new::<PathBuf>(None)
257 }
258}
259
260pub struct App {
264 pub store: Store,
265 pub metas: Vec<&'static ComponentMeta>,
266 pub shutdown_token: CancellationToken,
267 pub task_handle: RwLock<Option<JoinHandle<()>>>,
268}
269
270impl App {
271 pub fn inject<T: Component>(&self) -> Arc<T> {
273 self.store.inject_or_panic::<T>()
274 }
275
276 pub fn try_inject<T: Component>(&self) -> Option<Arc<T>> {
278 self.store.try_inject::<T>()
279 }
280
281 #[inline]
283 pub fn len(&self) -> usize {
284 self.store.len()
285 }
286
287 #[inline]
289 pub fn is_empty(&self) -> bool {
290 self.store.is_empty()
291 }
292
293 pub fn store(&self) -> &Store {
295 &self.store
296 }
297
298 fn init(app: &Arc<App>) -> RIE<()> {
302 for meta in &app.metas {
305 debug!("[di] init: {}", meta.name);
306 (meta.init_fn)(app)?;
307 }
308 Ok(())
309 }
310
311 async fn async_init(app: &Arc<App>) -> RIE<()> {
313 for meta in &app.metas {
315 debug!("[di] async_init: {}", meta.name);
316 (meta.async_init_fn)(app).await?;
317 }
318 Ok(())
319 }
320
321 async fn comp_run(app: Arc<App>, token: CancellationToken) -> RIE<()> {
323 let mut handles = Vec::new();
324
325 let metas: Vec<&'static ComponentMeta> = app.metas.clone();
327 for meta in metas {
328 let app_clone = app.clone();
329 let token_clone = token.clone();
330 let name = meta.name;
331 debug!("[di] async_run spawn: {}", name);
332
333 let handle = tokio::spawn(async move {
334 if let Err(e) = (meta.async_run_fn)(&app_clone, token_clone).await {
335 tracing::error!("[di] 组件 '{}' async_run 失败: {:?}", name, e);
336 }
337 });
338 handles.push(handle);
339 }
340 for handle in handles {
342 let _ = handle.await;
343 }
344 Ok(())
345 }
346
347 async fn run(app: Arc<App>, token: CancellationToken) -> RIE<()> {
349 App::init(&app)?;
350 App::async_init(&app).await?;
351 App::comp_run(app, token).await?;
352 Ok(())
353 }
354
355 pub async fn ins_run(self) -> RIE<Arc<App>> {
362 let app = Arc::new(App {
363 store: self.store,
364 metas: self.metas,
365 shutdown_token: self.shutdown_token,
366 task_handle: self.task_handle,
367 });
368
369 App::init(&app)?;
372 App::async_init(&app).await?;
373
374 let app_clone = app.clone();
376 let app_handler = tokio::spawn(async move {
377 if let Err(e) = App::comp_run(app_clone.clone(), app_clone.shutdown_token.clone()).await {
378 tracing::error!("[di] App 运行失败: {:?},将执行 shutdown", e);
379 app_clone.shutdown().await;
380 }
381 });
382
383 {
384 let mut guard = app.task_handle.write().await;
385 *guard = Some(app_handler);
386 }
387
388 Ok(app)
389 }
390
391 pub async fn shutdown(&self) {
393 let metas: Vec<&ComponentMeta> = self.metas.clone();
394 for meta in metas.iter().rev() {
396 debug!("[di] shutdown: {}", meta.name);
397 (meta.shutdown_fn)(&self.store);
398 }
399 }
400
401 pub async fn waiting_exit(&self) {
403 App::wait_for_exit_signal().await;
404 let start = Instant::now();
405 info!("正在等待退出...");
406 self.shutdown_token.cancel();
407
408 if let Some(handle) = self.task_handle.write().await.take() {
409 match tokio::time::timeout(std::time::Duration::from_secs(5), handle).await {
410 Ok(Ok(())) => {
411 info!("后台任务已正常关闭");
412 }
413 Ok(Err(e)) => {
414 tracing::error!("后台任务退出时发生错误: {:?}", e);
415 }
416 Err(_) => {
417 tracing::warn!("后台任务关闭超时(5秒),强制退出");
418 }
419 }
420 }
421
422 self.shutdown().await;
424
425 info!("app 已退出,耗时: {:?}", start.elapsed());
426 tokio::time::sleep(std::time::Duration::from_millis(200)).await;
427 }
428
429 async fn wait_for_exit_signal() {
431 #[cfg(unix)]
432 {
433 let mut sigterm = signal::unix::signal(signal::unix::SignalKind::terminate())
434 .expect("无法注册 SIGTERM 处理器");
435 let mut sighup = signal::unix::signal(signal::unix::SignalKind::hangup())
436 .expect("无法注册 SIGHUP 处理器");
437 tokio::select! {
438 _ = signal::ctrl_c() => {},
439 _ = sigterm.recv() => {},
440 _ = sighup.recv() => {},
441 }
442 }
443 #[cfg(windows)]
444 {
445 use tokio::signal::windows;
446 let ctrl_c = signal::ctrl_c();
447 let mut ctrl_break = windows::ctrl_break().expect("无法注册 Ctrl+Break 处理器");
448 let mut ctrl_close = windows::ctrl_close().expect("无法注册 Ctrl+Close 处理器");
449 let mut ctrl_shutdown =
450 windows::ctrl_shutdown().expect("无法注册 Ctrl+Shutdown 处理器");
451 tokio::select! {
452 _ = ctrl_c => {},
453 _ = ctrl_break.recv() => {},
454 _ = ctrl_close.recv() => {},
455 _ = ctrl_shutdown.recv() => {},
456 }
457 }
458 #[cfg(all(not(unix), not(windows)))]
459 {
460 let _ = signal::ctrl_c().await;
461 }
462 }
463}