diff --git a/.github/workflows/core-build-checks.yml b/.github/workflows/core-build-checks.yml index a738e619..e02c377e 100644 --- a/.github/workflows/core-build-checks.yml +++ b/.github/workflows/core-build-checks.yml @@ -51,6 +51,21 @@ jobs: working-directory: core run: cargo test --locked -p cortexbrain-common --test public_contracts --no-default-features --features network-structs,monitoring-structs + test-cli: + name: Test CLI + runs-on: ubuntu-latest + + steps: + - name: Checkout code + uses: actions/checkout@v4 + + - name: Install prerequisites + run: sudo apt-get install -y protobuf-compiler + + - name: Test CLI + working-directory: cli + run: cargo test --locked + build-core-components: runs-on: ubuntu-latest diff --git a/Examples/run-with-docker/docker-compose.yaml b/Examples/run-with-docker/docker-compose.yaml index 8d088a08..1bebcb37 100644 --- a/Examples/run-with-docker/docker-compose.yaml +++ b/Examples/run-with-docker/docker-compose.yaml @@ -49,7 +49,7 @@ services: # - cortexflow otel-collector: - image: otel/opentelemetry-collector:0.95.0 + image: otel/opentelemetry-collector-contrib:0.95.0 container_name: otel-collector command: - "--config=/conf/otel-collector-config.yaml" diff --git a/Examples/run-with-docker/otel-collector-config.yaml b/Examples/run-with-docker/otel-collector-config.yaml index 91cae017..b1e893be 100644 --- a/Examples/run-with-docker/otel-collector-config.yaml +++ b/Examples/run-with-docker/otel-collector-config.yaml @@ -11,6 +11,13 @@ processors: limit_mib: 1500 spike_limit_mib: 512 check_interval: 5s + transform/identity: + error_mode: ignore + metric_statements: + - context: datapoint + statements: + - set(attributes["container.name"], attributes["k8s.pod.name"]) + where attributes["container.name"] == nil and attributes["k8s.pod.name"] != nil exporters: logging: {} @@ -30,5 +37,5 @@ service: exporters: [logging] metrics: receivers: [otlp] - processors: [memory_limiter] + processors: [memory_limiter, transform/identity] exporters: [logging, prometheus] diff --git a/cli/Cargo.lock b/cli/Cargo.lock index e64cbf66..62be59d5 100644 --- a/cli/Cargo.lock +++ b/cli/Cargo.lock @@ -377,18 +377,23 @@ name = "cortexflow-cli" version = "0.1.5" dependencies = [ "anyhow", + "bytes", "clap", "colored", "cortexflow_agent_api", "directories", + "http", + "http-body-util", "k8s-openapi", "kube", "prost", "prost-types", "serde", + "serde_json", "tokio", "tonic", "tonic-reflection", + "tower", "tracing", ] diff --git a/cli/Cargo.toml b/cli/Cargo.toml index 11041c39..abd55f56 100644 --- a/cli/Cargo.toml +++ b/cli/Cargo.toml @@ -26,6 +26,13 @@ cortexflow_agent_api = {version="0.1.2",features = ["client"]} kube = "2.0.1" k8s-openapi = {version = "0.26.0", features = ["v1_34"]} +[dev-dependencies] +serde_json = "1.0" +tower = { version = "0.5", features = ["util"] } +http = "1" +http-body-util = "0.1" +bytes = "1" + [[bin]] name = "cfcli" path = "src/main.rs" diff --git a/cli/src/command_runner.rs b/cli/src/command_runner.rs new file mode 100644 index 00000000..2a6489c1 --- /dev/null +++ b/cli/src/command_runner.rs @@ -0,0 +1,76 @@ +use std::process::{Command, Output}; + +// docs: +// +// Abstraction over external process execution (kubectl, cargo, ...) so callers +// can inject a stub in tests instead of shelling out to a real binary/cluster. + +pub trait CommandRunner { + fn run(&self, program: &str, args: &[String]) -> std::io::Result; +} + +pub struct RealCommandRunner; + +impl CommandRunner for RealCommandRunner { + fn run(&self, program: &str, args: &[String]) -> std::io::Result { + Command::new(program).args(args).output() + } +} + +#[cfg(test)] +pub mod test_support { + use super::CommandRunner; + use std::process::{ExitStatus, Output}; + + #[cfg(unix)] + fn exit_status(success: bool) -> ExitStatus { + use std::os::unix::process::ExitStatusExt; + ExitStatus::from_raw(if success { 0 } else { 1 << 8 }) + } + + // Stubs process output so tests never shell out to a real binary or cluster. + pub struct StubCommandRunner { + result: std::io::Result, + } + + impl StubCommandRunner { + pub fn success(stdout: &str) -> Self { + Self { + result: Ok(Output { + status: exit_status(true), + stdout: stdout.as_bytes().to_vec(), + stderr: Vec::new(), + }), + } + } + + pub fn failure(stderr: &str) -> Self { + Self { + result: Ok(Output { + status: exit_status(false), + stdout: Vec::new(), + stderr: stderr.as_bytes().to_vec(), + }), + } + } + + pub fn io_error() -> Self { + Self { + result: Err(std::io::Error::other("command not found")), + } + } + } + + impl CommandRunner for StubCommandRunner { + fn run(&self, _program: &str, _args: &[String]) -> std::io::Result { + match &self.result { + Ok(output) => Ok(Output { + status: output.status, + stdout: output.stdout.clone(), + stderr: output.stderr.clone(), + }), + Err(_) => Err(std::io::Error::other("command not found")), + } + } + } +} diff --git a/cli/src/errors.rs b/cli/src/errors.rs index b813e35b..fde81563 100644 --- a/cli/src/errors.rs +++ b/cli/src/errors.rs @@ -122,3 +122,18 @@ impl fmt::Display for CliError { } } } + +# [cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_display_base_error() { + let err = CliError::BaseError { + reason: "some reason".to_string(), + }; + let err_str = format!("{}", err); + assert!(err_str.contains("An error occured. Reason:")); + assert!(err_str.contains("some reason")); + } +} \ No newline at end of file diff --git a/cli/src/essential.rs b/cli/src/essential.rs index 5ca01b9a..a386f3a3 100644 --- a/cli/src/essential.rs +++ b/cli/src/essential.rs @@ -1,3 +1,4 @@ +use crate::command_runner::{CommandRunner, RealCommandRunner}; use crate::errors::CliError; use std::borrow::Cow; use std::thread; @@ -130,10 +131,17 @@ pub fn update_cli() -> Result<(), CliError> { // // This function returns the latest version of the CLI from the crates.io registry pub fn get_latest_cfcli_version() -> Result { - let output = Command::new("cargo") - .args(["search", "cortexflow-cli", "--limit", "1"]) - .output() - .expect("Error"); + get_latest_cfcli_version_with(&RealCommandRunner) +} + +fn get_latest_cfcli_version_with(runner: &dyn CommandRunner) -> Result { + let args = [ + "search".to_string(), + "cortexflow-cli".to_string(), + "--limit".to_string(), + "1".to_string(), + ]; + let output = runner.run("cargo", &args).expect("Error"); if !output.status.success() { return Err(CliError::InstallerError { @@ -410,10 +418,18 @@ pub async fn update_configmap(config_struct: MetadataConfigFile) -> Result<(), C #[cfg(test)] mod tests { - use crate::essential::extract_version_from_output; + use crate::command_runner::test_support::StubCommandRunner; + use crate::essential::{create_configs, extract_version_from_output, get_latest_cfcli_version_with}; + + #[test] + fn creates_an_empty_blocklist_configuration() { + let configs = create_configs(); + + assert_eq!(configs.blocklist, vec![String::new()]); + } #[test] - fn test_version_extraction() { + fn extracts_the_version_from_cargo_search_output() { let command_stdout = String::from( r#"cortexflow-cli = "0.1.4-test_123" # CortexFlow command line interface made to interact with the CortexBrain core components... @@ -423,4 +439,19 @@ mod tests { let extracted_command = extract_version_from_output(command_stdout.into()); assert_eq!(extracted_command, "0.1.4-test_123"); } + + #[test] + fn get_latest_cfcli_version_returns_parsed_version_on_success() { + let runner = StubCommandRunner::success( + r#"cortexflow-cli = "0.1.5" # CortexFlow command line interface"#, + ); + let version = get_latest_cfcli_version_with(&runner); + assert_eq!(version.unwrap(), "0.1.5"); + } + + #[test] + fn get_latest_cfcli_version_errors_on_command_failure() { + let runner = StubCommandRunner::failure("network unreachable"); + assert!(get_latest_cfcli_version_with(&runner).is_err()); + } } diff --git a/cli/src/install.rs b/cli/src/install.rs index 105853c0..1553726a 100644 --- a/cli/src/install.rs +++ b/cli/src/install.rs @@ -1,3 +1,4 @@ +use crate::command_runner::{CommandRunner, RealCommandRunner}; use crate::errors::CliError; use crate::essential::{BASE_COMMAND, connect_to_client, create_config_file, create_configs}; use clap::{Args, Subcommand}; @@ -213,54 +214,59 @@ async fn install_simple_example_component() -> Result<(), CliError> { // docs: pub async fn install_blocklist_configmap() -> Result<(), CliError> { match connect_to_client().await { - Ok(client) => { - println!( - "{} {}", - "=====>".blue().bold(), - "Checking if the Blocklist configmap exists" - ); - sleep(Duration::from_secs(1)); - let blocklist_exists = check_if_blocklist_exists(client).await?; - if !blocklist_exists { + Ok(client) => install_blocklist_configmap_with_client(client).await, + Err(e) => Err(CliError::ClientError(Error::Api(ErrorResponse { + status: "failed".to_string(), + message: "Failed to connect to kubernetes client".to_string(), + reason: e.to_string(), + code: 404, + }))), + } +} + +// docs: +// +// Split out of install_blocklist_configmap so the blocklist-checking logic can be +// exercised in tests against a fake kube Client, without requiring a real cluster. + +async fn install_blocklist_configmap_with_client(client: Client) -> Result<(), CliError> { + println!( + "{} {}", + "=====>".blue().bold(), + "Checking if the Blocklist configmap exists" + ); + sleep(Duration::from_secs(1)); + let blocklist_exists = check_if_blocklist_exists(client).await?; + if !blocklist_exists { + println!( + "{} {}", + "=====>".blue().bold(), + "Blocklist configmap does not exist".red().bold() + ); + sleep(Duration::from_secs(1)); + println!("{} {}", "=====>".bold().blue(), "Creating configmap"); + let metdata_configs = create_configs(); + sleep(Duration::from_secs(1)); + match create_config_file(metdata_configs).await { + Ok(_) => { println!( "{} {}", - "=====>".blue().bold(), - "Blocklist configmap does not exist".red().bold() - ); - sleep(Duration::from_secs(1)); - println!("{} {}", "=====>".bold().blue(), "Creating configmap"); - let metdata_configs = create_configs(); - sleep(Duration::from_secs(1)); - match create_config_file(metdata_configs).await { - Ok(_) => { - println!( - "{} {}", - "=====>".bold().blue(), - "Configmap created/repaired successfully".bold().green() - ) - } - Err(e) => { - return Err(CliError::InstallerError { - reason: e.to_string(), - }); - } - } - return Ok(()); - } else { - println!() + "=====>".bold().blue(), + "Configmap created/repaired successfully".bold().green() + ) + } + Err(e) => { + return Err(CliError::InstallerError { + reason: e.to_string(), + }); } - - Ok(()) - } - Err(e) => { - return Err(CliError::ClientError(Error::Api(ErrorResponse { - status: "failed".to_string(), - message: "Failed to connect to kubernetes client".to_string(), - reason: e.to_string(), - code: 404, - }))); } + return Ok(()); + } else { + println!() } + + Ok(()) } // docs: @@ -297,6 +303,13 @@ async fn check_if_blocklist_exists(client: Client) -> Result { // fn install_components(components_type: &str) -> Result<(), CliError> { + install_components_with(&RealCommandRunner, components_type) +} + +fn install_components_with( + runner: &dyn CommandRunner, + components_type: &str, +) -> Result<(), CliError> { if components_type == "cortexbrain" { let files_to_install = vec![ "configmap-role.yaml", @@ -327,7 +340,7 @@ fn install_components(components_type: &str) -> Result<(), CliError> { "Applying", component.to_string().green().bold() ); - apply_component(component)?; + apply_component_with(runner, component)?; i = i + 1; } } else if components_type == "simple-example" { @@ -347,7 +360,7 @@ fn install_components(components_type: &str) -> Result<(), CliError> { "Applying", component.to_string().green().bold() ); - apply_component(component)?; + apply_component_with(runner, component)?; i = i + 1; } } else { @@ -367,10 +380,10 @@ fn install_components(components_type: &str) -> Result<(), CliError> { // // Returns an CliError if something fails -fn apply_component(file: &str) -> Result<(), CliError> { - let output = Command::new(BASE_COMMAND) - .args(["apply", "-f", file]) - .output() +fn apply_component_with(runner: &dyn CommandRunner, file: &str) -> Result<(), CliError> { + let args = ["apply".to_string(), "-f".to_string(), file.to_string()]; + let output = runner + .run(BASE_COMMAND, &args) .map_err(|e| CliError::InstallerError { reason: e.to_string(), })?; @@ -453,13 +466,16 @@ fn rm_installation_files(file_to_remove: InstallationType) -> Result<(), CliErro // Returns a CliError if something fails fn download_file(src: &str) -> Result<(), CliError> { - let output = - Command::new("wget") - .args([src]) - .output() - .map_err(|e| CliError::InstallerError { - reason: e.to_string(), - })?; + download_file_with(&RealCommandRunner, src) +} + +fn download_file_with(runner: &dyn CommandRunner, src: &str) -> Result<(), CliError> { + let args = [src.to_string()]; + let output = runner + .run("wget", &args) + .map_err(|e| CliError::InstallerError { + reason: e.to_string(), + })?; if !output.status.success() { return Err(CliError::InstallerError { @@ -486,6 +502,13 @@ fn download_file(src: &str) -> Result<(), CliError> { // Returns an CliError if something fails fn rm_file(file_to_remove: &str) -> Result<(), CliError> { + // rm -f never fails for a missing file, so check existence first + if !std::path::Path::new(file_to_remove).exists() { + return Err(CliError::InstallerError { + reason: format!("File not found: {}", file_to_remove), + }); + } + let output = Command::new("rm") .args(["-f", file_to_remove]) .output() @@ -508,3 +531,158 @@ fn rm_file(file_to_remove: &str) -> Result<(), CliError> { thread::sleep(Duration::from_secs(2)); Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + use clap::Parser; + + #[test] + fn test_download_file_with_success() { + let runner = crate::command_runner::test_support::StubCommandRunner::success(""); + assert!(download_file_with(&runner, "https://example.com/file.yaml").is_ok()); + } + + #[test] + fn test_download_file_with_command_failure() { + let runner = crate::command_runner::test_support::StubCommandRunner::failure("404 Not Found"); + assert!(download_file_with(&runner, "https://example.com/missing.yaml").is_err()); + } + + #[test] + fn test_download_file_with_io_error() { + let runner = crate::command_runner::test_support::StubCommandRunner::io_error(); + assert!(download_file_with(&runner, "https://example.com/file.yaml").is_err()); + } + + #[test] + fn test_rm_file_failure_cleanup() { + let file_name = "test_rm_file_failure_cleanup.txt"; + // Attempt to clean up a non-existent file + assert!(rm_file(file_name).is_err()); + } + + #[test] + fn test_rm_file_success_cleanup() { + let file_name = "test_rm_file_success_cleanup.txt"; + std::fs::write(file_name, "test content").unwrap(); + // Clean up a file that exists + assert!(rm_file(file_name).is_ok()); + } + + // docs: builds a fake kube Client backed by an in-memory tower service, so the + // Kubernetes API calls never leave the process or touch a real cluster. + fn mock_client(status: http::StatusCode, body: String) -> Client { + let service = tower::service_fn(move |_req: http::Request| { + let body = body.clone(); + async move { + let response = http::Response::builder() + .status(status) + .body(http_body_util::Full::new(bytes::Bytes::from(body))) + .unwrap(); + Ok::<_, std::convert::Infallible>(response) + } + }); + Client::new(service, "cortexflow") + } + + fn configmap_found_body() -> String { + serde_json::json!({ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": { + "name": "cortexbrain-client-config", + "namespace": "cortexflow" + }, + "data": {} + }) + .to_string() + } + + fn configmap_not_found_body() -> String { + serde_json::json!({ + "apiVersion": "v1", + "kind": "Status", + "status": "Failure", + "message": "configmaps \"cortexbrain-client-config\" not found", + "reason": "NotFound", + "code": 404 + }) + .to_string() + } + + #[tokio::test] + async fn test_check_if_blocklist_exists_true() { + let client = mock_client(http::StatusCode::OK, configmap_found_body()); + let exists = check_if_blocklist_exists(client).await.unwrap(); + assert!(exists); + } + + #[tokio::test] + async fn test_check_if_blocklist_exists_false() { + let client = mock_client(http::StatusCode::NOT_FOUND, configmap_not_found_body()); + let exists = check_if_blocklist_exists(client).await.unwrap(); + assert!(!exists); + } + + #[tokio::test] + async fn test_install_blocklist_configmap_with_client_already_exists() { + // When the configmap already exists no creation attempt is made, so this + // is safe to exercise end-to-end against the fake client. + let client = mock_client(http::StatusCode::OK, configmap_found_body()); + assert!(install_blocklist_configmap_with_client(client).await.is_ok()); + } + + #[test] + fn test_install_components_with_invalid_type() { + let runner = crate::command_runner::test_support::StubCommandRunner::success(""); + assert!(install_components_with(&runner, "not-a-real-type").is_err()); + } + + #[test] + fn test_install_components_with_simple_example_success() { + let runner = crate::command_runner::test_support::StubCommandRunner::success(""); + assert!(install_components_with(&runner, "simple-example").is_ok()); + } + + #[test] + fn test_install_components_with_cortexbrain_success() { + let runner = crate::command_runner::test_support::StubCommandRunner::success(""); + assert!(install_components_with(&runner, "cortexbrain").is_ok()); + } + + #[test] + fn test_install_components_with_apply_failure() { + let runner = crate::command_runner::test_support::StubCommandRunner::failure("kubectl error"); + assert!(install_components_with(&runner, "simple-example").is_err()); + } + + #[derive(clap::Parser)] + struct InstallCommandsHarness { + #[command(subcommand)] + cmd: InstallCommands, + } + + #[test] + fn test_install_commands_parses_cortexflow() { + let parsed = InstallCommandsHarness::try_parse_from(["cfcli", "cortexflow"]).unwrap(); + assert!(matches!(parsed.cmd, InstallCommands::All)); + } + + #[test] + fn test_install_commands_parses_simple_example() { + let parsed = InstallCommandsHarness::try_parse_from(["cfcli", "simple-example"]).unwrap(); + assert!(matches!(parsed.cmd, InstallCommands::TestPods)); + } + + #[test] + fn test_install_commands_parses_blocklist() { + let parsed = InstallCommandsHarness::try_parse_from(["cfcli", "blocklist"]).unwrap(); + assert!(matches!(parsed.cmd, InstallCommands::Blocklist)); + } + + #[test] + fn test_install_commands_rejects_unknown_subcommand() { + assert!(InstallCommandsHarness::try_parse_from(["cfcli", "not-a-command"]).is_err()); + } +} \ No newline at end of file diff --git a/cli/src/logs.rs b/cli/src/logs.rs index 102d97b6..863cd33b 100644 --- a/cli/src/logs.rs +++ b/cli/src/logs.rs @@ -1,3 +1,4 @@ +use crate::command_runner::{CommandRunner, RealCommandRunner}; use crate::errors::CliError; use crate::essential::{BASE_COMMAND, connect_to_client}; use clap::Args; @@ -5,6 +6,15 @@ use colored::Colorize; use kube::{Error, core::ErrorResponse}; use std::{process::Command, result::Result::Ok, str}; +fn parse_lines(stdout: &[u8]) -> Vec { + str::from_utf8(stdout) + .unwrap_or("") + .lines() + .map(|line| line.trim().to_string()) + .filter(|line| !line.is_empty()) + .collect() +} + #[derive(Args, Debug, Clone)] pub struct LogsArgs { #[arg(long)] @@ -182,16 +192,7 @@ pub async fn logs_command( pub async fn check_namespace_exists(namespace: &str) -> Result { match connect_to_client().await { - Ok(_) => { - let output = Command::new(BASE_COMMAND) - .args(["get", "namespace", namespace]) - .output(); - - match output { - Ok(output) => Ok(output.status.success()), - Err(_) => Ok(false), - } - } + Ok(_) => Ok(check_namespace_exists_with(&RealCommandRunner, namespace)), Err(e) => { return Err(CliError::ClientError(Error::Api(ErrorResponse { status: "failed".to_string(), @@ -203,6 +204,14 @@ pub async fn check_namespace_exists(namespace: &str) -> Result { } } +fn check_namespace_exists_with(runner: &dyn CommandRunner, namespace: &str) -> bool { + let args = ["get".to_string(), "namespace".to_string(), namespace.to_string()]; + match runner.run(BASE_COMMAND, &args) { + Ok(output) => output.status.success(), + Err(_) => false, + } +} + // docs: // // This function returns the available namespaces: @@ -215,30 +224,7 @@ pub async fn check_namespace_exists(namespace: &str) -> Result { pub async fn get_available_namespaces() -> Result, CliError> { match connect_to_client().await { - Ok(_) => { - let output = Command::new(BASE_COMMAND) - .args([ - "get", - "namespaces", - "--no-headers", - "-o", - "custom-columns=NAME:.metadata.name", - ]) - .output(); - - match output { - Ok(output) if output.status.success() => { - let stdout = str::from_utf8(&output.stdout).unwrap_or(""); - let ns = stdout - .lines() - .map(|line| line.trim().to_string()) - .filter(|line| !line.is_empty()) - .collect(); - Ok(ns) - } - _ => Ok(Vec::new()), - } - } + Ok(_) => Ok(get_available_namespaces_with(&RealCommandRunner)), Err(e) => { return Err(CliError::ClientError(Error::Api(ErrorResponse { status: "failed".to_string(), @@ -250,6 +236,21 @@ pub async fn get_available_namespaces() -> Result, CliError> { } } +fn get_available_namespaces_with(runner: &dyn CommandRunner) -> Vec { + let args = [ + "get".to_string(), + "namespaces".to_string(), + "--no-headers".to_string(), + "-o".to_string(), + "custom-columns=NAME:.metadata.name".to_string(), + ]; + + match runner.run(BASE_COMMAND, &args) { + Ok(output) if output.status.success() => parse_lines(&output.stdout), + _ => Vec::new(), + } +} + // docs: // // This function returns the pods: @@ -265,34 +266,11 @@ async fn get_pods_for_service( service_name: &str, ) -> Result, CliError> { match connect_to_client().await { - Ok(_) => { - let output = Command::new(BASE_COMMAND) - .args([ - "get", - "pods", - "-n", - namespace, - "-l", - &format!("app={}", service_name), - "--no-headers", - "-o", - "custom-columns=NAME:.metadata.name", - ]) - .output(); - - match output { - Ok(output) if output.status.success() => { - let stdout = str::from_utf8(&output.stdout).unwrap_or(""); - let pods = stdout - .lines() - .map(|line| line.trim().to_string()) - .filter(|line| !line.is_empty()) - .collect(); - Ok(pods) - } - _ => Ok(Vec::new()), - } - } + Ok(_) => Ok(get_pods_for_service_with( + &RealCommandRunner, + namespace, + service_name, + )), Err(e) => { return Err(CliError::ClientError(Error::Api(ErrorResponse { status: "failed".to_string(), @@ -304,6 +282,29 @@ async fn get_pods_for_service( } } +fn get_pods_for_service_with( + runner: &dyn CommandRunner, + namespace: &str, + service_name: &str, +) -> Vec { + let args = [ + "get".to_string(), + "pods".to_string(), + "-n".to_string(), + namespace.to_string(), + "-l".to_string(), + format!("app={}", service_name), + "--no-headers".to_string(), + "-o".to_string(), + "custom-columns=NAME:.metadata.name".to_string(), + ]; + + match runner.run(BASE_COMMAND, &args) { + Ok(output) if output.status.success() => parse_lines(&output.stdout), + _ => Vec::new(), + } +} + // docs: // // This function returns the pods: @@ -320,34 +321,11 @@ async fn get_pods_for_component( component: &Component, ) -> Result, CliError> { match connect_to_client().await { - Ok(_) => { - let output = Command::new(BASE_COMMAND) - .args([ - "get", - "pods", - "-n", - namespace, - "-l", - component.to_label_selector(), - "--no-headers", - "-o", - "custom-columns=NAME:.metadata.name", - ]) - .output(); - - match output { - Ok(output) if output.status.success() => { - let stdout = str::from_utf8(&output.stdout).unwrap_or(""); - let pods = stdout - .lines() - .map(|line| line.trim().to_string()) - .filter(|line| !line.is_empty()) - .collect(); - Ok(pods) - } - _ => Ok(Vec::new()), - } - } + Ok(_) => Ok(get_pods_for_component_with( + &RealCommandRunner, + namespace, + component, + )), Err(e) => { return Err(CliError::ClientError(Error::Api(ErrorResponse { status: "failed".to_string(), @@ -359,6 +337,29 @@ async fn get_pods_for_component( } } +fn get_pods_for_component_with( + runner: &dyn CommandRunner, + namespace: &str, + component: &Component, +) -> Vec { + let args = [ + "get".to_string(), + "pods".to_string(), + "-n".to_string(), + namespace.to_string(), + "-l".to_string(), + component.to_label_selector().to_string(), + "--no-headers".to_string(), + "-o".to_string(), + "custom-columns=NAME:.metadata.name".to_string(), + ]; + + match runner.run(BASE_COMMAND, &args) { + Ok(output) if output.status.success() => parse_lines(&output.stdout), + _ => Vec::new(), + } +} + // docs: // // This function returns the available namespaces: @@ -371,32 +372,7 @@ async fn get_pods_for_component( async fn get_all_pods(namespace: &str) -> Result, CliError> { match connect_to_client().await { - Ok(_) => { - let output = Command::new(BASE_COMMAND) - .args([ - "get", - "pods", - "-n", - namespace, - "--no-headers", - "-o", - "custom-columns=NAME:.metadata.name", - ]) - .output(); - - match output { - Ok(output) if output.status.success() => { - let stdout = str::from_utf8(&output.stdout).unwrap_or(""); - let pods = stdout - .lines() - .map(|line| line.trim().to_string()) - .filter(|line| !line.is_empty()) - .collect(); - Ok(pods) - } - _ => Ok(Vec::new()), - } - } + Ok(_) => Ok(get_all_pods_with(&RealCommandRunner, namespace)), Err(e) => { return Err(CliError::ClientError(Error::Api(ErrorResponse { status: "failed".to_string(), @@ -407,3 +383,216 @@ async fn get_all_pods(namespace: &str) -> Result, CliError> { } } } + +fn get_all_pods_with(runner: &dyn CommandRunner, namespace: &str) -> Vec { + let args = [ + "get".to_string(), + "pods".to_string(), + "-n".to_string(), + namespace.to_string(), + "--no-headers".to_string(), + "-o".to_string(), + "custom-columns=NAME:.metadata.name".to_string(), + ]; + + match runner.run(BASE_COMMAND, &args) { + Ok(output) if output.status.success() => parse_lines(&output.stdout), + _ => Vec::new(), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::command_runner::test_support::StubCommandRunner; + + // Component + + #[test] + fn test_component_from_control_plane() { + assert!(matches!( + Component::from("control-plane".to_string()), + Component::ControlPlane + )); + } + + #[test] + fn test_component_from_data_plane() { + assert!(matches!( + Component::from("data-plane".to_string()), + Component::DataPlane + )); + } + + #[test] + fn test_component_from_is_case_insensitive() { + assert!(matches!( + Component::from("Data-Plane".to_string()), + Component::DataPlane + )); + } + + #[test] + fn test_component_from_unknown_defaults_to_control_plane() { + assert!(matches!( + Component::from("unknown".to_string()), + Component::ControlPlane + )); + } + + #[test] + fn test_component_to_label_selector() { + assert_eq!( + Component::ControlPlane.to_label_selector(), + "component=control-plane" + ); + assert_eq!( + Component::DataPlane.to_label_selector(), + "component=data-plane" + ); + } + + // parse_lines + + #[test] + fn test_parse_lines_filters_empty_lines_and_trims() { + let stdout = b" pod-a \n\npod-b\n \n"; + let parsed = parse_lines(stdout); + assert_eq!(parsed, vec!["pod-a".to_string(), "pod-b".to_string()]); + } + + #[test] + fn test_parse_lines_empty_input() { + assert!(parse_lines(b"").is_empty()); + } + + // check_namespace_exists_with + + #[test] + fn test_check_namespace_exists_with_success() { + let runner = StubCommandRunner::success(""); + assert!(check_namespace_exists_with(&runner, "cortexflow")); + } + + #[test] + fn test_check_namespace_exists_with_failure() { + let runner = StubCommandRunner::failure("not found"); + assert!(!check_namespace_exists_with(&runner, "missing")); + } + + #[test] + fn test_check_namespace_exists_with_io_error() { + let runner = StubCommandRunner::io_error(); + assert!(!check_namespace_exists_with(&runner, "cortexflow")); + } + + // get_available_namespaces_with + + #[test] + fn test_get_available_namespaces_with_success() { + let runner = StubCommandRunner::success("default\ncortexflow\nkube-system\n"); + let namespaces = get_available_namespaces_with(&runner); + assert_eq!( + namespaces, + vec![ + "default".to_string(), + "cortexflow".to_string(), + "kube-system".to_string() + ] + ); + } + + #[test] + fn test_get_available_namespaces_with_empty_output() { + let runner = StubCommandRunner::success(""); + assert!(get_available_namespaces_with(&runner).is_empty()); + } + + #[test] + fn test_get_available_namespaces_with_command_failure() { + let runner = StubCommandRunner::failure("connection refused"); + assert!(get_available_namespaces_with(&runner).is_empty()); + } + + // get_pods_for_service_with + + #[test] + fn test_get_pods_for_service_with_success() { + let runner = StubCommandRunner::success("pod-a\npod-b\n"); + let pods = get_pods_for_service_with(&runner, "cortexflow", "my-service"); + assert_eq!(pods, vec!["pod-a".to_string(), "pod-b".to_string()]); + } + + #[test] + fn test_get_pods_for_service_with_no_matches() { + let runner = StubCommandRunner::success(""); + assert!(get_pods_for_service_with(&runner, "cortexflow", "unknown-service").is_empty()); + } + + #[test] + fn test_get_pods_for_service_with_command_failure() { + let runner = StubCommandRunner::failure("error"); + assert!(get_pods_for_service_with(&runner, "cortexflow", "my-service").is_empty()); + } + + // get_pods_for_component_with + + #[test] + fn test_get_pods_for_component_with_success() { + let runner = StubCommandRunner::success("agent-pod\n"); + let pods = get_pods_for_component_with(&runner, "cortexflow", &Component::DataPlane); + assert_eq!(pods, vec!["agent-pod".to_string()]); + } + + #[test] + fn test_get_pods_for_component_with_no_matches() { + let runner = StubCommandRunner::success(""); + assert!( + get_pods_for_component_with(&runner, "cortexflow", &Component::ControlPlane) + .is_empty() + ); + } + + #[test] + fn test_get_pods_for_component_with_command_failure() { + let runner = StubCommandRunner::failure("error"); + assert!( + get_pods_for_component_with(&runner, "cortexflow", &Component::ControlPlane) + .is_empty() + ); + } + + // get_all_pods_with + + #[test] + fn test_get_all_pods_with_success() { + let runner = StubCommandRunner::success("pod-a\npod-b\npod-c\n"); + let pods = get_all_pods_with(&runner, "cortexflow"); + assert_eq!( + pods, + vec![ + "pod-a".to_string(), + "pod-b".to_string(), + "pod-c".to_string() + ] + ); + } + + #[test] + fn test_get_all_pods_with_empty_namespace() { + let runner = StubCommandRunner::success(""); + assert!(get_all_pods_with(&runner, "empty_namespace").is_empty()); + } + + #[test] + fn test_get_all_pods_with_command_failure() { + let runner = StubCommandRunner::failure("namespace not found"); + assert!(get_all_pods_with(&runner, "non_existent_namespace").is_empty()); + } + + #[test] + fn test_get_all_pods_with_io_error() { + let runner = StubCommandRunner::io_error(); + assert!(get_all_pods_with(&runner, "cortexflow").is_empty()); + } +} \ No newline at end of file diff --git a/cli/src/main.rs b/cli/src/main.rs index 8d543cd1..89070ddc 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -1,3 +1,4 @@ +mod command_runner; mod errors; mod essential; mod install; @@ -194,3 +195,84 @@ async fn args_parser() -> Result<(), CliError> { async fn main() { let _ = args_parser().await.map_err(|e| eprintln!("{}", e)); } + +#[cfg(test)] +mod tests { + use super::*; + use clap::Parser; + + // docs: args_parser() dispatches to functions that reach a real kubernetes + // cluster / gRPC agent, so only the pure clap parsing logic is unit tested here. + + #[test] + fn test_parse_no_subcommand() { + let cli = Cli::try_parse_from(["cfcli"]).unwrap(); + assert!(cli.cmd.is_none()); + } + + #[test] + fn test_parse_uninstall() { + let cli = Cli::try_parse_from(["cfcli", "uninstall"]).unwrap(); + assert!(matches!(cli.cmd, Some(Commands::Uninstall))); + } + + #[test] + fn test_parse_update() { + let cli = Cli::try_parse_from(["cfcli", "update"]).unwrap(); + assert!(matches!(cli.cmd, Some(Commands::Update))); + } + + #[test] + fn test_parse_info() { + let cli = Cli::try_parse_from(["cfcli", "info"]).unwrap(); + assert!(matches!(cli.cmd, Some(Commands::Info))); + } + + #[test] + fn test_parse_install_cortexflow() { + let cli = Cli::try_parse_from(["cfcli", "install", "cortexflow"]).unwrap(); + match cli.cmd { + Some(Commands::Install(args)) => { + assert!(matches!(args.install_cmd, InstallCommands::All)); + } + _ => panic!("expected Install command"), + } + } + + #[test] + fn test_parse_logs_with_flags() { + let cli = Cli::try_parse_from([ + "cfcli", + "logs", + "--service", + "my-service", + "--namespace", + "cortexflow", + ]) + .unwrap(); + match cli.cmd { + Some(Commands::Logs(args)) => { + assert_eq!(args.service, Some("my-service".to_string())); + assert_eq!(args.namespace, Some("cortexflow".to_string())); + assert_eq!(args.component, None); + } + _ => panic!("expected Logs command"), + } + } + + #[test] + fn test_parse_status_with_output_format() { + let cli = Cli::try_parse_from(["cfcli", "status", "--output", "json"]).unwrap(); + match cli.cmd { + Some(Commands::Status(args)) => { + assert_eq!(args.output, Some("json".to_string())); + } + _ => panic!("expected Status command"), + } + } + + #[test] + fn test_parse_unknown_subcommand_fails() { + assert!(Cli::try_parse_from(["cfcli", "not-a-real-command"]).is_err()); + } +} \ No newline at end of file diff --git a/cli/src/mod.rs b/cli/src/mod.rs index fe7c8165..326b6c4d 100644 --- a/cli/src/mod.rs +++ b/cli/src/mod.rs @@ -6,4 +6,5 @@ pub mod status; pub mod logs; pub mod monitoring; pub mod policies; -pub mod errors; \ No newline at end of file +pub mod errors; +pub mod command_runner; \ No newline at end of file diff --git a/cli/src/monitoring.rs b/cli/src/monitoring.rs index eefae1c7..a27cc86f 100644 --- a/cli/src/monitoring.rs +++ b/cli/src/monitoring.rs @@ -10,7 +10,8 @@ use tonic_reflection::pb::v1::server_reflection_response::MessageResponse; use agent_api::client::{connect_to_client, connect_to_server_reflection}; use agent_api::requests::{ get_all_features, send_active_connection_request, send_dropped_packets_request, - send_latency_metrics_request, send_tracked_veth_request, send_veth_tracked_hashmap_req, + send_latency_metrics_request, + send_veth_tracked_hashmap_req, }; use crate::errors::CliError; @@ -345,3 +346,28 @@ fn convert_timestamp_to_date(timestamp: u64) -> String { .map(|dt| dt.to_string()) .unwrap_or_else(|| "Cannot convert timestamp to date".to_string()) } + + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_convert_timestamp_to_date_valid() { + // 1700000000000000 us = 2023-11-14T22:13:20Z + let result = convert_timestamp_to_date(1_700_000_000_000_000); + assert!(result.contains("2023")); + } + + #[test] + fn test_convert_timestamp_to_date_zero() { + let result = convert_timestamp_to_date(0); + assert!(result.contains("1970")); + } + + #[test] + fn test_convert_timestamp_to_date_invalid_returns_fallback() { + let result = convert_timestamp_to_date(i64::MAX as u64); + assert_eq!(result, "Cannot convert timestamp to date"); + } +} \ No newline at end of file diff --git a/cli/src/policies.rs b/cli/src/policies.rs index 03841e93..cdc962d8 100644 --- a/cli/src/policies.rs +++ b/cli/src/policies.rs @@ -127,3 +127,57 @@ pub async fn remove_ip(ip:&str) -> Result<(), Error> { } Ok(()) } + + +#[cfg(test)] +mod tests { + use super::*; + use clap::Parser; + + #[derive(Parser)] + struct TestCli { + #[command(flatten)] + args: PoliciesArgs, + } + + // docs: the gRPC-calling functions in this module have no pure logic to + // unit test without mocking the agent_api client; these tests instead + // verify the clap argument parsing, which is the only externally-free logic here. + + #[test] + fn test_parse_create_blocklist_with_ip() { + let cli = TestCli::try_parse_from(["cfcli", "--flags", "1.2.3.4", "create-blocklist"]) + .unwrap(); + assert!(matches!( + cli.args.policy_cmd, + PoliciesCommands::CreateBlocklist + )); + assert_eq!(cli.args.flags, Some("1.2.3.4".to_string())); + } + + #[test] + fn test_parse_check_blocklist() { + let cli = TestCli::try_parse_from(["cfcli", "check-blocklist"]).unwrap(); + assert!(matches!( + cli.args.policy_cmd, + PoliciesCommands::CheckBlocklist + )); + assert_eq!(cli.args.flags, None); + } + + #[test] + fn test_parse_remove_ip() { + let cli = + TestCli::try_parse_from(["cfcli", "--flags", "5.6.7.8", "remove-ip"]).unwrap(); + assert!(matches!( + cli.args.policy_cmd, + PoliciesCommands::RemoveIpFromBlocklist + )); + assert_eq!(cli.args.flags, Some("5.6.7.8".to_string())); + } + + #[test] + fn test_parse_unknown_subcommand_fails() { + assert!(TestCli::try_parse_from(["cfcli", "not-a-command"]).is_err()); + } +} \ No newline at end of file diff --git a/cli/src/service.rs b/cli/src/service.rs index 8cfebf10..9885afe7 100644 --- a/cli/src/service.rs +++ b/cli/src/service.rs @@ -7,6 +7,40 @@ use crate::errors::CliError; use crate::essential::{BASE_COMMAND, connect_to_client}; use crate::logs::{check_namespace_exists, get_available_namespaces}; +// docs: +// +// Pure formatting helper extracted from list_services so the table-row +// formatting logic can be unit tested with stubbed kubectl output. + +fn format_service_rows(stdout: &str) -> Vec { + stdout + .lines() + .filter_map(|line| { + let parts: Vec<&str> = line.split_whitespace().collect(); + if parts.len() >= 5 { + let name = parts[0]; + let ready = parts[1]; + let status = parts[2]; + let restarts = parts[3]; + let age = parts[4]; + + let full_status = if ready.contains('/') { + format!("{} ({})", status, ready) + } else { + status.to_string() + }; + + Some(format!( + "{:<40} {:<20} {:<10} {:<10}", + name, full_status, restarts, age + )) + } else { + None + } + }) + .collect() +} + //service subcommands #[derive(Subcommand, Debug, Clone)] pub enum ServiceCommands { @@ -104,26 +138,8 @@ pub async fn list_services(namespace: Option) -> Result<(), CliError> { println!("{}", "-".repeat(80)); // Display Each Pod. - for line in stdout.lines() { - let parts: Vec<&str> = line.split_whitespace().collect(); - if parts.len() >= 5 { - let name = parts[0]; - let ready = parts[1]; - let status = parts[2]; - let restarts = parts[3]; - let age = parts[4]; - - let full_status = if ready.contains('/') { - format!("{} ({})", status, ready) - } else { - status.to_string() - }; - - println!( - "{:<40} {:<20} {:<10} {:<10}", - name, full_status, restarts, age - ); - } + for row in format_service_rows(stdout) { + println!("{}", row); } Ok(()) } @@ -267,3 +283,39 @@ pub async fn describe_service( } } } + + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_format_service_rows_formats_matching_lines() { + let stdout = "my-svc 1/1 Running 0 5d\nother-svc 2/2 Pending 1 1h\n"; + let rows = format_service_rows(stdout); + assert_eq!(rows.len(), 2); + assert!(rows[0].contains("my-svc")); + assert!(rows[0].contains("Running (1/1)")); + assert!(rows[1].contains("other-svc")); + assert!(rows[1].contains("Pending (2/2)")); + } + + #[test] + fn test_format_service_rows_skips_short_lines() { + let stdout = "incomplete line\n"; + assert!(format_service_rows(stdout).is_empty()); + } + + #[test] + fn test_format_service_rows_without_ready_slash() { + let stdout = "svc-a Active 0 5d extra\n"; + let rows = format_service_rows(stdout); + assert_eq!(rows.len(), 1); + assert!(rows[0].contains("svc-a")); + } + + #[test] + fn test_format_service_rows_empty_input() { + assert!(format_service_rows("").is_empty()); + } +} \ No newline at end of file diff --git a/cli/src/status.rs b/cli/src/status.rs index ca5d43aa..ccb58aa1 100644 --- a/cli/src/status.rs +++ b/cli/src/status.rs @@ -7,6 +7,29 @@ use crate::logs::{ get_available_namespaces, check_namespace_exists }; use crate::essential::{ BASE_COMMAND, connect_to_client }; use crate::errors::CliError; +// docs: +// +// Pure parsing helper shared by get_pods_status/get_services_status so the +// column-parsing logic can be unit tested with stubbed kubectl output. + +fn parse_status_columns(stdout: &str, min_parts: usize) -> Vec<(String, String, String)> { + stdout + .lines() + .filter_map(|line| { + let parts: Vec<&str> = line.split_whitespace().collect(); + if parts.len() >= min_parts { + Some(( + parts[0].to_string(), + parts[1].to_string(), + parts[2].to_string(), + )) + } else { + None + } + }) + .collect() +} + #[derive(Debug)] pub enum OutputFormat { Text, @@ -165,23 +188,7 @@ async fn get_pods_status(namespace: &str) -> Result { let stdout = str::from_utf8(&output.stdout).unwrap_or(""); - Ok( - stdout - .lines() - .filter_map(|line| { - let parts: Vec<&str> = line.split_whitespace().collect(); - if parts.len() >= 3 { - Some(( - parts[0].to_string(), // name - parts[1].to_string(), // ready - parts[2].to_string(), // status - )) - } else { - None - } - }) - .collect() - ) + Ok(parse_status_columns(stdout, 3)) } _ => Ok(Vec::new()), } @@ -220,23 +227,7 @@ async fn get_services_status(namespace: &str) -> Result { let stdout = str::from_utf8(&output.stdout).unwrap_or(""); - Ok( - stdout - .lines() - .filter_map(|line| { - let parts: Vec<&str> = line.split_whitespace().collect(); - if parts.len() >= 4 { - Some(( - parts[0].to_string(), // name - parts[1].to_string(), // type - parts[2].to_string(), // cluster ips - )) - } else { - None - } - }) - .collect() - ) + Ok(parse_status_columns(stdout, 4)) } _ => Ok(Vec::new()), } @@ -360,3 +351,79 @@ fn display_yaml_format( println!(" cluster_ip: {}", cluster_ip); } } + + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_output_format_from_json() { + assert!(matches!(OutputFormat::from("json".to_string()), OutputFormat::Json)); + assert!(matches!(OutputFormat::from("JSON".to_string()), OutputFormat::Json)); + } + + #[test] + fn test_output_format_from_yaml() { + assert!(matches!(OutputFormat::from("yaml".to_string()), OutputFormat::Yaml)); + } + + #[test] + fn test_output_format_from_defaults_to_text() { + assert!(matches!(OutputFormat::from("unknown".to_string()), OutputFormat::Text)); + } + + #[test] + fn test_parse_status_columns_pods() { + let stdout = "pod-a 1/1 Running\npod-b 0/1 Pending\n"; + let parsed = parse_status_columns(stdout, 3); + assert_eq!( + parsed, + vec![ + ("pod-a".to_string(), "1/1".to_string(), "Running".to_string()), + ("pod-b".to_string(), "0/1".to_string(), "Pending".to_string()), + ] + ); + } + + #[test] + fn test_parse_status_columns_services() { + let stdout = "svc-a ClusterIP 10.0.0.1 443/TCP\n"; + let parsed = parse_status_columns(stdout, 4); + assert_eq!( + parsed, + vec![("svc-a".to_string(), "ClusterIP".to_string(), "10.0.0.1".to_string())] + ); + } + + #[test] + fn test_parse_status_columns_skips_short_lines() { + let stdout = "only-one-column\n"; + assert!(parse_status_columns(stdout, 3).is_empty()); + } + + #[test] + fn test_parse_status_columns_empty_input() { + assert!(parse_status_columns("", 3).is_empty()); + } + + #[test] + fn test_display_text_format_does_not_panic() { + display_text_format( + "cortexflow", + true, + vec![("pod-a".to_string(), "1/1".to_string(), "Running".to_string())], + vec![("svc-a".to_string(), "ClusterIP".to_string(), "10.0.0.1".to_string())], + ); + } + + #[test] + fn test_display_json_format_does_not_panic() { + display_json_format("cortexflow", false, Vec::new(), Vec::new()); + } + + #[test] + fn test_display_yaml_format_does_not_panic() { + display_yaml_format("cortexflow", false, Vec::new(), Vec::new()); + } +} \ No newline at end of file diff --git a/cli/src/uninstall.rs b/cli/src/uninstall.rs index b9558ddc..39949692 100644 --- a/cli/src/uninstall.rs +++ b/cli/src/uninstall.rs @@ -1,6 +1,7 @@ use colored::Colorize; -use std::{io::stdin, process::Command}; +use std::io::stdin; +use crate::command_runner::{CommandRunner, RealCommandRunner}; use crate::errors::CliError; use crate::essential::{BASE_COMMAND, connect_to_client}; use kube::{Error, core::ErrorResponse}; @@ -78,22 +79,7 @@ async fn uninstall_all() -> Result<(), CliError> { "=====>".blue().bold(), "Deleting cortexflow components".red().bold() ); - let output = Command::new(BASE_COMMAND) - .args(["delete", "namespace", "cortexflow"]) - .output() - .map_err(|e| CliError::InstallerError { - reason: format!("Failed to execute delete command: {}", e), - })?; - - if output.status.success() { - println!("✅ Removed cortexflow namespace"); - Ok(()) - } else { - let stderr = String::from_utf8_lossy(&output.stderr); - return Err(CliError::InstallerError { - reason: format!("Failed to delete cortexflow namespace. Error: {}", stderr), - }); - } + uninstall_all_with(&RealCommandRunner) } Err(e) => { return { @@ -108,6 +94,29 @@ async fn uninstall_all() -> Result<(), CliError> { } } +fn uninstall_all_with(runner: &dyn CommandRunner) -> Result<(), CliError> { + let args = [ + "delete".to_string(), + "namespace".to_string(), + "cortexflow".to_string(), + ]; + let output = runner + .run(BASE_COMMAND, &args) + .map_err(|e| CliError::InstallerError { + reason: format!("Failed to execute delete command: {}", e), + })?; + + if output.status.success() { + println!("✅ Removed cortexflow namespace"); + Ok(()) + } else { + let stderr = String::from_utf8_lossy(&output.stderr); + Err(CliError::InstallerError { + reason: format!("Failed to delete cortexflow namespace. Error: {}", stderr), + }) + } +} + //docs: // // This function manages the uninstall of given cortexflow components @@ -127,22 +136,7 @@ async fn uninstall_component(component_type: &str, component: &str) -> Result<() component ); - let output = Command::new(BASE_COMMAND) - .args(["delete", component_type, component, "-n", "cortexflow"]) - .output() - .map_err(|e| CliError::InstallerError { - reason: format!("Failed to execute delete command: {}", e), - })?; - - if output.status.success() { - println!("✅ Removed component {}", component); - Ok(()) - } else { - let stderr = String::from_utf8_lossy(&output.stderr); - return Err(CliError::InstallerError { - reason: format!("Failed to delete component '{}': {}", component, stderr), - }); - } + uninstall_component_with(&RealCommandRunner, component_type, component) } Err(e) => { return { @@ -155,4 +149,76 @@ async fn uninstall_component(component_type: &str, component: &str) -> Result<() }; } } +} + +fn uninstall_component_with( + runner: &dyn CommandRunner, + component_type: &str, + component: &str, +) -> Result<(), CliError> { + let args = [ + "delete".to_string(), + component_type.to_string(), + component.to_string(), + "-n".to_string(), + "cortexflow".to_string(), + ]; + let output = runner + .run(BASE_COMMAND, &args) + .map_err(|e| CliError::InstallerError { + reason: format!("Failed to execute delete command: {}", e), + })?; + + if output.status.success() { + println!("✅ Removed component {}", component); + Ok(()) + } else { + let stderr = String::from_utf8_lossy(&output.stderr); + Err(CliError::InstallerError { + reason: format!("Failed to delete component '{}': {}", component, stderr), + }) + } +} + + +#[cfg(test)] +mod tests { + use super::*; + use crate::command_runner::test_support::StubCommandRunner; + + #[test] + fn test_uninstall_all_with_success() { + let runner = StubCommandRunner::success(""); + assert!(uninstall_all_with(&runner).is_ok()); + } + + #[test] + fn test_uninstall_all_with_command_failure() { + let runner = StubCommandRunner::failure("namespace not found"); + assert!(uninstall_all_with(&runner).is_err()); + } + + #[test] + fn test_uninstall_all_with_io_error() { + let runner = StubCommandRunner::io_error(); + assert!(uninstall_all_with(&runner).is_err()); + } + + #[test] + fn test_uninstall_component_with_success() { + let runner = StubCommandRunner::success(""); + assert!(uninstall_component_with(&runner, "deployment", "cortexflow-identity").is_ok()); + } + + #[test] + fn test_uninstall_component_with_command_failure() { + let runner = StubCommandRunner::failure("component not found"); + assert!(uninstall_component_with(&runner, "deployment", "unknown").is_err()); + } + + #[test] + fn test_uninstall_component_with_io_error() { + let runner = StubCommandRunner::io_error(); + assert!(uninstall_component_with(&runner, "deployment", "cortexflow-identity").is_err()); + } } \ No newline at end of file diff --git a/core/api/Cargo.toml b/core/api/Cargo.toml index 07e15e41..a61719cc 100644 --- a/core/api/Cargo.toml +++ b/core/api/Cargo.toml @@ -29,7 +29,7 @@ tonic = "0.14.0" tonic-prost = "0.14.0" tracing = "0.1.41" aya = "0.13.1" -cortexbrain-common = { version="0.1.2", features = [ +cortexbrain-common = { version="0.1.4", features = [ "map-handlers", "network-structs", "buffer-reader", diff --git a/core/src/components/identity/Cargo.toml b/core/src/components/identity/Cargo.toml index 46667c61..4cf6c732 100644 --- a/core/src/components/identity/Cargo.toml +++ b/core/src/components/identity/Cargo.toml @@ -28,7 +28,7 @@ tokio = { version = "1.48.0", features = [ ] } tracing = "0.1.41" bytemuck = { version = "1.23.0", features = ["derive"] } -cortexbrain-common = { version="0.1.2", features = [ +cortexbrain-common = { version="0.1.4", features = [ "map-handlers", "program-handlers", "network-structs", diff --git a/core/src/testing/otel_agent.yaml b/core/src/testing/otel_agent.yaml index 7b624590..047f8341 100644 --- a/core/src/testing/otel_agent.yaml +++ b/core/src/testing/otel_agent.yaml @@ -17,6 +17,15 @@ data: http: endpoint: 0.0.0.0:4318 + processors: + transform/identity: + error_mode: ignore + metric_statements: + - context: datapoint + statements: + - set(attributes["container.name"], attributes["k8s.pod.name"]) + where attributes["container.name"] == nil and attributes["k8s.pod.name"] != nil + exporters: otlp: endpoint: ${OTEL_COLLECTOR_SERVICE_HOST}:4317 @@ -38,6 +47,7 @@ data: exporters: [otlp, logging] metrics: receivers: [otlp] + processors: [transform/identity] exporters: [otlp, prometheus] --- @@ -64,9 +74,9 @@ spec: dnsPolicy: ClusterFirstWithHostNet containers: - name: otel-agent - image: otel/opentelemetry-collector:0.95.0 + image: otel/opentelemetry-collector-contrib:0.95.0 command: - - "/otelcol" + - "/otelcol-contrib" - "--config=/conf/otel-agent-config.yaml" resources: limits: @@ -123,6 +133,13 @@ data: limit_mib: 1500 spike_limit_mib: 512 check_interval: 5s + transform/identity: + error_mode: ignore + metric_statements: + - context: datapoint + statements: + - set(attributes["container.name"], attributes["k8s.pod.name"]) + where attributes["container.name"] == nil and attributes["k8s.pod.name"] != nil exporters: # otlp: @@ -147,7 +164,7 @@ data: exporters: [logging] metrics: receivers: [otlp] - processors: [memory_limiter] + processors: [memory_limiter, transform/identity] exporters: [logging,prometheus] --- @@ -197,9 +214,9 @@ spec: spec: containers: - name: otel-collector - image: otel/opentelemetry-collector:0.95.0 + image: otel/opentelemetry-collector-contrib:0.95.0 command: - - "/otelcol" + - "/otelcol-contrib" - "--config=/conf/otel-collector-config.yaml" resources: limits: diff --git a/helm/Chart.yaml b/helm/Chart.yaml new file mode 100644 index 00000000..675510cd --- /dev/null +++ b/helm/Chart.yaml @@ -0,0 +1,5 @@ +apiVersion: v2 +name: CortexBrain +version: 0.1.0 +description: | + This chart installs CortexFlow to a kubernetes cluster, instead of using the cli installation method. diff --git a/helm/README.md b/helm/README.md new file mode 100644 index 00000000..892ad7ba --- /dev/null +++ b/helm/README.md @@ -0,0 +1,85 @@ +# CortexBrain + +![Version: 0.1.0](https://img.shields.io/badge/Version-0.1.0-informational?style=flat-square) + +This chart installs CortexFlow to a kubernetes cluster, instead of using the cli installation method. + +## Values + +| Key | Type | Default | Description | +|-----|------|---------|-------------| +| agent.image.repository | string | `"ghcr.io/cortexflow/agent"` | | +| agent.image.version | string | `"latest"` | | +| agent.priorityClassName | string | `""` | | +| agent.resources.limits.memory | string | `"200Mi"` | | +| agent.resources.requests.cpu | string | `"100m"` | | +| agent.resources.requests.memory | string | `"100Mi"` | | +| agent.securityContext.allowPrivilegeEscalation | bool | `true` | | +| agent.securityContext.capabilities.add[0] | string | `"SYS_ADMIN"` | | +| agent.securityContext.capabilities.add[1] | string | `"NET_ADMIN"` | | +| agent.securityContext.capabilities.add[2] | string | `"SYS_RESOURCE"` | | +| agent.securityContext.capabilities.add[3] | string | `"BPF"` | | +| agent.securityContext.capabilities.add[4] | string | `"SYS_PTRACE"` | | +| agent.securityContext.privileged | bool | `true` | | +| agent.tolerations | list | `[]` | | +| blocklist | string | `""` | | +| bpfMapPermissions.image.repository | string | `"ubuntu"` | | +| bpfMapPermissions.image.version | string | `"24.04"` | | +| bpfMapPermissions.securityContext.allowPrivilegeEscalation | bool | `true` | | +| bpfMapPermissions.securityContext.capabilities.add[0] | string | `"SYS_ADMIN"` | | +| bpfMapPermissions.securityContext.capabilities.add[1] | string | `"NET_ADMIN"` | | +| bpfMapPermissions.securityContext.capabilities.add[2] | string | `"SYS_RESOURCE"` | | +| bpfMapPermissions.securityContext.capabilities.add[3] | string | `"BPF"` | | +| bpfMapPermissions.securityContext.capabilities.add[4] | string | `"SYS_PTRACE"` | | +| bpfMapPermissions.securityContext.privileged | bool | `true` | | +| bpfMapPermissions.securityContext.runAsUser | int | `0` | | +| bpfTool.image.repository | string | `"danielpacak/bpftool-runner"` | | +| bpfTool.image.version | string | `"latest"` | | +| bpfTool.resources.limits.cpu | string | `"1"` | | +| bpfTool.resources.limits.memory | string | `"200Mi"` | | +| bpfTool.resources.requests.cpu | string | `"1"` | | +| bpfTool.resources.requests.memory | string | `"100Mi"` | | +| bpfTool.securityContext.allowPrivilegeEscalation | bool | `true` | | +| bpfTool.securityContext.capabilities.add[0] | string | `"SYS_ADMIN"` | | +| bpfTool.securityContext.capabilities.add[1] | string | `"NET_ADMIN"` | | +| bpfTool.securityContext.capabilities.add[2] | string | `"SYS_RESOURCE"` | | +| bpfTool.securityContext.capabilities.add[3] | string | `"BPF"` | | +| bpfTool.securityContext.capabilities.add[4] | string | `"SYS_PTRACE"` | | +| bpfTool.securityContext.privileged | bool | `true` | | +| global.otel.endpoint | string | `"http://localhost:4317"` | | +| global.otel.protocol | string | `"grpc"` | | +| global.priorityClassName | string | `""` | | +| global.tolerations | list | `[]` | | +| identity.image.repository | string | `"ghcr.io/cortexflow/identity"` | | +| identity.image.version | string | `"latest"` | | +| identity.priorityClassName | string | `""` | | +| identity.resources.limits.memory | string | `"200Mi"` | | +| identity.resources.requests.cpu | string | `"100m"` | | +| identity.resources.requests.memory | string | `"100Mi"` | | +| identity.securityContext.allowPrivilegeEscalation | bool | `true` | | +| identity.securityContext.capabilities.add[0] | string | `"SYS_ADMIN"` | | +| identity.securityContext.capabilities.add[1] | string | `"NET_ADMIN"` | | +| identity.securityContext.capabilities.add[2] | string | `"SYS_RESOURCE"` | | +| identity.securityContext.capabilities.add[3] | string | `"BPF"` | | +| identity.securityContext.capabilities.add[4] | string | `"SYS_PTRACE"` | | +| identity.securityContext.privileged | bool | `true` | | +| identity.tolerations | list | `[]` | | +| metrics.image.repository | string | `"ghcr.io/cortexflow/metrics"` | | +| metrics.image.version | string | `"latest"` | | +| metrics.priorityClassName | string | `""` | | +| metrics.resources.limits.cpu | string | `"1"` | | +| metrics.resources.limits.memory | string | `"200Mi"` | | +| metrics.resources.requests.cpu | string | `"1"` | | +| metrics.resources.requests.memory | string | `"100Mi"` | | +| metrics.securityContext.allowPrivilegeEscalation | bool | `true` | | +| metrics.securityContext.capabilities.add[0] | string | `"SYS_ADMIN"` | | +| metrics.securityContext.capabilities.add[1] | string | `"NET_ADMIN"` | | +| metrics.securityContext.capabilities.add[2] | string | `"SYS_RESOURCE"` | | +| metrics.securityContext.capabilities.add[3] | string | `"BPF"` | | +| metrics.securityContext.capabilities.add[4] | string | `"SYS_PTRACE"` | | +| metrics.securityContext.privileged | bool | `true` | | +| metrics.tolerations | list | `[]` | | +| serviceAccountName | string | `"cortexflow-sa"` | | + +---------------------------------------------- +Autogenerated from chart metadata using [helm-docs v1.14.2](https://github.com/norwoodj/helm-docs/releases/v1.14.2) diff --git a/helm/templates/_helper.tpl b/helm/templates/_helper.tpl new file mode 100644 index 00000000..815611d8 --- /dev/null +++ b/helm/templates/_helper.tpl @@ -0,0 +1,31 @@ +{{/* +Sets tolerations for daemonsets either from the global var or from individual values +*/}} +{{- define "common.tolerations" }} +{{- $ctx := .context }} +{{- $component := .component }} +{{- $local := index $ctx.Values $component "tolerations" }} +{{- $global := $ctx.Values.global.tolerations }} +{{- if and (not (empty $local)) }} +tolerations: +{{ toYaml $local | indent 2 }} +{{- else if and (not (empty $global)) }} +tolerations: +{{ toYaml $global | indent 2 }} +{{- end }} +{{- end }} + +{{/* +Sets priorityClassName for daemonsets either from the global var or from individual values +*/}} +{{- define "common.priorityClassName" }} +{{- $ctx := .context }} +{{- $component := .component }} +{{- $local := index $ctx.Values $component "priorityClassName" }} +{{- $global := $ctx.Values.global.priorityClassName }} +{{- if and (not (empty $local)) }} +priorityClassName: {{ toYaml $local }} +{{- else if and (not (empty $global)) }} +priorityClassName: {{ toYaml $global }} +{{- end }} +{{- end }} diff --git a/helm/templates/agent.yaml b/helm/templates/agent.yaml new file mode 100644 index 00000000..b8036c01 --- /dev/null +++ b/helm/templates/agent.yaml @@ -0,0 +1,106 @@ +apiVersion: apps/v1 +kind: DaemonSet +metadata: + name: cortexflow-agent + labels: + app: cortexflow-agent +spec: + selector: + matchLabels: + app: cortexflow-agent + template: + metadata: + labels: + app: cortexflow-agent + annotations: + checksum/config: {{ include (print $.Template.BasePath "/configmap.yaml") . | sha256sum }} + spec: + serviceAccountName: {{ .Values.serviceAccountName }} + hostPID: true + hostNetwork: true + {{- include "common.tolerations" (dict "context" . "component" "agent") | indent 6 }} + {{- include "common.priorityClassName" (dict "context" . "component" "agent") | indent 6 }} + containers: + - name: agent + image: "{{ .Values.agent.image.repository }}:{{ .Values.agent.image.version }}" + command: ["/bin/bash", "-c"] + args: + - | + echo "Running on kernel $(uname -r)" + if [ ! -d "/sys/fs/bpf" ]; then + echo "ERROR: BPF filesystem not mounted" + exit 1 + else + echo "Checking ebpf path..." + ls -l /sys/fs/bpf + fi + echo "checking privileges" + ls -ld /sys/fs/bpf + + echo "checking if conntracker path" + ls -l /usr/src/cortexbrain-agent/conntracker + + echo "checking if the bpf maps are reachable" + ls -l /sys/fs/bpf/maps + + echo "Running application..." + exec /usr/local/bin/agent-api || echo "Application exited with code $?" + env: + - name: OTEL_SERVICE_NAME + value: cortexflow-agent + - name: OTEL_EXPORTER_OTLP_ENDPOINT + value: {{ .Values.global.otel.endpoint }} + - name: OTEL_EXPORTER_OTLP_PROTOCOL + value: {{ .Values.global.otel.protocol }} + - name: OTEL_RESOURCE_ATTRIBUTES + value: service.namespace=cortexflow,service.version=0.1.5 + - name: AGENT_API_ENABLE_REFLECTION + value: "true" + volumeMounts: + - name: bpf + mountPath: /sys/fs/bpf + mountPropagation: Bidirectional + readOnly: false + - name: proc + mountPath: /host/proc + readOnly: false + - name: kernel-dev + mountPath: /lib/modules + readOnly: false + resources: + {{- .Values.identity.resources | toYaml | nindent 12 }} + securityContext: + {{- .Values.agent.securityContext | toYaml | nindent 12 }} + volumes: + - name: bpf + hostPath: + path: /sys/fs/bpf + type: Directory + - name: proc + hostPath: + path: /proc + type: Directory + - name: kernel-dev + hostPath: + path: /lib/modules + type: Directory + +--- + +apiVersion: v1 +kind: Service +metadata: + name: cortexflow-agent + namespace: cortexflow +spec: + selector: + app: cortexflow-agent + ports: + - protocol: TCP + name: agent-server-port + port: 9090 + targetPort: 9090 + appProtocol: grpc + type: ClusterIP + +--- diff --git a/helm/templates/configmap-role.yaml b/helm/templates/configmap-role.yaml new file mode 100644 index 00000000..70b4f76b --- /dev/null +++ b/helm/templates/configmap-role.yaml @@ -0,0 +1,21 @@ +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: configmap-reader +rules: + - apiGroups: [""] + resources: ["configmaps","services"] + verbs: ["get", "list","watch"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: configmap-reader-binding +subjects: + - kind: ServiceAccount + name: {{ .Values.serviceAccountName }} + namespace: {{ .Release.Namespace }} +roleRef: + kind: Role + name: configmap-reader + apiGroup: rbac.authorization.k8s.io diff --git a/helm/templates/configmap.yaml b/helm/templates/configmap.yaml new file mode 100644 index 00000000..9a56dc6c --- /dev/null +++ b/helm/templates/configmap.yaml @@ -0,0 +1,6 @@ +apiVersion: v1 +kind: ConfigMap +metadata: + name: cortexbrain-client-config +data: + blocklist: {{ .Values.blocklist | quote }} diff --git a/helm/templates/identity.yaml b/helm/templates/identity.yaml new file mode 100644 index 00000000..a4dbf611 --- /dev/null +++ b/helm/templates/identity.yaml @@ -0,0 +1,120 @@ +apiVersion: apps/v1 +kind: DaemonSet +metadata: + name: cortexflow-identity + labels: + app: cortexflow-identity +spec: + selector: + matchLabels: + app: cortexflow-identity + template: + metadata: + labels: + app: cortexflow-identity + annotations: + checksum/config: {{ include (print $.Template.BasePath "/configmap.yaml") . | sha256sum }} + spec: + serviceAccountName: {{ .Values.serviceAccountName }} + hostPID: true + hostNetwork: true + {{- include "common.tolerations" (dict "context" . "component" "identity") | indent 6 }} + {{- include "common.priorityClassName" (dict "context" . "component" "identity") | indent 6 }} + initContainers: + - name: bpf-map-permissions + image: "{{ .Values.bpfMapPermissions.image.repository }}:{{ .Values.bpfMapPermissions.image.version }}" + command: ["/bin/bash","-c"] + args: + - | + echo "mounting the bpf path " + mount -t bpf bpf /sys/fs/bpf + + echo "checking permissions" + ls -ld /sys/fs/bpf + volumeMounts: + - name: bpf + mountPath: /sys/fs/bpf + mountPropagation: Bidirectional + readOnly: false + - name: kernel-dev + mountPath: /lib/modules + readOnly: false + - name: cgroup + mountPath: /sys/fs/cgroup + readOnly: true + securityContext: + {{- .Values.bpfMapPermissions.securityContext | toYaml | nindent 12}} + containers: + - name: identity + image: "{{ .Values.identity.image.repository }}:{{ .Values.identity.image.version }}" + command: ["/bin/bash", "-c"] + args: + - | + echo "Running on kernel $(uname -r)" + if [ ! -d "/sys/fs/bpf" ]; then + echo "ERROR: BPF filesystem not mounted" + exit 1 + else + echo "Checking ebpf path..." + ls -l /sys/fs/bpf + fi + echo "checking privileges" + ls -ld /sys/fs/bpf + + echo "Running application..." + exec /usr/local/bin/cortexflow-identity-service || echo "Application exited with code $?" + env: + - name: OTEL_SERVICE_NAME + value: cortexflow-identity + - name: OTEL_EXPORTER_OTLP_ENDPOINT + value: {{ .Values.global.otel.endpoint }} + - name: OTEL_EXPORTER_OTLP_PROTOCOL + value: {{ .Values.global.otel.protocol }} + - name: OTEL_RESOURCE_ATTRIBUTES + value: service.namespace=cortexflow,service.version=0.1.5 + resources: + {{- .Values.identity.resources | toYaml | nindent 12 }} + volumeMounts: + - name: bpf + mountPath: /sys/fs/bpf + mountPropagation: Bidirectional + readOnly: false + - name: kernel-dev + mountPath: /lib/modules + readOnly: false + - name: cgroup + mountPath: /sys/fs/cgroup + readOnly: true + securityContext: + {{- .Values.identity.securityContext | toYaml | nindent 12 }} + - name: bpftool-control-manager + image: "{{ .Values.bpfTool.image.repository }}:{{ .Values.bpfTool.image.version }}" + command: ["/bin/bash", "-c","sleep infinity"] + volumeMounts: + - name: bpf + mountPath: /sys/fs/bpf + mountPropagation: Bidirectional + readOnly: false + - name: kernel-dev + mountPath: /lib/modules + readOnly: false + - name: cgroup + mountPath: /sys/fs/cgroup + readOnly: true + resources: + {{- .Values.bpfTool.resources | toYaml | nindent 12 }} + securityContext: + {{- .Values.bpfTool.securityContext | toYaml | nindent 12 }} + volumes: + - name: bpf + hostPath: + path: /sys/fs/bpf + type: Directory + - name: kernel-dev + hostPath: + path: /lib/modules + type: Directory + - name: cgroup + hostPath: + path: /sys/fs/cgroup + type: Directory diff --git a/helm/templates/metrics.yaml b/helm/templates/metrics.yaml new file mode 100644 index 00000000..1f2e71de --- /dev/null +++ b/helm/templates/metrics.yaml @@ -0,0 +1,105 @@ +apiVersion: apps/v1 +kind: DaemonSet +metadata: + name: cortexflow-metrics + namespace: cortexflow + labels: + app: cortexflow-metrics +spec: + selector: + matchLabels: + app: cortexflow-metrics + template: + metadata: + labels: + app: cortexflow-metrics + annotations: + checksum/config: {{ include (print $.Template.BasePath "/configmap.yaml") . | sha256sum }} + spec: + serviceAccountName: {{ .Values.serviceAccountName }} + hostPID: true + hostNetwork: true + {{- include "common.tolerations" (dict "context" . "component" "metrics") | indent 6 }} + {{- include "common.priorityClassName" (dict "context" . "component" "metrics") | indent 6 }} + containers: + - name: metrics + image: "{{ .Values.metrics.image.repository }}:{{ .Values.metrics.image.version }}" + command: ["/bin/bash", "-c"] + args: + - | + echo "Running on kernel $(uname -r)" + if [ ! -d "/sys/fs/bpf" ]; then + echo "ERROR: BPF filesystem not mounted" + exit 1 + else + echo "Checking ebpf path..." + ls -l /sys/fs/bpf + fi + echo "checking privileges" + ls -ld /sys/fs/bpf + + echo "Running application..." + exec /usr/local/bin/cortexflow-metrics || echo "Application exited with code $?" + env: + - name: OTEL_SERVICE_NAME + value: cortexflow-metrics + - name: OTEL_EXPORTER_OTLP_ENDPOINT + value: {{ .Values.global.otel.endpoint }} + - name: OTEL_EXPORTER_OTLP_PROTOCOL + value: {{ .Values.global.otel.protocol }} + - name: OTEL_RESOURCE_ATTRIBUTES + value: service.namespace=cortexflow,service.version=0.1.5 + volumeMounts: + - name: bpf + mountPath: /sys/fs/bpf + mountPropagation: Bidirectional + readOnly: false + - name: proc + mountPath: /host/proc + readOnly: false + - name: kernel-dev + mountPath: /lib/modules + readOnly: false + - name: tracefs + mountPath: /sys/kernel/debug + readOnly: false + securityContext: + {{- .Values.metrics.securityContext | toYaml | nindent 12 }} + - name: bpftool-control-manager + image: "{{ .Values.bpfTool.image.repository }}:{{ .Values.bpfTool.image.version }}" + command: ["/bin/bash", "-c", "sleep infinity"] + volumeMounts: + - name: bpf + mountPath: /sys/fs/bpf + mountPropagation: Bidirectional + readOnly: false + - name: proc + mountPath: /host/proc + readOnly: false + - name: kernel-dev + mountPath: /lib/modules + readOnly: false + - name: tracefs + mountPath: /sys/kernel/debug + readOnly: false + resources: + {{- .Values.bpfTool.resources | toYaml | nindent 12 }} + securityContext: + {{- .Values.bpfTool.securityContext | toYaml | nindent 12 }} + volumes: + - name: bpf + hostPath: + path: /sys/fs/bpf + type: Directory + - name: proc + hostPath: + path: /proc + type: Directory + - name: kernel-dev + hostPath: + path: /lib/modules + type: Directory + - name: tracefs + hostPath: + path: /sys/kernel/debug + type: Directory diff --git a/helm/templates/serviceAccount.yaml b/helm/templates/serviceAccount.yaml new file mode 100644 index 00000000..aa7977de --- /dev/null +++ b/helm/templates/serviceAccount.yaml @@ -0,0 +1,4 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: {{ .Values.serviceAccountName }} diff --git a/helm/values.yaml b/helm/values.yaml new file mode 100644 index 00000000..f50852ea --- /dev/null +++ b/helm/values.yaml @@ -0,0 +1,117 @@ +global: + otel: + endpoint: http://localhost:4317 + protocol: grpc + tolerations: [] + priorityClassName: "" + +blocklist: "" +serviceAccountName: cortexflow-sa + +agent: + image: + repository: ghcr.io/cortexflow/agent + version: latest + resources: + limits: + memory: "200Mi" + requests: + cpu: "100m" + memory: "100Mi" + securityContext: + privileged: true + allowPrivilegeEscalation: true + capabilities: + add: + - SYS_ADMIN + - NET_ADMIN + - SYS_RESOURCE + - BPF + - SYS_PTRACE + tolerations: [] + priorityClassName: "" + +identity: + image: + repository: ghcr.io/cortexflow/identity + version: latest + resources: + limits: + memory: "200Mi" + requests: + cpu: "100m" + memory: "100Mi" + securityContext: + privileged: true + allowPrivilegeEscalation: true + capabilities: + add: + - SYS_ADMIN + - NET_ADMIN + - SYS_RESOURCE + - BPF + - SYS_PTRACE + tolerations: [] + priorityClassName: "" + +metrics: + image: + repository: ghcr.io/cortexflow/metrics + version: latest + resources: + limits: + cpu: "1" + memory: "200Mi" + requests: + cpu: "1" + memory: "100Mi" + securityContext: + privileged: true + allowPrivilegeEscalation: true + capabilities: + add: + - SYS_ADMIN + - NET_ADMIN + - SYS_RESOURCE + - BPF + - SYS_PTRACE + tolerations: [] + priorityClassName: "" + +bpfTool: + image: + repository: danielpacak/bpftool-runner + version: latest + resources: + limits: + cpu: "1" + memory: "200Mi" + requests: + cpu: "1" + memory: "100Mi" + securityContext: + privileged: true + allowPrivilegeEscalation: true + capabilities: + add: + - SYS_ADMIN + - NET_ADMIN + - SYS_RESOURCE + - BPF + - SYS_PTRACE + +bpfMapPermissions: + image: + repository: ubuntu + version: "24.04" + securityContext: + runAsUser: 0 + privileged: true + allowPrivilegeEscalation: true + capabilities: + add: + - SYS_ADMIN + - NET_ADMIN + - SYS_RESOURCE + - BPF + - SYS_PTRACE diff --git a/mcp/Cargo.lock b/mcp/Cargo.lock index e3f3bb17..f03d11f4 100644 --- a/mcp/Cargo.lock +++ b/mcp/Cargo.lock @@ -17,12 +17,6 @@ dependencies = [ "memchr", ] -[[package]] -name = "allocator-api2" -version = "0.2.21" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" - [[package]] name = "android_system_properties" version = "0.1.5" @@ -38,12 +32,6 @@ version = "1.0.104" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "330a5ed07fa54e4702c9d6c4174f74427fc0ef6e214bbd677ae50a5099946470" -[[package]] -name = "assert_matches" -version = "1.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9b34d609dfbaf33d6889b2b7106d3ca345eacad44200913df5ba02bfd31d2ba9" - [[package]] name = "async-compression" version = "0.4.42" @@ -56,17 +44,6 @@ dependencies = [ "tokio", ] -[[package]] -name = "async-trait" -version = "0.1.91" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ae36dc4177970ef04fde5178d3e2429882def40e57a451f919c098f72baa6cec" -dependencies = [ - "proc-macro2", - "quote", - "syn 3.0.3", -] - [[package]] name = "atomic-waker" version = "1.1.2" @@ -102,37 +79,6 @@ dependencies = [ "pkg-config", ] -[[package]] -name = "aya" -version = "0.13.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d18bc4e506fbb85ab7392ed993a7db4d1a452c71b75a246af4a80ab8c9d2dd50" -dependencies = [ - "assert_matches", - "aya-obj", - "bitflags", - "bytes", - "libc", - "log", - "object", - "once_cell", - "thiserror 1.0.69", -] - -[[package]] -name = "aya-obj" -version = "0.2.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c51b96c5a8ed8705b40d655273bc4212cbbf38d4e3be2788f36306f154523ec7" -dependencies = [ - "bytes", - "core-error", - "hashbrown 0.15.5", - "log", - "object", - "thiserror 1.0.69", -] - [[package]] name = "base64" version = "0.22.1" @@ -160,23 +106,6 @@ version = "3.20.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" -[[package]] -name = "bytemuck" -version = "1.25.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "95832e849adfb21180ccb6826a99da14e5d266ae5c2e668e1602cf234f153797" - -[[package]] -name = "bytemuck_derive" -version = "1.11.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f65693059b6b9c588b9f62fed1cedbf0a8b805631457ea162d68f0de186f3de5" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.119", -] - [[package]] name = "bytes" version = "1.12.1" @@ -215,7 +144,7 @@ checksum = "d524456ba66e72eb8b115ff89e01e497f8e6d11d78b70b1aa13c0fbd97540a81" dependencies = [ "cfg-if", "cpufeatures", - "rand_core 0.10.1", + "rand_core", ] [[package]] @@ -277,15 +206,6 @@ dependencies = [ "unicode-segmentation", ] -[[package]] -name = "core-error" -version = "0.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "efcdb2972eb64230b4c50646d8498ff73f5128d196a90c7236eec4cbe8619b8f" -dependencies = [ - "version_check", -] - [[package]] name = "core-foundation" version = "0.9.4" @@ -312,29 +232,6 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" -[[package]] -name = "cortexbrain-common" -version = "0.1.4" -dependencies = [ - "anyhow", - "aya", - "bytemuck", - "bytemuck_derive", - "bytes", - "k8s-openapi", - "kube", - "opentelemetry", - "opentelemetry-appender-tracing", - "opentelemetry-otlp", - "opentelemetry-stdout", - "opentelemetry_sdk", - "serde", - "serde_json", - "tokio", - "tracing", - "tracing-subscriber", -] - [[package]] name = "cpufeatures" version = "0.3.0" @@ -483,12 +380,6 @@ version = "1.0.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d0881ea181b1df73ff77ffaaf9c7544ecc11e82fba9b5f27b262a3c73a332555" -[[package]] -name = "either" -version = "1.17.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e5e8f6c15a24b9a3ee5efec809ccd006d3b30e8b3bb63c39af737c7f87daa1d" - [[package]] name = "encoding_rs" version = "0.8.35" @@ -547,12 +438,6 @@ version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" -[[package]] -name = "foldhash" -version = "0.1.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" - [[package]] name = "form_urlencoded" version = "1.2.2" @@ -695,18 +580,6 @@ dependencies = [ "wasm-bindgen", ] -[[package]] -name = "getrandom" -version = "0.3.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" -dependencies = [ - "cfg-if", - "libc", - "r-efi 5.3.0", - "wasip2", -] - [[package]] name = "getrandom" version = "0.4.3" @@ -716,8 +589,8 @@ dependencies = [ "cfg-if", "js-sys", "libc", - "r-efi 6.0.0", - "rand_core 0.10.1", + "r-efi", + "rand_core", "wasm-bindgen", ] @@ -746,17 +619,6 @@ version = "0.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" -[[package]] -name = "hashbrown" -version = "0.15.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" -dependencies = [ - "allocator-api2", - "equivalent", - "foldhash", -] - [[package]] name = "hashbrown" version = "0.17.1" @@ -775,15 +637,6 @@ version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" -[[package]] -name = "home" -version = "0.5.12" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cc627f471c528ff0c4a49e1d5e60450c8f6461dd6d10ba9dcd3a61d3dff7728d" -dependencies = [ - "windows-sys 0.61.2", -] - [[package]] name = "http" version = "1.4.2" @@ -853,27 +706,12 @@ dependencies = [ "http", "hyper", "hyper-util", - "log", "rustls", - "rustls-native-certs", "tokio", "tokio-rustls", "tower-service", ] -[[package]] -name = "hyper-timeout" -version = "0.5.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2b90d566bffbce6a75bd8b09a05aa8c2cb1fabb6cb348f8840c9e4c90a0d83b0" -dependencies = [ - "hyper", - "hyper-util", - "pin-project-lite", - "tokio", - "tower-service", -] - [[package]] name = "hyper-util" version = "0.1.20" @@ -1061,15 +899,6 @@ version = "2.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d98f6fed1fde3f8c21bc40a1abb88dd75e67924f9cffc3ef95607bad8017f8e2" -[[package]] -name = "itertools" -version = "0.14.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2b192c782037fadd9cfa75548310488aabdbf3d2da73885b31bd0abd03351285" -dependencies = [ - "either", -] - [[package]] name = "itoa" version = "1.0.18" @@ -1088,7 +917,7 @@ dependencies = [ "jni-sys", "log", "simd_cesu8", - "thiserror 2.0.19", + "thiserror", "walkdir", "windows-link 0.2.1", ] @@ -1146,101 +975,6 @@ dependencies = [ "wasm-bindgen", ] -[[package]] -name = "jsonpath-rust" -version = "0.7.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0c00ae348f9f8fd2d09f82a98ca381c60df9e0820d8d79fce43e649b4dc3128b" -dependencies = [ - "pest", - "pest_derive", - "regex", - "serde_json", - "thiserror 2.0.19", -] - -[[package]] -name = "k8s-openapi" -version = "0.26.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "06d9e5e61dd037cdc51da0d7e2b2be10f497478ea7e120d85dad632adb99882b" -dependencies = [ - "base64", - "chrono", - "serde", - "serde_json", -] - -[[package]] -name = "kube" -version = "2.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "48e7bb0b6a46502cc20e4575b6ff401af45cfea150b34ba272a3410b78aa014e" -dependencies = [ - "k8s-openapi", - "kube-client", - "kube-core", -] - -[[package]] -name = "kube-client" -version = "2.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4987d57a184d2b5294fdad3d7fc7f278899469d21a4da39a8f6ca16426567a36" -dependencies = [ - "base64", - "bytes", - "chrono", - "either", - "futures", - "home", - "http", - "http-body", - "http-body-util", - "hyper", - "hyper-rustls", - "hyper-timeout", - "hyper-util", - "jsonpath-rust", - "k8s-openapi", - "kube-core", - "pem", - "rustls", - "secrecy", - "serde", - "serde_json", - "serde_yaml", - "thiserror 2.0.19", - "tokio", - "tokio-util", - "tower", - "tower-http", - "tracing", -] - -[[package]] -name = "kube-core" -version = "2.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "914bbb770e7bb721a06e3538c0edd2babed46447d128f7c21caa68747060ee73" -dependencies = [ - "chrono", - "derive_more", - "form_urlencoded", - "http", - "k8s-openapi", - "serde", - "serde-value", - "serde_json", - "thiserror 2.0.19", -] - -[[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.189" @@ -1274,21 +1008,11 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" -[[package]] -name = "matchers" -version = "0.2.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d1525a2a28c7f4fa0fc98bb91ae755d1e2d1505079e05539e35bc876b5d65ae9" -dependencies = [ - "regex-automata", -] - [[package]] name = "mcp" version = "0.1.0" dependencies = [ "anyhow", - "cortexbrain-common", "dotenv", "genai", "reqwest", @@ -1370,15 +1094,6 @@ dependencies = [ "minimal-lexical", ] -[[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" @@ -1394,18 +1109,6 @@ dependencies = [ "autocfg", ] -[[package]] -name = "object" -version = "0.36.7" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "62948e14d923ea95ea2c7c86c71013138b66525b86bdc08d2dcc262bdb497b87" -dependencies = [ - "crc32fast", - "hashbrown 0.15.5", - "indexmap 2.14.0", - "memchr", -] - [[package]] name = "once_cell" version = "1.21.4" @@ -1418,115 +1121,6 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" -[[package]] -name = "opentelemetry" -version = "0.32.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b0142c63252a9e054e68a4c61a5778f7b14f576274d593f8ce883d191a099682" -dependencies = [ - "futures-core", - "futures-sink", - "js-sys", - "pin-project-lite", - "thiserror 2.0.19", - "tracing", -] - -[[package]] -name = "opentelemetry-appender-tracing" -version = "0.32.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2c0080f0dc1d7c786f467cd85a4e395fcab11ee852004f39a29a18ab7c25d837" -dependencies = [ - "opentelemetry", - "tracing", - "tracing-core", - "tracing-subscriber", -] - -[[package]] -name = "opentelemetry-http" -version = "0.32.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5683015d09e2df236ef005b17f6f196f0d5f6313c4fa43a7b6a53b52776e4331" -dependencies = [ - "async-trait", - "bytes", - "http", - "opentelemetry", - "reqwest", -] - -[[package]] -name = "opentelemetry-otlp" -version = "0.32.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9966929966d17620d7c316c643ba62631826e10021409357772d5eea84f62c35" -dependencies = [ - "http", - "opentelemetry", - "opentelemetry-http", - "opentelemetry-proto", - "opentelemetry_sdk", - "prost", - "reqwest", - "thiserror 2.0.19", - "tokio", - "tonic", - "tonic-types", -] - -[[package]] -name = "opentelemetry-proto" -version = "0.32.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "56d658ba1faf63f7b9c492cfbe6e0ec365440a16132d3270c1065f7b33f1b638" -dependencies = [ - "opentelemetry", - "opentelemetry_sdk", - "prost", - "tonic", - "tonic-prost", -] - -[[package]] -name = "opentelemetry-stdout" -version = "0.32.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a1b1c6a247d79091f0062a5f4bd058589525cf987a8d4c169440d9c1be72f0ad" -dependencies = [ - "chrono", - "opentelemetry", - "opentelemetry_sdk", -] - -[[package]] -name = "opentelemetry_sdk" -version = "0.32.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9b59f80e1ac4d5ff7a2db8fb6c80badb7f0f3f858211fba08dd9aaec750894f9" -dependencies = [ - "futures-channel", - "futures-executor", - "futures-util", - "opentelemetry", - "percent-encoding", - "portable-atomic", - "rand 0.9.5", - "thiserror 2.0.19", - "tokio", - "tokio-stream", -] - -[[package]] -name = "ordered-float" -version = "2.10.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "68f19d67e5a2795c94e73e0bb1cc1a7edeb2e28efd39e2e1c9b7a40c1108b11c" -dependencies = [ - "num-traits", -] - [[package]] name = "parking_lot" version = "0.12.5" @@ -1556,84 +1150,12 @@ version = "1.0.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "57c0d7b74b563b49d38dae00a0c37d4d6de9b432382b2892f0574ddcae73fd0a" -[[package]] -name = "pem" -version = "3.0.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1d30c53c26bc5b31a98cd02d20f25a7c8567146caf63ed593a9d87b2775291be" -dependencies = [ - "base64", - "serde_core", -] - [[package]] name = "percent-encoding" version = "2.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" -[[package]] -name = "pest" -version = "2.8.8" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7df728be843c7070fab6ab7c328c4e9e9d78e23bf749c0669c86ee7ebfa050a2" -dependencies = [ - "memchr", - "ucd-trie", -] - -[[package]] -name = "pest_derive" -version = "2.8.8" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e2dd6fc3b26b3462ee188aac870f5a41d398f1cd5e2408d16531bd71c9591fd" -dependencies = [ - "pest", - "pest_generator", -] - -[[package]] -name = "pest_generator" -version = "2.8.8" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a7a9205cfb6f596a9e8b689c0a15f9ceb7a1aafae7aaf788150ac65b29975b6" -dependencies = [ - "pest", - "pest_meta", - "proc-macro2", - "quote", - "syn 2.0.119", -] - -[[package]] -name = "pest_meta" -version = "2.8.8" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "85abd351c0de1e8384fc791a0737111a350394937e92b956b743dac12429f57c" -dependencies = [ - "pest", -] - -[[package]] -name = "pin-project" -version = "1.1.13" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2466b2336ed02bcdca6b294417127b90ec92038d1d5c4fbeac971a922e0e0924" -dependencies = [ - "pin-project-internal", -] - -[[package]] -name = "pin-project-internal" -version = "1.1.13" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c96395f0a926bc13b1c17622aaddda1ecb55d49c8f1bf9777e4d877800a43f8b" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.119", -] - [[package]] name = "pin-project-lite" version = "0.2.17" @@ -1646,12 +1168,6 @@ version = "0.3.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e" -[[package]] -name = "portable-atomic" -version = "1.14.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3d20d5497ef88037a52ff98267d066e7f11fcc5e99bbfbd58a42336193aacec3" - [[package]] name = "potential_utf" version = "0.1.5" @@ -1667,15 +1183,6 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" -[[package]] -name = "ppv-lite86" -version = "0.2.21" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" -dependencies = [ - "zerocopy", -] - [[package]] name = "proc-macro2" version = "1.0.107" @@ -1699,38 +1206,6 @@ dependencies = [ "windows", ] -[[package]] -name = "prost" -version = "0.14.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "528ac67416ff8646872a3c02cad9cc4ee5dc9f9540c9b10771855c95cb2e5ae1" -dependencies = [ - "bytes", - "prost-derive", -] - -[[package]] -name = "prost-derive" -version = "0.14.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" -dependencies = [ - "anyhow", - "itertools", - "proc-macro2", - "quote", - "syn 2.0.119", -] - -[[package]] -name = "prost-types" -version = "0.14.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f94967dc7688f3054c7fac87473ffae4cc4c3904800e2d9f5b857246d8963b0a" -dependencies = [ - "prost", -] - [[package]] name = "quinn" version = "0.11.11" @@ -1745,7 +1220,7 @@ dependencies = [ "rustc-hash", "rustls", "socket2", - "thiserror 2.0.19", + "thiserror", "tokio", "tracing", "web-time", @@ -1761,14 +1236,14 @@ dependencies = [ "bytes", "getrandom 0.4.3", "lru-slab", - "rand 0.10.2", + "rand", "rand_pcg", "ring", "rustc-hash", "rustls", "rustls-pki-types", "slab", - "thiserror 2.0.19", + "thiserror", "tinyvec", "tracing", "web-time", @@ -1797,28 +1272,12 @@ dependencies = [ "proc-macro2", ] -[[package]] -name = "r-efi" -version = "5.3.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" - [[package]] name = "r-efi" version = "6.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" -[[package]] -name = "rand" -version = "0.9.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b9ef1d0d795eb7d84685bca4f72f3649f064e6641543d3a8c415898726a57b41" -dependencies = [ - "rand_chacha", - "rand_core 0.9.5", -] - [[package]] name = "rand" version = "0.10.2" @@ -1827,26 +1286,7 @@ checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" dependencies = [ "chacha20", "getrandom 0.4.3", - "rand_core 0.10.1", -] - -[[package]] -name = "rand_chacha" -version = "0.9.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" -dependencies = [ - "ppv-lite86", - "rand_core 0.9.5", -] - -[[package]] -name = "rand_core" -version = "0.9.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c" -dependencies = [ - "getrandom 0.3.4", + "rand_core", ] [[package]] @@ -1861,7 +1301,7 @@ version = "0.10.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "caa0f4137e1c0a72f4c651489402276c8e8e1cf081f3b0ba156d2cbeef09e86a" dependencies = [ - "rand_core 0.10.1", + "rand_core", ] [[package]] @@ -1931,7 +1371,6 @@ dependencies = [ "base64", "bytes", "encoding_rs", - "futures-channel", "futures-core", "futures-util", "h2", @@ -1997,7 +1436,7 @@ dependencies = [ "schemars 1.2.1", "serde", "serde_json", - "thiserror 2.0.19", + "thiserror", "tokio", "tokio-stream", "tokio-util", @@ -2039,9 +1478,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3c54fcab019b409d04215d3a17cb438fd7fbf192ee61461f20f4fe18704bc138" dependencies = [ "aws-lc-rs", - "log", "once_cell", - "ring", "rustls-pki-types", "rustls-webpki", "subtle", @@ -2183,15 +1620,6 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" -[[package]] -name = "secrecy" -version = "0.10.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e891af845473308773346dc847b2c23ee78fe442e0472ac50e22a18a93d3ae5a" -dependencies = [ - "zeroize", -] - [[package]] name = "security-framework" version = "3.7.0" @@ -2231,16 +1659,6 @@ dependencies = [ "serde_derive", ] -[[package]] -name = "serde-value" -version = "0.7.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f3a1a3341211875ef120e117ea7fd5228530ae7e7036a779fdc9117be6b3282c" -dependencies = [ - "ordered-float", - "serde", -] - [[package]] name = "serde_core" version = "1.0.229" @@ -2329,28 +1747,6 @@ dependencies = [ "syn 2.0.119", ] -[[package]] -name = "serde_yaml" -version = "0.9.34+deprecated" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a8b1a1a2ebf674015cc02edccce75287f1a0130d394307b36743c2f5d504b47" -dependencies = [ - "indexmap 2.14.0", - "itoa", - "ryu", - "serde", - "unsafe-libyaml", -] - -[[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" @@ -2513,33 +1909,13 @@ dependencies = [ "libc", ] -[[package]] -name = "thiserror" -version = "1.0.69" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b6aaf5339b578ea85b50e080feb250a3e8ae8cfcdff9a461c9ec2904bc923f52" -dependencies = [ - "thiserror-impl 1.0.69", -] - [[package]] name = "thiserror" version = "2.0.19" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09a43598840e33d5b0331f38c5e30d13bb11c11210a4b58f0d9b18a5a5eefcd9" dependencies = [ - "thiserror-impl 2.0.19", -] - -[[package]] -name = "thiserror-impl" -version = "1.0.69" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4fee6c4efc90059e10f81e6d42c60a18f76588c3d74cb83a0b242a2b6c7504c1" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.119", + "thiserror-impl", ] [[package]] @@ -2553,15 +1929,6 @@ dependencies = [ "syn 3.0.3", ] -[[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.54" @@ -2680,54 +2047,6 @@ dependencies = [ "tokio", ] -[[package]] -name = "tonic" -version = "0.14.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ac2a5518c70fa84342385732db33fb3f44bc4cc748936eb5833d2df34d6445ef" -dependencies = [ - "async-trait", - "base64", - "bytes", - "http", - "http-body", - "http-body-util", - "hyper", - "hyper-timeout", - "hyper-util", - "percent-encoding", - "pin-project", - "sync_wrapper", - "tokio", - "tokio-stream", - "tower", - "tower-layer", - "tower-service", - "tracing", -] - -[[package]] -name = "tonic-prost" -version = "0.14.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "50849f68853be452acf590cde0b146665b8d507b3b8af17261df47e02c209ea0" -dependencies = [ - "bytes", - "prost", - "tonic", -] - -[[package]] -name = "tonic-types" -version = "0.14.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "73ab1b02061f83d519bba3caa167f88f261ef05720ab8ebc954ade70de3348e8" -dependencies = [ - "prost", - "prost-types", - "tonic", -] - [[package]] name = "tower" version = "0.5.3" @@ -2736,15 +2055,11 @@ checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4" dependencies = [ "futures-core", "futures-util", - "indexmap 2.14.0", "pin-project-lite", - "slab", "sync_wrapper", "tokio", - "tokio-util", "tower-layer", "tower-service", - "tracing", ] [[package]] @@ -2754,7 +2069,6 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840" dependencies = [ "async-compression", - "base64", "bitflags", "bytes", "futures-core", @@ -2762,14 +2076,12 @@ dependencies = [ "http", "http-body", "http-body-util", - "mime", "pin-project-lite", "tokio", "tokio-util", "tower", "tower-layer", "tower-service", - "tracing", "url", ] @@ -2791,7 +2103,6 @@ version = "0.1.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100" dependencies = [ - "log", "pin-project-lite", "tracing-attributes", "tracing-core", @@ -2815,36 +2126,6 @@ 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]] @@ -2853,12 +2134,6 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" -[[package]] -name = "ucd-trie" -version = "0.1.7" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2896d95c02a80c6d6a5d6e953d479f5ddf2dfdb6a244441010e373ac0fb88971" - [[package]] name = "unicase" version = "2.9.0" @@ -2883,12 +2158,6 @@ version = "0.2.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ebc1c04c71510c7f702b52b7c350734c9ff1295c464a03335b00bb84fc54f853" -[[package]] -name = "unsafe-libyaml" -version = "0.2.11" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "673aac59facbab8a9007c7f6108d11f63b603f7cabff99fabf650fea5c32b861" - [[package]] name = "untrusted" version = "0.9.0" @@ -2924,12 +2193,6 @@ dependencies = [ "wasm-bindgen", ] -[[package]] -name = "valuable" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" - [[package]] name = "value-ext" version = "0.1.5" @@ -2941,12 +2204,6 @@ dependencies = [ "serde_json", ] -[[package]] -name = "version_check" -version = "0.9.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" - [[package]] name = "walkdir" version = "2.5.0" @@ -2972,15 +2229,6 @@ version = "0.11.1+wasi-snapshot-preview1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" -[[package]] -name = "wasip2" -version = "1.0.4+wasi-0.2.12" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b67efb37e106e55ce722a510d6b5f9c17f083e5fc79afc2badeb12cc313d9487" -dependencies = [ - "wit-bindgen", -] - [[package]] name = "wasm-bindgen" version = "0.2.126" @@ -3328,12 +2576,6 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" -[[package]] -name = "wit-bindgen" -version = "0.57.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" - [[package]] name = "writeable" version = "0.6.3" @@ -3363,26 +2605,6 @@ dependencies = [ "synstructure", ] -[[package]] -name = "zerocopy" -version = "0.8.55" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b5a105cd7b140f6eeec8acff2ea38135d3cab283ada58540f629fe51e46696eb" -dependencies = [ - "zerocopy-derive", -] - -[[package]] -name = "zerocopy-derive" -version = "0.8.55" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0fe976fb70c78cd64cccfe3a6fc142244e8a77b70959b30faf9d0ac37ee228eb" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.119", -] - [[package]] name = "zerofrom" version = "0.1.8" diff --git a/mcp/Cargo.toml b/mcp/Cargo.toml index 439903f7..c6ef42a6 100644 --- a/mcp/Cargo.toml +++ b/mcp/Cargo.toml @@ -13,4 +13,4 @@ serde_json = "1.0.151" reqwest = {version="0.13.4",features=["json","query"]} schemars = "1.2.1" serde = "1.0.229" -cortexbrain-common={path="../core/common"} +# TEMP-CHECK cortexbrain-common diff --git a/mcp/src/prometheus.rs b/mcp/src/prometheus.rs index 716e9cd8..99e2fe88 100644 --- a/mcp/src/prometheus.rs +++ b/mcp/src/prometheus.rs @@ -1,91 +1,574 @@ -use anyhow::{Ok, Result}; -/// main prometheus client -/// accept a request as a query +use anyhow::{bail, Context, Result}; use reqwest::Client; +use serde::{Deserialize, Serialize}; +use serde_json::{Map, Value}; +use std::collections::BTreeMap; +use std::sync::Arc; +use tokio::sync::OnceCell; + +/// Endpoint used when `PROMETHEUS_URL` is not set. +pub const DEFAULT_PROMETHEUS_SERVER: &str = "http://localhost:9090"; + +/// Labels able to identify the workload a series belongs to, most specific +/// first. Docker hosts expose `container_name`, Kubernetes exposes +/// `k8s_pod_name`; the `transform/identity` processor in the otel-collector +/// configs backfills the former from the latter, so on a current deployment +/// the first entry always wins. The remaining entries keep the server usable +/// against collectors that have not been rolled out yet. +pub const IDENTITY_LABELS: [&str; 3] = ["container_name", "k8s_pod_name", "k8s_namespace_name"]; + +/// Counter used for discovery queries on the busiest signal. +pub const EVENTS_TOTAL: &str = "cortexbrain_events_total"; + +/// Prometheus histogram of the OTel instrument kind, which decides the +/// aggregation applied to a metric. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MetricKind { + /// Monotonic counter, aggregated with `rate()`. + Counter, + /// Instantaneous value, aggregated with `avg_over_time()`. `rate()` is not + /// valid here: it is defined for counters and yields meaningless output + /// for gauges. + Gauge, +} + +/// How the workload argument is compared against the identity label. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum MatchMode { + /// `label=~".*workload.*"`, so `otel-agent` also matches + /// `otel-agent-qsm7c`. + Substring, + /// `label="workload"`, an exact match. + Exact, +} + +impl MatchMode { + pub fn parse(value: Option<&str>) -> Result { + match value.map(str::trim).filter(|v| !v.is_empty()) { + None | Some("substring") => Ok(Self::Substring), + Some("exact") => Ok(Self::Exact), + Some(other) => bail!("invalid match_mode `{other}`, expected `substring` or `exact`"), + } + } + + fn selector(self, label: &str, workload: &str) -> String { + match self { + Self::Substring => format!("{label}=~\".*{}.*\"", escape_regex(workload)), + Self::Exact => format!("{label}=\"{}\"", escape_string(workload)), + } + } +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct TimeSeriesPoint { + pub timestamp: i64, + pub value: f64, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct TimeSeries { + /// The label set of the series, including the identity label, so a + /// response is self-describing. + pub labels: Map, + pub values: Vec, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct TimeSeriesResponse { + pub metric_name: String, + /// Label the workload was matched on. + pub identity_label: String, + pub workload: String, + pub match_mode: String, + /// Range vector used for the aggregation, e.g. `5m`. + pub window: String, + pub start: String, + pub end: String, + pub step: String, + pub series: Vec, +} + +/// Workload names discoverable in Prometheus, keyed by label. +#[derive(Debug, Serialize, Deserialize)] +pub struct WorkloadListing { + /// Label the metric tools will use for the current deployment. + pub identity_label: String, + pub labels: BTreeMap>, + /// Workloads with at least one event inside the discovery window. + pub active: Vec, +} + #[derive(Clone)] pub struct PromClient { baseurl: String, client: Client, + /// Resolved once per process: the label set does not change while the + /// server runs, and a failed query would otherwise re-probe Prometheus on + /// every call. + identity_label: Arc>, } -const PROMETHEUS_SERVER: &str = "http://localhost:9090/api/v1/query"; - impl PromClient { pub fn new() -> Result { + let baseurl = std::env::var("PROMETHEUS_URL") + .unwrap_or_else(|_| DEFAULT_PROMETHEUS_SERVER.to_string()); + Ok(PromClient { client: Client::new(), - baseurl: PROMETHEUS_SERVER.to_string(), + baseurl, + identity_label: Arc::new(OnceCell::new()), }) } - /// creates a query using promQL language - pub async fn query(&self, promql_query: &str) -> Result { + + /// Label used to select workloads, probed once and cached. + pub async fn identity_label(&self) -> Result<&str> { + self.identity_label + .get_or_try_init(|| async { + for candidate in IDENTITY_LABELS { + if !self.label_values(candidate).await?.is_empty() { + return Ok(candidate.to_string()); + } + } + bail!( + "no cortexbrain series found in {} for any of {:?}; the metrics \ + service may not be running, or the otel-collector may be \ + missing the transform/identity processor", + self.baseurl, + IDENTITY_LABELS + ) + }) + .await + .map(String::as_str) + } + + /// Values currently present for a label. + pub async fn label_values(&self, label: &str) -> Result> { let response = self .client - .get(format!("{}", self.baseurl)) + .get(format!("{}/api/v1/label/{label}/values", self.baseurl)) + .send() + .await + .with_context(|| format!("GET /api/v1/label/{label}/values failed"))? + .error_for_status() + .with_context(|| format!("GET /api/v1/label/{label}/values returned an error status"))?; + + let body: Value = response + .json() + .await + .with_context(|| format!("GET /api/v1/label/{label}/values returned invalid JSON"))?; + + check_status(&body).with_context(|| format!("listing values of label `{label}`"))?; + + Ok(body + .get("data") + .and_then(Value::as_array) + .map(|values| { + values + .iter() + .filter_map(Value::as_str) + .map(str::to_string) + .collect() + }) + .unwrap_or_default()) + } + + /// Instant query built from a PromQL expression. + pub async fn query(&self, promql_query: &str) -> Result { + let response = self + .client + .get(format!("{}/api/v1/query", self.baseurl)) .query(&[("query", promql_query)]) .send() - .await?; - Ok(response.json().await?) + .await + .with_context(|| format!("GET /api/v1/query failed for `{promql_query}`"))? + .error_for_status() + .with_context(|| format!("query `{promql_query}` returned an error status"))?; + + let body: Value = response + .json() + .await + .with_context(|| format!("query `{promql_query}` returned invalid JSON"))?; + + check_status(&body).with_context(|| format!("running query `{promql_query}`"))?; + + Ok(body) } - pub async fn query_get_cpu_bytes(&self, container_name: &str, timeframe: &str) -> Result { - //query: sum by(container_name) (rate(cortexbrain_cpu_bytes_alloc[10m])) - let promql = format!( - r#"sum by(container_name) (rate(cortexbrain_cpu_bytes_alloc{{container_name=~".*{container_name}.*"}}[{timeframe}]))"# - ); - let res = serde_json::to_string_pretty(&self.query(&promql).await?)?; - Ok(res) + /// Range query built from a PromQL expression. + pub async fn query_range( + &self, + promql_query: &str, + start: &str, + end: &str, + step: &str, + ) -> Result { + let response = self + .client + .get(format!("{}/api/v1/query_range", self.baseurl)) + .query(&[ + ("query", promql_query), + ("start", start), + ("end", end), + ("step", step), + ]) + .send() + .await + .with_context(|| format!("GET /api/v1/query_range failed for `{promql_query}`"))? + .error_for_status() + .with_context(|| format!("range query `{promql_query}` returned an error status"))?; + + let body: Value = response + .json() + .await + .with_context(|| format!("range query `{promql_query}` returned invalid JSON"))?; + + check_status(&body).with_context(|| format!("running range query `{promql_query}`"))?; + + Ok(body) } - pub async fn query_get_memory_allocated_bytes( + + /// Run a windowed aggregation over one metric and normalise the matrix + /// into a typed response. + #[allow(clippy::too_many_arguments)] + pub async fn windowed_query( &self, - container_name: &str, - timeframe: &str, - ) -> Result { - //query: sum by(container_name) (rate(cortexbrain_enter_mem_alloc{container_name="..."}[10m])) - let promql = format!( - r#"sum by(container_name) (rate(cortexbrain_enter_mem_alloc{{container_name=~".*{container_name}.*"}}[{timeframe}]))"# + metric: &str, + kind: MetricKind, + identity_label: &str, + mode: MatchMode, + workload: &str, + window: &str, + start: &str, + end: &str, + step: &str, + ) -> Result { + let promql = build_promql(metric, kind, identity_label, mode, workload, window); + let raw = self.query_range(&promql, start, end, step).await?; + let series = parse_series(&raw, identity_label); + + Ok(TimeSeriesResponse { + metric_name: metric.to_string(), + identity_label: identity_label.to_string(), + workload: workload.to_string(), + match_mode: match mode { + MatchMode::Substring => "substring", + MatchMode::Exact => "exact", + } + .to_string(), + window: window.to_string(), + start: start.to_string(), + end: end.to_string(), + step: step.to_string(), + series, + }) + } + + /// Workloads visible in Prometheus, both as raw label values and as those + /// that reported events inside the discovery window. + pub async fn list_workloads(&self, lookback: &str) -> Result { + let identity_label = self.identity_label().await?.to_string(); + + let mut labels = BTreeMap::new(); + for label in IDENTITY_LABELS { + labels.insert(label.to_string(), self.label_values(label).await?); + } + + let active_promql = format!("group by({identity_label}) ({EVENTS_TOTAL}[{lookback}])"); + let mut active = Vec::new(); + if let Ok(raw) = self.query(&active_promql).await { + for series in result_vector(&raw) { + if let Some(name) = series.get(&identity_label).and_then(Value::as_str) { + active.push(name.to_string()); + } + } + } + active.sort(); + active.dedup(); + + Ok(WorkloadListing { + identity_label, + labels, + active, + }) + } +} + +/// Build the aggregation for one metric. Counters use `rate()`, gauges use +/// `avg_over_time()`, and both are summed by the identity label so a +/// multi-match regex collapses every matching workload into one series. +pub fn build_promql( + metric: &str, + kind: MetricKind, + identity_label: &str, + mode: MatchMode, + workload: &str, + window: &str, +) -> String { + let selector = mode.selector(identity_label, workload); + let inner = match kind { + MetricKind::Counter => format!("rate({metric}{{{selector}}}[{window}])"), + MetricKind::Gauge => format!("avg_over_time({metric}{{{selector}}}[{window}])"), + }; + format!("sum by({identity_label}) ({inner})") +} + +/// Escape a value interpolated into a PromQL string or a regex literal. +fn escape_regex(input: &str) -> String { + let mut escaped = String::with_capacity(input.len()); + for ch in input.chars() { + if r"\.+*?()|[]{}^$-".contains(ch) { + escaped.push('\\'); + } + escaped.push(ch); + } + escaped +} + +/// Escape a value interpolated into a double-quoted PromQL string literal. +fn escape_string(input: &str) -> String { + let mut escaped = String::with_capacity(input.len()); + for ch in input.chars() { + if ch == '\\' || ch == '"' || ch == '\n' { + escaped.push('\\'); + } + escaped.push(ch); + } + escaped +} + +/// Fail when Prometheus reports a query error, which is otherwise easy to +/// mistake for an empty result set. +fn check_status(body: &Value) -> Result<()> { + match body.get("status").and_then(Value::as_str) { + Some("success") => Ok(()), + _ => { + let kind = body + .get("errorType") + .and_then(Value::as_str) + .unwrap_or("unknown"); + let message = body + .get("error") + .and_then(Value::as_str) + .unwrap_or("no error message"); + bail!("prometheus returned `{kind}`: {message}") + } + } +} + +/// Every series of a matrix, or every point of the first series of an instant +/// vector. Unlike a first-series-only parse, a regex matching several +/// workloads keeps all of them. +fn parse_series(raw: &Value, identity_label: &str) -> Vec { + result_vector(raw) + .iter() + .map(|series| { + let labels = series + .get("metric") + .and_then(Value::as_object) + .cloned() + .unwrap_or_default(); + + // A matrix carries `values`, an instant vector a single `value`. + let points: Vec<&Value> = series + .get("values") + .and_then(Value::as_array) + .map(|values| values.iter().collect()) + .unwrap_or_else(|| series.get("value").into_iter().collect()); + + let values = points.into_iter().filter_map(parse_point).collect(); + + TimeSeries { + labels, + values, + } + }) + .filter(|series| { + // Drop series with no usable samples, and require the identity + // label so an unrelated series cannot be reported as a workload. + !series.values.is_empty() + && series + .labels + .get(identity_label) + .and_then(Value::as_str) + .is_some() + }) + .collect() +} + +fn result_vector(raw: &Value) -> Vec<&Value> { + raw.get("data") + .and_then(|v| v.get("result")) + .and_then(Value::as_array) + .map(|result| result.iter().collect()) + .unwrap_or_default() +} + +fn parse_point(point: &Value) -> Option { + let pair = point.as_array()?; + if pair.len() < 2 { + return None; + } + let timestamp = match &pair[0] { + Value::Number(n) => n.as_i64()?, + Value::String(s) => s.parse().ok()?, + _ => return None, + }; + let value = match &pair[1] { + Value::Number(n) => n.as_f64()?, + Value::String(s) => s.parse().ok()?, + _ => return None, + }; + Some(TimeSeriesPoint { timestamp, value }) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn counter_queries_use_rate() { + assert_eq!( + build_promql( + EVENTS_TOTAL, + MetricKind::Counter, + "container_name", + MatchMode::Substring, + "grafana", + "5m", + ), + r#"sum by(container_name) (rate(cortexbrain_events_total{container_name=~".*grafana.*"}[5m]))"# ); - let res = serde_json::to_string_pretty(&self.query(&promql).await?)?; - Ok(res) } - pub async fn query_get_events(&self, container_name: &str, timeframe: &str) -> Result { - //query: sum by(container_name) (rate(cortexbrain_events_total{container_name="..."}[1m])) - let promql = format!( - r#"sum by(container_name) (rate(cortexbrain_events_total{{container_name=~".*{container_name}.*"}}[{timeframe}]))"# + + #[test] + fn gauge_queries_do_not_use_rate() { + // rate() is undefined for gauges, this is the regression guard. + let query = build_promql( + "cortexbrain_enter_mem_alloc", + MetricKind::Gauge, + "k8s_pod_name", + MatchMode::Substring, + "grafana", + "5m", + ); + assert!(!query.contains("rate("), "{query}"); + assert_eq!( + query, + r#"sum by(k8s_pod_name) (avg_over_time(cortexbrain_enter_mem_alloc{k8s_pod_name=~".*grafana.*"}[5m]))"# ); - let res = serde_json::to_string_pretty(&self.query(&promql).await?)?; - Ok(res) } - pub async fn query_get_l4_events(&self, container_name: &str, timeframe: &str) -> Result { - //query: sum by(container_name) (rate(cortexbrain_socket_events_total{container_name="..."}[10m])) - let promql = format!( - r#"sum by(container_name) (rate(cortexbrain_socket_events_total{{container_name=~".*{container_name}.*"}}[{timeframe}]))"# + + #[test] + fn exact_match_is_not_a_regex() { + assert_eq!( + build_promql( + EVENTS_TOTAL, + MetricKind::Counter, + "container_name", + MatchMode::Exact, + "grafana", + "1m", + ), + r#"sum by(container_name) (rate(cortexbrain_events_total{container_name="grafana"}[1m]))"# ); - let res = serde_json::to_string_pretty(&self.query(&promql).await?)?; - Ok(res) } - pub async fn query_get_ssl_write_events( - &self, - container_name: &str, - timeframe: &str, - ) -> Result { - //query: sum by(container_name) (rate(cortexbrain_ssl_write_bytes{container_name="..."}[10m])) - let promql = format!( - r#"sum by(container_name) (rate(cortexbrain_ssl_write_bytes{{container_name=~".*{container_name}.*"}}[{timeframe}]))"# + + #[test] + fn match_mode_parsing_defaults_to_substring() { + assert_eq!(MatchMode::parse(None).unwrap(), MatchMode::Substring); + assert_eq!(MatchMode::parse(Some("")).unwrap(), MatchMode::Substring); + assert_eq!( + MatchMode::parse(Some("exact")).unwrap(), + MatchMode::Exact ); - let res = serde_json::to_string_pretty(&self.query(&promql).await?)?; - Ok(res) + assert!(MatchMode::parse(Some("fuzzy")).is_err()); } - pub async fn query_get_ssl_read_events( - &self, - container_name: &str, - timeframe: &str, - ) -> Result { - //query: sum by(container_name) (rate(cortexbrain_ssl_read_bytes{container_name="..."}[10m])) - let promql = format!( - r#"sum by(container_name) (rate(cortexbrain_ssl_read_bytes{{container_name=~".*{container_name}.*"}}[{timeframe}]))"# + + #[test] + fn workload_is_escaped() { + // A workload name must not be able to break out of the regex literal. + let query = build_promql( + EVENTS_TOTAL, + MetricKind::Counter, + "container_name", + MatchMode::Substring, + "a.b", + "1m", ); - let res = serde_json::to_string_pretty(&self.query(&promql).await?)?; - Ok(res) + assert!(query.contains(r#".*a\.b.*"#), "{query}"); + } + + #[test] + fn every_series_is_kept() { + let raw = json!({ + "status": "success", + "data": { "resultType": "matrix", "result": [ + { "metric": { "container_name": "grafana-1" }, + "values": [[1, "0.5"], [2, "0.7"]] }, + { "metric": { "container_name": "grafana-2" }, + "values": [[1, "0.1"]] }, + ]} + }); + + let series = parse_series(&raw, "container_name"); + assert_eq!(series.len(), 2); + assert_eq!(series[0].labels["container_name"], "grafana-1"); + assert_eq!(series[0].values.len(), 2); + assert_eq!(series[0].values[0].value, 0.5); + assert_eq!(series[1].values[0].timestamp, 1); + } + + #[test] + fn series_without_samples_or_identity_are_dropped() { + let raw = json!({ + "status": "success", + "data": { "resultType": "matrix", "result": [ + { "metric": { "container_name": "grafana" }, "values": [] }, + { "metric": { "command": "node" }, "values": [[1, "1"]] }, + { "metric": { "container_name": "otel" }, "values": [[1, "2"]] }, + ]} + }); + + let series = parse_series(&raw, "container_name"); + assert_eq!(series.len(), 1); + assert_eq!(series[0].labels["container_name"], "otel"); + } + + #[test] + fn instant_vector_is_parsed() { + let raw = json!({ + "status": "success", + "data": { "resultType": "vector", "result": [ + { "metric": { "container_name": "grafana" }, "value": [7, "3.5"] }, + ]} + }); + + let series = parse_series(&raw, "container_name"); + assert_eq!(series.len(), 1); + assert_eq!(series[0].values[0].timestamp, 7); + assert_eq!(series[0].values[0].value, 3.5); + } + + #[test] + fn non_numeric_samples_are_skipped() { + let raw = json!({ + "status": "success", + "data": { "resultType": "matrix", "result": [ + { "metric": { "container_name": "grafana" }, + "values": [[1, "None"], [2, "1.25"]] }, + ]} + }); + let series = parse_series(&raw, "container_name"); + assert_eq!(series[0].values.len(), 1); + assert_eq!(series[0].values[0].value, 1.25); + } + + #[test] + fn prometheus_errors_are_surfaced() { + let body = json!({ "status": "error", "errorType": "bad_data", "error": "parse error" }); + let err = check_status(&body).unwrap_err().to_string(); + assert!(err.contains("bad_data"), "{err}"); + assert!(err.contains("parse error"), "{err}"); } } diff --git a/mcp/src/tools.rs b/mcp/src/tools.rs index dde642b2..97942da0 100644 --- a/mcp/src/tools.rs +++ b/mcp/src/tools.rs @@ -1,5 +1,7 @@ -use crate::prometheus::PromClient; -use anyhow::Result; +use crate::prometheus::{ + MatchMode, MetricKind, PromClient, TimeSeriesResponse, EVENTS_TOTAL, +}; +use anyhow::Result as AnyResult; use rmcp::handler::server::ServerHandler; use rmcp::handler::server::router::tool::ToolRouter; use rmcp::handler::server::wrapper::Parameters; @@ -8,10 +10,60 @@ use rmcp::{tool, tool_handler, tool_router}; use schemars::JsonSchema; use serde::Deserialize; +/// Outcome of a tool call. An `Err` is reported to the assistant as a tool +/// error carrying the real cause, instead of the previous `Err(())` that made +/// a broken Prometheus indistinguishable from a workload with no activity. +type ToolResult = Result; + +/// Render an error with its context chain, e.g. +/// `GET /api/v1/label/... failed: connection refused`. +fn describe(err: anyhow::Error) -> String { + format!("{err:#}") +} + +/// Serialisation of an already built response should not fail; report it +/// plainly instead of pretending the query failed. +fn describe_serde(err: serde_json::Error) -> String { + err.to_string() +} + +/// Request contract shared by every metric tool. `window` is the range vector +/// of the aggregation and defaults to `step`; keeping it separate lets a query +/// return fine-grained samples over a stable rate window. #[derive(Deserialize, JsonSchema)] struct Params { - container_name: String, // example "grafana/grafana:13.1.0" - timeframe: String, //example "10m" + /// Workload to report on, matched against the identity label + /// (`container_name`, or `k8s_pod_name` on collectors that have not been + /// normalised yet). Use `list_workloads` to discover valid values. + container_name: String, + /// Range start, RFC 3339 timestamp or unix seconds, e.g. + /// `2026-09-27T00:00:00Z`. + start: String, + /// Range end, same format as `start`. + end: String, + /// Query resolution, e.g. `1m`. Also the default aggregation window. + step: String, + /// Aggregation window, e.g. `5m`. Defaults to `step`. + #[serde(default)] + window: Option, + /// `substring` (default) matches `container_name=~".*.*"`, so + /// `otel-agent` also matches `otel-agent-qsm7c`. `exact` requires a full + /// match. + #[serde(default)] + match_mode: Option, +} + +impl Params { + fn window(&self) -> &str { + self.window.as_deref().unwrap_or(&self.step) + } +} + +#[derive(Deserialize, JsonSchema, Default)] +struct ListParams { + /// How far back to look for active workloads, e.g. `1h`. Defaults to `1h`. + #[serde(default)] + lookback: Option, } #[derive(Clone)] @@ -21,7 +73,7 @@ pub struct PrometheusTool { } impl PrometheusTool { - pub fn new() -> Result { + pub fn new() -> AnyResult { Ok(Self { prometheus: PromClient::new()?, tool_router: Self::tool_router(), @@ -29,96 +81,135 @@ impl PrometheusTool { } } +impl PrometheusTool { + /// Resolve the identity label, run the aggregation and render it. + async fn run(&self, params: &Params, metric: &str, kind: MetricKind) -> ToolResult { + let label = self + .prometheus + .identity_label() + .await + .map_err(describe)? + .to_string(); + let mode = MatchMode::parse(params.match_mode.as_deref()).map_err(describe)?; + let response: TimeSeriesResponse = self + .prometheus + .windowed_query( + metric, + kind, + &label, + mode, + ¶ms.container_name, + params.window(), + ¶ms.start, + ¶ms.end, + ¶ms.step, + ) + .await + .map_err(describe)?; + serde_json::to_string_pretty(&response).map_err(describe_serde) + } +} + #[tool_router] impl PrometheusTool { - #[tool(name = "get_cpu_bytes", description = "CPU bytes allocation per event")] + #[tool( + name = "get_cpu_bytes", + description = "CPU bytes allocation per event, averaged over the window" + )] pub async fn get_cpu_bytes( &self, Parameters(params): Parameters, - ) -> Result { - let res = self - .prometheus - .query_get_cpu_bytes(¶ms.container_name, ¶ms.timeframe) - .await - .expect("An error occured"); - Ok(res) + ) -> ToolResult { + self.run( + ¶ms, + "cortexbrain_cpu_bytes_alloc", + MetricKind::Gauge, + ) + .await } #[tool( name = "get_memory_allocated_bytes", - description = "Bytes requested via mmap syscalls" + description = "Bytes requested via mmap syscalls, averaged over the window" )] pub async fn get_memory_allocated_bytes( &self, Parameters(params): Parameters, - ) -> Result { - let res = self - .prometheus - .query_get_memory_allocated_bytes(¶ms.container_name, ¶ms.timeframe) - .await - .expect("An error occured"); - Ok(res) + ) -> ToolResult { + self.run( + ¶ms, + "cortexbrain_enter_mem_alloc", + MetricKind::Gauge, + ) + .await } #[tool( name = "get_events", - description = "Total number of eBPF events processed across all perf buffers" + description = "Total number of eBPF events processed across all perf buffers, as a per-second rate" )] - pub async fn get_events(&self, Parameters(params): Parameters) -> Result { - let res = self - .prometheus - .query_get_events(¶ms.container_name, ¶ms.timeframe) - .await - .expect("An error occured"); - Ok(res) + pub async fn get_events( + &self, + Parameters(params): Parameters, + ) -> ToolResult { + self.run(¶ms, EVENTS_TOTAL, MetricKind::Counter).await } #[tool( name = "get_l4_events", - description = "Total number of socket state events processed" + description = "Total number of socket state events processed, as a per-second rate" )] pub async fn get_l4_events( &self, Parameters(params): Parameters, - ) -> Result { - let res = self - .prometheus - .query_get_l4_events(¶ms.container_name, ¶ms.timeframe) - .await - .expect("An error occured"); - Ok(res) + ) -> ToolResult { + self.run( + ¶ms, + "cortexbrain_socket_events_total", + MetricKind::Counter, + ) + .await } #[tool( name = "get_ssl_write_events", - description = "Total bytes requested by the ssl_write function" + description = "Total bytes requested by the ssl_write function, averaged over the window" )] pub async fn get_ssl_write_events( &self, Parameters(params): Parameters, - ) -> Result { - let res = self - .prometheus - .query_get_ssl_write_events(¶ms.container_name, ¶ms.timeframe) + ) -> ToolResult { + self.run(¶ms, "cortexbrain_ssl_write_bytes", MetricKind::Gauge) .await - .expect("An error occured"); - Ok(res) } #[tool( name = "get_ssl_read_events", - description = "Total bytes requested by the ssl_read function" + description = "Total bytes requested by the ssl_read function, averaged over the window" )] pub async fn get_ssl_read_events( &self, Parameters(params): Parameters, - ) -> Result { - let res = self + ) -> ToolResult { + self.run(¶ms, "cortexbrain_ssl_read_bytes", MetricKind::Gauge) + .await + } + + #[tool( + name = "list_workloads", + description = "List the workloads CortexBrain observes, with the identity label in use and which of them are active. Call this first instead of guessing container names." + )] + pub async fn list_workloads( + &self, + Parameters(params): Parameters, + ) -> ToolResult { + let lookback = params.lookback.as_deref().unwrap_or("1h"); + let listing = self .prometheus - .query_get_ssl_read_events(¶ms.container_name, ¶ms.timeframe) + .list_workloads(lookback) .await - .expect("An error occured"); - Ok(res) + .map_err(describe)?; + serde_json::to_string_pretty(&listing).map_err(describe_serde) } }