refactor: load service logs through agent

This commit is contained in:
Julien Denizot
2026-06-24 15:53:55 +02:00
parent a368bd5118
commit cebd8abe93
3 changed files with 129 additions and 116 deletions
+119 -2
View File
@@ -2,8 +2,9 @@ use std::{env, path::PathBuf, time::Duration};
use enuxia_aio_protocol::{
AgentHealthResponse, AgentRequest, AgentResponse, AgentStatus, BackupResult,
MAX_BACKUP_HISTORY_LIMIT, OperationKind, OperationParameters, OperationStatusResponse,
PROTOCOL_VERSION, ServiceStatusResponse, StartOperationRequest, StartOperationResponse,
MAX_BACKUP_HISTORY_LIMIT, MAX_SERVICE_LOG_TAIL, OperationKind, OperationParameters,
OperationStatusResponse, PROTOCOL_VERSION, ServiceLogsResponse, ServiceStatusResponse,
StartOperationRequest, StartOperationResponse,
};
use tokio::{
@@ -70,6 +71,70 @@ impl AgentClient {
}
}
pub async fn service_logs(
&self,
service: &str,
tail: u16,
) -> Result<ServiceLogsResponse, String> {
validate_service_name(service)?;
validate_service_log_tail(tail)?;
let request = AgentRequest::ServiceLogs {
protocol_version: PROTOCOL_VERSION,
service: service.to_owned(),
tail,
};
let response = match timeout(Duration::from_secs(5), self.request_inner(request)).await {
Ok(result) => result?,
Err(_) => {
return Err("Le délai de lecture des journaux \
a été dépassé."
.to_owned());
}
};
match response {
AgentResponse::ServiceLogs(response) => {
validate_protocol_version(response.protocol_version)?;
if response.service != service {
return Err(format!(
"Lagent a retourné les journaux \
du service {} au lieu de \
{service}.",
response.service,
));
}
if response.tail != tail {
return Err(format!(
"Lagent a utilisé une limite de \
{} lignes au lieu de {tail}.",
response.tail,
));
}
if response.container_name.trim().is_empty() {
return Err(format!(
"Le service {service} ne possède \
aucun nom de conteneur."
));
}
Ok(response)
}
AgentResponse::Error(error) => Err(format!("{} : {}", error.code, error.message,)),
unexpected => Err(format!(
"Réponse inattendue de lagent : \
{unexpected:?}",
)),
}
}
pub async fn service_status(&self) -> Result<ServiceStatusResponse, String> {
let response = self
.request(AgentRequest::ServiceStatus {
@@ -379,6 +444,31 @@ impl AgentClient {
}
}
fn validate_service_name(service: &str) -> Result<(), String> {
let valid = !service.is_empty()
&& service.chars().all(|character| {
character.is_ascii_alphanumeric() || matches!(character, '-' | '_' | '.')
});
if valid {
return Ok(());
}
Err(format!("Le service Compose {service:?} est invalide."))
}
fn validate_service_log_tail(tail: u16) -> Result<(), String> {
if (1..=MAX_SERVICE_LOG_TAIL).contains(&tail) {
return Ok(());
}
Err(format!(
"La limite des journaux doit être comprise \
entre 1 et {MAX_SERVICE_LOG_TAIL}, \
valeur reçue : {tail}.",
))
}
fn validate_backup_history_limit(limit: u16) -> Result<(), String> {
if (1..=MAX_BACKUP_HISTORY_LIMIT).contains(&limit) {
return Ok(());
@@ -450,4 +540,31 @@ mod tests {
assert!(error.contains("entre 1 et"));
}
#[test]
fn valid_service_log_name_is_accepted() {
validate_service_name("queue-short").expect("valid service");
}
#[test]
fn unsafe_service_log_name_is_rejected() {
assert!(validate_service_name("../backend").is_err());
}
#[test]
fn service_log_tail_boundaries_are_accepted() {
validate_service_log_tail(1).expect("minimum accepted");
validate_service_log_tail(MAX_SERVICE_LOG_TAIL).expect("maximum accepted");
}
#[test]
fn zero_service_log_tail_is_rejected() {
assert!(validate_service_log_tail(0).is_err());
}
#[test]
fn excessive_service_log_tail_is_rejected() {
assert!(validate_service_log_tail(MAX_SERVICE_LOG_TAIL + 1,).is_err());
}
}
+9 -113
View File
@@ -16,7 +16,7 @@ use backup_history::BackupHistoryController;
use enuxia_aio_protocol::{ServiceState, ServiceStatus};
use serde::Deserialize;
use std::{collections::HashMap, env, net::SocketAddr};
use tokio::{net::TcpListener, process::Command};
use tokio::net::TcpListener;
use verification::VerificationController;
const EXPECTED_SERVICES: [(&str, &str, &str); 9] = [
@@ -33,7 +33,6 @@ const EXPECTED_SERVICES: [(&str, &str, &str); 9] = [
#[derive(Clone)]
struct AppState {
docker_host: String,
compose_project: String,
instance_name: String,
site_name: String,
@@ -46,12 +45,6 @@ struct AppState {
verification: VerificationController,
}
#[derive(Debug, Deserialize)]
struct DockerContainer {
#[serde(rename = "Names", default)]
names: String,
}
struct ServiceView {
label: String,
description: String,
@@ -68,7 +61,6 @@ struct DashboardTemplate {
frappe_version: String,
build_version: String,
target_name: String,
docker_host: String,
agent_connection_label: String,
agent_status_label: String,
@@ -118,57 +110,6 @@ struct DashboardTemplate {
docker_error: String,
}
async fn load_docker_containers(state: &AppState) -> Result<Vec<DockerContainer>, String> {
let project_filter = format!("label=com.docker.compose.project={}", state.compose_project);
let output = Command::new("docker")
.arg("--host")
.arg(&state.docker_host)
.arg("ps")
.arg("--all")
.arg("--filter")
.arg(project_filter)
.arg("--format")
.arg("{{json .}}")
.output()
.await
.map_err(|error| format!("Impossible de lancer Docker : {error}"))?;
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr);
return Err(format!("Docker a retourné une erreur : {}", stderr.trim()));
}
let stdout = String::from_utf8_lossy(&output.stdout);
let mut containers = Vec::new();
for line in stdout.lines().filter(|line| !line.trim().is_empty()) {
let container: DockerContainer = serde_json::from_str(line)
.map_err(|error| format!("Réponse Docker JSON invalide : {error}"))?;
containers.push(container);
}
Ok(containers)
}
fn compose_service_name(container_name: &str, project_name: &str) -> String {
let prefix = format!("{project_name}-");
let without_project = container_name
.strip_prefix(&prefix)
.unwrap_or(container_name);
match without_project.rsplit_once('-') {
Some((service, replica)) if replica.chars().all(|character| character.is_ascii_digit()) => {
service.to_owned()
}
_ => without_project.to_owned(),
}
}
fn build_service_view(
label: &str,
description: &str,
@@ -280,8 +221,6 @@ async fn dashboard(State(state): State<AppState>) -> Html<String> {
frappe_version: state.frappe_version.clone(),
build_version: state.build_version.clone(),
target_name: state.target_name.clone(),
docker_host: state.docker_host.clone(),
agent_connection_label: agent_view.connection_label,
agent_status_label: agent_view.status_label,
@@ -396,53 +335,20 @@ fn log_service_label(service_name: &str) -> Option<&'static str> {
}
async fn load_container_logs(state: &AppState, service_name: &str) -> Result<String, String> {
let containers = load_docker_containers(state).await?;
let response = state.agent.service_logs(service_name, 200).await?;
let container = containers
.iter()
.find(|container| {
compose_service_name(&container.names, &state.compose_project) == service_name
})
.ok_or_else(|| format!("Aucun conteneur trouvé pour le service {service_name}."))?;
let output = Command::new("docker")
.arg("--host")
.arg(&state.docker_host)
.arg("logs")
.arg("--timestamps")
.arg("--tail")
.arg("200")
.arg(&container.names)
.output()
.await
.map_err(|error| format!("Impossible de lancer la commande Docker logs : {error}"))?;
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr);
return Err(format!(
"Docker logs a retourné une erreur : {}",
stderr.trim()
));
}
let stdout = String::from_utf8_lossy(&output.stdout);
let stderr = String::from_utf8_lossy(&output.stderr);
let mut logs = stdout.into_owned();
if !stderr.trim().is_empty() {
if !logs.is_empty() {
logs.push('\n');
}
logs.push_str(&stderr);
}
let mut logs = response.logs;
if logs.trim().is_empty() {
logs = "Aucune entrée récente pour ce service.".to_owned();
}
if response.truncated {
logs = format!(
"[Journal tronqué : seules les entrées les plus récentes sont affichées]\n\n{logs}"
);
}
Ok(logs)
}
@@ -530,13 +436,6 @@ async fn health() -> impl IntoResponse {
"ok"
}
fn default_docker_host() -> String {
let runtime_directory =
env::var("XDG_RUNTIME_DIR").unwrap_or_else(|_| "/run/user/1000".to_owned());
format!("unix://{runtime_directory}/enuxia-docker.sock")
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let agent = AgentClient::from_env();
@@ -546,8 +445,6 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
let verification = VerificationController::new(agent.clone());
let state = AppState {
docker_host: env::var("ENUXIA_DOCKER_HOST").unwrap_or_else(|_| default_docker_host()),
compose_project: env::var("ENUXIA_COMPOSE_PROJECT")
.unwrap_or_else(|_| "enuxia-aio-lab".to_owned()),
@@ -572,7 +469,6 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
};
println!("Enuxia Frappe AIO");
println!("Docker cible : {}", state.docker_host);
let app = Router::new()
.route("/", get(dashboard))
+1 -1
View File
@@ -1021,7 +1021,7 @@
État réel des conteneurs constituant
cette stack Frappe.
<br>
{{ docker_host }}
Agent local sécurisé
</div>
{% endif %}
</div>