From d1585d7483ac8012af72b293f7d3f0cde9193d6c Mon Sep 17 00:00:00 2001 From: Hare Date: Mon, 7 Sep 2026 18:14:07 +0900 Subject: [PATCH] fix: route worker failure logs through tracing --- Cargo.lock | 16 ++++ Cargo.toml | 1 + crates/worker-runtime/Cargo.toml | 2 + crates/worker-runtime/src/main.rs | 13 ++++ crates/worker-runtime/src/runtime.rs | 58 ++++++--------- crates/workspace-server/Cargo.toml | 1 + crates/workspace-server/src/main.rs | 13 ++++ crates/workspace-server/src/server.rs | 101 ++++++++++---------------- 8 files changed, 106 insertions(+), 99 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8c6a126e..4640a34f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5257,6 +5257,16 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-serde" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "704b1aeb7be0d0a84fc9828cae51dab5970fee5088f83d1dd7ee6f6246fc6ff1" +dependencies = [ + "serde", + "tracing-core", +] + [[package]] name = "tracing-subscriber" version = "0.3.23" @@ -5267,12 +5277,15 @@ dependencies = [ "nu-ansi-term", "once_cell", "regex-automata", + "serde", + "serde_json", "sharded-slab", "smallvec", "thread_local", "tracing", "tracing-core", "tracing-log", + "tracing-serde", ] [[package]] @@ -6695,6 +6708,8 @@ dependencies = [ "tokio-tungstenite 0.29.0", "toml", "tower", + "tracing", + "tracing-subscriber", "url", "uuid", "workdir", @@ -6840,6 +6855,7 @@ dependencies = [ "toml", "tower", "tracing", + "tracing-subscriber", "ts-rs", "url", "uuid", diff --git a/Cargo.toml b/Cargo.toml index 7bb06887..4b922b31 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -132,6 +132,7 @@ tokio-tungstenite = "0.29" tower = "0.5" toml = "1.1" tracing = "0.1" +tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] } url = "2.5" uuid = "1.23" zeroize = "1" diff --git a/crates/worker-runtime/Cargo.toml b/crates/worker-runtime/Cargo.toml index 16a2cb34..72220dee 100644 --- a/crates/worker-runtime/Cargo.toml +++ b/crates/worker-runtime/Cargo.toml @@ -40,6 +40,8 @@ ring.workspace = true tar.workspace = true thiserror = { workspace = true } tokio = { workspace = true, features = ["net", "rt", "sync", "time"] } +tracing.workspace = true +tracing-subscriber.workspace = true toml.workspace = true url.workspace = true uuid = { workspace = true, features = ["v7"] } diff --git a/crates/worker-runtime/src/main.rs b/crates/worker-runtime/src/main.rs index bc66229f..d064f711 100644 --- a/crates/worker-runtime/src/main.rs +++ b/crates/worker-runtime/src/main.rs @@ -53,6 +53,18 @@ fn main() -> ExitCode { } } +fn init_serve_tracing() { + let filter = tracing_subscriber::EnvFilter::try_from_default_env() + .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")); + let _ = tracing_subscriber::fmt() + .with_env_filter(filter) + .with_writer(std::io::stdout) + .with_ansi(false) + .json() + .flatten_event(true) + .try_init(); +} + fn run() -> Result<(), ProcessError> { let args = env::args().skip(1).collect::>(); if matches!(args.first().map(String::as_str), Some("migrate")) { @@ -68,6 +80,7 @@ fn run() -> Result<(), ProcessError> { println!("{}", usage()); return Ok(()); }; + init_serve_tracing(); config.http.auth = load_runtime_http_auth(&config)?; let runtime = tokio::runtime::Builder::new_current_thread() diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 4cfc7af3..62be7430 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -47,7 +47,6 @@ use protocol::{Event, Method}; use std::collections::BTreeMap; #[cfg(feature = "ws-server")] use std::collections::VecDeque; -use std::io::Write as _; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, MutexGuard, Weak}; #[cfg(feature = "ws-server")] @@ -3377,12 +3376,10 @@ fn validate_create_worker_request(request: &CreateWorkerRequest) -> Result<(), R Ok(()) } -fn runtime_worker_create_failure_log_line( - worker_id: WorkerId, - workspace_id: Option<&str>, +fn runtime_worker_create_failure_fields( error: &RuntimeError, -) -> String { - let (error_kind, operation, outcome) = match error { +) -> (&'static str, Option, Option) { + match error { RuntimeError::RuntimeStopped => ("runtime_stopped", None, None), RuntimeError::InvalidInitialInputKind { .. } => ("invalid_initial_input_kind", None, None), RuntimeError::WorkerNotFound { .. } => ("worker_not_found", None, None), @@ -3415,18 +3412,7 @@ fn runtime_worker_create_failure_log_line( RuntimeError::StoreMissing { .. } => ("store_missing", None, None), RuntimeError::StoreCorrupt { .. } => ("store_corrupt", None, None), RuntimeError::StatePoisoned => ("state_poisoned", None, None), - }; - serde_json::json!({ - "level": "ERROR", - "event": "worker_create_failed", - "component": "runtime", - "workspace_id": workspace_id, - "worker_id": worker_id.to_string(), - "error_kind": error_kind, - "operation": operation, - "outcome": outcome, - }) - .to_string() + } } fn write_runtime_worker_create_failure( @@ -3434,10 +3420,18 @@ fn write_runtime_worker_create_failure( workspace_id: Option<&str>, error: &RuntimeError, ) { - let line = runtime_worker_create_failure_log_line(worker_id, workspace_id, error); - let mut stdout = std::io::stdout().lock(); - let _ = writeln!(stdout, "{line}"); - let _ = stdout.flush(); + let (error_kind, operation, outcome) = runtime_worker_create_failure_fields(error); + tracing::error!( + target: "yoi::worker_create", + event = "worker_create_failed", + component = "runtime", + workspace_id = workspace_id.unwrap_or(""), + worker_id = %worker_id, + error_kind, + operation = operation.as_deref().unwrap_or(""), + outcome = outcome.as_deref().unwrap_or(""), + "Worker creation failed" + ); } fn validate_create_workspace_scope( @@ -3564,21 +3558,13 @@ mod tests { use std::sync::{Arc, Mutex}; #[test] - fn worker_create_failure_log_is_structured_for_stdout() { - let worker_id = WorkerId::now_v7(); - let line = runtime_worker_create_failure_log_line( - worker_id, - Some("workspace-a"), - &RuntimeError::InvalidRequest("rejected create".to_string()), + fn worker_create_failure_fields_exclude_raw_error_messages() { + let (error_kind, operation, outcome) = runtime_worker_create_failure_fields( + &RuntimeError::InvalidRequest("private-token".to_string()), ); - let event: serde_json::Value = serde_json::from_str(&line).unwrap(); - assert_eq!(event["level"], "ERROR"); - assert_eq!(event["event"], "worker_create_failed"); - assert_eq!(event["component"], "runtime"); - assert_eq!(event["workspace_id"], "workspace-a"); - assert_eq!(event["worker_id"], worker_id.to_string()); - assert_eq!(event["error_kind"], "invalid_request"); - assert!(event.get("message").is_none()); + assert_eq!(error_kind, "invalid_request"); + assert_eq!(operation, None); + assert_eq!(outcome, None); } fn test_command() -> protocol::WorkerCommandEnvelope { diff --git a/crates/workspace-server/Cargo.toml b/crates/workspace-server/Cargo.toml index ad365732..2b6dfff3 100644 --- a/crates/workspace-server/Cargo.toml +++ b/crates/workspace-server/Cargo.toml @@ -45,6 +45,7 @@ workdir = { workspace = true, features = ["http-client"] } worker-runtime.workspace = true toml.workspace = true tracing.workspace = true +tracing-subscriber.workspace = true ts-rs = { version = "12.0.1", optional = true } url.workspace = true uuid = { workspace = true, features = ["v7"] } diff --git a/crates/workspace-server/src/main.rs b/crates/workspace-server/src/main.rs index 80950223..2c44f905 100644 --- a/crates/workspace-server/src/main.rs +++ b/crates/workspace-server/src/main.rs @@ -553,7 +553,20 @@ fn run_migrate(options: MigrateOptions) -> Result<(), Box Ok(()) } +fn init_serve_tracing() { + let filter = tracing_subscriber::EnvFilter::try_from_default_env() + .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")); + let _ = tracing_subscriber::fmt() + .with_env_filter(filter) + .with_writer(std::io::stdout) + .with_ansi(false) + .json() + .flatten_event(true) + .try_init(); +} + async fn run_serve(options: ServeOptions) -> Result<(), Box> { + init_serve_tracing(); let database_path = ServerConfig::default_server_database_path(); if let Some(parent) = database_path.parent() { tokio::fs::create_dir_all(parent).await?; diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 0bad4e14..8530cdc3 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -1,5 +1,4 @@ use std::collections::{BTreeMap, HashMap, HashSet}; -use std::io::Write as _; use std::path::{Component, Path, PathBuf}; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex, RwLock, Weak}; @@ -3499,23 +3498,29 @@ async fn log_failed_api_response(request: Request, next: Next) -> Response { if uri.path().starts_with("/api/") && (status.is_client_error() || status.is_server_error()) { let error = response.extensions().get::(); - eprintln!( - "{} yoi-server {}", - Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true), - failed_api_log_json(&method, &uri, status, error) + let event = api_failure_log_event(&method, &uri, status, error); + tracing::error!( + target: "yoi::api", + event = event.event, + method = event.method, + path = event.path, + status = event.status, + kind = event.kind.unwrap_or("unknown"), + message = event.message.unwrap_or(""), + diagnostics = ?event.diagnostics.unwrap_or_default(), ); } response } -fn failed_api_log_json( - method: &Method, - uri: &Uri, +fn api_failure_log_event<'a>( + method: &'a Method, + uri: &'a Uri, status: StatusCode, - error: Option<&ApiErrorLog>, -) -> String { - let event = ApiFailureLogEvent { + error: Option<&'a ApiErrorLog>, +) -> ApiFailureLogEvent<'a> { + ApiFailureLogEvent { event: "api_error", method: method.as_str(), path: uri.path(), @@ -3523,7 +3528,17 @@ fn failed_api_log_json( kind: error.map(|error| error.kind.as_str()), message: error.map(|error| error.message.as_str()), diagnostics: error.map(|error| error.diagnostics.as_slice()), - }; + } +} + +#[cfg(test)] +fn failed_api_log_json( + method: &Method, + uri: &Uri, + status: StatusCode, + error: Option<&ApiErrorLog>, +) -> String { + let event = api_failure_log_event(method, uri, status, error); serde_json::to_string(&event).unwrap_or_else(|serialization_error| { format!( "{{\"event\":\"api_error\",\"status\":{},\"log_serialization_error\":{:?}}}", @@ -14292,35 +14307,23 @@ fn finalize_worker_spawn_stage( )) } -fn workspace_worker_create_failure_log_line( - runtime_id: &str, - worker_id: WorkerId, - phase: &str, - diagnostics: &[RuntimeDiagnostic], -) -> String { - serde_json::json!({ - "level": "ERROR", - "event": "worker_create_failed", - "component": "workspace_server", - "runtime_id": runtime_id, - "worker_id": worker_id.to_string(), - "phase": phase, - "cleanup_succeeded": diagnostics.is_empty(), - "cleanup_diagnostics": diagnostics, - }) - .to_string() -} - fn write_workspace_worker_create_failure( runtime_id: &str, worker_id: WorkerId, phase: &str, diagnostics: &[RuntimeDiagnostic], ) { - let line = workspace_worker_create_failure_log_line(runtime_id, worker_id, phase, diagnostics); - let mut stdout = std::io::stdout().lock(); - let _ = writeln!(stdout, "{line}"); - let _ = stdout.flush(); + tracing::error!( + target: "yoi::worker_create", + event = "worker_create_failed", + component = "workspace_server", + runtime_id, + worker_id = %worker_id, + phase, + cleanup_succeeded = diagnostics.is_empty(), + cleanup_diagnostics = ?diagnostics, + "Worker creation failed" + ); } fn api_error_with_additional_diagnostics( @@ -19612,34 +19615,6 @@ mod tests { assert_eq!(sanitized, "failed to open server database"); } - #[test] - fn worker_create_failure_log_is_structured_for_stdout() { - let worker_id = WorkerId::now_v7(); - let diagnostics = vec![RuntimeDiagnostic { - code: "worker_cleanup_failed".to_string(), - severity: DiagnosticSeverity::Error, - message: "cleanup failed".to_string(), - }]; - let line = workspace_worker_create_failure_log_line( - "runtime-a", - worker_id, - "runtime_spawn_transport", - &diagnostics, - ); - let event: serde_json::Value = serde_json::from_str(&line).unwrap(); - assert_eq!(event["level"], "ERROR"); - assert_eq!(event["event"], "worker_create_failed"); - assert_eq!(event["component"], "workspace_server"); - assert_eq!(event["runtime_id"], "runtime-a"); - assert_eq!(event["worker_id"], worker_id.to_string()); - assert_eq!(event["phase"], "runtime_spawn_transport"); - assert_eq!(event["cleanup_succeeded"], false); - assert_eq!( - event["cleanup_diagnostics"][0]["code"], - "worker_cleanup_failed" - ); - } - #[test] fn workdir_runtime_miss_uses_exact_typed_code() { let typed_not_found = [RuntimeDiagnostic {