diff --git a/.github/workflows/build-graph-node.yml b/.github/workflows/build-graph-node.yml index 7c132bda49..c8b016429e 100644 --- a/.github/workflows/build-graph-node.yml +++ b/.github/workflows/build-graph-node.yml @@ -165,4 +165,31 @@ jobs: working-directory: npm/packages/graph-node env: NODE_AUTH_TOKEN: ${{ secrets.NPM_TOKEN }} - run: npm publish --access public || echo "Package may already exist" + run: | + # Same class as #1007: never swallow a failed publish. + VERSION=$(node -p "require('./package.json').version") + if npm view "@ruvector/graph-node@${VERSION}" version >/dev/null 2>&1; then + echo "@ruvector/graph-node@${VERSION} already published" + else + npm publish --access public + fi + + - name: Verify all packages are on the registry + working-directory: npm/packages/graph-node + run: | + VERSION=$(node -p "require('./package.json').version") + for pkg in @ruvector/graph-node @ruvector/graph-node-linux-x64-gnu \ + @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 + 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)..." + sleep 10 + done + if [ "$ok" != 1 ]; then + echo "::error::${pkg}@${VERSION} is not resolvable on npm" + exit 1 + fi + echo "verified ${pkg}@${VERSION}" + done diff --git a/.github/workflows/ruvector-publish.yml b/.github/workflows/ruvector-publish.yml index f75fe2d1e5..1acc81f04f 100644 --- a/.github/workflows/ruvector-publish.yml +++ b/.github/workflows/ruvector-publish.yml @@ -77,3 +77,18 @@ jobs: run: npm publish --access public env: NODE_AUTH_TOKEN: ${{ secrets.NPM_TOKEN }} + + - name: Verify on the registry + if: ${{ github.event.inputs.dry_run != 'true' }} + 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 + 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)..." + sleep 10 + done + echo "::error::ruvector@${VERSION} is not resolvable on npm" + exit 1 diff --git a/crates/ruvector-graph-node/src/lib.rs b/crates/ruvector-graph-node/src/lib.rs index 5460a661df..808a215f4e 100644 --- a/crates/ruvector-graph-node/src/lib.rs +++ b/crates/ruvector-graph-node/src/lib.rs @@ -315,6 +315,30 @@ fn run_query( }) } +/// Build the stored property map for an edge: the caller's public `metadata` +/// (#984 — it was previously dropped) plus this binding's internal `__` keys. +/// Internal keys are inserted last so caller metadata can never spoof them. +fn edge_properties( + metadata: Option>, + confidence: f32, + embedding: Vec, +) -> HashMap { + let mut properties: HashMap = metadata + .unwrap_or_default() + .into_iter() + .map(|(k, v)| (k, PropertyValue::String(v))) + .collect(); + properties.insert( + "__confidence".to_string(), + PropertyValue::FloatArray(vec![confidence]), + ); + properties.insert( + "__embedding".to_string(), + PropertyValue::FloatArray(embedding), + ); + properties +} + /// Properties this binding stores for its own use rather than the caller's. /// /// Embeddings and edge confidences are persisted as ordinary properties under a @@ -487,6 +511,7 @@ impl GraphDatabase { let description = edge.description.clone(); let embedding = edge.embedding.to_vec(); let confidence = edge.confidence.unwrap_or(1.0) as f32; + let metadata = edge.metadata; let hydrated = self.hydrated.clone(); tokio::task::spawn_blocking(move || { @@ -499,16 +524,7 @@ impl GraphDatabase { .map_err(|e| Error::from_reason(format!("Failed to create edge: {}", e)))?; drop(hg); - let properties = HashMap::from([ - ( - "__confidence".to_string(), - PropertyValue::FloatArray(vec![confidence]), - ), - ( - "__embedding".to_string(), - PropertyValue::FloatArray(embedding), - ), - ]); + let properties = edge_properties(metadata, confidence, embedding); let graph_edge = GraphEdge::new(edge_id.clone(), from, to, description, properties); if let Some(storage_arc) = storage { storage_arc @@ -548,6 +564,7 @@ impl GraphDatabase { let description = hyperedge.description.clone(); let embedding = hyperedge.embedding.to_vec(); let confidence = hyperedge.confidence.unwrap_or(1.0) as f32; + let metadata = hyperedge.metadata; let graph_db = self.graph_db.clone(); let hydrated = self.hydrated.clone(); @@ -571,10 +588,18 @@ impl GraphDatabase { nodes, edge_type: "HYPEREDGE".to_string(), description: Some(description), - properties: HashMap::from([( - "__embedding".to_string(), - PropertyValue::FloatArray(embedding), - )]), + properties: { + let mut props: HashMap = metadata + .unwrap_or_default() + .into_iter() + .map(|(k, v)| (k, PropertyValue::String(v))) + .collect(); + props.insert( + "__embedding".to_string(), + PropertyValue::FloatArray(embedding), + ); + props + }, confidence, }; storage_arc @@ -814,16 +839,7 @@ impl GraphDatabase { edge.from, edge.to, edge.description, - HashMap::from([ - ( - "__confidence".to_string(), - PropertyValue::FloatArray(vec![confidence]), - ), - ( - "__embedding".to_string(), - PropertyValue::FloatArray(embedding), - ), - ]), + edge_properties(edge.metadata, confidence, embedding), ); if let Some(storage_arc) = storage.as_ref() { storage_arc @@ -1258,6 +1274,109 @@ mod tests { std::fs::remove_file(path).expect("remove test database"); } + /// #984: public edge metadata from createEdge and batchInsert must be + /// stored, returned by query(), survive a reopen, and never override the + /// binding's internal `__` properties. + #[tokio::test] + async fn edge_metadata_round_trips_through_query_and_reopen() { + let path = temp_storage_path("edge-metadata"); + let db = GraphDatabase::new(Some(JsGraphOptions { + distance_metric: Some(JsDistanceMetric::Cosine), + dimensions: Some(2), + storage_path: Some(path.clone()), + })) + .expect("create persistent database"); + for id in ["a", "b", "c"] { + db.create_node(JsNode { + id: id.to_string(), + embedding: Float32Array::new(vec![1.0, 0.0]), + labels: None, + properties: None, + }) + .await + .expect("create node"); + } + let metadata = |edge: &str| { + HashMap::from([ + ("sourceEdgeId".to_string(), edge.to_string()), + ("weight".to_string(), "0.8".to_string()), + ("note".to_string(), "témoin ✓".to_string()), + ("__confidence".to_string(), "spoofed".to_string()), + ]) + }; + db.create_edge(JsEdge { + from: "a".to_string(), + to: "b".to_string(), + description: "supports".to_string(), + embedding: Float32Array::new(vec![1.0, 0.0]), + confidence: Some(0.75), + metadata: Some(metadata("single")), + }) + .await + .expect("create edge"); + db.batch_insert(JsBatchInsert { + nodes: vec![], + edges: vec![JsEdge { + from: "a".to_string(), + to: "c".to_string(), + description: "supports".to_string(), + embedding: Float32Array::new(vec![1.0, 0.0]), + confidence: Some(0.75), + metadata: Some(metadata("batch")), + }], + }) + .await + .expect("batch insert"); + + let check = |result: JsQueryResult| { + assert_eq!(result.edges.len(), 2, "both edges returned"); + let mut seen: Vec = result + .edges + .iter() + .map(|e| { + assert_eq!(e.properties.get("weight").map(String::as_str), Some("0.8")); + assert_eq!( + e.properties.get("note").map(String::as_str), + Some("témoin ✓") + ); + assert!(!e.properties.contains_key("__confidence")); + e.properties["sourceEdgeId"].clone() + }) + .collect(); + seen.sort(); + assert_eq!(seen, vec!["batch".to_string(), "single".to_string()]); + }; + check( + db.query("MATCH (a)-[r]->(b) RETURN a,r,b".to_string()) + .await + .expect("query"), + ); + drop(db); + + let reopened = GraphDatabase::open(path.clone()).expect("reopen"); + check( + reopened + .query("MATCH (a)-[r]->(b) RETURN a,r,b".to_string()) + .await + .expect("query after reopen"), + ); + // The spoofed key must not have replaced the real confidence. + let stored = reopened + .graph_db + .read() + .expect("graph lock") + .get_edges_by_type("supports"); + assert_eq!(stored.len(), 2); + for edge in stored { + assert_eq!( + prop_to_f32_vec(edge.properties.get("__confidence")), + vec![0.75] + ); + } + drop(reopened); + std::fs::remove_file(path).expect("remove test database"); + } + #[tokio::test] async fn non_cascade_retains_durable_relationships_but_cascade_removes_them() { let non_cascade_path = temp_storage_path("non-cascade"); diff --git a/crates/ruvector-kge/src/ann.rs b/crates/ruvector-kge/src/ann.rs index 2a70eeed28..6aa37d0902 100644 --- a/crates/ruvector-kge/src/ann.rs +++ b/crates/ruvector-kge/src/ann.rs @@ -3,31 +3,23 @@ //! entity's [`Scorer::index_vector`], retrieve candidates for a query, then //! exact-rerank with the true score. //! -//! ## Why not `DistanceMetric::DotProduct` directly +//! ## Why the MIPS→L2 route rather than `DistanceMetric::DotProduct` //! -//! ADR-001 names a `DotProduct` HNSW, but router-core's `DotProduct` is a -//! *similarity*, not a distance, and the index treats it as one incorrectly for -//! MIPS. Two facts from `crates/ruvector-router-core/src/`: +//! ADR-001 names a `DotProduct` HNSW. router-core's `DotProduct` is already a +//! proper distance for that: `distance::dot_product` returns `-(a·b)`, so the +//! index's smallest-distance-first search is a maximum-inner-product search. +//! Do **not** negate the query on top of that — doing so searches for the +//! *minimum* inner product (recall@10 = 0.000, as reported in #1009). With the +//! query passed through unchanged the same table measures recall@10 = 1.000; see +//! `dotproduct_recall_matches_mips`. //! -//! - `distance.rs:18` returns the raw dot product for `DotProduct`, and -//! `index.rs`'s search/build keep the **smallest** "distance" (min-heap). -//! So a `DotProduct` index retrieves the *minimum* inner product, the -//! opposite of MIPS. Negating the query flips the search target but not the -//! graph: `insert` (`index.rs:135`) wires each node to its lowest-dot -//! neighbours, and that entity-to-entity orientation is symmetric under any -//! global sign flip — no vector transform fixes it. -//! -//! - The robust, exact fix is the standard MIPS→L2 reduction (Bachrach et al. -//! 2014): append one coordinate so Euclidean nearest-neighbour equals maximum -//! inner product. With `φ² = max_e ‖x_e‖²`, -//! `x_e' = [x_e ; √(φ² − ‖x_e‖²)]` and `q' = [q ; 0]`, then -//! `‖q' − x_e'‖² = ‖q‖² + φ² − 2 q·x_e`, so minimising L2 maximises `q·x_e`. -//! Euclidean **is** a proper distance, so the HNSW graph is built and searched -//! consistently. This is the "handle it and document" path the task allows. -//! -//! For the record, a raw-`DotProduct`-with-negated-query index measured well -//! below the L2 route on the same table — see the `#[ignore]`d -//! `raw_dotproduct_recall_for_the_record` test. +//! This module keeps the L2 reduction because Euclidean is a metric (non-negative, +//! zero on identity), which the neighbour-selection heuristic's pruning +//! comparisons are written against. The standard MIPS→L2 reduction (Bachrach +//! et al. 2014) appends one coordinate so Euclidean nearest-neighbour equals +//! maximum inner product. With `φ² = max_e ‖x_e‖²`, +//! `x_e' = [x_e ; √(φ² − ‖x_e‖²)]` and `q' = [q ; 0]`, then +//! `‖q' − x_e'‖² = ‖q‖² + φ² − 2 q·x_e`, so minimising L2 maximises `q·x_e`. use crate::scorer::Scorer; use crate::{Candidate, EntityId, KgeError, Result, Tables}; @@ -193,11 +185,10 @@ mod tests { assert!(recall >= 0.9, "recall@{k} = {recall:.4} < 0.9"); } - /// Not a gate — records how the raw `DotProduct` + negated-query route does - /// on the same table, so the deviation to L2 is backed by a number. + /// Gate for #1009: router-core's `DotProduct` HNSW, queried with the + /// *unmodified* query vector, is a maximum-inner-product search. #[test] - #[ignore = "diagnostic: documents the rejected DotProduct route"] - fn raw_dotproduct_recall_for_the_record() { + fn dotproduct_recall_matches_mips() { use ruvector_router_core::index::{HnswConfig, HnswIndex}; use ruvector_router_core::types::{DistanceMetric, SearchQuery}; @@ -238,11 +229,9 @@ mod tests { Side::Tail, ); let truth: Vec = batch.top_k(&q, k).into_iter().map(|c| c.entity).collect(); - // Negate the query to turn the min-heap into a max-inner-product search. - let neg: Vec = q.iter().map(|x| -x).collect(); let results = hnsw .search(&SearchQuery { - vector: neg, + vector: q.clone(), k: ef, filters: None, threshold: None, @@ -266,6 +255,7 @@ mod tests { hits += truth.iter().filter(|e| got.contains(e)).count(); } let recall = hits as f32 / (n_queries * k) as f32; - println!("RAW DotProduct+negated-query recall@{k} = {recall:.4} (rejected route)"); + println!("DotProduct recall@{k} = {recall:.4}"); + assert!(recall >= 0.9, "DotProduct recall@{k} = {recall:.4} < 0.9"); } } diff --git a/crates/ruvllm-cli/src/commands/quantize.rs b/crates/ruvllm-cli/src/commands/quantize.rs index 04909f307c..913d277554 100644 --- a/crates/ruvllm-cli/src/commands/quantize.rs +++ b/crates/ruvllm-cli/src/commands/quantize.rs @@ -1,33 +1,33 @@ //! Quantize command implementation //! -//! Quantizes models to GGUF format with K-quant or Q8 quantization. -//! Optimized for Apple Neural Engine inference on M4 Pro and other Apple Silicon. +//! Intended to quantize models to GGUF (K-quant / Q8). The GGUF tensor writer is +//! not implemented yet, so `run` validates its arguments and fails closed +//! instead of writing an empty or header-only file (#968). -use std::fs::{self, File}; -use std::io::{BufReader, BufWriter, Read, Seek, SeekFrom, Write}; -use std::path::{Path, PathBuf}; -use std::time::Instant; +use std::path::PathBuf; use colored::Colorize; -use indicatif::{ProgressBar, ProgressStyle}; -use ruvllm::{ - estimate_memory_q4, estimate_memory_q5, estimate_memory_q8, GgufFile, GgufQuantType, - MemoryEstimate, QuantConfig, RuvltraQuantizer, TargetFormat, -}; +use ruvllm::TargetFormat; /// Run the quantize command +/// +/// The GGUF tensor writer this command needs does not exist yet: the previous +/// implementation advanced a progress bar without loading any tensors and wrote +/// either a 75-byte GGUF header declaring zero tensors (SafeTensors/PyTorch +/// input) or an empty file (GGUF input), then reported success (#968). Until a +/// real writer lands the command validates its arguments and fails closed — +/// before any output file is created — pointing at working converters. pub async fn run( model: &str, output: &str, quant: &str, - ane_optimize: bool, - keep_embed_fp16: bool, - keep_output_fp16: bool, - verbose: bool, + _ane_optimize: bool, + _keep_embed_fp16: bool, + _keep_output_fp16: bool, + _verbose: bool, cache_dir: &str, ) -> anyhow::Result<()> { - // Parse target format let format = TargetFormat::from_str(quant).ok_or_else(|| { anyhow::anyhow!( "Unknown quantization format: {}. Supported: q4_k_m, q5_k_m, q8_0, f16", @@ -35,45 +35,7 @@ pub async fn run( ) })?; - println!("\n{} RuvLTRA Model Quantizer", "==>".bright_blue().bold()); - println!(" Target format: {}", format.name().bright_cyan()); - println!(" Bits per weight: {:.1}", format.bits_per_weight()); - println!( - " ANE optimization: {}", - if ane_optimize { "enabled" } else { "disabled" } - ); - - // Resolve input model path let input_path = resolve_model_path(model, cache_dir)?; - println!( - "\n{} Input model: {}", - "-->".bright_blue(), - input_path.display() - ); - - // Determine output path - let output_path = if output.is_empty() { - // Generate output name based on input - let stem = input_path - .file_stem() - .and_then(|s| s.to_str()) - .unwrap_or("model"); - let output_name = format!("{}-{}.gguf", stem, quant.to_lowercase()); - input_path - .parent() - .unwrap_or(Path::new(".")) - .join(output_name) - } else { - PathBuf::from(output) - }; - - println!( - "{} Output file: {}", - "-->".bright_blue(), - output_path.display() - ); - - // Check if input exists if !input_path.exists() { return Err(anyhow::anyhow!( "Input model not found: {}", @@ -81,99 +43,39 @@ pub async fn run( )); } - // Check if output already exists - if output_path.exists() { - println!( - "\n{} Output file already exists. Overwriting...", - "Warning:".yellow().bold() - ); - } - - // Get input file size - let input_metadata = fs::metadata(&input_path)?; - let input_size = input_metadata.len(); - println!( - "\n{} Input size: {:.2} MB", - "-->".bright_blue(), - input_size as f64 / (1024.0 * 1024.0) - ); - - // Estimate output size - let estimated_output = estimate_output_size(input_size, format); - println!( - "{} Estimated output: {:.2} MB ({:.1}x compression)", - "-->".bright_blue(), - estimated_output as f64 / (1024.0 * 1024.0), - input_size as f64 / estimated_output as f64 - ); - - // Memory estimates for common model sizes - print_memory_estimates(format); - - // Create quantizer configuration - let config = QuantConfig::default() - .with_format(format) - .with_ane_optimization(ane_optimize) - .with_verbose(verbose); - - let mut config = config; - config.keep_embed_fp16 = keep_embed_fp16; - config.keep_output_fp16 = keep_output_fp16; - - // Check if input is GGUF - let is_gguf = input_path - .extension() - .and_then(|e| e.to_str()) - .map(|e| e.to_lowercase() == "gguf") - .unwrap_or(false); - - println!("\n{} Starting quantization...", "==>".bright_blue().bold()); - - let start_time = Instant::now(); - - if is_gguf { - // Quantize GGUF to GGUF (re-quantization) - quantize_gguf_model(&input_path, &output_path, config, verbose).await?; + let kind = if input_path.is_dir() { + "a HuggingFace model directory" } else { - // Quantize from other formats (safetensors, etc.) - quantize_model(&input_path, &output_path, config, verbose).await?; - } - - let elapsed = start_time.elapsed(); - - // Verify output - let output_metadata = fs::metadata(&output_path)?; - let output_size = output_metadata.len(); - - println!("\n{} Quantization complete!", "==>".bright_green().bold()); - println!( - " Output size: {:.2} MB", - output_size as f64 / (1024.0 * 1024.0) - ); - println!( - " Compression: {:.1}x", - input_size as f64 / output_size as f64 - ); - println!(" Time: {:.1}s", elapsed.as_secs_f64()); - println!( - " Throughput: {:.1} MB/s", - input_size as f64 / (1024.0 * 1024.0) / elapsed.as_secs_f64() - ); - - println!( - "\n{} Output saved to: {}", - "-->".bright_green(), - output_path.display() - ); + match input_path + .extension() + .and_then(|e| e.to_str()) + .map(|e| e.to_lowercase()) + .as_deref() + { + Some("gguf") => "GGUF (re-quantization)", + Some("safetensors") => "SafeTensors", + Some("bin") | Some("pt") => "PyTorch", + _ => "this input format", + } + }; - // Usage hint - println!( - "\n{} To use the quantized model:", - "Tip:".bright_cyan().bold() - ); - println!(" ruvllm chat {} -q {}", output_path.display(), quant); + let target = if output.is_empty() { + "-.gguf".to_string() + } else { + output.to_string() + }; - Ok(()) + Err(anyhow::anyhow!( + "`ruvllm quantize` cannot yet produce {} from {} ({}): the GGUF tensor \ + writer is not implemented, so no output was written to {}. Use llama.cpp \ + (`convert_hf_to_gguf.py` then `llama-quantize {}`) or \ + candle's GGUF writer instead. Tracking: https://github.com/ruvnet/ruvector/issues/968", + format.name(), + kind, + input_path.display(), + target, + quant.to_uppercase(), + )) } /// Resolve model path from identifier or path @@ -208,249 +110,6 @@ fn resolve_model_path(model: &str, cache_dir: &str) -> anyhow::Result { Ok(path) } -/// Estimate output size based on format -fn estimate_output_size(input_bytes: u64, format: TargetFormat) -> u64 { - // Assume input is FP32 - let input_elements = input_bytes / 4; - let bits_per_weight = format.bits_per_weight() as f64; - - ((input_elements as f64 * bits_per_weight) / 8.0) as u64 -} - -/// Print memory estimates for common model sizes -fn print_memory_estimates(format: TargetFormat) { - println!( - "\n{} Memory estimates for {}:", - "-->".bright_blue(), - format.name() - ); - - // RuvLTRA-Small (0.5B) estimates - let estimate_fn: fn(f64, usize, usize, usize) -> MemoryEstimate = match format { - TargetFormat::Q4_K_M => estimate_memory_q4, - TargetFormat::Q5_K_M => estimate_memory_q5, - TargetFormat::Q8_0 => estimate_memory_q8, - TargetFormat::F16 => |p, v, h, l| { - let mut e = estimate_memory_q8(p, v, h, l); - e.total_bytes *= 2; - e.total_mb *= 2.0; - e - }, - // PiQ3: 3.0625 bits/weight vs Q4's ~4.5, so ~68% of Q4 size - TargetFormat::PiQ3 => |p, v, h, l| { - let mut e = estimate_memory_q4(p, v, h, l); - e.total_bytes = (e.total_bytes as f64 * 3.0625 / 4.5) as usize; - e.total_mb = e.total_bytes as f64 / (1024.0 * 1024.0); - e - }, - // PiQ2: 2.0625 bits/weight vs Q4's ~4.5, so ~46% of Q4 size - TargetFormat::PiQ2 => |p, v, h, l| { - let mut e = estimate_memory_q4(p, v, h, l); - e.total_bytes = (e.total_bytes as f64 * 2.0625 / 4.5) as usize; - e.total_mb = e.total_bytes as f64 / (1024.0 * 1024.0); - e - }, - }; - - // Qwen2.5-0.5B (RuvLTRA-Small) - let est_05b = estimate_fn(0.5, 151936, 896, 24); - println!( - " RuvLTRA-Small (0.5B): {:.0} MB ({:.1}x compression)", - est_05b.total_mb, est_05b.compression_ratio - ); - - // Also show for 1B and 3B for reference - let est_1b = estimate_fn(1.0, 151936, 1536, 28); - println!( - " 1B model: {:.0} MB ({:.1}x compression)", - est_1b.total_mb, est_1b.compression_ratio - ); - - let est_3b = estimate_fn(3.0, 151936, 2048, 36); - println!( - " 3B model: {:.0} MB ({:.1}x compression)", - est_3b.total_mb, est_3b.compression_ratio - ); -} - -/// Quantize a GGUF model (re-quantization) -async fn quantize_gguf_model( - input_path: &Path, - output_path: &Path, - config: QuantConfig, - verbose: bool, -) -> anyhow::Result<()> { - // Load input GGUF - let gguf = GgufFile::open_mmap(input_path)?; - - println!( - " Architecture: {}", - gguf.architecture().unwrap_or("unknown") - ); - println!(" Tensors: {}", gguf.tensors.len()); - - let total_size: usize = gguf.tensors.iter().map(|t| t.byte_size()).sum(); - - // Create progress bar - let pb = ProgressBar::new(total_size as u64); - pb.set_style( - ProgressStyle::default_bar() - .template("{spinner:.green} [{elapsed_precise}] [{bar:40.cyan/blue}] {bytes}/{total_bytes} ({eta})") - .unwrap() - .progress_chars("#>-"), - ); - - // Create quantizer - let mut quantizer = RuvltraQuantizer::new(config.clone())?; - - // Open output file - let output_file = File::create(output_path)?; - let mut writer = BufWriter::new(output_file); - - // Write GGUF header (we'll need to implement proper GGUF writing) - // For now, we'll process tensors and show progress - let mut processed = 0usize; - - for tensor_info in &gguf.tensors { - if verbose { - pb.set_message(format!("Processing: {}", tensor_info.name)); - } - - // Load tensor as FP32 - let tensor_data = gguf.load_tensor_f32(&tensor_info.name)?; - - // Quantize - let quantized = quantizer.quantize_tensor(&tensor_data, &tensor_info.name)?; - - // In a full implementation, we'd write this to the output GGUF - // For now, accumulate statistics - processed += tensor_info.byte_size(); - pb.set_position(processed as u64); - } - - pb.finish_with_message("Quantization complete"); - - // Write placeholder output (in production, write proper GGUF) - writer.write_all(&[0u8; 0])?; - - // Print stats - let stats = quantizer.stats(); - if verbose { - println!("\n Tensors quantized: {}", stats.tensors_quantized); - println!(" Elements processed: {}", stats.elements_processed); - } - - Ok(()) -} - -/// Quantize from other formats (safetensors, etc.) -async fn quantize_model( - input_path: &Path, - output_path: &Path, - config: QuantConfig, - verbose: bool, -) -> anyhow::Result<()> { - // Get file size - let input_size = fs::metadata(input_path)?.len(); - - // Create progress bar - let pb = ProgressBar::new(input_size); - pb.set_style( - ProgressStyle::default_bar() - .template("{spinner:.green} [{elapsed_precise}] [{bar:40.cyan/blue}] {bytes}/{total_bytes} ({eta})") - .unwrap() - .progress_chars("#>-"), - ); - - // Create quantizer - let mut quantizer = RuvltraQuantizer::new(config.clone())?; - - // For non-GGUF formats, we'd need to implement specific loaders - // This is a placeholder that shows the infrastructure - pb.set_message("Loading model..."); - - // Check file type and process accordingly - let extension = input_path - .extension() - .and_then(|e| e.to_str()) - .map(|e| e.to_lowercase()) - .unwrap_or_default(); - - match extension.as_str() { - "safetensors" => { - pb.set_message("Processing safetensors format..."); - // In production, use safetensors crate to load tensors - // For now, simulate processing - pb.set_position(input_size / 2); - tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; - pb.set_position(input_size); - } - "bin" | "pt" => { - pb.set_message("Processing PyTorch format..."); - // In production, use tch-rs or similar to load PyTorch tensors - pb.set_position(input_size / 2); - tokio::time::sleep(tokio::time::Duration::from_millis(100)).await; - pb.set_position(input_size); - } - _ => { - pb.set_message("Processing unknown format..."); - pb.set_position(input_size); - } - } - - pb.finish_with_message("Processing complete"); - - // Create output file - let output_file = File::create(output_path)?; - let mut writer = BufWriter::new(output_file); - - // Write minimal GGUF header for testing - // In production, this would be a proper GGUF file - write_gguf_header(&mut writer, &config)?; - - if verbose { - let stats = quantizer.stats(); - println!( - "\n Quantizer stats: {} tensors, {} elements", - stats.tensors_quantized, stats.elements_processed - ); - } - - Ok(()) -} - -/// Write a basic GGUF header -fn write_gguf_header(writer: &mut W, config: &QuantConfig) -> anyhow::Result<()> { - // GGUF magic: "GGUF" in little-endian - writer.write_all(&0x46554747u32.to_le_bytes())?; - - // Version: 3 - writer.write_all(&3u32.to_le_bytes())?; - - // Tensor count: 0 (placeholder) - writer.write_all(&0u64.to_le_bytes())?; - - // Metadata count: 1 - writer.write_all(&1u64.to_le_bytes())?; - - // Write one metadata entry for quantization type - let key = "general.quantization_type"; - let key_len = key.len() as u64; - writer.write_all(&key_len.to_le_bytes())?; - writer.write_all(key.as_bytes())?; - - // String type: 8 - writer.write_all(&8u32.to_le_bytes())?; - - // Value - let value = config.format.name(); - let value_len = value.len() as u64; - writer.write_all(&value_len.to_le_bytes())?; - writer.write_all(value.as_bytes())?; - - Ok(()) -} - /// Print detailed format comparison pub fn print_format_comparison() { println!( @@ -480,3 +139,61 @@ pub fn print_format_comparison() { "F16", "16", "~1000 MB", "Excellent", "No quant loss" ); } + +#[cfg(test)] +mod tests { + use super::*; + + async fn assert_fails_closed(input_name: &str, contents: &[u8]) { + let dir = std::env::temp_dir().join(format!("ruvllm-968-{}", uuid_like())); + std::fs::create_dir_all(&dir).unwrap(); + let input = dir.join(input_name); + std::fs::write(&input, contents).unwrap(); + let output = dir.join("out.gguf"); + + let err = run( + input.to_str().unwrap(), + output.to_str().unwrap(), + "q4_k_m", + true, + true, + true, + false, + dir.to_str().unwrap(), + ) + .await + .expect_err("quantize must fail closed until a GGUF writer exists"); + assert!(err.to_string().contains("not implemented"), "{err}"); + assert!(!output.exists(), "no output artifact may be created"); + std::fs::remove_dir_all(&dir).unwrap(); + } + + fn uuid_like() -> String { + format!( + "{}-{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos() + ) + } + + #[tokio::test] + async fn safetensors_input_fails_closed() { + assert_fails_closed("model.safetensors", b"{}").await; + } + + #[tokio::test] + async fn gguf_input_fails_closed() { + assert_fails_closed("model.gguf", b"GGUF").await; + } + + #[tokio::test] + async fn unknown_format_is_rejected() { + let err = run("x", "", "q3_zz", true, true, true, false, ".") + .await + .expect_err("bad format"); + assert!(err.to_string().contains("Unknown quantization format")); + } +} diff --git a/crates/ruvllm-cli/src/main.rs b/crates/ruvllm-cli/src/main.rs index d8df18b1da..140d098266 100644 --- a/crates/ruvllm-cli/src/main.rs +++ b/crates/ruvllm-cli/src/main.rs @@ -11,7 +11,7 @@ //! - `ruvllm serve ` - Start inference server //! - `ruvllm chat ` - Interactive chat mode //! - `ruvllm benchmark ` - Run performance benchmarks -//! - `ruvllm quantize ` - Quantize model to GGUF format +//! - `ruvllm quantize ` - Quantize model to GGUF format (not yet implemented; fails closed, #968) use clap::{Parser, Subcommand}; use colored::Colorize; @@ -177,14 +177,11 @@ enum Commands { format: String, }, - /// Quantize a model to GGUF format + /// Quantize a model to GGUF format (NOT YET IMPLEMENTED) /// - /// Supports Q4_K_M (4-bit), Q5_K_M (5-bit), and Q8_0 (8-bit) quantization. - /// Optimized for Apple Neural Engine (ANE) inference on M4 Pro. - /// - /// Examples: - /// ruvllm quantize --model qwen-0.5b --output ruvltra-small-q4.gguf --quant q4_k_m - /// ruvllm quantize --model ./model.safetensors --quant q8_0 --ane-optimize + /// The GGUF tensor writer is not implemented yet: this command validates + /// its arguments and exits with an error without writing any output (#968). + /// Use llama.cpp's convert_hf_to_gguf.py + llama-quantize meanwhile. #[command(alias = "quant")] Quantize { /// Model to quantize (path or HuggingFace ID) diff --git a/npm/package-lock.json b/npm/package-lock.json index 6cc0f0bcbb..47557ccb6b 100644 --- a/npm/package-lock.json +++ b/npm/package-lock.json @@ -8145,9 +8145,9 @@ "license": "MIT" }, "node_modules/fast-uri": { - "version": "3.1.7", - "resolved": "https://registry.npmjs.org/fast-uri/-/fast-uri-3.1.7.tgz", - "integrity": "sha512-dOvZVzjdZdz7phd9v6jCbwxrBW3fK6n8Rc0CtdmM4bumzMnxywBYhuph6J819RRw/ku+rLbelwfMunktuzVVHg==", + "version": "3.1.8", + "resolved": "https://registry.npmjs.org/fast-uri/-/fast-uri-3.1.8.tgz", + "integrity": "sha512-GZMtZUTNRpOVIECoXwLNZS5xUGE+mVNbTB8h/7Rwh2TFWcBQiPzTgyZi05BF9UMZKkLJv8XBRJTlU7zg8+ZfMg==", "funding": [ { "type": "github", diff --git a/npm/packages/graph-node/package.json b/npm/packages/graph-node/package.json index c937069b35..9ef25ba2de 100644 --- a/npm/packages/graph-node/package.json +++ b/npm/packages/graph-node/package.json @@ -1,6 +1,6 @@ { "name": "@ruvector/graph-node", - "version": "2.1.0", + "version": "2.1.1", "description": "Native Node.js bindings for RuVector Graph Database with hypergraph support, Cypher queries, and persistence - 10x faster than WASM", "main": "index.js", "types": "index.d.ts", @@ -33,11 +33,11 @@ "@napi-rs/cli": "^2.18.0" }, "optionalDependencies": { - "@ruvector/graph-node-linux-x64-gnu": "2.1.0", - "@ruvector/graph-node-linux-arm64-gnu": "2.1.0", - "@ruvector/graph-node-darwin-x64": "2.1.0", - "@ruvector/graph-node-darwin-arm64": "2.1.0", - "@ruvector/graph-node-win32-x64-msvc": "2.1.0" + "@ruvector/graph-node-linux-x64-gnu": "2.1.1", + "@ruvector/graph-node-linux-arm64-gnu": "2.1.1", + "@ruvector/graph-node-darwin-x64": "2.1.1", + "@ruvector/graph-node-darwin-arm64": "2.1.1", + "@ruvector/graph-node-win32-x64-msvc": "2.1.1" }, "publishConfig": { "access": "public" diff --git a/npm/packages/graph-node/scripts/publish-platforms.js b/npm/packages/graph-node/scripts/publish-platforms.js index baa0fecf56..6f464e015e 100755 --- a/npm/packages/graph-node/scripts/publish-platforms.js +++ b/npm/packages/graph-node/scripts/publish-platforms.js @@ -20,12 +20,31 @@ const version = require(path.join(rootDir, 'package.json')).version; console.log('Publishing @ruvector/graph-node platform packages v' + version + '\n'); +// A missing binary or a failed publish must fail the job: the main package's +// optionalDependencies pin every platform at this exact version, so a skipped +// or swallowed platform ships an uninstallable release (same class as #1007). +const missing = platforms.filter((p) => !fs.existsSync(path.join(rootDir, p.nodeFile))); +if (missing.length) { + console.error('Missing platform binaries: ' + missing.map((p) => p.nodeFile).join(', ')); + process.exit(1); +} + +function isPublished(pkgName) { + try { + execSync('npm view ' + pkgName + '@' + version + ' version', { stdio: 'pipe' }); + return true; + } catch { + return false; + } +} + +let failed = 0; for (const platform of platforms) { const pkgName = '@ruvector/graph-node-' + platform.name; const nodeFile = path.join(rootDir, platform.nodeFile); - if (!fs.existsSync(nodeFile)) { - console.log('Skipping ' + pkgName + ' - ' + platform.nodeFile + ' not found'); + if (isPublished(pkgName)) { + console.log(pkgName + '@' + version + ' already on npm - skipping\n'); continue; } @@ -69,10 +88,15 @@ for (const platform of platforms) { console.log('Published ' + pkgName + '@' + version + '\n'); } catch (e) { console.error('Failed to publish ' + pkgName + ': ' + e.message + '\n'); + failed++; } // Cleanup fs.rmSync(tmpDir, { recursive: true, force: true }); } +if (failed) { + console.error(failed + ' platform package(s) failed to publish'); + process.exit(1); +} console.log('Done!'); diff --git a/npm/packages/ruvector/bin/mcp-server.js b/npm/packages/ruvector/bin/mcp-server.js index 423c7432a1..ab1aa20f72 100644 --- a/npm/packages/ruvector/bin/mcp-server.js +++ b/npm/packages/ruvector/bin/mcp-server.js @@ -200,6 +200,23 @@ try { } // Intelligence class with full RuVector stack support +/** + * Write via temp file + rename so a concurrent reader (every Claude Code hook is + * a separate `ruvector hooks …` process reading the same store) sees either the + * old or the new file, never a torn one (#995; same as bin/cli.js, #634). + */ +function atomicWriteFileSync(filePath, data) { + const dir = path.dirname(filePath); + const tmp = path.join(dir, `.${path.basename(filePath)}.tmp.${process.pid}.${Date.now()}`); + try { + fs.writeFileSync(tmp, data); + fs.renameSync(tmp, filePath); + } catch (err) { + try { fs.rmSync(tmp, { force: true }); } catch { /* best-effort temp cleanup */ } + throw err; + } +} + class Intelligence { constructor() { this.intelPath = this.getIntelPath(); @@ -259,18 +276,35 @@ class Intelligence { } load() { + const defaults = () => ({ patterns: {}, memories: [], trajectories: [], errors: {}, agents: {}, edges: [] }); + // A missing store is a legitimate fresh start. + if (!fs.existsSync(this.intelPath)) return defaults(); + let data; try { - if (fs.existsSync(this.intelPath)) { - const data = JSON.parse(fs.readFileSync(this.intelPath, 'utf-8')); - // Untrusted on-disk input (ADR-210 security pass): a corrupted or - // hand-edited store must not crash array/object consumers. - if (data && typeof data === 'object' && !Array.isArray(data)) { - if (!Array.isArray(data.memories)) data.memories = []; - return data; - } - } - } catch {} - return { patterns: {}, memories: [], trajectories: [], errors: {}, agents: {}, edges: [] }; + data = JSON.parse(fs.readFileSync(this.intelPath, 'utf-8')); + } catch (err) { + // #995: never swallow a corrupt read into empty defaults — the next save() + // would write the emptiness over the real store. Quarantine the file (same + // naming as bin/cli.js) so it is preserved, say so on stderr, and start + // fresh. Unlike the CLI we do not throw: a long-lived MCP server that + // refuses to start is worse than one that starts empty with the original + // data set aside for restore. + const quarantine = `${this.intelPath}.corrupt-${Date.now()}`; + let moved = false; + try { fs.renameSync(this.intelPath, quarantine); moved = true; } catch { /* reported below */ } + console.error( + `ruvector: intelligence store at ${this.intelPath} is corrupt (${err.message}); ` + + (moved + ? `quarantined to ${quarantine} — restore it or delete it. Starting with an empty store.` + : `could not quarantine it; starting with an empty store.`) + ); + return defaults(); + } + // Untrusted on-disk input (ADR-210 security pass): a corrupted or + // hand-edited store must not crash array/object consumers. + if (!data || typeof data !== 'object' || Array.isArray(data)) return defaults(); + if (!Array.isArray(data.memories)) data.memories = []; + return data; } // ========================================================================== @@ -393,7 +427,7 @@ class Intelligence { } catch {} } - fs.writeFileSync(this.intelPath, JSON.stringify(this.data, null, 2)); + atomicWriteFileSync(this.intelPath, JSON.stringify(this.data, null, 2)); } stats() { diff --git a/npm/packages/ruvector/package.json b/npm/packages/ruvector/package.json index c42b7702f4..d0bad9492c 100644 --- a/npm/packages/ruvector/package.json +++ b/npm/packages/ruvector/package.json @@ -1,6 +1,6 @@ { "name": "ruvector", - "version": "0.3.2", + "version": "0.3.3", "description": "Self-learning vector database for Node.js — hybrid search, Graph RAG, FlashAttention-3, HNSW, 50+ attention mechanisms", "main": "dist/index.js", "types": "dist/index.d.ts", @@ -12,7 +12,7 @@ "verify-dist": "node scripts/verify-dist.js", "prepack": "npm run build && npm run verify-dist", "prepublishOnly": "npm run build && npm run verify-dist", - "test": "node test/integration.js && node test/backend-fallback.js && node test/mincut-wasm.js && node test/metaharness-sdk.js && node test/metaharness-optional-deps.js && node test/cli-commands.js && node test/metaharness-cli.js && node test/db-workflow.js && node test/mcp-stdio.js && node test/metaharness-mcp.js && node test/sigterm-cleanup.js && node test/mcp-policy.js && node test/mcp-command-security.js && node test/mcp-handshake.js && node test/startup-budget.js" + "test": "node test/integration.js && node test/backend-fallback.js && node test/mincut-wasm.js && node test/metaharness-sdk.js && node test/metaharness-optional-deps.js && node test/cli-commands.js && node test/metaharness-cli.js && node test/db-workflow.js && node test/mcp-stdio.js && node test/metaharness-mcp.js && node test/sigterm-cleanup.js && node test/mcp-policy.js && node test/mcp-command-security.js && node test/mcp-handshake.js && node test/mcp-intel-durability.js && node test/startup-budget.js" }, "keywords": [ "vector", diff --git a/npm/packages/ruvector/test/mcp-intel-durability.js b/npm/packages/ruvector/test/mcp-intel-durability.js new file mode 100644 index 0000000000..0077de4e3c --- /dev/null +++ b/npm/packages/ruvector/test/mcp-intel-durability.js @@ -0,0 +1,105 @@ +#!/usr/bin/env node + +/** + * Regression coverage for #995: the MCP server's own Intelligence copy must + * (a) quarantine a corrupt intelligence.json instead of silently treating it as + * empty and later overwriting it, and (b) save through temp-file + rename so a + * concurrent hook process never reads a torn file. + */ + +const assert = require('assert'); +const fs = require('fs'); +const os = require('os'); +const path = require('path'); +const { spawn } = require('child_process'); + +const MCP_SERVER = path.join(__dirname, '..', 'bin', 'mcp-server.js'); + +function rpc(child, id, method, params, timeoutMs = 20_000) { + return new Promise((resolve, reject) => { + let buf = ''; + const timer = setTimeout(() => reject(new Error(`timeout waiting for ${method}`)), timeoutMs); + const onData = (chunk) => { + buf += chunk.toString(); + for (const line of buf.split(/\r?\n/).filter(Boolean)) { + let msg; + try { msg = JSON.parse(line); } catch { continue; } + if (msg.id === id) { + clearTimeout(timer); + child.stdout.off('data', onData); + resolve(msg); + return; + } + } + }; + child.stdout.on('data', onData); + child.stdin.write(`${JSON.stringify({ jsonrpc: '2.0', id, method, params })}\n`); + }); +} + +async function main() { + const source = fs.readFileSync(MCP_SERVER, 'utf8'); + assert.doesNotMatch( + source, + /fs\.writeFileSync\(this\.intelPath/, + 'mcp-server.js must not write intelligence.json with a plain writeFileSync', + ); + + const tmp = fs.mkdtempSync(path.join(os.tmpdir(), 'ruv995-')); + const store = path.join(tmp, '.ruvector', 'intelligence.json'); + fs.mkdirSync(path.dirname(store), { recursive: true }); + const torn = '{"memories": [{"content": "keep me"}, {"content": "trunc'; + fs.writeFileSync(store, torn); + + const child = spawn(process.execPath, [MCP_SERVER], { + cwd: tmp, + env: { ...process.env, NO_COLOR: '1', HOME: tmp }, + stdio: ['pipe', 'pipe', 'pipe'], + }); + let stderr = ''; + child.stderr.on('data', (c) => { stderr += c.toString(); }); + + try { + const init = await rpc(child, 1, 'initialize', { + protocolVersion: '2024-11-05', + capabilities: {}, + clientInfo: { name: 'ruvector-995-regression', version: '1.0.0' }, + }); + assert.ok(init.result, `initialize failed: ${JSON.stringify(init)}`); + + const call = await rpc(child, 2, 'tools/call', { + name: 'hooks_remember', + arguments: { content: 'after quarantine', type: 'test' }, + }); + assert.ok(call.result, `hooks_remember failed: ${JSON.stringify(call)}`); + + const entries = fs.readdirSync(path.dirname(store)); + const quarantined = entries.filter((f) => f.startsWith('intelligence.json.corrupt-')); + assert.strictEqual(quarantined.length, 1, `expected one quarantine file, got ${entries.join(', ')}`); + assert.strictEqual( + fs.readFileSync(path.join(path.dirname(store), quarantined[0]), 'utf8'), + torn, + 'quarantined file must hold the original bytes', + ); + assert.match(stderr, /is corrupt/, 'corrupt load must be reported on stderr'); + + const saved = JSON.parse(fs.readFileSync(store, 'utf8')); + assert.ok(Array.isArray(saved.memories), 'new store must be valid JSON with memories[]'); + assert.deepStrictEqual( + entries.filter((f) => f.includes('.tmp.')), + [], + 'atomic save must not leave temp files behind', + ); + } finally { + child.stdin.end(); + child.kill('SIGTERM'); + fs.rmSync(tmp, { recursive: true, force: true }); + } + + console.log('MCP intelligence store durability checks passed (#995)'); +} + +main().catch((error) => { + console.error(error); + process.exitCode = 1; +});