use std::collections::HashMap;
use std::io::{BufRead, Write};
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use serde_json::{json, Value};
use crate::cache::{hash_bytes, DiskCache};
use crate::cli::args::OutputFormat;
use crate::cli::check_rulesets;
use crate::cli::pipeline::{
aggregate_results, attach_source_spans, process_schemas, render_output,
};
use crate::normalize::normalize;
use crate::profile::load;
use crate::rules::{Diagnostic, DiagnosticSeverity, RuleSet};
const MAX_PAYLOAD_BYTES: usize = 10_000_000;
const MAX_CHECK_SECONDS: u64 = 30;
const MAX_SCHEMA_BYTES: usize = 5 * 1024 * 1024; const MAX_SCHEMA_NODES: usize = 200_000; const MAX_SCHEMA_DEPTH: usize = 1_000;
pub fn run_server() {
let cache = Arc::new(DiskCache::new());
let profile_cache: Arc<Mutex<HashMap<String, crate::profile::Profile>>> =
Arc::new(Mutex::new(HashMap::new()));
let stdin = std::io::stdin();
let stdout = std::io::stdout();
let mut stdout_lock = stdout.lock();
for line in stdin.lock().lines() {
let line = match line {
Ok(l) => l,
Err(e) => {
eprintln!("error: failed to read line from stdin: {}", e);
break;
}
};
if line.len() > MAX_PAYLOAD_BYTES {
let error_response = json!({
"jsonrpc": "2.0",
"error": {
"code": -32600,
"message": "Request payload exceeds 10 MB limit"
},
"id": null
});
if writeln!(stdout_lock, "{}", error_response).is_err() {
break;
}
continue;
}
let request: Value = match serde_json::from_str(&line) {
Ok(v) => v,
Err(e) => {
let error_response = json!({
"jsonrpc": "2.0",
"error": {
"code": -32700,
"message": format!("Parse error: {e}")
},
"id": null
});
if writeln!(stdout_lock, "{}", error_response).is_err() {
break;
}
continue;
}
};
if request.get("jsonrpc") != Some(&json!("2.0")) {
let id = request.get("id").cloned().unwrap_or(json!(null));
let error_response = json!({
"jsonrpc": "2.0",
"error": {
"code": -32600,
"message": "Invalid JSON-RPC request: missing or incorrect jsonrpc field"
},
"id": id
});
if writeln!(stdout_lock, "{}", error_response).is_err() {
break;
}
continue;
}
let method = request.get("method").and_then(|v| v.as_str()).unwrap_or("");
let id = request.get("id").cloned().unwrap_or(json!(null));
match method {
"check" => {
let params = request.get("params").cloned().unwrap_or(json!({}));
let result = handle_check(params, &cache, &profile_cache);
let response = json!({
"jsonrpc": "2.0",
"result": result,
"id": id
});
if writeln!(stdout_lock, "{}", response).is_err() {
break;
}
}
"checkNode" => {
let params = request.get("params").cloned().unwrap_or(json!({}));
let result = handle_check_node(params, &profile_cache);
let response = json!({
"jsonrpc": "2.0",
"result": result,
"id": id
});
if writeln!(stdout_lock, "{}", response).is_err() {
break;
}
}
"checkPython" => {
let params = request.get("params").cloned().unwrap_or(json!({}));
let result = handle_check_python(params, &profile_cache);
let response = json!({
"jsonrpc": "2.0",
"result": result,
"id": id
});
if writeln!(stdout_lock, "{}", response).is_err() {
break;
}
}
"shutdown" => {
let response = json!({
"jsonrpc": "2.0",
"result": null,
"id": id
});
if writeln!(stdout_lock, "{}", response).is_err() {
break;
}
break;
}
"" => {
let error_response = json!({
"jsonrpc": "2.0",
"error": {
"code": -32600,
"message": "Invalid JSON-RPC request: missing method"
},
"id": id
});
if writeln!(stdout_lock, "{}", error_response).is_err() {
break;
}
}
_ => {
let error_response = json!({
"jsonrpc": "2.0",
"error": {
"code": -32601,
"message": format!("Method not found: {method}")
},
"id": id
});
if writeln!(stdout_lock, "{}", error_response).is_err() {
break;
}
}
}
}
}
fn handle_check(
params: Value,
cache: &Arc<DiskCache>,
profile_cache: &Arc<Mutex<HashMap<String, crate::profile::Profile>>>,
) -> Value {
let schema = match params.get("schema") {
Some(v) => v.clone(),
None => {
return json!({
"success": false,
"error": "Missing 'schema' parameter"
});
}
};
let profiles = match params.get("profiles").and_then(|v| v.as_array()) {
Some(arr) => arr.iter().filter_map(|v| v.as_str()).collect::<Vec<_>>(),
None => {
return json!({
"success": false,
"error": "Missing 'profiles' parameter"
});
}
};
let format_str = params
.get("format")
.and_then(|v| v.as_str())
.unwrap_or("json");
let format = match format_str {
"human" => OutputFormat::Human,
"json" => OutputFormat::Json,
"sarif" => OutputFormat::Sarif,
"gha" => OutputFormat::Gha,
"junit" => OutputFormat::Junit,
other => {
return json!({
"success": false,
"error": format!("Unknown format '{}'; expected one of: human, json, sarif, gha, junit", other)
});
}
};
let mut loaded_profiles = Vec::new();
{
let mut cache_guard = profile_cache.lock().unwrap_or_else(|e| e.into_inner());
for &profile_id in &profiles {
let profile = if let Some(cached) = cache_guard.get(profile_id) {
cached.clone()
} else {
let bytes = match crate::cli::resolve_builtin_profile(profile_id) {
Ok(b) => b,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to resolve profile '{profile_id}': {e}")
});
}
};
let profile = match load(&bytes) {
Ok(p) => p,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to load profile '{profile_id}': {e}")
});
}
};
cache_guard.insert(profile_id.to_string(), profile.clone());
profile
};
loaded_profiles.push(profile);
}
}
loaded_profiles.sort_by(|a, b| a.name.cmp(&b.name));
loaded_profiles.dedup_by_key(|p| p.name.clone());
let profile_rulesets: Vec<(&crate::profile::Profile, RuleSet)> = loaded_profiles
.iter()
.map(|p| (p, RuleSet::from_profile(p)))
.collect();
let profile_names: Vec<String> = loaded_profiles.iter().map(|p| p.name.clone()).collect();
let schema_bytes = serde_json::to_vec(&schema).unwrap_or_default();
if schema_bytes.len() > MAX_SCHEMA_BYTES {
return json!({
"success": false,
"error": format!(
"Schema serialized size ({} bytes) exceeds the {} byte limit",
schema_bytes.len(),
MAX_SCHEMA_BYTES
)
});
}
fn count_nodes_bounded(value: &Value, remaining: &mut usize, depth: usize) -> bool {
if depth > MAX_SCHEMA_DEPTH {
return false;
}
if *remaining == 0 {
return false;
}
*remaining -= 1;
match value {
Value::Array(arr) => arr
.iter()
.all(|v| count_nodes_bounded(v, remaining, depth + 1)),
Value::Object(map) => map
.values()
.all(|v| count_nodes_bounded(v, remaining, depth + 1)),
_ => true,
}
}
let mut budget = MAX_SCHEMA_NODES;
if !count_nodes_bounded(&schema, &mut budget, 0) {
return json!({
"success": false,
"error": format!(
"Schema exceeds complexity limits (max depth {MAX_SCHEMA_DEPTH}, \
max nodes {MAX_SCHEMA_NODES}); rejected to prevent resource exhaustion"
)
});
}
let start = Instant::now();
let hash = hash_bytes(&schema_bytes);
let normalized = match cache.get(hash, &schema_bytes) {
Some(n) => n,
None => {
let n = match normalize(schema) {
Ok(n) => n,
Err(e) => {
return json!({
"success": false,
"error": format!("Normalization failed: {e}")
});
}
};
cache.insert(hash, schema_bytes.clone(), n.clone());
n
}
};
let mut total_errors = 0usize;
let mut total_warnings = 0usize;
let diags = check_rulesets(&normalized.arena, &profile_rulesets);
for d in &diags {
match d.severity {
DiagnosticSeverity::Error => total_errors += 1,
DiagnosticSeverity::Warning => total_warnings += 1,
}
}
let all_diagnostics: Vec<(PathBuf, Vec<Diagnostic>)> = vec![(PathBuf::from("<inline>"), diags)];
if start.elapsed() > Duration::from_secs(MAX_CHECK_SECONDS) {
return json!({
"success": false,
"error": "Check execution exceeded 30 second limit"
});
}
let duration_ms = Some(start.elapsed().as_millis() as u64);
let output_text = render_output(
format,
&all_diagnostics,
total_errors,
total_warnings,
&profile_names,
duration_ms,
);
json!({
"success": true,
"output": output_text,
"total_errors": total_errors,
"total_warnings": total_warnings,
})
}
fn handle_check_node(
params: Value,
profile_cache: &Arc<Mutex<HashMap<String, crate::profile::Profile>>>,
) -> Value {
let sources = match params.get("sources").and_then(|v| v.as_array()) {
Some(arr) => arr
.iter()
.filter_map(|v| v.as_str().map(|s| s.to_string()))
.collect::<Vec<_>>(),
None => {
return json!({
"success": false,
"error": "Missing 'sources' parameter (expected array of glob strings)"
});
}
};
if sources.is_empty() {
return json!({
"success": false,
"error": "Empty 'sources' array; at least one source glob is required"
});
}
let profiles = match params.get("profiles").and_then(|v| v.as_array()) {
Some(arr) => arr.iter().filter_map(|v| v.as_str()).collect::<Vec<_>>(),
None => {
return json!({
"success": false,
"error": "Missing 'profiles' parameter (expected array of built-in profile IDs)"
});
}
};
let format_str = params
.get("format")
.and_then(|v| v.as_str())
.unwrap_or("json");
let format = match format_str {
"human" => OutputFormat::Human,
"json" => OutputFormat::Json,
"sarif" => OutputFormat::Sarif,
"gha" => OutputFormat::Gha,
"junit" => OutputFormat::Junit,
other => {
return json!({
"success": false,
"error": format!("Unknown format '{}'; expected one of: human, json, sarif, gha, junit", other)
});
}
};
let mut loaded_profiles = Vec::new();
{
let mut cache_guard = profile_cache.lock().unwrap_or_else(|e| e.into_inner());
for &profile_id in &profiles {
let profile = if let Some(cached) = cache_guard.get(profile_id) {
cached.clone()
} else {
let bytes = match crate::cli::resolve_builtin_profile(profile_id) {
Ok(b) => b,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to resolve profile '{profile_id}': {e}")
});
}
};
let profile = match load(&bytes) {
Ok(p) => p,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to load profile '{profile_id}': {e}")
});
}
};
cache_guard.insert(profile_id.to_string(), profile.clone());
profile
};
loaded_profiles.push(profile);
}
}
loaded_profiles.sort_by(|a, b| a.name.cmp(&b.name));
loaded_profiles.dedup_by_key(|p| p.name.clone());
let profile_rulesets: Vec<(&crate::profile::Profile, RuleSet)> = loaded_profiles
.iter()
.map(|p| (p, RuleSet::from_profile(p)))
.collect();
let profile_names: Vec<String> = loaded_profiles.iter().map(|p| p.name.clone()).collect();
let start = Instant::now();
let mut helper = match crate::node::NodeHelper::spawn(None) {
Ok(h) => h,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to spawn Node helper: {e}")
});
}
};
let mut discovered_models: Vec<crate::ingest::DiscoveredModel> = Vec::new();
let mut discovery_errors: Vec<String> = Vec::new();
let mut discovery_warnings: Vec<String> = Vec::new();
for source in &sources {
match helper.discover(source) {
Ok(resp) => {
for model in resp.models {
discovered_models.push(model);
}
for warning in &resp.warnings {
eprintln!(
"[checkNode] warning: discovery warning for '{}' in source '{}': {}",
warning.model, source, warning.message
);
discovery_warnings.push(format!(
"source '{}', model '{}': {}",
source, warning.model, warning.message
));
}
}
Err(e) => {
discovery_errors.push(format!("discovery failed for source '{}': {}", source, e));
}
}
}
helper.shutdown();
if discovered_models.is_empty() && !discovery_errors.is_empty() {
return json!({
"success": false,
"error": discovery_errors.join("; ")
});
}
if discovered_models.is_empty() {
let output_text = render_output(format, &[], 0, 0, &profile_names, Some(0));
let mut response = json!({
"success": true,
"output": output_text,
"total_errors": 0,
"total_warnings": 0,
});
if !discovery_warnings.is_empty() {
response["discovery_warnings"] = json!(discovery_warnings);
}
return response;
}
let schema_entries: Vec<(PathBuf, String, serde_json::Value)> = discovered_models
.iter()
.map(|m| {
(
PathBuf::from(&m.module_path),
m.name.clone(),
m.schema.clone(),
)
})
.collect();
let results = process_schemas(schema_entries, &profile_rulesets);
let results_with_spans = attach_source_spans(results, &discovered_models);
let (all_diagnostics, total_errors, total_warnings, fatal_errors) =
aggregate_results(results_with_spans);
let duration_ms = Some(start.elapsed().as_millis() as u64);
if fatal_errors > 0 || (!discovery_errors.is_empty() && discovered_models.is_empty()) {
return json!({
"success": false,
"error": format!("{} schema(s) failed normalization/checking", fatal_errors)
});
}
let output_text = render_output(
format,
&all_diagnostics,
total_errors,
total_warnings,
&profile_names,
duration_ms,
);
let mut response = json!({
"success": true,
"output": output_text,
"total_errors": total_errors,
"total_warnings": total_warnings,
});
if !discovery_errors.is_empty() {
response["discovery_errors"] = json!(discovery_errors);
}
if !discovery_warnings.is_empty() {
response["discovery_warnings"] = json!(discovery_warnings);
}
response
}
fn handle_check_python(
params: Value,
profile_cache: &Arc<Mutex<HashMap<String, crate::profile::Profile>>>,
) -> Value {
let packages = match params.get("packages").and_then(|v| v.as_array()) {
Some(arr) => arr
.iter()
.filter_map(|v| v.as_str().map(|s| s.to_string()))
.collect::<Vec<_>>(),
None => {
return json!({
"success": false,
"error": "Missing 'packages' parameter (expected array of Python package names)"
});
}
};
if packages.is_empty() {
return json!({
"success": false,
"error": "Empty 'packages' array; at least one package name is required"
});
}
let profiles = match params.get("profiles").and_then(|v| v.as_array()) {
Some(arr) => arr.iter().filter_map(|v| v.as_str()).collect::<Vec<_>>(),
None => {
return json!({
"success": false,
"error": "Missing 'profiles' parameter (expected array of built-in profile IDs)"
});
}
};
let format_str = params
.get("format")
.and_then(|v| v.as_str())
.unwrap_or("json");
let format = match format_str {
"human" => OutputFormat::Human,
"json" => OutputFormat::Json,
"sarif" => OutputFormat::Sarif,
"gha" => OutputFormat::Gha,
"junit" => OutputFormat::Junit,
other => {
return json!({
"success": false,
"error": format!("Unknown format '{}'; expected one of: human, json, sarif, gha, junit", other)
});
}
};
let mut loaded_profiles = Vec::new();
{
let mut cache_guard = profile_cache.lock().unwrap_or_else(|e| e.into_inner());
for &profile_id in &profiles {
let profile = if let Some(cached) = cache_guard.get(profile_id) {
cached.clone()
} else {
let bytes = match crate::cli::resolve_builtin_profile(profile_id) {
Ok(b) => b,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to resolve profile '{profile_id}': {e}")
});
}
};
let profile = match load(&bytes) {
Ok(p) => p,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to load profile '{profile_id}': {e}")
});
}
};
cache_guard.insert(profile_id.to_string(), profile.clone());
profile
};
loaded_profiles.push(profile);
}
}
loaded_profiles.sort_by(|a, b| a.name.cmp(&b.name));
loaded_profiles.dedup_by_key(|p| p.name.clone());
let profile_rulesets: Vec<(&crate::profile::Profile, RuleSet)> = loaded_profiles
.iter()
.map(|p| (p, RuleSet::from_profile(p)))
.collect();
let profile_names: Vec<String> = loaded_profiles.iter().map(|p| p.name.clone()).collect();
let start = Instant::now();
let mut helper = match crate::python::PythonHelper::spawn(None) {
Ok(h) => h,
Err(e) => {
return json!({
"success": false,
"error": format!("Failed to spawn Python helper: {e}")
});
}
};
let mut discovered_models: Vec<crate::ingest::DiscoveredModel> = Vec::new();
let mut discovery_errors: Vec<String> = Vec::new();
let mut discovery_warnings: Vec<String> = Vec::new();
for package in &packages {
match helper.discover(package) {
Ok(resp) => {
for model in resp.models {
discovered_models.push(model);
}
for warning in &resp.warnings {
discovery_warnings.push(format!(
"package '{}', model '{}': {}",
package, warning.model, warning.message
));
}
}
Err(e) => {
discovery_errors.push(format!("discovery failed for package '{}': {}", package, e));
}
}
}
helper.shutdown();
if discovered_models.is_empty() && !discovery_errors.is_empty() {
return json!({
"success": false,
"error": discovery_errors.join("; ")
});
}
if discovered_models.is_empty() {
let output_text = render_output(format, &[], 0, 0, &profile_names, Some(0));
let mut response = json!({
"success": true,
"output": output_text,
"total_errors": 0,
"total_warnings": 0,
});
if !discovery_warnings.is_empty() {
response["discovery_warnings"] = json!(discovery_warnings);
}
return response;
}
let schema_entries: Vec<(PathBuf, String, serde_json::Value)> = discovered_models
.iter()
.map(|m| {
(
PathBuf::from(&m.module_path),
m.name.clone(),
m.schema.clone(),
)
})
.collect();
let results = process_schemas(schema_entries, &profile_rulesets);
let results_with_spans = attach_source_spans(results, &discovered_models);
let (all_diagnostics, total_errors, total_warnings, fatal_errors) =
aggregate_results(results_with_spans);
let duration_ms = Some(start.elapsed().as_millis() as u64);
if fatal_errors > 0 || (!discovery_errors.is_empty() && discovered_models.is_empty()) {
return json!({
"success": false,
"error": format!("{} schema(s) failed normalization/checking", fatal_errors)
});
}
let output_text = render_output(
format,
&all_diagnostics,
total_errors,
total_warnings,
&profile_names,
duration_ms,
);
let mut response = json!({
"success": true,
"output": output_text,
"total_errors": total_errors,
"total_warnings": total_warnings,
});
if !discovery_errors.is_empty() {
response["discovery_errors"] = json!(discovery_errors);
}
if !discovery_warnings.is_empty() {
response["discovery_warnings"] = json!(discovery_warnings);
}
response
}