diff --git a/.github/images/architecture.png b/.github/images/architecture.png new file mode 100644 index 0000000..88f56a7 Binary files /dev/null and b/.github/images/architecture.png differ diff --git a/.github/images/mantis-logo.png b/.github/images/mantis-logo.png new file mode 100644 index 0000000..d2da023 Binary files /dev/null and b/.github/images/mantis-logo.png differ diff --git a/LICENSE b/LICENSE index f288702..0ad25db 100644 --- a/LICENSE +++ b/LICENSE @@ -1,5 +1,5 @@ - GNU GENERAL PUBLIC LICENSE - Version 3, 29 June 2007 + GNU AFFERO GENERAL PUBLIC LICENSE + Version 3, 19 November 2007 Copyright (C) 2007 Free Software Foundation, Inc. Everyone is permitted to copy and distribute verbatim copies @@ -7,17 +7,15 @@ Preamble - The GNU General Public License is a free, copyleft license for -software and other kinds of works. + The GNU Affero General Public License is a free, copyleft license for +software and other kinds of works, specifically designed to ensure +cooperation with the community in the case of network server software. The licenses for most software and other practical works are designed to take away your freedom to share and change the works. By contrast, -the GNU General Public License is intended to guarantee your freedom to +our General Public Licenses are intended to guarantee your freedom to share and change all versions of a program--to make sure it remains free -software for all its users. We, the Free Software Foundation, use the -GNU General Public License for most of our software; it applies also to -any other work released this way by its authors. You can apply it to -your programs, too. +software for all its users. When we speak of free software, we are referring to freedom, not price. Our General Public Licenses are designed to make sure that you @@ -26,44 +24,34 @@ them if you wish), that you receive source code or can get it if you want it, that you can change the software or use pieces of it in new free programs, and that you know you can do these things. - To protect your rights, we need to prevent others from denying you -these rights or asking you to surrender the rights. Therefore, you have -certain responsibilities if you distribute copies of the software, or if -you modify it: responsibilities to respect the freedom of others. + Developers that use our General Public Licenses protect your rights +with two steps: (1) assert copyright on the software, and (2) offer +you this License which gives you legal permission to copy, distribute +and/or modify the software. - For example, if you distribute copies of such a program, whether -gratis or for a fee, you must pass on to the recipients the same -freedoms that you received. You must make sure that they, too, receive -or can get the source code. And you must show them these terms so they -know their rights. + A secondary benefit of defending all users' freedom is that +improvements made in alternate versions of the program, if they +receive widespread use, become available for other developers to +incorporate. Many developers of free software are heartened and +encouraged by the resulting cooperation. However, in the case of +software used on network servers, this result may fail to come about. +The GNU General Public License permits making a modified version and +letting the public access it on a server without ever releasing its +source code to the public. - Developers that use the GNU GPL protect your rights with two steps: -(1) assert copyright on the software, and (2) offer you this License -giving you legal permission to copy, distribute and/or modify it. + The GNU Affero General Public License is designed specifically to +ensure that, in such cases, the modified source code becomes available +to the community. It requires the operator of a network server to +provide the source code of the modified version running there to the +users of that server. Therefore, public use of a modified version, on +a publicly accessible server, gives the public access to the source +code of the modified version. - For the developers' and authors' protection, the GPL clearly explains -that there is no warranty for this free software. For both users' and -authors' sake, the GPL requires that modified versions be marked as -changed, so that their problems will not be attributed erroneously to -authors of previous versions. - - Some devices are designed to deny users access to install or run -modified versions of the software inside them, although the manufacturer -can do so. This is fundamentally incompatible with the aim of -protecting users' freedom to change the software. The systematic -pattern of such abuse occurs in the area of products for individuals to -use, which is precisely where it is most unacceptable. Therefore, we -have designed this version of the GPL to prohibit the practice for those -products. If such problems arise substantially in other domains, we -stand ready to extend this provision to those domains in future versions -of the GPL, as needed to protect the freedom of users. - - Finally, every program is threatened constantly by software patents. -States should not allow patents to restrict development and use of -software on general-purpose computers, but in those that do, we wish to -avoid the special danger that patents applied to a free program could -make it effectively proprietary. To prevent this, the GPL assures that -patents cannot be used to render the program non-free. + An older license, called the Affero General Public License and +published by Affero, was designed to accomplish similar goals. This is +a different license, not a version of the Affero GPL, but Affero has +released a new version of the Affero GPL which permits relicensing under +this license. The precise terms and conditions for copying, distribution and modification follow. @@ -72,7 +60,7 @@ modification follow. 0. Definitions. - "This License" refers to version 3 of the GNU General Public License. + "This License" refers to version 3 of the GNU Affero General Public License. "Copyright" also means copyright-like laws that apply to other kinds of works, such as semiconductor masks. @@ -549,35 +537,45 @@ to collect a royalty for further conveying from those to whom you convey the Program, the only way you could satisfy both those terms and this License would be to refrain entirely from conveying the Program. - 13. Use with the GNU Affero General Public License. + 13. Remote Network Interaction; Use with the GNU General Public License. + + Notwithstanding any other provision of this License, if you modify the +Program, your modified version must prominently offer all users +interacting with it remotely through a computer network (if your version +supports such interaction) an opportunity to receive the Corresponding +Source of your version by providing access to the Corresponding Source +from a network server at no charge, through some standard or customary +means of facilitating copying of software. This Corresponding Source +shall include the Corresponding Source for any work covered by version 3 +of the GNU General Public License that is incorporated pursuant to the +following paragraph. Notwithstanding any other provision of this License, you have permission to link or combine any covered work with a work licensed -under version 3 of the GNU Affero General Public License into a single +under version 3 of the GNU General Public License into a single combined work, and to convey the resulting work. The terms of this License will continue to apply to the part which is the covered work, -but the special requirements of the GNU Affero General Public License, -section 13, concerning interaction through a network will apply to the -combination as such. +but the work with which it is combined will remain governed by version +3 of the GNU General Public License. 14. Revised Versions of this License. The Free Software Foundation may publish revised and/or new versions of -the GNU General Public License from time to time. Such new versions will -be similar in spirit to the present version, but may differ in detail to +the GNU Affero General Public License from time to time. Such new versions +will be similar in spirit to the present version, but may differ in detail to address new problems or concerns. Each version is given a distinguishing version number. If the -Program specifies that a certain numbered version of the GNU General +Program specifies that a certain numbered version of the GNU Affero General Public License "or any later version" applies to it, you have the option of following the terms and conditions either of that numbered version or of any later version published by the Free Software Foundation. If the Program does not specify a version number of the -GNU General Public License, you may choose any version ever published +GNU Affero General Public License, you may choose any version ever published by the Free Software Foundation. If the Program specifies that a proxy can decide which future -versions of the GNU General Public License can be used, that proxy's +versions of the GNU Affero General Public License can be used, that proxy's public statement of acceptance of a version permanently authorizes you to choose that version for the Program. @@ -635,40 +633,29 @@ the "copyright" line and a pointer to where the full notice is found. Copyright (C) This program is free software: you can redistribute it and/or modify - it under the terms of the GNU General Public License as published by - the Free Software Foundation, either version 3 of the License, or + it under the terms of the GNU Affero General Public License as published + by the Free Software Foundation, either version 3 of the License, or (at your option) any later version. This program is distributed in the hope that it will be useful, but WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the - GNU General Public License for more details. + GNU Affero General Public License for more details. - You should have received a copy of the GNU General Public License + You should have received a copy of the GNU Affero General Public License along with this program. If not, see . Also add information on how to contact you by electronic and paper mail. - If the program does terminal interaction, make it output a short -notice like this when it starts in an interactive mode: - - Copyright (C) - This program comes with ABSOLUTELY NO WARRANTY; for details type `show w'. - This is free software, and you are welcome to redistribute it - under certain conditions; type `show c' for details. - -The hypothetical commands `show w' and `show c' should show the appropriate -parts of the General Public License. Of course, your program's commands -might be different; for a GUI interface, you would use an "about box". + If your software can interact with users remotely through a computer +network, you should also make sure that it provides a way for users to +get its source. For example, if your program is a web application, its +interface could display a "Source" link that leads users to an archive +of the code. There are many ways you could offer source, and different +solutions will be better for different programs; see section 13 for the +specific requirements. You should also get your employer (if you work as a programmer) or school, if any, to sign a "copyright disclaimer" for the program, if necessary. -For more information on this, and how to apply and follow the GNU GPL, see +For more information on this, and how to apply and follow the GNU AGPL, see . - - The GNU General Public License does not permit incorporating your program -into proprietary programs. If your program is a subroutine library, you -may consider it more useful to permit linking proprietary applications with -the library. If this is what you want to do, use the GNU Lesser General -Public License instead of this License. But first, please read -. diff --git a/README.md b/README.md index eb07164..ca0644c 100644 --- a/README.md +++ b/README.md @@ -2,6 +2,10 @@ ## Project Overview +

+ log +

+ **Mantis** is a high-performance network security solution that combines eBPF XDP technology with deep learning models to provide advanced network protection. The system operates as a standalone network appliance that can run on any Ubuntu-based system with compatible network hardware. ## Core Technologies @@ -10,6 +14,9 @@ - **Deep Learning Models** - Identifies and predicts potential network attacks with intelligent threat detection - **Hardware Integration** - Designed to work with Intel i350 T2 and similar enterprise-grade network interface cards +## System Architecture +![Architecture](.github/images/architecture.png) + [//]: # (## Functional Modules) [//]: # () diff --git a/mantis-frontend b/mantis-frontend index 66ae9de..2b9c85d 160000 --- a/mantis-frontend +++ b/mantis-frontend @@ -1 +1 @@ -Subproject commit 66ae9de967702ef02ab47baa9597ab8f16e3ab91 +Subproject commit 2b9c85de173c2810a8ffca183ddddcc751f4e227 diff --git a/mantis/src/core/app_state.rs b/mantis/src/core/app_state.rs index 4be2a92..e580104 100644 --- a/mantis/src/core/app_state.rs +++ b/mantis/src/core/app_state.rs @@ -6,6 +6,7 @@ use crate::core::infrastructure::app_config::AppConfig; use crate::core::infrastructure::app_db::AppDB; use crate::core::infrastructure::detection_alert::DetectionAlert; use crate::core::infrastructure::health::SystemHealth; +use crate::core::infrastructure::log_broadcaster::LogBroadcaster; use crate::detection::ml::config_loader::InferenceConfig; #[derive(Clone)] @@ -17,4 +18,5 @@ pub struct AppState { pub health: Arc, pub detection_alert: Arc, pub app_db: Option>, + pub log_broadcaster: Arc, } diff --git a/mantis/src/core/infrastructure/app_db.rs b/mantis/src/core/infrastructure/app_db.rs index 7501d7d..ace4148 100644 --- a/mantis/src/core/infrastructure/app_db.rs +++ b/mantis/src/core/infrastructure/app_db.rs @@ -5,6 +5,7 @@ use argon2::Argon2; use argon2::password_hash::{PasswordHash, PasswordHasher, PasswordVerifier, SaltString, rand_core::OsRng}; use macros::log; use rusqlite::{Connection, params}; +use uuid::Uuid; use crate::model::error::Error; use crate::model::error::auth::AuthError; @@ -55,22 +56,24 @@ impl AppDB { Ok(Self { conn: Mutex::new(conn) }) } - pub fn ensure_default_admin(&self, default_password: &str) -> Result<(), Error> { + pub fn has_any_account(&self) -> Result { let conn = self.conn.lock().unwrap(); let count: i64 = conn .query_row("SELECT COUNT(*) FROM accounts", [], |row| row.get(0)) .map_err(|e| AuthError::DBError { msg: e.to_string() })?; + Ok(count > 0) + } - if count == 0 { - let hash = hash_password(default_password)?; - let now = chrono::Utc::now().timestamp(); - conn.execute( - "INSERT INTO accounts (id, username, password_hash, role, created_at) VALUES (?1, ?2, ?3, ?4, ?5)", - params!["admin", "admin", hash, "admin", now], - ) - .map_err(|e| AuthError::DBError { msg: e.to_string() })?; - log!(AuthLog::DefaultAdminCreated); - } + pub fn create_account(&self, username: &str, password: &str, role: &str) -> Result<(), Error> { + let hash = hash_password(password)?; + let id = Uuid::new_v4().to_string(); + let now = chrono::Utc::now().timestamp(); + let conn = self.conn.lock().unwrap(); + conn.execute( + "INSERT INTO accounts (id, username, password_hash, role, created_at) VALUES (?1, ?2, ?3, ?4, ?5)", + params![id, username, hash, role, now], + ) + .map_err(|e| AuthError::DBError { msg: e.to_string() })?; Ok(()) } diff --git a/mantis/src/core/infrastructure/log_broadcaster.rs b/mantis/src/core/infrastructure/log_broadcaster.rs new file mode 100644 index 0000000..a8eaedb --- /dev/null +++ b/mantis/src/core/infrastructure/log_broadcaster.rs @@ -0,0 +1,53 @@ +use std::collections::VecDeque; +use std::sync::{Arc, Mutex}; + +use serde::{Deserialize, Serialize}; +use tokio::sync::broadcast; + +const RING_BUFFER_SIZE: usize = 500; +const CHANNEL_CAPACITY: usize = 512; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogRecord { + pub timestamp: String, + pub level: String, + pub message: String, +} + +pub struct LogBroadcaster { + tx: broadcast::Sender, + recent: Mutex>, +} + +impl LogBroadcaster { + pub fn new() -> Arc { + let (tx, _) = broadcast::channel(CHANNEL_CAPACITY); + Arc::new(Self { + tx, + recent: Mutex::new(VecDeque::with_capacity(RING_BUFFER_SIZE)), + }) + } + + pub fn emit(&self, record: LogRecord) { + if let Ok(mut buf) = self.recent.lock() { + if buf.len() >= RING_BUFFER_SIZE { + buf.pop_front(); + } + buf.push_back(record.clone()); + } + if self.tx.receiver_count() > 0 { + let _ = self.tx.send(record); + } + } + + pub fn subscribe(&self) -> broadcast::Receiver { + self.tx.subscribe() + } + + pub fn recent_logs(&self) -> Vec { + self.recent + .lock() + .map(|buf| buf.iter().cloned().collect()) + .unwrap_or_default() + } +} diff --git a/mantis/src/core/infrastructure/mod.rs b/mantis/src/core/infrastructure/mod.rs index 4227c78..c24fe84 100644 --- a/mantis/src/core/infrastructure/mod.rs +++ b/mantis/src/core/infrastructure/mod.rs @@ -3,6 +3,7 @@ pub mod app_db; pub mod detection_alert; pub mod geoip; pub mod health; +pub mod log_broadcaster; use std::path::PathBuf; use std::sync::Arc; @@ -37,11 +38,16 @@ pub struct AppServices { pub ml_engine: Arc, pub suricata_engine: Option>, pub app_db: Option>, + pub log_broadcaster: Arc, shutdowns: SegQueue>, } impl AppServices { - pub fn new(app_config: Arc, inference_config: Arc) -> Result { + pub fn new( + app_config: Arc, + inference_config: Arc, + log_broadcaster: Arc, + ) -> Result { let health = SystemHealth::new(app_config.clone())?; let ml_models = Arc::new(MLModels::load_models(&app_config)?); @@ -105,7 +111,6 @@ impl AppServices { let app_db = if let Some(ref auth) = app_config.auth { let db_path = PathBuf::from(env!("DB_PATH")).join("app.db"); let db = AppDB::open(db_path, &auth.db_key)?; - db.ensure_default_admin(&auth.default_admin_password)?; Some(Arc::new(db)) } else { None @@ -119,6 +124,7 @@ impl AppServices { ml_engine, suricata_engine, app_db, + log_broadcaster, shutdowns: SegQueue::new(), }) } diff --git a/mantis/src/core/system.rs b/mantis/src/core/system.rs index 69d0798..f8ee8b9 100644 --- a/mantis/src/core/system.rs +++ b/mantis/src/core/system.rs @@ -22,8 +22,9 @@ use crate::model::error::misc::MiscError; use crate::model::log::ml::MLLog; use crate::model::log::system::SystemLog; use crate::utils::logging::Logging; +use crate::core::infrastructure::log_broadcaster::LogBroadcaster; use crate::web::api::default::default_route; -use crate::web::api::{auth, control, detection_alert, health, misc}; +use crate::web::api::{auth, control, detection_alert, health, logs, misc}; pub struct System { pub app_config: Arc, @@ -40,7 +41,8 @@ pub struct System { impl System { pub async fn new() -> Result { - Logging::initialize()?; + let log_broadcaster = LogBroadcaster::new(); + Logging::initialize(log_broadcaster.clone())?; log!(SystemLog::Initializing); @@ -59,7 +61,11 @@ impl System { &mut egress_ebpf, )?); - let app_services = Arc::new(AppServices::new(app_config.clone(), inference_config.clone())?); + let app_services = Arc::new(AppServices::new( + app_config.clone(), + inference_config.clone(), + log_broadcaster, + )?); let system = System { app_config, @@ -93,9 +99,10 @@ impl System { .run(app_services.ml_engine.clone(), app_services.suricata_engine.clone()) .await?; app_services.run().await?; - self.run_http_server().await?; log!(SystemLog::InitializeComplete); + + self.run_http_server().await?; Ok(()) } @@ -153,12 +160,14 @@ impl System { health: self.app_services.health.clone(), detection_alert: self.app_services.detection_alert.clone(), app_db: self.app_services.app_db.clone(), + log_broadcaster: self.app_services.log_broadcaster.clone(), }; let app = Router::new() .nest("/ebpf", control::router()) .nest("/detection", detection_alert::router()) .nest("/health", health::router()) + .nest("/logs", logs::router()) .nest("/misc", misc::router()) .nest("/auth", auth::router()) .fallback(default_route) diff --git a/mantis/src/model/log/auth.rs b/mantis/src/model/log/auth.rs index bffa47a..8b33089 100644 --- a/mantis/src/model/log/auth.rs +++ b/mantis/src/model/log/auth.rs @@ -6,8 +6,8 @@ loggable! { #[error("Auth DB initialized at {path}")] DbInitialized { path: String } => tracing::Level::INFO, - #[error("Default admin account created")] - DefaultAdminCreated => tracing::Level::INFO, + #[error("First admin account registered: {username}")] + FirstAdminRegistered { username: String } => tracing::Level::INFO, #[error("Login successful for user: {username}")] LoginSuccess { username: String } => tracing::Level::INFO, diff --git a/mantis/src/utils/logging.rs b/mantis/src/utils/logging.rs index 18b0904..5c967dc 100644 --- a/mantis/src/utils/logging.rs +++ b/mantis/src/utils/logging.rs @@ -1,4 +1,5 @@ use std::fs; +use std::sync::Arc; use tracing::Level; use tracing_appender::rolling::{RollingFileAppender, Rotation}; @@ -6,13 +7,15 @@ use tracing_subscriber::filter::EnvFilter; use tracing_subscriber::layer::SubscriberExt; use tracing_subscriber::util::SubscriberInitExt; +use crate::core::infrastructure::log_broadcaster::LogBroadcaster; use crate::model::error::Error; use crate::model::error::io::IOError; +use crate::utils::tracing_layer::BroadcastLayer; pub struct Logging; impl Logging { - pub fn initialize() -> Result<(), Error> { + pub fn initialize(broadcaster: Arc) -> Result<(), Error> { let log_directory = "logs"; fs::create_dir_all(log_directory).map_err(|err| IOError::CreateDirectoryFailed(log_directory, err))?; @@ -33,6 +36,8 @@ impl Logging { .with_ansi(false) .with_writer(file_appender); + let broadcast_layer = BroadcastLayer::new(broadcaster); + let level = if cfg!(debug_assertions) { Level::DEBUG } else { @@ -42,6 +47,7 @@ impl Logging { tracing_subscriber::registry() .with(stdout_layer) .with(file_layer) + .with(broadcast_layer) .with(EnvFilter::from_default_env().add_directive(level.into())) .init(); diff --git a/mantis/src/utils/mod.rs b/mantis/src/utils/mod.rs index b3ee894..6e26a27 100644 --- a/mantis/src/utils/mod.rs +++ b/mantis/src/utils/mod.rs @@ -3,5 +3,6 @@ pub mod cpu_affinity; pub mod logging; pub mod packet_parser; pub mod static_files; +pub mod tracing_layer; pub mod ip_address; diff --git a/mantis/src/utils/tracing_layer.rs b/mantis/src/utils/tracing_layer.rs new file mode 100644 index 0000000..19d08b6 --- /dev/null +++ b/mantis/src/utils/tracing_layer.rs @@ -0,0 +1,55 @@ +use std::sync::Arc; + +use tracing::Subscriber; +use tracing_subscriber::Layer; +use tracing_subscriber::layer::Context; + +use crate::core::infrastructure::log_broadcaster::{LogBroadcaster, LogRecord}; + +pub struct BroadcastLayer { + broadcaster: Arc, +} + +impl BroadcastLayer { + pub fn new(broadcaster: Arc) -> Self { + Self { broadcaster } + } +} + +struct MessageVisitor { + message: String, +} + +impl tracing::field::Visit for MessageVisitor { + fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) { + if field.name() == "message" { + self.message = format!("{value:?}"); + } + } + + fn record_str(&mut self, field: &tracing::field::Field, value: &str) { + if field.name() == "message" { + self.message = value.to_string(); + } + } +} + +impl Layer for BroadcastLayer { + fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) { + let level = event.metadata().level().to_string(); + let mut visitor = MessageVisitor { message: String::new() }; + event.record(&mut visitor); + + if visitor.message.is_empty() { + return; + } + + let record = LogRecord { + timestamp: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true), + level, + message: visitor.message, + }; + + self.broadcaster.emit(record); + } +} diff --git a/mantis/src/web/api/auth.rs b/mantis/src/web/api/auth.rs index 328d055..e58419c 100644 --- a/mantis/src/web/api/auth.rs +++ b/mantis/src/web/api/auth.rs @@ -18,11 +18,24 @@ use crate::web::middleware::auth::{AuthenticatedUser, Claims}; pub fn router() -> Router { Router::new() + .route("/status", get(status)) + .route("/register", post(register)) .route("/login", post(login)) .route("/me", get(me)) .route("/logout", post(logout)) } +#[derive(Serialize)] +struct StatusResponse { + initialized: bool, +} + +#[derive(Deserialize)] +struct RegisterRequest { + username: String, + password: String, +} + #[derive(Deserialize)] struct LoginRequest { username: String, @@ -41,6 +54,71 @@ struct MeResponse { role: String, } +async fn status(State(state): State) -> impl IntoResponse { + if state.app_config.auth.is_none() { + return (StatusCode::NOT_IMPLEMENTED, "Auth not configured").into_response(); + } + let db = match &state.app_db { + Some(db) => db, + None => return (StatusCode::INTERNAL_SERVER_ERROR, "Database not available").into_response(), + }; + match db.has_any_account() { + Ok(initialized) => Json(StatusResponse { initialized }).into_response(), + Err(_) => (StatusCode::INTERNAL_SERVER_ERROR, "Database error").into_response(), + } +} + +async fn register(State(state): State, Json(body): Json) -> impl IntoResponse { + let auth_cfg = match state.app_config.auth.as_ref() { + Some(c) => c, + None => return (StatusCode::NOT_IMPLEMENTED, "Auth not configured").into_response(), + }; + let db = match &state.app_db { + Some(db) => db, + None => return (StatusCode::INTERNAL_SERVER_ERROR, "Database not available").into_response(), + }; + + match db.has_any_account() { + Ok(true) => return (StatusCode::FORBIDDEN, "System already initialized").into_response(), + Ok(false) => {} + Err(_) => return (StatusCode::INTERNAL_SERVER_ERROR, "Database error").into_response(), + } + + if body.username.trim().is_empty() || body.password.len() < 8 { + return (StatusCode::BAD_REQUEST, "Username required and password must be at least 8 characters").into_response(); + } + + if db.create_account(body.username.trim(), &body.password, "admin").is_err() { + return (StatusCode::INTERNAL_SERVER_ERROR, "Failed to create account").into_response(); + } + + log!(AuthLog::FirstAdminRegistered { + username: body.username.trim().to_string() + }); + + let account = match db.find_account_by_username(body.username.trim()) { + Ok(Some(a)) => a, + _ => return (StatusCode::INTERNAL_SERVER_ERROR, "Account lookup failed").into_response(), + }; + + let exp = (Utc::now().timestamp() as usize) + (auth_cfg.token_ttl_secs as usize); + let claims = Claims { + sub: account.id, + username: account.username, + role: account.role, + exp, + }; + + match encode( + &Header::default(), + &claims, + &EncodingKey::from_secret(auth_cfg.jwt_secret.as_bytes()), + ) { + Ok(token) => Json(LoginResponse { token }).into_response(), + Err(_) => (StatusCode::INTERNAL_SERVER_ERROR, "Token generation failed").into_response(), + } +} + async fn login(State(state): State, Json(body): Json) -> impl IntoResponse { let auth_cfg = match state.app_config.auth.as_ref() { Some(c) => c, diff --git a/mantis/src/web/api/logs.rs b/mantis/src/web/api/logs.rs new file mode 100644 index 0000000..716066d --- /dev/null +++ b/mantis/src/web/api/logs.rs @@ -0,0 +1,25 @@ +use axum::Router; +use axum::extract::{State, WebSocketUpgrade}; +use axum::response::IntoResponse; +use axum::routing::get; + +use crate::core::app_state::AppState; +use crate::web::websocket::log_websocket::handle_log_stream; + +pub fn router() -> Router { + Router::new() + .route("/recent", get(recent_logs)) + .route("/websocket", get(log_websocket)) +} + +async fn recent_logs(State(state): State) -> impl IntoResponse { + let logs = state.log_broadcaster.recent_logs(); + axum::Json(logs) +} + +async fn log_websocket( + ws: WebSocketUpgrade, + State(state): State, +) -> impl IntoResponse { + ws.on_upgrade(move |socket| handle_log_stream(socket, state.log_broadcaster)) +} diff --git a/mantis/src/web/api/mod.rs b/mantis/src/web/api/mod.rs index 37af96a..c6341a5 100644 --- a/mantis/src/web/api/mod.rs +++ b/mantis/src/web/api/mod.rs @@ -3,4 +3,5 @@ pub mod control; pub mod default; pub mod detection_alert; pub mod health; +pub mod logs; pub mod misc; diff --git a/mantis/src/web/websocket/log_websocket.rs b/mantis/src/web/websocket/log_websocket.rs new file mode 100644 index 0000000..9048ca9 --- /dev/null +++ b/mantis/src/web/websocket/log_websocket.rs @@ -0,0 +1,64 @@ +use std::sync::Arc; + +use axum::extract::ws::{Message, WebSocket}; +use futures_util::{SinkExt, StreamExt}; +use macros::log; +use tokio::sync::broadcast; + +use crate::core::infrastructure::log_broadcaster::LogBroadcaster; +use crate::model::error::http::HttpError; +use crate::model::error::misc::MiscError; +use crate::model::log::http::HttpLog; + +pub async fn handle_log_stream(socket: WebSocket, broadcaster: Arc) { + let (mut sender, mut receiver) = socket.split(); + + // Send buffered recent logs first + let recent = broadcaster.recent_logs(); + for record in recent { + match serde_json::to_string(&record) { + Ok(json) => { + if sender.send(Message::Text(json.into())).await.is_err() { + return; + } + } + Err(e) => log!(MiscError::SerializeError(e)), + } + } + + let mut rx = broadcaster.subscribe(); + + loop { + tokio::select! { + msg = receiver.next() => { + match msg { + Some(Ok(Message::Ping(data))) => { + if sender.send(Message::Pong(data)).await.is_err() { break; } + } + Some(Ok(Message::Close(_))) | None => break, + Some(Err(e)) => { + log!(HttpError::WebSocketError { msg: e.to_string() }); + break; + } + _ => {} + } + } + result = rx.recv() => { + match result { + Ok(record) => { + match serde_json::to_string(&record) { + Ok(json) => { + if sender.send(Message::Text(json.into())).await.is_err() { break; } + } + Err(e) => log!(MiscError::SerializeError(e)), + } + } + Err(broadcast::error::RecvError::Lagged(n)) => { + log!(HttpLog::WebSocketLaged { skipped: n }); + } + Err(broadcast::error::RecvError::Closed) => break, + } + } + } + } +} diff --git a/mantis/src/web/websocket/mod.rs b/mantis/src/web/websocket/mod.rs index 8ae1755..84b14cc 100644 --- a/mantis/src/web/websocket/mod.rs +++ b/mantis/src/web/websocket/mod.rs @@ -1,3 +1,4 @@ pub mod alert_websocket; pub mod flow_websocket; pub mod health_websocket; +pub mod log_websocket;