1#![deny(unsafe_code)]
2
3use std::sync::Arc;
4
5use anyhow::anyhow;
6use async_trait::async_trait;
7use tokio::{self, sync::oneshot, sync::RwLock};
8
9use rmqtt::{
10 context::ServerContext,
11 hook::{Register, Type},
12 macros::Plugin,
13 plugin::{PackageInfo, Plugin},
14 register, Result,
15};
16
17use config::PluginConfig;
18
19mod api;
20mod clients;
21mod config;
22mod handler;
23mod plugin;
24mod prome;
25mod subs;
26mod types;
27
28type ShutdownTX = oneshot::Sender<()>;
29type PluginConfigType = Arc<RwLock<PluginConfig>>;
30
31register!(HttpApiPlugin::new);
32
33#[derive(Plugin)]
34struct HttpApiPlugin {
35 scx: ServerContext,
36 register: Box<dyn Register>,
37 cfg: PluginConfigType,
38 shutdown_tx: Option<ShutdownTX>,
39}
40
41impl HttpApiPlugin {
42 #[inline]
43 async fn new<S: Into<String>>(scx: ServerContext, name: S) -> Result<Self> {
44 let name = name.into();
45 let cfg = scx.plugins.read_config_default::<PluginConfig>(&name);
46 log::info!("{name} HttpApiPlugin cfg: {cfg:?}");
47 let cfg = Arc::new(RwLock::new(cfg?));
48 let register = scx.extends.hook_mgr().register();
49 let shutdown_tx = Some(Self::start(scx.clone(), cfg.clone()).await?);
50 Ok(Self { scx, register, cfg, shutdown_tx })
51 }
52
53 async fn start(scx: ServerContext, cfg: PluginConfigType) -> Result<ShutdownTX> {
54 let (shutdown_tx, shutdown_rx) = oneshot::channel();
55 let (started_tx, started_rx) = oneshot::channel();
56 let http_laddr = cfg.read().await.http_laddr;
57 tokio::spawn(async move {
58 if let Err(e) = api::listen_and_serve(scx, http_laddr, cfg, shutdown_rx, started_tx).await {
59 log::error!("{e:?}");
60 }
61 log::info!("Exit HTTP API Server, ..., http://{http_laddr:?}");
62 });
63
64 started_rx.await.map_err(|_| anyhow!("HTTP API server failed to bind on {http_laddr}"))?;
65
66 Ok(shutdown_tx)
67 }
68}
69
70#[async_trait]
71impl Plugin for HttpApiPlugin {
72 #[inline]
73 async fn init(&mut self) -> Result<()> {
74 log::info!("{} init", self.name());
75 let mgs_type = self.cfg.read().await.message_type;
76 self.register
77 .add(Type::GrpcMessageReceived, Box::new(handler::HookHandler::new(self.scx.clone(), mgs_type)))
78 .await;
79 Ok(())
80 }
81
82 #[inline]
83 async fn get_config(&self) -> Result<serde_json::Value> {
84 self.cfg.read().await.to_json()
85 }
86
87 #[inline]
88 async fn load_config(&mut self) -> Result<()> {
89 let new_cfg = self.scx.plugins.read_config::<PluginConfig>(self.name())?;
90 if !self.cfg.read().await.changed(&new_cfg) {
91 return Ok(());
92 }
93 let restart_enable = self.cfg.read().await.restart_enable(&new_cfg);
94 if restart_enable {
95 let new_cfg = Arc::new(RwLock::new(new_cfg));
96 match Self::start(self.scx.clone(), new_cfg.clone()).await {
97 Ok(tx) => {
98 if let Some(old_tx) = self.shutdown_tx.take() {
99 let _ = old_tx.send(());
100 }
101 self.shutdown_tx = Some(tx);
102 self.cfg = new_cfg;
103 }
104 Err(e) => {
105 return Err(e);
106 }
107 }
108 } else {
109 *self.cfg.write().await = new_cfg;
110 }
111
112 log::debug!("load_config ok, {:?}", self.cfg);
113 Ok(())
114 }
115
116 #[inline]
117 async fn start(&mut self) -> Result<()> {
118 log::info!("{} start", self.name());
119 self.register.start().await;
120 Ok(())
121 }
122
123 #[inline]
124 async fn stop(&mut self) -> Result<bool> {
125 log::info!("{} stop", self.name());
126 Ok(false)
128 }
129}