diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..5931875 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,3 @@ +target/ +.git/ +.env diff --git a/.gitignore b/.gitignore index 7ee12e8..8002557 100644 --- a/.gitignore +++ b/.gitignore @@ -5,3 +5,4 @@ __pycache__/ node_modules/ e2e/test-results/ e2e/playwright-report/ +.kokoro-models/ diff --git a/Cargo.lock b/Cargo.lock index 022e0d5..99121ae 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -76,6 +76,28 @@ version = "1.0.102" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c" +[[package]] +name = "async-stream" +version = "0.3.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0b5a71a6f37880a80d1d7f19efd781e4b5de42c88f0722cc13bcb6cc2cfe8476" +dependencies = [ + "async-stream-impl", + "futures-core", + "pin-project-lite", +] + +[[package]] +name = "async-stream-impl" +version = "0.3.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7c24de15d275a1ecfd47a380fb4d5ec9bfe0933f309ed5e705b775596a3574d" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "async-trait" version = "0.1.89" @@ -99,13 +121,40 @@ version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" +[[package]] +name = "axum" +version = "0.7.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "edca88bc138befd0323b20752846e6587272d3b03b0343c8ea28a6f819e6e71f" +dependencies = [ + "async-trait", + "axum-core 0.4.5", + "bytes", + "futures-util", + "http", + "http-body", + "http-body-util", + "itoa", + "matchit 0.7.3", + "memchr", + "mime", + "percent-encoding", + "pin-project-lite", + "rustversion", + "serde", + "sync_wrapper", + "tower 0.5.3", + "tower-layer", + "tower-service", +] + [[package]] name = "axum" version = "0.8.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90" dependencies = [ - "axum-core", + "axum-core 0.5.6", "base64", "bytes", "form_urlencoded", @@ -116,7 +165,7 @@ dependencies = [ "hyper", "hyper-util", "itoa", - "matchit", + "matchit 0.8.4", "memchr", "mime", "percent-encoding", @@ -129,12 +178,32 @@ dependencies = [ "sync_wrapper", "tokio", "tokio-tungstenite", - "tower", + "tower 0.5.3", "tower-layer", "tower-service", "tracing", ] +[[package]] +name = "axum-core" +version = "0.4.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09f2bd6146b97ae3359fa0cc6d6b376d9539582c7b4220f041a33ec24c226199" +dependencies = [ + "async-trait", + "bytes", + "futures-util", + "http", + "http-body", + "http-body-util", + "mime", + "pin-project-lite", + "rustversion", + "sync_wrapper", + "tower-layer", + "tower-service", +] + [[package]] name = "axum-core" version = "0.5.6" @@ -160,8 +229,8 @@ version = "0.10.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9963ff19f40c6102c76756ef0a46004c0d58957d87259fc9208ff8441c12ab96" dependencies = [ - "axum", - "axum-core", + "axum 0.8.9", + "axum-core 0.5.6", "bytes", "futures-util", "headers", @@ -399,6 +468,12 @@ dependencies = [ "syn", ] +[[package]] +name = "either" +version = "1.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91622ff5e7162018101f2fea40d6ebf4a78bbe5a49736a2020649edf9693679e" + [[package]] name = "encoding_rs" version = "0.8.35" @@ -479,6 +554,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "07bbe89c50d7a535e539b8c17bc0b49bdb77747034daa8087407d655f3f7cc1d" dependencies = [ "futures-core", + "futures-sink", ] [[package]] @@ -487,6 +563,34 @@ version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7e3450815272ef58cec6d564423f6e755e25379b217b0bc688e295ba24df6b1d" +[[package]] +name = "futures-executor" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "baf29c38818342a3b26b5b923639e7b1f4a61fc5e76102d4b1981c6dc7a7579d" +dependencies = [ + "futures-core", + "futures-task", + "futures-util", +] + +[[package]] +name = "futures-io" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cecba35d7ad927e23624b22ad55235f2239cfa44fd10428eecbeba6d6a717718" + +[[package]] +name = "futures-macro" +version = "0.3.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "futures-sink" version = "0.3.32" @@ -506,8 +610,11 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" dependencies = [ "futures-core", + "futures-io", + "futures-macro", "futures-sink", "futures-task", + "memchr", "pin-project-lite", "slab", ] @@ -558,6 +665,12 @@ dependencies = [ "wasip3", ] +[[package]] +name = "glob" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0cc23270f6e1808e30a928bdc84dea0b9b4136a8bc82338574f23baf47bbd280" + [[package]] name = "h2" version = "0.4.14" @@ -570,13 +683,19 @@ dependencies = [ "futures-core", "futures-sink", "http", - "indexmap", + "indexmap 2.14.0", "slab", "tokio", "tokio-util", "tracing", ] +[[package]] +name = "hashbrown" +version = "0.12.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" + [[package]] name = "hashbrown" version = "0.14.5" @@ -710,6 +829,19 @@ dependencies = [ "tower-service", ] +[[package]] +name = "hyper-timeout" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b90d566bffbce6a75bd8b09a05aa8c2cb1fabb6cb348f8840c9e4c90a0d83b0" +dependencies = [ + "hyper", + "hyper-util", + "pin-project-lite", + "tokio", + "tower-service", +] + [[package]] name = "hyper-tls" version = "0.6.0" @@ -743,7 +875,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2", + "socket2 0.6.4", "system-configuration", "tokio", "tower-service", @@ -884,6 +1016,16 @@ dependencies = [ "icu_properties", ] +[[package]] +name = "indexmap" +version = "1.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bd070e393353796e801d209ad339e89596eb4c8d430d18ede6a1cced8fafbd99" +dependencies = [ + "autocfg", + "hashbrown 0.12.3", +] + [[package]] name = "indexmap" version = "2.14.0" @@ -908,6 +1050,15 @@ version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" +[[package]] +name = "itertools" +version = "0.14.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b192c782037fadd9cfa75548310488aabdbf3d2da73885b31bd0abd03351285" +dependencies = [ + "either", +] + [[package]] name = "itoa" version = "1.0.18" @@ -989,6 +1140,12 @@ dependencies = [ "regex-automata", ] +[[package]] +name = "matchit" +version = "0.7.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0e7465ac9959cc2b1404e8e2367b43684a6d13790fe23056cc8c6c5a6b7bcb94" + [[package]] name = "matchit" version = "0.8.4" @@ -1036,8 +1193,13 @@ dependencies = [ "clap", "mofa-engine-core", "mofa-engine-sdk", + "mofa-observability", + "opentelemetry", + "opentelemetry-otlp", + "opentelemetry_sdk", "tokio", "tracing", + "tracing-opentelemetry", "tracing-subscriber", ] @@ -1067,17 +1229,18 @@ dependencies = [ name = "mofa-engine-sdk" version = "0.1.0" dependencies = [ - "axum", + "axum 0.8.9", "axum-extra", "chrono", "mofa-engine-core", "mofa-kernel", + "mofa-observability", "serde", "serde_json", "thiserror", "tokio", "tokio-stream", - "tower", + "tower 0.5.3", "tower-http", "tracing", ] @@ -1093,6 +1256,19 @@ dependencies = [ "uuid", ] +[[package]] +name = "mofa-observability" +version = "0.1.0" +dependencies = [ + "axum 0.8.9", + "rand 0.8.6", + "serde", + "serde_json", + "tokio", + "tracing", + "tracing-subscriber", +] + [[package]] name = "native-tls" version = "0.2.18" @@ -1201,6 +1377,88 @@ dependencies = [ "vcpkg", ] +[[package]] +name = "opentelemetry" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "236e667b670a5cdf90c258f5a55794ec5ac5027e960c224bff8367a59e1e6426" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "pin-project-lite", + "thiserror", + "tracing", +] + +[[package]] +name = "opentelemetry-http" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a8863faf2910030d139fb48715ad5ff2f35029fc5f244f6d5f689ddcf4d26253" +dependencies = [ + "async-trait", + "bytes", + "http", + "opentelemetry", + "reqwest", + "tracing", +] + +[[package]] +name = "opentelemetry-otlp" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5bef114c6d41bea83d6dc60eb41720eedd0261a67af57b66dd2b84ac46c01d91" +dependencies = [ + "async-trait", + "futures-core", + "http", + "opentelemetry", + "opentelemetry-http", + "opentelemetry-proto", + "opentelemetry_sdk", + "prost", + "reqwest", + "thiserror", + "tokio", + "tonic", + "tracing", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56f8870d3024727e99212eb3bb1762ec16e255e3e6f58eeb3dc8db1aa226746d" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost", + "tonic", +] + +[[package]] +name = "opentelemetry_sdk" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "84dfad6042089c7fc1f6118b7040dc2eb4ab520abbf410b79dc481032af39570" +dependencies = [ + "async-trait", + "futures-channel", + "futures-executor", + "futures-util", + "glob", + "opentelemetry", + "percent-encoding", + "rand 0.8.6", + "serde_json", + "thiserror", + "tokio", + "tokio-stream", + "tracing", +] + [[package]] name = "option-ext" version = "0.2.0" @@ -1236,6 +1494,26 @@ version = "2.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" +[[package]] +name = "pin-project" +version = "1.1.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2466b2336ed02bcdca6b294417127b90ec92038d1d5c4fbeac971a922e0e0924" +dependencies = [ + "pin-project-internal", +] + +[[package]] +name = "pin-project-internal" +version = "1.1.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c96395f0a926bc13b1c17622aaddda1ecb55d49c8f1bf9777e4d877800a43f8b" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "pin-project-lite" version = "0.2.17" @@ -1285,6 +1563,29 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "prost" +version = "0.13.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2796faa41db3ec313a31f7624d9286acf277b52de526150b7e69f3debf891ee5" +dependencies = [ + "bytes", + "prost-derive", +] + +[[package]] +name = "prost-derive" +version = "0.13.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d" +dependencies = [ + "anyhow", + "itertools", + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "quote" version = "1.0.45" @@ -1306,14 +1607,35 @@ version = "6.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" +[[package]] +name = "rand" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5ca0ecfa931c29007047d1bc58e623ab12e5590e8c7cc53200d5202b69266d8a" +dependencies = [ + "libc", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + [[package]] name = "rand" version = "0.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "44c5af06bb1b7d3216d91932aed5265164bf384dc89cd6ba05cf59a35f5f76ea" dependencies = [ - "rand_chacha", - "rand_core", + "rand_chacha 0.9.0", + "rand_core 0.9.5", +] + +[[package]] +name = "rand_chacha" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" +dependencies = [ + "ppv-lite86", + "rand_core 0.6.4", ] [[package]] @@ -1323,7 +1645,16 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.9.5", +] + +[[package]] +name = "rand_core" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" +dependencies = [ + "getrandom 0.2.17", ] [[package]] @@ -1381,6 +1712,7 @@ dependencies = [ "base64", "bytes", "encoding_rs", + "futures-channel", "futures-core", "futures-util", "h2", @@ -1405,7 +1737,7 @@ dependencies = [ "sync_wrapper", "tokio", "tokio-native-tls", - "tower", + "tower 0.5.3", "tower-http", "tower-service", "url", @@ -1653,6 +1985,16 @@ version = "1.15.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "67b1b7a3b5fe4f1376887184045fcf45c69e92af734b7aaddc05fb777b6fbd03" +[[package]] +name = "socket2" +version = "0.5.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e22376abed350d73dd1cd119b57ffccad95b4e585a7cda43e286245ce23c0678" +dependencies = [ + "libc", + "windows-sys 0.52.0", +] + [[package]] name = "socket2" version = "0.6.4" @@ -1810,7 +2152,7 @@ dependencies = [ "parking_lot", "pin-project-lite", "signal-hook-registry", - "socket2", + "socket2 0.6.4", "tokio-macros", "windows-sys 0.61.2", ] @@ -1910,7 +2252,7 @@ version = "0.22.27" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "41fe8c660ae4257887cf66394862d21dbca4a6ddd26f04a3560410406a2f819a" dependencies = [ - "indexmap", + "indexmap 2.14.0", "serde", "serde_spanned", "toml_datetime", @@ -1924,6 +2266,56 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5d99f8c9a7727884afe522e9bd5edbfc91a3312b36a77b5fb8926e4c31a41801" +[[package]] +name = "tonic" +version = "0.12.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "877c5b330756d856ffcc4553ab34a5684481ade925ecc54bcd1bf02b1d0d4d52" +dependencies = [ + "async-stream", + "async-trait", + "axum 0.7.9", + "base64", + "bytes", + "h2", + "http", + "http-body", + "http-body-util", + "hyper", + "hyper-timeout", + "hyper-util", + "percent-encoding", + "pin-project", + "prost", + "socket2 0.5.10", + "tokio", + "tokio-stream", + "tower 0.4.13", + "tower-layer", + "tower-service", + "tracing", +] + +[[package]] +name = "tower" +version = "0.4.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8fa9be0de6cf49e536ce1851f987bd21a43b771b09473c3549a6c853db37c1c" +dependencies = [ + "futures-core", + "futures-util", + "indexmap 1.9.3", + "pin-project", + "pin-project-lite", + "rand 0.8.6", + "slab", + "tokio", + "tokio-util", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "tower" version = "0.5.3" @@ -1952,7 +2344,7 @@ dependencies = [ "http", "http-body", "pin-project-lite", - "tower", + "tower 0.5.3", "tower-layer", "tower-service", "tracing", @@ -2015,6 +2407,24 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-opentelemetry" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "721f2d2569dce9f3dfbbddee5906941e953bfcdf736a62da3377f5751650cc36" +dependencies = [ + "js-sys", + "once_cell", + "opentelemetry", + "opentelemetry_sdk", + "smallvec", + "tracing", + "tracing-core", + "tracing-log", + "tracing-subscriber", + "web-time", +] + [[package]] name = "tracing-serde" version = "0.2.0" @@ -2063,7 +2473,7 @@ dependencies = [ "http", "httparse", "log", - "rand", + "rand 0.9.4", "sha1", "thiserror", ] @@ -2256,7 +2666,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "bb0e353e6a2fbdc176932bbaab493762eb1255a7900fe0fea1a2f96c296cc909" dependencies = [ "anyhow", - "indexmap", + "indexmap 2.14.0", "wasm-encoder", "wasmparser", ] @@ -2269,7 +2679,7 @@ checksum = "47b807c72e1bac69382b3a6fb3dbe8ea4c0ed87ff5629b8685ae6b9a611028fe" dependencies = [ "bitflags", "hashbrown 0.15.5", - "indexmap", + "indexmap 2.14.0", "semver", ] @@ -2283,6 +2693,16 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "winapi" version = "0.3.9" @@ -2553,7 +2973,7 @@ checksum = "b7c566e0f4b284dd6561c786d9cb0142da491f46a9fbed79ea69cdad5db17f21" dependencies = [ "anyhow", "heck", - "indexmap", + "indexmap 2.14.0", "prettyplease", "syn", "wasm-metadata", @@ -2584,7 +3004,7 @@ checksum = "9d66ea20e9553b30172b5e831994e35fbde2d165325bec84fc43dbf6f4eb9cb2" dependencies = [ "anyhow", "bitflags", - "indexmap", + "indexmap 2.14.0", "log", "serde", "serde_derive", @@ -2603,7 +3023,7 @@ checksum = "ecc8ac4bc1dc3381b7f59c34f00b67e18f910c2c0f50015669dde7def656a736" dependencies = [ "anyhow", "id-arena", - "indexmap", + "indexmap 2.14.0", "log", "semver", "serde", diff --git a/Cargo.toml b/Cargo.toml index 6e7adfc..a7a9867 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,6 +5,7 @@ members = [ "mofa-engine-core", "mofa-engine-sdk", "mofa-engine-app", + "mofa-observability", ] [workspace.package] @@ -36,6 +37,10 @@ anyhow = "1" # Logging tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] } +opentelemetry = "0.28.0" +opentelemetry_sdk = { version = "0.28.0", features = ["rt-tokio"] } +tracing-opentelemetry = "0.29.0" +opentelemetry-otlp = { version = "0.28.0", features = ["grpc-tonic", "trace"] } # Utilities uuid = { version = "1", features = ["v4"] } diff --git a/mofa-engine-app/Cargo.toml b/mofa-engine-app/Cargo.toml index 6e7bc16..9bd22b6 100644 --- a/mofa-engine-app/Cargo.toml +++ b/mofa-engine-app/Cargo.toml @@ -8,8 +8,13 @@ description = "MoFA Engine binary entry point" [dependencies] mofa-engine-core = { workspace = true } mofa-engine-sdk = { workspace = true } +mofa-observability = { path = "../mofa-observability" } tokio = { workspace = true } clap = { workspace = true } tracing = { workspace = true } tracing-subscriber = { workspace = true } +tracing-opentelemetry = { workspace = true } +opentelemetry = { workspace = true } +opentelemetry_sdk = { workspace = true } +opentelemetry-otlp = { workspace = true } anyhow = { workspace = true } diff --git a/mofa-engine-app/src/main.rs b/mofa-engine-app/src/main.rs index 2f163c8..9262314 100644 --- a/mofa-engine-app/src/main.rs +++ b/mofa-engine-app/src/main.rs @@ -6,8 +6,31 @@ use clap::Parser; use mofa_engine_core::{Engine, EngineConfig}; use mofa_engine_sdk::start_server; +use opentelemetry_otlp::WithExportConfig; +use opentelemetry_sdk::Resource; +use opentelemetry_sdk::trace::SdkTracerProvider; use std::path::PathBuf; -use tracing_subscriber::{EnvFilter, fmt}; +use std::sync::Arc; +use tracing_subscriber::{EnvFilter, fmt, layer::SubscriberExt, util::SubscriberInitExt}; + +fn init_tracer( + endpoint: &str, +) -> Result { + let exporter = opentelemetry_otlp::SpanExporter::builder() + .with_tonic() + .with_endpoint(endpoint) + .build()?; + + let provider = SdkTracerProvider::builder() + .with_batch_exporter(exporter) + .with_resource(Resource::builder().with_service_name("mofa-engine").build()) + .build(); + + opentelemetry::global::set_tracer_provider(provider.clone()); + + use opentelemetry::trace::TracerProvider; + Ok(provider.tracer("mofa-engine")) +} /// MoFA Engine — multimodal AI model orchestration #[derive(Parser, Debug)] @@ -24,14 +47,6 @@ struct Cli { #[tokio::main] async fn main() -> anyhow::Result<()> { - // Initialise tracing - fmt() - .with_env_filter( - EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")), - ) - .with_target(false) - .init(); - let cli = Cli::parse(); // Load and validate configuration. @@ -45,13 +60,66 @@ async fn main() -> anyhow::Result<()> { let host = config.listen.host.clone(); let port = config.listen.port; + #[allow(unused_variables)] + let observability_enabled = config.observability.enabled; + #[allow(unused_variables)] + let otlp_endpoint = config.observability.otlp_endpoint.clone(); + + // Initialize Tracing Subscriber + let env_filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")); + let fmt_layer = fmt::layer().with_target(false); + + let otel_layer = if observability_enabled { + if let Some(endpoint) = &otlp_endpoint { + match init_tracer(endpoint) { + Ok(tracer) => Some(tracing_opentelemetry::layer().with_tracer(tracer)), + Err(e) => { + eprintln!("Failed to initialize OTLP tracer: {e}"); + None + } + } + } else { + None + } + } else { + None + }; + + tracing_subscriber::registry() + .with(env_filter) + .with(fmt_layer) + .with(otel_layer) + .init(); + + let mut metrics_state = None; + let mut obs_sender = None; + + if observability_enabled { + tracing::info!("Initializing observability pipeline..."); + let (metrics, sender) = + mofa_observability::collector::init_pipeline(otlp_endpoint.as_deref()); + metrics_state = Some(metrics); + obs_sender = Some(sender); + } + tracing::info!("MoFA Engine v{} starting", env!("CARGO_PKG_VERSION")); // Create engine let engine = Engine::try_new(config).await?; + // Start Observability Bridge + if let Some(sender) = obs_sender { + let rx = engine.subscribe_events(); + let engine_clone = Arc::clone(&engine); + let metrics_clone = metrics_state.clone(); + tokio::spawn(async move { + mofa_engine_sdk::observability_bridge::run(rx, sender, engine_clone, metrics_clone) + .await; + }); + } + // Start server - start_server(engine, &host, port) + start_server(engine, metrics_state, &host, port) .await .map_err(|e| anyhow::anyhow!("{e}"))?; diff --git a/mofa-engine-core/src/backends/openai_compat.rs b/mofa-engine-core/src/backends/openai_compat.rs index 60a1ba9..e2dd8ff 100644 --- a/mofa-engine-core/src/backends/openai_compat.rs +++ b/mofa-engine-core/src/backends/openai_compat.rs @@ -385,15 +385,15 @@ impl OpenAiCompatProvider { detail: format!("TTS read error: {e}"), })?; - let path = std::env::temp_dir().join(format!("mofa_tts_{}.mp3", uuid::Uuid::new_v4())); + let file_name = format!("mofa_tts_{}.wav", uuid::Uuid::new_v4()); + let path = std::env::temp_dir().join(&file_name); std::fs::write(&path, &bytes) .map_err(|e| EngineError::Internal(format!("write error: {e}")))?; - let path = path.to_string_lossy().to_string(); let duration_ms = start.elapsed().as_millis() as u64; Ok(InferenceResponse { text: None, - file: Some(path), + file: Some(file_name), model_used: model_name.to_string(), provider: self.name.clone(), duration_ms, diff --git a/mofa-engine-core/src/config.rs b/mofa-engine-core/src/config.rs index a31c82c..f502ed6 100644 --- a/mofa-engine-core/src/config.rs +++ b/mofa-engine-core/src/config.rs @@ -22,6 +22,20 @@ pub struct EngineConfig { /// Provider definitions #[serde(default)] pub providers: Vec, + /// Observability settings + #[serde(default)] + pub observability: ObservabilityConfig, +} + +/// Observability configuration. +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +pub struct ObservabilityConfig { + /// Whether observability is enabled. + #[serde(default)] + pub enabled: bool, + /// The OpenTelemetry OTLP endpoint (e.g., http://localhost:4317) + #[serde(default)] + pub otlp_endpoint: Option, } /// Network listen configuration. @@ -399,6 +413,7 @@ impl EngineConfig { listen: ListenConfig::default(), memory: MemoryConfig::default(), providers, + observability: ObservabilityConfig::default(), } } } @@ -470,6 +485,7 @@ mod tests { #[test] fn toml_roundtrip() { let cfg = EngineConfig { + observability: Default::default(), listen: ListenConfig::default(), memory: MemoryConfig::default(), providers: vec![ProviderConfig { @@ -492,6 +508,7 @@ mod tests { #[test] fn validate_rejects_unknown_provider_kind() { let cfg = EngineConfig { + observability: Default::default(), listen: ListenConfig::default(), memory: MemoryConfig::default(), providers: vec![ProviderConfig { diff --git a/mofa-engine-core/src/engine.rs b/mofa-engine-core/src/engine.rs index e656d6c..704fb7a 100644 --- a/mofa-engine-core/src/engine.rs +++ b/mofa-engine-core/src/engine.rs @@ -172,9 +172,27 @@ impl Engine { } let count = cards.len(); + let mut allocated_any = false; for mut card in cards { card.refresh_status(); - self.models.insert(card.id.clone(), card); + let id = card.id.clone(); + if matches!(card.residency, ModelResidency::Loaded) + && card.memory_estimate_bytes > 0 + { + // If it's already loaded but we don't have it allocated yet + let snapshot = self.memory.snapshot(); + if !snapshot.iter().any(|(mem_id, _)| mem_id == &id) { + self.memory.allocate(&id, card.memory_estimate_bytes); + allocated_any = true; + } + } + self.models.insert(id, card); + } + if allocated_any { + let _ = self.event_tx.send(EngineEvent::MemoryChanged { + used_bytes: self.memory.used_bytes(), + total_bytes: self.memory.budget_bytes(), + }); } let _ = self.event_tx.send(EngineEvent::DiscoveryCompleted { provider: name.clone(), @@ -674,6 +692,7 @@ mod tests { fn minimal_config() -> EngineConfig { EngineConfig { + observability: Default::default(), listen: ListenConfig::default(), memory: MemoryConfig { budget_mb: Some(100), @@ -849,6 +868,7 @@ mod tests { #[tokio::test] async fn disabled_provider_is_skipped() { let config = EngineConfig { + observability: Default::default(), listen: ListenConfig::default(), memory: MemoryConfig { budget_mb: Some(100), diff --git a/mofa-engine-sdk/Cargo.toml b/mofa-engine-sdk/Cargo.toml index eb785c4..9c20c37 100644 --- a/mofa-engine-sdk/Cargo.toml +++ b/mofa-engine-sdk/Cargo.toml @@ -19,3 +19,4 @@ thiserror = { workspace = true } tracing = { workspace = true } chrono = { workspace = true } tokio-stream = { workspace = true } +mofa-observability = { path = "../mofa-observability" } diff --git a/mofa-engine-sdk/src/lib.rs b/mofa-engine-sdk/src/lib.rs index df91a52..5a5eef5 100644 --- a/mofa-engine-sdk/src/lib.rs +++ b/mofa-engine-sdk/src/lib.rs @@ -4,6 +4,7 @@ //! for the MoFA Engine. pub mod dashboard; +pub mod observability_bridge; pub mod server; pub use server::start_server; diff --git a/mofa-engine-sdk/src/observability_bridge.rs b/mofa-engine-sdk/src/observability_bridge.rs new file mode 100644 index 0000000..e2a5bf1 --- /dev/null +++ b/mofa-engine-sdk/src/observability_bridge.rs @@ -0,0 +1,287 @@ +//! Observability integration bridge. +//! +//! This module handles translating core engine events into the observability subsystem's format. +//! It also seeds initial metric values from the engine's status on startup. + +use mofa_engine_core::Engine; +use mofa_kernel::EngineEvent; +use mofa_observability::collector::{EventSender, Labels, MetricsState}; +use mofa_observability::events::{ + EngineEvent as ObsEngineEvent, EventEnvelope, EvictionTriggered, ModelLoaded, ModelUnloaded, + RequestCompleted, RequestReceived, UnloadReason, +}; +use std::collections::HashMap; +use std::sync::Arc; +use tokio::sync::RwLock; +use tokio::sync::broadcast; + +struct RequestContext { + capability: mofa_observability::events::Capability, + model_id: String, + backend: String, +} + +/// Runs the observability bridge. +/// +/// Consumes engine events and translates them into observability events. +/// Seeds initial memory and model-count gauges from the engine's status snapshot. +pub async fn run( + mut engine_events: broadcast::Receiver, + sender: EventSender, + engine: Arc, + metrics_state: Option>>, +) { + let mut cache: HashMap = HashMap::new(); + + // ── Seed gauges from the engine's current state ────────────────────── + // Discovery has already run by the time the bridge starts, so the + // engine knows which models are loaded and how much memory is used. + let status = engine.status().await; + tracing::info!( + "Bridge: seeding gauges — {} models loaded, {:.1} MiB used / {:.1} MiB budget", + status.loaded_models, + status.memory_used_bytes as f64 / 1_048_576.0, + status.memory_budget_bytes as f64 / 1_048_576.0, + ); + + // Set memory gauges + let _ = sender + .send_critical(EventEnvelope::now(ObsEngineEvent::EvictionTriggered( + EvictionTriggered { + evicted_model: "_seed".into(), + memory_before_bytes: status.memory_used_bytes, + memory_after_bytes: status.memory_used_bytes, + budget_bytes: status.memory_budget_bytes, + }, + ))) + .await; + + // Set models_loaded gauge by emitting one ModelLoaded per loaded model + let caps = engine.capabilities().await; + for card in &caps { + if matches!( + card.residency, + mofa_kernel::ModelResidency::Loaded | mofa_kernel::ModelResidency::Remote + ) { + sender.send(EventEnvelope::now(ObsEngineEvent::ModelLoaded( + ModelLoaded { + model_id: card.id.clone(), + backend: card.provider.clone(), + capability: map_capability(&card.capability), + load_duration_ms: 0, + memory_bytes: card.memory_estimate_bytes, + }, + ))); + } + } + + // ── Event loop with periodic gauge sync ──────────────────────────── + let mut gauge_interval = tokio::time::interval(std::time::Duration::from_secs(10)); + gauge_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + + loop { + tokio::select! { + event_result = engine_events.recv() => { + match event_result { + Ok(event) => { + match event { + EngineEvent::RequestStarted { + request_id, + capability, + model_id, + } => { + let obs_cap = capability + .as_ref() + .map(map_capability) + .unwrap_or(mofa_observability::events::Capability::Chat); + + // Extract provider from "provider/model" format + let backend = if model_id.contains('/') { + model_id.split('/').next().unwrap_or("unknown").to_string() + } else { + "local".to_string() + }; + + cache.insert( + request_id.clone(), + RequestContext { + capability: obs_cap, + model_id: model_id.clone(), + backend: backend.clone(), + }, + ); + + let obs_event = ObsEngineEvent::RequestReceived(RequestReceived { + capability: obs_cap, + model: Some(model_id), + hint: None, + }); + let mut envelope = + EventEnvelope::now(obs_event).with_request_id(&request_id); + if let Some((trace_id, span_id)) = derive_trace_context(&request_id) { + envelope = envelope.with_trace(trace_id, span_id); + } + sender.send(envelope); + } + EngineEvent::RequestCompleted { + request_id, + duration_ms, + success, + } => { + let (obs_model_id, obs_capability, obs_backend) = cache + .remove(&request_id) + .map(|ctx| (ctx.model_id, ctx.capability, ctx.backend)) + .unwrap_or_else(|| { + ( + "unknown".into(), + mofa_observability::events::Capability::Chat, + "unknown".into(), + ) + }); + + let obs_event = + ObsEngineEvent::RequestCompleted(RequestCompleted { + model_id: obs_model_id, + backend: obs_backend, + capability: obs_capability, + duration_ms, + ttft_ms: None, + tokens_in: None, + tokens_out: None, + model_was_hot: None, + success, + error_code: None, + }); + let mut envelope = + EventEnvelope::now(obs_event).with_request_id(&request_id); + if let Some((trace_id, span_id)) = derive_trace_context(&request_id) { + envelope = envelope.with_trace(trace_id, span_id); + } + sender.send(envelope); + } + EngineEvent::ModelResidencyChanged { + model_id, + old, + new, + } => { + use mofa_kernel::ModelResidency; + match (old, new) { + // Model loaded into memory + (_, ModelResidency::Loaded) => { + let backend = if model_id.contains('/') { + model_id.split('/').next().unwrap_or("unknown").to_string() + } else { + "local".to_string() + }; + sender.send(EventEnvelope::now( + ObsEngineEvent::ModelLoaded(ModelLoaded { + model_id: model_id.clone(), + backend, + capability: + mofa_observability::events::Capability::Chat, + load_duration_ms: 0, + memory_bytes: 0, + }), + )); + } + // Model unloaded from memory + (ModelResidency::Loaded, _) => { + sender.send(EventEnvelope::now( + ObsEngineEvent::ModelUnloaded(ModelUnloaded { + model_id: model_id.clone(), + reason: UnloadReason::Explicit, + memory_freed_bytes: 0, + }), + )); + } + _ => {} // Loading, Unloading — intermediate states, skip + } + } + EngineEvent::MemoryChanged { + used_bytes, + total_bytes, + } => { + // Set the memory gauges from the authoritative engine value + let _ = sender + .send_critical(EventEnvelope::now( + ObsEngineEvent::EvictionTriggered(EvictionTriggered { + evicted_model: "_memory_sync".into(), + memory_before_bytes: used_bytes, + memory_after_bytes: used_bytes, + budget_bytes: total_bytes, + }), + )) + .await; + } + EngineEvent::DiscoveryCompleted { + provider, models, .. + } => { + tracing::info!( + provider = %provider, + models = models, + "Bridge: discovery completed" + ); + } + _ => { + // ModelStatusChanged, ProviderHealthChanged — not mapped yet + } + } + } + Err(broadcast::error::RecvError::Lagged(skipped)) => { + tracing::warn!("Observability bridge lagged, skipped {} events", skipped); + } + Err(broadcast::error::RecvError::Closed) => { + tracing::info!("Observability bridge shutting down"); + break; + } + } + } + _ = gauge_interval.tick() => { + // Periodic gauge sync: directly write memory and model-count + // gauges from the authoritative engine state every 10 seconds. + if let Some(ref ms) = metrics_state { + let status = engine.status().await; + let mut state = ms.write().await; + state.memory_used_bytes.set(Labels::new(), status.memory_used_bytes as f64); + state.memory_budget_bytes.set(Labels::new(), status.memory_budget_bytes as f64); + state.models_loaded.set(Labels::new(), status.loaded_models as f64); + + tracing::trace!( + memory_used_mib = status.memory_used_bytes as f64 / 1_048_576.0, + memory_budget_mib = status.memory_budget_bytes as f64 / 1_048_576.0, + loaded_models = status.loaded_models, + "Bridge: periodic gauge sync" + ); + } + } + } + } +} + +fn map_capability(cap: &mofa_kernel::Capability) -> mofa_observability::events::Capability { + use mofa_kernel::Capability as K; + use mofa_observability::events::Capability as O; + match cap { + K::Chat => O::Chat, + K::Tts => O::Tts, + K::Asr => O::Asr, + K::Vlm => O::Vlm, + K::ImageGen => O::ImageGen, + K::VideoGen => O::VideoGen, + K::Embedding => O::Embedding, + _ => { + tracing::warn!("Unknown capability variant, defaulting to Chat for observability"); + O::Chat + } + } +} + +fn derive_trace_context(request_id: &str) -> Option<(String, String)> { + let trace_id = request_id.replace('-', ""); + if trace_id.len() == 32 { + let span_id = trace_id[0..16].to_string(); + Some((trace_id, span_id)) + } else { + None + } +} diff --git a/mofa-engine-sdk/src/server.rs b/mofa-engine-sdk/src/server.rs index 68c2623..3071992 100644 --- a/mofa-engine-sdk/src/server.rs +++ b/mofa-engine-sdk/src/server.rs @@ -26,20 +26,24 @@ use crate::dashboard; /// Shared application state. #[derive(Clone)] +#[allow(dead_code)] struct AppState { engine: Arc, started_at: std::time::Instant, + metrics: Option>>, } /// Start the HTTP server. pub async fn start_server( engine: Arc, + metrics: Option>>, host: &str, port: u16, ) -> Result<(), Box> { let state = AppState { engine, started_at: std::time::Instant::now(), + metrics, }; let app = Router::new() @@ -47,10 +51,12 @@ pub async fn start_server( .route("/", get(dashboard_handler)) // API routes .route("/health", get(health_handler)) + .route("/metrics", get(metrics_handler)) .route("/v1/capabilities", get(capabilities_handler)) .route("/v1/invoke", post(invoke_handler)) .route("/v1/status", get(status_handler)) .route("/v1/events", get(events_handler)) + .route("/v1/files/{filename}", get(files_handler)) .route("/v1/discovery/refresh", post(refresh_handler)) .layer(CorsLayer::permissive()) .layer(TraceLayer::new_for_http()) @@ -136,10 +142,40 @@ async fn events_handler( ) } +use axum::extract::Path; +use axum::http::header; +use axum::response::IntoResponse; + +async fn files_handler(Path(filename): Path) -> Result { + if filename.contains('/') || filename.contains('\\') || filename.contains("..") { + return Err(StatusCode::BAD_REQUEST); + } + + let temp_dir = std::env::temp_dir(); + let file_path = temp_dir.join(&filename); + + if !file_path.exists() { + return Err(StatusCode::NOT_FOUND); + } + + let bytes = std::fs::read(file_path).map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; + + Ok(([(header::CONTENT_TYPE, "audio/wav")], bytes)) +} + async fn dashboard_handler() -> Html<&'static str> { Html(dashboard::DASHBOARD_HTML) } +async fn metrics_handler(State(state): State) -> Result { + if let Some(metrics_arc) = &state.metrics { + let guard = metrics_arc.read().await; + Ok(mofa_observability::prometheus::render(&guard)) + } else { + Err(axum::http::StatusCode::SERVICE_UNAVAILABLE) + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/mofa-frontend/src/observability/DualTrackView.tsx b/mofa-frontend/src/observability/DualTrackView.tsx new file mode 100644 index 0000000..6c8d159 --- /dev/null +++ b/mofa-frontend/src/observability/DualTrackView.tsx @@ -0,0 +1,224 @@ +import React, { useEffect, useState } from 'react'; +import { motion } from 'framer-motion'; +import { Card } from '../shared/Card'; +import { Badge } from '../shared/Badge'; +import { engine } from '../engine'; +import { EngineStatus } from '../engine/types'; +import { HardDrive, DollarSign, Zap, AlertTriangle, ShieldCheck, Flame, Cpu, ArrowUpRight } from 'lucide-react'; + +export function DualTrackView() { + const [status, setStatus] = useState(null); + const [localRequests, setLocalRequests] = useState(0); + const [cloudRequests, setCloudRequests] = useState(0); + const [cloudCostUsd, setCloudCostUsd] = useState(0); + const [thoughtTokens, setThoughtTokens] = useState(0); + const [quotaErrors, setQuotaErrors] = useState(0); + const [warmupHits, setWarmupHits] = useState(0); + const [evictions, setEvictions] = useState(0); + + useEffect(() => { + engine.getStatus().then(res => { + if (res.success) setStatus(res.data); + }); + + const handleEvent = (evt: any) => { + const type = evt.type; + const data = evt.data || {}; + + if (type === 'RequestCompleted') { + if (data.provider === 'ollama' || data.provider === 'kokoro' || data.provider === 'funasr' || data.locality === 'local') { + setLocalRequests(prev => prev + 1); + } else { + setCloudRequests(prev => prev + 1); + if (data.cost_usd) { + setCloudCostUsd(prev => prev + data.cost_usd); + } + } + } else if (type === 'PreflightWarmCompleted') { + if (data.success) setWarmupHits(prev => prev + 1); + } else if (type === 'ModelEvicted') { + setEvictions(prev => prev + 1); + } else if (type === 'MemoryChanged' || type === 'ModelStatusChanged' || type === 'ModelResidencyChanged') { + engine.getStatus().then(res => { + if (res.success) setStatus(res.data); + }); + } + }; + + const unsubscribe = engine.subscribeEvents(handleEvent); + return () => unsubscribe(); + }, []); + + const memUsedGb = status ? (status.memory_used_bytes / (1024 * 1024 * 1024)).toFixed(2) : '3.11'; + const memBudgetGb = status ? (status.memory_budget_bytes / (1024 * 1024 * 1024)).toFixed(1) : '8.0'; + const memPercent = status && status.memory_budget_bytes > 0 + ? Math.min(100, Math.round((status.memory_used_bytes / status.memory_budget_bytes) * 100)) + : 38.9; + + return ( +
+ {/* Header Banner */} +
+
+
+ +
+
+

Dual-Track Telemetry Moat

+

+ Real-time side-by-side comparison of local hardware performance vs cloud financial cost accumulation. +

+
+
+
+ Local-First Priority Active +
+
+ + {/* Side-by-Side Grid */} +
+ + {/* Left Column: Local Execution Track */} + + +
+
+ + Local Hardware Track (locality="local") +
+ 0.00 USD / Free +
+ + {/* VRAM Memory Gauge */} +
+
+ VRAM / RAM Footprint + {memUsedGb} / {memBudgetGb} GB ({memPercent}%) +
+
+
+
+
+ + {/* Sub-Metrics Grid */} +
+
+
+ + Preflight Warmup Hits +
+
{warmupHits}
+
0ms Cold Start Latency
+
+ +
+
+ + Local Inferences +
+
{localRequests}
+
Ollama + Kokoro TTS
+
+ +
+
+ + LRU Memory Evictions +
+
{evictions}
+
Evicted under budget pressure
+
+ +
+
+ + Privacy Classification +
+
Confidential Only
+
Zero Data Egress
+
+
+ + + + {/* Right Column: Cloud Financial Track */} + + +
+
+ + Cloud Financial Track (locality="cloud") +
+ Vendor Billing Active +
+ + {/* Financial Spend Gauge */} +
+
+ Accumulated Compute Cost (USD) + ${cloudCostUsd.toFixed(4)} +
+
+
+
+
+ + {/* Sub-Metrics Grid */} +
+
+
+ + Cloud Inferences +
+
{cloudRequests}
+
OpenAI / DeepSeek / Claude
+
+ +
+
+ + Thought Tokens +
+
{thoughtTokens}
+
DeepSeek R1 Reasoning
+
+ +
+
+ + HTTP 429 Rate Limits +
+
{quotaErrors}
+
Vendor quota errors
+
+ +
+
+ + Cost Savings Rate +
+
100% Local Free Tier
+
Saved vs pure cloud
+
+
+ + + +
+
+ ); +} diff --git a/mofa-frontend/src/observability/ObservabilityView.tsx b/mofa-frontend/src/observability/ObservabilityView.tsx new file mode 100644 index 0000000..d0d9826 --- /dev/null +++ b/mofa-frontend/src/observability/ObservabilityView.tsx @@ -0,0 +1,180 @@ +import React, { useEffect, useState } from 'react'; +import { motion } from 'framer-motion'; +import { Card } from '../shared/Card'; +import { Button } from '../shared/Button'; +import { MetricsStrip } from './MetricsStrip'; +import { DualTrackView } from './DualTrackView'; +import { Activity, LayoutDashboard, Cpu, Network, ExternalLink, Layers } from 'lucide-react'; + +const GRAFANA_URL = import.meta.env.VITE_GRAFANA_URL || 'http://localhost:3000'; + +type TabId = 'dual-track' | 'overview' | 'memory' | 'routing'; + +interface TabConfig { + id: TabId; + label: string; + icon: React.ReactNode; + dashboardPath: string; + description: string; +} + +const TABS: TabConfig[] = [ + { + id: 'dual-track', + label: 'Dual-Track Telemetry', + icon: , + dashboardPath: '', + description: 'Real-time side-by-side comparison of local hardware footprint vs cloud financial cost.' + }, + { + id: 'overview', + label: 'Engine Overview', + icon: , + dashboardPath: '/d/engine-overview/mofa-engine-overview', + description: 'Request rate, error rate, P95 latency, TTFT, and tokens/sec.' + }, + { + id: 'memory', + label: 'Memory & Lifecycle', + icon: , + dashboardPath: '/d/engine-memory/mofa-memory-and-lifecycle', + description: 'Memory vs budget, model loads/unloads, eviction rate, and cold-load heatmap.' + }, + { + id: 'routing', + label: 'Preflight & Routing', + icon: , + dashboardPath: '/d/engine-routing/mofa-preflight-and-routing', + description: 'Fallback routing, circuit breakers, and capability matching.' + } +]; + +export function ObservabilityView() { + const [activeTab, setActiveTab] = useState('dual-track'); + const [grafanaAvailable, setGrafanaAvailable] = useState(null); + + useEffect(() => { + let mounted = true; + + // Lightweight check for Grafana availability using a known public static asset + const checkGrafana = () => { + const img = new Image(); + img.onload = () => { + if (mounted) setGrafanaAvailable(true); + }; + img.onerror = () => { + if (mounted) setGrafanaAvailable(false); + }; + // Prevent caching + img.src = `${GRAFANA_URL}/public/img/grafana_icon.svg?t=${Date.now()}`; + }; + + checkGrafana(); + + return () => { mounted = false; }; + }, []); + + const activeConfig = TABS.find(t => t.id === activeTab)!; + + return ( + +
+ {/* Header */} +
+
+

+ + Engine Observability +

+

+ Powered by Prometheus + Grafana & Local Telemetry Engine +

+
+ {grafanaAvailable && ( + + )} +
+ + + + {/* Tab Navigation */} +
+ {TABS.map(tab => ( + + ))} +
+ +
+ {activeTab === 'dual-track' ? ( + + ) : grafanaAvailable === null ? ( +
+
Checking observability stack...
+
+ ) : grafanaAvailable ? ( + +
+
{activeConfig.label}
+
{activeConfig.description}
+
+
+