diff --git a/Cargo.lock b/Cargo.lock index ea4b029ce..96d42bbfa 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -43,6 +43,17 @@ dependencies = [ "subtle", ] +[[package]] +name = "ahash" +version = "0.7.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "891477e0c6a8957309ee5c45a6368af3ae14bb510732d2684ffa19af310920f9" +dependencies = [ + "getrandom 0.2.17", + "once_cell", + "version_check", +] + [[package]] name = "ahash" version = "0.8.12" @@ -66,6 +77,12 @@ dependencies = [ "memchr", ] +[[package]] +name = "aliasable" +version = "0.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "250f629c0161ad8107cf89319e990051fae62832fd343083bea452d93e2205fd" + [[package]] name = "aligned" version = "0.4.3" @@ -542,6 +559,20 @@ dependencies = [ "zeroize", ] +[[package]] +name = "bigdecimal" +version = "0.4.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4d6867f1565b3aad85681f1015055b087fcfd840d6aeee6eee7f2da317603695" +dependencies = [ + "autocfg", + "libm", + "num-bigint", + "num-integer", + "num-traits", + "serde", +] + [[package]] name = "bit-set" version = "0.8.0" @@ -634,6 +665,29 @@ version = "0.2.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dc0b364ead1874514c8c2855ab558056ebfeb775653e7ae45ff72f28f8f3166c" +[[package]] +name = "borsh" +version = "1.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d1da5ab77c1437701eeff7c88d968729e7766172279eab0676857b3d63af7a6f" +dependencies = [ + "borsh-derive", + "cfg_aliases", +] + +[[package]] +name = "borsh-derive" +version = "1.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0686c856aa6aac0c4498f936d7d6a02df690f614c03e4d906d1018062b5c5e2c" +dependencies = [ + "once_cell", + "proc-macro-crate", + "proc-macro2", + "quote", + "syn 2.0.114", +] + [[package]] name = "bstr" version = "1.12.1" @@ -656,6 +710,28 @@ version = "3.19.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5dd9dc738b7a8311c7ade152424974d8115f2cdad61e8dab8dac9f2362298510" +[[package]] +name = "bytecheck" +version = "0.6.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "23cdc57ce23ac53c931e88a43d06d070a6fd142f2617be5855eb75efc9beb1c2" +dependencies = [ + "bytecheck_derive", + "ptr_meta", + "simdutf8", +] + +[[package]] +name = "bytecheck_derive" +version = "0.6.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3db406d29fbcd95542e92559bed4d8ad92636d1ca8b3b72ede10b4bcc010e659" +dependencies = [ + "proc-macro2", + "quote", + "syn 1.0.109", +] + [[package]] name = "bytecount" version = "0.6.9" @@ -2175,6 +2251,9 @@ name = "hashbrown" version = "0.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" +dependencies = [ + "ahash 0.7.8", +] [[package]] name = "hashbrown" @@ -2743,6 +2822,17 @@ dependencies = [ "web-time", ] +[[package]] +name = "inherent" +version = "1.0.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c727f80bfa4a6c6e2508d2f05b6f4bfce242030bd88ed15ae5331c5b5d30fba7" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.114", +] + [[package]] name = "inlinable_string" version = "0.1.15" @@ -2898,7 +2988,7 @@ version = "0.42.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6103dcd3815c9bd04239a6d0343dac2b7a18ee260e800d0a6e3eb5b9333504fb" dependencies = [ - "ahash", + "ahash 0.8.12", "bytecount", "data-encoding", "email_address", @@ -3240,7 +3330,6 @@ dependencies = [ "linkme", "mockall", "regex", - "reqwest 0.13.2", "rstest", "schemars 1.2.1", "serde", @@ -3363,6 +3452,7 @@ dependencies = [ "reqwest 0.13.2", "rstest", "schemars 1.2.1", + "sea-orm", "serde", "serde_json", "sha2", @@ -4055,6 +4145,15 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "04744f49eae99ab78e0d5c0b603ab218f515ea8cfe5a456d7629ad883a3b6e7d" +[[package]] +name = "ordered-float" +version = "4.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7bb71e1b3fa6ca1c61f383464aaf2bb0e2f8e772a1f01d486832464de363b951" +dependencies = [ + "num-traits", +] + [[package]] name = "ort" version = "2.0.0-rc.11" @@ -4080,6 +4179,30 @@ dependencies = [ "ureq 3.2.0", ] +[[package]] +name = "ouroboros" +version = "0.18.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e0f050db9c44b97a94723127e6be766ac5c340c48f2c4bb3ffa11713744be59" +dependencies = [ + "aliasable", + "ouroboros_macro", + "static_assertions", +] + +[[package]] +name = "ouroboros_macro" +version = "0.18.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3c7028bdd3d43083f6d8d4d5187680d0d3560d54df4cc9d752005268b41e64d0" +dependencies = [ + "heck 0.4.1", + "proc-macro2", + "proc-macro2-diagnostics", + "quote", + "syn 2.0.114", +] + [[package]] name = "outref" version = "0.5.2" @@ -4259,6 +4382,15 @@ dependencies = [ "serde", ] +[[package]] +name = "pgvector" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc58e2d255979a31caa7cabfa7aac654af0354220719ab7a68520ae7a91e8c0b" +dependencies = [ + "serde", +] + [[package]] name = "pin-project" version = "1.1.10" @@ -4570,6 +4702,26 @@ dependencies = [ "prost", ] +[[package]] +name = "ptr_meta" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0738ccf7ea06b608c10564b31debd4f5bc5e197fc8bfe088f68ae5ce81e7a4f1" +dependencies = [ + "ptr_meta_derive", +] + +[[package]] +name = "ptr_meta_derive" +version = "0.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "16b845dbfca988fa33db069c0e230574d15a3088f147a87b64c7589eb662c9ac" +dependencies = [ + "proc-macro2", + "quote", + "syn 1.0.109", +] + [[package]] name = "pulldown-cmark" version = "0.13.0" @@ -4945,7 +5097,7 @@ version = "0.42.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1a76b0c93f3dd49da2900cfb3c6b883602471a0e3e34b1578a3310ca1269d418" dependencies = [ - "ahash", + "ahash 0.8.12", "fluent-uri", "getrandom 0.3.4", "hashbrown 0.16.1", @@ -4989,6 +5141,15 @@ version = "1.9.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba39f3699c378cd8970968dcbff9c43159ea4cfbd88d43c00b22f2ef10a435d2" +[[package]] +name = "rend" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "71fe3824f5629716b1589be05dacd749f6aa084c87e00e016714a8cdfccc997c" +dependencies = [ + "bytecheck", +] + [[package]] name = "reqwest" version = "0.12.28" @@ -5104,6 +5265,35 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rkyv" +version = "0.7.46" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2297bf9c81a3f0dc96bc9521370b88f054168c29826a75e89c55ff196e7ed6a1" +dependencies = [ + "bitvec", + "bytecheck", + "bytes", + "hashbrown 0.12.3", + "ptr_meta", + "rend", + "rkyv_derive", + "seahash", + "tinyvec", + "uuid", +] + +[[package]] +name = "rkyv_derive" +version = "0.7.46" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "84d7b42d4b8d06048d3ac8db0eb31bcb942cbeb709f0b5f2b2ebde398d3038f5" +dependencies = [ + "proc-macro2", + "quote", + "syn 1.0.109", +] + [[package]] name = "rmcp" version = "0.16.0" @@ -5243,6 +5433,22 @@ dependencies = [ "thiserror 2.0.18", ] +[[package]] +name = "rust_decimal" +version = "1.40.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "61f703d19852dbf87cbc513643fa81428361eb6940f1ac14fd58155d295a3eb0" +dependencies = [ + "arrayvec", + "borsh", + "bytes", + "num-traits", + "rand 0.8.5", + "rkyv", + "serde", + "serde_json", +] + [[package]] name = "rustc-hash" version = "2.1.1" @@ -5496,6 +5702,101 @@ version = "3.0.10" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "490dcfcbfef26be6800d11870ff2df8774fa6e86d047e3e8c8a76b25655e41ca" +[[package]] +name = "sea-bae" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f694a6ab48f14bc063cfadff30ab551d3c7e46d8f81836c51989d548f44a2a25" +dependencies = [ + "heck 0.4.1", + "proc-macro-error2", + "proc-macro2", + "quote", + "syn 2.0.114", +] + +[[package]] +name = "sea-orm" +version = "1.1.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6d945f62558fac19e5988680d2fdf747b734c2dbc6ce2cb81ba33ed8dde5b103" +dependencies = [ + "async-stream", + "async-trait", + "bigdecimal", + "chrono", + "derive_more", + "futures-util", + "log", + "ouroboros", + "pgvector", + "rust_decimal", + "sea-orm-macros", + "sea-query", + "sea-query-binder", + "serde", + "serde_json", + "sqlx", + "strum 0.26.3", + "thiserror 2.0.18", + "time", + "tracing", + "url", + "uuid", +] + +[[package]] +name = "sea-orm-macros" +version = "1.1.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "84c2e64a50a9cc8339f10a27577e10062c7f995488e469f2c95762c5ee847832" +dependencies = [ + "heck 0.5.0", + "proc-macro2", + "quote", + "sea-bae", + "syn 2.0.114", + "unicode-ident", +] + +[[package]] +name = "sea-query" +version = "0.32.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8a5d1c518eaf5eda38e5773f902b26ab6d5e9e9e2bb2349ca6c64cf96f80448c" +dependencies = [ + "bigdecimal", + "chrono", + "inherent", + "ordered-float", + "rust_decimal", + "serde_json", + "time", + "uuid", +] + +[[package]] +name = "sea-query-binder" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0019f47430f7995af63deda77e238c17323359af241233ec768aba1faea7608" +dependencies = [ + "bigdecimal", + "chrono", + "rust_decimal", + "sea-query", + "serde_json", + "sqlx", + "time", + "uuid", +] + +[[package]] +name = "seahash" +version = "4.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1c107b6f4780854c8b126e228ea8869f4d7b71260f962fefb57b996b8959ba6b" + [[package]] name = "security-framework" version = "2.11.1" @@ -5839,6 +6140,12 @@ dependencies = [ "quote", ] +[[package]] +name = "simdutf8" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e3a9fe34e3e7a50316060351f37187a3f546bce95496156754b601a5fa71b76e" + [[package]] name = "similar" version = "2.7.0" @@ -5942,7 +6249,9 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ee6798b1838b6a0f69c007c133b8df5866302197e404e8b6ee8ed3e3a5e68dc6" dependencies = [ "base64 0.22.1", + "bigdecimal", "bytes", + "chrono", "crc", "crossbeam-queue", "either", @@ -5958,15 +6267,20 @@ dependencies = [ "memchr", "once_cell", "percent-encoding", + "rust_decimal", + "rustls 0.23.36", "serde", "serde_json", "sha2", "smallvec", "thiserror 2.0.18", + "time", "tokio", "tokio-stream", "tracing", "url", + "uuid", + "webpki-roots 0.26.11", ] [[package]] @@ -6015,9 +6329,11 @@ checksum = "aa003f0038df784eb8fecbbac13affe3da23b45194bd57dba231c8f48199c526" dependencies = [ "atoi", "base64 0.22.1", + "bigdecimal", "bitflags 2.10.0", "byteorder", "bytes", + "chrono", "crc", "digest", "dotenvy", @@ -6038,6 +6354,7 @@ dependencies = [ "percent-encoding", "rand 0.8.5", "rsa", + "rust_decimal", "serde", "sha1", "sha2", @@ -6045,7 +6362,9 @@ dependencies = [ "sqlx-core", "stringprep", "thiserror 2.0.18", + "time", "tracing", + "uuid", "whoami", ] @@ -6057,8 +6376,10 @@ checksum = "db58fcd5a53cf07c184b154801ff91347e4c30d17a3562a635ff028ad5deda46" dependencies = [ "atoi", "base64 0.22.1", + "bigdecimal", "bitflags 2.10.0", "byteorder", + "chrono", "crc", "dotenvy", "etcetera", @@ -6073,8 +6394,10 @@ dependencies = [ "log", "md-5", "memchr", + "num-bigint", "once_cell", "rand 0.8.5", + "rust_decimal", "serde", "serde_json", "sha2", @@ -6082,7 +6405,9 @@ dependencies = [ "sqlx-core", "stringprep", "thiserror 2.0.18", + "time", "tracing", + "uuid", "whoami", ] @@ -6093,6 +6418,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2d12fe70b2c1b4401038055f90f151b78208de1f9f89a7dbfd41587a10c3eea" dependencies = [ "atoi", + "chrono", "flume", "futures-channel", "futures-core", @@ -6106,8 +6432,10 @@ dependencies = [ "serde_urlencoded", "sqlx-core", "thiserror 2.0.18", + "time", "tracing", "url", + "uuid", ] [[package]] @@ -6164,6 +6492,12 @@ version = "0.24.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "063e6045c0e62079840579a7e47a355ae92f60eb74daaf156fb1e84ba164e63f" +[[package]] +name = "strum" +version = "0.26.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8fec0f0aef304996cf250b31b5a10dee7980c85da9d759361292b8bca5a18f06" + [[package]] name = "strum" version = "0.27.2" @@ -6449,7 +6783,7 @@ version = "0.22.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b238e22d44a15349529690fb07bd645cf58149a1b1e44d6cb5bd1641ff1a6223" dependencies = [ - "ahash", + "ahash 0.8.12", "aho-corasick", "compact_str", "dary_heap", diff --git a/Cargo.toml b/Cargo.toml index 5b5c54a5d..9d78b081c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -319,6 +319,9 @@ jsonpath-rust = "1.0.4" # Lazy static values once_cell = "1.19" +# ORM +sea-orm = { version = "1", features = ["sqlx-sqlite", "runtime-tokio-rustls"] } + # ============================================ # Build Profiles # ============================================ diff --git a/crates/mcb-application/src/use_cases/indexing_service.rs b/crates/mcb-application/src/use_cases/indexing_service.rs index bff1459fe..a2f62cb10 100644 --- a/crates/mcb-application/src/use_cases/indexing_service.rs +++ b/crates/mcb-application/src/use_cases/indexing_service.rs @@ -26,13 +26,13 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Instant; -use ignore::WalkBuilder; use mcb_domain::constants::{INDEXING_STATUS_COMPLETED, INDEXING_STATUS_STARTED}; use mcb_domain::error::Result; use mcb_domain::events::DomainEvent; use mcb_domain::ports::{ - ContextServiceInterface, EventBusProvider, FileHashRepository, IndexingOperationsInterface, - IndexingResult, IndexingServiceInterface, IndexingStatus, LanguageChunkingProvider, + ContextServiceInterface, EventBusProvider, FileHashRepository, FileSystemProvider, + IndexingOperationsInterface, IndexingResult, IndexingServiceInterface, IndexingStatus, + LanguageChunkingProvider, TaskRunnerProvider, }; use mcb_domain::value_objects::{CollectionId, OperationId}; use tracing::{error, info, warn}; @@ -102,11 +102,14 @@ pub struct IndexingServiceImpl { language_chunker: Arc, indexing_ops: Arc, event_bus: Arc, + file_system_provider: Arc, + task_runner_provider: Arc, file_hash_repository: Option>, supported_extensions: Vec, } /// `IndexingServiceDeps` struct. +#[allow(missing_docs)] pub struct IndexingServiceDeps { /// Service for Context operations pub context_service: Arc, @@ -116,6 +119,8 @@ pub struct IndexingServiceDeps { pub indexing_ops: Arc, /// Event bus pub event_bus: Arc, + pub file_system_provider: Arc, + pub task_runner_provider: Arc, /// Supported file extensions pub supported_extensions: Vec, } @@ -139,6 +144,8 @@ impl IndexingServiceImpl { language_chunker: Arc, indexing_ops: Arc, event_bus: Arc, + file_system_provider: Arc, + task_runner_provider: Arc, supported_extensions: Vec, ) -> Self { Self { @@ -146,6 +153,8 @@ impl IndexingServiceImpl { language_chunker, indexing_ops, event_bus, + file_system_provider, + task_runner_provider, file_hash_repository: None, supported_extensions: Self::normalize_supported_extensions(supported_extensions), } @@ -163,6 +172,8 @@ impl IndexingServiceImpl { language_chunker: service.language_chunker, indexing_ops: service.indexing_ops, event_bus: service.event_bus, + file_system_provider: service.file_system_provider, + task_runner_provider: service.task_runner_provider, file_hash_repository: Some(file_hash_repository), supported_extensions: Self::normalize_supported_extensions( service.supported_extensions, @@ -185,31 +196,32 @@ impl IndexingServiceImpl { progress: &mut IndexingProgress, ) -> Vec { let mut files = Vec::new(); - let walker = WalkBuilder::new(path) - .hidden(false) - .filter_entry(|entry| { - if !entry.file_type().is_some_and(|ft| ft.is_dir()) { - return true; - } - - entry - .file_name() - .to_str() - .is_none_or(|name| !SKIP_DIRS.contains(&name)) - }) - .build(); - - for entry_result in walker { - match entry_result { - Ok(entry) => { - if entry.file_type().is_some_and(|ft| ft.is_file()) - && self.is_supported_file(entry.path()) - { - files.push(entry.path().to_path_buf()); + let mut dirs = vec![path.to_path_buf()]; + + while let Some(dir) = dirs.pop() { + match self.file_system_provider.read_dir_entries(&dir).await { + Ok(entries) => { + for entry in entries { + if entry.is_dir { + let should_skip = entry + .path + .file_name() + .and_then(|n| n.to_str()) + .is_some_and(|name| SKIP_DIRS.contains(&name)); + + if !should_skip { + dirs.push(entry.path); + } + continue; + } + + if entry.is_file && self.is_supported_file(&entry.path) { + files.push(entry.path); + } } } Err(e) => { - progress.record_error("Failed to read directory entry", path, e); + progress.record_error("Failed to read directory entries", &dir, e); } } } @@ -286,10 +298,14 @@ impl IndexingServiceInterface for IndexingServiceImpl { // Fire-and-forget: caller gets operation_id immediately, polling for completion. // Sync execution path available via run_indexing_task() directly in tests. - let _handle = tokio::spawn(async move { + let task = Box::pin(async move { Self::run_indexing_task(service, files, workspace_root, collection_id, op_id).await; }); + if let Err(e) = self.task_runner_provider.spawn(task) { + warn!("Failed to spawn indexing background task: {}", e); + } + // Return immediately with operation_id Ok(IndexingResult { files_processed: 0, @@ -424,19 +440,16 @@ impl IndexingServiceImpl { .update_progress(operation_id, Some(relative_path.clone()), index); // Read file content - let content = std::fs::read_to_string(file_path) - .map_err(|e| mcb_domain::error::Error::internal(format!("Failed to read file: {e}")))?; + let content = self.file_system_provider.read_to_string(file_path).await?; // Incremental check using file hashes let current_hash = mcb_domain::utils::compute_content_hash(&content); - #[allow(clippy::collapsible_if)] - if let Some(repo) = &self.file_hash_repository { - if !repo + if let Some(repo) = &self.file_hash_repository + && !repo .has_changed(&collection.to_string(), &relative_path, ¤t_hash) .await? - { - return Ok(ProcessResult::Skipped); - } + { + return Ok(ProcessResult::Skipped); } // Generate semantic chunks diff --git a/crates/mcb-domain/Cargo.toml b/crates/mcb-domain/Cargo.toml index e65de3c00..714af28e1 100644 --- a/crates/mcb-domain/Cargo.toml +++ b/crates/mcb-domain/Cargo.toml @@ -62,8 +62,6 @@ futures = { workspace = true } # Plugin registration linkme = { workspace = true } -# HTTP Client types for providers -reqwest = { workspace = true } toml = { workspace = true, optional = true } [features] diff --git a/crates/mcb-domain/src/entities/api_key.rs b/crates/mcb-domain/src/entities/api_key.rs index b5bbd453b..ce56947a5 100644 --- a/crates/mcb-domain/src/entities/api_key.rs +++ b/crates/mcb-domain/src/entities/api_key.rs @@ -26,3 +26,27 @@ crate::define_entity! { pub revoked_at: Option, } } + +crate::impl_table_schema!(ApiKey, "api_keys", + columns: [ + ("id", Text, pk), + ("user_id", Text), + ("org_id", Text), + ("key_hash", Text), + ("name", Text), + ("scopes_json", Text), + ("expires_at", Integer, nullable), + ("created_at", Integer), + ("revoked_at", Integer, nullable), + ], + indexes: [ + "idx_api_keys_user" => ["user_id"], + "idx_api_keys_org" => ["org_id"], + "idx_api_keys_key_hash" => ["key_hash"], + ], + foreign_keys: [ + ("user_id", "users", "id"), + ("org_id", "organizations", "id"), + ], + unique_constraints: [], +); diff --git a/crates/mcb-domain/src/entities/organization.rs b/crates/mcb-domain/src/entities/organization.rs index a2b5dc471..4b20e5435 100644 --- a/crates/mcb-domain/src/entities/organization.rs +++ b/crates/mcb-domain/src/entities/organization.rs @@ -44,3 +44,17 @@ pub enum OrgStatus { } crate::impl_as_str_from_as_ref!(OrgStatus); + +crate::impl_table_schema!(Organization, "organizations", + columns: [ + ("id", Text, pk), + ("name", Text), + ("slug", Text, unique), + ("settings_json", Text), + ("created_at", Integer), + ("updated_at", Integer), + ], + indexes: [ + "idx_organizations_name" => ["name"], + ], +); diff --git a/crates/mcb-domain/src/entities/team.rs b/crates/mcb-domain/src/entities/team.rs index 867fd8555..a65278628 100644 --- a/crates/mcb-domain/src/entities/team.rs +++ b/crates/mcb-domain/src/entities/team.rs @@ -58,3 +58,41 @@ pub enum TeamMemberRole { } crate::impl_as_str_from_as_ref!(TeamMemberRole); + +crate::impl_table_schema!(Team, "teams", + columns: [ + ("id", Text, pk), + ("org_id", Text), + ("name", Text), + ("created_at", Integer), + ], + indexes: [ + "idx_teams_org" => ["org_id"], + ], + foreign_keys: [ + ("org_id", "organizations", "id"), + ], + unique_constraints: [ + ["org_id", "name"], + ], +); + +crate::impl_table_schema!(TeamMember, "team_members", + columns: [ + ("team_id", Text, pk), + ("user_id", Text, pk), + ("role", Text), + ("joined_at", Integer), + ], + indexes: [ + "idx_team_members_team" => ["team_id"], + "idx_team_members_user" => ["user_id"], + ], + foreign_keys: [ + ("team_id", "teams", "id"), + ("user_id", "users", "id"), + ], + unique_constraints: [ + ["team_id", "user_id"], + ], +); diff --git a/crates/mcb-domain/src/entities/user.rs b/crates/mcb-domain/src/entities/user.rs index 20729b9c2..79731a3ad 100644 --- a/crates/mcb-domain/src/entities/user.rs +++ b/crates/mcb-domain/src/entities/user.rs @@ -50,3 +50,27 @@ pub enum UserRole { } crate::impl_as_str_from_as_ref!(UserRole); + +crate::impl_table_schema!(User, "users", + columns: [ + ("id", Text, pk), + ("org_id", Text), + ("email", Text), + ("display_name", Text), + ("role", Text), + ("api_key_hash", Text, nullable), + ("created_at", Integer), + ("updated_at", Integer), + ], + indexes: [ + "idx_users_org" => ["org_id"], + "idx_users_email" => ["email"], + "idx_users_api_key_hash" => ["api_key_hash"], + ], + foreign_keys: [ + ("org_id", "organizations", "id"), + ], + unique_constraints: [ + ["org_id", "email"], + ], +); diff --git a/crates/mcb-domain/src/macros/di.rs b/crates/mcb-domain/src/macros/di.rs new file mode 100644 index 000000000..ea7d5bf89 --- /dev/null +++ b/crates/mcb-domain/src/macros/di.rs @@ -0,0 +1,41 @@ +/// Generates `Arc`-cloning getter methods for DI container fields. +#[macro_export] +macro_rules! arc_getters { + () => {}; + + ($name:ident : $ty:ty => $field:ident $(= $impl:ty)? , $($rest:tt)*) => { + #[doc = concat!("Get `", stringify!($name), "`.")] + #[must_use] + pub fn $name(&self) -> std::sync::Arc<$ty> { + std::sync::Arc::clone(&self.$field) + } + + $crate::arc_getters! { $($rest)* } + }; + + ($name:ident : $ty:ty => $field:ident $(= $impl:ty)? $(,)?) => { + #[doc = concat!("Get `", stringify!($name), "`.")] + #[must_use] + pub fn $name(&self) -> std::sync::Arc<$ty> { + std::sync::Arc::clone(&self.$field) + } + }; + + ($name:ident : $ty:ty $(= $impl:ty)? , $($rest:tt)*) => { + #[doc = concat!("Get `", stringify!($name), "`.")] + #[must_use] + pub fn $name(&self) -> std::sync::Arc<$ty> { + std::sync::Arc::clone(&self.$name) + } + + $crate::arc_getters! { $($rest)* } + }; + + ($name:ident : $ty:ty $(= $impl:ty)? $(,)?) => { + #[doc = concat!("Get `", stringify!($name), "`.")] + #[must_use] + pub fn $name(&self) -> std::sync::Arc<$ty> { + std::sync::Arc::clone(&self.$name) + } + }; +} diff --git a/crates/mcb-domain/src/macros/mod.rs b/crates/mcb-domain/src/macros/mod.rs index c2eadade9..b3e0b003f 100644 --- a/crates/mcb-domain/src/macros/mod.rs +++ b/crates/mcb-domain/src/macros/mod.rs @@ -19,3 +19,5 @@ mod ports; mod schema; #[macro_use] mod registry; +#[macro_use] +mod di; diff --git a/crates/mcb-domain/src/macros/schema.rs b/crates/mcb-domain/src/macros/schema.rs index 9994c7f38..f11fbb916 100644 --- a/crates/mcb-domain/src/macros/schema.rs +++ b/crates/mcb-domain/src/macros/schema.rs @@ -119,3 +119,99 @@ macro_rules! unique { } }; } + +/// Implement [`HasTableSchema`] for an entity type using compact column specs. +/// +/// Co-locate this invocation in the entity file so the schema lives next to +/// the struct it describes. The macro reuses the existing `col!`, `table!`, +/// `index!`, `fk!` and `unique!` helpers internally. +/// +/// # Example +/// +/// ```ignore +/// impl_table_schema!(Organization, "organizations", +/// columns: [ +/// ("id", Text, pk), +/// ("name", Text), +/// ("slug", Text, unique), +/// ("settings_json", Text), +/// ("created_at", Integer), +/// ("updated_at", Integer), +/// ], +/// indexes: [ +/// "idx_organizations_name" => ["name"], +/// ], +/// ); +/// ``` +#[macro_export] +macro_rules! impl_table_schema { + // Full form: columns + indexes + foreign_keys + unique_constraints + ($entity:ty, $table_name:expr, + columns: [ $( ($col_name:expr, $col_type:ident $(, $flag:ident)?) ),* $(,)? ], + indexes: [ $( $idx_name:expr => [ $($idx_col:expr),* $(,)? ] ),* $(,)? ], + foreign_keys: [ $( ($fk_col:expr, $fk_table:expr, $fk_ref:expr) ),* $(,)? ], + unique_constraints: [ $( [ $($uc_col:expr),* $(,)? ] ),* $(,)? ], + ) => { + impl $crate::schema::types::HasTableSchema for $entity { + fn table_def() -> $crate::schema::types::TableDef { + $crate::table!($table_name, [ + $( $crate::col!($col_name, $col_type $(, $flag)?) ),* + ]) + } + + fn indexes() -> Vec<$crate::schema::types::IndexDef> { + vec![ + $( $crate::index!($idx_name, $table_name, [ $($idx_col),* ]) ),* + ] + } + + fn foreign_keys() -> Vec<$crate::schema::types::ForeignKeyDef> { + vec![ + $( $crate::fk!($table_name, $fk_col, $fk_table, $fk_ref) ),* + ] + } + + fn unique_constraints() -> Vec<$crate::schema::types::UniqueConstraintDef> { + vec![ + $( $crate::unique!($table_name, [ $($uc_col),* ]) ),* + ] + } + } + }; + // Shorthand: columns only (no indexes, no FKs, no UCs) + ($entity:ty, $table_name:expr, + columns: [ $( ($col_name:expr, $col_type:ident $(, $flag:ident)?) ),* $(,)? ], + ) => { + $crate::impl_table_schema!($entity, $table_name, + columns: [ $( ($col_name, $col_type $(, $flag)?) ),* ], + indexes: [], + foreign_keys: [], + unique_constraints: [], + ); + }; + // Shorthand: columns + indexes (no FKs, no UCs) + ($entity:ty, $table_name:expr, + columns: [ $( ($col_name:expr, $col_type:ident $(, $flag:ident)?) ),* $(,)? ], + indexes: [ $( $idx_name:expr => [ $($idx_col:expr),* $(,)? ] ),* $(,)? ], + ) => { + $crate::impl_table_schema!($entity, $table_name, + columns: [ $( ($col_name, $col_type $(, $flag)?) ),* ], + indexes: [ $( $idx_name => [ $($idx_col),* ] ),* ], + foreign_keys: [], + unique_constraints: [], + ); + }; + // Shorthand: columns + foreign_keys (no indexes, no UCs) + ($entity:ty, $table_name:expr, + columns: [ $( ($col_name:expr, $col_type:ident $(, $flag:ident)?) ),* $(,)? ], + foreign_keys: [ $( ($fk_col:expr, $fk_table:expr, $fk_ref:expr) ),* $(,)? ], + ) => { + $crate::impl_table_schema!($entity, $table_name, + columns: [ $( ($col_name, $col_type $(, $flag)?) ),* ], + indexes: [], + foreign_keys: [ $( ($fk_col, $fk_table, $fk_ref) ),* ], + unique_constraints: [], + ); + }; + // Shorthand: columns + indexes + foreign_keys (no UCs) +} diff --git a/crates/mcb-domain/src/ports/infrastructure/database.rs b/crates/mcb-domain/src/ports/infrastructure/database.rs index 00fa3e5f0..c3cc12edb 100644 --- a/crates/mcb-domain/src/ports/infrastructure/database.rs +++ b/crates/mcb-domain/src/ports/infrastructure/database.rs @@ -12,6 +12,10 @@ use std::sync::Arc; use async_trait::async_trait; use crate::error::Result; +use crate::ports::{ + AgentRepository, FileHashRepository, IssueEntityRepository, MemoryRepository, + OrgEntityRepository, PlanEntityRepository, ProjectRepository, VcsEntityRepository, +}; /// Parameter for prepared statement binding (driver-agnostic). #[derive(Debug, Clone)] @@ -75,7 +79,49 @@ pub trait DatabaseExecutor: Send + Sync { /// Provider factory for database connections with schema initialization. #[async_trait] +#[allow(missing_docs)] pub trait DatabaseProvider: Send + Sync { /// Performs the connect operation. async fn connect(&self, path: &std::path::Path) -> Result>; + + fn create_memory_repository( + &self, + executor: Arc, + ) -> Arc; + + fn create_agent_repository( + &self, + executor: Arc, + ) -> Arc; + + fn create_project_repository( + &self, + executor: Arc, + ) -> Arc; + + fn create_file_hash_repository( + &self, + executor: Arc, + project_id: String, + ) -> Arc; + + fn create_vcs_entity_repository( + &self, + executor: Arc, + ) -> Arc; + + fn create_plan_entity_repository( + &self, + executor: Arc, + ) -> Arc; + + fn create_issue_entity_repository( + &self, + executor: Arc, + ) -> Arc; + + fn create_org_entity_repository( + &self, + executor: Arc, + ) -> Arc; } diff --git a/crates/mcb-domain/src/ports/mod.rs b/crates/mcb-domain/src/ports/mod.rs index 6cc7b61e5..e751716bd 100644 --- a/crates/mcb-domain/src/ports/mod.rs +++ b/crates/mcb-domain/src/ports/mod.rs @@ -59,11 +59,11 @@ pub use providers::vector_store::{VectorStoreAdmin, VectorStoreBrowser}; pub use providers::{ CacheEntryConfig, CacheProvider, CacheStats, ComplexityAnalyzer, ComplexityFinding, CryptoProvider, DEFAULT_CACHE_NAMESPACE, DEFAULT_CACHE_TTL_SECS, DeadCodeDetector, - DeadCodeFinding, EmbeddingProvider, EncryptedData, FileMetrics, FunctionMetrics, - HalsteadMetrics, HttpClientConfig, HttpClientProvider, HybridSearchProvider, - HybridSearchResult, LanguageChunkingProvider, MetricLabels, MetricsAnalysisProvider, - MetricsError, MetricsProvider, MetricsResult, ProjectDetector, ProjectDetectorConfig, - ProjectDetectorEntry, ProviderConfigManagerInterface, TdgFinding, TdgScorer, ValidationOptions, + DeadCodeFinding, DirEntry, EmbeddingProvider, EncryptedData, FileMetrics, FileSystemProvider, + FunctionMetrics, HalsteadMetrics, HybridSearchProvider, HybridSearchResult, + LanguageChunkingProvider, MetricLabels, MetricsAnalysisProvider, MetricsError, MetricsProvider, + MetricsResult, ProjectDetector, ProjectDetectorConfig, ProjectDetectorEntry, + ProviderConfigManagerInterface, TaskRunnerProvider, TdgFinding, TdgScorer, ValidationOptions, ValidationProvider, ValidatorInfo, VcsProvider, VectorStoreProvider, }; diff --git a/crates/mcb-domain/src/ports/providers/fs.rs b/crates/mcb-domain/src/ports/providers/fs.rs new file mode 100644 index 000000000..295eaee53 --- /dev/null +++ b/crates/mcb-domain/src/ports/providers/fs.rs @@ -0,0 +1,23 @@ +use std::path::{Path, PathBuf}; + +use async_trait::async_trait; + +use crate::error::Error; + +#[allow(missing_docs)] +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DirEntry { + pub path: PathBuf, + pub is_file: bool, + pub is_dir: bool, +} + +#[allow(missing_docs)] +#[async_trait] +pub trait FileSystemProvider: Send + Sync { + async fn read_to_string(&self, path: &Path) -> std::result::Result; + + async fn read_dir_entries(&self, path: &Path) -> std::result::Result, Error>; + + async fn canonicalize_path(&self, path: &Path) -> std::result::Result; +} diff --git a/crates/mcb-domain/src/ports/providers/http.rs b/crates/mcb-domain/src/ports/providers/http.rs deleted file mode 100644 index de1868310..000000000 --- a/crates/mcb-domain/src/ports/providers/http.rs +++ /dev/null @@ -1,62 +0,0 @@ -//! -//! **Documentation**: [docs/modules/domain.md](../../../../../docs/modules/domain.md#provider-ports) -//! -use std::time::Duration; - -use reqwest::Client; -use serde::{Deserialize, Serialize}; - -/// HTTP client configuration -/// -/// Controls connection pooling, timeouts, and other HTTP client behavior. -/// Used by `HttpClientProvider` to configure HTTP requests. -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct HttpClientConfig { - /// Maximum idle connections per host - pub max_idle_per_host: usize, - /// Idle connection timeout - pub idle_timeout: Duration, - /// TCP keep-alive duration - pub keepalive: Duration, - /// Total timeout for requests - pub timeout: Duration, - /// User agent string - pub user_agent: String, -} - -impl Default for HttpClientConfig { - fn default() -> Self { - Self { - max_idle_per_host: 10, - idle_timeout: Duration::from_secs(90), - keepalive: Duration::from_secs(60), - timeout: Duration::from_secs(30), - user_agent: "mcb/domain-client".to_owned(), - } - } -} - -/// HTTP client provider trait -/// -/// Defines the interface for HTTP client operations used by API-based providers. -pub trait HttpClientProvider: Send + Sync { - /// Get a reference to the underlying reqwest Client - fn client(&self) -> &Client; - - /// Get the configuration - fn config(&self) -> &HttpClientConfig; - - /// Create a new client with custom timeout for specific operations. - /// - /// # Errors - /// - /// Returns an error if the HTTP client cannot be constructed with the - /// given timeout configuration. - fn client_with_timeout( - &self, - timeout: Duration, - ) -> Result>; - - /// Check if the client pool is enabled - fn is_enabled(&self) -> bool; -} diff --git a/crates/mcb-domain/src/ports/providers/mod.rs b/crates/mcb-domain/src/ports/providers/mod.rs index dd86088a5..c475128b1 100644 --- a/crates/mcb-domain/src/ports/providers/mod.rs +++ b/crates/mcb-domain/src/ports/providers/mod.rs @@ -30,8 +30,7 @@ pub mod config; pub mod crypto; /// Embedding provider port pub mod embedding; -/// HTTP client provider port -pub mod http; +pub mod fs; /// Hybrid search provider port pub mod hybrid_search; /// Language chunking provider port @@ -42,6 +41,8 @@ pub mod metrics; pub mod metrics_analysis; /// Project detection provider port pub mod project_detection; +/// Background task runner provider port +pub mod task; /// Validation provider port pub mod validation; /// Version control system provider port @@ -59,7 +60,7 @@ pub use cache::{ pub use config::ProviderConfigManagerInterface; pub use crypto::{CryptoProvider, EncryptedData}; pub use embedding::EmbeddingProvider; -pub use http::{HttpClientConfig, HttpClientProvider}; +pub use fs::{DirEntry, FileSystemProvider}; pub use hybrid_search::{HybridSearchProvider, HybridSearchResult}; pub use language_chunking::LanguageChunkingProvider; pub use metrics::{MetricLabels, MetricsError, MetricsProvider, MetricsResult}; @@ -67,6 +68,7 @@ pub use metrics_analysis::{ FileMetrics, FunctionMetrics, HalsteadMetrics, MetricsAnalysisProvider, }; pub use project_detection::{ProjectDetector, ProjectDetectorConfig, ProjectDetectorEntry}; +pub use task::TaskRunnerProvider; pub use validation::{ValidationOptions, ValidationProvider, ValidatorInfo}; pub use vcs::VcsProvider; pub use vector_store::VectorStoreProvider; diff --git a/crates/mcb-domain/src/ports/providers/task.rs b/crates/mcb-domain/src/ports/providers/task.rs new file mode 100644 index 000000000..8c7780db7 --- /dev/null +++ b/crates/mcb-domain/src/ports/providers/task.rs @@ -0,0 +1,12 @@ +use futures::future::BoxFuture; + +use crate::error::Result; + +/// Port for spawning background tasks (used by DI; implementations in mcb-providers). +pub trait TaskRunnerProvider: Send + Sync { + /// Spawns a background task. + /// + /// # Errors + /// Returns an error if the task could not be spawned. + fn spawn(&self, task: BoxFuture<'static, ()>) -> Result<()>; +} diff --git a/crates/mcb-domain/src/registry/event_bus.rs b/crates/mcb-domain/src/registry/event_bus.rs new file mode 100644 index 000000000..36e05cbae --- /dev/null +++ b/crates/mcb-domain/src/registry/event_bus.rs @@ -0,0 +1,20 @@ +#![allow(missing_docs)] + +use std::collections::HashMap; + +#[derive(Debug, Clone, Default)] +pub struct EventBusProviderConfig { + pub provider: String, + pub extra: HashMap, +} + +crate::impl_config_builder!(EventBusProviderConfig {}); + +crate::impl_registry!( + provider_trait: crate::ports::EventBusProvider, + config_type: EventBusProviderConfig, + entry_type: EventBusProviderEntry, + slice_name: EVENT_BUS_PROVIDERS, + resolve_fn: resolve_event_bus_provider, + list_fn: list_event_bus_providers +); diff --git a/crates/mcb-domain/src/registry/fs.rs b/crates/mcb-domain/src/registry/fs.rs new file mode 100644 index 000000000..d3da532b3 --- /dev/null +++ b/crates/mcb-domain/src/registry/fs.rs @@ -0,0 +1,20 @@ +#![allow(missing_docs)] + +use std::collections::HashMap; + +#[derive(Debug, Clone, Default)] +pub struct FileSystemProviderConfig { + pub provider: String, + pub extra: HashMap, +} + +crate::impl_config_builder!(FileSystemProviderConfig {}); + +crate::impl_registry!( + provider_trait: crate::ports::FileSystemProvider, + config_type: FileSystemProviderConfig, + entry_type: FileSystemProviderEntry, + slice_name: FILE_SYSTEM_PROVIDERS, + resolve_fn: resolve_file_system_provider, + list_fn: list_file_system_providers +); diff --git a/crates/mcb-domain/src/registry/mod.rs b/crates/mcb-domain/src/registry/mod.rs index 43f602c27..c2d1260f2 100644 --- a/crates/mcb-domain/src/registry/mod.rs +++ b/crates/mcb-domain/src/registry/mod.rs @@ -56,6 +56,10 @@ pub mod cache; pub mod database; pub mod embedding; +pub mod event_bus; +pub mod fs; pub mod language; +pub mod task_runner; pub mod validation; +pub mod vcs; pub mod vector_store; diff --git a/crates/mcb-domain/src/registry/task_runner.rs b/crates/mcb-domain/src/registry/task_runner.rs new file mode 100644 index 000000000..bbef29adc --- /dev/null +++ b/crates/mcb-domain/src/registry/task_runner.rs @@ -0,0 +1,20 @@ +#![allow(missing_docs)] + +use std::collections::HashMap; + +#[derive(Debug, Clone, Default)] +pub struct TaskRunnerProviderConfig { + pub provider: String, + pub extra: HashMap, +} + +crate::impl_config_builder!(TaskRunnerProviderConfig {}); + +crate::impl_registry!( + provider_trait: crate::ports::TaskRunnerProvider, + config_type: TaskRunnerProviderConfig, + entry_type: TaskRunnerProviderEntry, + slice_name: TASK_RUNNER_PROVIDERS, + resolve_fn: resolve_task_runner_provider, + list_fn: list_task_runner_providers +); diff --git a/crates/mcb-domain/src/registry/vcs.rs b/crates/mcb-domain/src/registry/vcs.rs new file mode 100644 index 000000000..14406a513 --- /dev/null +++ b/crates/mcb-domain/src/registry/vcs.rs @@ -0,0 +1,20 @@ +#![allow(missing_docs)] + +use std::collections::HashMap; + +#[derive(Debug, Clone, Default)] +pub struct VcsProviderConfig { + pub provider: String, + pub extra: HashMap, +} + +crate::impl_config_builder!(VcsProviderConfig {}); + +crate::impl_registry!( + provider_trait: crate::ports::VcsProvider, + config_type: VcsProviderConfig, + entry_type: VcsProviderEntry, + slice_name: VCS_PROVIDERS, + resolve_fn: resolve_vcs_provider, + list_fn: list_vcs_providers +); diff --git a/crates/mcb-domain/src/schema/api_keys.rs b/crates/mcb-domain/src/schema/api_keys.rs deleted file mode 100644 index b2e43cda1..000000000 --- a/crates/mcb-domain/src/schema/api_keys.rs +++ /dev/null @@ -1,40 +0,0 @@ -//! -//! **Documentation**: [docs/modules/domain.md](../../../../docs/modules/domain.md) -//! -use crate::schema::types::{ForeignKeyDef, IndexDef, TableDef, UniqueConstraintDef}; - -pub fn table() -> TableDef { - crate::table!( - "api_keys", - [ - crate::col!("id", Text, pk), - crate::col!("user_id", Text), - crate::col!("org_id", Text), - crate::col!("key_hash", Text), - crate::col!("name", Text), - crate::col!("scopes_json", Text), - crate::col!("expires_at", Integer, nullable), - crate::col!("created_at", Integer), - crate::col!("revoked_at", Integer, nullable), - ] - ) -} - -pub fn indexes() -> Vec { - vec![ - crate::index!("idx_api_keys_user", "api_keys", ["user_id"]), - crate::index!("idx_api_keys_org", "api_keys", ["org_id"]), - crate::index!("idx_api_keys_key_hash", "api_keys", ["key_hash"]), - ] -} - -pub fn foreign_keys() -> Vec { - vec![ - crate::fk!("api_keys", "user_id", "users", "id"), - crate::fk!("api_keys", "org_id", "organizations", "id"), - ] -} - -pub fn unique_constraints() -> Vec { - Vec::new() -} diff --git a/crates/mcb-domain/src/schema/definition.rs b/crates/mcb-domain/src/schema/definition.rs index 3ab276424..b678e5304 100644 --- a/crates/mcb-domain/src/schema/definition.rs +++ b/crates/mcb-domain/src/schema/definition.rs @@ -1,65 +1,81 @@ //! //! **Documentation**: [docs/modules/domain.md](../../../../docs/modules/domain.md) //! -use super::types::{ForeignKeyDef, FtsDef, IndexDef, Schema, TableDef, UniqueConstraintDef}; +use super::types::{ + ForeignKeyDef, FtsDef, HasTableSchema, IndexDef, Schema, TableDef, UniqueConstraintDef, +}; + +// Schema modules not yet migrated to HasTableSchema use super::{ - agent_sessions, agent_worktree_assignments, api_keys, branches, checkpoints, collections, - delegations, error_pattern_matches, error_patterns, file_hashes, issue_comments, - issue_label_assignments, issue_labels, observations, organizations, plan_reviews, - plan_versions, plans, project_issues, projects, repositories, session_summaries, team_members, - teams, tool_calls, users, worktrees, + agent_sessions, agent_worktree_assignments, branches, checkpoints, collections, delegations, + error_pattern_matches, error_patterns, file_hashes, issue_comments, issue_label_assignments, + issue_labels, observations, plan_reviews, plan_versions, plans, project_issues, projects, + repositories, session_summaries, tool_calls, worktrees, }; -struct SchemaModule { +struct SchemaEntry { table: fn() -> TableDef, indexes: fn() -> Vec, foreign_keys: fn() -> Vec, unique_constraints: fn() -> Vec, } -macro_rules! schema_modules { - ($($module:ident),+ $(,)?) => { - &[ - $( - SchemaModule { - table: $module::table, - indexes: $module::indexes, - foreign_keys: $module::foreign_keys, - unique_constraints: $module::unique_constraints, - }, - )+ - ] +/// Build a [`SchemaEntry`] from a type implementing [`HasTableSchema`]. +macro_rules! from_entity { + ($entity:ty) => { + SchemaEntry { + table: <$entity as HasTableSchema>::table_def, + indexes: <$entity as HasTableSchema>::indexes, + foreign_keys: <$entity as HasTableSchema>::foreign_keys, + unique_constraints: <$entity as HasTableSchema>::unique_constraints, + } + }; +} + +/// Build a [`SchemaEntry`] from a legacy schema module (4 free functions). +macro_rules! from_module { + ($module:ident) => { + SchemaEntry { + table: $module::table, + indexes: $module::indexes, + foreign_keys: $module::foreign_keys, + unique_constraints: $module::unique_constraints, + } }; } -const SCHEMA_MODULES: &[SchemaModule] = schema_modules![ - organizations, - users, - teams, - team_members, - api_keys, - projects, - collections, - observations, - session_summaries, - file_hashes, - agent_sessions, - delegations, - tool_calls, - checkpoints, - error_patterns, - error_pattern_matches, - project_issues, - issue_comments, - issue_labels, - issue_label_assignments, - plans, - plan_versions, - plan_reviews, - repositories, - branches, - worktrees, - agent_worktree_assignments, +use crate::entities::{ApiKey, Organization, Team, TeamMember, User}; + +const SCHEMA_ENTRIES: &[SchemaEntry] = &[ + // ── Migrated to HasTableSchema (entity is the source of truth) ── + from_entity!(Organization), + from_entity!(User), + from_entity!(Team), + from_entity!(TeamMember), + from_entity!(ApiKey), + // ── Legacy schema modules (pending migration) ── + from_module!(projects), + from_module!(collections), + from_module!(observations), + from_module!(session_summaries), + from_module!(file_hashes), + from_module!(agent_sessions), + from_module!(delegations), + from_module!(tool_calls), + from_module!(checkpoints), + from_module!(error_patterns), + from_module!(error_pattern_matches), + from_module!(project_issues), + from_module!(issue_comments), + from_module!(issue_labels), + from_module!(issue_label_assignments), + from_module!(plans), + from_module!(plan_versions), + from_module!(plan_reviews), + from_module!(repositories), + from_module!(branches), + from_module!(worktrees), + from_module!(agent_worktree_assignments), ]; impl Schema { @@ -76,10 +92,7 @@ impl Schema { } fn tables() -> Vec { - SCHEMA_MODULES - .iter() - .map(|module| (module.table)()) - .collect() + SCHEMA_ENTRIES.iter().map(|entry| (entry.table)()).collect() } fn fts_def() -> Option { @@ -92,23 +105,23 @@ impl Schema { } fn indexes() -> Vec { - SCHEMA_MODULES + SCHEMA_ENTRIES .iter() - .flat_map(|module| (module.indexes)().into_iter()) + .flat_map(|entry| (entry.indexes)().into_iter()) .collect() } fn foreign_keys() -> Vec { - SCHEMA_MODULES + SCHEMA_ENTRIES .iter() - .flat_map(|module| (module.foreign_keys)().into_iter()) + .flat_map(|entry| (entry.foreign_keys)().into_iter()) .collect() } fn unique_constraints() -> Vec { - SCHEMA_MODULES + SCHEMA_ENTRIES .iter() - .flat_map(|module| (module.unique_constraints)().into_iter()) + .flat_map(|entry| (entry.unique_constraints)().into_iter()) .collect() } } diff --git a/crates/mcb-domain/src/schema/mod.rs b/crates/mcb-domain/src/schema/mod.rs index 905679f9b..7119f38ef 100644 --- a/crates/mcb-domain/src/schema/mod.rs +++ b/crates/mcb-domain/src/schema/mod.rs @@ -1,9 +1,11 @@ //! //! **Documentation**: [docs/modules/domain.md](../../../../docs/modules/domain.md) //! + +// Legacy schema modules — pending migration to impl_table_schema! on entities. +// Migrated: organizations, users, teams, team_members, api_keys (deleted). mod agent_sessions; mod agent_worktree_assignments; -mod api_keys; mod branches; mod checkpoints; mod collections; @@ -16,7 +18,6 @@ mod issue_comments; mod issue_label_assignments; mod issue_labels; mod observations; -mod organizations; mod plan_reviews; mod plan_versions; mod plans; @@ -24,12 +25,9 @@ mod project_issues; mod projects; mod repositories; mod session_summaries; -mod team_members; -mod teams; mod tool_calls; /// Canonical schema model types and DDL generation traits. pub mod types; -mod users; mod worktrees; pub use types::*; diff --git a/crates/mcb-domain/src/schema/organizations.rs b/crates/mcb-domain/src/schema/organizations.rs deleted file mode 100644 index 0c3ad70ab..000000000 --- a/crates/mcb-domain/src/schema/organizations.rs +++ /dev/null @@ -1,34 +0,0 @@ -//! -//! **Documentation**: [docs/modules/domain.md](../../../../docs/modules/domain.md) -//! -use crate::schema::types::{ForeignKeyDef, IndexDef, TableDef, UniqueConstraintDef}; - -pub fn table() -> TableDef { - crate::table!( - "organizations", - [ - crate::col!("id", Text, pk), - crate::col!("name", Text), - crate::col!("slug", Text, unique), - crate::col!("settings_json", Text), - crate::col!("created_at", Integer), - crate::col!("updated_at", Integer), - ] - ) -} - -pub fn indexes() -> Vec { - vec![crate::index!( - "idx_organizations_name", - "organizations", - ["name"] - )] -} - -pub fn foreign_keys() -> Vec { - Vec::new() -} - -pub fn unique_constraints() -> Vec { - Vec::new() -} diff --git a/crates/mcb-domain/src/schema/team_members.rs b/crates/mcb-domain/src/schema/team_members.rs deleted file mode 100644 index 1ec2230d7..000000000 --- a/crates/mcb-domain/src/schema/team_members.rs +++ /dev/null @@ -1,34 +0,0 @@ -//! -//! **Documentation**: [docs/modules/domain.md](../../../../docs/modules/domain.md) -//! -use crate::schema::types::{ForeignKeyDef, IndexDef, TableDef, UniqueConstraintDef}; - -pub fn table() -> TableDef { - crate::table!( - "team_members", - [ - crate::col!("team_id", Text, pk), - crate::col!("user_id", Text, pk), - crate::col!("role", Text), - crate::col!("joined_at", Integer), - ] - ) -} - -pub fn indexes() -> Vec { - vec![ - crate::index!("idx_team_members_team", "team_members", ["team_id"]), - crate::index!("idx_team_members_user", "team_members", ["user_id"]), - ] -} - -pub fn foreign_keys() -> Vec { - vec![ - crate::fk!("team_members", "team_id", "teams", "id"), - crate::fk!("team_members", "user_id", "users", "id"), - ] -} - -pub fn unique_constraints() -> Vec { - vec![crate::unique!("team_members", ["team_id", "user_id"])] -} diff --git a/crates/mcb-domain/src/schema/teams.rs b/crates/mcb-domain/src/schema/teams.rs deleted file mode 100644 index be4757b8b..000000000 --- a/crates/mcb-domain/src/schema/teams.rs +++ /dev/null @@ -1,28 +0,0 @@ -//! -//! **Documentation**: [docs/modules/domain.md](../../../../docs/modules/domain.md) -//! -use crate::schema::types::{ForeignKeyDef, IndexDef, TableDef, UniqueConstraintDef}; - -pub fn table() -> TableDef { - crate::table!( - "teams", - [ - crate::col!("id", Text, pk), - crate::col!("org_id", Text), - crate::col!("name", Text), - crate::col!("created_at", Integer), - ] - ) -} - -pub fn indexes() -> Vec { - vec![crate::index!("idx_teams_org", "teams", ["org_id"])] -} - -pub fn foreign_keys() -> Vec { - vec![crate::fk!("teams", "org_id", "organizations", "id")] -} - -pub fn unique_constraints() -> Vec { - vec![crate::unique!("teams", ["org_id", "name"])] -} diff --git a/crates/mcb-domain/src/schema/types.rs b/crates/mcb-domain/src/schema/types.rs index 0b7eb1a57..6f4d24f3b 100644 --- a/crates/mcb-domain/src/schema/types.rs +++ b/crates/mcb-domain/src/schema/types.rs @@ -116,6 +116,28 @@ pub struct Schema { pub unique_constraints: Vec, } +/// Trait for entities that carry their own table schema definition. +/// +/// Implementors provide the canonical DDL metadata for their corresponding +/// database table. The [`Schema::definition`] aggregator collects all +/// implementations to build the full canonical schema. +pub trait HasTableSchema { + /// Returns the canonical table definition (columns, types, PK, nullable). + fn table_def() -> TableDef; + /// Returns secondary indexes for this table. + fn indexes() -> Vec { + Vec::new() + } + /// Returns foreign key relationships for this table. + fn foreign_keys() -> Vec { + Vec::new() + } + /// Returns multi-column uniqueness constraints for this table. + fn unique_constraints() -> Vec { + Vec::new() + } +} + /// Port for generating DDL from the canonical schema. pub trait SchemaDdlGenerator: Send + Sync { /// Generates backend-specific DDL statements from the canonical schema. diff --git a/crates/mcb-domain/src/schema/users.rs b/crates/mcb-domain/src/schema/users.rs deleted file mode 100644 index 61291064b..000000000 --- a/crates/mcb-domain/src/schema/users.rs +++ /dev/null @@ -1,36 +0,0 @@ -//! -//! **Documentation**: [docs/modules/domain.md](../../../../docs/modules/domain.md) -//! -use crate::schema::types::{ForeignKeyDef, IndexDef, TableDef, UniqueConstraintDef}; - -pub fn table() -> TableDef { - crate::table!( - "users", - [ - crate::col!("id", Text, pk), - crate::col!("org_id", Text), - crate::col!("email", Text), - crate::col!("display_name", Text), - crate::col!("role", Text), - crate::col!("api_key_hash", Text, nullable), - crate::col!("created_at", Integer), - crate::col!("updated_at", Integer), - ] - ) -} - -pub fn indexes() -> Vec { - vec![ - crate::index!("idx_users_org", "users", ["org_id"]), - crate::index!("idx_users_email", "users", ["email"]), - crate::index!("idx_users_api_key_hash", "users", ["api_key_hash"]), - ] -} - -pub fn foreign_keys() -> Vec { - vec![crate::fk!("users", "org_id", "organizations", "id")] -} - -pub fn unique_constraints() -> Vec { - vec![crate::unique!("users", ["org_id", "email"])] -} diff --git a/crates/mcb-domain/tests/unit/mod.rs b/crates/mcb-domain/tests/unit/mod.rs index c1a0567aa..2d9a5096c 100644 --- a/crates/mcb-domain/tests/unit/mod.rs +++ b/crates/mcb-domain/tests/unit/mod.rs @@ -8,5 +8,6 @@ pub mod error; pub mod events; pub mod ports; pub mod repositories; +pub mod schema_entity_sync_tests; pub mod utils; pub mod value_objects; diff --git a/crates/mcb-domain/tests/unit/schema_entity_sync_tests.rs b/crates/mcb-domain/tests/unit/schema_entity_sync_tests.rs new file mode 100644 index 000000000..35d6d3944 --- /dev/null +++ b/crates/mcb-domain/tests/unit/schema_entity_sync_tests.rs @@ -0,0 +1,233 @@ +//! Schema-Entity synchronization tests. +//! +//! Detects divergences between entity struct fields and their corresponding schema definitions. +//! Uses the canonical Schema::definition() to access table metadata. + +#[cfg(test)] +mod tests { + use mcb_domain::schema::Schema; + use std::collections::HashMap; + + #[test] + fn print_full_schema_report() { + let schema = Schema::definition(); + println!("\n=== FULL SCHEMA ANALYSIS ==="); + println!("Total tables: {}", schema.tables.len()); + println!("Total indexes: {}", schema.indexes.len()); + println!("Total foreign keys: {}", schema.foreign_keys.len()); + println!("\nTable summary:"); + + for table in &schema.tables { + println!( + " {} ({} cols): {}", + table.name, + table.columns.len(), + table + .columns + .iter() + .map(|c| &c.name) + .cloned() + .collect::>() + .join(", ") + ); + } + } + + #[test] + fn detect_schema_only_fields() { + let schema = Schema::definition(); + + // Build a map of table_name -> column_names + let mut tables_by_name: HashMap> = HashMap::new(); + for table in &schema.tables { + let cols: Vec = table.columns.iter().map(|c| c.name.clone()).collect(); + tables_by_name.insert(table.name.clone(), cols); + } + + println!("\n=== SCHEMA-ONLY FIELDS (Not in entities) ==="); + + // Check repositories + if let Some(repo_cols) = tables_by_name.get("repositories") { + if repo_cols.contains(&"origin_context".to_string()) { + println!(" ⚠ Repository.origin_context - schema-only, needs entity field"); + } + } + + // Check worktrees + if let Some(wt_cols) = tables_by_name.get("worktrees") { + let mut extra = Vec::new(); + if wt_cols.contains(&"org_id".to_string()) { + extra.push("org_id"); + } + if wt_cols.contains(&"project_id".to_string()) { + extra.push("project_id"); + } + if wt_cols.contains(&"origin_context".to_string()) { + extra.push("origin_context"); + } + if !extra.is_empty() { + println!(" ⚠ Worktree has schema-only fields: {}", extra.join(", ")); + } + } + + // Check branches + if let Some(br_cols) = tables_by_name.get("branches") { + if br_cols.contains(&"origin_context".to_string()) { + println!(" ⚠ Branch.origin_context - schema-only, needs entity field"); + } + } + } + + #[test] + fn test_organizations_table() { + let schema = Schema::definition(); + let org_table = schema + .tables + .iter() + .find(|t| t.name == "organizations") + .expect("organizations table not found"); + + assert_eq!(org_table.columns.len(), 6); + let col_names: Vec<&str> = org_table.columns.iter().map(|c| c.name.as_str()).collect(); + assert_eq!( + col_names, + vec![ + "id", + "name", + "slug", + "settings_json", + "created_at", + "updated_at" + ] + ); + println!("✓ organizations table has correct schema"); + } + + #[test] + fn test_users_table() { + let schema = Schema::definition(); + let users_table = schema + .tables + .iter() + .find(|t| t.name == "users") + .expect("users table not found"); + + let col_names: Vec<&str> = users_table + .columns + .iter() + .map(|c| c.name.as_str()) + .collect(); + assert!(col_names.contains(&"id")); + assert!(col_names.contains(&"org_id")); + assert!(col_names.contains(&"email")); + assert!(col_names.contains(&"display_name")); + assert!(col_names.contains(&"role")); + println!("✓ users table has correct columns"); + } + + #[test] + fn test_plans_table() { + let schema = Schema::definition(); + let plans_table = schema + .tables + .iter() + .find(|t| t.name == "plans") + .expect("plans table not found"); + + let col_names: Vec<&str> = plans_table + .columns + .iter() + .map(|c| c.name.as_str()) + .collect(); + assert!(col_names.contains(&"id")); + assert!(col_names.contains(&"org_id")); + assert!(col_names.contains(&"project_id")); + assert!(col_names.contains(&"title")); + assert!(col_names.contains(&"status")); + assert!(col_names.contains(&"created_by")); + println!("✓ plans table has correct columns"); + } + + #[test] + fn test_project_issues_table() { + let schema = Schema::definition(); + let issues_table = schema + .tables + .iter() + .find(|t| t.name == "project_issues") + .expect("project_issues table not found"); + + let col_names: Vec<&str> = issues_table + .columns + .iter() + .map(|c| c.name.as_str()) + .collect(); + assert!(col_names.contains(&"id")); + assert!( + col_names.contains(&"labels"), + "labels column should be Text (JSON)" + ); + assert!(col_names.contains(&"issue_type")); + assert!(col_names.contains(&"status")); + println!("✓ project_issues table has correct columns"); + } + + #[test] + fn test_org_has_table_schema_matches_definition() { + use mcb_domain::entities::Organization; + use mcb_domain::schema::HasTableSchema; + + let trait_table = Organization::table_def(); + let schema = Schema::definition(); + let def_table = schema + .tables + .iter() + .find(|t| t.name == "organizations") + .expect("organizations table not found in definition"); + + assert_eq!(trait_table.name, def_table.name); + assert_eq!(trait_table.columns.len(), def_table.columns.len()); + + for (t_col, d_col) in trait_table.columns.iter().zip(def_table.columns.iter()) { + assert_eq!(t_col.name, d_col.name, "column name mismatch"); + assert_eq!(t_col.type_, d_col.type_, "type mismatch for {}", t_col.name); + assert_eq!( + t_col.primary_key, d_col.primary_key, + "pk mismatch for {}", + t_col.name + ); + assert_eq!( + t_col.not_null, d_col.not_null, + "not_null mismatch for {}", + t_col.name + ); + assert_eq!( + t_col.unique, d_col.unique, + "unique mismatch for {}", + t_col.name + ); + } + + let trait_indexes = Organization::indexes(); + let def_indexes: Vec<_> = schema + .indexes + .iter() + .filter(|i| i.table == "organizations") + .collect(); + assert_eq!(trait_indexes.len(), def_indexes.len()); + println!("✓ Organization HasTableSchema matches Schema::definition() exactly"); + } + + #[test] + fn report_divergence_mitigation() { + println!("\n=== DIVERGENCE MITIGATION STRATEGY ==="); + println!("\nFor each schema-only field:"); + println!(" 1. Add to entity struct: pub field_name: Type"); + println!(" 2. Add to derive(TableSchema) anotation"); + println!(" 3. OR use extra_columns for temporary escape hatch"); + println!("\nFor each entity field not in schema:"); + println!(" 1. Review if field should be persisted"); + println!(" 2. If yes, add to schema definition"); + println!(" 3. If no, mark with #[schema(skip)]"); + } +} diff --git a/crates/mcb-infrastructure/src/config/types/system.rs b/crates/mcb-infrastructure/src/config/types/system.rs index 7094331c9..2b48c356d 100644 --- a/crates/mcb-infrastructure/src/config/types/system.rs +++ b/crates/mcb-infrastructure/src/config/types/system.rs @@ -10,10 +10,9 @@ use std::path::PathBuf; use serde::{Deserialize, Serialize}; -use mcb_providers::constants::EVENTS_TOKIO_DEFAULT_CAPACITY; - use crate::constants::events::{ - DEFAULT_NATS_CLIENT_NAME, EVENT_BUS_CONNECTION_TIMEOUT_MS, EVENT_BUS_MAX_RECONNECT_ATTEMPTS, + DEFAULT_NATS_CLIENT_NAME, EVENT_BUS_BUFFER_SIZE, EVENT_BUS_CONNECTION_TIMEOUT_MS, + EVENT_BUS_MAX_RECONNECT_ATTEMPTS, }; // ============================================================================ @@ -121,7 +120,7 @@ impl EventBusConfig { pub fn tokio() -> Self { Self { provider: EventBusProvider::Tokio, - capacity: EVENTS_TOKIO_DEFAULT_CAPACITY, + capacity: EVENT_BUS_BUFFER_SIZE, nats_url: None, nats_client_name: Some(DEFAULT_NATS_CLIENT_NAME.to_owned()), connection_timeout_ms: EVENT_BUS_CONNECTION_TIMEOUT_MS, @@ -146,7 +145,7 @@ impl EventBusConfig { pub fn nats(url: impl Into) -> Self { Self { provider: EventBusProvider::Nats, - capacity: EVENTS_TOKIO_DEFAULT_CAPACITY, + capacity: EVENT_BUS_BUFFER_SIZE, nats_url: Some(url.into()), nats_client_name: Some(DEFAULT_NATS_CLIENT_NAME.to_owned()), connection_timeout_ms: EVENT_BUS_CONNECTION_TIMEOUT_MS, diff --git a/crates/mcb-infrastructure/src/di/bootstrap.rs b/crates/mcb-infrastructure/src/di/bootstrap.rs index 127ddd0f1..5994f58ed 100644 --- a/crates/mcb-infrastructure/src/di/bootstrap.rs +++ b/crates/mcb-infrastructure/src/di/bootstrap.rs @@ -11,11 +11,11 @@ use std::sync::Arc; use mcb_domain::error::Result; use mcb_domain::ports::{ AgentRepository, CacheAdminInterface, CryptoProvider, EmbeddingAdminInterface, - EventBusProvider, FileHashRepository, HighlightServiceInterface, IndexingOperationsInterface, - IssueEntityRepository, LanguageAdminInterface, LifecycleManaged, MemoryRepository, - OrgEntityRepository, PerformanceMetricsInterface, PlanEntityRepository, ProjectDetectorService, - ProjectRepository, ShutdownCoordinator, VcsEntityRepository, VcsProvider, - VectorStoreAdminInterface, + EventBusProvider, FileHashRepository, FileSystemProvider, HighlightServiceInterface, + IndexingOperationsInterface, IssueEntityRepository, LanguageAdminInterface, LifecycleManaged, + MemoryRepository, OrgEntityRepository, PerformanceMetricsInterface, PlanEntityRepository, + ProjectDetectorService, ProjectRepository, ShutdownCoordinator, TaskRunnerProvider, + VcsEntityRepository, VcsProvider, VectorStoreAdminInterface, }; use crate::config::{AppConfig, ConfigLoader}; @@ -32,18 +32,14 @@ use crate::di::handles::{ CacheProviderHandle, EmbeddingProviderHandle, LanguageProviderHandle, VectorStoreProviderHandle, }; use crate::di::provider_resolvers::{ - CacheProviderResolver, EmbeddingProviderResolver, LanguageProviderResolver, - VectorStoreProviderResolver, + CacheProviderResolver, EmbeddingProviderResolver, EventBusProviderResolver, + FileSystemProviderResolver, LanguageProviderResolver, TaskRunnerProviderResolver, + VcsProviderResolver, VectorStoreProviderResolver, }; use crate::infrastructure::admin::{AtomicPerformanceMetrics, DefaultIndexingOperations}; use crate::infrastructure::lifecycle::DefaultShutdownCoordinator; use crate::project::ProjectService; use crate::services::HighlightServiceImpl; -use mcb_providers::database::{ - SqliteFileHashConfig, SqliteFileHashRepository, SqliteMemoryRepository, - create_agent_repository_from_executor, create_project_repository_from_executor, -}; -use mcb_providers::events::TokioEventBusProvider; /// Application context with provider handles and infrastructure services pub struct AppContext { @@ -81,6 +77,8 @@ pub struct AppContext { shutdown_coordinator: Arc, performance_metrics: Arc, indexing_operations: Arc, + file_system_provider: Arc, + task_runner_provider: Arc, /// Services eligible for lifecycle management pub lifecycle_services: Vec>, @@ -106,172 +104,37 @@ pub struct AppContext { } impl AppContext { - /// Get embedding provider handle - #[must_use] - pub fn embedding_handle(&self) -> Arc { - Arc::clone(&self.embedding_handle) - } - - /// Get vector store provider handle - #[must_use] - pub fn vector_store_handle(&self) -> Arc { - Arc::clone(&self.vector_store_handle) - } - - /// Get cache provider handle - #[must_use] - pub fn cache_handle(&self) -> Arc { - Arc::clone(&self.cache_handle) - } - - /// Get language provider handle - #[must_use] - pub fn language_handle(&self) -> Arc { - Arc::clone(&self.language_handle) - } - - /// Get embedding provider resolver - #[must_use] - pub fn embedding_resolver(&self) -> Arc { - Arc::clone(&self.embedding_resolver) - } - - /// Get vector store provider resolver - #[must_use] - pub fn vector_store_resolver(&self) -> Arc { - Arc::clone(&self.vector_store_resolver) - } - - /// Get cache provider resolver - #[must_use] - pub fn cache_resolver(&self) -> Arc { - Arc::clone(&self.cache_resolver) - } - - /// Get language provider resolver - #[must_use] - pub fn language_resolver(&self) -> Arc { - Arc::clone(&self.language_resolver) - } - - /// Get embedding admin service - #[must_use] - pub fn embedding_admin(&self) -> Arc { - Arc::clone(&self.embedding_admin) - } - - /// Get vector store admin service - #[must_use] - pub fn vector_store_admin(&self) -> Arc { - Arc::clone(&self.vector_store_admin) - } - - /// Get cache admin service - #[must_use] - pub fn cache_admin(&self) -> Arc { - Arc::clone(&self.cache_admin) - } - - /// Get language admin service - #[must_use] - pub fn language_admin(&self) -> Arc { - Arc::clone(&self.language_admin) - } - - /// Get event bus - #[must_use] - pub fn event_bus(&self) -> Arc { - Arc::clone(&self.event_bus) - } - - /// Get shutdown coordinator - #[must_use] - pub fn shutdown(&self) -> Arc { - Arc::clone(&self.shutdown_coordinator) - } - - /// Get performance metrics - #[must_use] - pub fn performance(&self) -> Arc { - Arc::clone(&self.performance_metrics) - } - - /// Get indexing operations - #[must_use] - pub fn indexing(&self) -> Arc { - Arc::clone(&self.indexing_operations) - } - - /// Get memory repository - #[must_use] - pub fn memory_repository(&self) -> Arc { - Arc::clone(&self.memory_repository) - } - - /// Get agent repository - #[must_use] - pub fn agent_repository(&self) -> Arc { - Arc::clone(&self.agent_repository) - } - - /// Get project repository - #[must_use] - pub fn project_repository(&self) -> Arc { - Arc::clone(&self.project_repository) - } - - /// Get VCS provider - #[must_use] - pub fn vcs_provider(&self) -> Arc { - Arc::clone(&self.vcs_provider) - } - - /// Get project service - #[must_use] - pub fn project_service(&self) -> Arc { - Arc::clone(&self.project_service) - } - - /// Get VCS entity repository - #[must_use] - pub fn vcs_entity_repository(&self) -> Arc { - Arc::clone(&self.vcs_entity_repository) - } - - /// Get plan entity repository - #[must_use] - pub fn plan_entity_repository(&self) -> Arc { - Arc::clone(&self.plan_entity_repository) - } - - /// Get issue entity repository - #[must_use] - pub fn issue_entity_repository(&self) -> Arc { - Arc::clone(&self.issue_entity_repository) - } - - /// Get org entity repository - #[must_use] - pub fn org_entity_repository(&self) -> Arc { - Arc::clone(&self.org_entity_repository) - } - - /// Get file hash repository - #[must_use] - pub fn file_hash_repository(&self) -> Arc { - Arc::clone(&self.file_hash_repository) - } - - /// Get highlight service - #[must_use] - pub fn highlight_service(&self) -> Arc { - Arc::clone(&self.highlight_service) - } - - /// Get crypto service - #[must_use] - pub fn crypto_service(&self) -> Arc { - Arc::clone(&self.crypto_service) + mcb_domain::arc_getters! { + embedding_handle: EmbeddingProviderHandle, + vector_store_handle: VectorStoreProviderHandle, + cache_handle: CacheProviderHandle, + language_handle: LanguageProviderHandle, + embedding_resolver: EmbeddingProviderResolver, + vector_store_resolver: VectorStoreProviderResolver, + cache_resolver: CacheProviderResolver, + language_resolver: LanguageProviderResolver, + embedding_admin: dyn EmbeddingAdminInterface, + vector_store_admin: dyn VectorStoreAdminInterface, + cache_admin: dyn CacheAdminInterface, + language_admin: dyn LanguageAdminInterface, + event_bus: dyn EventBusProvider, + shutdown: dyn ShutdownCoordinator => shutdown_coordinator, + performance: dyn PerformanceMetricsInterface => performance_metrics, + indexing: dyn IndexingOperationsInterface => indexing_operations, + file_system_provider: dyn FileSystemProvider, + task_runner_provider: dyn TaskRunnerProvider, + memory_repository: dyn MemoryRepository, + agent_repository: dyn AgentRepository, + project_repository: dyn ProjectRepository, + vcs_provider: dyn VcsProvider, + project_service: dyn ProjectDetectorService, + vcs_entity_repository: dyn VcsEntityRepository, + plan_entity_repository: dyn PlanEntityRepository, + issue_entity_repository: dyn IssueEntityRepository, + org_entity_repository: dyn OrgEntityRepository, + file_hash_repository: dyn FileHashRepository, + highlight_service: dyn HighlightServiceInterface, + crypto_service: dyn CryptoProvider, } /// Build domain services for the server layer @@ -293,6 +156,8 @@ impl AppContext { let indexing_ops = self.indexing(); let event_bus = self.event_bus(); + let file_system_provider = self.file_system_provider(); + let task_runner_provider = self.task_runner_provider(); let project_id = current_project_id()?; @@ -312,6 +177,8 @@ impl AppContext { language_chunker, indexing_ops, event_bus, + file_system_provider, + task_runner_provider, memory_repository, agent_repository, file_hash_repository, @@ -435,13 +302,19 @@ pub async fn init_app(config: AppConfig) -> Result { // Create Infrastructure Services // ======================================================================== - let event_bus: Arc = Arc::new(TokioEventBusProvider::new()); + let event_bus_resolver = EventBusProviderResolver::new(Arc::clone(&config)); + let event_bus: Arc = event_bus_resolver.resolve_from_config()?; let shutdown_coordinator: Arc = Arc::new(DefaultShutdownCoordinator::new()); let performance_metrics: Arc = Arc::new(AtomicPerformanceMetrics::new()); let indexing_operations: Arc = Arc::new(DefaultIndexingOperations::new()); + let fs_resolver = FileSystemProviderResolver::new(Arc::clone(&config)); + let file_system_provider: Arc = fs_resolver.resolve_from_config()?; + let task_runner_resolver = TaskRunnerProviderResolver::new(Arc::clone(&config)); + let task_runner_provider: Arc = + task_runner_resolver.resolve_from_config()?; mcb_domain::info!("bootstrap", "Created infrastructure services"); @@ -461,6 +334,9 @@ pub async fn init_app(config: AppConfig) -> Result { })?; let db_resolver = DatabaseProviderResolver::new(Arc::clone(&config)); + let db_provider = db_resolver.resolve_provider().map_err(|e| { + mcb_domain::error::Error::internal(format!("Failed to resolve database provider: {e}")) + })?; let db_executor = db_resolver .resolve_and_connect(memory_db_path.as_path()) .await @@ -469,32 +345,25 @@ pub async fn init_app(config: AppConfig) -> Result { })?; let memory_repository: Arc = - Arc::new(SqliteMemoryRepository::new(Arc::clone(&db_executor))); - let agent_repository = create_agent_repository_from_executor(Arc::clone(&db_executor)); - let project_repository = create_project_repository_from_executor(Arc::clone(&db_executor)); + db_provider.create_memory_repository(Arc::clone(&db_executor)); + let agent_repository = db_provider.create_agent_repository(Arc::clone(&db_executor)); + let project_repository = db_provider.create_project_repository(Arc::clone(&db_executor)); let project_id = current_project_id()?; let file_hash_repository: Arc = - Arc::new(SqliteFileHashRepository::new( - Arc::clone(&db_executor), - SqliteFileHashConfig::default(), - project_id, - )); + db_provider.create_file_hash_repository(Arc::clone(&db_executor), project_id); - let vcs_provider = crate::di::vcs::default_vcs_provider(); + let vcs_resolver = VcsProviderResolver::new(Arc::clone(&config)); + let vcs_provider = vcs_resolver.resolve_from_config()?; let project_service: Arc = Arc::new(ProjectService::new()); - let vcs_entity_repository: Arc = Arc::new( - mcb_providers::database::SqliteVcsEntityRepository::new(Arc::clone(&db_executor)), - ); - let plan_entity_repository: Arc = Arc::new( - mcb_providers::database::SqlitePlanEntityRepository::new(Arc::clone(&db_executor)), - ); - let issue_entity_repository: Arc = Arc::new( - mcb_providers::database::SqliteIssueEntityRepository::new(Arc::clone(&db_executor)), - ); - let org_entity_repository: Arc = Arc::new( - mcb_providers::database::SqliteOrgEntityRepository::new(Arc::clone(&db_executor)), - ); + let vcs_entity_repository: Arc = + db_provider.create_vcs_entity_repository(Arc::clone(&db_executor)); + let plan_entity_repository: Arc = + db_provider.create_plan_entity_repository(Arc::clone(&db_executor)); + let issue_entity_repository: Arc = + db_provider.create_issue_entity_repository(Arc::clone(&db_executor)); + let org_entity_repository: Arc = + db_provider.create_org_entity_repository(Arc::clone(&db_executor)); let highlight_service: Arc = Arc::new(HighlightServiceImpl::new()); @@ -525,6 +394,8 @@ pub async fn init_app(config: AppConfig) -> Result { shutdown_coordinator, performance_metrics, indexing_operations, + file_system_provider, + task_runner_provider, memory_repository, agent_repository, project_repository, diff --git a/crates/mcb-infrastructure/src/di/database_resolver.rs b/crates/mcb-infrastructure/src/di/database_resolver.rs index 1e62cefdb..11be782af 100644 --- a/crates/mcb-infrastructure/src/di/database_resolver.rs +++ b/crates/mcb-infrastructure/src/di/database_resolver.rs @@ -35,14 +35,22 @@ impl DatabaseProviderResolver { /// /// Returns an error if the provider resolution or database connection fails. pub async fn resolve_and_connect(&self, path: &Path) -> Result> { - let provider_name = self.config.providers.database.provider.as_str(); - - let config = DatabaseProviderConfig::new(provider_name); - let provider = self.resolve_from_override(&config)?; + let provider = self.resolve_provider()?; provider.connect(path).await } + /// Resolve the configured database provider instance from linkme registry. + /// + /// # Errors + /// + /// Returns an error when the configured provider name cannot be resolved. + pub fn resolve_provider(&self) -> Result> { + let provider_name = self.config.providers.database.provider.as_str(); + let config = DatabaseProviderConfig::new(provider_name); + self.resolve_from_override(&config) + } + /// Resolve a provider from specific configuration /// /// # Errors diff --git a/crates/mcb-infrastructure/src/di/mod.rs b/crates/mcb-infrastructure/src/di/mod.rs index dc7fafc96..890973f10 100644 --- a/crates/mcb-infrastructure/src/di/mod.rs +++ b/crates/mcb-infrastructure/src/di/mod.rs @@ -46,6 +46,7 @@ pub mod provider_resolvers; pub mod repositories; pub mod resolver; pub mod test_factory; +#[allow(missing_docs)] pub mod vcs; pub use admin::{ @@ -60,8 +61,9 @@ pub use handles::{ }; pub use modules::{DomainServicesContainer, DomainServicesFactory, ServiceDependencies}; pub use provider_resolvers::{ - CacheProviderResolver, EmbeddingProviderResolver, LanguageProviderResolver, - VectorStoreProviderResolver, + CacheProviderResolver, EmbeddingProviderResolver, EventBusProviderResolver, + FileSystemProviderResolver, LanguageProviderResolver, TaskRunnerProviderResolver, + VcsProviderResolver, VectorStoreProviderResolver, }; pub use repositories::{ create_memory_repository, create_memory_repository_with_executor, create_vcs_entity_repository, diff --git a/crates/mcb-infrastructure/src/di/modules/domain_services.rs b/crates/mcb-infrastructure/src/di/modules/domain_services.rs index b97f876af..47b3a70a0 100644 --- a/crates/mcb-infrastructure/src/di/modules/domain_services.rs +++ b/crates/mcb-infrastructure/src/di/modules/domain_services.rs @@ -24,11 +24,12 @@ use mcb_application::use_cases::{ use mcb_domain::error::Result; use mcb_domain::ports::{ AgentRepository, AgentSessionServiceInterface, ContextServiceInterface, CryptoProvider, - EmbeddingProvider, EventBusProvider, FileHashRepository, IndexingOperationsInterface, - IndexingServiceInterface, IssueEntityRepository, LanguageChunkingProvider, MemoryRepository, - MemoryServiceInterface, OrgEntityRepository, PlanEntityRepository, ProjectDetectorService, - ProjectRepository, SearchServiceInterface, ValidationServiceInterface, VcsEntityRepository, - VcsProvider, VectorStoreProvider, + EmbeddingProvider, EventBusProvider, FileHashRepository, FileSystemProvider, + IndexingOperationsInterface, IndexingServiceInterface, IssueEntityRepository, + LanguageChunkingProvider, MemoryRepository, MemoryServiceInterface, OrgEntityRepository, + PlanEntityRepository, ProjectDetectorService, ProjectRepository, SearchServiceInterface, + TaskRunnerProvider, ValidationServiceInterface, VcsEntityRepository, VcsProvider, + VectorStoreProvider, }; use super::super::bootstrap::AppContext; @@ -96,6 +97,7 @@ pub struct DomainServicesContainer { /// * `agent_repository` - Repository for agent session data /// * `vcs_provider` - Version control system provider /// * `project_service` - Service for project detection and management +#[allow(missing_docs)] pub struct ServiceDependencies { /// Unique identifier for the current project pub project_id: String, @@ -115,6 +117,8 @@ pub struct ServiceDependencies { pub indexing_ops: Arc, /// Event bus for domain events pub event_bus: Arc, + pub file_system_provider: Arc, + pub task_runner_provider: Arc, /// Repository for memory persistence pub memory_repository: Arc, /// Repository for agent session data @@ -145,6 +149,8 @@ struct IndexingServiceInputs { language_chunker: Arc, indexing_ops: Arc, event_bus: Arc, + file_system_provider: Arc, + task_runner_provider: Arc, file_hash_repository: Arc, supported_extensions: Vec, } @@ -170,6 +176,8 @@ impl DomainServicesFactory { language_chunker: inputs.language_chunker, indexing_ops: inputs.indexing_ops, event_bus: inputs.event_bus, + file_system_provider: inputs.file_system_provider, + task_runner_provider: inputs.task_runner_provider, supported_extensions: inputs.supported_extensions, }, file_hash_repository: inputs.file_hash_repository, @@ -209,6 +217,8 @@ impl DomainServicesFactory { language_chunker: deps.language_chunker, indexing_ops: deps.indexing_ops, event_bus: deps.event_bus, + file_system_provider: deps.file_system_provider, + task_runner_provider: deps.task_runner_provider, file_hash_repository: deps.file_hash_repository, supported_extensions: deps.config.mcp.indexing.supported_extensions.clone(), }); @@ -253,6 +263,8 @@ impl DomainServicesFactory { ) -> Result> { let indexing_ops = app_context.indexing(); let event_bus = app_context.event_bus(); + let file_system_provider = app_context.file_system_provider(); + let task_runner_provider = app_context.task_runner_provider(); let language_chunker = app_context.language_handle().get(); let context_service = Self::create_context_service(app_context).await?; let file_hash_repository = app_context.file_hash_repository(); @@ -263,6 +275,8 @@ impl DomainServicesFactory { language_chunker, indexing_ops, event_bus, + file_system_provider, + task_runner_provider, file_hash_repository, supported_extensions, })) diff --git a/crates/mcb-infrastructure/src/di/provider_resolvers.rs b/crates/mcb-infrastructure/src/di/provider_resolvers.rs index b2610f5ca..03fbace49 100644 --- a/crates/mcb-infrastructure/src/di/provider_resolvers.rs +++ b/crates/mcb-infrastructure/src/di/provider_resolvers.rs @@ -16,17 +16,23 @@ use std::collections::HashMap; use std::sync::Arc; use mcb_domain::ports::{ - CacheProvider, EmbeddingProvider, LanguageChunkingProvider, VectorStoreProvider, + CacheProvider, EmbeddingProvider, EventBusProvider, FileSystemProvider, + LanguageChunkingProvider, TaskRunnerProvider, VcsProvider, VectorStoreProvider, }; use mcb_domain::registry::cache::{CacheProviderConfig, resolve_cache_provider}; use mcb_domain::registry::embedding::{EmbeddingProviderConfig, resolve_embedding_provider}; +use mcb_domain::registry::event_bus::{EventBusProviderConfig, resolve_event_bus_provider}; +use mcb_domain::registry::fs::{FileSystemProviderConfig, resolve_file_system_provider}; use mcb_domain::registry::language::{LanguageProviderConfig, resolve_language_provider}; +use mcb_domain::registry::task_runner::{TaskRunnerProviderConfig, resolve_task_runner_provider}; +use mcb_domain::registry::vcs::{VcsProviderConfig, resolve_vcs_provider}; use mcb_domain::registry::vector_store::{ VectorStoreProviderConfig, resolve_vector_store_provider, }; use mcb_domain::value_objects::{EmbeddingConfig, VectorStoreConfig}; use crate::config::AppConfig; +use crate::config::types::EventBusProvider as InfraEventBusProvider; use crate::constants::providers::{ DEFAULT_DB_CONFIG_NAME, FALLBACK_EMBEDDING_PROVIDER, FALLBACK_VECTOR_STORE_PROVIDER, }; @@ -330,6 +336,164 @@ impl LanguageProviderResolver { } impl_resolver_common!(LanguageProviderResolver); +#[allow(missing_docs)] +pub struct EventBusProviderResolver { + config: Arc, +} + +#[allow(missing_docs)] +impl EventBusProviderResolver { + /// Resolves the event bus provider from application config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved or configured. + pub fn resolve_from_config(&self) -> mcb_domain::error::Result> { + let mut cfg = EventBusProviderConfig::new( + match self.config.system.infrastructure.event_bus.provider { + InfraEventBusProvider::Tokio => "tokio", + InfraEventBusProvider::Nats => "nats", + }, + ); + + cfg.extra.insert( + "capacity".to_owned(), + self.config + .system + .infrastructure + .event_bus + .capacity + .to_string(), + ); + if let Some(url) = &self.config.system.infrastructure.event_bus.nats_url { + cfg.extra.insert("url".to_owned(), url.clone()); + } + if let Some(name) = &self.config.system.infrastructure.event_bus.nats_client_name { + cfg.extra.insert("client_name".to_owned(), name.clone()); + } + + resolve_event_bus_provider(&cfg) + } + + /// Resolves the event bus provider from an override config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved. + pub fn resolve_from_override( + &self, + override_config: &EventBusProviderConfig, + ) -> mcb_domain::error::Result> { + resolve_event_bus_provider(override_config) + } + + #[must_use] + pub fn list_available(&self) -> Vec<(&'static str, &'static str)> { + mcb_domain::registry::event_bus::list_event_bus_providers() + } +} +impl_resolver_common!(EventBusProviderResolver); + +#[allow(missing_docs)] +pub struct VcsProviderResolver { + config: Arc, +} + +#[allow(missing_docs)] +impl VcsProviderResolver { + /// Resolves the VCS provider from application config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved. + pub fn resolve_from_config(&self) -> mcb_domain::error::Result> { + let _ = self.config.as_ref(); + crate::di::vcs::default_vcs_provider() + } + + /// Resolves the VCS provider from an override config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved. + pub fn resolve_from_override( + &self, + override_config: &VcsProviderConfig, + ) -> mcb_domain::error::Result> { + resolve_vcs_provider(override_config) + } + + #[must_use] + pub fn list_available(&self) -> Vec<(&'static str, &'static str)> { + mcb_domain::registry::vcs::list_vcs_providers() + } +} +impl_resolver_common!(VcsProviderResolver); + +#[allow(missing_docs)] +pub struct FileSystemProviderResolver { + config: Arc, +} + +#[allow(missing_docs)] +impl FileSystemProviderResolver { + /// Resolves the file system provider from application config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved. + pub fn resolve_from_config(&self) -> mcb_domain::error::Result> { + let _ = self.config.as_ref(); + resolve_file_system_provider(&FileSystemProviderConfig::new("local")) + } + + /// Resolves the file system provider from an override config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved. + pub fn resolve_from_override( + &self, + override_config: &FileSystemProviderConfig, + ) -> mcb_domain::error::Result> { + resolve_file_system_provider(override_config) + } + + #[must_use] + pub fn list_available(&self) -> Vec<(&'static str, &'static str)> { + mcb_domain::registry::fs::list_file_system_providers() + } +} +impl_resolver_common!(FileSystemProviderResolver); + +#[allow(missing_docs)] +pub struct TaskRunnerProviderResolver { + config: Arc, +} + +#[allow(missing_docs)] +impl TaskRunnerProviderResolver { + /// Resolves the task runner provider from application config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved. + pub fn resolve_from_config(&self) -> mcb_domain::error::Result> { + let _ = self.config.as_ref(); + resolve_task_runner_provider(&TaskRunnerProviderConfig::new("tokio")) + } + + /// Resolves the task runner provider from an override config. + /// + /// # Errors + /// Returns an error if the provider cannot be resolved. + pub fn resolve_from_override( + &self, + override_config: &TaskRunnerProviderConfig, + ) -> mcb_domain::error::Result> { + resolve_task_runner_provider(override_config) + } + + #[must_use] + pub fn list_available(&self) -> Vec<(&'static str, &'static str)> { + mcb_domain::registry::task_runner::list_task_runner_providers() + } +} +impl_resolver_common!(TaskRunnerProviderResolver); + // ============================================================================ // Helper Functions // ============================================================================ diff --git a/crates/mcb-infrastructure/src/di/test_factory.rs b/crates/mcb-infrastructure/src/di/test_factory.rs index 09101b658..130625961 100644 --- a/crates/mcb-infrastructure/src/di/test_factory.rs +++ b/crates/mcb-infrastructure/src/di/test_factory.rs @@ -61,6 +61,8 @@ pub fn create_test_dependencies( language_chunker: app_context.language_handle().get(), indexing_ops: app_context.indexing(), event_bus: app_context.event_bus(), + file_system_provider: app_context.file_system_provider(), + task_runner_provider: app_context.task_runner_provider(), memory_repository, agent_repository, file_hash_repository, diff --git a/crates/mcb-infrastructure/src/di/vcs.rs b/crates/mcb-infrastructure/src/di/vcs.rs index af00fc2c7..718bfab7d 100644 --- a/crates/mcb-infrastructure/src/di/vcs.rs +++ b/crates/mcb-infrastructure/src/di/vcs.rs @@ -1,20 +1,128 @@ -//! -//! **Documentation**: [docs/modules/infrastructure.md](../../../../docs/modules/infrastructure.md#dependency-injection) -//! -//! VCS provider factory for standalone/server composition. -//! -//! Provides a default VCS provider so that the server layer does not -//! import concrete providers directly (CA006). - +use std::collections::HashSet; +use std::path::{Path, PathBuf}; use std::sync::Arc; +use async_trait::async_trait; + +use mcb_domain::entities::vcs::{RefDiff, RepositoryId, VcsBranch, VcsCommit, VcsRepository}; +use mcb_domain::error::{Error, Result}; use mcb_domain::ports::VcsProvider; -use mcb_providers::vcs; - -/// Returns the default VCS provider for standalone and server modes. -/// -/// Centralizes provider instantiation in the infrastructure layer. -#[must_use] -pub fn default_vcs_provider() -> Arc { - vcs::default_vcs_provider() +use mcb_domain::registry::vcs::{VcsProviderConfig, list_vcs_providers, resolve_vcs_provider}; + +#[allow(missing_docs)] +pub struct DynamicVcsProvider { + providers: Vec>, +} + +impl DynamicVcsProvider { + fn from_registry() -> Result { + let mut providers = Vec::new(); + for (name, _) in list_vcs_providers() { + providers.push(resolve_vcs_provider(&VcsProviderConfig::new(name))?); + } + + if providers.is_empty() { + return Err(Error::configuration( + "VCS: no providers registered in linkme registry", + )); + } + + Ok(Self { providers }) + } + + async fn provider_and_repo_for_path( + &self, + path: &Path, + ) -> Result<(Arc, VcsRepository)> { + let mut last_error: Option = None; + + for provider in &self.providers { + match provider.open_repository(path).await { + Ok(repo) => return Ok((Arc::clone(provider), repo)), + Err(e) => last_error = Some(e), + } + } + + Err(last_error.unwrap_or_else(|| { + Error::vcs(format!( + "No registered VCS provider can open path: {}", + path.display() + )) + })) + } +} + +#[async_trait] +impl VcsProvider for DynamicVcsProvider { + async fn open_repository(&self, path: &Path) -> Result { + let (_, repo) = self.provider_and_repo_for_path(path).await?; + Ok(repo) + } + + fn repository_id(&self, repo: &VcsRepository) -> RepositoryId { + self.providers.first().map_or_else( + || repo.id().clone(), + |provider| provider.repository_id(repo), + ) + } + + async fn list_branches(&self, repo: &VcsRepository) -> Result> { + let (provider, concrete_repo) = self.provider_and_repo_for_path(repo.path()).await?; + provider.list_branches(&concrete_repo).await + } + + async fn commit_history( + &self, + repo: &VcsRepository, + branch: &str, + limit: Option, + ) -> Result> { + let (provider, concrete_repo) = self.provider_and_repo_for_path(repo.path()).await?; + provider.commit_history(&concrete_repo, branch, limit).await + } + + async fn list_files(&self, repo: &VcsRepository, branch: &str) -> Result> { + let (provider, concrete_repo) = self.provider_and_repo_for_path(repo.path()).await?; + provider.list_files(&concrete_repo, branch).await + } + + async fn read_file(&self, repo: &VcsRepository, branch: &str, path: &Path) -> Result { + let (provider, concrete_repo) = self.provider_and_repo_for_path(repo.path()).await?; + provider.read_file(&concrete_repo, branch, path).await + } + + fn vcs_name(&self) -> &str { + "dynamic" + } + + async fn diff_refs( + &self, + repo: &VcsRepository, + base_ref: &str, + head_ref: &str, + ) -> Result { + let (provider, concrete_repo) = self.provider_and_repo_for_path(repo.path()).await?; + provider.diff_refs(&concrete_repo, base_ref, head_ref).await + } + + async fn list_repositories(&self, root: &Path) -> Result> { + let mut out = Vec::new(); + let mut seen = HashSet::new(); + + for provider in &self.providers { + for repo in provider.list_repositories(root).await? { + let key = repo.path().to_path_buf(); + if seen.insert(key) { + out.push(repo); + } + } + } + + Ok(out) + } +} + +#[allow(missing_docs)] +pub fn default_vcs_provider() -> Result> { + Ok(Arc::new(DynamicVcsProvider::from_registry()?)) } diff --git a/crates/mcb-infrastructure/tests/integration/di/dispatch_tests.rs b/crates/mcb-infrastructure/tests/integration/di/dispatch_tests.rs index 137236cc0..b27c233af 100644 --- a/crates/mcb-infrastructure/tests/integration/di/dispatch_tests.rs +++ b/crates/mcb-infrastructure/tests/integration/di/dispatch_tests.rs @@ -6,6 +6,7 @@ //! providers are registered via linkme distributed slices. use mcb_infrastructure::config::{AppConfig, ConfigLoader}; +use mcb_infrastructure::di::VcsProviderResolver; use mcb_infrastructure::di::bootstrap::init_app; use serial_test::serial; @@ -195,3 +196,57 @@ async fn test_infrastructure_services_from_app_context() { "Indexing service should have valid Arc reference" ); } + +#[tokio::test] +#[serial] +async fn test_vcs_provider_resolver_detects_git_repository_by_structure() +-> Result<(), Box> { + let default_path = + std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../config/default.toml"); + let config = ConfigLoader::new().with_config_path(default_path).load()?; + let resolver = VcsProviderResolver::new(std::sync::Arc::new(config)); + + let vcs_provider = resolver.resolve_from_config()?; + assert!(!vcs_provider.vcs_name().is_empty()); + + let temp_dir = tempfile::tempdir()?; + let git_init = std::process::Command::new("git") + .args(["init", "--quiet"]) + .current_dir(temp_dir.path()) + .output()?; + + assert!( + git_init.status.success(), + "git init must succeed to validate VCS detection" + ); + + std::fs::write(temp_dir.path().join("README.md"), "mcb test")?; + let add = std::process::Command::new("git") + .args(["add", "."]) + .current_dir(temp_dir.path()) + .output()?; + assert!(add.status.success(), "git add must succeed"); + let commit = std::process::Command::new("git") + .args([ + "-c", + "user.name=MCB Test", + "-c", + "user.email=test@mcb.local", + "commit", + "-m", + "init", + "--quiet", + ]) + .current_dir(temp_dir.path()) + .output()?; + assert!(commit.status.success(), "git commit must succeed"); + + let opened_repo = vcs_provider.open_repository(temp_dir.path()).await?; + assert_eq!(opened_repo.path(), temp_dir.path()); + + let repos = vcs_provider.list_repositories(temp_dir.path()).await?; + assert_eq!(repos.len(), 1, "expected exactly one git repository"); + assert_eq!(repos[0].path(), temp_dir.path()); + + Ok(()) +} diff --git a/crates/mcb-infrastructure/tests/unit/error/error_ext_tests.rs b/crates/mcb-infrastructure/tests/unit/error/error_ext_tests.rs index dcd3635fd..7016170c9 100644 --- a/crates/mcb-infrastructure/tests/unit/error/error_ext_tests.rs +++ b/crates/mcb-infrastructure/tests/unit/error/error_ext_tests.rs @@ -22,15 +22,23 @@ fn test_error_context_extension() { } #[rstest] -#[allow(clippy::wildcard_enum_match_arm)] fn test_infra_error_creation() { let error = infra::infrastructure_error_msg("test error message"); - match error { + match &error { Error::Infrastructure { message, source } => { assert_eq!(message, "test error message"); assert!(source.is_none()); } - _ => panic!("Expected Infrastructure error"), + Error::Vcs(_) + | Error::Database(_) + | Error::Validation(_) + | Error::Embedding(_) + | Error::VectorStore(_) + | Error::Cache(_) + | Error::Config(_) + | Error::NotFound(_) + | Error::Language(_) + | Error::Other(_) => panic!("Expected Infrastructure error, got {error:?}"), } } diff --git a/crates/mcb-infrastructure/tests/unit/routing/router_tests.rs b/crates/mcb-infrastructure/tests/unit/routing/router_tests.rs index 33464a132..eead64dc9 100644 --- a/crates/mcb-infrastructure/tests/unit/routing/router_tests.rs +++ b/crates/mcb-infrastructure/tests/unit/routing/router_tests.rs @@ -96,11 +96,10 @@ fn test_get_all_health(monitor: InMemoryHealthMonitor) { // ============================================================================= #[fixture] -#[allow(clippy::clone_on_ref_ptr)] fn router_setup() -> (Arc, DefaultProviderRouter) { let monitor = Arc::new(InMemoryHealthMonitor::new()); let router = DefaultProviderRouter::new( - monitor.clone(), + Arc::clone(&monitor), vec!["provider-a".to_owned(), "provider-b".to_owned()], vec![], ); diff --git a/crates/mcb-infrastructure/tests/utils/fs_guards.rs b/crates/mcb-infrastructure/tests/utils/fs_guards.rs index db6db611b..45a9187bc 100644 --- a/crates/mcb-infrastructure/tests/utils/fs_guards.rs +++ b/crates/mcb-infrastructure/tests/utils/fs_guards.rs @@ -7,7 +7,10 @@ pub struct CurrentDirGuard { } impl CurrentDirGuard { - #[allow(clippy::missing_errors_doc)] + /// Changes process current directory to `new_dir` and returns a guard that restores it on drop. + /// + /// # Errors + /// Fails if `env::current_dir` or `env::set_current_dir` fails. pub fn new(new_dir: &Path) -> Result> { let original = env::current_dir()?; env::set_current_dir(new_dir)?; @@ -27,7 +30,10 @@ pub struct RestoreFileGuard { } impl RestoreFileGuard { - #[allow(clippy::missing_errors_doc)] + /// Moves `target` to `backup` and returns a guard that restores it on drop. + /// + /// # Errors + /// Fails if `fs::rename(target, backup)` fails. pub fn move_out(target: &Path, backup: &Path) -> Result> { fs::rename(target, backup)?; Ok(Self { diff --git a/crates/mcb-infrastructure/tests/utils/workspace.rs b/crates/mcb-infrastructure/tests/utils/workspace.rs index 76f9fa016..7d674e886 100644 --- a/crates/mcb-infrastructure/tests/utils/workspace.rs +++ b/crates/mcb-infrastructure/tests/utils/workspace.rs @@ -1,6 +1,9 @@ use std::path::{Path, PathBuf}; -#[allow(clippy::missing_errors_doc)] +/// Returns the workspace root (directory containing Cargo.lock) from CARGO_MANIFEST_DIR. +/// +/// # Errors +/// Fails if no ancestor of CARGO_MANIFEST_DIR contains Cargo.lock. pub fn workspace_root() -> Result> { let manifest_dir = PathBuf::from(env!("CARGO_MANIFEST_DIR")); for dir in manifest_dir.ancestors() { diff --git a/crates/mcb-providers/Cargo.toml b/crates/mcb-providers/Cargo.toml index d25416d27..04941042e 100644 --- a/crates/mcb-providers/Cargo.toml +++ b/crates/mcb-providers/Cargo.toml @@ -127,6 +127,9 @@ async-nats = { workspace = true } # SQLite for memory repository (uses generic schema from domain) sqlx = { workspace = true } + +# SeaORM for SQLite CRUD (simplifies repositories and row mapping) +sea-orm = { workspace = true } sha2.workspace = true walkdir.workspace = true diff --git a/crates/mcb-providers/src/database/sqlite/backend.rs b/crates/mcb-providers/src/database/sqlite/backend.rs new file mode 100644 index 000000000..600d76223 --- /dev/null +++ b/crates/mcb-providers/src/database/sqlite/backend.rs @@ -0,0 +1,64 @@ +//! +//! SQLite backend that combines the port executor with a SeaORM connection. +//! +//! Used by repositories that can use SeaORM for simpler CRUD while still +//! exposing a single `DatabaseExecutor` to the rest of the stack. + +use std::sync::Arc; + +use async_trait::async_trait; +use mcb_domain::error::Result; +use mcb_domain::ports::{DatabaseExecutor, SqlParam}; +use sea_orm::DatabaseConnection; + +use super::executor::SqliteExecutor; + +/// SQLite executor plus SeaORM connection for the same database. +/// +/// Implements [`DatabaseExecutor`] by delegating to the inner executor. +/// Repositories can downcast via `as_any()` and use [`sea_conn`](Self::sea_conn) for ORM operations. +pub struct SqliteBackend { + executor: SqliteExecutor, + sea_conn: DatabaseConnection, +} + +impl SqliteBackend { + /// Creates a backend with the given executor and SeaORM connection. + #[must_use] + pub fn new(executor: SqliteExecutor, sea_conn: DatabaseConnection) -> Self { + Self { executor, sea_conn } + } + + /// Returns the SeaORM connection for the same database. + #[must_use] + pub fn sea_conn(&self) -> &DatabaseConnection { + &self.sea_conn + } +} + +#[async_trait] +impl DatabaseExecutor for SqliteBackend { + async fn execute(&self, sql: &str, params: &[SqlParam]) -> Result<()> { + self.executor.execute(sql, params).await + } + + async fn query_one( + &self, + sql: &str, + params: &[SqlParam], + ) -> Result>> { + self.executor.query_one(sql, params).await + } + + async fn query_all( + &self, + sql: &str, + params: &[SqlParam], + ) -> Result>> { + self.executor.query_all(sql, params).await + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } +} diff --git a/crates/mcb-providers/src/database/sqlite/macros.rs b/crates/mcb-providers/src/database/sqlite/macros.rs new file mode 100644 index 000000000..7a0da13bf --- /dev/null +++ b/crates/mcb-providers/src/database/sqlite/macros.rs @@ -0,0 +1,51 @@ +//! +//! Macros for SQLite provider (row conversion, SeaORM entities). +//! +//! **Documentation**: [docs/modules/providers.md](../../../../../../docs/modules/providers.md#database) + +/// Generates a `FromRow` impl from a list of (field, extractor) pairs. +/// +/// Column name is the field name. Extractors: `req_str`, `req_i64`, `req_parsed`, `opt_str`, `opt_i64`. +/// Use manual `impl FromRow` for custom logic (e.g. computed fields, JSON, custom types). +#[macro_export] +macro_rules! from_row_simple { + ($type:ty, { $($field:ident : $ext:ident),* $(,)? }) => { + impl $crate::database::sqlite::row_convert::FromRow for $type { + fn from_row(row: &dyn mcb_domain::ports::SqlRow) -> mcb_domain::error::Result { + Ok(Self { + $($field: $ext(row, stringify!($field))?),* + }) + } + } + }; +} + +/// Generates a SeaORM entity from a table name and column list. +/// +/// First column is the primary key. List must match `mcb_domain::schema::::table()`. +/// Emits `Model`, `Relation`, `ActiveModelBehavior` and `SCHEMA_COLUMNS` (for tests). +#[macro_export] +macro_rules! sea_entity { + ($table:expr, [ ($first:ident : $first_ty:ty), $( ($f:ident : $ty:ty) ),* $(,)? ]) => { + #[allow(dead_code)] + #[derive(Clone, Debug, PartialEq, Eq, sea_orm::DeriveEntityModel)] + #[sea_orm(table_name = $table)] + pub struct Model { + #[sea_orm(primary_key)] + pub $first: $first_ty, + $( pub $f: $ty ),* + } + + #[allow(dead_code)] + #[derive(Copy, Clone, Debug, sea_orm::EnumIter, sea_orm::DeriveRelation)] + pub enum Relation {} + + impl sea_orm::ActiveModelBehavior for ActiveModel {} + + #[cfg(test)] + pub const SCHEMA_COLUMNS: &[(&str, &str)] = &[ + (stringify!($first), stringify!($first_ty)), + $( (stringify!($f), stringify!($ty)) ),* + ]; + }; +} diff --git a/crates/mcb-providers/src/database/sqlite/mod.rs b/crates/mcb-providers/src/database/sqlite/mod.rs index adab7495f..7bd9a3f41 100644 --- a/crates/mcb-providers/src/database/sqlite/mod.rs +++ b/crates/mcb-providers/src/database/sqlite/mod.rs @@ -8,7 +8,10 @@ //! (port `MemoryRepository`), and factory functions for DI. mod agent_repository; +mod backend; mod ddl; +#[macro_use] +mod macros; pub(crate) mod ensure_parent; pub mod executor; mod file_hash_repository; @@ -19,9 +22,11 @@ mod plan_entity_repository; mod project_repository; mod provider; mod row_convert; +mod sea_entities; mod vcs_entity_repository; pub use agent_repository::SqliteAgentRepository; +pub use backend::SqliteBackend; pub use ddl::SqliteSchemaDdlGenerator; pub use executor::SqliteExecutor; pub use file_hash_repository::{SqliteFileHashConfig, SqliteFileHashRepository}; diff --git a/crates/mcb-providers/src/database/sqlite/org_entity_repository.rs b/crates/mcb-providers/src/database/sqlite/org_entity_repository.rs index ca95f48eb..76d8bfb43 100644 --- a/crates/mcb-providers/src/database/sqlite/org_entity_repository.rs +++ b/crates/mcb-providers/src/database/sqlite/org_entity_repository.rs @@ -9,16 +9,22 @@ use mcb_domain::error::{Error, Result}; use mcb_domain::ports::{ ApiKeyRegistry, OrgRegistry, TeamMemberManager, TeamRegistry, UserRegistry, }; -use mcb_domain::ports::{DatabaseExecutor, SqlParam, SqlRow}; +use mcb_domain::ports::{DatabaseExecutor, SqlParam}; +use sea_orm::ActiveValue::Set; +use sea_orm::{ActiveModelTrait, EntityTrait}; +use crate::database::sqlite::row_convert::FromRow; +use crate::database::sqlite::sea_entities::organization; use crate::utils::sqlite::query as query_helpers; -use crate::utils::sqlite::row::{ - opt_i64, opt_i64_param, opt_str, opt_str_param, req_i64, req_parsed, req_str, -}; +use crate::utils::sqlite::row::{opt_i64_param, opt_str_param}; /// SQLite-backed repository for organization, user, team, and API key entities. +/// +/// When constructed with [`Self::new_with_sea`], organization CRUD uses SeaORM; +/// otherwise it uses the executor port. Other entities (users, teams, etc.) always use the executor. pub struct SqliteOrgEntityRepository { executor: Arc, + sea_conn: Option, } const INSERT_ORG_SQL: &str = "INSERT INTO organizations (id, name, slug, settings_json, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?)"; @@ -34,82 +40,48 @@ fn org_insert_params(org: &Organization) -> [SqlParam; 6] { ] } -impl SqliteOrgEntityRepository { - /// Creates a new repository using the provided database executor. - pub fn new(executor: Arc) -> Self { - Self { executor } +fn org_to_model(org: &Organization) -> organization::ActiveModel { + organization::ActiveModel { + id: Set(org.id.clone()), + name: Set(org.name.clone()), + slug: Set(org.slug.clone()), + settings_json: Set(org.settings_json.clone()), + created_at: Set(org.created_at), + updated_at: Set(org.updated_at), + ..Default::default() } } -/// Converts a SQL row to an Organization. -fn row_to_org(row: &dyn SqlRow) -> Result { - Ok(Organization { - id: req_str(row, "id")?, - name: req_str(row, "name")?, - slug: req_str(row, "slug")?, - settings_json: req_str(row, "settings_json")?, - created_at: req_i64(row, "created_at")?, - updated_at: req_i64(row, "updated_at")?, - }) -} - -/// Converts a SQL row to a User. -fn row_to_user(row: &dyn SqlRow) -> Result { - Ok(User { - id: req_str(row, "id")?, - org_id: req_str(row, "org_id")?, - email: req_str(row, "email")?, - display_name: req_str(row, "display_name")?, - role: req_parsed(row, "role")?, - api_key_hash: opt_str(row, "api_key_hash")?, - created_at: req_i64(row, "created_at")?, - updated_at: req_i64(row, "updated_at")?, - }) -} - -/// Converts a SQL row to a Team. -fn row_to_team(row: &dyn SqlRow) -> Result { - Ok(Team { - id: req_str(row, "id")?, - org_id: req_str(row, "org_id")?, - name: req_str(row, "name")?, - created_at: req_i64(row, "created_at")?, - }) +fn model_to_org(m: organization::Model) -> Organization { + Organization { + id: m.id, + name: m.name, + slug: m.slug, + settings_json: m.settings_json, + created_at: m.created_at, + updated_at: m.updated_at, + } } -use mcb_domain::utils::id; -use mcb_domain::value_objects::ids::TeamMemberId; - -// ... - -/// Converts a SQL row to a `TeamMember`. -fn row_to_team_member(row: &dyn SqlRow) -> Result { - let team_id = req_str(row, "team_id")?; - let user_id = req_str(row, "user_id")?; - let id_uuid = id::deterministic("team_member", &format!("{team_id}:{user_id}")); - - Ok(TeamMember { - id: TeamMemberId::from_uuid(id_uuid), - team_id, - user_id, - role: req_parsed(row, "role")?, - joined_at: req_i64(row, "joined_at")?, - }) -} +impl SqliteOrgEntityRepository { + /// Creates a new repository using the provided database executor (no SeaORM). + pub fn new(executor: Arc) -> Self { + Self { + executor, + sea_conn: None, + } + } -/// Converts a SQL row to an `ApiKey`. -fn row_to_api_key(row: &dyn SqlRow) -> Result { - Ok(ApiKey { - id: req_str(row, "id")?, - user_id: req_str(row, "user_id")?, - org_id: req_str(row, "org_id")?, - key_hash: req_str(row, "key_hash")?, - name: req_str(row, "name")?, - scopes_json: req_str(row, "scopes_json")?, - expires_at: opt_i64(row, "expires_at")?, - created_at: req_i64(row, "created_at")?, - revoked_at: opt_i64(row, "revoked_at")?, - }) + /// Creates a repository with a SeaORM connection for organization CRUD. + pub fn new_with_sea( + executor: Arc, + sea_conn: sea_orm::DatabaseConnection, + ) -> Self { + Self { + executor, + sea_conn: Some(sea_conn), + } + } } #[async_trait] @@ -117,17 +89,31 @@ fn row_to_api_key(row: &dyn SqlRow) -> Result { impl OrgRegistry for SqliteOrgEntityRepository { /// Creates a new organization. async fn create_org(&self, org: &Organization) -> Result<()> { + if let Some(ref db) = self.sea_conn { + let am = org_to_model(org); + am.insert(db) + .await + .map_err(|e| Error::memory_with_source("SeaORM insert organization", e))?; + return Ok(()); + } let params = org_insert_params(org); query_helpers::execute(&self.executor, INSERT_ORG_SQL, ¶ms).await } /// Retrieves an organization by ID. async fn get_org(&self, id: &str) -> Result { + if let Some(ref db) = self.sea_conn { + let opt = organization::Entity::find_by_id(id.to_string()) + .one(db) + .await + .map_err(|e| Error::memory_with_source("SeaORM find organization", e))?; + return Error::not_found_or(opt.map(model_to_org), "Organization", id); + } let org = query_helpers::query_one( &self.executor, "SELECT * FROM organizations WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_org, + Organization::from_row, ) .await?; Error::not_found_or(org, "Organization", id) @@ -135,11 +121,18 @@ impl OrgRegistry for SqliteOrgEntityRepository { /// Lists all organizations. async fn list_orgs(&self) -> Result> { + if let Some(ref db) = self.sea_conn { + let models = organization::Entity::find() + .all(db) + .await + .map_err(|e| Error::memory_with_source("SeaORM list organizations", e))?; + return Ok(models.into_iter().map(model_to_org).collect()); + } query_helpers::query_all( &self.executor, "SELECT * FROM organizations", &[], - row_to_org, + Organization::from_row, "org entity", ) .await @@ -147,6 +140,13 @@ impl OrgRegistry for SqliteOrgEntityRepository { /// Updates an existing organization. async fn update_org(&self, org: &Organization) -> Result<()> { + if let Some(ref db) = self.sea_conn { + let am = org_to_model(org); + am.update(db) + .await + .map_err(|e| Error::memory_with_source("SeaORM update organization", e))?; + return Ok(()); + } self.executor .execute( "UPDATE organizations SET name = ?, slug = ?, settings_json = ?, updated_at = ? WHERE id = ?", @@ -163,6 +163,20 @@ impl OrgRegistry for SqliteOrgEntityRepository { /// Deletes an organization. async fn delete_org(&self, id: &str) -> Result<()> { + if let Some(ref db) = self.sea_conn { + if let Some(active) = organization::Entity::find_by_id(id.to_string()) + .one(db) + .await + .map_err(|e| Error::memory_with_source("SeaORM find for delete", e))? + .map(organization::ActiveModel::from) + { + active + .delete(db) + .await + .map_err(|e| Error::memory_with_source("SeaORM delete organization", e))?; + } + return Ok(()); + } self.executor .execute( "DELETE FROM organizations WHERE id = ?", @@ -200,7 +214,7 @@ impl UserRegistry for SqliteOrgEntityRepository { &self.executor, "SELECT * FROM users WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_user, + User::from_row, ) .await?; Error::not_found_or(user, "User", id) @@ -215,7 +229,7 @@ impl UserRegistry for SqliteOrgEntityRepository { SqlParam::String(org_id.to_owned()), SqlParam::String(email.to_owned()), ], - row_to_user, + User::from_row, ) .await?; Error::not_found_or(user, "User", email) @@ -227,7 +241,7 @@ impl UserRegistry for SqliteOrgEntityRepository { &self.executor, "SELECT * FROM users WHERE org_id = ?", &[SqlParam::String(org_id.to_owned())], - row_to_user, + User::from_row, "org entity", ) .await @@ -286,7 +300,7 @@ impl TeamRegistry for SqliteOrgEntityRepository { &self.executor, "SELECT * FROM teams WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_team, + Team::from_row, ) .await?; Error::not_found_or(team, "Team", id) @@ -298,7 +312,7 @@ impl TeamRegistry for SqliteOrgEntityRepository { &self.executor, "SELECT * FROM teams WHERE org_id = ?", &[SqlParam::String(org_id.to_owned())], - row_to_team, + Team::from_row, "org entity", ) .await @@ -352,7 +366,7 @@ impl TeamMemberManager for SqliteOrgEntityRepository { &self.executor, "SELECT * FROM team_members WHERE team_id = ?", &[SqlParam::String(team_id.to_owned())], - row_to_team_member, + TeamMember::from_row, "org entity", ) .await @@ -388,7 +402,7 @@ impl ApiKeyRegistry for SqliteOrgEntityRepository { &self.executor, "SELECT * FROM api_keys WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_api_key, + ApiKey::from_row, ) .await?; Error::not_found_or(key, "ApiKey", id) @@ -400,7 +414,7 @@ impl ApiKeyRegistry for SqliteOrgEntityRepository { &self.executor, "SELECT * FROM api_keys WHERE org_id = ?", &[SqlParam::String(org_id.to_owned())], - row_to_api_key, + ApiKey::from_row, "org entity", ) .await diff --git a/crates/mcb-providers/src/database/sqlite/plan_entity_repository.rs b/crates/mcb-providers/src/database/sqlite/plan_entity_repository.rs index 2fbccaff6..74a2523fc 100644 --- a/crates/mcb-providers/src/database/sqlite/plan_entity_repository.rs +++ b/crates/mcb-providers/src/database/sqlite/plan_entity_repository.rs @@ -4,13 +4,13 @@ use std::sync::Arc; use async_trait::async_trait; -use mcb_domain::entities::plan::{Plan, PlanReview, PlanStatus, PlanVersion, ReviewVerdict}; +use mcb_domain::entities::plan::{Plan, PlanReview, PlanVersion}; use mcb_domain::error::{Error, Result}; -use mcb_domain::ports::{DatabaseExecutor, SqlParam, SqlRow}; +use mcb_domain::ports::{DatabaseExecutor, SqlParam}; use mcb_domain::ports::{PlanRegistry, PlanReviewRegistry, PlanVersionRegistry}; +use crate::database::sqlite::row_convert::FromRow; use crate::utils::sqlite::query as query_helpers; -use crate::utils::sqlite::row::{req_i64, req_parsed, req_str}; /// SQLite-backed repository for plan, version, and review entities. pub struct SqlitePlanEntityRepository { @@ -24,52 +24,6 @@ impl SqlitePlanEntityRepository { } } -/// Converts a SQL row to a Plan. -fn row_to_plan(row: &dyn SqlRow) -> Result { - let status: PlanStatus = req_parsed(row, "status")?; - - Ok(Plan { - id: req_str(row, "id")?, - org_id: req_str(row, "org_id")?, - project_id: req_str(row, "project_id")?, - title: req_str(row, "title")?, - description: req_str(row, "description")?, - status, - created_by: req_str(row, "created_by")?, - created_at: req_i64(row, "created_at")?, - updated_at: req_i64(row, "updated_at")?, - }) -} - -/// Converts a SQL row to a `PlanVersion`. -fn row_to_plan_version(row: &dyn SqlRow) -> Result { - Ok(PlanVersion { - id: req_str(row, "id")?, - org_id: req_str(row, "org_id")?, - plan_id: req_str(row, "plan_id")?, - version_number: req_i64(row, "version_number")?, - content_json: req_str(row, "content_json")?, - change_summary: req_str(row, "change_summary")?, - created_by: req_str(row, "created_by")?, - created_at: req_i64(row, "created_at")?, - }) -} - -/// Converts a SQL row to a `PlanReview`. -fn row_to_plan_review(row: &dyn SqlRow) -> Result { - let verdict: ReviewVerdict = req_parsed(row, "verdict")?; - - Ok(PlanReview { - id: req_str(row, "id")?, - org_id: req_str(row, "org_id")?, - plan_version_id: req_str(row, "plan_version_id")?, - reviewer_id: req_str(row, "reviewer_id")?, - verdict, - feedback: req_str(row, "feedback")?, - created_at: req_i64(row, "created_at")?, - }) -} - #[async_trait] /// Persistent plan registry using `SQLite`. impl PlanRegistry for SqlitePlanEntityRepository { @@ -102,7 +56,7 @@ impl PlanRegistry for SqlitePlanEntityRepository { SqlParam::String(org_id.to_owned()), SqlParam::String(id.to_owned()), ], - row_to_plan, + Plan::from_row, ) .await?; Error::not_found_or(plan, "Plan", id) @@ -117,7 +71,7 @@ impl PlanRegistry for SqlitePlanEntityRepository { SqlParam::String(org_id.to_owned()), SqlParam::String(project_id.to_owned()), ], - row_to_plan, + Plan::from_row, "plan entity", ) .await @@ -182,7 +136,7 @@ impl PlanVersionRegistry for SqlitePlanEntityRepository { &self.executor, "SELECT * FROM plan_versions WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_plan_version, + PlanVersion::from_row, ) .await?; Error::not_found_or(version, "PlanVersion", id) @@ -194,7 +148,7 @@ impl PlanVersionRegistry for SqlitePlanEntityRepository { &self.executor, "SELECT * FROM plan_versions WHERE plan_id = ?", &[SqlParam::String(plan_id.to_owned())], - row_to_plan_version, + PlanVersion::from_row, "plan entity", ) .await @@ -228,7 +182,7 @@ impl PlanReviewRegistry for SqlitePlanEntityRepository { &self.executor, "SELECT * FROM plan_reviews WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_plan_review, + PlanReview::from_row, ) .await?; Error::not_found_or(review, "PlanReview", id) @@ -240,7 +194,7 @@ impl PlanReviewRegistry for SqlitePlanEntityRepository { &self.executor, "SELECT * FROM plan_reviews WHERE plan_version_id = ?", &[SqlParam::String(plan_version_id.to_owned())], - row_to_plan_review, + PlanReview::from_row, "plan entity", ) .await diff --git a/crates/mcb-providers/src/database/sqlite/provider.rs b/crates/mcb-providers/src/database/sqlite/provider.rs index 9531b2c24..944ac1f00 100644 --- a/crates/mcb-providers/src/database/sqlite/provider.rs +++ b/crates/mcb-providers/src/database/sqlite/provider.rs @@ -31,13 +31,16 @@ use std::sync::Arc; use async_trait::async_trait; use mcb_domain::error::Result; use mcb_domain::ports::{ - AgentRepository, MemoryRepository, ProjectRepository, VcsEntityRepository, + AgentRepository, FileHashRepository, IssueEntityRepository, MemoryRepository, + OrgEntityRepository, PlanEntityRepository, ProjectRepository, VcsEntityRepository, }; use mcb_domain::ports::{DatabaseExecutor, DatabaseProvider}; use mcb_domain::schema::{Schema, SchemaDdlGenerator}; use super::{ - SqliteAgentRepository, SqliteExecutor, SqliteMemoryRepository, SqliteProjectRepository, + SqliteAgentRepository, SqliteBackend, SqliteExecutor, SqliteFileHashConfig, + SqliteFileHashRepository, SqliteIssueEntityRepository, SqliteMemoryRepository, + SqliteOrgEntityRepository, SqlitePlanEntityRepository, SqliteProjectRepository, SqliteSchemaDdlGenerator, SqliteVcsEntityRepository, }; use mcb_domain::registry::database::{ @@ -67,7 +70,80 @@ static SQLITE_DATABASE_PROVIDER: DatabaseProviderEntry = DatabaseProviderEntry { impl DatabaseProvider for SqliteDatabaseProvider { async fn connect(&self, path: &Path) -> Result> { let pool = connect_and_init(path.to_path_buf()).await?; - Ok(Arc::new(SqliteExecutor::new(pool))) + let db_url = format!("sqlite:{}?mode=rwc", path.display()); + let sea_conn = sea_orm::Database::connect(&db_url) + .await + .map_err(|e| mcb_domain::error::Error::memory_with_source("SeaORM connect", e))?; + let backend = SqliteBackend::new(SqliteExecutor::new(pool), sea_conn); + Ok(Arc::new(backend)) + } + + fn create_memory_repository( + &self, + executor: Arc, + ) -> Arc { + Arc::new(SqliteMemoryRepository::new(executor)) + } + + fn create_agent_repository( + &self, + executor: Arc, + ) -> Arc { + Arc::new(SqliteAgentRepository::new(executor)) + } + + fn create_project_repository( + &self, + executor: Arc, + ) -> Arc { + Arc::new(SqliteProjectRepository::new(executor)) + } + + fn create_file_hash_repository( + &self, + executor: Arc, + project_id: String, + ) -> Arc { + Arc::new(SqliteFileHashRepository::new( + executor, + SqliteFileHashConfig::default(), + project_id, + )) + } + + fn create_vcs_entity_repository( + &self, + executor: Arc, + ) -> Arc { + Arc::new(SqliteVcsEntityRepository::new(executor)) + } + + fn create_plan_entity_repository( + &self, + executor: Arc, + ) -> Arc { + Arc::new(SqlitePlanEntityRepository::new(executor)) + } + + fn create_issue_entity_repository( + &self, + executor: Arc, + ) -> Arc { + Arc::new(SqliteIssueEntityRepository::new(executor)) + } + + fn create_org_entity_repository( + &self, + executor: Arc, + ) -> Arc { + let sea_conn = executor + .as_any() + .downcast_ref::() + .map(|b| b.sea_conn().clone()); + match sea_conn { + Some(conn) => Arc::new(SqliteOrgEntityRepository::new_with_sea(executor, conn)), + None => Arc::new(SqliteOrgEntityRepository::new(executor)), + } } } diff --git a/crates/mcb-providers/src/database/sqlite/row_convert.rs b/crates/mcb-providers/src/database/sqlite/row_convert.rs index 67e91efe8..d9019ae39 100644 --- a/crates/mcb-providers/src/database/sqlite/row_convert.rs +++ b/crates/mcb-providers/src/database/sqlite/row_convert.rs @@ -3,15 +3,23 @@ //! //! Row-to-entity conversion using the domain port [`SqlRow`]. -use crate::utils::sqlite::row::{json_opt, json_vec, req_i64, req_parsed, req_str}; +use crate::utils::sqlite::row::{ + json_opt, json_vec, opt_i64, opt_str, req_i64, req_parsed, req_str, +}; use mcb_domain::constants::keys as schema; use mcb_domain::entities::agent::{AgentSession, Checkpoint, CheckpointType}; use mcb_domain::entities::issue::{IssueComment, IssueLabel}; use mcb_domain::entities::memory::{Observation, ObservationMetadata, SessionSummary}; +use mcb_domain::entities::plan::{Plan, PlanReview, PlanStatus, PlanVersion, ReviewVerdict}; use mcb_domain::entities::project::{Project, ProjectIssue}; +use mcb_domain::entities::repository::{Branch, Repository, VcsType}; +use mcb_domain::entities::worktree::{AgentWorktreeAssignment, Worktree, WorktreeStatus}; +use mcb_domain::entities::{ApiKey, Organization, Team, TeamMember, User}; use mcb_domain::error::{Error, Result}; use mcb_domain::ports::SqlRow; use mcb_domain::schema::COL_OBSERVATION_TYPE; +use mcb_domain::utils::id; +use mcb_domain::value_objects::ids::TeamMemberId; /// Trait for converting a database row to an entity type. #[allow(dead_code)] @@ -233,3 +241,169 @@ impl FromRow for IssueLabel { row_to_label(row) } } + +from_row_simple!(Organization, { + id: req_str, + name: req_str, + slug: req_str, + settings_json: req_str, + created_at: req_i64, + updated_at: req_i64, +}); + +from_row_simple!(User, { + id: req_str, + org_id: req_str, + email: req_str, + display_name: req_str, + role: req_parsed, + api_key_hash: opt_str, + created_at: req_i64, + updated_at: req_i64, +}); + +from_row_simple!(Team, { + id: req_str, + org_id: req_str, + name: req_str, + created_at: req_i64, +}); + +impl FromRow for TeamMember { + fn from_row(row: &dyn SqlRow) -> Result { + let team_id = req_str(row, "team_id")?; + let user_id = req_str(row, "user_id")?; + let id_uuid = id::deterministic("team_member", &format!("{team_id}:{user_id}")); + Ok(TeamMember { + id: TeamMemberId::from_uuid(id_uuid), + team_id, + user_id, + role: req_parsed(row, "role")?, + joined_at: req_i64(row, "joined_at")?, + }) + } +} + +impl FromRow for ApiKey { + fn from_row(row: &dyn SqlRow) -> Result { + Ok(ApiKey { + id: req_str(row, "id")?, + user_id: req_str(row, "user_id")?, + org_id: req_str(row, "org_id")?, + key_hash: req_str(row, "key_hash")?, + name: req_str(row, "name")?, + scopes_json: req_str(row, "scopes_json")?, + expires_at: opt_i64(row, "expires_at")?, + created_at: req_i64(row, "created_at")?, + revoked_at: opt_i64(row, "revoked_at")?, + }) + } +} + +impl FromRow for Plan { + fn from_row(row: &dyn SqlRow) -> Result { + let status: PlanStatus = req_parsed(row, "status")?; + Ok(Plan { + id: req_str(row, "id")?, + org_id: req_str(row, "org_id")?, + project_id: req_str(row, "project_id")?, + title: req_str(row, "title")?, + description: req_str(row, "description")?, + status, + created_by: req_str(row, "created_by")?, + created_at: req_i64(row, "created_at")?, + updated_at: req_i64(row, "updated_at")?, + }) + } +} + +impl FromRow for PlanVersion { + fn from_row(row: &dyn SqlRow) -> Result { + Ok(PlanVersion { + id: req_str(row, "id")?, + org_id: req_str(row, "org_id")?, + plan_id: req_str(row, "plan_id")?, + version_number: req_i64(row, "version_number")?, + content_json: req_str(row, "content_json")?, + change_summary: req_str(row, "change_summary")?, + created_by: req_str(row, "created_by")?, + created_at: req_i64(row, "created_at")?, + }) + } +} + +impl FromRow for PlanReview { + fn from_row(row: &dyn SqlRow) -> Result { + let verdict: ReviewVerdict = req_parsed(row, "verdict")?; + Ok(PlanReview { + id: req_str(row, "id")?, + org_id: req_str(row, "org_id")?, + plan_version_id: req_str(row, "plan_version_id")?, + reviewer_id: req_str(row, "reviewer_id")?, + verdict, + feedback: req_str(row, "feedback")?, + created_at: req_i64(row, "created_at")?, + }) + } +} + +impl FromRow for Repository { + fn from_row(row: &dyn SqlRow) -> Result { + let vcs_type: VcsType = req_parsed(row, "vcs_type")?; + Ok(Repository { + id: req_str(row, "id")?, + org_id: req_str(row, "org_id")?, + project_id: req_str(row, "project_id")?, + name: req_str(row, "name")?, + url: req_str(row, "url")?, + local_path: req_str(row, "local_path")?, + vcs_type, + created_at: req_i64(row, "created_at")?, + updated_at: req_i64(row, "updated_at")?, + }) + } +} + +impl FromRow for Branch { + fn from_row(row: &dyn SqlRow) -> Result { + let is_default_i = req_i64(row, "is_default")?; + Ok(Branch { + id: req_str(row, "id")?, + org_id: req_str(row, "org_id")?, + repository_id: req_str(row, "repository_id")?, + name: req_str(row, "name")?, + is_default: is_default_i != 0, + head_commit: req_str(row, "head_commit")?, + upstream: row.try_get_string("upstream")?, + created_at: req_i64(row, "created_at")?, + }) + } +} + +impl FromRow for Worktree { + fn from_row(row: &dyn SqlRow) -> Result { + let status: WorktreeStatus = req_parsed(row, "status")?; + Ok(Worktree { + id: req_str(row, "id")?, + repository_id: req_str(row, "repository_id")?, + branch_id: req_str(row, "branch_id")?, + path: req_str(row, "path")?, + status, + assigned_agent_id: row.try_get_string("assigned_agent_id")?, + created_at: req_i64(row, "created_at")?, + updated_at: req_i64(row, "updated_at")?, + }) + } +} + +impl FromRow for AgentWorktreeAssignment { + fn from_row(row: &dyn SqlRow) -> Result { + Ok(AgentWorktreeAssignment { + id: req_str(row, "id")?, + agent_session_id: req_str(row, "agent_session_id")?, + worktree_id: req_str(row, "worktree_id")?, + assigned_at: req_i64(row, "assigned_at")?, + released_at: row.try_get_i64("released_at")?, + }) + } +} diff --git a/crates/mcb-providers/src/database/sqlite/sea_entities/mod.rs b/crates/mcb-providers/src/database/sqlite/sea_entities/mod.rs new file mode 100644 index 000000000..14bb123aa --- /dev/null +++ b/crates/mcb-providers/src/database/sqlite/sea_entities/mod.rs @@ -0,0 +1,55 @@ +//! +//! SeaORM entities for SQLite tables, generated from domain schema. +//! +//! Column lists stay in sync with `mcb_domain::schema::*::table()`; tests enforce it. + +#[allow(dead_code)] +pub mod organization; + +#[cfg(test)] +mod schema_sync_tests { + use mcb_domain::schema::Schema; + use mcb_domain::schema::types::ColumnType; + + fn column_type_to_rust(t: &ColumnType) -> &'static str { + match t { + ColumnType::Text => "String", + ColumnType::Integer => "i64", + ColumnType::Real => "f64", + ColumnType::Boolean => "bool", + ColumnType::Blob => "Vec", + ColumnType::Json => "String", + ColumnType::Uuid => "String", + ColumnType::Timestamp => "i64", + } + } + + #[test] + fn organizations_entity_matches_domain_schema() { + let schema = Schema::definition(); + let table = schema + .tables + .iter() + .find(|t| t.name == "organizations") + .expect("organizations table in schema"); + let expected: Vec<(&str, &str)> = table + .columns + .iter() + .map(|c| (c.name.as_str(), column_type_to_rust(&c.type_))) + .collect(); + let actual = super::organization::SCHEMA_COLUMNS; + assert_eq!( + expected.len(), + actual.len(), + "column count must match schema" + ); + for (i, (name, ty)) in expected.iter().enumerate() { + assert_eq!( + (*name, *ty), + (actual[i].0, actual[i].1), + "column {} must match schema", + i + ); + } + } +} diff --git a/crates/mcb-providers/src/database/sqlite/sea_entities/organization.rs b/crates/mcb-providers/src/database/sqlite/sea_entities/organization.rs new file mode 100644 index 000000000..c10023b72 --- /dev/null +++ b/crates/mcb-providers/src/database/sqlite/sea_entities/organization.rs @@ -0,0 +1,17 @@ +//! +//! SeaORM entity for `organizations`. Generated from schema; keep in sync with +//! `mcb_domain::schema::organizations::table()`. + +use sea_orm::entity::prelude::*; + +crate::sea_entity!( + "organizations", + [ + (id: String), + (name: String), + (slug: String), + (settings_json: String), + (created_at: i64), + (updated_at: i64), + ] +); diff --git a/crates/mcb-providers/src/database/sqlite/vcs_entity_repository.rs b/crates/mcb-providers/src/database/sqlite/vcs_entity_repository.rs index aae296c5d..6f8b662f2 100644 --- a/crates/mcb-providers/src/database/sqlite/vcs_entity_repository.rs +++ b/crates/mcb-providers/src/database/sqlite/vcs_entity_repository.rs @@ -6,15 +6,15 @@ use std::sync::Arc; use async_trait::async_trait; -use mcb_domain::entities::repository::{Branch, Repository, VcsType}; -use mcb_domain::entities::worktree::{AgentWorktreeAssignment, Worktree, WorktreeStatus}; +use mcb_domain::entities::repository::{Branch, Repository}; +use mcb_domain::entities::worktree::{AgentWorktreeAssignment, Worktree}; use mcb_domain::error::{Error, Result}; use mcb_domain::ports::VcsEntityRepository; -use mcb_domain::ports::{DatabaseExecutor, SqlParam, SqlRow}; +use mcb_domain::ports::{DatabaseExecutor, SqlParam}; use serde_json::json; +use crate::database::sqlite::row_convert::FromRow; use crate::utils::sqlite::query as query_helpers; -use crate::utils::sqlite::row::{req_i64, req_parsed, req_str}; /// SQLite-backed repository for VCS repositories, branches, worktrees, and assignments. pub struct SqliteVcsEntityRepository { @@ -28,65 +28,6 @@ impl SqliteVcsEntityRepository { } } -/// Converts a SQL row to a Repository. -fn row_to_repository(row: &dyn SqlRow) -> Result { - let vcs_type: VcsType = req_parsed(row, "vcs_type")?; - - Ok(Repository { - id: req_str(row, "id")?, - org_id: req_str(row, "org_id")?, - project_id: req_str(row, "project_id")?, - name: req_str(row, "name")?, - url: req_str(row, "url")?, - local_path: req_str(row, "local_path")?, - vcs_type, - created_at: req_i64(row, "created_at")?, - updated_at: req_i64(row, "updated_at")?, - }) -} - -/// Converts a SQL row to a Branch. -fn row_to_branch(row: &dyn SqlRow) -> Result { - let is_default_i = req_i64(row, "is_default")?; - Ok(Branch { - id: req_str(row, "id")?, - org_id: req_str(row, "org_id")?, - repository_id: req_str(row, "repository_id")?, - name: req_str(row, "name")?, - is_default: is_default_i != 0, - head_commit: req_str(row, "head_commit")?, - upstream: row.try_get_string("upstream")?, - created_at: req_i64(row, "created_at")?, - }) -} - -/// Converts a SQL row to a Worktree. -fn row_to_worktree(row: &dyn SqlRow) -> Result { - let status: WorktreeStatus = req_parsed(row, "status")?; - - Ok(Worktree { - id: req_str(row, "id")?, - repository_id: req_str(row, "repository_id")?, - branch_id: req_str(row, "branch_id")?, - path: req_str(row, "path")?, - status, - assigned_agent_id: row.try_get_string("assigned_agent_id")?, - created_at: req_i64(row, "created_at")?, - updated_at: req_i64(row, "updated_at")?, - }) -} - -/// Converts a SQL row to an `AgentWorktreeAssignment`. -fn row_to_assignment(row: &dyn SqlRow) -> Result { - Ok(AgentWorktreeAssignment { - id: req_str(row, "id")?, - agent_session_id: req_str(row, "agent_session_id")?, - worktree_id: req_str(row, "worktree_id")?, - assigned_at: req_i64(row, "assigned_at")?, - released_at: row.try_get_i64("released_at")?, - }) -} - #[async_trait] /// Repository for VCS entities. impl VcsEntityRepository for SqliteVcsEntityRepository { @@ -131,7 +72,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { SqlParam::String(org_id.to_owned()), SqlParam::String(id.to_owned()), ], - row_to_repository, + Repository::from_row, ) .await?; Error::not_found_or(repo, "Repository", id) @@ -146,7 +87,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { SqlParam::String(org_id.to_owned()), SqlParam::String(project_id.to_owned()), ], - row_to_repository, + Repository::from_row, "vcs entity", ) .await @@ -229,7 +170,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { &self.executor, "SELECT * FROM branches WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_branch, + Branch::from_row, ) .await?; Error::not_found_or(branch, "Branch", id) @@ -241,7 +182,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { &self.executor, "SELECT * FROM branches WHERE repository_id = ?", &[SqlParam::String(repository_id.to_owned())], - row_to_branch, + Branch::from_row, "vcs entity", ) .await @@ -321,7 +262,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { &self.executor, "SELECT * FROM worktrees WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_worktree, + Worktree::from_row, ) .await?; Error::not_found_or(wt, "Worktree", id) @@ -333,7 +274,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { &self.executor, "SELECT * FROM worktrees WHERE repository_id = ?", &[SqlParam::String(repository_id.to_owned())], - row_to_worktree, + Worktree::from_row, "vcs entity", ) .await @@ -410,7 +351,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { &self.executor, "SELECT * FROM agent_worktree_assignments WHERE id = ?", &[SqlParam::String(id.to_owned())], - row_to_assignment, + AgentWorktreeAssignment::from_row, ) .await?; Error::not_found_or(asgn, "Assignment", id) @@ -425,7 +366,7 @@ impl VcsEntityRepository for SqliteVcsEntityRepository { &self.executor, "SELECT * FROM agent_worktree_assignments WHERE worktree_id = ?", &[SqlParam::String(worktree_id.to_owned())], - row_to_assignment, + AgentWorktreeAssignment::from_row, "vcs entity", ) .await diff --git a/crates/mcb-providers/src/events/nats.rs b/crates/mcb-providers/src/events/nats.rs index 025011b56..210525401 100644 --- a/crates/mcb-providers/src/events/nats.rs +++ b/crates/mcb-providers/src/events/nats.rs @@ -37,12 +37,46 @@ use futures::{StreamExt, stream}; use mcb_domain::error::{Error, Result}; use mcb_domain::events::DomainEvent; use mcb_domain::ports::{DomainEventStream, EventBusProvider}; +use mcb_domain::registry::event_bus::{ + EVENT_BUS_PROVIDERS, EventBusProviderConfig, EventBusProviderEntry, +}; use mcb_domain::utils::id; use tokio::sync::RwLock; use tracing::{debug, info, warn}; use crate::constants::NATS_DEFAULT_SUBJECT; +fn create_nats_event_bus_provider( + config: &EventBusProviderConfig, +) -> std::result::Result, String> { + let url = config + .extra + .get("url") + .cloned() + .unwrap_or_else(|| "nats://127.0.0.1:4222".to_owned()); + let subject = config + .extra + .get("subject") + .cloned() + .unwrap_or_else(|| NATS_DEFAULT_SUBJECT.to_owned()); + let client_name = config.extra.get("client_name").cloned(); + + let rt = tokio::runtime::Handle::try_current().map_err(|e| e.to_string())?; + rt.block_on(async move { + NatsEventBusProvider::with_options(&url, &subject, client_name.as_deref()) + .await + .map(|provider| Arc::new(provider) as Arc) + .map_err(|e| e.to_string()) + }) +} + +#[linkme::distributed_slice(EVENT_BUS_PROVIDERS)] +static NATS_EVENT_BUS_PROVIDER: EventBusProviderEntry = EventBusProviderEntry { + name: "nats", + description: "NATS distributed event bus", + build: create_nats_event_bus_provider, +}; + /// Event bus provider using NATS for distributed systems /// /// Provides distributed event broadcasting across multiple processes/nodes. diff --git a/crates/mcb-providers/src/events/tokio.rs b/crates/mcb-providers/src/events/tokio.rs index c693893f2..acfba0da9 100644 --- a/crates/mcb-providers/src/events/tokio.rs +++ b/crates/mcb-providers/src/events/tokio.rs @@ -37,12 +37,34 @@ use futures::stream; use mcb_domain::error::Result; use mcb_domain::events::DomainEvent; use mcb_domain::ports::{DomainEventStream, EventBusProvider}; +use mcb_domain::registry::event_bus::{ + EVENT_BUS_PROVIDERS, EventBusProviderConfig, EventBusProviderEntry, +}; use mcb_domain::utils::id; use tokio::sync::broadcast; use tracing::{debug, warn}; use crate::constants::EVENTS_TOKIO_DEFAULT_CAPACITY; +fn create_tokio_event_bus_provider( + config: &EventBusProviderConfig, +) -> std::result::Result, String> { + let capacity = config + .extra + .get("capacity") + .and_then(|v| v.parse::().ok()) + .unwrap_or(EVENTS_TOKIO_DEFAULT_CAPACITY); + + Ok(Arc::new(TokioEventBusProvider::with_capacity(capacity))) +} + +#[linkme::distributed_slice(EVENT_BUS_PROVIDERS)] +static TOKIO_EVENT_BUS_PROVIDER: EventBusProviderEntry = EventBusProviderEntry { + name: "tokio", + description: "Tokio broadcast event bus", + build: create_tokio_event_bus_provider, +}; + /// Event bus provider using tokio broadcast channels /// /// Provides in-process event distribution with multiple subscribers. diff --git a/crates/mcb-providers/src/fs/mod.rs b/crates/mcb-providers/src/fs/mod.rs new file mode 100644 index 000000000..653a69bba --- /dev/null +++ b/crates/mcb-providers/src/fs/mod.rs @@ -0,0 +1,73 @@ +use std::path::{Path, PathBuf}; + +use async_trait::async_trait; +use mcb_domain::error::Error; +use mcb_domain::ports::{DirEntry, FileSystemProvider}; +use mcb_domain::registry::fs::{ + FILE_SYSTEM_PROVIDERS, FileSystemProviderConfig, FileSystemProviderEntry, +}; + +fn create_local_file_system_provider( + _config: &FileSystemProviderConfig, +) -> std::result::Result, String> { + Ok(std::sync::Arc::new(LocalFileSystemProvider::new())) +} + +#[linkme::distributed_slice(FILE_SYSTEM_PROVIDERS)] +static LOCAL_FILE_SYSTEM_PROVIDER: FileSystemProviderEntry = FileSystemProviderEntry { + name: "local", + description: "Local file-system provider", + build: create_local_file_system_provider, +}; + +#[allow(missing_docs)] +#[derive(Debug, Default)] +pub struct LocalFileSystemProvider; + +#[allow(missing_docs)] +impl LocalFileSystemProvider { + #[must_use] + pub fn new() -> Self { + Self + } + + fn read_dir_sync(path: &Path) -> Result, Error> { + let entries = std::fs::read_dir(path).map_err(|e| { + Error::internal(format!("Failed to read directory {}: {e}", path.display())) + })?; + + let mut out = Vec::new(); + for entry in entries { + let entry = entry + .map_err(|e| Error::internal(format!("Failed to read directory entry: {e}")))?; + let file_type = entry + .file_type() + .map_err(|e| Error::internal(format!("Failed to read file type: {e}")))?; + + out.push(DirEntry { + path: entry.path(), + is_file: file_type.is_file(), + is_dir: file_type.is_dir(), + }); + } + + Ok(out) + } +} + +#[async_trait] +impl FileSystemProvider for LocalFileSystemProvider { + async fn read_to_string(&self, path: &Path) -> Result { + std::fs::read_to_string(path) + .map_err(|e| Error::internal(format!("Failed to read file {}: {e}", path.display()))) + } + + async fn read_dir_entries(&self, path: &Path) -> Result, Error> { + Self::read_dir_sync(path) + } + + async fn canonicalize_path(&self, path: &Path) -> Result { + std::fs::canonicalize(path) + .map_err(|e| Error::internal(format!("Failed to canonicalize {}: {e}", path.display()))) + } +} diff --git a/crates/mcb-providers/src/lib.rs b/crates/mcb-providers/src/lib.rs index 8668ffab1..6ed240a87 100644 --- a/crates/mcb-providers/src/lib.rs +++ b/crates/mcb-providers/src/lib.rs @@ -88,6 +88,12 @@ pub mod language; /// Provides BM25 text ranking algorithm and hybrid score fusion. pub mod hybrid_search; +#[allow(missing_docs)] +pub mod fs; + +#[allow(missing_docs)] +pub mod task; + // Re-export hybrid search providers pub use hybrid_search::HybridSearchEngine; diff --git a/crates/mcb-providers/src/task/mod.rs b/crates/mcb-providers/src/task/mod.rs new file mode 100644 index 000000000..3f4b74f31 --- /dev/null +++ b/crates/mcb-providers/src/task/mod.rs @@ -0,0 +1,38 @@ +use futures::future::BoxFuture; +use mcb_domain::error::Result; +use mcb_domain::ports::TaskRunnerProvider; +use mcb_domain::registry::task_runner::{ + TASK_RUNNER_PROVIDERS, TaskRunnerProviderConfig, TaskRunnerProviderEntry, +}; + +fn create_tokio_task_runner_provider( + _config: &TaskRunnerProviderConfig, +) -> std::result::Result, String> { + Ok(std::sync::Arc::new(TokioTaskRunnerProvider::new())) +} + +#[linkme::distributed_slice(TASK_RUNNER_PROVIDERS)] +static TOKIO_TASK_RUNNER_PROVIDER: TaskRunnerProviderEntry = TaskRunnerProviderEntry { + name: "tokio", + description: "Tokio task runner provider", + build: create_tokio_task_runner_provider, +}; + +#[allow(missing_docs)] +#[derive(Debug, Default)] +pub struct TokioTaskRunnerProvider; + +#[allow(missing_docs)] +impl TokioTaskRunnerProvider { + #[must_use] + pub fn new() -> Self { + Self + } +} + +impl TaskRunnerProvider for TokioTaskRunnerProvider { + fn spawn(&self, task: BoxFuture<'static, ()>) -> Result<()> { + tokio::spawn(task); + Ok(()) + } +} diff --git a/crates/mcb-providers/src/vcs/git2_provider.rs b/crates/mcb-providers/src/vcs/git2_provider.rs index d4a042046..fc0781b0f 100644 --- a/crates/mcb-providers/src/vcs/git2_provider.rs +++ b/crates/mcb-providers/src/vcs/git2_provider.rs @@ -7,6 +7,7 @@ use std::path::{Path, PathBuf}; use async_trait::async_trait; use git2::{BranchType, Repository, Sort}; +use mcb_domain::registry::vcs::{VCS_PROVIDERS, VcsProviderConfig, VcsProviderEntry}; use mcb_domain::utils::id; use mcb_domain::{ entities::vcs::{ @@ -16,6 +17,20 @@ use mcb_domain::{ error::{Error, Result}, ports::VcsProvider, }; +use std::sync::Arc; + +fn create_git2_provider( + _config: &VcsProviderConfig, +) -> std::result::Result, String> { + Ok(Arc::new(Git2Provider::new())) +} + +#[linkme::distributed_slice(VCS_PROVIDERS)] +static GIT2_VCS_PROVIDER: VcsProviderEntry = VcsProviderEntry { + name: "git2", + description: "libgit2-based VCS provider", + build: create_git2_provider, +}; /// Git implementation of `VcsProvider` using libgit2. /// diff --git a/crates/mcb-providers/src/vector_store/milvus.rs b/crates/mcb-providers/src/vector_store/milvus.rs index 3be755166..7d027bcbc 100644 --- a/crates/mcb-providers/src/vector_store/milvus.rs +++ b/crates/mcb-providers/src/vector_store/milvus.rs @@ -424,15 +424,12 @@ impl MilvusVectorStoreProvider { ] } - #[allow(clippy::str_to_string)] // False positive: iter yields &i64, not &str fn parse_milvus_ids(result: &milvus::proto::milvus::MutationResult) -> Vec { match &result.i_ds { Some(ids) => match &ids.id_field { - Some(milvus::proto::schema::i_ds::IdField::IntId(int_ids)) => int_ids - .data - .iter() - .map(std::string::ToString::to_string) - .collect(), + Some(milvus::proto::schema::i_ds::IdField::IntId(int_ids)) => { + int_ids.data.iter().map(|id: &i64| id.to_string()).collect() + } Some(milvus::proto::schema::i_ds::IdField::StrId(str_ids)) => str_ids.data.clone(), None => Vec::new(), }, diff --git a/crates/mcb/src/cli/validate.rs b/crates/mcb/src/cli/validate.rs index e6bb39f0c..ac49edd15 100644 --- a/crates/mcb/src/cli/validate.rs +++ b/crates/mcb/src/cli/validate.rs @@ -1,6 +1,6 @@ //! Validate command - runs architecture validation -#![allow(clippy::print_stdout)] +use std::io::Write; use std::path::PathBuf; use clap::Args; @@ -112,7 +112,7 @@ impl ValidateArgs { /// Print report as JSON fn print_json(report: &mcb_validate::GenericReport) -> Result<(), Box> { let json = serde_json::to_string_pretty(report)?; - println!("{json}"); + writeln!(std::io::stdout(), "{json}")?; Ok(()) } @@ -150,7 +150,7 @@ impl ValidateArgs { } if has_violations { - println!(); + let _ = writeln!(std::io::stdout()); } } @@ -167,26 +167,34 @@ impl ValidateArgs { let file_display = violation.file.as_deref().unwrap_or("-"); let line = violation.line.unwrap_or(0); - println!( + let _ = writeln!( + std::io::stdout(), "[{}] {}: {} ({}:{})", - violation.severity, violation.id, violation.message, file_display, line + violation.severity, + violation.id, + violation.message, + file_display, + line ); if let Some(ref suggestion) = violation.suggestion { - println!(" → {suggestion}"); + let _ = writeln!(std::io::stdout(), " → {suggestion}"); } } fn print_summary(&self, report: &mcb_validate::GenericReport) { - println!( + let _ = writeln!( + std::io::stdout(), "Validation complete: {} error(s), {} warning(s), {} info(s)", - report.summary.errors, report.summary.warnings, report.summary.infos + report.summary.errors, + report.summary.warnings, + report.summary.infos ); // Print category breakdown (unless quick mode) if !self.quick && !report.summary.by_category.is_empty() { - println!("\nBy category:"); + let _ = writeln!(std::io::stdout(), "\nBy category:"); for (category, count) in &report.summary.by_category { - println!(" {category}: {count}"); + let _ = writeln!(std::io::stdout(), " {category}: {count}"); } } }