diff --git a/Cargo.lock b/Cargo.lock index 00b3e0b6fbd..c0205a65751 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4466,6 +4466,7 @@ dependencies = [ "blake3", "deadpool-postgres", "ironclaw_host_api", + "ironclaw_observability", "ironclaw_safety", "libsql", "serde", @@ -4657,6 +4658,7 @@ dependencies = [ "ironclaw_memory", "ironclaw_memory_native", "ironclaw_network", + "ironclaw_observability", "ironclaw_process_sandbox", "ironclaw_processes", "ironclaw_product_adapter_registry", @@ -4850,6 +4852,13 @@ dependencies = [ "urlencoding", ] +[[package]] +name = "ironclaw_observability" +version = "0.1.0" +dependencies = [ + "tracing", +] + [[package]] name = "ironclaw_outbound" version = "0.1.0" @@ -5056,6 +5065,7 @@ dependencies = [ "ironclaw_host_runtime", "ironclaw_llm", "ironclaw_loop_support", + "ironclaw_observability", "ironclaw_processes", "ironclaw_reborn_event_store", "ironclaw_resources", @@ -5146,6 +5156,7 @@ dependencies = [ "ironclaw_loop_support", "ironclaw_mcp", "ironclaw_network", + "ironclaw_observability", "ironclaw_outbound", "ironclaw_processes", "ironclaw_product_adapter_registry", @@ -5660,6 +5671,7 @@ dependencies = [ "hex", "ironclaw_filesystem", "ironclaw_host_api", + "ironclaw_observability", "serde", "serde_jcs", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index 17ad8e93410..523254bd0c0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "crates/ironclaw_filesystem", "crates/ironclaw_attachments", "crates/ironclaw_extractors", "crates/ironclaw_memory", "crates/ironclaw_memory_native", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_event_streams", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_process_sandbox", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_wasm_sandbox_core", "crates/ironclaw_wasm_limiter", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_auth", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_prompt_envelope", "crates/ironclaw_hooks", "crates/ironclaw_hooks_postgres", "crates/ironclaw_hooks_libsql", "crates/ironclaw_hooks_parity", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_reborn_identity", "crates/ironclaw_first_party_extensions", "crates/ironclaw_reborn_cli", "crates/ironclaw_reborn_traces", "crates/ironclaw_reborn_webui_ingress", "crates/ironclaw_reborn_openai_compat", "crates/ironclaw_reborn_openai_compat_storage", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_context", "crates/ironclaw_product_workflow", "crates/ironclaw_product_workflow_storage", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_slack_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_triggers", "crates/ironclaw_projects", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_oauth", "crates/ironclaw_llm", "crates/ironclaw_embeddings", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui", "crates/ironclaw_webui_v2", "crates/ironclaw_webui_v2_static", "crates/ironclaw_skill_learning", "tools/ironclaw_stress"] +members = [".", "crates/ironclaw_common", "crates/ironclaw_observability", "crates/ironclaw_host_api", "crates/ironclaw_filesystem", "crates/ironclaw_attachments", "crates/ironclaw_extractors", "crates/ironclaw_memory", "crates/ironclaw_memory_native", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_event_streams", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_process_sandbox", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_wasm_sandbox_core", "crates/ironclaw_wasm_limiter", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_auth", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_prompt_envelope", "crates/ironclaw_hooks", "crates/ironclaw_hooks_postgres", "crates/ironclaw_hooks_libsql", "crates/ironclaw_hooks_parity", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_reborn_identity", "crates/ironclaw_first_party_extensions", "crates/ironclaw_reborn_cli", "crates/ironclaw_reborn_traces", "crates/ironclaw_reborn_webui_ingress", "crates/ironclaw_reborn_openai_compat", "crates/ironclaw_reborn_openai_compat_storage", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_context", "crates/ironclaw_product_workflow", "crates/ironclaw_product_workflow_storage", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_slack_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_triggers", "crates/ironclaw_projects", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_oauth", "crates/ironclaw_llm", "crates/ironclaw_embeddings", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui", "crates/ironclaw_webui_v2", "crates/ironclaw_webui_v2_static", "crates/ironclaw_skill_learning", "tools/ironclaw_stress"] exclude = [ "channels-src/discord", "channels-src/feishu", diff --git a/crates/ironclaw_filesystem/Cargo.toml b/crates/ironclaw_filesystem/Cargo.toml index 5ec14273992..84c54988879 100644 --- a/crates/ironclaw_filesystem/Cargo.toml +++ b/crates/ironclaw_filesystem/Cargo.toml @@ -20,6 +20,7 @@ async-trait = "0.1" blake3 = "1" deadpool-postgres = { version = "0.14", optional = true } ironclaw_host_api = { path = "../ironclaw_host_api", version = "0.1.0" } +ironclaw_observability = { path = "../ironclaw_observability" } ironclaw_safety = { path = "../ironclaw_safety", version = "0.2.2" } libsql = { version = "0.9", optional = true, default-features = false, features = ["core", "replication", "remote", "tls"] } serde = { version = "1", features = ["derive"] } diff --git a/crates/ironclaw_filesystem/src/scoped.rs b/crates/ironclaw_filesystem/src/scoped.rs index 8f03b3a1f4a..7450a4815da 100644 --- a/crates/ironclaw_filesystem/src/scoped.rs +++ b/crates/ironclaw_filesystem/src/scoped.rs @@ -1,8 +1,9 @@ -use std::sync::Arc; +use std::{sync::Arc, time::Instant}; use ironclaw_host_api::{ HostApiError, MountPermissions, MountView, ResourceScope, ScopedPath, VirtualPath, }; +use ironclaw_observability::live_latency_started_at; use crate::backend::{EventRecord, StorageTxn}; use crate::{ @@ -50,6 +51,89 @@ impl std::fmt::Debug for ScopedFilesystem { } } +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +#[serde(rename_all = "snake_case")] +enum PathClass { + Workspace, + Memory, + Artifacts, + Turns, + Other, +} + +impl PathClass { + fn as_str(self) -> &'static str { + match self { + Self::Workspace => "workspace", + Self::Memory => "memory", + Self::Artifacts => "artifacts", + Self::Turns => "turns", + Self::Other => "other", + } + } +} + +fn scoped_path_class(path: &ScopedPath) -> PathClass { + match path.as_str().split('/').nth(1) { + Some("workspace") => PathClass::Workspace, + Some("memory") => PathClass::Memory, + Some("artifacts") => PathClass::Artifacts, + Some("turns") => PathClass::Turns, + _ => PathClass::Other, + } +} + +fn filesystem_error_kind(error: &FilesystemError) -> &'static str { + match error { + FilesystemError::Contract(_) => "contract", + FilesystemError::PermissionDenied { .. } => "permission_denied", + FilesystemError::MountNotFound { .. } => "mount_not_found", + FilesystemError::NotFound { .. } => "not_found", + FilesystemError::PathOutsideMount { .. } => "path_outside_mount", + FilesystemError::SymlinkEscape { .. } => "symlink_escape", + FilesystemError::MountConflict { .. } => "mount_conflict", + FilesystemError::Backend { .. } => "backend", + FilesystemError::VersionMismatch { .. } => "version_mismatch", + FilesystemError::Unsupported { .. } => "unsupported", + FilesystemError::IndexConflict { .. } => "index_conflict", + FilesystemError::DescriptorOverclaims { .. } => "descriptor_overclaims", + FilesystemError::SerializeIndexed { .. } => "serialize_indexed", + FilesystemError::DeserializeIndexed { .. } => "deserialize_indexed", + FilesystemError::CorruptRecordVersion { .. } => "corrupt_record_version", + FilesystemError::IndexSpecMissingAfterUpsert { .. } => "index_spec_missing_after_upsert", + FilesystemError::BackendInfrastructure { .. } => "backend_infrastructure", + } +} + +fn trace_fs_latency( + operation: &'static str, + path: &ScopedPath, + started_at: Option, + result: &Result, + bytes: Option, +) { + let path_class = scoped_path_class(path); + match result { + Ok(_) => ironclaw_observability::live_latency_trace_ok!( + "filesystem", + operation, + started_at, + path_class = path_class.as_str(), + bytes = bytes.unwrap_or(0), + "filesystem operation completed", + ), + Err(error) => ironclaw_observability::live_latency_trace_error!( + "filesystem", + operation, + started_at, + filesystem_error_kind(error), + path_class = path_class.as_str(), + bytes = bytes.unwrap_or(0), + "filesystem operation failed", + ), + } +} + impl ScopedFilesystem where F: RootFilesystem + ?Sized, @@ -104,9 +188,13 @@ where entry: Entry, cas: CasExpectation, ) -> Result { + let started_at = live_latency_started_at(); + let bytes = entry.body.len(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::WriteFile)?; - self.root.put(&virtual_path, entry, cas).await + let result = self.root.put(&virtual_path, entry, cas).await; + trace_fs_latency("put", path, started_at, &result, Some(bytes)); + result } /// Read the entry at `path`, returning `None` if absent. @@ -115,9 +203,12 @@ where scope: &ResourceScope, path: &ScopedPath, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::ReadFile)?; - self.root.get(&virtual_path).await + let result = self.root.get(&virtual_path).await; + trace_fs_latency("get", path, started_at, &result, None); + result } /// Filtered query over `prefix`. @@ -128,9 +219,12 @@ where filter: &Filter, page: Page, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, prefix, FilesystemOperation::Query)?; - self.root.query(&virtual_path, filter, page).await + let result = self.root.query(&virtual_path, filter, page).await; + trace_fs_latency("query", prefix, started_at, &result, None); + result } /// Declare an index on the mount under `prefix`. @@ -140,9 +234,12 @@ where prefix: &ScopedPath, spec: &IndexSpec, ) -> Result<(), FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, prefix, FilesystemOperation::EnsureIndex)?; - self.root.ensure_index(&virtual_path, spec).await + let result = self.root.ensure_index(&virtual_path, spec).await; + trace_fs_latency("ensure_index", prefix, started_at, &result, None); + result } /// Begin a multi-key transaction (capability-gated). @@ -157,10 +254,13 @@ where scope: &ResourceScope, prefix: &ScopedPath, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let view = self.mount_view(scope)?; let virtual_path = resolve_with_permission_view(&view, prefix, FilesystemOperation::BeginTxn)?; - let inner = self.root.begin(&virtual_path).await?; + let result = self.root.begin(&virtual_path).await; + trace_fs_latency("begin", prefix, started_at, &result, None); + let inner = result?; let permissions = view.resolve_with_grant(prefix)?.1.permissions.clone(); Ok(Box::new(ScopedStorageTxn { inner, @@ -178,9 +278,13 @@ where path: &ScopedPath, payload: Vec, ) -> Result { + let started_at = live_latency_started_at(); + let bytes = payload.len(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::Append)?; - self.root.append(&virtual_path, payload).await + let result = self.root.append(&virtual_path, payload).await; + trace_fs_latency("append", path, started_at, &result, Some(bytes)); + result } /// Append multiple `payloads` to the event log at `path` in one backend @@ -193,9 +297,13 @@ where path: &ScopedPath, payloads: Vec>, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); + let bytes = payloads.iter().map(Vec::len).sum(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::Append)?; - self.root.append_batch(&virtual_path, payloads).await + let result = self.root.append_batch(&virtual_path, payloads).await; + trace_fs_latency("append_batch", path, started_at, &result, Some(bytes)); + result } /// Read events at `path` starting just after `from`. @@ -205,8 +313,11 @@ where path: &ScopedPath, from: SeqNo, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::Tail)?; - self.root.tail(&virtual_path, from).await + let result = self.root.tail(&virtual_path, from).await; + trace_fs_latency("tail", path, started_at, &result, None); + result } /// Read at most `max_records` events at `path` starting just after `from`. @@ -217,10 +328,14 @@ where from: SeqNo, max_records: usize, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::Tail)?; - self.root + let result = self + .root .tail_bounded(&virtual_path, from, max_records) - .await + .await; + trace_fs_latency("tail_bounded", path, started_at, &result, None); + result } /// Return the highest seq present at `path` with `seq > from`, or `None` @@ -232,9 +347,12 @@ where path: &ScopedPath, from: SeqNo, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::HeadSeq)?; - self.root.head_seq(&virtual_path, from).await + let result = self.root.head_seq(&virtual_path, from).await; + trace_fs_latency("head_seq", path, started_at, &result, None); + result } /// Reserve a path-local monotonic sequence number. @@ -243,9 +361,12 @@ where scope: &ResourceScope, path: &ScopedPath, ) -> Result { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::ReserveSeq)?; - self.root.reserve_sequence(&virtual_path).await + let result = self.root.reserve_sequence(&virtual_path).await; + trace_fs_latency("reserve_sequence", path, started_at, &result, None); + result } // ─── Legacy bytes-plane methods (DEPRECATED — transitional) ─────────── @@ -257,9 +378,12 @@ where scope: &ResourceScope, path: &ScopedPath, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::ReadFile)?; - self.root.read_file(&virtual_path).await + let result = self.root.read_file(&virtual_path).await; + trace_fs_latency("read_file", path, started_at, &result, None); + result } /// **DEPRECATED — use [`write_bytes`](Self::write_bytes) or @@ -270,9 +394,12 @@ where path: &ScopedPath, bytes: &[u8], ) -> Result<(), FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::WriteFile)?; - self.root.write_file(&virtual_path, bytes).await + let result = self.root.write_file(&virtual_path, bytes).await; + trace_fs_latency("write_file", path, started_at, &result, Some(bytes.len())); + result } /// Write bytes using an already-authorized mount view instead of the @@ -286,16 +413,26 @@ where path: &ScopedPath, bytes: &[u8], ) -> Result<(), FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = resolve_with_permission_view(view, path, FilesystemOperation::WriteFile)?; - self.root + let result = self + .root .put( &virtual_path, Entry::bytes(bytes.to_vec()), CasExpectation::Any, ) .await - .map(|_| ()) + .map(|_| ()); + trace_fs_latency( + "write_bytes_with_mount_view", + path, + started_at, + &result, + Some(bytes.len()), + ); + result } /// **DEPRECATED — no direct replacement on the unified surface.** Use @@ -307,9 +444,12 @@ where path: &ScopedPath, bytes: &[u8], ) -> Result<(), FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::AppendFile)?; - self.root.append_file(&virtual_path, bytes).await + let result = self.root.append_file(&virtual_path, bytes).await; + trace_fs_latency("append_file", path, started_at, &result, Some(bytes.len())); + result } pub async fn list_dir( @@ -317,9 +457,12 @@ where scope: &ResourceScope, path: &ScopedPath, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::ListDir)?; - self.root.list_dir(&virtual_path).await + let result = self.root.list_dir(&virtual_path).await; + trace_fs_latency("list_dir", path, started_at, &result, None); + result } pub async fn list_dir_bounded( @@ -328,9 +471,12 @@ where path: &ScopedPath, max_entries: usize, ) -> Result, FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::ListDir)?; - self.root.list_dir_bounded(&virtual_path, max_entries).await + let result = self.root.list_dir_bounded(&virtual_path, max_entries).await; + trace_fs_latency("list_dir_bounded", path, started_at, &result, None); + result } pub async fn stat( @@ -338,8 +484,11 @@ where scope: &ResourceScope, path: &ScopedPath, ) -> Result { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::Stat)?; - self.root.stat(&virtual_path).await + let result = self.root.stat(&virtual_path).await; + trace_fs_latency("stat", path, started_at, &result, None); + result } pub async fn delete( @@ -347,9 +496,12 @@ where scope: &ResourceScope, path: &ScopedPath, ) -> Result<(), FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::Delete)?; - self.root.delete(&virtual_path).await + let result = self.root.delete(&virtual_path).await; + trace_fs_latency("delete", path, started_at, &result, None); + result } /// **DEPRECATED — the unified entry plane infers directories from path @@ -359,9 +511,12 @@ where scope: &ResourceScope, path: &ScopedPath, ) -> Result<(), FilesystemError> { + let started_at = live_latency_started_at(); let virtual_path = self.resolve_with_permission(scope, path, FilesystemOperation::CreateDirAll)?; - self.root.create_dir_all(&virtual_path).await + let result = self.root.create_dir_all(&virtual_path).await; + trace_fs_latency("create_dir_all", path, started_at, &result, None); + result } // ─── Convenience helpers for byte-only callers ──────────────────────── diff --git a/crates/ironclaw_filesystem/src/scoped/tests.rs b/crates/ironclaw_filesystem/src/scoped/tests.rs index 1d12c4ac68a..a02ba92e59e 100644 --- a/crates/ironclaw_filesystem/src/scoped/tests.rs +++ b/crates/ironclaw_filesystem/src/scoped/tests.rs @@ -39,6 +39,23 @@ fn expect_err(result: Result) -> FilesystemError { } } +#[test] +fn scoped_path_class_buckets_known_segments_and_redacts_unknowns() { + let cases = [ + ("/workspace/project/file.txt", PathClass::Workspace), + ("/memory/profile.json", PathClass::Memory), + ("/artifacts/run/output.json", PathClass::Artifacts), + ("/turns/state.json", PathClass::Turns), + ("/users/alice/private.txt", PathClass::Other), + ("/tenants/acme/users/alice/secrets", PathClass::Other), + ]; + + for (raw, expected) in cases { + let path = ScopedPath::new(raw).unwrap(); + assert_eq!(scoped_path_class(&path), expected); + } +} + fn scoped_in_memory(permissions: MountPermissions) -> ScopedFilesystem { ScopedFilesystem::with_fixed_view( Arc::new(InMemoryBackend::new()), diff --git a/crates/ironclaw_host_runtime/Cargo.toml b/crates/ironclaw_host_runtime/Cargo.toml index 8cb618a8c5f..7c47d5f59a9 100644 --- a/crates/ironclaw_host_runtime/Cargo.toml +++ b/crates/ironclaw_host_runtime/Cargo.toml @@ -35,6 +35,7 @@ ironclaw_memory = { path = "../ironclaw_memory" } ironclaw_memory_native = { path = "../ironclaw_memory_native" } ironclaw_mcp = { path = "../ironclaw_mcp" } ironclaw_network = { path = "../ironclaw_network" } +ironclaw_observability = { path = "../ironclaw_observability" } ironclaw_processes = { path = "../ironclaw_processes" } ironclaw_product_adapter_registry = { path = "../ironclaw_product_adapter_registry" } ironclaw_prompt_envelope = { path = "../ironclaw_prompt_envelope" } diff --git a/crates/ironclaw_host_runtime/src/production.rs b/crates/ironclaw_host_runtime/src/production.rs index 6ffef8b999d..ce45b590912 100644 --- a/crates/ironclaw_host_runtime/src/production.rs +++ b/crates/ironclaw_host_runtime/src/production.rs @@ -13,7 +13,7 @@ //! claims. The default fail-closed policy denies authority until composition //! supplies a concrete host policy. -use std::sync::Arc; +use std::{sync::Arc, time::Instant}; use async_trait::async_trait; use futures_util::future::join_all; @@ -35,6 +35,7 @@ use ironclaw_host_api::{ RuntimeCredentialAuthRequirement, RuntimeDispatchErrorKind, RuntimeKind, SecretHandle, runtime_policy::EffectiveRuntimePolicy, sha256_digest_token, }; +use ironclaw_observability::live_latency_started_at; use ironclaw_process_sandbox::{ PROCESS_SANDBOX_CAPABILITY_ID, SandboxProcessPlan, ValidatedSandboxProcessPlan, }; @@ -49,6 +50,52 @@ use ironclaw_secrets::SecretStore; use ironclaw_trust::{HostTrustPolicy, TrustDecision, TrustError, TrustPolicy, TrustProvenance}; use ironclaw_turns::run_profile::LoopSafeSummary; +fn trace_capability_latency_ok( + operation: &'static str, + capability_id: &CapabilityId, + scope: &ResourceScope, + started_at: Option, +) { + ironclaw_observability::live_latency_trace_ok!( + "host_runtime", + operation, + started_at, + capability_id = %capability_id, + tenant_id = %scope.tenant_id, + user_id = %scope.user_id, + agent_id = scope.agent_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + project_id = scope.project_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + mission_id = scope.mission_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + thread_id = scope.thread_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + invocation_id = %scope.invocation_id, + "host runtime capability operation completed", + ); +} + +fn trace_capability_latency_error( + operation: &'static str, + capability_id: &CapabilityId, + scope: &ResourceScope, + started_at: Option, + _error: &E, +) { + ironclaw_observability::live_latency_trace_error!( + "host_runtime", + operation, + started_at, + "error", + capability_id = %capability_id, + tenant_id = %scope.tenant_id, + user_id = %scope.user_id, + agent_id = scope.agent_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + project_id = scope.project_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + mission_id = scope.mission_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + thread_id = scope.thread_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + invocation_id = %scope.invocation_id, + "host runtime capability operation failed", + ); +} + use crate::{ BuiltinObligationHandler, BuiltinObligationServices, CancelRuntimeWorkOutcome, CancelRuntimeWorkRequest, CapabilitySurfaceVersion, HostRuntime, HostRuntimeError, @@ -377,6 +424,7 @@ impl HostRuntime for DefaultHostRuntime { } = request; let scope = context.resource_scope.clone(); let invocation_id = context.invocation_id; + let total_started_at = live_latency_started_at(); // Forward the (currently advisory) idempotency key into spans for // audit/tracing only — dedupe enforcement is not yet implemented at // this layer (see `RuntimeCapabilityRequest::idempotency_key`). @@ -395,6 +443,12 @@ impl HostRuntime for DefaultHostRuntime { runtime_policy_error_kind = error.kind(), "capability runtime policy rejected invocation before dispatch" ); + trace_capability_latency_ok( + "invoke_capability_policy_rejected", + &capability_id, + &scope, + total_started_at, + ); return Ok(runtime_policy_failure(capability_id, error)); } @@ -406,6 +460,12 @@ impl HostRuntime for DefaultHostRuntime { trust_error_kind = error.kind(), "capability trust evaluation failed before dispatch" ); + trace_capability_latency_ok( + "invoke_capability_trust_rejected", + &capability_id, + &scope, + total_started_at, + ); return Ok(trust_evaluation_failure(capability_id, error)); } }; @@ -429,13 +489,33 @@ impl HostRuntime for DefaultHostRuntime { // before the authorizer and trust/authorization checks. The dispatch-time // obligation check (which runs after those checks) is the enforcing layer. // The pre-flight provides ordering only (credentials before approval gate). + let credential_preflight_started_at = live_latency_started_at(); if let Some(auth_required) = self .credential_preflight_check(&capability_id, &scope, ®istry) .await { + trace_capability_latency_ok( + "credential_preflight_check", + &capability_id, + &scope, + credential_preflight_started_at, + ); + trace_capability_latency_ok( + "invoke_capability_auth_required", + &capability_id, + &scope, + total_started_at, + ); return Ok(auth_required); } + trace_capability_latency_ok( + "credential_preflight_check", + &capability_id, + &scope, + credential_preflight_started_at, + ); + let approval_started_at = live_latency_started_at(); self.apply_persistent_approval_policy( &mut context, ®istry, @@ -445,6 +525,12 @@ impl HostRuntime for DefaultHostRuntime { &trust_decision, ) .await; + trace_capability_latency_ok( + "persistent_approval_policy", + &capability_id, + &scope, + approval_started_at, + ); let host = self.capability_host(®istry); let invocation = CapabilityInvocationRequest { @@ -455,19 +541,63 @@ impl HostRuntime for DefaultHostRuntime { trust_decision, }; + let dispatch_started_at = live_latency_started_at(); match host.invoke_json(invocation).await { - Ok(result) => Ok(RuntimeCapabilityOutcome::Completed(Box::new( - completed_outcome_from(result, capability_id), - ))), + Ok(result) => { + trace_capability_latency_ok( + "capability_host_invoke_json", + &capability_id, + &scope, + dispatch_started_at, + ); + trace_capability_latency_ok( + "invoke_capability", + &capability_id, + &scope, + total_started_at, + ); + Ok(RuntimeCapabilityOutcome::Completed(Box::new( + completed_outcome_from(result, capability_id), + ))) + } Err(error) => { + trace_capability_latency_error( + "capability_host_invoke_json", + &capability_id, + &scope, + dispatch_started_at, + &error, + ); tracing::debug!( capability_id = %capability_id, error_kind = failure_kind_from(&error).as_str(), idempotency_key = idempotency_key.as_deref().unwrap_or(""), "capability invocation failed" ); - self.translate_invocation_error(error, capability_id, scope, invocation_id) - .await + let translated = self + .translate_invocation_error( + error, + capability_id.clone(), + scope.clone(), + invocation_id, + ) + .await; + match &translated { + Ok(_) => trace_capability_latency_ok( + "invoke_capability", + &capability_id, + &scope, + total_started_at, + ), + Err(error) => trace_capability_latency_error( + "invoke_capability", + &capability_id, + &scope, + total_started_at, + error, + ), + } + translated } } } diff --git a/crates/ironclaw_host_runtime/src/services/process_executor.rs b/crates/ironclaw_host_runtime/src/services/process_executor.rs index 9dcae91b85c..74a9a079c26 100644 --- a/crates/ironclaw_host_runtime/src/services/process_executor.rs +++ b/crates/ironclaw_host_runtime/src/services/process_executor.rs @@ -1,11 +1,118 @@ -use std::sync::Arc; +use std::{sync::Arc, time::Instant}; use async_trait::async_trait; use ironclaw_host_api::{CapabilityDispatchRequest, CapabilityDispatcher, RuntimeKind}; +use ironclaw_observability::live_latency_started_at; use ironclaw_processes::{ ProcessExecutionError, ProcessExecutionRequest, ProcessExecutionResult, ProcessExecutor, }; +struct ProcessLatencyFields { + capability_id: String, + runtime: String, + tenant_id: String, + user_id: String, + agent_id: String, + project_id: String, + mission_id: String, + thread_id: String, + invocation_id: String, +} + +impl ProcessLatencyFields { + fn from_request( + started_at: Option, + request: &ProcessExecutionRequest, + ) -> Option { + started_at?; + Some(Self { + capability_id: request.capability_id.to_string(), + runtime: format!("{:?}", request.runtime), + tenant_id: request.scope.tenant_id.as_str().to_string(), + user_id: request.scope.user_id.as_str().to_string(), + agent_id: request + .scope + .agent_id + .as_ref() + .map(|id| id.as_str().to_string()) + .unwrap_or_default(), + project_id: request + .scope + .project_id + .as_ref() + .map(|id| id.as_str().to_string()) + .unwrap_or_default(), + mission_id: request + .scope + .mission_id + .as_ref() + .map(|id| id.as_str().to_string()) + .unwrap_or_default(), + thread_id: request + .scope + .thread_id + .as_ref() + .map(|id| id.as_str().to_string()) + .unwrap_or_default(), + invocation_id: request.scope.invocation_id.to_string(), + }) + } +} + +fn trace_process_latency_ok( + operation: &'static str, + fields: Option<&ProcessLatencyFields>, + started_at: Option, +) { + let (Some(fields), Some(started_at)) = (fields, started_at) else { + return; + }; + + ironclaw_observability::live_latency_trace_ok!( + "process_executor", + operation, + Some(started_at), + capability_id = fields.capability_id.as_str(), + runtime = fields.runtime.as_str(), + tenant_id = fields.tenant_id.as_str(), + user_id = fields.user_id.as_str(), + agent_id = fields.agent_id.as_str(), + project_id = fields.project_id.as_str(), + mission_id = fields.mission_id.as_str(), + thread_id = fields.thread_id.as_str(), + invocation_id = fields.invocation_id.as_str(), + "process execution operation completed", + ); +} + +fn trace_process_latency_error( + operation: &'static str, + fields: Option<&ProcessLatencyFields>, + started_at: Option, + _error: &E, +) { + let (Some(fields), Some(started_at)) = (fields, started_at) else { + return; + }; + + ironclaw_observability::live_latency_trace_error!( + "process_executor", + operation, + Some(started_at), + "process_execution_error", + capability_id = fields.capability_id.as_str(), + runtime = fields.runtime.as_str(), + tenant_id = fields.tenant_id.as_str(), + user_id = fields.user_id.as_str(), + agent_id = fields.agent_id.as_str(), + project_id = fields.project_id.as_str(), + mission_id = fields.mission_id.as_str(), + thread_id = fields.thread_id.as_str(), + invocation_id = fields.invocation_id.as_str(), + "process execution operation failed", + ); +} + #[derive(Clone)] pub(super) struct HostProcessExecutor { dispatch_executor: Arc, @@ -30,15 +137,44 @@ impl ProcessExecutor for HostProcessExecutor { &self, request: ProcessExecutionRequest, ) -> Result { + let started_at = live_latency_started_at(); + let fields = ProcessLatencyFields::from_request(started_at, &request); if is_process_sandbox_request(&request) { let Some(executor) = &self.process_sandbox_executor else { - return Err(ProcessExecutionError::new( - "missing_process_sandbox_executor", - )); + let error = ProcessExecutionError::new("missing_process_sandbox_executor"); + trace_process_latency_error( + "host_process_execute", + fields.as_ref(), + started_at, + &error, + ); + return Err(error); }; - return executor.execute(request).await; + let result = executor.execute(request).await; + match &result { + Ok(_) => { + trace_process_latency_ok("host_process_execute", fields.as_ref(), started_at) + } + Err(error) => trace_process_latency_error( + "host_process_execute", + fields.as_ref(), + started_at, + error, + ), + } + return result; } - self.dispatch_executor.execute(request).await + let result = self.dispatch_executor.execute(request).await; + match &result { + Ok(_) => trace_process_latency_ok("host_process_execute", fields.as_ref(), started_at), + Err(error) => trace_process_latency_error( + "host_process_execute", + fields.as_ref(), + started_at, + error, + ), + } + result } } @@ -64,8 +200,17 @@ impl ProcessExecutor for RuntimeDispatchProcessExecutor { &self, request: ProcessExecutionRequest, ) -> Result { + let started_at = live_latency_started_at(); + let fields = ProcessLatencyFields::from_request(started_at, &request); if request.cancellation.is_cancelled() { - return Err(ProcessExecutionError::new("cancelled")); + let error = ProcessExecutionError::new("cancelled"); + trace_process_latency_error( + "runtime_dispatch_execute", + fields.as_ref(), + started_at, + &error, + ); + return Err(error); } let result = self .dispatcher @@ -80,11 +225,20 @@ impl ProcessExecutor for RuntimeDispatchProcessExecutor { .await .map_err(|error| ProcessExecutionError::new(error.event_kind()))?; if request.cancellation.is_cancelled() { - return Err(ProcessExecutionError::new("cancelled")); + let error = ProcessExecutionError::new("cancelled"); + trace_process_latency_error( + "runtime_dispatch_execute", + fields.as_ref(), + started_at, + &error, + ); + return Err(error); } - Ok(ProcessExecutionResult { + let result = Ok(ProcessExecutionResult { output: result.output, - }) + }); + trace_process_latency_ok("runtime_dispatch_execute", fields.as_ref(), started_at); + result } } diff --git a/crates/ironclaw_host_runtime/src/turn_scheduler.rs b/crates/ironclaw_host_runtime/src/turn_scheduler.rs index 5d3346e0966..c68f493cc66 100644 --- a/crates/ironclaw_host_runtime/src/turn_scheduler.rs +++ b/crates/ironclaw_host_runtime/src/turn_scheduler.rs @@ -5,6 +5,7 @@ use std::{ use async_trait::async_trait; use chrono::Utc; use futures_util::FutureExt; +use ironclaw_observability::live_latency_started_at; use ironclaw_turns::{ SanitizedFailure, TurnError, TurnLeaseToken, TurnRunId, TurnRunWake, TurnRunWakeNotifier, TurnRunWakeNotifyError, TurnRunnerId, TurnScope, @@ -22,6 +23,10 @@ use tokio_util::sync::CancellationToken; use tracing::Instrument; use tracing::debug; +mod executor_task; +mod latency; +use self::executor_task::ExecutorTaskOutcome; + #[derive(Debug, Clone)] pub struct TurnRunSchedulerConfig { max_concurrent_runs: usize, @@ -158,9 +163,7 @@ pub struct TurnRunExecutorError { impl TurnRunExecutorError { pub fn new(failure_category: impl Into) -> Result { - Ok(Self { - failure: SanitizedFailure::new(failure_category)?, - }) + SanitizedFailure::new(failure_category).map(|failure| Self { failure }) } pub fn failure(&self) -> &SanitizedFailure { @@ -303,9 +306,14 @@ impl fmt::Debug for SchedulerTurnRunWakeNotifier { impl TurnRunWakeNotifier for SchedulerTurnRunWakeNotifier { fn notify_queued_run(&self, wake: TurnRunWake) -> Result<(), TurnRunWakeNotifyError> { - self.command_tx + let started_at = live_latency_started_at(); + let trace_fields = latency::run_fields_from_wake(started_at, &wake); + let result = self + .command_tx .try_send(SchedulerCommand::Wake(wake)) - .map_err(|_| TurnRunWakeNotifyError::DeliveryUnavailable) + .map_err(|_| TurnRunWakeNotifyError::DeliveryUnavailable); + latency::notify_queued_run_result(trace_fields.as_ref(), started_at, &result); + result } } @@ -594,6 +602,8 @@ async fn drain_queued_runs( let Ok(permit) = Arc::clone(&context.semaphore).try_acquire_owned() else { return false; }; + let claim_started_at = live_latency_started_at(); + let scope_filter_fields = latency::scope_fields(claim_started_at, scope_filter.as_ref()); let claim = context .transitions .claim_next_run(ClaimRunRequest { @@ -602,6 +612,7 @@ async fn drain_queued_runs( scope_filter: scope_filter.clone(), }) .await; + latency::claim_next_run_result(scope_filter_fields.as_ref(), claim_started_at, &claim); match claim { Ok(Some(claimed)) => { let run_id = claimed.state.run_id; @@ -644,11 +655,6 @@ async fn drain_queued_runs( } } -enum ExecutorTaskOutcome { - Completed, - TerminalFailure(Option), -} - #[derive(Clone, Copy)] struct ExecutorTaskConfig { runner_heartbeat_interval: Duration, @@ -681,7 +687,7 @@ fn spawn_executor_task( // the event self-contained and allows test layers to find them without // relying on span registration timing (which can be racy under parallel // test execution when using `tracing::dispatcher::set_default`). - let recovery_thread_id = claimed.state.scope.thread_id.clone(); + let recovery_scope = claimed.state.scope.clone(); let recovery_run_id_for_start = claimed.state.run_id; executor_tasks.spawn( async move { @@ -689,10 +695,11 @@ fn spawn_executor_task( let recovery_runner_id = claimed.runner_id; let recovery_lease_token = claimed.lease_token; tracing::debug!( - thread_id = %recovery_thread_id, + thread_id = %recovery_scope.thread_id, run_id = %recovery_run_id_for_start, "turn run started", ); + let executor_started_at = live_latency_started_at(); let mut heartbeat_tick = interval(task_config.runner_heartbeat_interval); heartbeat_tick.set_missed_tick_behavior(MissedTickBehavior::Delay); // Consume the immediate first tick so the heartbeat loop never fires @@ -711,15 +718,12 @@ fn spawn_executor_task( tokio::select! { biased; result = &mut executor_result => { - break match result { - Ok(Ok(())) => ExecutorTaskOutcome::Completed, - Ok(Err(error)) => ExecutorTaskOutcome::TerminalFailure(Some( - error.failure().clone(), - )), - Err(_) => ExecutorTaskOutcome::TerminalFailure(scheduler_failure( - "scheduler_executor_panic", - )), - }; + break executor_task::result_to_outcome( + &recovery_scope, + recovery_run_id, + executor_started_at, + result, + ); } _ = heartbeat_tick.tick(), if heartbeats.is_idle() => { heartbeats.spawn( diff --git a/crates/ironclaw_host_runtime/src/turn_scheduler/executor_task.rs b/crates/ironclaw_host_runtime/src/turn_scheduler/executor_task.rs new file mode 100644 index 00000000000..98974ee6332 --- /dev/null +++ b/crates/ironclaw_host_runtime/src/turn_scheduler/executor_task.rs @@ -0,0 +1,39 @@ +use std::{any::Any, time::Instant}; + +use ironclaw_turns::{SanitizedFailure, TurnRunId, TurnScope}; + +use super::{TurnRunExecutorError, latency, scheduler_failure}; + +pub(super) enum ExecutorTaskOutcome { + Completed, + TerminalFailure(Option), +} + +pub(super) fn result_to_outcome( + scope: &TurnScope, + run_id: TurnRunId, + started_at: Option, + result: Result, Box>, +) -> ExecutorTaskOutcome { + match result { + Ok(Ok(())) => { + latency::operation_ok("execute_claimed_run", scope, run_id, started_at); + ExecutorTaskOutcome::Completed + } + Ok(Err(error)) => { + latency::operation_error( + "execute_claimed_run", + scope, + run_id, + started_at, + "executor_error", + ); + ExecutorTaskOutcome::TerminalFailure(Some(error.failure().clone())) + } + Err(_) => { + let reason = "scheduler_executor_panic"; + latency::operation_error("execute_claimed_run", scope, run_id, started_at, reason); + ExecutorTaskOutcome::TerminalFailure(scheduler_failure(reason)) + } + } +} diff --git a/crates/ironclaw_host_runtime/src/turn_scheduler/latency.rs b/crates/ironclaw_host_runtime/src/turn_scheduler/latency.rs new file mode 100644 index 00000000000..13017796e59 --- /dev/null +++ b/crates/ironclaw_host_runtime/src/turn_scheduler/latency.rs @@ -0,0 +1,201 @@ +use std::time::Instant; + +use ironclaw_turns::{TurnRunId, TurnRunWake, TurnScope, runner::ClaimedTurnRun}; + +pub(super) struct RunFields { + scope: ScopeFields, + run_id: TurnRunId, +} + +pub(super) struct ScopeFields { + tenant_id: String, + agent_id: String, + project_id: String, + thread_id: String, + owner_user_id: String, +} + +impl ScopeFields { + fn from_scope(scope: &TurnScope) -> Self { + Self { + tenant_id: scope.tenant_id.as_str().to_string(), + agent_id: scope + .agent_id + .as_ref() + .map(|id| id.as_str().to_string()) + .unwrap_or_default(), + project_id: scope + .project_id + .as_ref() + .map(|id| id.as_str().to_string()) + .unwrap_or_default(), + thread_id: scope.thread_id.as_str().to_string(), + owner_user_id: scope + .explicit_owner_user_id() + .map(|id| id.as_str().to_string()) + .unwrap_or_default(), + } + } +} + +pub(super) fn run_fields_from_wake( + started_at: Option, + wake: &TurnRunWake, +) -> Option { + started_at?; + Some(RunFields { + scope: ScopeFields::from_scope(&wake.scope), + run_id: wake.run_id, + }) +} + +pub(super) fn scope_fields( + started_at: Option, + scope: Option<&TurnScope>, +) -> Option { + started_at?; + scope.map(ScopeFields::from_scope) +} + +pub(super) fn operation_ok( + operation: &'static str, + scope: &TurnScope, + run_id: TurnRunId, + started_at: Option, +) { + if started_at.is_none() { + return; + } + trace_ok( + operation, + Some(&ScopeFields::from_scope(scope)), + Some(run_id), + started_at, + ); +} + +pub(super) fn operation_error( + operation: &'static str, + scope: &TurnScope, + run_id: TurnRunId, + started_at: Option, + error_kind: &'static str, +) { + if started_at.is_none() { + return; + } + trace_error( + operation, + Some(&ScopeFields::from_scope(scope)), + Some(run_id), + started_at, + error_kind, + ); +} + +pub(super) fn notify_queued_run_result( + fields: Option<&RunFields>, + started_at: Option, + result: &Result<(), E>, +) { + let Some(fields) = fields else { + return; + }; + + match result { + Ok(()) => trace_ok( + "notify_queued_run", + Some(&fields.scope), + Some(fields.run_id), + started_at, + ), + Err(_) => trace_error( + "notify_queued_run", + Some(&fields.scope), + Some(fields.run_id), + started_at, + "notify_error", + ), + } +} + +pub(super) fn claim_next_run_result( + scope_filter: Option<&ScopeFields>, + started_at: Option, + claim: &Result, E>, +) { + match claim { + Ok(Some(claimed)) => trace_ok( + "claim_next_run", + Some(&ScopeFields::from_scope(&claimed.state.scope)), + Some(claimed.state.run_id), + started_at, + ), + Ok(None) => trace_ok("claim_next_run_empty", scope_filter, None, started_at), + Err(_) => trace_error( + "claim_next_run", + scope_filter, + None, + started_at, + "claim_error", + ), + } +} + +fn trace_ok( + operation: &'static str, + scope: Option<&ScopeFields>, + run_id: Option, + started_at: Option, +) { + let run_id = run_id.map(|id| id.to_string()).unwrap_or_default(); + let tenant_id = scope.map(|scope| scope.tenant_id.as_str()).unwrap_or(""); + let agent_id = scope.map(|scope| scope.agent_id.as_str()).unwrap_or(""); + let project_id = scope.map(|scope| scope.project_id.as_str()).unwrap_or(""); + let thread_id = scope.map(|scope| scope.thread_id.as_str()).unwrap_or(""); + let owner_user_id = scope + .map(|scope| scope.owner_user_id.as_str()) + .unwrap_or(""); + ironclaw_observability::live_latency_trace_ok!( + "turn_scheduler", + operation, + started_at, + tenant_id = tenant_id, + agent_id = agent_id, + project_id = project_id, + thread_id = thread_id, + owner_user_id = owner_user_id, + run_id = run_id.as_str(), + "turn scheduler operation completed", + ); +} + +fn trace_error( + operation: &'static str, + scope: Option<&ScopeFields>, + run_id: Option, + started_at: Option, + error_kind: &'static str, +) { + let run_id = run_id.map(|id| id.to_string()).unwrap_or_default(); + let tenant_id = scope.map(|scope| scope.tenant_id.as_str()).unwrap_or(""); + let agent_id = scope.map(|scope| scope.agent_id.as_str()).unwrap_or(""); + let project_id = scope.map(|scope| scope.project_id.as_str()).unwrap_or(""); + let thread_id = scope.map(|scope| scope.thread_id.as_str()).unwrap_or(""); + let owner_user_id = scope + .map(|scope| scope.owner_user_id.as_str()) + .unwrap_or(""); + ironclaw_observability::live_latency_trace_error!( + "turn_scheduler", + operation, + started_at, + error_kind, + tenant_id = tenant_id, + agent_id = agent_id, + project_id = project_id, + thread_id = thread_id, + owner_user_id = owner_user_id, + run_id = run_id.as_str(), + "turn scheduler operation failed", + ); +} diff --git a/crates/ironclaw_observability/Cargo.toml b/crates/ironclaw_observability/Cargo.toml new file mode 100644 index 00000000000..1aa35180eac --- /dev/null +++ b/crates/ironclaw_observability/Cargo.toml @@ -0,0 +1,14 @@ +[package] +name = "ironclaw_observability" +version = "0.1.0" +edition = "2024" +rust-version.workspace = true +description = "Low-level observability helpers for IronClaw" +authors = ["NEAR AI "] +license = "MIT OR Apache-2.0" +homepage = "https://github.com/nearai/ironclaw" +repository = "https://github.com/nearai/ironclaw" +publish = false + +[dependencies] +tracing = "0.1" diff --git a/crates/ironclaw_observability/src/lib.rs b/crates/ironclaw_observability/src/lib.rs new file mode 100644 index 00000000000..16ed0cd5b0f --- /dev/null +++ b/crates/ironclaw_observability/src/lib.rs @@ -0,0 +1,65 @@ +//! Shared low-level observability helpers. +#![warn(unreachable_pub)] + +use std::time::Instant; + +pub use tracing; + +#[inline] +pub fn elapsed_ms(started_at: Instant) -> u64 { + started_at + .elapsed() + .as_millis() + .try_into() + .unwrap_or(u64::MAX) +} + +#[inline] +pub fn live_latency_enabled() -> bool { + tracing::enabled!(target: "ironclaw_latency", tracing::Level::TRACE) +} + +#[inline] +pub fn live_latency_started_at() -> Option { + live_latency_enabled().then(Instant::now) +} + +#[macro_export] +macro_rules! live_latency_trace { + ($($fields:tt)*) => { + $crate::tracing::trace!(target: "ironclaw_latency", $($fields)*) + }; +} + +#[macro_export] +macro_rules! live_latency_trace_ok { + ($component:expr, $operation:expr, $started_at:expr, $($fields:tt)*) => { + if let Some(started_at) = $started_at { + let elapsed_ms = $crate::elapsed_ms(started_at); + $crate::live_latency_trace!( + component = $component, + operation = $operation, + elapsed_ms, + outcome = "ok", + $($fields)* + ); + } + }; +} + +#[macro_export] +macro_rules! live_latency_trace_error { + ($component:expr, $operation:expr, $started_at:expr, $error_kind:expr, $($fields:tt)*) => { + if let Some(started_at) = $started_at { + let elapsed_ms = $crate::elapsed_ms(started_at); + $crate::live_latency_trace!( + component = $component, + operation = $operation, + elapsed_ms, + outcome = "error", + error_kind = $error_kind, + $($fields)* + ); + } + }; +} diff --git a/crates/ironclaw_reborn/Cargo.toml b/crates/ironclaw_reborn/Cargo.toml index 2d520cd1156..684ae2d5929 100644 --- a/crates/ironclaw_reborn/Cargo.toml +++ b/crates/ironclaw_reborn/Cargo.toml @@ -56,6 +56,7 @@ ironclaw_host_api = { path = "../ironclaw_host_api", version = "0.1.0" } ironclaw_host_runtime = { path = "../ironclaw_host_runtime", version = "0.1.0" } ironclaw_llm = { path = "../ironclaw_llm", version = "0.1.0", optional = true, default-features = false } ironclaw_loop_support = { path = "../ironclaw_loop_support", version = "0.1.0" } +ironclaw_observability = { path = "../ironclaw_observability" } ironclaw_safety = { path = "../ironclaw_safety", version = "0.2.2" } ironclaw_secrets = { path = "../ironclaw_secrets", version = "0.1.0", optional = true } ironclaw_filesystem = { path = "../ironclaw_filesystem", version = "0.1.0", optional = true } diff --git a/crates/ironclaw_reborn/src/model_gateway.rs b/crates/ironclaw_reborn/src/model_gateway.rs index d7f24978298..8d06e68a483 100644 --- a/crates/ironclaw_reborn/src/model_gateway.rs +++ b/crates/ironclaw_reborn/src/model_gateway.rs @@ -10,6 +10,7 @@ use std::{ Arc, atomic::{AtomicU64, Ordering}, }, + time::Instant, }; use async_trait::async_trait; @@ -29,6 +30,7 @@ use ironclaw_loop_support::{ ModelCost, StaticModelCostTable, ThreadBackedLoopContextPort, ThreadBackedLoopModelPort, ThreadContextWindowCache, }; +use ironclaw_observability::live_latency_started_at; use ironclaw_safety::{ is_provider_arguments_too_large_summary, provider_arguments_exceed_max_bytes, }; @@ -61,6 +63,42 @@ const PROVIDER_TOOL_ARGUMENTS_OMITTED_MARKER: &str = "arguments omitted because they exceeded the host provider-tool limit"; const UNAVAILABLE_CAPABILITY_REPLY: &str = "That capability is unavailable or disabled for this request, so I will not route it through another tool."; +fn trace_model_latency_ok( + operation: &'static str, + replay_identity: &ProviderReplayIdentity, + provider_turn_scope: Option<&str>, + started_at: Option, +) { + ironclaw_observability::live_latency_trace_ok!( + "model_gateway", + operation, + started_at, + provider_id = %replay_identity.provider_id, + provider_model_id = %replay_identity.provider_model_id, + provider_turn_scope = provider_turn_scope.unwrap_or(""), + "model gateway operation completed", + ); +} + +fn trace_model_latency_error( + operation: &'static str, + replay_identity: &ProviderReplayIdentity, + provider_turn_scope: Option<&str>, + started_at: Option, + _error: &E, +) { + ironclaw_observability::live_latency_trace_error!( + "model_gateway", + operation, + started_at, + "model_gateway_error", + provider_id = %replay_identity.provider_id, + provider_model_id = %replay_identity.provider_model_id, + provider_turn_scope = provider_turn_scope.unwrap_or(""), + "model gateway operation failed", + ); +} + /// Fail-closed routing policy from resolved Reborn model profile ids to the /// host-selected provider/model envelope. #[derive(Debug, Clone, Default)] @@ -847,12 +885,31 @@ where let tool_request = ToolCompletionRequest::from_completion_request(completion, llm_tool_definitions); debug!("reborn model gateway dispatching tool-capable provider request"); - let response = provider - .complete_with_tools(tool_request.clone()) - .await - .map_err(map_provider_error)?; + let provider_started_at = live_latency_started_at(); + let response = match provider.complete_with_tools(tool_request.clone()).await { + Ok(response) => { + trace_model_latency_ok( + "provider_complete_with_tools", + &replay_identity, + provider_turn_scope.as_deref(), + provider_started_at, + ); + response + } + Err(error) => { + trace_model_latency_error( + "provider_complete_with_tools", + &replay_identity, + provider_turn_scope.as_deref(), + provider_started_at, + &error, + ); + return Err(map_provider_error(error)); + } + }; let response = recover_textual_tool_calls_from_tool_response(response, &recovery_tool_names)?; + let host_response_started_at = live_latency_started_at(); match tool_response_to_host( response.clone(), Arc::clone(&capabilities), @@ -864,8 +921,23 @@ where ) .await { - Ok(response) => return Ok(response), + Ok(response) => { + trace_model_latency_ok( + "tool_response_to_host", + &replay_identity, + provider_turn_scope.as_deref(), + host_response_started_at, + ); + return Ok(response); + } Err(error) if is_repairable_provider_tool_output_error(&error) => { + trace_model_latency_error( + "tool_response_to_host", + &replay_identity, + provider_turn_scope.as_deref(), + host_response_started_at, + &error, + ); debug!( safe_summary = error.safe_summary.as_str(), "reborn model gateway retrying after repairable provider tool output" @@ -878,16 +950,35 @@ where error.safe_summary.as_str(), )); let rejected_response = response; - let response = provider - .complete_with_tools(repair_request) - .await - .map_err(map_provider_error)?; + let retry_started_at = live_latency_started_at(); + let response = match provider.complete_with_tools(repair_request).await { + Ok(response) => { + trace_model_latency_ok( + "provider_complete_with_tools_repair", + &replay_identity, + provider_turn_scope.as_deref(), + retry_started_at, + ); + response + } + Err(error) => { + trace_model_latency_error( + "provider_complete_with_tools_repair", + &replay_identity, + provider_turn_scope.as_deref(), + retry_started_at, + &error, + ); + return Err(map_provider_error(error)); + } + }; let mut response = recover_textual_tool_calls_from_tool_response( response, &recovery_tool_names, )?; accumulate_tool_response_usage(&mut response, &rejected_response); - return tool_response_to_host( + let repair_host_started_at = live_latency_started_at(); + let result = tool_response_to_host( response, capabilities, provider_turn_scope @@ -897,8 +988,33 @@ where unavailable_capability_guard.as_ref(), ) .await; + match &result { + Ok(_) => trace_model_latency_ok( + "tool_response_to_host_repair", + &replay_identity, + provider_turn_scope.as_deref(), + repair_host_started_at, + ), + Err(error) => trace_model_latency_error( + "tool_response_to_host_repair", + &replay_identity, + provider_turn_scope.as_deref(), + repair_host_started_at, + error, + ), + } + return result; + } + Err(error) => { + trace_model_latency_error( + "tool_response_to_host", + &replay_identity, + provider_turn_scope.as_deref(), + host_response_started_at, + &error, + ); + return Err(error); } - Err(error) => return Err(error), } } debug!( @@ -910,10 +1026,28 @@ where ); } - let response = provider - .complete(completion) - .await - .map_err(map_provider_error)?; + let provider_started_at = live_latency_started_at(); + let response = match provider.complete(completion).await { + Ok(response) => { + trace_model_latency_ok( + "provider_complete", + &replay_identity, + provider_turn_scope.as_deref(), + provider_started_at, + ); + response + } + Err(error) => { + trace_model_latency_error( + "provider_complete", + &replay_identity, + provider_turn_scope.as_deref(), + provider_started_at, + &error, + ); + return Err(map_provider_error(error)); + } + }; debug!( finish_reason = ?response.finish_reason, content_bytes = response.content.len(), diff --git a/crates/ironclaw_reborn/src/turn_run_executor.rs b/crates/ironclaw_reborn/src/turn_run_executor.rs index 787177972f8..6b876d58f41 100644 --- a/crates/ironclaw_reborn/src/turn_run_executor.rs +++ b/crates/ironclaw_reborn/src/turn_run_executor.rs @@ -3,10 +3,14 @@ //! Adapts `RebornLoopDriverHostFactory` + `DriverRegistry` + `LoopExitApplier` //! to the `TurnRunExecutor` trait consumed by `TurnRunScheduler`. -use std::sync::{Arc, OnceLock}; +use std::{ + sync::{Arc, OnceLock}, + time::Instant, +}; use async_trait::async_trait; use ironclaw_host_runtime::{TurnRunExecutor, TurnRunExecutorError}; +use ironclaw_observability::live_latency_started_at; use ironclaw_turns::{ AgentLoopDriverError, AgentLoopDriverResumeRequest, AgentLoopDriverRunRequest, LoopExit, TurnStatus, @@ -24,6 +28,46 @@ use crate::{ turn_runner::{HostFactory, sanitized_driver_failure, sanitized_failure}, }; +fn trace_executor_latency_ok( + operation: &'static str, + claimed: &ClaimedTurnRun, + started_at: Option, +) { + ironclaw_observability::live_latency_trace_ok!( + "reborn_turn_executor", + operation, + started_at, + tenant_id = %claimed.state.scope.tenant_id, + agent_id = claimed.state.scope.agent_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + project_id = claimed.state.scope.project_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + thread_id = %claimed.state.scope.thread_id, + owner_user_id = claimed.state.scope.explicit_owner_user_id().map(|id| id.as_str()).unwrap_or(""), + run_id = %claimed.state.run_id, + "reborn turn executor operation completed", + ); +} + +fn trace_executor_latency_error( + operation: &'static str, + claimed: &ClaimedTurnRun, + started_at: Option, + _error: &E, +) { + ironclaw_observability::live_latency_trace_error!( + "reborn_turn_executor", + operation, + started_at, + "executor_error", + tenant_id = %claimed.state.scope.tenant_id, + agent_id = claimed.state.scope.agent_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + project_id = claimed.state.scope.project_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + thread_id = %claimed.state.scope.thread_id, + owner_user_id = claimed.state.scope.explicit_owner_user_id().map(|id| id.as_str()).unwrap_or(""), + run_id = %claimed.state.run_id, + "reborn turn executor operation failed", + ); +} + /// A `TurnRunExecutorError` for the static category `"unknown_failure"`. /// /// Built once on first access via `OnceLock`. Used as a guaranteed-valid @@ -91,11 +135,27 @@ impl TurnRunExecutor for RebornTurnRunExecutor { claimed: ClaimedTurnRun, transitions: Arc, ) -> Result<(), TurnRunExecutorError> { + let started_at = live_latency_started_at(); match self.invoke_driver(&claimed, &transitions).await { - Ok(exit) => self - .apply_exit(&claimed, exit, &transitions) - .await - .map_err(|()| unknown_failure_error().clone()), + Ok(exit) => { + let result = self.apply_exit(&claimed, exit, &transitions).await; + match result { + Ok(()) => { + trace_executor_latency_ok("execute_claimed_run", &claimed, started_at); + Ok(()) + } + Err(()) => { + let error = unknown_failure_error().clone(); + trace_executor_latency_error( + "execute_claimed_run", + &claimed, + started_at, + &error, + ); + Err(error) + } + } + } Err(err) => { let sanitized = match &err { DriverInvocationError::DriverError(AgentLoopDriverError::Failed { @@ -123,8 +183,10 @@ impl TurnRunExecutor for RebornTurnRunExecutor { // guard that is never reached in practice. let failure = sanitized.unwrap_or_else(|| unknown_failure_error().failure().clone()); - Err(TurnRunExecutorError::new(failure.category()) - .unwrap_or_else(|_| unknown_failure_error().clone())) + let error = TurnRunExecutorError::new(failure.category()) + .unwrap_or_else(|_| unknown_failure_error().clone()); + trace_executor_latency_error("execute_claimed_run", &claimed, started_at, &error); + Err(error) } } } @@ -157,24 +219,52 @@ impl RebornTurnRunExecutor { "reborn executor resolved loop driver" ); - let host = self - .host_factory - .create_host(claimed) + let host_started_at = live_latency_started_at(); + let host = match self.host_factory.create_host(claimed).await { + Ok(host) => { + trace_executor_latency_ok("create_loop_host", claimed, host_started_at); + host + } + Err(err) => { + trace_executor_latency_error("create_loop_host", claimed, host_started_at, &err); + // Use the error's full `Display` (`err.to_string()`) rather than a single + // field, so whatever context the host factory embedded in its message + // survives into `reason` (HostFactoryError is a flat message with no + // `source()` chain of its own). + return Err(DriverInvocationError::HostCreationFailed { + reason: err.to_string(), + }); + } + }; + let route_snapshot_started_at = live_latency_started_at(); + if let Err(error) = self + .persist_model_route_snapshot(claimed, host.as_ref(), transitions) .await - // Use the error's full `Display` (`err.to_string()`) rather than a single - // field, so whatever context the host factory embedded in its message - // survives into `reason` (HostFactoryError is a flat message with no - // `source()` chain of its own). - .map_err(|err| DriverInvocationError::HostCreationFailed { - reason: err.to_string(), - })?; - self.persist_model_route_snapshot(claimed, host.as_ref(), transitions) - .await?; + { + trace_executor_latency_error( + "persist_model_route_snapshot", + claimed, + route_snapshot_started_at, + &error, + ); + return Err(error); + } + trace_executor_latency_ok( + "persist_model_route_snapshot", + claimed, + route_snapshot_started_at, + ); let turn_id = claimed.state.turn_id; let run_id = claimed.state.run_id; - match (claimed.state.status, claimed.state.checkpoint_id) { + let driver_started_at = live_latency_started_at(); + let driver_operation = if claimed.state.checkpoint_id.is_some() { + "driver_resume" + } else { + "driver_run" + }; + let driver_result = match (claimed.state.status, claimed.state.checkpoint_id) { // Requeued blocked runs keep their checkpoint while returning to // `Queued`; checkpoint identity is the resume signal. (_, Some(checkpoint_id)) => driver @@ -213,7 +303,14 @@ impl RebornTurnRunExecutor { ) .await .map_err(DriverInvocationError::DriverError), + }; + match &driver_result { + Ok(_) => trace_executor_latency_ok(driver_operation, claimed, driver_started_at), + Err(error) => { + trace_executor_latency_error(driver_operation, claimed, driver_started_at, error) + } } + driver_result } async fn persist_model_route_snapshot( @@ -255,12 +352,14 @@ impl RebornTurnRunExecutor { exit: LoopExit, transitions: &Arc, ) -> Result<(), ()> { + let started_at = live_latency_started_at(); let run_id = claimed.state.run_id; let runner_id = claimed.runner_id; let lease_token = claimed.lease_token; match self.loop_exit_applier.apply(claimed, exit).await { Ok(state) => { + trace_executor_latency_ok("apply_loop_exit", claimed, started_at); debug!( runner_id = ?runner_id, run_id = ?run_id, @@ -270,6 +369,7 @@ impl RebornTurnRunExecutor { Ok(()) } Err(err) => { + trace_executor_latency_error("apply_loop_exit", claimed, started_at, &err); error!( runner_id = ?runner_id, run_id = ?run_id, diff --git a/crates/ironclaw_reborn_composition/Cargo.toml b/crates/ironclaw_reborn_composition/Cargo.toml index 07a093903b9..8fed69dc03e 100644 --- a/crates/ironclaw_reborn_composition/Cargo.toml +++ b/crates/ironclaw_reborn_composition/Cargo.toml @@ -115,6 +115,7 @@ ironclaw_llm = { path = "../ironclaw_llm", optional = true, default-features = f ironclaw_loop_support = { path = "../ironclaw_loop_support" } ironclaw_mcp = { path = "../ironclaw_mcp" } ironclaw_network = { path = "../ironclaw_network" } +ironclaw_observability = { path = "../ironclaw_observability" } ironclaw_outbound = { path = "../ironclaw_outbound" } ironclaw_processes = { path = "../ironclaw_processes" } ironclaw_product_adapters = { path = "../ironclaw_product_adapters" } diff --git a/crates/ironclaw_reborn_composition/src/runtime.rs b/crates/ironclaw_reborn_composition/src/runtime.rs index 8d4dbb1b275..4e809e8119c 100644 --- a/crates/ironclaw_reborn_composition/src/runtime.rs +++ b/crates/ironclaw_reborn_composition/src/runtime.rs @@ -51,6 +51,7 @@ use ironclaw_loop_support::{ LoopCapabilityInputResolver, LoopCapabilityPortFactory, LoopCapabilityResultWriter, ModelGatewayBackedSystemInferencePort, }; +use ironclaw_observability::live_latency_started_at; use ironclaw_product_adapters::ProjectionStream; use ironclaw_product_workflow::{ ApprovalBlockedTurnRun, ApprovalInteractionScope, ApprovalInteractionService, @@ -99,6 +100,7 @@ use ironclaw_product_workflow::{ }; use ironclaw_turns::run_profile::UserProfileContext; +use self::latency::{trace_runtime_latency_error, trace_runtime_latency_ok}; use self::runtime_turn_scheduler::RuntimeTurnScheduler; use crate::default_system_prompt::DefaultSystemPromptIdentitySource; use crate::factory::{LocalDevRootFilesystem, LocalDevTurnStateStore, builtin_extension_registry}; @@ -339,6 +341,7 @@ mod auth_interaction_tests; #[cfg(test)] #[path = "runtime/tests/default_system_prompt.rs"] mod default_system_prompt_tests; +mod latency; mod local_dev; #[cfg(test)] #[path = "runtime/tests/outbound_delivery.rs"] @@ -1689,15 +1692,46 @@ impl RebornRuntime { cancellation: CancellationToken, capture_skill_execution_plan: bool, ) -> Result { - let submitted = self + let total_started_at = live_latency_started_at(); + let submit_started_at = total_started_at; + let submitted = match self .submit_user_turn( conversation, text, &cancellation, capture_skill_execution_plan, ) - .await?; + .await + { + Ok(submitted) => { + trace_runtime_latency_ok( + "submit_user_turn", + &conversation.0, + Some(submitted.run_id), + submit_started_at, + ); + submitted + } + Err(error) => { + trace_runtime_latency_error( + "submit_user_turn", + &conversation.0, + None, + submit_started_at, + &error, + ); + trace_runtime_latency_error( + "send_user_message", + &conversation.0, + None, + total_started_at, + &error, + ); + return Err(error); + } + }; + let wait_started_at = live_latency_started_at(); let reply = async { let terminal_state = self .wait_for_terminal(&submitted.scope, submitted.run_id, &cancellation) @@ -1718,6 +1752,21 @@ impl RebornRuntime { }) } .await; + match &reply { + Ok(_) => trace_runtime_latency_ok( + "wait_for_terminal_and_read_reply", + &conversation.0, + Some(submitted.run_id), + wait_started_at, + ), + Err(error) => trace_runtime_latency_error( + "wait_for_terminal_and_read_reply", + &conversation.0, + Some(submitted.run_id), + wait_started_at, + error, + ), + } if let Some(skill_activation_source) = &self.skill_activation_source && let Err(clear_error) = skill_activation_source @@ -1726,6 +1775,13 @@ impl RebornRuntime { if reply.is_ok() { // Primary turn succeeded, so the cleanup failure is the only // error to surface. + trace_runtime_latency_error( + "send_user_message", + &conversation.0, + Some(submitted.run_id), + total_started_at, + &clear_error, + ); return Err(RebornRuntimeError::TurnSubmission(clear_error.to_string())); } // Primary turn already failed: don't mask it with the cleanup @@ -1737,6 +1793,21 @@ impl RebornRuntime { ); } + match &reply { + Ok(_) => trace_runtime_latency_ok( + "send_user_message", + &conversation.0, + Some(submitted.run_id), + total_started_at, + ), + Err(error) => trace_runtime_latency_error( + "send_user_message", + &conversation.0, + Some(submitted.run_id), + total_started_at, + error, + ), + } reply } @@ -1754,14 +1825,30 @@ impl RebornRuntime { capture_skill_execution_plan: bool, ) -> Result { let send_lock = self.send_lock_for(conversation).await; + let send_lock_started_at = live_latency_started_at(); let _send_guard = send_lock.lock_owned().await; + trace_runtime_latency_ok( + "send_lock_wait", + &conversation.0, + None, + send_lock_started_at, + ); // Stopped only when every worker has exited; a single crashed worker must not // reject submissions while others run. if self.turn_scheduler.is_stopped() { - return Err(RebornRuntimeError::WorkerStopped); + let error = RebornRuntimeError::WorkerStopped; + trace_runtime_latency_error( + "submit_user_turn_preflight", + &conversation.0, + None, + send_lock_started_at, + &error, + ); + return Err(error); } let scope = self.turn_scope_for(&conversation.0); - let accepted = self + let accept_started_at = live_latency_started_at(); + let accepted = match self .thread_service .accept_inbound_message(AcceptInboundMessageRequest { scope: self.thread_scope.clone(), @@ -1780,7 +1867,27 @@ impl RebornRuntime { content: MessageContent::text(text.to_string()), }) .await - .map_err(|error| RebornRuntimeError::ThreadService(error.to_string()))?; + { + Ok(accepted) => { + trace_runtime_latency_ok( + "accept_inbound_message", + &conversation.0, + None, + accept_started_at, + ); + accepted + } + Err(error) => { + trace_runtime_latency_error( + "accept_inbound_message", + &conversation.0, + None, + accept_started_at, + &error, + ); + return Err(RebornRuntimeError::ThreadService(error.to_string())); + } + }; let accepted_message_ref = AcceptedMessageRef::new(format!("msg:{}", accepted.message_id)) .map_err(|reason| RebornRuntimeError::InvalidArgument { reason })?; @@ -1796,19 +1903,52 @@ impl RebornRuntime { .skill_execution_adapter .as_ref() .ok_or(RebornRuntimeError::SkillExecutionUnavailable)?; - adapter - .record_user_message_for_execution( - scope.clone(), - accepted_message_ref.clone(), - text, - ) - .map_err(|error| RebornRuntimeError::TurnSubmission(error.to_string()))?; + let skill_record_started_at = live_latency_started_at(); + if let Err(error) = adapter.record_user_message_for_execution( + scope.clone(), + accepted_message_ref.clone(), + text, + ) { + trace_runtime_latency_error( + "record_skill_execution_message", + &conversation.0, + None, + skill_record_started_at, + &error, + ); + return Err(RebornRuntimeError::TurnSubmission(error.to_string())); + } + trace_runtime_latency_ok( + "record_skill_execution_message", + &conversation.0, + None, + skill_record_started_at, + ); } else if let Some(skill_activation_source) = &self.skill_activation_source { - skill_activation_source - .record_user_message(scope.clone(), accepted_message_ref.clone(), text) - .map_err(|error| RebornRuntimeError::TurnSubmission(error.to_string()))?; + let skill_record_started_at = live_latency_started_at(); + if let Err(error) = skill_activation_source.record_user_message( + scope.clone(), + accepted_message_ref.clone(), + text, + ) { + trace_runtime_latency_error( + "record_skill_activation_message", + &conversation.0, + None, + skill_record_started_at, + &error, + ); + return Err(RebornRuntimeError::TurnSubmission(error.to_string())); + } + trace_runtime_latency_ok( + "record_skill_activation_message", + &conversation.0, + None, + skill_record_started_at, + ); } + let turn_submit_started_at = live_latency_started_at(); let response = match self .turn_coordinator .submit_turn(SubmitTurnRequest { @@ -1830,8 +1970,24 @@ impl RebornRuntime { }) .await { - Ok(response) => response, + Ok(response) => { + let SubmitTurnResponse::Accepted { run_id, .. } = &response; + trace_runtime_latency_ok( + "turn_coordinator_submit_turn", + &conversation.0, + Some(*run_id), + turn_submit_started_at, + ); + response + } Err(error) => { + trace_runtime_latency_error( + "turn_coordinator_submit_turn", + &conversation.0, + None, + turn_submit_started_at, + &error, + ); if let Some(skill_activation_source) = &self.skill_activation_source { skill_activation_source .clear_accepted_message(&scope, &accepted_message_ref) @@ -1864,12 +2020,19 @@ impl RebornRuntime { .await?; return Err(RebornRuntimeError::OperationCancelled); } + let notify_started_at = live_latency_started_at(); self.turn_scheduler.notify(TurnRunWake { scope: scope.clone(), run_id, status: submit_status, event_cursor: submit_cursor, }); + trace_runtime_latency_ok( + "turn_scheduler_notify", + &conversation.0, + Some(run_id), + notify_started_at, + ); Ok(SubmittedTurn { _send_guard, diff --git a/crates/ironclaw_reborn_composition/src/runtime/latency.rs b/crates/ironclaw_reborn_composition/src/runtime/latency.rs new file mode 100644 index 00000000000..fe2ef9941d3 --- /dev/null +++ b/crates/ironclaw_reborn_composition/src/runtime/latency.rs @@ -0,0 +1,40 @@ +use std::time::Instant; + +use ironclaw_host_api::ThreadId; +use ironclaw_turns::TurnRunId; + +pub(super) fn trace_runtime_latency_ok( + operation: &'static str, + thread_id: &ThreadId, + run_id: Option, + started_at: Option, +) { + let run_id = run_id.map(|id| id.to_string()).unwrap_or_default(); + ironclaw_observability::live_latency_trace_ok!( + "reborn_runtime", + operation, + started_at, + thread_id = %thread_id, + run_id = run_id.as_str(), + "reborn runtime operation completed", + ); +} + +pub(super) fn trace_runtime_latency_error( + operation: &'static str, + thread_id: &ThreadId, + run_id: Option, + started_at: Option, + _error: &E, +) { + let run_id = run_id.map(|id| id.to_string()).unwrap_or_default(); + ironclaw_observability::live_latency_trace_error!( + "reborn_runtime", + operation, + started_at, + "runtime_error", + thread_id = %thread_id, + run_id = run_id.as_str(), + "reborn runtime operation failed", + ); +} diff --git a/crates/ironclaw_turns/Cargo.toml b/crates/ironclaw_turns/Cargo.toml index dca67000155..735ce29767a 100644 --- a/crates/ironclaw_turns/Cargo.toml +++ b/crates/ironclaw_turns/Cargo.toml @@ -21,6 +21,7 @@ chrono = { version = "0.4", features = ["serde"] } chrono-tz = "0.10" ironclaw_filesystem = { path = "../ironclaw_filesystem", version = "0.1.0" } ironclaw_host_api = { path = "../ironclaw_host_api", version = "0.1.0" } +ironclaw_observability = { path = "../ironclaw_observability" } hex = "0.4.3" serde = { version = "1", features = ["derive"] } serde_json = "1" diff --git a/crates/ironclaw_turns/src/coordinator.rs b/crates/ironclaw_turns/src/coordinator.rs index a32d21ea50e..b471b9ced4c 100644 --- a/crates/ironclaw_turns/src/coordinator.rs +++ b/crates/ironclaw_turns/src/coordinator.rs @@ -1,14 +1,60 @@ use async_trait::async_trait; +use ironclaw_observability::live_latency_started_at; use std::{ collections::HashMap, fmt, panic::{AssertUnwindSafe, catch_unwind}, sync::{Arc, Mutex}, + time::Instant, }; use tracing::debug; const MAX_PREPARED_RUN_IDS: usize = 4096; +fn trace_coordinator_latency_ok( + operation: &'static str, + scope: &TurnScope, + run_id: Option, + started_at: Option, +) { + let run_id = run_id.map(|id| id.to_string()).unwrap_or_default(); + ironclaw_observability::live_latency_trace_ok!( + "turn_coordinator", + operation, + started_at, + tenant_id = %scope.tenant_id, + agent_id = scope.agent_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + project_id = scope.project_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + thread_id = %scope.thread_id, + owner_user_id = scope.explicit_owner_user_id().map(|id| id.as_str()).unwrap_or(""), + run_id = run_id.as_str(), + "turn coordinator operation completed", + ); +} + +fn trace_coordinator_latency_error( + operation: &'static str, + scope: &TurnScope, + run_id: Option, + started_at: Option, + _error: &E, +) { + let run_id = run_id.map(|id| id.to_string()).unwrap_or_default(); + ironclaw_observability::live_latency_trace_error!( + "turn_coordinator", + operation, + started_at, + "turn_error", + tenant_id = %scope.tenant_id, + agent_id = scope.agent_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + project_id = scope.project_id.as_ref().map(|id| id.as_str()).unwrap_or(""), + thread_id = %scope.thread_id, + owner_user_id = scope.explicit_owner_user_id().map(|id| id.as_str()).unwrap_or(""), + run_id = run_id.as_str(), + "turn coordinator operation failed", + ); +} + use crate::{ AdmissionRejection, CancelRunRequest, CancelRunResponse, GetRunStateRequest, InMemoryRunProfileResolver, ResumeTurnRequest, ResumeTurnResponse, RunProfileResolver, @@ -257,6 +303,7 @@ where &self, request: SubmitTurnRequest, ) -> Result { + let started_at = live_latency_started_at(); // If the caller passed a run id that came out of prepare_turn, verify // it is being submitted under the same scope it was prepared under, // unless this is a child run (parent_run_id set). Subagent spawn @@ -272,14 +319,36 @@ where return Err(TurnError::Unauthorized); } let scope = request.scope.clone(); - let response = self + let response = match self .store .submit_turn( request, self.admission_policy.as_ref(), self.run_profile_resolver.as_ref(), ) - .await?; + .await + { + Ok(response) => { + let SubmitTurnResponse::Accepted { run_id, .. } = &response; + trace_coordinator_latency_ok( + "store_submit_turn", + &scope, + Some(*run_id), + started_at, + ); + response + } + Err(error) => { + trace_coordinator_latency_error( + "store_submit_turn", + &scope, + None, + started_at, + &error, + ); + return Err(error); + } + }; let SubmitTurnResponse::Accepted { run_id, status, @@ -294,7 +363,12 @@ where resolved_run_profile_version = resolved_run_profile_version.as_u64(), "turn coordinator accepted turn with resolved run profile" ); - notify_queued_run_best_effort(self.wake_notifier.as_ref(), submit_wake(scope, &response)); + let wake_started_at = live_latency_started_at(); + notify_queued_run_best_effort( + self.wake_notifier.as_ref(), + submit_wake(scope.clone(), &response), + ); + trace_coordinator_latency_ok("notify_queued_run", &scope, Some(*run_id), wake_started_at); deferred_publication_error( self.publication_error_port.as_ref(), submit_event_cursor(&response), @@ -306,28 +380,84 @@ where &self, request: ResumeTurnRequest, ) -> Result { + let started_at = live_latency_started_at(); let scope = request.scope.clone(); - let response = self.store.resume_turn(request).await?; + let response = match self.store.resume_turn(request).await { + Ok(response) => { + trace_coordinator_latency_ok( + "store_resume_turn", + &scope, + Some(response.run_id), + started_at, + ); + response + } + Err(error) => { + trace_coordinator_latency_error( + "store_resume_turn", + &scope, + None, + started_at, + &error, + ); + return Err(error); + } + }; + let wake_started_at = live_latency_started_at(); notify_queued_run_best_effort( self.wake_notifier.as_ref(), resume_wake(scope.clone(), &response), ); + trace_coordinator_latency_ok( + "notify_resumed_run", + &scope, + Some(response.run_id), + wake_started_at, + ); deferred_publication_error(self.publication_error_port.as_ref(), response.event_cursor)?; Ok(response) } async fn cancel_run(&self, request: CancelRunRequest) -> Result { + let started_at = live_latency_started_at(); let scope = request.scope.clone(); - let response = self.store.request_cancel(request).await?; + let response = match self.store.request_cancel(request).await { + Ok(response) => { + trace_coordinator_latency_ok( + "store_cancel_run", + &scope, + Some(response.run_id), + started_at, + ); + response + } + Err(error) => { + trace_coordinator_latency_error( + "store_cancel_run", + &scope, + None, + started_at, + &error, + ); + return Err(error); + } + }; // Wake on `CancelRequested` (the cooperative case) AND on any terminal // transition. Registered handles otherwise rely solely on the polling // fallback to discover a direct-to-terminal cancellation, which would // leave them in the requester map until the polling task next ticks. if response.status == TurnStatus::CancelRequested || response.status.is_terminal() { + let wake_started_at = live_latency_started_at(); notify_queued_run_best_effort( self.wake_notifier.as_ref(), cancel_wake(scope.clone(), &response), ); + trace_coordinator_latency_ok( + "notify_cancelled_run", + &scope, + Some(response.run_id), + wake_started_at, + ); } if !response.already_terminal { deferred_publication_error( @@ -352,18 +482,49 @@ where &self, request: SubmitChildRunRequest, ) -> Result { + let started_at = live_latency_started_at(); let child_scope = request.child_scope.clone(); - let response = self + let response = match self .store .submit_child_turn( request, self.admission_policy.as_ref(), self.run_profile_resolver.as_ref(), ) - .await?; + .await + { + Ok(response) => { + let SubmitTurnResponse::Accepted { run_id, .. } = &response; + trace_coordinator_latency_ok( + "store_submit_child_turn", + &child_scope, + Some(*run_id), + started_at, + ); + response + } + Err(error) => { + trace_coordinator_latency_error( + "store_submit_child_turn", + &child_scope, + None, + started_at, + &error, + ); + return Err(error); + } + }; + let SubmitTurnResponse::Accepted { run_id, .. } = &response; + let wake_started_at = live_latency_started_at(); notify_queued_run_best_effort( self.wake_notifier.as_ref(), - submit_wake(child_scope, &response), + submit_wake(child_scope.clone(), &response), + ); + trace_coordinator_latency_ok( + "notify_child_run", + &child_scope, + Some(*run_id), + wake_started_at, ); deferred_publication_error( self.publication_error_port.as_ref(), diff --git a/docs/internal/live-latency-instrumentation.md b/docs/internal/live-latency-instrumentation.md new file mode 100644 index 00000000000..acdb0b65e04 --- /dev/null +++ b/docs/internal/live-latency-instrumentation.md @@ -0,0 +1,51 @@ +# Live Latency Instrumentation + +This branch adds low-level `TRACE` events under the `ironclaw_latency` target for +debugging live instances that stall together. The events are intended for short +diagnostic windows on affected instances, not for always-on production logging. + +Enable the latency trace target alongside normal logs: + +```bash +RUST_LOG=info,ironclaw_latency=trace +``` + +Each latency event includes: + +- `component`: subsystem that emitted the event. +- `operation`: measured operation inside that subsystem. +- `elapsed_ms`: wall-clock duration. +- `outcome`: `ok`, `error`, or a more specific terminal status. +- Correlation fields when available: `tenant_id`, `user_id`, `thread_id`, + `run_id`, `invocation_id`, `capability_id`, and provider identifiers. + +Filesystem events intentionally log only `path_class`, such as `/threads`, +`/turns`, `/resources`, or `/runs`, rather than full virtual paths. + +## Reading Stalls + +Use the correlated events to locate where requests are waiting: + +- `reborn_runtime submit_user_turn` or `accept_inbound_message` is slow, and + filesystem `/threads` operations are also slow: thread persistence is likely + on the hot path. +- `turn_coordinator store_submit_turn` is fast but + `turn_scheduler claim_next_run` or `execute_claimed_run` is delayed: scheduler + wake/claim or worker availability is suspect. +- `reborn_turn_executor driver_run` is slow, but model gateway events are not: + the delay is before provider execution, usually host setup, turn state, or + driver/tool-loop overhead. +- `model_gateway provider_complete*` is slow: the provider request is the main + latency contributor. +- `host_runtime capability_host_invoke_json`, + `process_executor host_process_execute`, or + `runtime_dispatch_execute` is slow: capability dispatch or tool process + startup/runtime is the likely bottleneck. +- Many unrelated components pause and resolve together: look at host CPU, + blocking I/O, DB connection pool saturation, global runtime starvation, or + filesystem backend stalls. + +The first symptom described by operators is lag before the typing indicator +starts. For that path, start by comparing `send_user_message`, +`submit_user_turn`, `accept_inbound_message`, `store_submit_turn`, and the +filesystem `/threads` and `/turns` events for the same `thread_id` and `run_id`.