diff --git a/core/Cargo.lock b/core/Cargo.lock index 501f74f..5cf0292 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -1342,6 +1342,12 @@ dependencies = [ "tracing", ] +[[package]] +name = "lazy_static" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" + [[package]] name = "libc" version = "0.2.186" @@ -1396,6 +1402,15 @@ dependencies = [ "pkg-config", ] +[[package]] +name = "matchers" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d1525a2a28c7f4fa0fc98bb91ae755d1e2d1505079e05539e35bc876b5d65ae9" +dependencies = [ + "regex-automata", +] + [[package]] name = "memchr" version = "2.8.3" @@ -1446,6 +1461,15 @@ dependencies = [ "tempfile", ] +[[package]] +name = "nu-ansi-term" +version = "0.50.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "num-conv" version = "0.2.2" @@ -1557,6 +1581,9 @@ dependencies = [ "serde_yaml", "tokio", "tokio-tungstenite", + "tracing", + "tracing-log", + "tracing-subscriber", "urlencoding", "uuid", "winres", @@ -2198,6 +2225,15 @@ dependencies = [ "digest", ] +[[package]] +name = "sharded-slab" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f40ca3c46823713e0d4209592e8d6e826aa57e928f09752619fc696c499637f6" +dependencies = [ + "lazy_static", +] + [[package]] name = "shlex" version = "2.0.1" @@ -2364,6 +2400,15 @@ dependencies = [ "syn", ] +[[package]] +name = "thread_local" +version = "1.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ad99c4c6d32803332c548b1af0540b357b3f5fc0be8f6c6bfe8b2e6ae784070" +dependencies = [ + "cfg-if", +] + [[package]] name = "time" version = "0.3.53" @@ -2569,6 +2614,36 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "db97caf9d906fbde555dd62fa95ddba9eecfd14cb388e4f491a66d74cd5fb79a" dependencies = [ "once_cell", + "valuable", +] + +[[package]] +name = "tracing-log" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee855f1f400bd0e5c02d150ae5de3840039a3f54b025156404e34c23c03f47c3" +dependencies = [ + "log", + "once_cell", + "tracing-core", +] + +[[package]] +name = "tracing-subscriber" +version = "0.3.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb7f578e5945fb242538965c2d0b04418d38ec25c79d160cd279bf0731c8d319" +dependencies = [ + "matchers", + "nu-ansi-term", + "once_cell", + "regex-automata", + "sharded-slab", + "smallvec", + "thread_local", + "tracing", + "tracing-core", + "tracing-log", ] [[package]] @@ -2671,6 +2746,12 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "valuable" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" + [[package]] name = "vcpkg" version = "0.2.15" diff --git a/core/engine/Cargo.toml b/core/engine/Cargo.toml index ca22862..febb78d 100644 --- a/core/engine/Cargo.toml +++ b/core/engine/Cargo.toml @@ -14,6 +14,9 @@ serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" futures-util = "0.3" log = "0.4" +tracing = "0.1" +tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] } +tracing-log = "0.2" uuid = { version = "1.0", features = ["v4"] } urlencoding = "2.1" async-trait = "0.1" diff --git a/core/engine/src/ipc/bridge.rs b/core/engine/src/ipc/bridge.rs index 6553487..6be44e0 100644 --- a/core/engine/src/ipc/bridge.rs +++ b/core/engine/src/ipc/bridge.rs @@ -6,6 +6,9 @@ use tokio_tungstenite::{connect_async, tungstenite::Message, MaybeTlsStream, Web use std::sync::Arc; use tokio::sync::Mutex; +/// How often the engine sends a WebSocket-level Ping to keep the TCP connection alive. +pub const PING_INTERVAL_SECS: u64 = 30; + /// Auth info received from Neutralino via stdin or CLI flags. #[derive(Debug, Default, Deserialize)] pub struct AuthInfo { @@ -86,7 +89,7 @@ impl Bridge { urlencoding::encode(&auth.nl_connect_token) ); - log::info!("Connecting to WebSocket: {}", url); + tracing::info!("Connecting to Neutralino WebSocket server at port {}", auth.nl_port); let (ws, _) = connect_async(&url).await?; let (writer, reader) = ws.split(); @@ -114,7 +117,19 @@ impl Bridge { }; let text = serde_json::to_string(&msg)?; - log::info!("Sending typed event"); + let ev_name = event.event_name(); + match event { + super::events::OrbitEvent::ErrorOccurred { message } => { + tracing::error!(event = ev_name, error = %message, "Response to UI (Error)"); + } + super::events::OrbitEvent::CommandSucceeded { message } => { + tracing::info!(event = ev_name, message = %message, "Response to UI (Command Succeeded)"); + } + _ => { + tracing::debug!(event = ev_name, "Response to UI (Event Broadcast)"); + } + } + let mut w = writer.lock().await; w.send(Message::Text(text.into())).await?; Ok(()) @@ -125,6 +140,11 @@ impl Bridge { match reader.next().await { Some(Ok(Message::Text(text))) => { let msg: WsMessage = serde_json::from_str(&text)?; + if let Some(ref ev) = msg.event { + tracing::info!(event = %ev, id = ?msg.id, "Received request/event from UI"); + } else if let Some(ref method) = msg.method { + tracing::info!(method = %method, id = ?msg.id, "Received method from UI"); + } return Ok(msg); } Some(Ok(Message::Ping(data))) => { @@ -132,8 +152,14 @@ impl Bridge { w.send(Message::Pong(data)).await?; } Some(Ok(_)) => continue, - Some(Err(e)) => return Err(Box::new(e)), - None => return Err("WebSocket closed".into()), + Some(Err(e)) => { + tracing::error!("WebSocket error reading message: {:?}", e); + return Err(Box::new(e)); + } + None => { + tracing::warn!("WebSocket connection closed by server"); + return Err("WebSocket closed".into()); + } } } } diff --git a/core/engine/src/ipc/events.rs b/core/engine/src/ipc/events.rs index c4c6437..284b34c 100644 --- a/core/engine/src/ipc/events.rs +++ b/core/engine/src/ipc/events.rs @@ -151,3 +151,44 @@ pub enum OrbitEvent { data: serde_json::Value, }, } + +impl OrbitEvent { + pub fn event_name(&self) -> &'static str { + match self { + OrbitEvent::EngineConnected { .. } => "engineConnected", + OrbitEvent::Ping { .. } => "ping", + OrbitEvent::Pong { .. } => "pong", + OrbitEvent::NamespacesUpdated { .. } => "namespacesUpdated", + OrbitEvent::PodsUpdated { .. } => "podsUpdated", + OrbitEvent::DeploymentsUpdated { .. } => "deploymentsUpdated", + OrbitEvent::StatefulSetsUpdated { .. } => "statefulSetsUpdated", + OrbitEvent::DaemonSetsUpdated { .. } => "daemonSetsUpdated", + OrbitEvent::ReplicaSetsUpdated { .. } => "replicaSetsUpdated", + OrbitEvent::JobsUpdated { .. } => "jobsUpdated", + OrbitEvent::CronJobsUpdated { .. } => "cronJobsUpdated", + OrbitEvent::ClustersUpdated { .. } => "clustersUpdated", + OrbitEvent::ActiveClusterChanged { .. } => "activeClusterChanged", + OrbitEvent::UserProfileUpdated { .. } => "userProfileUpdated", + OrbitEvent::NodesUpdated { .. } => "nodesUpdated", + OrbitEvent::ServicesUpdated { .. } => "servicesUpdated", + OrbitEvent::IngressesUpdated { .. } => "ingressesUpdated", + OrbitEvent::ConfigMapsUpdated { .. } => "configMapsUpdated", + OrbitEvent::SecretsUpdated { .. } => "secretsUpdated", + OrbitEvent::EventsUpdated { .. } => "eventsUpdated", + OrbitEvent::PersistentVolumesUpdated { .. } => "persistentVolumesUpdated", + OrbitEvent::PersistentVolumeClaimsUpdated { .. } => "persistentVolumeClaimsUpdated", + OrbitEvent::StorageClassesUpdated { .. } => "storageClassesUpdated", + OrbitEvent::PoliciesUpdated { .. } => "policiesUpdated", + OrbitEvent::ResourceUpdated { .. } => "resourceUpdated", + OrbitEvent::PodMetricsUpdated { .. } => "podMetricsUpdated", + OrbitEvent::ErrorOccurred { .. } => "errorOccurred", + OrbitEvent::LogLineReceived { .. } => "logLineReceived", + OrbitEvent::LogLinesChunkReceived { .. } => "logLinesChunkReceived", + OrbitEvent::UpdateCheckFinished { .. } => "updateCheckFinished", + OrbitEvent::CommandSucceeded { .. } => "commandSucceeded", + OrbitEvent::UpdateDownloadProgress { .. } => "updateDownloadProgress", + OrbitEvent::UpdateReady { .. } => "updateReady", + OrbitEvent::ResourceRawData { .. } => "resourceRawData", + } + } +} diff --git a/core/engine/src/ipc/handlers/mod.rs b/core/engine/src/ipc/handlers/mod.rs index d0f3908..6302f69 100644 --- a/core/engine/src/ipc/handlers/mod.rs +++ b/core/engine/src/ipc/handlers/mod.rs @@ -27,6 +27,7 @@ pub fn dispatch( token: String, manager: Arc>, ) { + tracing::info!(event = %event_name, "Dispatching UI request"); match event_name { "getClusters" => cluster::get_clusters(writer, token, manager), "getUserProfile" => cluster::get_user_profile(writer, token, manager), @@ -65,6 +66,8 @@ pub fn dispatch( "cloneIngress" => network::clone_ingress(data, writer, token, manager), "cloneDeployment" => workloads::clone_deployment(data, writer, token, manager), "rollbackDeployment" => workloads::rollback_deployment(data, writer, token, manager), - _ => {} + other => { + tracing::debug!(event = %other, "Unhandled UI event in dispatcher"); + } } } diff --git a/core/engine/src/kubernetes/workloads.rs b/core/engine/src/kubernetes/workloads.rs index a3db7fd..111abe1 100644 --- a/core/engine/src/kubernetes/workloads.rs +++ b/core/engine/src/kubernetes/workloads.rs @@ -385,6 +385,7 @@ pub async fn scale_resource( name: &str, replicas: i32, ) -> Result<(), kube::Error> { + tracing::info!(kind = %kind, namespace = %namespace, name = %name, replicas = replicas, "Sending request to Kubernetes: scale resource"); let patch = serde_json::json!({ "spec": { "replicas": replicas @@ -412,6 +413,7 @@ pub async fn scale_resource( code: 400, })), } + tracing::info!(kind = %kind, namespace = %namespace, name = %name, "Kubernetes request completed: scale resource"); Ok(()) } @@ -421,6 +423,7 @@ pub async fn redeploy_resource( kind: &str, name: &str, ) -> Result<(), kube::Error> { + tracing::info!(kind = %kind, namespace = %namespace, name = %name, "Sending request to Kubernetes: redeploy/restart resource"); let now = chrono::Utc::now().to_rfc3339(); let patch = serde_json::json!({ "spec": { @@ -455,6 +458,7 @@ pub async fn redeploy_resource( code: 400, })), } + tracing::info!(kind = %kind, namespace = %namespace, name = %name, "Kubernetes request completed: redeploy resource"); Ok(()) } @@ -463,8 +467,10 @@ pub async fn delete_pod( namespace: &str, name: &str, ) -> Result<(), kube::Error> { + tracing::info!(namespace = %namespace, name = %name, "Sending request to Kubernetes: delete pod"); let api: Api = Api::namespaced(client.clone(), namespace); api.delete(name, &DeleteParams::default()).await?; + tracing::info!(namespace = %namespace, name = %name, "Kubernetes request completed: delete pod"); Ok(()) } @@ -475,6 +481,7 @@ pub async fn update_images_resource( name: &str, containers: Vec, ) -> Result<(), kube::Error> { + tracing::info!(kind = %kind, namespace = %namespace, name = %name, "Sending request to Kubernetes: update container images"); let containers_patch: Vec = containers .into_iter() .map(|c| { @@ -518,6 +525,7 @@ pub async fn update_images_resource( })) } } + tracing::info!(kind = %kind, namespace = %namespace, name = %name, "Kubernetes request completed: update container images"); Ok(()) } @@ -528,6 +536,7 @@ pub async fn clone_deployment( new_name: &str, new_namespace: &str, ) -> Result<(), kube::Error> { + tracing::info!(source_namespace = %source_namespace, source_name = %source_name, new_namespace = %new_namespace, new_name = %new_name, "Sending request to Kubernetes: clone deployment"); let source_api: Api = Api::namespaced(client.clone(), source_namespace); let mut cloned = source_api.get(source_name).await?; @@ -574,6 +583,7 @@ pub async fn clone_deployment( let target_api: Api = Api::namespaced(client.clone(), new_namespace); target_api.create(&PostParams::default(), &cloned).await?; + tracing::info!(new_namespace = %new_namespace, new_name = %new_name, "Kubernetes request completed: clone deployment"); Ok(()) } @@ -656,6 +666,7 @@ pub async fn rollback_deployment( name: &str, target_revision: Option, ) -> Result<(), kube::Error> { + tracing::info!(namespace = %namespace, name = %name, target_revision = ?target_revision, "Sending request to Kubernetes: rollback deployment"); let deploy_api: Api = Api::namespaced(client.clone(), namespace); let deployment = deploy_api.get(name).await?; @@ -693,7 +704,7 @@ pub async fn rollback_deployment( let patch_params = PatchParams::default(); deploy_api.patch(name, &patch_params, &Patch::Merge(&patch)).await?; - + tracing::info!(namespace = %namespace, name = %name, "Kubernetes request completed: rollback deployment"); Ok(()) } diff --git a/core/engine/src/logger.rs b/core/engine/src/logger.rs new file mode 100644 index 0000000..871ffb5 --- /dev/null +++ b/core/engine/src/logger.rs @@ -0,0 +1,206 @@ +use std::fs::{File, OpenOptions}; +use std::io::{self, Write}; +use std::path::PathBuf; +use std::sync::{Arc, Mutex}; +use chrono::{DateTime, Local}; +use tracing_subscriber::fmt::writer::MakeWriterExt; +use tracing_subscriber::EnvFilter; + +struct Inner { + dir: PathBuf, + active_date: String, + file: Option, +} + +impl Inner { + fn new(dir: PathBuf) -> io::Result { + let today = Local::now().format("%Y-%m-%d").to_string(); + let log_path = dir.join("orbit.log"); + + // If orbit.log already exists, check if it was last modified on a previous day + if log_path.exists() { + if let Ok(metadata) = std::fs::metadata(&log_path) { + if let Ok(modified) = metadata.modified() { + let mod_datetime: DateTime = modified.into(); + let mod_date = mod_datetime.format("%Y-%m-%d").to_string(); + if mod_date != today { + // Rename existing orbit.log to orbit.log. + let archive_path = dir.join(format!("orbit.log.{}", mod_date)); + let _ = std::fs::rename(&log_path, &archive_path); + } + } + } + } + + let file = OpenOptions::new() + .create(true) + .append(true) + .open(&log_path)?; + + Ok(Self { + dir, + active_date: today, + file: Some(file), + }) + } + + fn check_rotate(&mut self) -> io::Result<()> { + let today = Local::now().format("%Y-%m-%d").to_string(); + if today != self.active_date { + // Drop current file handle so it can be safely renamed + self.file = None; + + let log_path = self.dir.join("orbit.log"); + let archive_path = self.dir.join(format!("orbit.log.{}", self.active_date)); + if log_path.exists() { + let _ = std::fs::rename(&log_path, &archive_path); + } + + self.active_date = today; + let file = OpenOptions::new() + .create(true) + .append(true) + .open(&log_path)?; + self.file = Some(file); + } + Ok(()) + } + + fn write_all(&mut self, buf: &[u8]) -> io::Result<()> { + self.check_rotate()?; + if let Some(ref mut f) = self.file { + f.write_all(buf)?; + f.flush()?; + } + Ok(()) + } +} + +#[derive(Clone)] +pub struct RotatingFileAppender { + inner: Arc>, +} + +impl RotatingFileAppender { + pub fn new(dir: PathBuf) -> Self { + let inner = Inner::new(dir).expect("Failed to initialize rotating file appender"); + Self { + inner: Arc::new(Mutex::new(inner)), + } + } +} + +impl Write for RotatingFileAppender { + fn write(&mut self, buf: &[u8]) -> io::Result { + let mut inner = self + .inner + .lock() + .map_err(|_| io::Error::new(io::ErrorKind::Other, "Lock poisoned"))?; + inner.write_all(buf)?; + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + let mut inner = self + .inner + .lock() + .map_err(|_| io::Error::new(io::ErrorKind::Other, "Lock poisoned"))?; + if let Some(ref mut f) = inner.file { + f.flush()?; + } + Ok(()) + } +} + +impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for RotatingFileAppender { + type Writer = Self; + + fn make_writer(&'a self) -> Self::Writer { + self.clone() + } +} + +/// Initialize tracing with file logging in `~/.orbit/orbit.log` and stdout. +pub fn init() -> Result<(), Box> { + let config_dir = crate::config::OrbitConfig::config_dir() + .ok_or_else(|| "Could not determine Orbit config directory for logs")?; + std::fs::create_dir_all(&config_dir)?; + + let appender = RotatingFileAppender::new(config_dir.clone()); + let writer = appender.and(io::stdout); + + let env_filter = EnvFilter::try_from_default_env() + .unwrap_or_else(|_| EnvFilter::new("info,orbit_engine=info")); + + tracing_subscriber::fmt() + .with_env_filter(env_filter) + .with_writer(writer) + .with_ansi(false) + .init(); + + let _ = tracing_log::LogTracer::init(); + + tracing::info!( + "Orbit Activity Logger initialized. Active log: {:?}", + config_dir.join("orbit.log") + ); + + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_rotating_file_appender_writes() { + let temp_dir = std::env::temp_dir().join(format!("orbit_test_log_{}", uuid::Uuid::new_v4())); + std::fs::create_dir_all(&temp_dir).unwrap(); + + let mut appender = RotatingFileAppender::new(temp_dir.clone()); + writeln!(appender, "Test log line 1").unwrap(); + writeln!(appender, "Test log line 2").unwrap(); + + let log_file = temp_dir.join("orbit.log"); + assert!(log_file.exists()); + + let content = std::fs::read_to_string(&log_file).unwrap(); + assert!(content.contains("Test log line 1")); + assert!(content.contains("Test log line 2")); + + let _ = std::fs::remove_dir_all(temp_dir); + } + + #[test] + fn test_rotating_file_appender_rotates_on_date_change() { + let temp_dir = std::env::temp_dir().join(format!("orbit_test_log_rotate_{}", uuid::Uuid::new_v4())); + std::fs::create_dir_all(&temp_dir).unwrap(); + + let mut appender = RotatingFileAppender::new(temp_dir.clone()); + writeln!(appender, "Day 1 content").unwrap(); + + // Simulate previous active date + { + let mut inner = appender.inner.lock().unwrap(); + inner.active_date = "2026-08-20".to_string(); + } + + // Write new line, triggering rotation + writeln!(appender, "Day 2 content").unwrap(); + + // Check that orbit.log.2026-08-20 was created with Day 1 content + let rotated_file = temp_dir.join("orbit.log.2026-08-20"); + assert!(rotated_file.exists()); + let rotated_content = std::fs::read_to_string(&rotated_file).unwrap(); + assert!(rotated_content.contains("Day 1 content")); + + // Check that current orbit.log has Day 2 content + let active_file = temp_dir.join("orbit.log"); + assert!(active_file.exists()); + let active_content = std::fs::read_to_string(&active_file).unwrap(); + assert!(active_content.contains("Day 2 content")); + assert!(!active_content.contains("Day 1 content")); + + let _ = std::fs::remove_dir_all(temp_dir); + } +} diff --git a/core/engine/src/main.rs b/core/engine/src/main.rs index 2d553d5..9d49272 100644 --- a/core/engine/src/main.rs +++ b/core/engine/src/main.rs @@ -2,17 +2,76 @@ mod ipc; mod kubernetes; pub mod updater; pub mod config; +pub mod logger; use std::error::Error; use std::sync::Arc; -use tokio::sync::RwLock; -use ipc::bridge::{AuthInfo, Bridge}; +use std::time::Duration; +use futures_util::SinkExt; +use tokio::sync::{Mutex, RwLock}; +use tokio::time::MissedTickBehavior; +use tokio_tungstenite::tungstenite::Message; +use ipc::bridge::{AuthInfo, Bridge, PING_INTERVAL_SECS, WsWriter}; use ipc::events::OrbitEvent; use kubernetes::manager::KubeManager; +async fn broadcast_engine_ready( + writer: &Arc>, + token: &str, + kube_manager: &Arc>, +) { + let _ = Bridge::send_event( + writer, + token, + &OrbitEvent::EngineConnected { + status: "ready".to_string(), + message: "Orbit Engine is connected and ready.".to_string(), + }, + ).await; + + let r_manager = kube_manager.read().await; + let clusters = r_manager.get_clusters(); + let active_cluster_id = r_manager.active_context.clone(); + drop(r_manager); + + let _ = Bridge::send_event(writer, token, &OrbitEvent::ClustersUpdated { clusters }).await; + let _ = Bridge::send_event(writer, token, &OrbitEvent::ActiveClusterChanged { active_cluster_id }).await; +} + +async fn restart_watchers( + bridge: &Bridge, + kube_manager: &Arc>, +) { + let mut w_manager = kube_manager.write().await; + + // Cancel any existing watcher tasks + if let Some(cancel) = w_manager.watch_cancel.take() { + let _ = cancel.send(true); + } + + let client = w_manager.active_client.clone(); + let (tx, rx) = tokio::sync::watch::channel(false); + w_manager.watch_cancel = Some(tx); + drop(w_manager); + + if let Some(ref client) = client { + ipc::handlers::spawn_watchers( + client, + bridge.writer.clone(), + bridge.token.clone(), + rx, + ); + } +} + #[tokio::main] async fn main() -> Result<(), Box> { - println!("Orbit Core Engine starting up..."); + // Initialize persistent activity logging + if let Err(e) = logger::init() { + eprintln!("Warning: Failed to initialize file logger: {}", e); + } + + tracing::info!("Orbit Core Engine starting up..."); // Retrieve authentication details from stdin let mut auth = AuthInfo::from_stdin(); @@ -21,102 +80,112 @@ async fn main() -> Result<(), Box> { println!("Connecting to port: {}", auth.nl_port); - let mut bridge = Bridge::connect(&auth).await?; - println!("Orbit Core Engine connected to Neutralinojs WebSocket server."); - // Initialize KubeManager let kube_manager = Arc::new(RwLock::new(KubeManager::new().await)); - // Broadcast that the core is connected and ready - Bridge::send_event( - &bridge.writer, - &bridge.token, - &OrbitEvent::EngineConnected { - status: "ready".to_string(), - message: "Orbit Engine is connected and ready.".to_string(), - }, - ).await?; - - // Also send initial clusters list and active context - { - let mut manager = kube_manager.write().await; - let clusters = manager.get_clusters(); - let active_cluster_id = manager.active_context.clone(); - let client = manager.active_client.clone(); - - let _ = Bridge::send_event( - &bridge.writer, - &bridge.token, - &OrbitEvent::ClustersUpdated { clusters }, - ).await; - - let _ = Bridge::send_event( - &bridge.writer, - &bridge.token, - &OrbitEvent::ActiveClusterChanged { active_cluster_id }, - ).await; - - if let Some(ref client) = client { - if let Some(cancel) = manager.watch_cancel.take() { - let _ = cancel.send(true); - } - let (tx, rx) = tokio::sync::watch::channel(false); - manager.watch_cancel = Some(tx); - ipc::handlers::spawn_watchers(client, bridge.writer.clone(), bridge.token.clone(), rx.clone()); - } - } + let mut backoff_secs = 1u64; - // Message processing loop - loop { - let msg = match Bridge::read_message(&mut bridge.reader, &bridge.writer).await { - Ok(msg) => msg, + 'reconnect: loop { + let mut bridge = match Bridge::connect(&auth).await { + Ok(b) => { + tracing::info!("Orbit Core Engine connected to Neutralinojs WebSocket server."); + backoff_secs = 1; + b + } Err(e) => { - eprintln!("WebSocket error occurred or connection closed: {:?}", e); - break; + tracing::error!("Failed to connect to Neutralino WebSocket server: {:?}", e); + tokio::time::sleep(Duration::from_secs(backoff_secs)).await; + backoff_secs = next_backoff(backoff_secs); + continue 'reconnect; } }; - if msg.event.as_deref() == Some("windowClose") { - log::info!("Received windowClose, shutting down."); - break; + // Broadcast that the core is connected and ready along with initial clusters & context + broadcast_engine_ready(&bridge.writer, &bridge.token, &kube_manager).await; + + // Restart watchers with the new bridge writer + restart_watchers(&bridge, &kube_manager).await; + + let mut ping_interval = tokio::time::interval(Duration::from_secs(PING_INTERVAL_SECS)); + ping_interval.set_missed_tick_behavior(MissedTickBehavior::Skip); + + // Message processing loop + loop { + tokio::select! { + _ = ping_interval.tick() => { + let mut w = bridge.writer.lock().await; + if let Err(e) = w.send(Message::Ping(vec![].into())).await { + tracing::warn!("Ping failed (bridge dead), reconnecting: {:?}", e); + break; + } + } + result = Bridge::read_message(&mut bridge.reader, &bridge.writer) => { + let msg = match result { + Ok(msg) => msg, + Err(e) => { + tracing::warn!("WebSocket error occurred or connection closed: {:?}. Reconnecting...", e); + break; + } + }; + + if msg.event.as_deref() == Some("windowClose") { + tracing::info!("Received windowClose, shutting down."); + break 'reconnect; + } + + // Re-broadcast connection status when a client connects to ensure the frontend receives it + if msg.event.as_deref() == Some("appClientConnect") || msg.event.as_deref() == Some("clientConnect") { + let writer = bridge.writer.clone(); + let token = bridge.token.clone(); + let manager = kube_manager.clone(); + tokio::spawn(async move { + broadcast_engine_ready(&writer, &token, &manager).await; + }); + } + + // Dispatch all Kubernetes resource events to the handler module + if let Some(event_name) = msg.event.as_deref() { + ipc::handlers::dispatch( + event_name, + msg.data.clone(), + bridge.writer.clone(), + bridge.token.clone(), + kube_manager.clone(), + ); + } + } + } } - // Re-broadcast connection status when a client connects to ensure the frontend receives it - if msg.event.as_deref() == Some("appClientConnect") || msg.event.as_deref() == Some("clientConnect") { - let writer = bridge.writer.clone(); - let token = bridge.token.clone(); - let manager = kube_manager.clone(); - tokio::spawn(async move { - let _ = Bridge::send_event( - &writer, - &token, - &OrbitEvent::EngineConnected { - status: "ready".to_string(), - message: "Orbit Engine is connected and ready.".to_string(), - }, - ).await; - - let r_manager = manager.read().await; - let clusters = r_manager.get_clusters(); - let active_cluster_id = r_manager.active_context.clone(); - - let _ = Bridge::send_event(&writer, &token, &OrbitEvent::ClustersUpdated { clusters }).await; - let _ = Bridge::send_event(&writer, &token, &OrbitEvent::ActiveClusterChanged { active_cluster_id }).await; - }); - } + tracing::info!("Bridge disconnected. Reconnecting in {}s...", backoff_secs); + tokio::time::sleep(Duration::from_secs(backoff_secs)).await; + backoff_secs = next_backoff(backoff_secs); + } - // Dispatch all Kubernetes resource events to the handler module - if let Some(event_name) = msg.event.as_deref() { - ipc::handlers::dispatch( - event_name, - msg.data.clone(), - bridge.writer.clone(), - bridge.token.clone(), - kube_manager.clone(), - ); - } + Ok(()) +} +fn next_backoff(current_backoff: u64) -> u64 { + (current_backoff * 2).min(30) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_next_backoff() { + assert_eq!(next_backoff(1), 2); + assert_eq!(next_backoff(2), 4); + assert_eq!(next_backoff(4), 8); + assert_eq!(next_backoff(8), 16); + assert_eq!(next_backoff(16), 30); + assert_eq!(next_backoff(30), 30); } - Ok(()) + #[test] + fn test_ping_interval_constant() { + assert_eq!(PING_INTERVAL_SECS, 30); + } } +