diff --git a/.github/workflows/build-graph-node.yml b/.github/workflows/build-graph-node.yml index c8b016429e..2069cdfa67 100644 --- a/.github/workflows/build-graph-node.yml +++ b/.github/workflows/build-graph-node.yml @@ -182,9 +182,9 @@ jobs: @ruvector/graph-node-linux-arm64-gnu @ruvector/graph-node-darwin-x64 \ @ruvector/graph-node-darwin-arm64 @ruvector/graph-node-win32-x64-msvc; do ok=0 - for i in $(seq 1 18); do + for i in $(seq 1 72); do if npm view "${pkg}@${VERSION}" version >/dev/null 2>&1; then ok=1; break; fi - echo "waiting for ${pkg}@${VERSION} to propagate (attempt ${i}/18)..." + echo "waiting for ${pkg}@${VERSION} to propagate (attempt ${i}/72)..." sleep 10 done if [ "$ok" != 1 ]; then diff --git a/.github/workflows/build-kge.yml b/.github/workflows/build-kge.yml index 3ed59e86cd..07af882f20 100644 --- a/.github/workflows/build-kge.yml +++ b/.github/workflows/build-kge.yml @@ -202,17 +202,17 @@ jobs: local url="https://registry.npmjs.org/${enc}" local tmp="$RUNNER_TEMP/reg.json" local i - for i in $(seq 1 18); do + for i in $(seq 1 72); do if curl -sfL "$url" -o "$tmp"; then if REG_JSON="$tmp" REG_VER="$ver" node -e 'const d=require(process.env.REG_JSON); process.exit(d.versions && d.versions[process.env.REG_VER] ? 0 : 1)'; then echo "verified ${pkg}@${ver} on the npm registry" return 0 fi fi - echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/18)..." + echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/72)..." sleep 10 done - echo "::error::${pkg}@${ver} did not appear on the npm registry within ~3 minutes of publish" + echo "::error::${pkg}@${ver} did not appear on the npm registry within ~12 minutes of publish" return 1 } VERIFY diff --git a/.github/workflows/build-router.yml b/.github/workflows/build-router.yml index 2d27e41b11..75b55ae3f7 100644 --- a/.github/workflows/build-router.yml +++ b/.github/workflows/build-router.yml @@ -170,17 +170,17 @@ jobs: local url="https://registry.npmjs.org/${enc}" local tmp="$RUNNER_TEMP/reg.json" local i - for i in $(seq 1 18); do + for i in $(seq 1 72); do if curl -sfL "$url" -o "$tmp"; then if REG_JSON="$tmp" REG_VER="$ver" node -e 'const d=require(process.env.REG_JSON); process.exit(d.versions && d.versions[process.env.REG_VER] ? 0 : 1)'; then echo "verified ${pkg}@${ver} on the npm registry" return 0 fi fi - echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/18)..." + echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/72)..." sleep 10 done - echo "::error::${pkg}@${ver} did not appear on the npm registry within ~3 minutes of publish" + echo "::error::${pkg}@${ver} did not appear on the npm registry within ~12 minutes of publish" return 1 } @@ -284,17 +284,17 @@ jobs: local url="https://registry.npmjs.org/${enc}" local tmp="$RUNNER_TEMP/reg.json" local i - for i in $(seq 1 18); do + for i in $(seq 1 72); do if curl -sfL "$url" -o "$tmp"; then if REG_JSON="$tmp" REG_VER="$ver" node -e 'const d=require(process.env.REG_JSON); process.exit(d.versions && d.versions[process.env.REG_VER] ? 0 : 1)'; then echo "verified ${pkg}@${ver} on the npm registry" return 0 fi fi - echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/18)..." + echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/72)..." sleep 10 done - echo "::error::${pkg}@${ver} did not appear on the npm registry within ~3 minutes of publish" + echo "::error::${pkg}@${ver} did not appear on the npm registry within ~12 minutes of publish" return 1 } if npm publish --access public; then diff --git a/.github/workflows/build-typesafe.yml b/.github/workflows/build-typesafe.yml index dd14a4eb6c..8ed749a70e 100644 --- a/.github/workflows/build-typesafe.yml +++ b/.github/workflows/build-typesafe.yml @@ -220,17 +220,17 @@ jobs: local url="https://registry.npmjs.org/${enc}" local tmp="$RUNNER_TEMP/reg.json" local i - for i in $(seq 1 18); do + for i in $(seq 1 72); do if curl -sfL "$url" -o "$tmp"; then if REG_JSON="$tmp" REG_VER="$ver" node -e 'const d=require(process.env.REG_JSON); process.exit(d.versions && d.versions[process.env.REG_VER] ? 0 : 1)'; then echo "verified ${pkg}@${ver} on the npm registry" return 0 fi fi - echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/18)..." + echo "waiting for ${pkg}@${ver} to propagate (attempt ${i}/72)..." sleep 10 done - echo "::error::${pkg}@${ver} did not appear on the npm registry within ~3 minutes of publish" + echo "::error::${pkg}@${ver} did not appear on the npm registry within ~12 minutes of publish" return 1 } VERIFY diff --git a/.github/workflows/ruvector-publish.yml b/.github/workflows/ruvector-publish.yml index 1acc81f04f..d0c9f21839 100644 --- a/.github/workflows/ruvector-publish.yml +++ b/.github/workflows/ruvector-publish.yml @@ -83,11 +83,11 @@ jobs: run: | # #1007: the job's conclusion must mean "it is on the registry". VERSION=$(node -p "require('./package.json').version") - for i in $(seq 1 18); do + for i in $(seq 1 72); do if npm view "ruvector@${VERSION}" version >/dev/null 2>&1; then echo "verified ruvector@${VERSION}"; exit 0 fi - echo "waiting for ruvector@${VERSION} to propagate (attempt ${i}/18)..." + echo "waiting for ruvector@${VERSION} to propagate (attempt ${i}/72)..." sleep 10 done echo "::error::ruvector@${VERSION} is not resolvable on npm" diff --git a/Cargo.lock b/Cargo.lock index ef3f7f7e64..f59bdf2483 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11407,7 +11407,7 @@ dependencies = [ [[package]] name = "ruvllm-cli" -version = "2.3.0" +version = "2.3.1" dependencies = [ "anyhow", "assert_cmd", diff --git a/crates/neural-trader-strategies/src/coherence_bridge.rs b/crates/neural-trader-strategies/src/coherence_bridge.rs index f613a7aff9..aca4cffe59 100644 --- a/crates/neural-trader-strategies/src/coherence_bridge.rs +++ b/crates/neural-trader-strategies/src/coherence_bridge.rs @@ -8,7 +8,9 @@ //! //! This module provides a thin wrapper so the full actuation path is: //! -//! Strategy → Intent → CoherenceChecker.check → RiskGate → Order +//! ```text +//! Strategy → Intent → CoherenceChecker.check → RiskGate → Order +//! ``` //! //! The wrapper is lightweight by design: operators may want to skip the //! coherence gate in paper-only flows (which they can, by not diff --git a/crates/neural-trader-strategies/src/ev_kelly.rs b/crates/neural-trader-strategies/src/ev_kelly.rs index 1507147103..ac69a83531 100644 --- a/crates/neural-trader-strategies/src/ev_kelly.rs +++ b/crates/neural-trader-strategies/src/ev_kelly.rs @@ -4,8 +4,10 @@ //! and a current market price `m` in cents (0..=100), the canonical //! fractional Kelly sizing for a YES bet is: //! -//! edge = p - m / 100 -//! kelly = edge / (1 - m / 100) +//! ```text +//! edge = p - m / 100 +//! kelly = edge / (1 - m / 100) +//! ``` //! //! The strategy multiplies `kelly` by a conservative `kelly_fraction` //! (default 0.25) and by `intent.confidence`, then converts the resulting diff --git a/crates/ruvector-diskann/tests/search_allocations.rs b/crates/ruvector-diskann/tests/search_allocations.rs index e09a428563..6fbde0698c 100644 --- a/crates/ruvector-diskann/tests/search_allocations.rs +++ b/crates/ruvector-diskann/tests/search_allocations.rs @@ -2,6 +2,7 @@ use rand::rngs::StdRng; use rand::{Rng, SeedableRng}; use ruvector_diskann::{DiskAnnConfig, DiskAnnIndex}; use std::alloc::{GlobalAlloc, Layout, System}; +use std::cell::Cell; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; struct CountingAllocator; @@ -11,9 +12,20 @@ static LARGE_ALLOCATION_THRESHOLD: AtomicUsize = AtomicUsize::new(usize::MAX); static LARGE_ALLOCATED_BYTES: AtomicUsize = AtomicUsize::new(0); static TOTAL_ALLOCATED_BYTES: AtomicUsize = AtomicUsize::new(0); +thread_local! { + // Search is single-threaded; count only the searching thread. Idle rayon + // workers left over from build() could otherwise add unrelated bytes and + // make the total budget flaky (386 700 vs 300 000 in run 35908988894). + static ON_SEARCH_THREAD: Cell = const { Cell::new(false) }; +} + +fn tracking() -> bool { + TRACKING.load(Ordering::Relaxed) && ON_SEARCH_THREAD.try_with(Cell::get).unwrap_or(false) +} + unsafe impl GlobalAlloc for CountingAllocator { unsafe fn alloc(&self, layout: Layout) -> *mut u8 { - if TRACKING.load(Ordering::Relaxed) { + if tracking() { TOTAL_ALLOCATED_BYTES.fetch_add(layout.size(), Ordering::Relaxed); if layout.size() >= LARGE_ALLOCATION_THRESHOLD.load(Ordering::Relaxed) { LARGE_ALLOCATED_BYTES.fetch_add(layout.size(), Ordering::Relaxed); @@ -27,7 +39,7 @@ unsafe impl GlobalAlloc for CountingAllocator { } unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 { - if TRACKING.load(Ordering::Relaxed) { + if tracking() { TOTAL_ALLOCATED_BYTES.fetch_add(new_size, Ordering::Relaxed); if new_size >= LARGE_ALLOCATION_THRESHOLD.load(Ordering::Relaxed) { LARGE_ALLOCATED_BYTES.fetch_add(new_size, Ordering::Relaxed); @@ -72,11 +84,13 @@ fn pooled_search_does_not_allocate_visited_set_storage() { LARGE_ALLOCATION_THRESHOLD.store(visited_generation_bytes, Ordering::Relaxed); LARGE_ALLOCATED_BYTES.store(0, Ordering::Relaxed); TOTAL_ALLOCATED_BYTES.store(0, Ordering::Relaxed); + ON_SEARCH_THREAD.with(|t| t.set(true)); TRACKING.store(true, Ordering::Relaxed); for _ in 0..SEARCHES { std::hint::black_box(index.search(&query, 10).unwrap()); } TRACKING.store(false, Ordering::Relaxed); + ON_SEARCH_THREAD.with(|t| t.set(false)); let allocated = LARGE_ALLOCATED_BYTES.load(Ordering::Relaxed); let total_allocated = TOTAL_ALLOCATED_BYTES.load(Ordering::Relaxed); diff --git a/crates/ruvector-filter/src/expression.rs b/crates/ruvector-filter/src/expression.rs index c19257f24c..4b1719cda6 100644 --- a/crates/ruvector-filter/src/expression.rs +++ b/crates/ruvector-filter/src/expression.rs @@ -2,8 +2,15 @@ use serde::{Deserialize, Serialize}; use serde_json::Value; /// Filter expression for querying vectors by payload +/// +/// Serialized through [`FilterExpressionWire`]: logical operators use struct +/// variants on the wire (`{"type":"and","filters":[...]}`, +/// `{"type":"not","filter":{...}}`), while leaf variants keep their original +/// shape. Deriving serde directly on this internally tagged enum made the +/// recursive `And`/`Or`/`Not` newtype variants unserializable (serde cannot tag a +/// sequence) and sent rustc into unbounded `TaggedSerializer` nesting (E0275). #[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(tag = "type", rename_all = "snake_case")] +#[serde(into = "FilterExpressionWire", from = "FilterExpressionWire")] pub enum FilterExpression { // Comparison operators Eq { @@ -77,6 +84,159 @@ pub enum FilterExpression { }, } +/// Serde representation of [`FilterExpression`]; see its docs. +#[derive(Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +enum FilterExpressionWire { + Eq { + field: String, + value: Value, + }, + Ne { + field: String, + value: Value, + }, + Gt { + field: String, + value: Value, + }, + Gte { + field: String, + value: Value, + }, + Lt { + field: String, + value: Value, + }, + Lte { + field: String, + value: Value, + }, + Range { + field: String, + gte: Option, + lte: Option, + }, + In { + field: String, + values: Vec, + }, + Match { + field: String, + text: String, + }, + GeoRadius { + field: String, + lat: f64, + lon: f64, + radius_m: f64, + }, + GeoBoundingBox { + field: String, + top_left: (f64, f64), + bottom_right: (f64, f64), + }, + And { + filters: Vec, + }, + Or { + filters: Vec, + }, + Not { + filter: Box, + }, + Exists { + field: String, + }, + IsNull { + field: String, + }, +} + +impl From for FilterExpressionWire { + fn from(e: FilterExpression) -> Self { + use FilterExpression as F; + match e { + F::Eq { field, value } => Self::Eq { field, value }, + F::Ne { field, value } => Self::Ne { field, value }, + F::Gt { field, value } => Self::Gt { field, value }, + F::Gte { field, value } => Self::Gte { field, value }, + F::Lt { field, value } => Self::Lt { field, value }, + F::Lte { field, value } => Self::Lte { field, value }, + F::Range { field, gte, lte } => Self::Range { field, gte, lte }, + F::In { field, values } => Self::In { field, values }, + F::Match { field, text } => Self::Match { field, text }, + F::GeoRadius { + field, + lat, + lon, + radius_m, + } => Self::GeoRadius { + field, + lat, + lon, + radius_m, + }, + F::GeoBoundingBox { + field, + top_left, + bottom_right, + } => Self::GeoBoundingBox { + field, + top_left, + bottom_right, + }, + F::And(filters) => Self::And { filters }, + F::Or(filters) => Self::Or { filters }, + F::Not(filter) => Self::Not { filter }, + F::Exists { field } => Self::Exists { field }, + F::IsNull { field } => Self::IsNull { field }, + } + } +} + +impl From for FilterExpression { + fn from(w: FilterExpressionWire) -> Self { + use FilterExpressionWire as W; + match w { + W::Eq { field, value } => Self::Eq { field, value }, + W::Ne { field, value } => Self::Ne { field, value }, + W::Gt { field, value } => Self::Gt { field, value }, + W::Gte { field, value } => Self::Gte { field, value }, + W::Lt { field, value } => Self::Lt { field, value }, + W::Lte { field, value } => Self::Lte { field, value }, + W::Range { field, gte, lte } => Self::Range { field, gte, lte }, + W::In { field, values } => Self::In { field, values }, + W::Match { field, text } => Self::Match { field, text }, + W::GeoRadius { + field, + lat, + lon, + radius_m, + } => Self::GeoRadius { + field, + lat, + lon, + radius_m, + }, + W::GeoBoundingBox { + field, + top_left, + bottom_right, + } => Self::GeoBoundingBox { + field, + top_left, + bottom_right, + }, + W::And { filters } => Self::And(filters), + W::Or { filters } => Self::Or(filters), + W::Not { filter } => Self::Not(filter), + W::Exists { field } => Self::Exists { field }, + W::IsNull { field } => Self::IsNull { field }, + } + } +} + impl FilterExpression { /// Create an equality filter pub fn eq(field: impl Into, value: Value) -> Self { @@ -280,5 +440,33 @@ mod tests { let json = serde_json::to_string(&filter).unwrap(); let deserialized: FilterExpression = serde_json::from_str(&json).unwrap(); assert!(matches!(deserialized, FilterExpression::Eq { .. })); + // Leaf wire format is unchanged by the wire-enum indirection. + assert_eq!( + serde_json::to_value(&filter).unwrap(), + json!({"type": "eq", "field": "status", "value": "active"}) + ); + } + + #[test] + fn test_serialization_of_logical_operators() { + let filter = FilterExpression::and(vec![ + FilterExpression::eq("status", json!("active")), + FilterExpression::or(vec![ + FilterExpression::gte("age", json!(18)), + FilterExpression::not(FilterExpression::exists("banned")), + ]), + ]); + let value = serde_json::to_value(&filter).unwrap(); + assert_eq!(value["type"], "and"); + assert_eq!(value["filters"][1]["type"], "or"); + assert_eq!(value["filters"][1]["filters"][1]["type"], "not"); + assert_eq!( + value["filters"][1]["filters"][1]["filter"]["type"], + "exists" + ); + + let back: FilterExpression = serde_json::from_value(value.clone()).unwrap(); + assert_eq!(serde_json::to_value(&back).unwrap(), value); + assert_eq!(back.get_fields(), vec!["age", "banned", "status"]); } } diff --git a/crates/ruvector-metrics/src/lib.rs b/crates/ruvector-metrics/src/lib.rs index 154f0d4af7..167d686ecf 100644 --- a/crates/ruvector-metrics/src/lib.rs +++ b/crates/ruvector-metrics/src/lib.rs @@ -75,6 +75,12 @@ lazy_static! { /// Gather all metrics in Prometheus text format pub fn gather_metrics() -> String { + // The metrics are lazy statics, registered on first use. Initialize the + // unlabeled ones so a fresh process exports them (at zero) instead of an + // empty page until something happens to touch them. + lazy_static::initialize(&COLLECTIONS_TOTAL); + lazy_static::initialize(&MEMORY_USAGE_BYTES); + lazy_static::initialize(&UPTIME_SECONDS); let encoder = TextEncoder::new(); let metric_families = prometheus::gather(); let mut buffer = Vec::new(); diff --git a/crates/ruvector-mincut-gated-transformer/tests/verification.rs b/crates/ruvector-mincut-gated-transformer/tests/verification.rs index e841f1f6e6..5ce07e5d44 100644 --- a/crates/ruvector-mincut-gated-transformer/tests/verification.rs +++ b/crates/ruvector-mincut-gated-transformer/tests/verification.rs @@ -482,12 +482,20 @@ fn test_flash_attention_memory_efficiency() { println!("FlashAttention 1024 seq_len: {:?}", elapsed); - // Should complete without OOM and in reasonable time + // Should complete without OOM. The wall-clock bound only means something + // for an optimized build: CI runs this unoptimized on shared runners, where + // 1024x1024 took 1.27s (run 35900536715). assert!( - elapsed.as_millis() < 1000, - "FlashAttention too slow: {:?}", - elapsed + output.iter().all(|x| x.is_finite()), + "non-finite attention output" ); + if !cfg!(debug_assertions) { + assert!( + elapsed.as_millis() < 1000, + "FlashAttention too slow: {:?}", + elapsed + ); + } } // ============================================================================ diff --git a/crates/ruvllm-cli/Cargo.toml b/crates/ruvllm-cli/Cargo.toml index 4925bfe54e..bd85e2fc32 100644 --- a/crates/ruvllm-cli/Cargo.toml +++ b/crates/ruvllm-cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "ruvllm-cli" -version.workspace = true +version = "2.3.1" edition.workspace = true rust-version.workspace = true license.workspace = true diff --git a/crates/ruvllm/src/tests/attention_tests.rs b/crates/ruvllm/src/tests/attention_tests.rs index e1270cbfcb..70d9a5b831 100644 --- a/crates/ruvllm/src/tests/attention_tests.rs +++ b/crates/ruvllm/src/tests/attention_tests.rs @@ -678,11 +678,15 @@ fn test_attention_benchmark_short_sequence() { let duration = start.elapsed(); let avg_us = duration.as_micros() as f64 / iterations as f64; - assert!( - avg_us < 1000.0, - "Short sequence attention should be fast: {}us", - avg_us - ); + // cargo-llvm-cov instruments every branch (the Test Coverage job measured + // 1099us); a wall-clock bound is only meaningful uninstrumented. + if std::env::var_os("LLVM_PROFILE_FILE").is_none() { + assert!( + avg_us < 1000.0, + "Short sequence attention should be fast: {}us", + avg_us + ); + } } #[test]