use std::time::Duration;
use rmcp::ServiceExt;
use velesdb_memory::mcp::McpServer;
use velesdb_memory::{DynEmbedder, ExtractorSelection, MemoryService, NativeStore};
fn main() -> Result<(), Box<dyn std::error::Error>> {
let args: Vec<String> = std::env::args().collect();
if args
.get(1)
.is_some_and(|arg| arg == "--version" || arg == "-V")
{
println!("velesdb-memory {}", env!("CARGO_PKG_VERSION"));
return Ok(());
}
if args.get(1).is_some_and(|arg| arg == "compile-stdin") {
return run_compile_stdin(&args[2..]);
}
if args.get(1).is_some_and(|arg| arg == "migrate-embeddings") {
return run_migrate_embeddings(&args, &args[2..]);
}
if args.get(1).is_some_and(|arg| arg == "export") {
return run_export(&args, &args[2..]);
}
apply_logging();
#[cfg(unix)]
let original_parent = std::os::unix::process::parent_id();
#[cfg(not(unix))]
let original_parent = 0_u32;
let configured = build_configured_service(&args)?;
let http_bind = requested_http_bind(&args);
let ConfiguredService {
service,
store_path,
embedder_model,
embedder_dimension,
} = configured;
let server = apply_ingest_roots(apply_default_ttl(
build_server(service)?
.with_embedder_identity(embedder_model, embedder_dimension)
.with_store_dir(store_path),
)?)?;
tokio::runtime::Runtime::new()?.block_on(async move {
match http_bind {
#[cfg(feature = "http")]
Some(request) => serve_http(server, request).await,
#[cfg(not(feature = "http"))]
Some(_never) => unreachable!(
"requested_http_bind only returns Some when built with --features http"
),
None => {
spawn_orphan_watchdog(original_parent);
let running = server
.serve((tokio::io::stdin(), tokio::io::stdout()))
.await?;
running.waiting().await?;
Ok::<(), Box<dyn std::error::Error>>(())
}
}
})
}
fn apply_logging() {
if let Err(err) = velesdb_memory::logging::init_from_env() {
eprintln!("[velesdb-memory] {err}");
std::process::exit(1);
}
}
#[cfg(feature = "http")]
struct HttpServeRequest {
bind_addr: String,
insecure: bool,
}
#[cfg(feature = "http")]
fn requested_http_bind(args: &[String]) -> Option<HttpServeRequest> {
let http_flag = args.iter().any(|arg| arg == "--http");
let http_env = std::env::var("VELESDB_MEMORY_HTTP").as_deref() == Ok("1");
if !http_flag && !http_env {
return None;
}
let port_override = args
.iter()
.position(|arg| arg == "--http-port")
.and_then(|flag_index| args.get(flag_index + 1));
let default_bind = std::env::var("VELESDB_MEMORY_HTTP_BIND")
.unwrap_or_else(|_| velesdb_memory::http::DEFAULT_HTTP_BIND.to_owned());
let bind_addr = match port_override {
Some(port) => match default_bind.rsplit_once(':') {
Some((host, _existing_port)) => format!("{host}:{port}"),
None => format!("127.0.0.1:{port}"),
},
None => default_bind,
};
if !is_loopback_host(&bind_addr)
&& std::env::var("VELESDB_MEMORY_HTTP_ALLOW_REMOTE").as_deref() != Ok("1")
{
eprintln!(
"[velesdb-memory] refusing to bind the HTTP transport to '{bind_addr}': it is not a \
loopback address, and the streamable-HTTP transport has no authentication — anyone \
who can reach that socket gets full read/write access to the store. Set \
VELESDB_MEMORY_HTTP_ALLOW_REMOTE=1 to override (put an authenticating reverse proxy \
in front first)."
);
std::process::exit(1);
}
let insecure_flag = args.iter().any(|arg| arg == "--http-insecure");
let insecure_env = std::env::var("VELESDB_MEMORY_HTTP_INSECURE").as_deref() == Ok("1");
let insecure = insecure_flag || insecure_env;
Some(HttpServeRequest {
bind_addr,
insecure,
})
}
#[cfg(feature = "http")]
fn is_loopback_host(bind_addr: &str) -> bool {
let host = bind_addr
.rsplit_once(':')
.map_or(bind_addr, |(host, _port)| host)
.trim_start_matches('[')
.trim_end_matches(']');
host.parse::<std::net::IpAddr>()
.is_ok_and(|ip| ip.is_loopback())
}
#[cfg(not(feature = "http"))]
fn requested_http_bind(args: &[String]) -> Option<String> {
let http_flag = args.iter().any(|arg| arg == "--http");
let http_env = std::env::var("VELESDB_MEMORY_HTTP").as_deref() == Ok("1");
if http_flag || http_env {
eprintln!(
"[velesdb-memory] --http / VELESDB_MEMORY_HTTP=1 requires a binary built with \
`--features http` (e.g. `cargo install velesdb-memory --features http`) — \
this binary was built without it"
);
std::process::exit(1);
}
None
}
#[cfg(feature = "http")]
async fn serve_http(
server: McpServer,
request: HttpServeRequest,
) -> Result<(), Box<dyn std::error::Error>> {
let HttpServeRequest {
bind_addr,
insecure,
} = request;
let ct = tokio_util::sync::CancellationToken::new();
let app = velesdb_memory::http::router(server, ct.child_token());
let listener = tokio::net::TcpListener::bind(&bind_addr).await?;
spawn_shutdown_signals(ct.clone());
if insecure {
eprintln!(
"[velesdb-memory] WARNING: --http-insecure / VELESDB_MEMORY_HTTP_INSECURE=1 is set — \
serving PLAIN HTTP (no TLS) on http://{bind_addr}/mcp. Every request is readable by \
anyone who can reach that socket (loopback-only by default — see \
VELESDB_MEMORY_HTTP_ALLOW_REMOTE above). Use this only for local debugging, or when \
a trusted TLS-terminating proxy already sits in front."
);
eprintln!("[velesdb-memory] HTTP server listening on http://{bind_addr}/mcp");
axum::serve(listener, app)
.with_graceful_shutdown(async move { ct.cancelled_owned().await })
.await?;
return Ok(());
}
let tls_dir = velesdb_memory::tls::tls_dir_from_env();
let material = velesdb_memory::tls::ensure_tls_material(&tls_dir)?;
let acceptor = velesdb_memory::tls::tls_acceptor_from_material(&material)?;
eprintln!("[velesdb-memory] HTTPS server listening on https://{bind_addr}/mcp");
eprintln!(
"[velesdb-memory] Local CA: {} — a client only needs to trust this once (see \
./scripts/install-memory-daemon.sh, which does this automatically on macOS); every \
future leaf certificate this daemon issues is signed by the same CA and is trusted \
automatically after that.",
material.ca_cert_path.display()
);
velesdb_memory::http::serve_tls(app, listener, acceptor, ct).await;
Ok(())
}
#[cfg(unix)]
const ORPHAN_CHECK_INTERVAL: std::time::Duration = std::time::Duration::from_secs(2);
#[cfg(unix)]
fn spawn_orphan_watchdog(original_parent: u32) {
use std::os::unix::process::parent_id;
tokio::spawn(async move {
loop {
tokio::time::sleep(ORPHAN_CHECK_INTERVAL).await;
let current_parent = parent_id();
if current_parent != original_parent {
eprintln!(
"[velesdb-memory] parent process (pid {original_parent}) is gone \
(now reparented under pid {current_parent}) — exiting to release \
the store lock rather than leak a zombie session (#1448)"
);
std::process::exit(0);
}
}
});
}
#[cfg(not(unix))]
fn spawn_orphan_watchdog(_original_parent: u32) {}
const LOCK_RETRY_ATTEMPTS: u32 = 3;
const LOCK_RETRY_DELAY: Duration = Duration::from_millis(500);
fn open_store_with_actionable_lock_error(
store_path: &str,
configured: ConfiguredEmbedder,
) -> Result<MemoryService<DynEmbedder>, Box<dyn std::error::Error>> {
use velesdb_memory::MemoryError;
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();
let mut last_locked_path: Option<String> = None;
for attempt in 0..LOCK_RETRY_ATTEMPTS {
match NativeStore::open(store_path, dimension) {
Ok(store) => {
if unrecorded {
record_embedding_model(store_dir, &store, &model, dimension);
}
return Ok(MemoryService::with_store(store, embedder));
}
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);
}
}
Err(other) if unrecorded => {
return Err(format!(
"{other}\n{}",
velesdb_memory::embedding_provenance::unrecorded_model_note(&model)
)
.into())
}
Err(other) => return Err(other.into()),
}
}
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);
}
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::MemoryStore 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"
);
}
}
}
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 configured = build_embedder()?;
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,
})
}
struct ConfiguredService {
service: MemoryService<DynEmbedder>,
store_path: String,
embedder_model: String,
embedder_dimension: usize,
}
fn run_export(argv: &[String], flags: &[String]) -> Result<(), Box<dyn std::error::Error>> {
apply_config_file(argv)?;
let options = ExportOptions::parse(flags)?;
let store = options
.store_path
.or_else(|| std::env::var("VELESDB_MEMORY_PATH").ok())
.unwrap_or_else(default_store_path);
let store = std::path::Path::new(&store);
let written = write_export(store, options.output.as_deref(), options.include_internal)?;
eprintln!("[velesdb-memory] exported {written} memories");
Ok(())
}
fn write_export(
store: &std::path::Path,
output: Option<&str>,
include_internal: bool,
) -> Result<u64, Box<dyn std::error::Error>> {
if let Some(path) = output {
let mut file = std::io::BufWriter::new(std::fs::File::create(path)?);
let written = velesdb_memory::export::export_jsonl(store, &mut file, include_internal)?;
std::io::Write::flush(&mut file)?;
Ok(written)
} else {
let stdout = std::io::stdout();
let mut lock = stdout.lock();
Ok(velesdb_memory::export::export_jsonl(
store,
&mut lock,
include_internal,
)?)
}
}
struct ExportOptions {
store_path: Option<String>,
output: Option<String>,
include_internal: bool,
}
impl ExportOptions {
fn parse(flags: &[String]) -> Result<Self, Box<dyn std::error::Error>> {
let mut options = Self {
store_path: None,
output: None,
include_internal: false,
};
let mut it = flags.iter();
while let Some(flag) = it.next() {
match flag.as_str() {
"--include-internal" => options.include_internal = true,
"--store" => options.store_path = Some(Self::value_of(&mut it, "--store")?),
"--output" => options.output = Some(Self::value_of(&mut it, "--output")?),
other => return Err(format!("unknown export flag '{other}'").into()),
}
}
Ok(options)
}
fn value_of(
it: &mut std::slice::Iter<'_, String>,
flag: &str,
) -> Result<String, Box<dyn std::error::Error>> {
Ok(it
.next()
.ok_or_else(|| format!("{flag} requires a path argument"))?
.clone())
}
}
fn run_migrate_embeddings(
argv: &[String],
flags: &[String],
) -> Result<(), Box<dyn std::error::Error>> {
use velesdb_memory::migration;
let options = migration::parse_migrate_args(flags)?;
apply_config_file(argv)?;
let store_path = migrate_store_path(&options);
let ConfiguredEmbedder { embedder, model } = build_embedder()?;
let target = migration::TargetContract {
model,
dimension: embedder.dimension(),
strategy: options.strategy,
};
let scratch = migrate_scratch_parent(&options, &store_path)?;
if options.dry_run {
let report = migration::dry_run(
&store_path,
&scratch,
&target,
options.destination.as_deref(),
)?;
print!("{}", migration::render(&report));
if migration::refuses(&report) {
std::process::exit(2);
}
return Ok(());
}
run_migrate_rebuild(&options, &store_path, &scratch, &target, embedder.as_ref())
}
fn run_migrate_rebuild(
options: &velesdb_memory::migration::MigrateOptions,
store_path: &std::path::Path,
scratch: &std::path::Path,
target: &velesdb_memory::migration::TargetContract,
embedder: &dyn velesdb_memory::Embedder,
) -> Result<(), Box<dyn std::error::Error>> {
use velesdb_memory::migration;
let destination = migration::require_destination(options)?;
let outcome = migration::migrate(
store_path,
scratch,
target,
&destination,
embedder,
MIGRATE_BATCH,
)?;
if let Some(executed) = &outcome.executed {
print!("{}", migration::render(&executed.report));
println!(
"rebuild: {} facts written, {} already present, {} edges, journal at {}",
executed.rebuild.facts,
executed.rebuild.collisions,
executed.rebuild.edges,
executed.workspace.display(),
);
}
if let Some(validated) = &outcome.validated {
println!(
"validated: {} facts and {} edges compared, {} divergence(s) explained by expiry",
validated.facts, validated.edges, validated.explained_by_expiry,
);
}
println!("activated: {}", outcome.switched.activated.display());
println!("{}", migration::migration_complete_notice());
Ok(())
}
const MIGRATE_BATCH: usize = 1024;
fn migrate_store_path(options: &velesdb_memory::migration::MigrateOptions) -> std::path::PathBuf {
options.store.clone().unwrap_or_else(|| {
std::path::PathBuf::from(
std::env::var("VELESDB_MEMORY_PATH").unwrap_or_else(|_| default_store_path()),
)
})
}
fn migrate_scratch_parent(
options: &velesdb_memory::migration::MigrateOptions,
store_path: &std::path::Path,
) -> Result<std::path::PathBuf, String> {
if let Some(dir) = options.scratch.clone() {
return Ok(dir);
}
let resolved = std::fs::canonicalize(store_path).unwrap_or_else(|_| store_path.to_path_buf());
velesdb_memory::migration::default_scratch_parent(&resolved)
}
struct ConfiguredEmbedder {
embedder: DynEmbedder,
model: String,
}
#[cfg(feature = "http")]
fn spawn_shutdown_signals(ct: tokio_util::sync::CancellationToken) {
let interrupt = ct.clone();
tokio::spawn(async move {
if tokio::signal::ctrl_c().await.is_ok() {
interrupt.cancel();
}
});
#[cfg(unix)]
tokio::spawn(async move {
if let Ok(mut term) =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
{
term.recv().await;
ct.cancel();
}
});
#[cfg(not(unix))]
drop(ct);
}
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))
}
}
}
#[cfg(not(feature = "extract"))]
fn warn_if_extraction_backend_is_unreachable(_backend: &str) {}
#[cfg(feature = "extract")]
const EXTRACTION_PROBE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3);
#[cfg(feature = "extract")]
fn warn_if_extraction_backend_is_unreachable(backend: &str) {
use velesdb_memory::reachability::{probe_openai, warning_line, Reachability};
if std::env::var_os("VELESDB_MEMORY_QUIET").is_some() {
return;
}
let Some((url, model)) = extraction_endpoint_for_probe(backend) else {
return;
};
let token = env_opt("VELESDB_MEMORY_EXTRACTOR_API_TOKEN");
let outcome = probe_openai(&url, &model, token.as_deref(), EXTRACTION_PROBE_TIMEOUT);
if outcome == Reachability::Reachable {
return;
}
if let Some(line) = warning_line("extraction", &url, &model, &outcome) {
eprintln!("{line}");
}
}
#[cfg(feature = "extract")]
fn extraction_endpoint_for_probe(backend: &str) -> Option<(String, String)> {
let endpoint = extractor_endpoint().ok()?;
let url = match backend {
"ollama" => Some(
endpoint
.url
.unwrap_or_else(|| velesdb_memory::extract::DEFAULT_OLLAMA_URL.to_owned()),
),
"openai" => endpoint.url,
_ => None,
}?;
Some((url, endpoint.model?))
}
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(())
}
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(),
}
}
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")]
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)]
fn apply_ingest_roots(server: McpServer) -> Result<McpServer, Box<dyn std::error::Error>> {
Ok(server)
}
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)
}
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)?)?)
}
}
}
#[cfg(feature = "extract")]
fn build_remote_extractor(
backend: &str,
) -> Result<velesdb_memory::DynExtractor, Box<dyn std::error::Error>> {
match backend {
"ollama" => build_ollama_extractor(),
"openai" => build_openai_extractor(),
other => Err(unwired_backend("extraction", other).into()),
}
}
#[cfg(not(feature = "extract"))]
fn build_remote_extractor(
backend: &str,
) -> Result<velesdb_memory::DynExtractor, Box<dyn std::error::Error>> {
Err(format!(
"VELESDB_MEMORY_EXTRACTOR={backend} needs a build with `--features extract`; \
for an offline deterministic graph with no rebuild, set \
VELESDB_MEMORY_EXTRACTOR=outline instead"
)
.into())
}
#[cfg(feature = "ollama")]
fn embedder_endpoint() -> Result<velesdb_memory::RemoteEndpoint, Box<dyn std::error::Error>> {
let (endpoint, notice) = velesdb_memory::embedder_env_endpoint()?;
if let Some(notice) = notice {
if std::env::var_os("VELESDB_MEMORY_QUIET").is_none() {
eprintln!("{notice}");
}
}
Ok(endpoint)
}
#[cfg(feature = "extract")]
fn extractor_endpoint() -> Result<velesdb_memory::RemoteEndpoint, Box<dyn std::error::Error>> {
Ok(velesdb_memory::RemoteEndpoint {
url: env_opt("VELESDB_MEMORY_EXTRACTOR_URL"),
model: env_opt("VELESDB_MEMORY_EXTRACTOR_MODEL"),
auth: velesdb_memory::role_auth("VELESDB_MEMORY_EXTRACTOR_API_TOKEN")?,
})
}
#[cfg(any(feature = "ollama", feature = "extract"))]
fn env_opt(name: &str) -> Option<String> {
std::env::var(name).ok()
}
#[cfg(any(feature = "ollama", feature = "extract"))]
fn unwired_backend(role: &str, backend: &str) -> String {
format!(
"the {role} backend '{backend}' is accepted by velesdb-memory's selector but \
the daemon has no builder wired for it — this is a bug in velesdb-memory, \
not a configuration error; please report it quoting this message"
)
}
#[cfg(feature = "extract")]
fn build_ollama_extractor() -> Result<velesdb_memory::DynExtractor, Box<dyn std::error::Error>> {
use std::sync::Arc;
use velesdb_memory::extract::DEFAULT_OLLAMA_URL;
use velesdb_memory::OllamaExtractor;
let endpoint = extractor_endpoint()?;
let url = endpoint
.url
.unwrap_or_else(|| DEFAULT_OLLAMA_URL.to_owned());
let model = endpoint.model.ok_or(
"VELESDB_MEMORY_EXTRACTOR=ollama requires VELESDB_MEMORY_EXTRACTOR_MODEL \
(e.g. qwen3.6:35b-mlx)",
)?;
Ok(Arc::new(OllamaExtractor::new(url, model)))
}
#[cfg(feature = "extract")]
fn build_openai_extractor() -> Result<velesdb_memory::DynExtractor, Box<dyn std::error::Error>> {
use std::sync::Arc;
use velesdb_memory::OpenAiExtractor;
let (url, model, auth) = extractor_endpoint()?.require("VELESDB_MEMORY_EXTRACTOR")?;
Ok(Arc::new(OpenAiExtractor::new(url, model, auth)))
}
fn build_embedder() -> Result<ConfiguredEmbedder, Box<dyn std::error::Error>> {
let backend = std::env::var("VELESDB_MEMORY_EMBEDDER");
let selection = velesdb_memory::select_embedder(backend.as_deref().ok())
.map_err(|err| format!("VELESDB_MEMORY_EMBEDDER: {err}"))?;
match selection {
velesdb_memory::EmbedderSelection::Ready("hash", embedder) => {
warn_hash_embedder_not_semantic();
Ok(ConfiguredEmbedder {
embedder,
model: "hash".to_owned(),
})
}
velesdb_memory::EmbedderSelection::Ready(name, embedder) => Ok(ConfiguredEmbedder {
embedder,
model: name.to_owned(),
}),
velesdb_memory::EmbedderSelection::NeedsRemoteConfig(backend) => {
build_remote_embedder(backend)
}
}
}
#[cfg(feature = "ollama")]
fn build_remote_embedder(backend: &str) -> Result<ConfiguredEmbedder, Box<dyn std::error::Error>> {
match backend {
"ollama" => build_ollama_embedder(),
"openai" => build_openai_embedder(),
other => Err(unwired_backend("embedding", other).into()),
}
}
#[cfg(not(feature = "ollama"))]
fn build_remote_embedder(backend: &str) -> Result<ConfiguredEmbedder, Box<dyn std::error::Error>> {
Err(format!(
"the '{backend}' embedder requires building with `--features ollama` \
(that feature carries the HTTP dependency for both remote embedding \
backends); VELESDB_MEMORY_EMBEDDER=hash needs no rebuild"
)
.into())
}
fn warn_hash_embedder_not_semantic() {
if std::env::var_os("VELESDB_MEMORY_QUIET").is_some() {
return;
}
eprintln!(
"[velesdb-memory] Using the default 'hash' embedder: deterministic and \
fully offline, but NOT semantic — recall matches surface form, not meaning. \
For real semantic recall set VELESDB_MEMORY_EMBEDDER=ollama or =openai \
(no rebuild needed; see crates/velesdb-memory/README.md for the model \
to pull). Set VELESDB_MEMORY_QUIET=1 to silence this notice."
);
}
#[cfg(feature = "ollama")]
fn build_ollama_embedder() -> Result<ConfiguredEmbedder, Box<dyn std::error::Error>> {
use velesdb_memory::{OllamaEmbedder, DEFAULT_OLLAMA_MODEL, DEFAULT_OLLAMA_URL};
let endpoint = embedder_endpoint()?;
let url = endpoint
.url
.unwrap_or_else(|| DEFAULT_OLLAMA_URL.to_owned());
let model = endpoint
.model
.unwrap_or_else(|| DEFAULT_OLLAMA_MODEL.to_owned());
Ok(ConfiguredEmbedder {
embedder: Box::new(OllamaEmbedder::new(&url, &model)?),
model,
})
}
#[cfg(feature = "ollama")]
fn build_openai_embedder() -> Result<ConfiguredEmbedder, Box<dyn std::error::Error>> {
use velesdb_memory::OpenAiEmbedder;
let (url, model, auth) = embedder_endpoint()?.require("VELESDB_MEMORY_EMBEDDER")?;
Ok(ConfiguredEmbedder {
embedder: Box::new(OpenAiEmbedder::new(url, &model, auth)?),
model,
})
}
#[cfg(feature = "context")]
const DEFAULT_COMPILE_STDIN_BUDGET: u64 = 2_000;
#[cfg(feature = "context")]
#[derive(Debug, PartialEq, Eq)]
struct CompileStdinOptions {
token_budget: u64,
query: String,
}
#[cfg(feature = "context")]
impl Default for CompileStdinOptions {
fn default() -> Self {
Self {
token_budget: DEFAULT_COMPILE_STDIN_BUDGET,
query: String::new(),
}
}
}
#[cfg(feature = "context")]
#[derive(serde::Serialize)]
struct CompileStdinOutput {
content: String,
tokens_in: u64,
tokens_out: u64,
tokens_saved: u64,
risk: String,
}
#[cfg(feature = "context")]
fn parse_compile_stdin_budget(value: Option<&String>) -> Result<u64, String> {
let raw = value.ok_or_else(|| "--budget requires a value".to_owned())?;
let parsed: u64 = raw
.parse()
.map_err(|_| format!("--budget expects a positive integer, got {raw:?}"))?;
if parsed == 0 {
return Err("--budget must be greater than 0".to_owned());
}
Ok(parsed)
}
#[cfg(feature = "context")]
fn parse_compile_stdin_args(args: &[String]) -> Result<CompileStdinOptions, String> {
let mut options = CompileStdinOptions::default();
let mut index = 0;
while index < args.len() {
let flag = args[index].as_str();
let value = args.get(index + 1);
match flag {
"--budget" => {
options.token_budget = parse_compile_stdin_budget(value)?;
index += 2;
}
"--query" => {
options
.query
.clone_from(value.ok_or_else(|| "--query requires a value".to_owned())?);
index += 2;
}
other => return Err(format!("unknown compile-stdin flag {other:?}")),
}
}
Ok(options)
}
#[cfg(feature = "context")]
fn compile_stdin_json(
text: &str,
options: &CompileStdinOptions,
) -> Result<String, Box<dyn std::error::Error>> {
use velesdb_memory::context::{
segment_transcript, CompilePolicy, CompileRequest, ContextCompiler, SegmentationPolicy,
};
if text.trim().is_empty() {
return Err("compile-stdin received empty input on stdin".into());
}
let outcome = segment_transcript(text, &SegmentationPolicy::default())?;
let request = CompileRequest {
query: options.query.clone(),
fragments: outcome
.segments
.into_iter()
.map(|segment| segment.fragment)
.collect(),
project: None,
target_model: None,
token_budget: options.token_budget,
memory_scope: None,
policy: None,
};
let compiled = ContextCompiler::new(CompilePolicy::default()).compile(&request)?;
if compiled.content.is_empty() {
return Err(format!(
"a budget of {} tokens fits none of the {} input tokens — every fragment was \
externalized and the compiled context is empty; raise --budget",
options.token_budget, compiled.insights.tokens_in
)
.into());
}
let output = CompileStdinOutput {
content: compiled.content,
tokens_in: compiled.insights.tokens_in,
tokens_out: compiled.insights.tokens_out,
tokens_saved: compiled.insights.tokens_saved,
risk: format!("{:?}", compiled.risk).to_lowercase(),
};
Ok(serde_json::to_string(&output)?)
}
#[cfg(feature = "context")]
fn run_compile_stdin(args: &[String]) -> Result<(), Box<dyn std::error::Error>> {
use std::io::Read as _;
let options = parse_compile_stdin_args(args)?;
let mut text = String::new();
std::io::stdin().read_to_string(&mut text)?;
println!("{}", compile_stdin_json(&text, &options)?);
Ok(())
}
#[cfg(not(feature = "context"))]
fn run_compile_stdin(_args: &[String]) -> Result<(), Box<dyn std::error::Error>> {
Err("`compile-stdin` requires building with `--features context`".into())
}
#[cfg(all(test, feature = "context"))]
mod compile_stdin_tests {
use super::{
compile_stdin_json, parse_compile_stdin_args, CompileStdinOptions,
DEFAULT_COMPILE_STDIN_BUDGET,
};
fn noisy_tool_output() -> String {
use std::fmt::Write as _;
let mut text = String::new();
for i in 0..120 {
let _ = writeln!(
text,
"[2026-07-25T01:0{}:00Z] INFO worker: processing batch {} of 120 — retry=0 status=ok",
i % 10,
i
);
}
text
}
fn parse(value: &str) -> serde_json::Value {
serde_json::from_str(value).expect("compile-stdin must emit valid JSON")
}
#[test]
fn tight_budget_actually_shrinks_the_payload() {
let options = CompileStdinOptions {
token_budget: 1_500,
query: "what did the worker do".to_owned(),
};
let compiled = parse(&compile_stdin_json(&noisy_tool_output(), &options).unwrap());
let tokens_in = compiled["tokens_in"].as_u64().unwrap();
let tokens_out = compiled["tokens_out"].as_u64().unwrap();
assert!(tokens_in > 0, "tokens_in must be measured, got {tokens_in}");
assert!(
tokens_out < tokens_in,
"a 200-token budget over {tokens_in} tokens of logs must compress: got {tokens_out}"
);
assert_eq!(
compiled["tokens_saved"].as_u64().unwrap(),
tokens_in - tokens_out
);
let content = compiled["content"].as_str().unwrap();
assert!(
!content.is_empty(),
"an empty compilation is worse than no compilation — the caller would replace a \
real tool result with nothing"
);
assert!(
content.len() < noisy_tool_output().len(),
"the compiled content must be shorter than the raw tool output"
);
}
#[test]
fn budget_too_small_for_any_fragment_is_an_error() {
let options = CompileStdinOptions {
token_budget: 50,
query: String::new(),
};
let error = compile_stdin_json(&noisy_tool_output(), &options).unwrap_err();
let message = error.to_string();
assert!(
message.contains("budget"),
"the error must point at the budget, got {message}"
);
}
#[test]
fn compilation_is_byte_identical_across_runs() {
let options = CompileStdinOptions {
token_budget: 1_500,
query: "worker batches".to_owned(),
};
let first = compile_stdin_json(&noisy_tool_output(), &options).unwrap();
let second = compile_stdin_json(&noisy_tool_output(), &options).unwrap();
assert_eq!(first, second, "the compiler must be deterministic");
}
#[test]
fn empty_stdin_is_rejected() {
let error = compile_stdin_json(" \n\t ", &CompileStdinOptions::default()).unwrap_err();
assert!(
error.to_string().contains("empty"),
"the error must name the cause, got {error}"
);
}
#[test]
fn flags_default_and_override() {
assert_eq!(
parse_compile_stdin_args(&[]).unwrap(),
CompileStdinOptions {
token_budget: DEFAULT_COMPILE_STDIN_BUDGET,
query: String::new(),
}
);
let parsed = parse_compile_stdin_args(&[
"--budget".to_owned(),
"512".to_owned(),
"--query".to_owned(),
"why did it fail".to_owned(),
])
.unwrap();
assert_eq!(parsed.token_budget, 512);
assert_eq!(parsed.query, "why did it fail");
}
#[test]
fn malformed_flags_are_rejected() {
for bad in [
vec!["--budget".to_owned()],
vec!["--budget".to_owned(), "zero".to_owned()],
vec!["--budget".to_owned(), "0".to_owned()],
vec!["--nope".to_owned()],
] {
assert!(
parse_compile_stdin_args(&bad).is_err(),
"must reject {bad:?}"
);
}
}
}
#[cfg(all(test, feature = "http"))]
mod tests {
use super::is_loopback_host;
#[test]
fn loopback_v4_and_v6_are_recognized() {
assert!(is_loopback_host("127.0.0.1:18090"));
assert!(is_loopback_host("127.0.0.5:18090"));
assert!(is_loopback_host("[::1]:18090"));
}
#[test]
fn non_loopback_hosts_are_rejected() {
assert!(!is_loopback_host("0.0.0.0:18090"));
assert!(!is_loopback_host("192.168.1.10:18090"));
assert!(!is_loopback_host("[::]:18090"));
assert!(!is_loopback_host("mcp.example.com:18090"));
}
}