diff --git a/crates/noxa-rag/Cargo.toml b/crates/noxa-rag/Cargo.toml new file mode 100644 index 0000000..d7e8287 --- /dev/null +++ b/crates/noxa-rag/Cargo.toml @@ -0,0 +1,66 @@ +[package] +name = "noxa-rag" +description = "RAG pipeline for noxa — TEI embeddings + Qdrant vector store" +version.workspace = true +edition.workspace = true +license.workspace = true + +[[bin]] +name = "noxa-rag-daemon" +path = "src/bin/noxa-rag-daemon.rs" + +[dependencies] +noxa-core = { workspace = true } + +# Async runtime +tokio = { workspace = true } + +# Serialization +serde = { workspace = true } +serde_json = { workspace = true } +toml = "0.8" + +# Error handling +thiserror = { workspace = true } + +# Tracing +tracing = { workspace = true } +tracing-subscriber = { workspace = true } + +# Async traits +async-trait = "0.1" + +# HTTP client (plain reqwest — no primp patches needed for TEI/Qdrant) +reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } + +# Vector store — REST transport only (reqwest feature, no protoc required) +qdrant-client = { version = "1", default-features = false, features = ["reqwest"] } + +# Chunking +text-splitter = { version = "0.25", features = ["markdown", "tokenizers"] } +tokenizers = "0.21" + +# UUID v5 for deterministic point IDs +uuid = { version = "1", features = ["v5", "serde"] } + +# Filesystem watcher +notify = "6" +notify-debouncer-mini = "0.4" + +# Concurrent data structures +dashmap = "6" + +# URL parsing +url = "2" + +# CLI args +clap = { workspace = true } + +# Date/time for failed-jobs log +chrono = { version = "0.4", features = ["serde"] } + +# CancellationToken for coordinated shutdown +tokio-util = { version = "0.7", features = ["io"] } + +[dev-dependencies] +tokio = { workspace = true } diff --git a/crates/noxa-rag/src/bin/noxa-rag-daemon.rs b/crates/noxa-rag/src/bin/noxa-rag-daemon.rs new file mode 100644 index 0000000..6f8c76e --- /dev/null +++ b/crates/noxa-rag/src/bin/noxa-rag-daemon.rs @@ -0,0 +1,5 @@ +// noxa-rag-daemon — implemented in noxa-68r.8 +fn main() { + eprintln!("[noxa-rag] daemon stub — not yet implemented (noxa-68r.8)"); + std::process::exit(1); +} diff --git a/crates/noxa-rag/src/chunker.rs b/crates/noxa-rag/src/chunker.rs new file mode 100644 index 0000000..76b1ab5 --- /dev/null +++ b/crates/noxa-rag/src/chunker.rs @@ -0,0 +1,13 @@ +// Chunker — implemented in noxa-68r.2 +use noxa_core::types::ExtractionResult; +use crate::config::ChunkerConfig; +use crate::types::Chunk; + +pub fn chunk( + _result: &ExtractionResult, + _config: &ChunkerConfig, + _tokenizer: &tokenizers::Tokenizer, +) -> Vec { + // Full implementation in noxa-68r.2 + vec![] +} diff --git a/crates/noxa-rag/src/config.rs b/crates/noxa-rag/src/config.rs new file mode 100644 index 0000000..5ba6814 --- /dev/null +++ b/crates/noxa-rag/src/config.rs @@ -0,0 +1,145 @@ +use serde::Deserialize; +use std::path::{Path, PathBuf}; + +use crate::error::RagError; + +/// Top-level configuration deserialized from noxa-rag.toml. +#[derive(Debug, Clone, Deserialize)] +pub struct RagConfig { + pub source: SourceConfig, + pub embed_provider: EmbedProviderConfig, + pub vector_store: VectorStoreConfig, + pub chunker: ChunkerConfig, + pub pipeline: PipelineConfig, + /// UUID namespace for deterministic point IDs. + /// Default: 6ba7b810-9dad-11d1-80b4-00c04fd430c8 + #[serde(default = "default_uuid_namespace")] + pub uuid_namespace: uuid::Uuid, +} + +fn default_uuid_namespace() -> uuid::Uuid { + uuid::Uuid::parse_str("6ba7b810-9dad-11d1-80b4-00c04fd430c8").unwrap() +} + +#[derive(Debug, Clone, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum SourceConfig { + FsWatcher { + watch_dir: PathBuf, + #[serde(default = "default_debounce_ms")] + debounce_ms: u64, + }, +} + +fn default_debounce_ms() -> u64 { + 500 +} + +#[derive(Debug, Clone, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum EmbedProviderConfig { + Tei { + url: String, + model: String, + /// Optional: load tokenizer from local path (avoids HF Hub at startup). + local_path: Option, + }, + OpenAi { + api_key: String, + model: String, + }, + VoyageAi { + api_key: String, + model: String, + }, +} + +#[derive(Debug, Clone, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum VectorStoreConfig { + Qdrant { + /// REST URL — port 6333, NOT 6334 (gRPC). + url: String, + collection: String, + /// Optional API key. Override with NOXA_RAG_QDRANT_API_KEY env var. + api_key: Option, + }, + /// Dev/test only — factory returns RagError::Config("not implemented"). + InMemory, +} + +#[derive(Debug, Clone, Deserialize)] +pub struct ChunkerConfig { + #[serde(default = "default_target_tokens")] + pub target_tokens: usize, + #[serde(default = "default_overlap_tokens")] + pub overlap_tokens: usize, + #[serde(default = "default_min_words")] + pub min_words: usize, + #[serde(default = "default_max_chunks_per_page")] + pub max_chunks_per_page: usize, +} + +impl Default for ChunkerConfig { + fn default() -> Self { + Self { + target_tokens: default_target_tokens(), + overlap_tokens: default_overlap_tokens(), + min_words: default_min_words(), + max_chunks_per_page: default_max_chunks_per_page(), + } + } +} + +fn default_target_tokens() -> usize { 512 } +fn default_overlap_tokens() -> usize { 64 } +fn default_min_words() -> usize { 50 } +fn default_max_chunks_per_page() -> usize { 100 } + +#[derive(Debug, Clone, Deserialize)] +pub struct PipelineConfig { + #[serde(default = "default_embed_concurrency")] + pub embed_concurrency: usize, + /// MUST be an absolute path — systemd daemon runs with CWD = /. + pub failed_jobs_log: Option, +} + +impl Default for PipelineConfig { + fn default() -> Self { + Self { + embed_concurrency: default_embed_concurrency(), + failed_jobs_log: None, + } + } +} + +fn default_embed_concurrency() -> usize { 4 } + +/// Load and validate config from a TOML file. +pub fn load_config(path: &Path) -> Result { + let content = std::fs::read_to_string(path) + .map_err(|e| RagError::Config(format!("cannot read config file {}: {}", path.display(), e)))?; + + let config: RagConfig = toml::from_str(&content) + .map_err(|e| RagError::Config(format!("config parse error: {}", e)))?; + + // Validate embed_concurrency > 0 + if config.pipeline.embed_concurrency == 0 { + return Err(RagError::Config( + "pipeline.embed_concurrency must be > 0 (0 causes Semaphore deadlock)".to_string(), + )); + } + + // Validate failed_jobs_log is absolute if set + if let Some(ref log_path) = config.pipeline.failed_jobs_log { + if !log_path.is_absolute() { + return Err(RagError::Config(format!( + "pipeline.failed_jobs_log must be an absolute path (got: {}). \ + systemd daemon runs with CWD = / and relative paths resolve there.", + log_path.display() + ))); + } + } + + Ok(config) +} diff --git a/crates/noxa-rag/src/embed/mod.rs b/crates/noxa-rag/src/embed/mod.rs new file mode 100644 index 0000000..b275add --- /dev/null +++ b/crates/noxa-rag/src/embed/mod.rs @@ -0,0 +1,19 @@ +use async_trait::async_trait; +use std::sync::Arc; + +use crate::error::RagError; + +/// Pluggable embedding provider. +/// +/// Trait surface is minimal by design — only what ALL impls share. +/// `is_available()` and `dimensions()` are concrete methods on each provider struct, +/// called during factory startup probes (not via dyn dispatch). +#[async_trait] +pub trait EmbedProvider: Send + Sync { + async fn embed(&self, texts: &[String]) -> Result>, RagError>; +} + +pub type DynEmbedProvider = Arc; + +pub mod tei; +pub use tei::TeiProvider; diff --git a/crates/noxa-rag/src/embed/tei.rs b/crates/noxa-rag/src/embed/tei.rs new file mode 100644 index 0000000..1004a26 --- /dev/null +++ b/crates/noxa-rag/src/embed/tei.rs @@ -0,0 +1,39 @@ +// TeiProvider — implemented in noxa-68r.3 +use async_trait::async_trait; +use crate::embed::EmbedProvider; +use crate::error::RagError; + +pub struct TeiProvider { + pub(crate) client: reqwest::Client, + pub(crate) url: String, + pub(crate) model: String, + pub(crate) dimensions: usize, +} + +impl TeiProvider { + pub async fn is_available(&self) -> bool { + self.client + .get(format!("{}/health", self.url)) + .timeout(std::time::Duration::from_secs(2)) + .send() + .await + .map(|r| r.status().is_success()) + .unwrap_or(false) + } + + pub fn dimensions(&self) -> usize { + self.dimensions + } + + pub fn name(&self) -> &str { + "tei" + } +} + +#[async_trait] +impl EmbedProvider for TeiProvider { + async fn embed(&self, _texts: &[String]) -> Result>, RagError> { + // Full implementation in noxa-68r.3 + Err(RagError::Embed("TeiProvider not yet implemented".to_string())) + } +} diff --git a/crates/noxa-rag/src/error.rs b/crates/noxa-rag/src/error.rs new file mode 100644 index 0000000..df6dbce --- /dev/null +++ b/crates/noxa-rag/src/error.rs @@ -0,0 +1,22 @@ +use thiserror::Error; + +#[non_exhaustive] +#[derive(Debug, Error)] +pub enum RagError { + #[error("embed error: {0}")] + Embed(String), + #[error("store error: {0}")] + Store(String), + #[error("chunk error: {0}")] + Chunk(String), + #[error("config error: {0}")] + Config(String), + #[error("io error: {0}")] + Io(#[from] std::io::Error), + #[error("http error: {0}")] + Http(#[from] reqwest::Error), + #[error("json error: {0}")] + Json(#[from] serde_json::Error), + #[error("error: {0}")] + Generic(String), +} diff --git a/crates/noxa-rag/src/factory.rs b/crates/noxa-rag/src/factory.rs new file mode 100644 index 0000000..ce33ab6 --- /dev/null +++ b/crates/noxa-rag/src/factory.rs @@ -0,0 +1,18 @@ +// Factory — implemented in noxa-68r.6 +use crate::config::RagConfig; +use crate::embed::DynEmbedProvider; +use crate::error::RagError; +use crate::store::DynVectorStore; + +pub async fn build_embed_provider(_config: &RagConfig) -> Result { + // Full implementation in noxa-68r.6 + Err(RagError::Config("factory not yet implemented".to_string())) +} + +pub async fn build_vector_store( + _config: &RagConfig, + _embed_dims: usize, +) -> Result { + // Full implementation in noxa-68r.6 + Err(RagError::Config("factory not yet implemented".to_string())) +} diff --git a/crates/noxa-rag/src/lib.rs b/crates/noxa-rag/src/lib.rs new file mode 100644 index 0000000..f37a7a4 --- /dev/null +++ b/crates/noxa-rag/src/lib.rs @@ -0,0 +1,40 @@ +/// noxa-rag — RAG pipeline crate. +/// +/// Watches noxa output directory for ExtractionResult JSON files, +/// chunks them, embeds via TEI, and upserts to Qdrant. +/// +/// # Crate structure +/// - `embed` — EmbedProvider trait + TeiProvider impl +/// - `store` — VectorStore trait + QdrantStore impl +/// - `chunker` — ExtractionResult → Vec +/// - `config` — RagConfig (TOML deserialization) +/// - `factory` — build_embed_provider / build_vector_store +/// - `pipeline` — filesystem watcher orchestration +/// - `error` — RagError enum + +// Tokenizer Sync compile-time assertion. +// tokenizers::Tokenizer must be Sync to be used across tokio workers. +// If this fails to compile, workers cannot safely share the tokenizer. +const _: () = { + fn assert_sync() {} + fn _check() { + assert_sync::(); + } +}; + +pub mod chunker; +pub mod config; +pub mod embed; +pub mod error; +pub mod factory; +pub mod pipeline; +pub mod store; +pub mod types; + +// Re-export most-used types at crate root +pub use config::{load_config, RagConfig}; +pub use embed::{DynEmbedProvider, EmbedProvider}; +pub use error::RagError; +pub use factory::{build_embed_provider, build_vector_store}; +pub use store::{DynVectorStore, VectorStore}; +pub use types::{Chunk, Point, PointPayload, SearchResult}; diff --git a/crates/noxa-rag/src/pipeline.rs b/crates/noxa-rag/src/pipeline.rs new file mode 100644 index 0000000..c911bae --- /dev/null +++ b/crates/noxa-rag/src/pipeline.rs @@ -0,0 +1,24 @@ +// Pipeline — implemented in noxa-68r.7 +use std::sync::Arc; +use tokio_util::sync::CancellationToken; + +use crate::config::RagConfig; +use crate::embed::DynEmbedProvider; +use crate::error::RagError; +use crate::store::DynVectorStore; + +pub struct Pipeline { + pub config: RagConfig, + pub embed: DynEmbedProvider, + pub store: DynVectorStore, + pub tokenizer: Arc, + pub shutdown: CancellationToken, +} + +impl Pipeline { + pub async fn run(&self) -> Result<(), RagError> { + // Full implementation in noxa-68r.7 + self.shutdown.cancelled().await; + Ok(()) + } +} diff --git a/crates/noxa-rag/src/store/mod.rs b/crates/noxa-rag/src/store/mod.rs new file mode 100644 index 0000000..f9d7aaf --- /dev/null +++ b/crates/noxa-rag/src/store/mod.rs @@ -0,0 +1,23 @@ +use async_trait::async_trait; +use std::sync::Arc; + +use crate::error::RagError; +use crate::types::{Point, SearchResult}; + +/// Pluggable vector store backend. +/// +/// Trait surface is minimal — only what ALL impls share. +/// Collection lifecycle (create_collection, collection_exists) lives in factory.rs +/// as concrete methods on each store struct, called during startup probes. +#[async_trait] +pub trait VectorStore: Send + Sync { + async fn upsert(&self, points: Vec) -> Result<(), RagError>; + async fn delete_by_url(&self, url: &str) -> Result<(), RagError>; + async fn search(&self, vector: &[f32], limit: usize) -> Result, RagError>; + fn name(&self) -> &str; +} + +pub type DynVectorStore = Arc; + +pub mod qdrant; +pub use qdrant::QdrantStore; diff --git a/crates/noxa-rag/src/store/qdrant.rs b/crates/noxa-rag/src/store/qdrant.rs new file mode 100644 index 0000000..cadb063 --- /dev/null +++ b/crates/noxa-rag/src/store/qdrant.rs @@ -0,0 +1,30 @@ +// QdrantStore — implemented in noxa-68r.4 +use async_trait::async_trait; +use crate::error::RagError; +use crate::store::VectorStore; +use crate::types::{Point, SearchResult}; + +pub struct QdrantStore { + pub(crate) client: qdrant_client::Qdrant, + pub(crate) collection: String, +} + +#[async_trait] +impl VectorStore for QdrantStore { + async fn upsert(&self, _points: Vec) -> Result<(), RagError> { + // Full implementation in noxa-68r.4 + Err(RagError::Store("QdrantStore not yet implemented".to_string())) + } + + async fn delete_by_url(&self, _url: &str) -> Result<(), RagError> { + Err(RagError::Store("QdrantStore not yet implemented".to_string())) + } + + async fn search(&self, _vector: &[f32], _limit: usize) -> Result, RagError> { + Err(RagError::Store("QdrantStore not yet implemented".to_string())) + } + + fn name(&self) -> &str { + "qdrant" + } +} diff --git a/crates/noxa-rag/src/types.rs b/crates/noxa-rag/src/types.rs new file mode 100644 index 0000000..bb26127 --- /dev/null +++ b/crates/noxa-rag/src/types.rs @@ -0,0 +1,45 @@ +use serde::{Deserialize, Serialize}; +use uuid::Uuid; + +/// A chunk produced from an ExtractionResult. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Chunk { + pub text: String, + pub source_url: String, + pub domain: String, + pub chunk_index: usize, + pub total_chunks: usize, + pub char_offset: usize, + pub token_estimate: usize, +} + +/// A point ready for upsert into the vector store. +#[derive(Debug, Clone)] +pub struct Point { + /// UUID v5 deterministic ID: url#chunkN + pub id: Uuid, + pub vector: Vec, + pub payload: PointPayload, +} + +/// Payload stored alongside each vector in the store. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PointPayload { + pub text: String, + /// Normalized URL (strip fragment, trailing slash, lowercase scheme+host). + pub url: String, + pub domain: String, + pub chunk_index: usize, + pub total_chunks: usize, + pub token_estimate: usize, +} + +/// A search result returned by VectorStore::search(). +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SearchResult { + pub text: String, + pub url: String, + pub score: f32, + pub chunk_index: usize, + pub token_estimate: usize, +}