Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ members = [
"integrations/nvtx/example",
"integrations/nvtx/analyzer",
"integrations/nvtx/ui",
"integrations/nvtx/server",
# Domain-specific crates
"domains/query_engine/analyzer",
"domains/query_engine/model",
Expand Down
53 changes: 48 additions & 5 deletions crates/io/src/filesystem/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,10 +31,10 @@ impl TryFrom<&str> for Format {
}

impl Format {
/// Detect the format of a context directory from the first recognized
/// `*.<ext>` event stream in any of its per-entity subdirectories. Returns
/// `None` if no readable stream with a known extension is present.
/// Detect the common format of a context's event streams. Returns `None` if
/// no recognized stream is present or streams use different formats.
pub fn detect(context_dir: &std::path::Path) -> Option<Self> {
let mut detected = None;
for entry in std::fs::read_dir(context_dir).ok()?.flatten() {
if !entry.file_type().map(|t| t.is_dir()).unwrap_or(false) {
continue;
Expand All @@ -48,10 +48,53 @@ impl Format {
.and_then(|ext| ext.to_str())
.and_then(|ext| Self::try_from(ext).ok())
{
return Some(format);
match detected {
Some(existing) if existing != format => return None,
None => detected = Some(format),
_ => {}
}
}
}
}
None
detected
}
}

#[cfg(all(test, feature = "ndjson"))]
mod tests {
use std::path::PathBuf;

use super::*;
use uuid::Uuid;

fn context_with_streams(streams: &[(&str, &str)]) -> PathBuf {
let context = std::env::temp_dir().join(format!("quent-format-{}", Uuid::now_v7()));
for &(entity, extension) in streams {
let stream_dir = context.join(entity);
std::fs::create_dir_all(&stream_dir).unwrap();
std::fs::write(stream_dir.join(format!("events.{extension}")), []).unwrap();
}
context
}

#[test]
fn detects_one_context_wide_format() {
let context =
context_with_streams(&[("EngineEvent", "ndjson"), ("NvtxEventEntity", "ndjson")]);
let detected = Format::detect(&context);
std::fs::remove_dir_all(context).unwrap();

assert_eq!(detected, Some(Format::Ndjson));
}

#[test]
#[cfg(feature = "msgpack")]
fn rejects_mixed_context_formats() {
let context =
context_with_streams(&[("EngineEvent", "ndjson"), ("NvtxEventEntity", "msgpack")]);
let detected = Format::detect(&context);
std::fs::remove_dir_all(context).unwrap();

assert_eq!(detected, None);
}
}
23 changes: 21 additions & 2 deletions crates/open/src/viewer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ async fn build_one(group: ViewerGroup) -> Result<BuiltViewer> {
println!("building: {label}");

let crate_dir = build_dir(&spec)?;
wrapper::generate(&spec, &crate_dir, wrapper::IO_PACKAGE)?;
wrapper::generate(&spec, &crate_dir, wrapper::IO_PACKAGE, true)?;
let bin = match cargo_build(&crate_dir).await {
// The pinned quent revision predates the `quent-exporter` → `quent-io`
// rename (cargo found no `quent-io` package there, failing resolution
Expand All @@ -102,7 +102,26 @@ async fn build_one(group: ViewerGroup) -> Result<BuiltViewer> {
wrapper::IO_PACKAGE,
wrapper::LEGACY_IO_PACKAGE
);
wrapper::generate(&spec, &crate_dir, wrapper::LEGACY_IO_PACKAGE)?;
wrapper::generate(&spec, &crate_dir, wrapper::LEGACY_IO_PACKAGE, true)?;
cargo_build(&crate_dir).await?
}
Err(error) if missing_package(&error, "nvtx-server") => {
println!(
"note: pinned quent has no `nvtx-server` package; retrying without NVTX routes"
);
wrapper::generate(&spec, &crate_dir, wrapper::IO_PACKAGE, false)?;
cargo_build(&crate_dir).await?
}
Err(error)
if missing_package(&error, wrapper::IO_PACKAGE)
|| missing_package(&error, wrapper::LEGACY_IO_PACKAGE) =>
{
let io_package = if missing_package(&error, wrapper::LEGACY_IO_PACKAGE) {
wrapper::LEGACY_IO_PACKAGE
} else {
wrapper::IO_PACKAGE
};
wrapper::generate(&spec, &crate_dir, io_package, false)?;
cargo_build(&crate_dir).await?
}
result => result?,
Expand Down
148 changes: 104 additions & 44 deletions crates/open/src/wrapper.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,10 +35,21 @@ pub const ADDR_ENV: &str = "QUENT_OPEN_ADDR";
/// Write the wrapper crate (`Cargo.toml` + `src/main.rs`) into `crate_dir`.
/// `io_package` is the name of quent's I/O crate at the pinned revision
/// ([`IO_PACKAGE`], or [`LEGACY_IO_PACKAGE`] for revisions predating the rename).
pub fn generate(spec: &ViewerSpec, crate_dir: &Path, io_package: &str) -> Result<()> {
pub fn generate(
spec: &ViewerSpec,
crate_dir: &Path,
io_package: &str,
with_nvtx_routes: bool,
) -> Result<()> {
std::fs::create_dir_all(crate_dir.join("src"))?;
std::fs::write(crate_dir.join("Cargo.toml"), cargo_toml(spec, io_package))?;
std::fs::write(crate_dir.join("src/main.rs"), main_rs(spec))?;
std::fs::write(
crate_dir.join("Cargo.toml"),
cargo_toml(spec, io_package, with_nvtx_routes),
)?;
std::fs::write(
crate_dir.join("src/main.rs"),
main_rs(spec, with_nvtx_routes),
)?;
Ok(())
}

Expand All @@ -55,10 +66,10 @@ fn git_dep(url: String, rev: &str, features: &[&str]) -> Dependency {
/// Wrapper `Cargo.toml`, built with `cargo-manifest`: pin quent crates to
/// `quent.{remote,commit}` and the analyzer to `analyzer.{remote,commit}`; the
/// empty `[workspace]` keeps the generated crate out of any parent workspace.
fn cargo_toml(spec: &ViewerSpec, io_package: &str) -> String {
fn cargo_toml(spec: &ViewerSpec, io_package: &str, with_nvtx_routes: bool) -> String {
let quent = spec.quent.cargo_url();
let q_rev = spec.quent.commit.as_str();
let dependencies = BTreeMap::from([
let mut dependencies = BTreeMap::from([
(
"quent-query-engine-server".to_string(),
git_dep(quent.clone(), q_rev, &["ui"]),
Expand All @@ -70,7 +81,7 @@ fn cargo_toml(spec: &ViewerSpec, io_package: &str) -> String {
(
// All formats enabled so the analyzer can detect the artifact's format at runtime.
io_package.to_string(),
git_dep(quent, q_rev, &["ndjson", "msgpack", "postcard"]),
git_dep(quent.clone(), q_rev, &["ndjson", "msgpack", "postcard"]),
),
(
spec.analyzer_package.clone(),
Expand All @@ -92,6 +103,9 @@ fn cargo_toml(spec: &ViewerSpec, io_package: &str) -> String {
),
("uuid".to_string(), Dependency::Simple("1".to_string())),
]);
if with_nvtx_routes {
dependencies.insert("nvtx-server".to_string(), git_dep(quent, q_rev, &[]));
}

let mut package = Package::new(WRAPPER_PACKAGE.to_string(), "0.0.0".to_string());
package.edition = Some(MaybeInherited::Local(Edition::E2024));
Expand All @@ -115,43 +129,85 @@ fn cargo_toml(spec: &ViewerSpec, io_package: &str) -> String {
/// Wrapper `src/main.rs`: wire `<analyzer>::Viewer`'s analyzer/importer into
/// `analyzer_service_router` and serve it. Root (`<context-uuid>/` subdirs) and
/// bind address come from env so one built binary serves any artifacts.
fn main_rs(spec: &ViewerSpec) -> String {
fn main_rs(spec: &ViewerSpec, with_nvtx_routes: bool) -> String {
let analyzer_crate = format_ident!("{}", spec.analyzer_crate());
let (root_env, addr_env) = (ROOT_ENV, ADDR_ENV);
let tokens = quote! {
use std::net::SocketAddr;
use std::path::PathBuf;

use quent_query_engine_analyzer::ui::QuentViewer;
use quent_query_engine_server::analyzer_cache::index_query_engines;
use quent_query_engine_server::analyzer_service_router;
use #analyzer_crate::Viewer;

type Analyzer = <Viewer as QuentViewer>::Analyzer;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let root = PathBuf::from(std::env::var(#root_env)?);
let addr: SocketAddr = std::env::var(#addr_env)?.parse()?;

let import_root = root.clone();
let importer = move |id: uuid::Uuid| {
Ok(<Viewer as QuentViewer>::import_events(
&import_root.join(id.to_string()),
)?)
};
let lister_root = root.clone();
let lister = move || index_query_engines(&lister_root);

let router = analyzer_service_router::<Analyzer>(
Box::new(importer),
Box::new(lister),
None,
)?;

let listener = tokio::net::TcpListener::bind(addr).await?;
axum::serve(listener, router.into_make_service()).await?;
Ok(())
let tokens = if with_nvtx_routes {
quote! {
use std::net::SocketAddr;
use std::path::PathBuf;

use quent_query_engine_analyzer::ui::QuentViewer;
use quent_query_engine_server::analyzer_cache::index_query_engines;
use quent_query_engine_server::analyzer_service_router_with_routes;
use nvtx_server::{import_context_events, routes as nvtx_routes};
use #analyzer_crate::Viewer;

type Analyzer = <Viewer as QuentViewer>::Analyzer;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let root = PathBuf::from(std::env::var(#root_env)?);
let addr: SocketAddr = std::env::var(#addr_env)?.parse()?;

let import_root = root.clone();
let importer = move |id: uuid::Uuid| {
Ok(<Viewer as QuentViewer>::import_events(
&import_root.join(id.to_string()),
&import_root.join(id.to_string()),
)?)
};
let lister_root = root.clone();
let lister = move || index_query_engines(&lister_root);
let nvtx_root = root.clone();
let nvtx_importer = move |id: uuid::Uuid| import_context_events(&nvtx_root, id);

let router = analyzer_service_router_with_routes::<Analyzer>(
Box::new(importer),
Box::new(lister),
None,
nvtx_routes(Box::new(nvtx_importer)),
)?;

let listener = tokio::net::TcpListener::bind(addr).await?;
axum::serve(listener, router.into_make_service()).await?;
Ok(())
}
}
} else {
quote! {
use std::net::SocketAddr;
use std::path::PathBuf;

use quent_query_engine_analyzer::ui::QuentViewer;
use quent_query_engine_server::analyzer_cache::index_query_engines;
use quent_query_engine_server::analyzer_service_router;
use #analyzer_crate::Viewer;

type Analyzer = <Viewer as QuentViewer>::Analyzer;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let root = PathBuf::from(std::env::var(#root_env)?);
let addr: SocketAddr = std::env::var(#addr_env)?.parse()?;

let import_root = root.clone();
let importer = move |id: uuid::Uuid| {
Ok(<Viewer as QuentViewer>::import_events(&import_root.join(id.to_string()))?)
};
let lister_root = root.clone();
let lister = move || index_query_engines(&lister_root);

let router = analyzer_service_router::<Analyzer>(
Box::new(importer),
Box::new(lister),
None,
)?;

let listener = tokio::net::TcpListener::bind(addr).await?;
axum::serve(listener, router.into_make_service()).await?;
Ok(())
}
}
};
let file = syn::parse2(tokens).expect("generated wrapper main.rs is valid Rust");
Expand Down Expand Up @@ -182,13 +238,14 @@ mod tests {

#[test]
fn cargo_toml_pins_quent_and_analyzer() {
let manifest: toml::Value = toml::from_str(&cargo_toml(&spec(), IO_PACKAGE)).unwrap();
let manifest: toml::Value = toml::from_str(&cargo_toml(&spec(), IO_PACKAGE, true)).unwrap();
assert!(manifest.get("workspace").is_some(), "standalone workspace");
let deps = &manifest["dependencies"];
let server = &deps["quent-query-engine-server"];
assert_eq!(server["git"].as_str().unwrap(), "https://example.com/quent");
assert_eq!(server["rev"].as_str().unwrap(), "quentcommit");
assert_eq!(server["features"][0].as_str().unwrap(), "ui");
assert_eq!(deps["nvtx-server"]["rev"].as_str().unwrap(), "quentcommit");
// The exporter enables all formats so the analyzer detects the artifact's format at runtime.
let exporter_features = deps["quent-io"]["features"].as_array().unwrap();
for format in ["ndjson", "msgpack", "postcard"] {
Expand All @@ -207,10 +264,11 @@ mod tests {
// Artifacts pinned to quent revisions predating the `quent-exporter` →
// `quent-io` rename depend on the legacy package instead, same features.
let manifest: toml::Value =
toml::from_str(&cargo_toml(&spec(), LEGACY_IO_PACKAGE)).unwrap();
toml::from_str(&cargo_toml(&spec(), LEGACY_IO_PACKAGE, false)).unwrap();
let deps = &manifest["dependencies"];
assert!(deps.get("quent-io").is_none());
let exporter = &deps["quent-exporter"];
assert!(deps.get("nvtx-server").is_none());
assert_eq!(
exporter["git"].as_str().unwrap(),
"https://example.com/quent"
Expand All @@ -224,9 +282,11 @@ mod tests {

#[test]
fn main_rs_wires_the_viewer() {
let main = main_rs(&spec());
let main = main_rs(&spec(), true);
assert!(main.contains("use quent_simulator_analyzer::Viewer;"));
assert!(main.contains("import_events"));
assert!(main.contains("import_context_events"));
assert!(main.contains("analyzer_service_router_with_routes"));
assert!(main.contains("QUENT_OPEN_ADDR")); // bind address is configurable
}
}
Loading