Skip to main content

rmqtt_http_api/
lib.rs

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        //self.register.stop().await;
127        Ok(false)
128    }
129}