fix: route worker failure logs through tracing
This commit is contained in:
Generated
+16
@@ -5257,6 +5257,16 @@ dependencies = [
|
|||||||
"tracing-core",
|
"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]]
|
[[package]]
|
||||||
name = "tracing-subscriber"
|
name = "tracing-subscriber"
|
||||||
version = "0.3.23"
|
version = "0.3.23"
|
||||||
@@ -5267,12 +5277,15 @@ dependencies = [
|
|||||||
"nu-ansi-term",
|
"nu-ansi-term",
|
||||||
"once_cell",
|
"once_cell",
|
||||||
"regex-automata",
|
"regex-automata",
|
||||||
|
"serde",
|
||||||
|
"serde_json",
|
||||||
"sharded-slab",
|
"sharded-slab",
|
||||||
"smallvec",
|
"smallvec",
|
||||||
"thread_local",
|
"thread_local",
|
||||||
"tracing",
|
"tracing",
|
||||||
"tracing-core",
|
"tracing-core",
|
||||||
"tracing-log",
|
"tracing-log",
|
||||||
|
"tracing-serde",
|
||||||
]
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
@@ -6695,6 +6708,8 @@ dependencies = [
|
|||||||
"tokio-tungstenite 0.29.0",
|
"tokio-tungstenite 0.29.0",
|
||||||
"toml",
|
"toml",
|
||||||
"tower",
|
"tower",
|
||||||
|
"tracing",
|
||||||
|
"tracing-subscriber",
|
||||||
"url",
|
"url",
|
||||||
"uuid",
|
"uuid",
|
||||||
"workdir",
|
"workdir",
|
||||||
@@ -6840,6 +6855,7 @@ dependencies = [
|
|||||||
"toml",
|
"toml",
|
||||||
"tower",
|
"tower",
|
||||||
"tracing",
|
"tracing",
|
||||||
|
"tracing-subscriber",
|
||||||
"ts-rs",
|
"ts-rs",
|
||||||
"url",
|
"url",
|
||||||
"uuid",
|
"uuid",
|
||||||
|
|||||||
@@ -132,6 +132,7 @@ tokio-tungstenite = "0.29"
|
|||||||
tower = "0.5"
|
tower = "0.5"
|
||||||
toml = "1.1"
|
toml = "1.1"
|
||||||
tracing = "0.1"
|
tracing = "0.1"
|
||||||
|
tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] }
|
||||||
url = "2.5"
|
url = "2.5"
|
||||||
uuid = "1.23"
|
uuid = "1.23"
|
||||||
zeroize = "1"
|
zeroize = "1"
|
||||||
|
|||||||
@@ -40,6 +40,8 @@ ring.workspace = true
|
|||||||
tar.workspace = true
|
tar.workspace = true
|
||||||
thiserror = { workspace = true }
|
thiserror = { workspace = true }
|
||||||
tokio = { workspace = true, features = ["net", "rt", "sync", "time"] }
|
tokio = { workspace = true, features = ["net", "rt", "sync", "time"] }
|
||||||
|
tracing.workspace = true
|
||||||
|
tracing-subscriber.workspace = true
|
||||||
toml.workspace = true
|
toml.workspace = true
|
||||||
url.workspace = true
|
url.workspace = true
|
||||||
uuid = { workspace = true, features = ["v7"] }
|
uuid = { workspace = true, features = ["v7"] }
|
||||||
|
|||||||
@@ -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> {
|
fn run() -> Result<(), ProcessError> {
|
||||||
let args = env::args().skip(1).collect::<Vec<_>>();
|
let args = env::args().skip(1).collect::<Vec<_>>();
|
||||||
if matches!(args.first().map(String::as_str), Some("migrate")) {
|
if matches!(args.first().map(String::as_str), Some("migrate")) {
|
||||||
@@ -68,6 +80,7 @@ fn run() -> Result<(), ProcessError> {
|
|||||||
println!("{}", usage());
|
println!("{}", usage());
|
||||||
return Ok(());
|
return Ok(());
|
||||||
};
|
};
|
||||||
|
init_serve_tracing();
|
||||||
config.http.auth = load_runtime_http_auth(&config)?;
|
config.http.auth = load_runtime_http_auth(&config)?;
|
||||||
|
|
||||||
let runtime = tokio::runtime::Builder::new_current_thread()
|
let runtime = tokio::runtime::Builder::new_current_thread()
|
||||||
|
|||||||
@@ -47,7 +47,6 @@ use protocol::{Event, Method};
|
|||||||
use std::collections::BTreeMap;
|
use std::collections::BTreeMap;
|
||||||
#[cfg(feature = "ws-server")]
|
#[cfg(feature = "ws-server")]
|
||||||
use std::collections::VecDeque;
|
use std::collections::VecDeque;
|
||||||
use std::io::Write as _;
|
|
||||||
use std::sync::atomic::{AtomicBool, Ordering};
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
use std::sync::{Arc, Mutex, MutexGuard, Weak};
|
use std::sync::{Arc, Mutex, MutexGuard, Weak};
|
||||||
#[cfg(feature = "ws-server")]
|
#[cfg(feature = "ws-server")]
|
||||||
@@ -3377,12 +3376,10 @@ fn validate_create_worker_request(request: &CreateWorkerRequest) -> Result<(), R
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn runtime_worker_create_failure_log_line(
|
fn runtime_worker_create_failure_fields(
|
||||||
worker_id: WorkerId,
|
|
||||||
workspace_id: Option<&str>,
|
|
||||||
error: &RuntimeError,
|
error: &RuntimeError,
|
||||||
) -> String {
|
) -> (&'static str, Option<String>, Option<String>) {
|
||||||
let (error_kind, operation, outcome) = match error {
|
match error {
|
||||||
RuntimeError::RuntimeStopped => ("runtime_stopped", None, None),
|
RuntimeError::RuntimeStopped => ("runtime_stopped", None, None),
|
||||||
RuntimeError::InvalidInitialInputKind { .. } => ("invalid_initial_input_kind", None, None),
|
RuntimeError::InvalidInitialInputKind { .. } => ("invalid_initial_input_kind", None, None),
|
||||||
RuntimeError::WorkerNotFound { .. } => ("worker_not_found", 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::StoreMissing { .. } => ("store_missing", None, None),
|
||||||
RuntimeError::StoreCorrupt { .. } => ("store_corrupt", None, None),
|
RuntimeError::StoreCorrupt { .. } => ("store_corrupt", None, None),
|
||||||
RuntimeError::StatePoisoned => ("state_poisoned", 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(
|
fn write_runtime_worker_create_failure(
|
||||||
@@ -3434,10 +3420,18 @@ fn write_runtime_worker_create_failure(
|
|||||||
workspace_id: Option<&str>,
|
workspace_id: Option<&str>,
|
||||||
error: &RuntimeError,
|
error: &RuntimeError,
|
||||||
) {
|
) {
|
||||||
let line = runtime_worker_create_failure_log_line(worker_id, workspace_id, error);
|
let (error_kind, operation, outcome) = runtime_worker_create_failure_fields(error);
|
||||||
let mut stdout = std::io::stdout().lock();
|
tracing::error!(
|
||||||
let _ = writeln!(stdout, "{line}");
|
target: "yoi::worker_create",
|
||||||
let _ = stdout.flush();
|
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(
|
fn validate_create_workspace_scope(
|
||||||
@@ -3564,21 +3558,13 @@ mod tests {
|
|||||||
use std::sync::{Arc, Mutex};
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn worker_create_failure_log_is_structured_for_stdout() {
|
fn worker_create_failure_fields_exclude_raw_error_messages() {
|
||||||
let worker_id = WorkerId::now_v7();
|
let (error_kind, operation, outcome) = runtime_worker_create_failure_fields(
|
||||||
let line = runtime_worker_create_failure_log_line(
|
&RuntimeError::InvalidRequest("private-token".to_string()),
|
||||||
worker_id,
|
|
||||||
Some("workspace-a"),
|
|
||||||
&RuntimeError::InvalidRequest("rejected create".to_string()),
|
|
||||||
);
|
);
|
||||||
let event: serde_json::Value = serde_json::from_str(&line).unwrap();
|
assert_eq!(error_kind, "invalid_request");
|
||||||
assert_eq!(event["level"], "ERROR");
|
assert_eq!(operation, None);
|
||||||
assert_eq!(event["event"], "worker_create_failed");
|
assert_eq!(outcome, None);
|
||||||
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());
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn test_command() -> protocol::WorkerCommandEnvelope {
|
fn test_command() -> protocol::WorkerCommandEnvelope {
|
||||||
|
|||||||
@@ -45,6 +45,7 @@ workdir = { workspace = true, features = ["http-client"] }
|
|||||||
worker-runtime.workspace = true
|
worker-runtime.workspace = true
|
||||||
toml.workspace = true
|
toml.workspace = true
|
||||||
tracing.workspace = true
|
tracing.workspace = true
|
||||||
|
tracing-subscriber.workspace = true
|
||||||
ts-rs = { version = "12.0.1", optional = true }
|
ts-rs = { version = "12.0.1", optional = true }
|
||||||
url.workspace = true
|
url.workspace = true
|
||||||
uuid = { workspace = true, features = ["v7"] }
|
uuid = { workspace = true, features = ["v7"] }
|
||||||
|
|||||||
@@ -553,7 +553,20 @@ fn run_migrate(options: MigrateOptions) -> Result<(), Box<dyn std::error::Error>
|
|||||||
Ok(())
|
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<dyn std::error::Error>> {
|
async fn run_serve(options: ServeOptions) -> Result<(), Box<dyn std::error::Error>> {
|
||||||
|
init_serve_tracing();
|
||||||
let database_path = ServerConfig::default_server_database_path();
|
let database_path = ServerConfig::default_server_database_path();
|
||||||
if let Some(parent) = database_path.parent() {
|
if let Some(parent) = database_path.parent() {
|
||||||
tokio::fs::create_dir_all(parent).await?;
|
tokio::fs::create_dir_all(parent).await?;
|
||||||
|
|||||||
@@ -1,5 +1,4 @@
|
|||||||
use std::collections::{BTreeMap, HashMap, HashSet};
|
use std::collections::{BTreeMap, HashMap, HashSet};
|
||||||
use std::io::Write as _;
|
|
||||||
use std::path::{Component, Path, PathBuf};
|
use std::path::{Component, Path, PathBuf};
|
||||||
use std::sync::atomic::{AtomicU64, Ordering};
|
use std::sync::atomic::{AtomicU64, Ordering};
|
||||||
use std::sync::{Arc, Mutex, RwLock, Weak};
|
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()) {
|
if uri.path().starts_with("/api/") && (status.is_client_error() || status.is_server_error()) {
|
||||||
let error = response.extensions().get::<ApiErrorLog>();
|
let error = response.extensions().get::<ApiErrorLog>();
|
||||||
eprintln!(
|
let event = api_failure_log_event(&method, &uri, status, error);
|
||||||
"{} yoi-server {}",
|
tracing::error!(
|
||||||
Utc::now().to_rfc3339_opts(SecondsFormat::Millis, true),
|
target: "yoi::api",
|
||||||
failed_api_log_json(&method, &uri, status, error)
|
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
|
response
|
||||||
}
|
}
|
||||||
|
|
||||||
fn failed_api_log_json(
|
fn api_failure_log_event<'a>(
|
||||||
method: &Method,
|
method: &'a Method,
|
||||||
uri: &Uri,
|
uri: &'a Uri,
|
||||||
status: StatusCode,
|
status: StatusCode,
|
||||||
error: Option<&ApiErrorLog>,
|
error: Option<&'a ApiErrorLog>,
|
||||||
) -> String {
|
) -> ApiFailureLogEvent<'a> {
|
||||||
let event = ApiFailureLogEvent {
|
ApiFailureLogEvent {
|
||||||
event: "api_error",
|
event: "api_error",
|
||||||
method: method.as_str(),
|
method: method.as_str(),
|
||||||
path: uri.path(),
|
path: uri.path(),
|
||||||
@@ -3523,7 +3528,17 @@ fn failed_api_log_json(
|
|||||||
kind: error.map(|error| error.kind.as_str()),
|
kind: error.map(|error| error.kind.as_str()),
|
||||||
message: error.map(|error| error.message.as_str()),
|
message: error.map(|error| error.message.as_str()),
|
||||||
diagnostics: error.map(|error| error.diagnostics.as_slice()),
|
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| {
|
serde_json::to_string(&event).unwrap_or_else(|serialization_error| {
|
||||||
format!(
|
format!(
|
||||||
"{{\"event\":\"api_error\",\"status\":{},\"log_serialization_error\":{:?}}}",
|
"{{\"event\":\"api_error\",\"status\":{},\"log_serialization_error\":{:?}}}",
|
||||||
@@ -14292,35 +14307,23 @@ fn finalize_worker_spawn_stage<T>(
|
|||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
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(
|
fn write_workspace_worker_create_failure(
|
||||||
runtime_id: &str,
|
runtime_id: &str,
|
||||||
worker_id: WorkerId,
|
worker_id: WorkerId,
|
||||||
phase: &str,
|
phase: &str,
|
||||||
diagnostics: &[RuntimeDiagnostic],
|
diagnostics: &[RuntimeDiagnostic],
|
||||||
) {
|
) {
|
||||||
let line = workspace_worker_create_failure_log_line(runtime_id, worker_id, phase, diagnostics);
|
tracing::error!(
|
||||||
let mut stdout = std::io::stdout().lock();
|
target: "yoi::worker_create",
|
||||||
let _ = writeln!(stdout, "{line}");
|
event = "worker_create_failed",
|
||||||
let _ = stdout.flush();
|
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(
|
fn api_error_with_additional_diagnostics(
|
||||||
@@ -19612,34 +19615,6 @@ mod tests {
|
|||||||
assert_eq!(sanitized, "failed to open server database");
|
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]
|
#[test]
|
||||||
fn workdir_runtime_miss_uses_exact_typed_code() {
|
fn workdir_runtime_miss_uses_exact_typed_code() {
|
||||||
let typed_not_found = [RuntimeDiagnostic {
|
let typed_not_found = [RuntimeDiagnostic {
|
||||||
|
|||||||
Reference in New Issue
Block a user