Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 81 additions & 0 deletions core/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions core/engine/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
34 changes: 30 additions & 4 deletions core/engine/src/ipc/bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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();

Expand Down Expand Up @@ -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(())
Expand All @@ -125,15 +140,26 @@ 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))) => {
let mut w = writer.lock().await;
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());
}
}
}
}
Expand Down
41 changes: 41 additions & 0 deletions core/engine/src/ipc/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",
}
}
}
5 changes: 4 additions & 1 deletion core/engine/src/ipc/handlers/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ pub fn dispatch(
token: String,
manager: Arc<RwLock<KubeManager>>,
) {
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),
Expand Down Expand Up @@ -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");
}
}
}
13 changes: 12 additions & 1 deletion core/engine/src/kubernetes/workloads.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(())
}

Expand All @@ -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": {
Expand Down Expand Up @@ -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(())
}

Expand All @@ -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<Pod> = Api::namespaced(client.clone(), namespace);
api.delete(name, &DeleteParams::default()).await?;
tracing::info!(namespace = %namespace, name = %name, "Kubernetes request completed: delete pod");
Ok(())
}

Expand All @@ -475,6 +481,7 @@ pub async fn update_images_resource(
name: &str,
containers: Vec<models::ContainerImageInfo>,
) -> Result<(), kube::Error> {
tracing::info!(kind = %kind, namespace = %namespace, name = %name, "Sending request to Kubernetes: update container images");
let containers_patch: Vec<serde_json::Value> = containers
.into_iter()
.map(|c| {
Expand Down Expand Up @@ -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(())
}

Expand All @@ -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<Deployment> = Api::namespaced(client.clone(), source_namespace);
let mut cloned = source_api.get(source_name).await?;

Expand Down Expand Up @@ -574,6 +583,7 @@ pub async fn clone_deployment(

let target_api: Api<Deployment> = 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(())
}

Expand Down Expand Up @@ -656,6 +666,7 @@ pub async fn rollback_deployment(
name: &str,
target_revision: Option<i64>,
) -> Result<(), kube::Error> {
tracing::info!(namespace = %namespace, name = %name, target_revision = ?target_revision, "Sending request to Kubernetes: rollback deployment");
let deploy_api: Api<Deployment> = Api::namespaced(client.clone(), namespace);
let deployment = deploy_api.get(name).await?;

Expand Down Expand Up @@ -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(())
}

Expand Down
Loading