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