use std::time::Duration;
use rmcp::ServiceExt;
use velesdb_memory::mcp::McpServer;
use velesdb_memory::{DynEmbedder, HashEmbedder, MemoryService, NativeStore, DEFAULT_DIMENSION};
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..]);
}
#[cfg(unix)]
let original_parent = std::os::unix::process::parent_id();
#[cfg(not(unix))]
let original_parent = 0_u32;
let service = build_configured_service(&args)?;
let http_bind = requested_http_bind(&args);
let server = apply_ingest_roots(apply_default_ttl(build_server(service)?)?)?;
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>>(())
}
}
})
}
#[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,
embedder: DynEmbedder,
) -> Result<MemoryService<DynEmbedder>, Box<dyn std::error::Error>> {
use velesdb_memory::MemoryError;
let dimension = embedder.dimension();
let mut last_locked_path: Option<String> = None;
for attempt in 0..LOCK_RETRY_ATTEMPTS {
match NativeStore::open(store_path, dimension) {
Ok(store) => 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) => 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 build_configured_service(
args: &[String],
) -> Result<MemoryService<DynEmbedder>, 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 embedder = build_embedder()?;
apply_autograph(open_store_with_actionable_lock_error(
&store_path,
embedder,
)?)
}
#[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);
}
#[cfg(feature = "extract")]
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);
}
if std::env::var("VELESDB_MEMORY_EXTRACTOR").as_deref() != Ok("ollama") {
return Err(
"autograph is on ([graph] autograph = true / VELESDB_MEMORY_AUTOGRAPH=1) but no \
extraction backend is configured — set [extractor] backend = \"ollama\" (and a \
model), or turn autograph off"
.into(),
);
}
Ok(service.with_autograph(build_ollama_extractor()?))
}
#[cfg(not(feature = "extract"))]
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 Err(
"autograph is on but this binary was built without --features extract, so no \
extraction backend exists"
.into(),
);
}
Ok(service)
}
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)
}
#[cfg(feature = "extract")]
fn build_server(
service: MemoryService<DynEmbedder>,
) -> Result<McpServer, Box<dyn std::error::Error>> {
let server = McpServer::new(service);
match std::env::var("VELESDB_MEMORY_EXTRACTOR").as_deref() {
Ok("ollama") => Ok(server.with_extractor(build_ollama_extractor()?)),
Ok("none") | Err(_) => Ok(server),
Ok(other) => {
Err(format!("unknown VELESDB_MEMORY_EXTRACTOR '{other}' (expected 'ollama')").into())
}
}
}
#[cfg(not(feature = "extract"))]
#[allow(clippy::unnecessary_wraps)]
fn build_server(
service: MemoryService<DynEmbedder>,
) -> Result<McpServer, Box<dyn std::error::Error>> {
Ok(McpServer::new(service))
}
#[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 url = std::env::var("VELESDB_MEMORY_EXTRACTOR_URL")
.unwrap_or_else(|_| DEFAULT_OLLAMA_URL.to_owned());
let model = std::env::var("VELESDB_MEMORY_EXTRACTOR_MODEL").map_err(|_| {
"VELESDB_MEMORY_EXTRACTOR=ollama requires VELESDB_MEMORY_EXTRACTOR_MODEL \
(e.g. qwen3.6:35b-mlx)"
})?;
Ok(Arc::new(OllamaExtractor::new(url, model)))
}
fn build_embedder() -> Result<DynEmbedder, Box<dyn std::error::Error>> {
match std::env::var("VELESDB_MEMORY_EMBEDDER").as_deref() {
Ok("ollama") => build_ollama_embedder(),
Ok("hash") | Err(_) => {
warn_hash_embedder_not_semantic();
Ok(Box::new(HashEmbedder::new(DEFAULT_DIMENSION)))
}
Ok(other) => Err(format!(
"unknown VELESDB_MEMORY_EMBEDDER '{other}' (expected 'hash' or 'ollama')"
)
.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, run an Ollama build with \
VELESDB_MEMORY_EMBEDDER=ollama (see crates/velesdb-memory/README.md). \
Set VELESDB_MEMORY_QUIET=1 to silence this notice."
);
}
#[cfg(feature = "ollama")]
fn build_ollama_embedder() -> Result<DynEmbedder, Box<dyn std::error::Error>> {
use velesdb_memory::{OllamaEmbedder, DEFAULT_OLLAMA_MODEL, DEFAULT_OLLAMA_URL};
let url = std::env::var("VELESDB_MEMORY_OLLAMA_URL")
.unwrap_or_else(|_| DEFAULT_OLLAMA_URL.to_owned());
let model = std::env::var("VELESDB_MEMORY_OLLAMA_MODEL")
.unwrap_or_else(|_| DEFAULT_OLLAMA_MODEL.to_owned());
Ok(Box::new(OllamaEmbedder::new(url, model)?))
}
#[cfg(not(feature = "ollama"))]
fn build_ollama_embedder() -> Result<DynEmbedder, Box<dyn std::error::Error>> {
Err("the 'ollama' embedder requires building with `--features ollama`".into())
}
#[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"));
}
}