SHA256
Generated
+1825
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,17 @@
|
||||
[package]
|
||||
name = "semios-build-agent"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
axum = { version = "0.8", features = ["macros"] }
|
||||
tokio = { version = "1", features = ["full"] }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }
|
||||
tracing = "0.1"
|
||||
tracing-subscriber = { version = "0.3", features = ["env-filter", "registry", "fmt", "std"] }
|
||||
toml = "0.8"
|
||||
sysinfo = "0.36"
|
||||
anyhow = "1"
|
||||
base64 = "0.22"
|
||||
@@ -0,0 +1,7 @@
|
||||
[agent]
|
||||
master_url = "http://your-primary-machine:3000"
|
||||
name = "build-worker-1"
|
||||
hostname = "192.168.1.100"
|
||||
port = 3001
|
||||
token = "change-me-to-a-secret-token"
|
||||
ports_repo_path = "/var/lib/semios-build/ports"
|
||||
@@ -0,0 +1,376 @@
|
||||
use axum::{routing::post, Json, Router};
|
||||
use base64::Engine;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::Arc;
|
||||
use sysinfo::System;
|
||||
use tokio::io::{AsyncBufReadExt, BufReader};
|
||||
use tokio::process::Command;
|
||||
use tracing::{info, error};
|
||||
|
||||
#[derive(Debug, Deserialize, Clone)]
|
||||
struct AgentConfig {
|
||||
master_url: String,
|
||||
name: String,
|
||||
hostname: String,
|
||||
port: u16,
|
||||
token: String,
|
||||
ports_repo_path: String,
|
||||
#[serde(default = "default_clean_interval")]
|
||||
clean_interval_secs: u64,
|
||||
}
|
||||
|
||||
fn default_clean_interval() -> u64 { 21600 }
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
struct BuildRequest {
|
||||
build_id: i64,
|
||||
package_name: String,
|
||||
target: String,
|
||||
atombuild: bool,
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize)]
|
||||
struct BuildResponse {
|
||||
status: String,
|
||||
message: String,
|
||||
artifacts: Vec<String>,
|
||||
}
|
||||
|
||||
fn get_system_stats() -> (f64, f64, f64, f64) {
|
||||
let mut sys = System::new_all();
|
||||
sys.refresh_all();
|
||||
let cpu = sys.global_cpu_usage();
|
||||
let mem_total = sys.total_memory() as f64 / 1024.0 / 1024.0;
|
||||
let mem_used = sys.used_memory() as f64 / 1024.0 / 1024.0;
|
||||
let mem_pct = if mem_total > 0.0 { mem_used / mem_total * 100.0 } else { 0.0 };
|
||||
|
||||
// Get disk usage for the partition containing the ports repo
|
||||
let disk_pct = std::process::Command::new("df")
|
||||
.args(["--output=pcent", "/"])
|
||||
.output()
|
||||
.ok()
|
||||
.and_then(|o| {
|
||||
let out = String::from_utf8_lossy(&o.stdout).to_string();
|
||||
out.lines().nth(1)?.trim().trim_end_matches('%').parse::<f64>().ok()
|
||||
})
|
||||
.unwrap_or(0.0);
|
||||
|
||||
(cpu as f64, mem_pct as f64, mem_total, disk_pct)
|
||||
}
|
||||
|
||||
fn get_host_arch() -> String {
|
||||
std::process::Command::new("packie")
|
||||
.args(["print", "profile.host_arch"])
|
||||
.output()
|
||||
.map(|output| {
|
||||
if output.status.success() {
|
||||
String::from_utf8_lossy(&output.stdout).trim().to_string()
|
||||
} else {
|
||||
error!("Failed to run packie print profile.host_arch");
|
||||
std::env::consts::ARCH.to_string()
|
||||
}
|
||||
})
|
||||
.unwrap_or_else(|e| {
|
||||
error!("Failed to execute packie: {}", e);
|
||||
std::env::consts::ARCH.to_string()
|
||||
})
|
||||
}
|
||||
|
||||
fn run_clean(ports_path: &str) {
|
||||
info!("Running ./x clean --all");
|
||||
let output = std::process::Command::new("./x")
|
||||
.args(["clean", "--all"])
|
||||
.current_dir(ports_path)
|
||||
.output();
|
||||
match output {
|
||||
Ok(o) => {
|
||||
if o.status.success() {
|
||||
info!("Clean succeeded");
|
||||
} else {
|
||||
error!("Clean failed: {}", String::from_utf8_lossy(&o.stderr));
|
||||
}
|
||||
}
|
||||
Err(e) => error!("Failed to execute clean: {}", e),
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_build(
|
||||
axum::extract::State(state): axum::extract::State<AgentState>,
|
||||
Json(req): Json<BuildRequest>,
|
||||
) -> Json<BuildResponse> {
|
||||
info!("Build request received: {:?}", req);
|
||||
|
||||
state.building.store(true, Ordering::SeqCst);
|
||||
|
||||
let ports_path = std::path::PathBuf::from(&state.config.ports_repo_path);
|
||||
let pkg_dir = ports_path.join("packages").join(&req.package_name);
|
||||
|
||||
if !pkg_dir.exists() {
|
||||
state.building.store(false, Ordering::SeqCst);
|
||||
return Json(BuildResponse {
|
||||
status: "error".to_string(),
|
||||
message: format!("Package directory not found: {}", pkg_dir.display()),
|
||||
artifacts: Vec::new(),
|
||||
});
|
||||
}
|
||||
|
||||
let mut cmd = std::process::Command::new("./x");
|
||||
if req.atombuild {
|
||||
cmd.arg("--atombuild");
|
||||
}
|
||||
cmd.arg("--target").arg(&req.target);
|
||||
cmd.arg("build").arg(&req.package_name);
|
||||
cmd.current_dir(&ports_path);
|
||||
cmd.env("PYTHONUNBUFFERED", "1");
|
||||
|
||||
info!("Executing: ./x {} build {}", if req.atombuild { "--atombuild" } else { "" }, req.package_name);
|
||||
|
||||
let log_url = format!("{}/api/builds/{}/log", state.config.master_url, req.build_id);
|
||||
let client = reqwest::Client::new();
|
||||
|
||||
let mut child = match Command::new("./x")
|
||||
.arg("--atombuild")
|
||||
.arg("--target").arg(&req.target)
|
||||
.arg("build").arg(&req.package_name)
|
||||
.current_dir(&ports_path)
|
||||
.env("PYTHONUNBUFFERED", "1")
|
||||
.stdout(std::process::Stdio::piped())
|
||||
.stderr(std::process::Stdio::piped())
|
||||
.spawn()
|
||||
{
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
error!("Failed to spawn build: {}", e);
|
||||
state.building.store(false, Ordering::SeqCst);
|
||||
return Json(BuildResponse {
|
||||
status: "error".to_string(),
|
||||
message: e.to_string(),
|
||||
artifacts: Vec::new(),
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
// Stream stdout and stderr concurrently via tokio tasks
|
||||
let stdout = child.stdout.take().unwrap();
|
||||
let stderr = child.stderr.take().unwrap();
|
||||
|
||||
let stdout_url = log_url.clone();
|
||||
let stderr_url = log_url.clone();
|
||||
let client_out = client.clone();
|
||||
let client_err = client.clone();
|
||||
|
||||
let t1 = tokio::spawn(async move {
|
||||
let mut lines = BufReader::new(stdout).lines();
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
let _ = client_out.post(&stdout_url)
|
||||
.json(&serde_json::json!({ "line": &line }))
|
||||
.send()
|
||||
.await;
|
||||
}
|
||||
});
|
||||
|
||||
let t2 = tokio::spawn(async move {
|
||||
let mut lines = BufReader::new(stderr).lines();
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
let _ = client_err.post(&stderr_url)
|
||||
.json(&serde_json::json!({ "line": &line }))
|
||||
.send()
|
||||
.await;
|
||||
}
|
||||
});
|
||||
|
||||
let status = child.wait().await;
|
||||
let _ = tokio::join!(t1, t2);
|
||||
|
||||
let result = match status {
|
||||
Ok(exit_status) => {
|
||||
if exit_status.success() {
|
||||
info!("Build {} succeeded", req.package_name);
|
||||
|
||||
let target_dir = pkg_dir.join("target");
|
||||
let mut artifacts = Vec::new();
|
||||
if target_dir.exists() {
|
||||
if let Ok(entries) = std::fs::read_dir(&target_dir) {
|
||||
for entry in entries.flatten() {
|
||||
let name = entry.file_name().to_string_lossy().to_string();
|
||||
if name.ends_with(".pkg") {
|
||||
artifacts.push(name);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if !artifacts.is_empty() {
|
||||
let mut files = Vec::new();
|
||||
for name in &artifacts {
|
||||
let file_path = target_dir.join(name);
|
||||
if let Ok(data) = std::fs::read(&file_path) {
|
||||
files.push(serde_json::json!({
|
||||
"name": name,
|
||||
"content": base64::engine::general_purpose::STANDARD.encode(&data),
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
let upload_url = format!("{}/api/builds/{}/artifacts", state.config.master_url, req.build_id);
|
||||
match client.post(&upload_url).json(&serde_json::json!({ "files": files })).send().await {
|
||||
Ok(resp) => {
|
||||
if resp.status().is_success() {
|
||||
info!("Uploaded {} artifacts for build {}", artifacts.len(), req.build_id);
|
||||
} else {
|
||||
error!("Failed to upload artifacts: {}", resp.status());
|
||||
}
|
||||
}
|
||||
Err(e) => error!("Failed to upload artifacts: {}", e),
|
||||
}
|
||||
}
|
||||
|
||||
Json(BuildResponse {
|
||||
status: "success".to_string(),
|
||||
message: String::new(),
|
||||
artifacts,
|
||||
})
|
||||
} else {
|
||||
error!("Build {} failed", req.package_name);
|
||||
Json(BuildResponse {
|
||||
status: "failed".to_string(),
|
||||
message: String::new(),
|
||||
artifacts: Vec::new(),
|
||||
})
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Failed to wait for build: {}", e);
|
||||
Json(BuildResponse {
|
||||
status: "error".to_string(),
|
||||
message: e.to_string(),
|
||||
artifacts: Vec::new(),
|
||||
})
|
||||
}
|
||||
};
|
||||
|
||||
state.building.store(false, Ordering::SeqCst);
|
||||
result
|
||||
}
|
||||
|
||||
async fn send_heartbeat(config: &AgentConfig) {
|
||||
let client = reqwest::Client::new();
|
||||
let (cpu, mem, mem_total, disk) = get_system_stats();
|
||||
let arch = get_host_arch();
|
||||
|
||||
let body = serde_json::json!({
|
||||
"name": config.name,
|
||||
"hostname": config.hostname,
|
||||
"port": config.port,
|
||||
"arch": arch,
|
||||
"cpu_usage": cpu,
|
||||
"memory_usage": mem,
|
||||
"memory_total_mb": mem_total,
|
||||
"disk_usage": disk,
|
||||
});
|
||||
|
||||
let url = format!("{}/api/heartbeat", config.master_url);
|
||||
match client.post(&url).json(&body).send().await {
|
||||
Ok(_) => info!("Heartbeat sent"),
|
||||
Err(e) => error!("Failed to send heartbeat: {}", e),
|
||||
}
|
||||
}
|
||||
|
||||
async fn register(config: &AgentConfig) {
|
||||
let client = reqwest::Client::new();
|
||||
let (cpu, mem, mem_total, disk) = get_system_stats();
|
||||
let arch = get_host_arch();
|
||||
|
||||
let body = serde_json::json!({
|
||||
"name": config.name,
|
||||
"hostname": config.hostname,
|
||||
"port": config.port,
|
||||
"arch": arch,
|
||||
"cpu_usage": cpu,
|
||||
"memory_usage": mem,
|
||||
"memory_total_mb": mem_total,
|
||||
"disk_usage": disk,
|
||||
});
|
||||
|
||||
let url = format!("{}/api/register", config.master_url);
|
||||
match client.post(&url).json(&body).send().await {
|
||||
Ok(_) => info!("Registered with master"),
|
||||
Err(e) => error!("Failed to register: {}", e),
|
||||
}
|
||||
}
|
||||
|
||||
struct AgentState {
|
||||
config: AgentConfig,
|
||||
building: Arc<AtomicBool>,
|
||||
}
|
||||
|
||||
impl Clone for AgentState {
|
||||
fn clone(&self) -> Self {
|
||||
Self {
|
||||
config: self.config.clone(),
|
||||
building: self.building.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> anyhow::Result<()> {
|
||||
tracing_subscriber::fmt()
|
||||
.with_env_filter(
|
||||
tracing_subscriber::EnvFilter::try_from_default_env()
|
||||
.unwrap_or_else(|_| "info".into())
|
||||
)
|
||||
.init();
|
||||
|
||||
let config_path = std::env::var("AGENT_CONFIG").unwrap_or_else(|_| "agent.toml".to_string());
|
||||
let config: AgentConfig = toml::from_str(&std::fs::read_to_string(&config_path)?)?;
|
||||
|
||||
info!("Starting build agent: {} on port {}", config.name, config.port);
|
||||
info!("Clean interval: {}s", config.clean_interval_secs);
|
||||
|
||||
register(&config).await;
|
||||
|
||||
let building = Arc::new(AtomicBool::new(false));
|
||||
|
||||
// Heartbeat loop
|
||||
let config_clone = config.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval(std::time::Duration::from_secs(10));
|
||||
loop {
|
||||
interval.tick().await;
|
||||
send_heartbeat(&config_clone).await;
|
||||
}
|
||||
});
|
||||
|
||||
// Idle clean loop
|
||||
let clean_building = building.clone();
|
||||
let clean_config = config.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval(std::time::Duration::from_secs(clean_config.clean_interval_secs));
|
||||
loop {
|
||||
interval.tick().await;
|
||||
if !clean_building.load(Ordering::SeqCst) {
|
||||
run_clean(&clean_config.ports_repo_path);
|
||||
} else {
|
||||
info!("Skipping clean: agent is building");
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
let state = AgentState {
|
||||
config: config.clone(),
|
||||
building,
|
||||
};
|
||||
|
||||
let app = Router::new()
|
||||
.route("/build", post(handle_build))
|
||||
.with_state(state);
|
||||
|
||||
let addr = format!("0.0.0.0:{}", config.port);
|
||||
info!("Agent listening on {}", addr);
|
||||
let listener = tokio::net::TcpListener::bind(&addr).await?;
|
||||
axum::serve(listener, app).await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
Reference in New Issue
Block a user