Skip to main content

rmqtt_http_api/
lib.rs

1//! HTTP API plugin for RMQTT.
2//!
3//! Provides a RESTful HTTP API for broker management and monitoring.
4//! Supports operations like querying client information, managing
5//! subscriptions, publishing messages, and retrieving metrics.
6//!
7//! # Features
8//!
9//! - Client management: list, kick, inspect sessions.
10//! - Message operations: publish, query retained messages.
11//! - Subscription management: list, add, remove subscriptions.
12//! - Metrics: broker statistics, node info, Prometheus endpoint.
13//! - Cluster-aware: queries span all nodes via gRPC.
14//!
15#![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        //self.register.stop().await;
141        Ok(false)
142    }
143}