use crate::backends::{
build_embedder, build_migration_target, build_remote_extractor,
warn_if_extraction_backend_is_unreachable, ConfiguredEmbedder,
};
use std::time::Duration;
use velesdb_memory::mcp::McpServer;
use velesdb_memory::{DynEmbedder, ExtractorSelection, MemoryService, NativeStore};
pub(crate) fn build_configured_server(
configured: ConfiguredService,
) -> Result<McpServer, Box<dyn std::error::Error>> {
let ConfiguredService {
service,
store_path,
embedder_model,
embedder_dimension,
} = configured;
let store_path = std::path::PathBuf::from(store_path);
apply_ingest_roots(apply_default_ttl(
build_server(service)?
.with_embedder_identity(embedder_model, embedder_dimension)
.with_store_dir(&store_path)
.with_online_migration(&store_path, build_migration_target)?
.with_extraction_jobs(&store_path)
.map_err(std::io::Error::other)?,
)?)
}
pub(crate) const LOCK_RETRY_ATTEMPTS: u32 = 3;
pub(crate) const LOCK_RETRY_DELAY: Duration = Duration::from_millis(500);
pub(crate) fn open_store_with_actionable_lock_error(
store_path: &str,
configured: ConfiguredEmbedder,
) -> Result<MemoryService<DynEmbedder>, Box<dyn std::error::Error>> {
let ConfiguredEmbedder { embedder, model } = configured;
let dimension = embedder.dimension();
let store_dir = std::path::Path::new(store_path);
let recorded = velesdb_memory::embedding_provenance::read(store_dir)?;
velesdb_memory::embedding_provenance::check(recorded.as_ref(), &model, dimension)?;
let unrecorded = recorded.is_none();
match open_through_lock_retries(store_path, dimension) {
Ok(store) => {
if unrecorded {
record_embedding_model(store_dir, &store, &model, dimension);
}
Ok(MemoryService::with_store(store, embedder))
}
Err(other) if unrecorded => Err(format!(
"{other}\n{}",
velesdb_memory::embedding_provenance::unrecorded_model_note(&model)
)
.into()),
Err(other) => Err(other.into()),
}
}
pub(crate) fn open_through_lock_retries(
store_path: &str,
dimension: usize,
) -> Result<NativeStore, velesdb_memory::MemoryError> {
use velesdb_memory::MemoryError;
let mut last_locked_path: Option<String> = None;
for attempt in 0..LOCK_RETRY_ATTEMPTS {
match NativeStore::open(store_path, dimension) {
Err(MemoryError::Storage(velesdb_core::Error::DatabaseLocked(locked_path))) => {
last_locked_path = Some(locked_path);
if attempt + 1 < LOCK_RETRY_ATTEMPTS {
std::thread::sleep(LOCK_RETRY_DELAY);
}
}
outcome => return outcome,
}
}
let locked_path = last_locked_path.unwrap_or_else(|| store_path.to_owned());
eprintln!(
"[velesdb-memory] another velesdb-memory process holds {locked_path} — \
kill it (pkill velesdb-memory) or point VELESDB_MEMORY_PATH elsewhere"
);
std::process::exit(1);
}
pub(crate) fn record_embedding_model(
store_dir: &std::path::Path,
store: &NativeStore,
model: &str,
dimension: usize,
) {
use velesdb_memory::embedding_provenance::{write, EmbeddingProvenance};
use velesdb_memory::FactStore as _;
if store.count() != 0 {
return;
}
if let Err(err) = write(store_dir, &EmbeddingProvenance::new(model, dimension)) {
if std::env::var_os("VELESDB_MEMORY_QUIET").is_none() {
eprintln!(
"[velesdb-memory] could not record the embedding model ({err}) — the store works, \
but a later model change will only be checked against the vector dimension"
);
}
}
}
pub(crate) fn build_configured_service(
args: &[String],
) -> Result<ConfiguredService, Box<dyn std::error::Error>> {
apply_config_file(args)?;
let store_path = std::env::var("VELESDB_MEMORY_PATH").unwrap_or_else(|_| default_store_path());
let recovery = velesdb_memory::migration::recover_online_migration_startup(
std::path::Path::new(&store_path),
build_migration_target,
)?;
let configured = embedder_after_startup_recovery(recovery)?;
let embedder_model = configured.model.clone();
let embedder_dimension = configured.embedder.dimension();
let service = apply_autograph(open_store_with_actionable_lock_error(
&store_path,
configured,
)?)?;
Ok(ConfiguredService {
service,
store_path,
embedder_model,
embedder_dimension,
})
}
pub(crate) fn embedder_after_startup_recovery(
recovery: velesdb_memory::migration::OnlineMigrationStartup,
) -> Result<ConfiguredEmbedder, Box<dyn std::error::Error>> {
match recovery {
velesdb_memory::migration::OnlineMigrationStartup::None => build_embedder(),
velesdb_memory::migration::OnlineMigrationStartup::SourceRestored { source_model } => {
let configured = build_embedder()?;
if configured.model != source_model {
return Err(format!(
"online migration restored source model '{source_model}', but startup configured '{}'",
configured.model
)
.into());
}
Ok(configured)
}
velesdb_memory::migration::OnlineMigrationStartup::TargetActivated { embedder, model } => {
Ok(ConfiguredEmbedder { embedder, model })
}
}
}
pub(crate) struct ConfiguredService {
service: MemoryService<DynEmbedder>,
store_path: String,
embedder_model: String,
embedder_dimension: usize,
}
pub(crate) fn apply_autograph(
service: MemoryService<DynEmbedder>,
) -> Result<MemoryService<DynEmbedder>, Box<dyn std::error::Error>> {
if std::env::var("VELESDB_MEMORY_AUTOGRAPH").as_deref() != Ok("1") {
return Ok(service);
}
let backend = std::env::var("VELESDB_MEMORY_EXTRACTOR").unwrap_or_default();
match velesdb_memory::select_extractor(&backend)? {
ExtractorSelection::Disabled => Err(
"autograph is on ([graph] autograph = true / VELESDB_MEMORY_AUTOGRAPH=1) but no \
extraction backend is configured — set [extractor] backend = \"outline\" for the \
offline deterministic reader (no rebuild, no model to run), or \"ollama\" with a \
model, or turn autograph off"
.into(),
),
ExtractorSelection::Ready(extractor) => Ok(service.with_autograph(extractor)),
ExtractorSelection::NeedsRemoteConfig(backend) => {
let extractor = build_remote_extractor(backend)?;
warn_if_extraction_backend_is_unreachable(backend);
Ok(service.with_autograph(extractor))
}
}
}
pub(crate) fn apply_config_file(args: &[String]) -> Result<(), Box<dyn std::error::Error>> {
let explicit = args
.iter()
.position(|arg| arg == "--config")
.and_then(|at| args.get(at + 1))
.map(String::as_str);
let store_dir = std::env::var("VELESDB_MEMORY_PATH").unwrap_or_else(|_| default_store_path());
let Some(path) =
velesdb_memory::config::resolve_path(explicit, Some(std::path::Path::new(&store_dir)))
else {
return Ok(());
};
let loaded = velesdb_memory::config::load(&path)?;
let applied = velesdb_memory::config::apply(&loaded.values);
if !applied.is_empty() && std::env::var_os("VELESDB_MEMORY_QUIET").is_none() {
eprintln!(
"velesdb-memory: {} setting(s) from {}",
applied.len(),
path.display()
);
}
Ok(())
}
pub(crate) fn default_store_path() -> String {
let home = std::env::var_os("HOME")
.or_else(|| std::env::var_os("USERPROFILE"))
.filter(|h| !h.is_empty());
match home {
Some(home) => std::path::Path::new(&home)
.join(".velesdb-memory")
.to_string_lossy()
.into_owned(),
None => "./velesdb-memory-store".to_owned(),
}
}
pub(crate) fn apply_default_ttl(
server: McpServer,
) -> Result<McpServer, Box<dyn std::error::Error>> {
match std::env::var("VELESDB_MEMORY_DEFAULT_TTL") {
Ok(raw) => {
let ttl_seconds: u64 = raw.trim().parse().map_err(|_| {
format!(
"VELESDB_MEMORY_DEFAULT_TTL must be a non-negative integer (seconds), got '{raw}'"
)
})?;
Ok(server.with_default_ttl(ttl_seconds))
}
Err(_) => Ok(server),
}
}
#[cfg(feature = "context")]
pub(crate) fn apply_ingest_roots(
server: McpServer,
) -> Result<McpServer, Box<dyn std::error::Error>> {
match std::env::var("VELESDB_MEMORY_INGEST_ROOTS") {
Ok(raw) if !raw.trim().is_empty() => {
let roots = velesdb_memory::context::IngestRoots::parse(&raw)?;
Ok(server.with_ingest_roots(roots))
}
_ => Ok(server),
}
}
#[cfg(not(feature = "context"))]
#[allow(clippy::unnecessary_wraps)]
pub(crate) fn apply_ingest_roots(
server: McpServer,
) -> Result<McpServer, Box<dyn std::error::Error>> {
Ok(server)
}
pub(crate) fn build_server(
service: MemoryService<DynEmbedder>,
) -> Result<McpServer, Box<dyn std::error::Error>> {
let backend = std::env::var("VELESDB_MEMORY_EXTRACTOR").unwrap_or_default();
attach_extractor(McpServer::new(service), &backend)
}
pub(crate) fn attach_extractor(
server: McpServer,
backend: &str,
) -> Result<McpServer, Box<dyn std::error::Error>> {
match velesdb_memory::select_extractor(backend)? {
ExtractorSelection::Disabled => Ok(server),
ExtractorSelection::Ready(extractor) => {
Ok(server.with_named_extractor(backend, extractor)?)
}
ExtractorSelection::NeedsRemoteConfig(backend) => {
Ok(server.with_named_extractor(backend, build_remote_extractor(backend)?)?)
}
}
}