//! MCP server: JSON-RPC 2.0 over newline-delimited stdio.
//!
//! Framing is **one JSON object per line**. LSP-style `Content-Length`
//! headers are not accepted — a header line is a parse error (`-32700`)
//! and the loop continues. Blank lines are skipped.
//!
//! # Methods
//!
//! - `initialize` — `protocolVersion` `"2024-11-05"`, `capabilities.tools`,
//! `serverInfo.name` `"mushroomdb"`, `serverInfo.version` (crate version)
//! - `notifications/initialized` — ignored
//! - `tools/list` — the default listing is
//! [`ASSOCIATION_TOOLS`], the tools that answer a question about an entity
//! graph, in that order, whatever store the server opened. Graph-tool
//! descriptions carry the prefix `Advanced: ` so a host ranking tools by
//! description puts the task tools in front. `mushroomdb mcp --all-tools`
//! lists every served tool; the rest are callable either way, just not
//! advertised
//! - `tools/call` — dispatch; success for a graph tool is
//! `{content:[{type:"text", text:<json string>}]}`, and for a task tool one
//! text block holding the rendered digest — or, with `json: true`, the
//! serialised report. No task tool returns `structuredContent`
//!
//! Unknown methods on a **request** (has `id`) → `-32601`. A notification
//! (no `id` member) never writes a response, including unknown methods.
//!
//! # Error split
//!
//! Protocol errors are JSON-RPC `error` objects:
//! - `-32700` parse — unparseable line (invalid JSON / invalid UTF-8)
//! - `-32600` invalid request — parsed JSON that is not an object, or a
//! request with missing / non-string `method`
//! - `-32601` method — unknown `method` on a request
//! - `-32602` params — `tools/call` envelope invalid: `params` not an object,
//! missing / non-string `name`, `arguments` present but not an object,
//! or unknown tool name
//!
//! Tool-level failures are JSON-RPC **results** with `isError: true` and a
//! text message: missing or wrong-typed fields inside a known tool's
//! `arguments`, and every [`GraphError`] from core-api.
//!
//! # Deadlock
//!
//! [`SharedDb::read`] / [`SharedDb::write`] guards are held only for the
//! public core-api call, then dropped before serializing or writing. Do not
//! nest a second lock on the same handle (the `RwLock` is not re-entrant).
//!
//! EOF on `reader` returns `Ok(())`. Read/write I/O errors propagate.
use crate::json::{
edge_history_result_json, namespace_arg, node_history_json, node_info_json, params_from_json,
parse_ingest_edges, result_set_json, rule_def_from_json, stamp_namespace, stamp_namespace_row,
};
use core_api::{
json_to_rows, json_to_value, memory, AsOfScope, AutoFk, GraphError, IngestOptions, MaskMode,
NodeMask, PropPredicate, SharedDb, Value,
};
use serde_json::{json, Value as Js};
use std::collections::BTreeMap;
use std::io::{self, BufRead, Write};
use std::path::{Path, PathBuf};
/// Run the MCP loop until `reader` hits EOF.
///
/// `db_dir` is where the store lives on disk. `mushroomdb mcp <db>` passes it;
/// a caller that has only a handle passes `None`. Its one consumer is
/// `recall`, which names the store by this path in its digest header and in
/// the `schema apply` command it suggests when there is no text index; with
/// `None` it says "store" instead.
pub fn run_mcp_stdio(
db: SharedDb,
db_dir: Option<PathBuf>,
reader: impl BufRead,
writer: impl Write,
) -> io::Result<()> {
run_mcp_stdio_with(db, db_dir, false, reader, writer)
}
/// [`run_mcp_stdio`], with the tool list chosen by the caller.
///
/// `all_tools` false lists [`ASSOCIATION_TOOLS`]; true lists every served
/// tool. Either way every tool remains callable — the flag decides
/// what is advertised, not what is served.
pub fn run_mcp_stdio_with(
db: SharedDb,
db_dir: Option<PathBuf>,
all_tools: bool,
mut reader: impl BufRead,
mut writer: impl Write,
) -> io::Result<()> {
let mut buf = Vec::new();
loop {
buf.clear();
let n = reader.read_until(b'\n', &mut buf)?;
if n == 0 {
return Ok(());
}
match std::str::from_utf8(&buf) {
Ok(s) if s.trim().is_empty() => continue,
Ok(s) => handle_line(&db, db_dir.as_deref(), all_tools, s.trim(), &mut writer)?,
Err(_) => write_error(&mut writer, None, -32700, "Parse error")?,
}
}
}
fn handle_line(
db: &SharedDb,
db_dir: Option<&Path>,
all_tools: bool,
line: &str,
writer: &mut impl Write,
) -> io::Result<()> {
let msg: Js = match serde_json::from_str(line) {
Ok(v) => v,
Err(_) => return write_error(writer, None, -32700, "Parse error"),
};
let Some(obj) = msg.as_object() else {
return write_error(writer, None, -32600, "Invalid Request");
};
let is_request = obj.contains_key("id");
let id = obj.get("id").cloned();
let method = match obj.get("method").and_then(Js::as_str) {
Some(m) => m,
None => {
if is_request {
write_error(writer, id, -32600, "Invalid Request")?;
}
return Ok(());
}
};
match method {
"initialize" => {
if is_request {
write_result(writer, id, initialize_result())?;
}
}
"notifications/initialized" => {
if is_request {
write_result(writer, id, json!({}))?;
}
}
"tools/list" => {
if is_request {
write_result(writer, id, tools_list(all_tools))?;
}
}
"tools/call" => {
if is_request {
match dispatch_call(db, db_dir, obj.get("params")) {
CallOutcome::Protocol { code, message } => {
write_error(writer, id, code, &message)?;
}
CallOutcome::ToolOk(payload) => {
write_result(writer, id, tool_ok(payload))?;
}
CallOutcome::TaskOk { text } => {
write_result(writer, id, task_ok(&text))?;
}
CallOutcome::ToolErr(message) => {
write_result(writer, id, tool_err(&message))?;
}
}
}
}
_ => {
if is_request {
write_error(writer, id, -32601, "Method not found")?;
}
}
}
Ok(())
}
pub(crate) enum CallOutcome {
Protocol {
code: i64,
message: String,
},
/// A graph tool's JSON payload, returned as a JSON string in `content`.
ToolOk(Js),
/// A task tool's answer, as text and nothing else: the rendered digest, or
/// the serialised report when the call passed `json: true`.
TaskOk {
text: String,
},
ToolErr(String),
}
fn dispatch_call(db: &SharedDb, db_dir: Option<&Path>, params: Option<&Js>) -> CallOutcome {
let Some(params) = params.and_then(Js::as_object) else {
return protocol_invalid();
};
let Some(name) = params.get("name").and_then(Js::as_str) else {
return protocol_invalid();
};
let empty = json!({});
let args = match params.get("arguments") {
None => &empty,
Some(a) if a.is_object() => a,
Some(_) => return protocol_invalid(),
};
// The task tools first, in the order `tools/list` advertises.
if let Some(outcome) = crate::mcp_tasks::dispatch(db, db_dir, name, args) {
return outcome;
}
match name {
"query" => tool_query(db, args),
"ingest_json" => tool_ingest(db, args),
"create_rule" => tool_create_rule(db, args),
"explain" => tool_explain(db, args),
"stats" => tool_stats(db, args),
"node_info" => tool_node_info(db, args),
"upsert_entity" => tool_upsert_entity(db, args),
"find_similar" => tool_find_similar(db, args),
"pairwise_similar" => tool_pairwise_similar(db, args),
"hybrid_search" => tool_hybrid_search(db, args),
"node_history" => tool_node_history(db, args),
"edge_history" => tool_edge_history(db, args),
"was_linked" => tool_was_linked(db, args),
"rename_node" => tool_rename_node(db, args),
_ => protocol_invalid(),
}
}
fn protocol_invalid() -> CallOutcome {
CallOutcome::Protocol {
code: -32602,
message: "Invalid params".into(),
}
}
fn tool_query(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(cypher) = args.get("cypher").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing cypher".into());
};
let params = match params_from_json(args.get("params")) {
Ok(p) => p,
Err(e) => return CallOutcome::ToolErr(e),
};
// Two ways to ask the same restricted question: a `role` names one the
// store already defines, a `mask` writes the allow-list out by hand. Both
// route to `query_masked` (read-only). Passing both is not a merge of the
// two — it is a caller that has not decided which restriction applies, so
// it is refused rather than silently resolved one way.
let role = match args.get("role") {
None | Some(Js::Null) => None,
Some(Js::String(s)) if !s.is_empty() => Some(s.as_str()),
Some(_) => return CallOutcome::ToolErr("role must be a non-empty string".into()),
};
let mask_keys = match args.get("mask") {
None => None,
Some(v) => match mask_key_list(v) {
Ok(keys) => Some(keys),
Err(e) => return CallOutcome::ToolErr(e),
},
};
if role.is_some() && mask_keys.is_some() {
return CallOutcome::ToolErr("pass role or mask, not both".into());
}
// The second visibility axis. `namespace` is not a third way to say what
// `role` and `mask` say — it is a leg that **intersects** whichever of them
// is present (and stands alone when neither is), so it can only narrow what
// they already allow. A role bound to namespaces honours them with no
// argument here; passing one outside the binding is the empty intersection,
// never the union.
let namespace = match namespace_arg(args.get("namespace")) {
Ok(n) => n,
Err(e) => return CallOutcome::ToolErr(e),
};
// Optional time travel: a 0-based WAL commit index. The graph is read as
// of that commit; a `role` is still the role the store defines now, since
// `roles.json` is a sidecar and is never a WAL record.
let as_of = match args.get("as_of") {
None | Some(Js::Null) => None,
Some(v) => match v.as_u64() {
Some(n) => Some(n),
None => {
return CallOutcome::ToolErr(
"as_of must be a non-negative integer commit index".into(),
)
}
},
};
if let Some(commit) = as_of {
// Stub mode discloses node existence, which is exactly the question an
// as-of read is asking. The two do not compose.
if args
.get("stub_hidden")
.and_then(|v| v.as_bool())
.unwrap_or(false)
{
return CallOutcome::ToolErr(
"as_of (time-travel) does not compose with stub_hidden".into(),
);
}
let scope = match (role, &mask_keys) {
(Some(role), _) => AsOfScope::Role(role),
(None, Some(keys)) => AsOfScope::Keys(keys),
(None, None) => match namespace.as_deref() {
// A namespace alone is its own as-of scope.
Some(ns) => AsOfScope::Namespace(ns),
None => {
return match db.read().query_at(commit, cypher, ¶ms) {
Ok(rs) => CallOutcome::ToolOk(result_set_json(&rs)),
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
},
};
let g = db.read();
let out = match (namespace.as_deref(), role.is_some() || mask_keys.is_some()) {
// Both legs: the namespace intersects the scope at that commit.
(Some(ns), true) => g.query_at_scoped_in_namespace(commit, cypher, ¶ms, scope, ns),
_ => g.query_at_scoped(commit, cypher, ¶ms, scope),
};
return match out {
Ok(rs) => CallOutcome::ToolOk(result_set_json(&rs)),
Err(GraphError::KeyNotFound { key }) if key.starts_with("role:") => {
CallOutcome::ToolErr(format!("unknown role '{}'", &key["role:".len()..]))
}
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
};
}
if role.is_some() || mask_keys.is_some() || namespace.is_some() {
let stub_hidden = args
.get("stub_hidden")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let g = db.read();
let mask = match (role, &mask_keys) {
(Some(role), _) => match g.mask_for_role(role) {
Ok(m) => m,
// The one error a caller can fix by rereading `roles.json`,
// told apart from a store whose roles never loaded at all.
Err(GraphError::KeyNotFound { .. }) => {
return CallOutcome::ToolErr(format!("unknown role '{role}'"))
}
Err(e) => return CallOutcome::ToolErr(graph_err_msg(e)),
},
(None, Some(keys)) => NodeMask::from_keys(&*g, keys.iter().map(String::as_str)),
// A namespace alone: the namespace leg is the whole mask.
(None, None) => g.mask_for_namespace(
namespace
.as_deref()
.expect("one of the three is Some in this branch"),
),
};
// With a role or a client mask present, the namespace is a second leg
// intersected into it — the same `NodeMask::intersect` the
// role-plus-client-mask path uses, so never-widen holds by construction.
let mask = match (namespace.as_deref(), role.is_some() || mask_keys.is_some()) {
(Some(ns), true) => mask.intersect(&g.mask_for_namespace(ns)),
_ => mask,
};
let mask = if stub_hidden {
mask.with_mode(MaskMode::Stub)
} else {
mask
};
return match g.query_masked(cypher, ¶ms, &mask) {
Ok(rs) => CallOutcome::ToolOk(result_set_json(&rs)),
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
};
}
let is_write = match core_api::is_write_query(cypher) {
Ok(b) => b,
Err(e) => return CallOutcome::ToolErr(e),
};
let rs = if is_write {
let mut g = db.write();
g.query_write(cypher, ¶ms)
} else {
let g = db.read();
g.query(cypher, ¶ms)
};
match rs {
Ok(rs) => CallOutcome::ToolOk(result_set_json(&rs)),
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
/// A `mask` argument as a key list. `Err` when it is anything but an array of
/// strings — including `null`, which is a caller that meant to pass one.
fn mask_key_list(mask: &Js) -> Result<Vec<String>, String> {
let arr = mask
.as_array()
.ok_or_else(|| "mask must be an array of strings".to_string())?;
arr.iter()
.map(|v| {
v.as_str()
.map(str::to_string)
.ok_or_else(|| "mask must be an array of strings".to_string())
})
.collect()
}
fn tool_ingest(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(label) = args.get("label").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing label".into());
};
let Some(rows_json) = args.get("rows_json").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing rows_json".into());
};
let mut opts = IngestOptions::default();
if let Some(kf) = args.get("key_field") {
match kf.as_str() {
Some(s) => opts.key_field = s.to_string(),
None => return CallOutcome::ToolErr("key_field must be a string".into()),
}
}
if let Some(suf) = args.get("auto_fk_suffix") {
match suf.as_str() {
Some(s) => {
opts.auto_fk = AutoFk::Auto {
suffix: s.to_string(),
}
}
None => return CallOutcome::ToolErr("auto_fk_suffix must be a string".into()),
}
}
let edges = match args.get("edges") {
None | Some(Js::Null) => Vec::new(),
Some(raw) => match parse_ingest_edges(raw) {
Ok(e) => e,
Err(e) => return CallOutcome::ToolErr(e),
},
};
let parsed: Js = match serde_json::from_str(rows_json) {
Ok(v) => v,
Err(e) => {
return CallOutcome::ToolErr(graph_err_msg(GraphError::IngestError {
detail: e.to_string(),
}))
}
};
let mut converted = match json_to_rows(&parsed) {
Ok(c) => c,
Err(e) => return CallOutcome::ToolErr(graph_err_msg(e)),
};
// `namespace` applies to every node this call creates.
let namespace = match namespace_arg(args.get("namespace")) {
Ok(n) => n,
Err(e) => return CallOutcome::ToolErr(e),
};
if let Err(e) = stamp_namespace(&mut converted.rows, namespace.as_deref()) {
return CallOutcome::ToolErr(e);
}
let taken = std::mem::take(&mut converted.rows);
let report = {
let mut g = db.write();
g.ingest_with_edges(label, taken, &opts, &edges)
};
match report.map(|r| converted.into_report(r)) {
Ok(r) => match serde_json::to_value(&r) {
Ok(v) => CallOutcome::ToolOk(v),
Err(e) => CallOutcome::ToolErr(e.to_string()),
},
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
fn tool_create_rule(db: &SharedDb, args: &Js) -> CallOutcome {
let def = match rule_def_from_json(args.clone()) {
Ok(d) => d,
Err(e) => return CallOutcome::ToolErr(e),
};
let name = def.name.clone();
let res = {
let mut g = db.write();
g.create_rule(def)
};
if let Err(e) = res {
return CallOutcome::ToolErr(graph_err_msg(e));
}
// A rule over a corpus too large to index in one commit is installed but
// derives nothing yet. Saying "ok" there would tell the caller to go and
// query edges that do not exist, so report the build instead — `stats`
// carries the same progress under each rule's `building`.
let building = db
.read()
.builds_in_progress()
.into_iter()
.find(|b| b.rule == name);
match building {
Some(b) => CallOutcome::ToolOk(json!({
"ok": true,
"name": name,
"building": {"indexed": b.indexed, "total": b.total},
"note": format!(
"the vector index for {name:?} is still being built ({}/{} vectors); \
this rule derives no edges until it finishes. Every write advances it, \
and `mushroomdb build-index <db-dir>` finishes it now. Poll `stats` — \
the rule's `building` field disappears when its edges are in.",
b.indexed, b.total
),
})),
None => CallOutcome::ToolOk(json!({"ok": true, "name": name})),
}
}
fn tool_explain(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(a) = args.get("a").and_then(Js::as_str).filter(|s| !s.is_empty()) else {
return CallOutcome::ToolErr("missing a".into());
};
let Some(b) = args.get("b").and_then(Js::as_str).filter(|s| !s.is_empty()) else {
return CallOutcome::ToolErr("missing b".into());
};
let out = {
let g = db.read();
g.explain(a, b)
};
match out {
Ok(v) => match serde_json::to_value(&v) {
Ok(j) => CallOutcome::ToolOk(j),
Err(e) => CallOutcome::ToolErr(e.to_string()),
},
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
/// `stats`, with the namespace roster narrowed when the caller names a role or a
/// namespace.
///
/// The roster is the one part of `stats` that is a list of *other tenants*:
/// every namespace and its live count. A caller answering as a role should be
/// told about its own namespaces and no others, so `role` narrows the roster to
/// the role's binding and `namespace` to that one name. The store-wide counts
/// beside it are unchanged — they were never per-namespace and narrowing them
/// would make the two halves of one body disagree.
fn tool_stats(db: &SharedDb, args: &Js) -> CallOutcome {
let role = match args.get("role") {
None | Some(Js::Null) => None,
Some(Js::String(s)) if !s.is_empty() => Some(s.clone()),
Some(_) => return CallOutcome::ToolErr("role must be a non-empty string".into()),
};
let namespace = match namespace_arg(args.get("namespace")) {
Ok(n) => n,
Err(e) => return CallOutcome::ToolErr(e),
};
let (snap, role_def) = {
let g = db.read();
let def = match &role {
Some(r) => {
// Resolve through the same resolver `query` uses, so a store
// whose `roles.json` was corrupt at open says so here too
// instead of reporting the role simply unknown — one answer per
// cause, the same one on both tools. The mask is memoised per
// (role, commit_seq), so asking costs nothing a `query` with the
// same role would not already have paid.
if let Err(e) = g.mask_for_role(r) {
return match e {
GraphError::KeyNotFound { .. } => {
CallOutcome::ToolErr(format!("unknown role '{r}'"))
}
other => CallOutcome::ToolErr(graph_err_msg(other)),
};
}
g.roles().into_iter().find(|d| &d.name == r)
}
None => None,
};
(g.stats(), def)
};
let mut snap = snap;
// Scoped on the ARGUMENTS, not on `role_def`: naming a role that resolved
// is a scoped call even in the window where the definition lookup misses,
// and the safe direction there is to omit the roster rather than send all
// of it.
let scoped = role.is_some() || namespace.is_some();
if scoped {
snap.namespaces.retain(|n| {
role_def.as_ref().is_none_or(|d| d.sees_namespace(&n.name))
&& namespace.as_deref().is_none_or(|ns| ns == n.name)
});
}
match serde_json::to_value(&snap) {
Ok(mut v) => {
// An unscoped caller gets the store-wide counts without the roster.
// Omitted, not emptied: `"namespaces": []` still discloses that the
// roster exists and invites a guess at its size.
if !scoped {
if let Some(obj) = v.as_object_mut() {
obj.remove("namespaces");
}
}
CallOutcome::ToolOk(v)
}
Err(e) => CallOutcome::ToolErr(e.to_string()),
}
}
fn tool_node_info(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(key) = args.get("key").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing key".into());
};
let info = {
let g = db.read();
g.node_info(key)
};
match info {
Some(info) => CallOutcome::ToolOk(node_info_json(&info)),
None => CallOutcome::ToolErr(graph_err_msg(GraphError::KeyNotFound {
key: key.to_string(),
})),
}
}
/// Insert a new node or update an existing node's properties, keyed by `key`.
///
/// The policy and the write are [`memory::remember::upsert_entity`] — the one
/// implementation this tool and the Python binding share — which also clears
/// a provisional mark (`memory_schema::PROVISIONAL_PROP`) on `key` if it had
/// one, whoever set it. This function parses the arguments and shapes the
/// reply.
///
/// If the node exists: every supplied property is checked (reserved names, the
/// `ns` rule, a view-owned field, type) and then all of them are written in one
/// engine commit — a refusal leaves the node unchanged. If the node does not
/// exist: `label` is required; the node is created and then given the
/// supplied props, `id` among them.
///
/// `namespace` is the namespace a node this call **creates** is created in. On a
/// node that already exists it is written like any other property, which is what
/// makes naming the namespace the node is already in a no-op and naming another
/// one the engine's `NamespaceImmutable` refusal — one rule, stated once, in the
/// place that owns it. A no-op `ns` writes nothing and is not counted in
/// `updated_fields`, because nothing was updated.
///
/// `id` in `props` is **dropped on both paths**: it is the node's key. The create
/// path stores `id` from `key`, and `rename_node` is the only way to change it.
/// One row builder now serves the create and the update path, so the rule is the
/// same on both — before v0.6.6 the update path wrote `props.id` straight through
/// `set_prop`, which could leave a stored `id` disagreeing with the key the node
/// is reached by, while the create path had always ignored it.
///
/// A node's **label never changes on an update** — the engine has no
/// label-mutation path at all, and `describe_entity` uses `label` only on
/// create. A `label` that names something else on an existing `key` is a
/// refusal, checked before anything is written: nothing in `props` is
/// touched either, matching every other pre-write check this tool makes (a
/// namespace mismatch, a view-owned field). Before fix round 2 this was
/// silently accepted and dropped — every update, not only a provisional
/// stub's — which is worse than it sounds for the one case the feature
/// exists for: a provisional entity `remember` created is *always* `Entity`,
/// permanently, and `{"created":false,"ok":true}` beside a quietly-ignored
/// `label` is not a signal an agent acts on. There is a working path instead
/// — `remember`'s `entities`, at first mention, produces a correctly
/// labelled, non-provisional node with its `ABOUT` edge intact — so refusing
/// costs nothing a caller could not already avoid.
///
/// Returns `{ok, key, label, created, updated_fields?}`.
fn tool_upsert_entity(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(key) = args.get("key").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing key".into());
};
let label_opt = args.get("label").and_then(Js::as_str);
let Some(props_obj) = args.get("props").and_then(Js::as_object) else {
return CallOutcome::ToolErr("missing props".into());
};
let namespace = match namespace_arg(args.get("namespace")) {
Ok(n) => n,
Err(e) => return CallOutcome::ToolErr(e),
};
let aliases: Vec<String> = match args.get("aliases") {
None | Some(Js::Null) => Vec::new(),
Some(Js::Array(items)) => match items
.iter()
.map(|v| v.as_str().map(str::to_string))
.collect::<Option<Vec<_>>>()
{
Some(a) => a,
None => return CallOutcome::ToolErr("aliases must be an array of strings".into()),
},
Some(_) => return CallOutcome::ToolErr("aliases must be an array of strings".into()),
};
// One row, stamped with the namespace, whichever path takes it: the
// conflict rule between an explicit `props.ns` and `namespace` is then the
// same one `ingest_json` applies.
let mut row: BTreeMap<String, Value> = BTreeMap::new();
for (field, json_val) in props_obj {
if field == "id" {
continue;
}
match json_to_value(json_val.clone()) {
Some(v) => {
row.insert(field.clone(), v);
}
None => {
return CallOutcome::ToolErr(format!("prop {field} is not a supported value type"))
}
}
}
if let Some(ns) = namespace.as_deref() {
if let Err(e) = stamp_namespace_row(&mut row, ns) {
return CallOutcome::ToolErr(e);
}
}
// One write guard for the existence check and the write: the policy —
// label on create, no relabel, the namespace no-op, `id` from the key —
// is `core-api`'s, shared with the Python binding.
let outcome = {
let mut g = db.write();
match memory::remember::upsert_entity(&mut *g, key, label_opt, row, &aliases) {
Ok(o) => o,
Err(e) => return CallOutcome::ToolErr(graph_err_msg(e)),
}
};
if outcome.created {
return CallOutcome::ToolOk(json!({
"ok": true,
"key": key,
"label": outcome.label,
"created": true
}));
}
let mut reply = json!({
"ok": true,
"key": key,
"label": outcome.label,
"created": false,
"updated_fields": outcome.updated_fields
});
// Absent when nothing was lost, so an ordinary update's reply is
// unchanged.
let lost = &outcome.same_as_lost;
if !lost.is_empty() {
reply["same_as_lost"] = crate::mcp_tasks::same_as_lost_json(lost);
reply["same_as_lost_total"] = json!(lost.len());
reply["same_as_lost_note"] = json!(format!(
"this write retracted {} SAME_AS link(s): {}",
lost.len(),
crate::mcp_tasks::SAME_AS_LOST_REMEDY
));
}
CallOutcome::ToolOk(reply)
}
/// Return neighbors connected by a given edge type (default `"SIMILAR"`).
///
/// Results are read from edges already materialized by a derivation rule
/// (e.g. a `VectorSimilar` rule). Without a matching rule the returned list
/// is empty — no live cosine computation is performed here.
/// Returns up to `limit` (default 10) neighbor entries.
fn tool_find_similar(db: &SharedDb, args: &Js) -> CallOutcome {
// Parse the optional mask once — it applies to both vector and edge paths.
// An invalid mask value (non-array or non-string element) fails closed.
let mask_keys: Option<Vec<String>> = if let Some(mask_val) = args.get("mask") {
match mask_val.as_array() {
Some(arr) => {
let mut ks: Vec<String> = Vec::with_capacity(arr.len());
for v in arr {
match v.as_str() {
Some(s) => ks.push(s.to_string()),
None => {
return CallOutcome::ToolErr("mask must be an array of strings".into())
}
}
}
Some(ks)
}
None => return CallOutcome::ToolErr("mask must be an array of strings".into()),
}
} else {
None
};
// When a `vector` array is provided, use the HNSW / brute-force vector
// similarity path instead of looking up pre-derived edges.
if let Some(vec_js) = args.get("vector").and_then(Js::as_array) {
let q: Vec<f64> = vec_js.iter().filter_map(|v| v.as_f64()).collect();
if q.is_empty() {
return CallOutcome::ToolErr("vector must be a non-empty array of numbers".into());
}
let field = args
.get("field")
.and_then(Js::as_str)
.unwrap_or("embedding");
let label_str = args.get("label").and_then(Js::as_str).unwrap_or("");
let label = if label_str.is_empty() {
None
} else {
Some(label_str)
};
let k = args
.get("k")
.and_then(Js::as_u64)
.map(|n| n as usize)
.unwrap_or(10);
let min = args
.get("min")
.and_then(Js::as_f64)
.unwrap_or(core_api::FIND_SIMILAR_DEFAULT_MIN);
let where_pred = match parse_where_arg(args) {
Ok(p) => p,
Err(e) => return CallOutcome::ToolErr(e),
};
let exact = match args.get("exact") {
None => false,
Some(v) => match v.as_bool() {
Some(b) => b,
None => return CallOutcome::ToolErr("exact must be a boolean".into()),
},
};
let exact = exact || where_pred.is_some();
let hits = {
let g = db.read();
let node_mask = mask_keys
.as_ref()
.map(|keys| NodeMask::from_keys(&*g, keys.iter().map(String::as_str)));
match g.find_similar_vector_filtered(
field,
label,
&q,
k,
min,
node_mask.as_ref(),
where_pred.as_ref(),
exact,
) {
Ok(h) => h,
Err(e) => return CallOutcome::ToolErr(graph_err_msg(e)),
}
};
let results: Vec<Js> = hits
.into_iter()
.map(|(key, score)| json!({ "key": key, "score": score }))
.collect();
return CallOutcome::ToolOk(json!({
"mode": "vector",
"field": field,
"label": label,
"k": k,
"min": min,
"results": results
}));
}
// Edge-traversal path: return neighbors connected by the given edge type.
let Some(key) = args.get("key").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing key (or provide vector for vector search)".into());
};
let edge_type = args
.get("edge_type")
.and_then(Js::as_str)
.unwrap_or("SIMILAR");
let limit = args
.get("limit")
.and_then(Js::as_u64)
.map(|n| n as usize)
.unwrap_or(10);
// When a mask is present, a hidden query key behaves identically to a
// nonexistent key — we do not confirm its existence.
if let Some(ref mask) = mask_keys {
let mask_set: std::collections::HashSet<&str> = mask.iter().map(String::as_str).collect();
if !mask_set.contains(key) {
return CallOutcome::ToolErr(graph_err_msg(GraphError::KeyNotFound {
key: key.into(),
}));
}
let out = {
let g = db.read();
g.node_edges(key)
};
return match out {
Ok(edges) => {
let similar: Vec<Js> = edges
.iter()
.filter(|e| e.edge_type == edge_type)
.filter(|e| {
// Keep only edges where the neighbor is also visible.
let neighbor_key = if e.src_key == key {
&e.dst_key
} else {
&e.src_key
};
mask_set.contains(neighbor_key.as_str())
})
.take(limit)
.map(|e| {
let neighbor_key = if e.src_key == key {
&e.dst_key
} else {
&e.src_key
};
let direction = if e.src_key == key { "out" } else { "in" };
json!({
"neighbor_key": neighbor_key,
"direction": direction,
"edge_type": e.edge_type,
"derived": e.derived,
})
})
.collect();
CallOutcome::ToolOk(json!({
"key": key,
"edge_type": edge_type,
"similar": similar
}))
}
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
};
}
let out = {
let g = db.read();
g.node_edges(key)
};
match out {
Ok(edges) => {
let similar: Vec<Js> = edges
.iter()
.filter(|e| e.edge_type == edge_type)
.take(limit)
.map(|e| {
let neighbor_key = if e.src_key == key {
&e.dst_key
} else {
&e.src_key
};
let direction = if e.src_key == key { "out" } else { "in" };
json!({
"neighbor_key": neighbor_key,
"direction": direction,
"edge_type": e.edge_type,
"derived": e.derived,
})
})
.collect();
CallOutcome::ToolOk(json!({
"key": key,
"edge_type": edge_type,
"similar": similar
}))
}
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
/// Exact cosine top-k among a caller key set. Self excluded. Never HNSW.
fn tool_pairwise_similar(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(keys_js) = args.get("keys").and_then(Js::as_array) else {
return CallOutcome::ToolErr("missing required field: keys".into());
};
let mut keys: Vec<String> = Vec::with_capacity(keys_js.len());
for v in keys_js {
match v.as_str() {
Some(s) => keys.push(s.to_string()),
None => return CallOutcome::ToolErr("keys must be an array of strings".into()),
}
}
let Some(field) = args.get("field").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing required field: field".into());
};
let k = args
.get("k")
.and_then(Js::as_u64)
.map(|n| n as usize)
.unwrap_or(10);
let min = args.get("min").and_then(Js::as_f64).unwrap_or(0.0);
let refs: Vec<&str> = keys.iter().map(String::as_str).collect();
let out = {
let g = db.read();
g.pairwise_similar(&refs, field, k, min)
};
match out {
Ok(pairs) => {
let results: Vec<Js> = pairs
.into_iter()
.map(|(key, neighbors)| {
json!({
"key": key,
"neighbors": neighbors
.into_iter()
.map(|(n, score)| json!({ "key": n, "score": score }))
.collect::<Vec<_>>(),
})
})
.collect();
CallOutcome::ToolOk(json!({
"field": field,
"k": k,
"min": min,
"results": results
}))
}
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
fn tool_hybrid_search(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(query_text) = args.get("query_text").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing required field: query_text".into());
};
let Some(text_field) = args.get("text_field").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing required field: text_field".into());
};
let vector_field = args
.get("vector_field")
.and_then(Js::as_str)
.unwrap_or("embedding");
let label = args.get("label").and_then(Js::as_str);
let k = args
.get("k")
.and_then(Js::as_u64)
.map(|n| n as usize)
.unwrap_or(10);
let query_vec: Vec<f64> = args
.get("vector")
.and_then(Js::as_array)
.map(|arr| arr.iter().filter_map(|v| v.as_f64()).collect())
.unwrap_or_default();
let hits = {
let g = db.read();
g.search_hybrid(text_field, query_text, vector_field, &query_vec, label, k)
};
let results: Vec<Js> = hits
.into_iter()
.map(|(key, score)| json!({ "key": key, "score": score }))
.collect();
CallOutcome::ToolOk(json!({
"query_text": query_text,
"text_field": text_field,
"vector_field": vector_field,
"label": label,
"k": k,
"results": results
}))
}
fn tool_node_history(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(key) = args.get("key").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing key".into());
};
let g = db.read();
let result = match g.node_history(key) {
Ok(e) => e,
Err(e) => return CallOutcome::ToolErr(graph_err_msg(e)),
};
CallOutcome::ToolOk(node_history_json(key, &result))
}
fn tool_edge_history(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(a) = args.get("a").and_then(Js::as_str).filter(|s| !s.is_empty()) else {
return CallOutcome::ToolErr("missing a".into());
};
let Some(b) = args.get("b").and_then(Js::as_str).filter(|s| !s.is_empty()) else {
return CallOutcome::ToolErr("missing b".into());
};
let result = {
let g = db.read();
g.edge_history(a, b)
};
match result {
Ok(hr) => CallOutcome::ToolOk(edge_history_result_json(a, b, &hr)),
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
fn tool_was_linked(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(a) = args.get("a").and_then(Js::as_str).filter(|s| !s.is_empty()) else {
return CallOutcome::ToolErr("missing a".into());
};
let Some(b) = args.get("b").and_then(Js::as_str).filter(|s| !s.is_empty()) else {
return CallOutcome::ToolErr("missing b".into());
};
let Some(edge_type) = args
.get("edge_type")
.and_then(Js::as_str)
.filter(|s| !s.is_empty())
else {
return CallOutcome::ToolErr("missing edge_type".into());
};
// Like `edges_at`, `at_commit` takes a date as well as an index. A caller
// asking "were they linked on 2026-06-19" should not have to find the
// commit themselves.
let at_commit = match args.get("at_commit") {
Some(Js::String(date)) => match db.read().resolve_date(date) {
Ok(c) => c,
Err(e) => return CallOutcome::ToolErr(graph_err_msg(e)),
},
Some(v) => match v.as_u64() {
Some(n) => n,
None => {
return CallOutcome::ToolErr(
"at_commit must be a non-negative commit index or an RFC 3339 date".into(),
)
}
},
None => return CallOutcome::ToolErr("missing or invalid at_commit".into()),
};
let result = {
let g = db.read();
g.was_linked(a, b, edge_type, at_commit)
};
match result {
Ok(linked) => CallOutcome::ToolOk(json!({
"a": a,
"b": b,
"edge_type": edge_type,
"at_commit": at_commit,
"linked": linked,
})),
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
fn tool_rename_node(db: &SharedDb, args: &Js) -> CallOutcome {
let Some(old_key) = args.get("old_key").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing old_key".into());
};
let Some(new_key) = args.get("new_key").and_then(Js::as_str) else {
return CallOutcome::ToolErr("missing new_key".into());
};
let mut g = db.write();
match g.rename_node(old_key, new_key) {
Ok(()) => CallOutcome::ToolOk(json!({
"ok": true,
"old_key": old_key,
"new_key": new_key,
})),
Err(e) => CallOutcome::ToolErr(graph_err_msg(e)),
}
}
fn parse_where_arg(args: &Js) -> std::result::Result<Option<PropPredicate>, String> {
let Some(w) = args.get("where") else {
return Ok(None);
};
if w.is_null() {
return Ok(None);
}
let pred: PropPredicate =
serde_json::from_value(w.clone()).map_err(|e| format!("where: {e}"))?;
pred.validate_named("where")?;
Ok(Some(pred))
}
pub(crate) fn graph_err_msg(e: GraphError) -> String {
match e {
GraphError::QueryError { detail } | GraphError::IngestError { detail } => detail,
other => other.to_string(),
}
}
fn initialize_result() -> Js {
json!({
"protocolVersion": "2024-11-05",
"capabilities": { "tools": {} },
"serverInfo": { "name": "mushroomdb", "version": env!("CARGO_PKG_VERSION") }
})
}
/// The prefix every graph tool's description carries.
///
/// A host that ranks tools by their description now has one signal that the
/// task tools are the ones to reach for first, and that everything
/// under this prefix is the lower-level surface beneath them.
const ADVANCED_PREFIX: &str = "Advanced: ";
/// The tools every store advertises, in the order it lists them.
///
/// A store with no repository in it used to be handed the code door's own task
/// tools, which answered from a code graph there was none of. Those tools are
/// gone, and a store `ingest-git` built is an entity store like any other, so
/// it is served this same listing. These are
/// the questions an entity graph *can* answer: what is there (`query` — now
/// with a `role`), why two things are associated, what is around a node, what
/// it is and what it is joined to, whether a link held at a commit and when it
/// changed, what is like it, and what was written down about it.
///
/// Listing order is ranking: a host that defers schemas shows this list in
/// order, so the two questions this door exists for come first.
///
/// `edges_at` and `what_if` are the two the first association benchmark run
/// showed missing: a run asked what a node's relationships were at a past
/// commit and spent twenty to sixty-seven turns replaying `edge_history` for
/// it, and had no way at all to ask what a change would do. They sit after
/// `was_linked`, which is the narrowest form of the same time question.
///
/// [`tools_list`] decides what is *advertised*, never what is answered: the
/// two served tools this leaves out, `explain` and `rename_node`, stay callable.
///
/// `upsert_entity`, `ingest_json` and `create_rule` were added in 0.6.12, and
/// the reason is the shape of a first session. A new install opens on an empty
/// store; the skill's first instruction is to offer to fill it, naming
/// `ingest_json` for a batch and `upsert_entity` for one; and the README sells
/// `upsert_entity -> create_rule -> find_similar -> explain_association` as the
/// minimal workflow. Two of those four were served and not advertised — and a
/// host builds its tool set from `tools/list`, so "served" did not help. The
/// listing described a store you could question but never populate.
///
/// They sit after the readers deliberately. Listing order is ranking and the
/// questions remain the point; writing is what a session does once, at the
/// start, before it has anything to ask.
pub const ASSOCIATION_TOOLS: [&str; 23] = [
"query",
"explain_association",
"neighborhood",
"node_info",
"node_edges",
"was_linked",
"edges_at",
"what_if",
"node_history",
"edge_history",
"find_similar",
"pairwise_similar",
"hybrid_search",
"remember",
"recall",
"upsert_entity",
"ingest_json",
"create_rule",
"stats",
"schema",
"analyze",
"suggest_rules",
"forget",
];
/// The tools `tools/list` advertises: the task tools, then the
/// graph tools with their descriptions prefixed.
///
/// `all` false — the default — lists [`ASSOCIATION_TOOLS`], **in the order it
/// names them**. The order is the point. A host that defers tool schemas makes
/// a model search for them, and the list it searches is read top-down.
///
/// `all` true lists every served tool, task tools first and graph tools after,
/// which is what `mushroomdb mcp --all-tools` runs and what the published
/// server card documents.
///
/// Either way every tool stays callable: the flag decides what is advertised,
/// not what is served.
fn tools_list(all: bool) -> Js {
let mut served: Vec<Js> = crate::mcp_tasks::task_tools();
for mut tool in graph_tools() {
if let Some(d) = tool.get("description").and_then(Js::as_str) {
let prefixed = format!("{ADVANCED_PREFIX}{d}");
tool["description"] = Js::String(prefixed);
}
served.push(tool);
}
if all {
return json!({ "tools": served });
}
let mut tools: Vec<Js> = Vec::with_capacity(ASSOCIATION_TOOLS.len());
for name in ASSOCIATION_TOOLS {
let Some(tool) = served
.iter()
.find(|t| t.get("name").and_then(Js::as_str) == Some(name))
else {
debug_assert!(false, "ASSOCIATION_TOOLS lists {name}, which is not served");
continue;
};
tools.push(tool.clone());
}
json!({ "tools": tools })
}
/// The fourteen graph tools, in the order they have always been listed, with
/// their descriptions unprefixed. [`tools_list`] adds the prefix.
fn graph_tools() -> Vec<Js> {
let Js::Array(tools) = json!([
{
"name": "query",
"description": "Who may see this, and anything else one pattern can answer — run a Cypher query (read or write) against the graph. Pass 'role' to answer as one of the store's roles: only the nodes that role may see, writes refused. 'mask' is the same restriction written out as an explicit key allow-list. Pass 'as_of' to answer from a past commit; it composes with 'role' or with 'mask'. Pass 'namespace' to answer from one namespace only. Cypher dialect: MATCH/WHERE/RETURN, CREATE, MERGE, SET, DELETE, with $named parameters in 'params'. A node's key and label read as properties (n.key, n.label) or as key(n)/labels(n). One MATCH takes comma-separated patterns that share variables — MATCH (t)-[:A]->(c), (t)-[:B]->(c) is the intersection of both, and count(DISTINCT t) after WITH counts each t once. WHERE takes STARTS WITH, ENDS WITH, CONTAINS, IN, and a list subscript (n.location[0]) — which is null when the index is out of range, the property is not a list, or the index is not an integer, so a subscript never errors and never matches.",
"inputSchema": {
"type": "object",
"properties": {
"cypher": { "type": "string", "description": "Cypher query text." },
"params": {
"type": "object",
"description": "Named JSON-scalar query parameters."
},
"mask": {
"type": "array",
"items": { "type": "string" },
"description": "Optional node key allow-list. When present, only these nodes are visible; write statements are rejected."
},
"role": {
"type": "string",
"description": "Answer as this role from the store's roles: only the nodes it may see. A role may also be narrowed by one property test (`status in [...]`), declared in the store's roles."
},
"as_of": {
"type": "integer",
"minimum": 0,
"description": "0-based WAL commit index: answer from the graph as it was at that commit. Composes with 'role' or with 'mask' — never both, which is refused as it is without 'as_of' — and whichever is passed is resolved against the graph as it was then. Deleting a node does not remove it from a role's past, and a role's 'keys' resolve to whichever node held the key at that commit. Writes are refused."
},
"namespace": {
"type": "string",
"description": "Answer only from this namespace. Intersects with 'role' and 'mask' — it can only narrow what they already allow. A role bound to namespaces honours them with no argument here. 'default' is the namespace of every node that names none; a name no node uses answers with nothing."
}
},
"required": ["cypher"]
}
},
{
"name": "ingest_json",
"description": "Fill the store from a batch — ingest a JSON array of objects as nodes of one label. Each object becomes a node; declare the rules that should link them with 'create_rule'.",
"inputSchema": {
"type": "object",
"properties": {
"label": { "type": "string" },
"rows_json": {
"type": "string",
"description": "JSON text of an array of objects."
},
"key_field": { "type": "string" },
"auto_fk_suffix": { "type": "string" },
"edges": {
"type": "array",
"description": "Optional user edges [{edge_type, src, dst}]."
},
"namespace": {
"type": "string",
"description": "Namespace for every node this call creates. Omitted means the 'default' namespace. A row that carries its own 'ns' must name the same namespace. A namespace is set at insert and cannot be changed afterwards."
}
},
"required": ["label", "rows_json"]
}
},
{
"name": "create_rule",
"description": "How should this kind of relationship be derived from now on — declare a rule (RuleDef JSON) and the engine maintains its edges as the data changes. Propose it and show the edges it would derive before creating one.",
"inputSchema": {
"type": "object",
"properties": {
"name": { "type": "string" },
"src_label": { "type": "string" },
"dst_label": { "type": "string" },
"predicate": { "type": "object" },
"edge_type": { "type": "string" },
"weight_prop": {
"type": ["string", "null"],
"description": "Edge property that stores the score (default: weight)."
},
"max_edges": { "type": ["integer", "null"] },
"namespace": {
"type": "string",
"description": "Scope the rule to one namespace: it sees only that namespace's nodes — source, via hop and destination — so every edge it derives stays inside. Omitted means a global rule, which is the only kind that may derive an edge across a boundary."
}
},
"required": ["name", "src_label", "dst_label", "predicate", "edge_type"]
}
},
{
"name": "explain",
"description": "Why are A and B related, as a raw array — the same rule-derived edges explain_association renders, for a caller that wants the JSON without asking.",
"inputSchema": {
"type": "object",
"properties": {
"a": { "type": "string", "minLength": 1 },
"b": { "type": "string", "minLength": 1 }
},
"required": ["a", "b"]
}
},
{
"name": "stats",
"description": "How big is this store — live node, edge and rule counts, plus `history_floor`, the oldest commit history still reaches (0 when nothing has been pruned). Pass 'role' or 'namespace' to also get `namespaces`, the namespaces that argument may see with a live-node count each; without either argument the roster is omitted entirely.",
"inputSchema": {
"type": "object",
"properties": {
"role": {
"type": "string",
"description": "Report only the namespaces this role may see. The store-wide counts beside them are unchanged."
},
"namespace": {
"type": "string",
"description": "Report only this namespace. Intersects with 'role'."
}
}
}
},
{
"name": "node_info",
"description": "What is K — its label and every property it holds.",
"inputSchema": {
"type": "object",
"properties": {
"key": { "type": "string" }
},
"required": ["key"]
}
},
{
"name": "upsert_entity",
"description": "Record what is now true about K — insert or update a node by key. If the key exists, updates the supplied properties atomically: every property is checked before any is written, so a refusal leaves the node unchanged. If not, creates a new node with the given label and properties. 'id' in 'props' is ignored on both paths: a created node stores 'id' as its key, and 'rename_node' is the only way to change it. Useful for agent memory: store or refresh an entity without checking existence first.",
"inputSchema": {
"type": "object",
"properties": {
"key": { "type": "string", "description": "Unique node key." },
"label": { "type": "string", "description": "Node label (required when creating a new entity)." },
"props": {
"type": "object",
"description": "Properties to set. Values must be scalars (string, number, bool) or arrays of scalars."
},
"namespace": {
"type": "string",
"description": "Namespace for a node this call creates. Omitted means the 'default' namespace. On a node that already exists, naming the namespace it is in is a no-op and naming another one is refused — a namespace is set at insert and cannot be changed."
},
"aliases": {
"type": "array",
"items": { "type": "string" },
"description": "Other names this entity goes by, kept as written in 'alias_keys': with the identity preset, one equal to a provisional stub's key (exact, case-sensitive) links that stub. The store derives 'aliases' from the key and the name only. Do not set either in props."
}
},
"required": ["key", "props"]
}
},
{
"name": "find_similar",
"description": "What is most like this — two modes: (1) Vector search — provide `vector` (and optionally `field`, `label`, `k`, `min`, `where`, `exact`) to find the k most similar nodes by cosine similarity using the HNSW index when available, brute-force otherwise. `where` is a property predicate (`{field, eq}` or `{field, in}`) and implies exact search. `exact` true skips HNSW. (2) Edge traversal — provide `key` (and optionally `edge_type`, `limit`) to return neighbors previously connected by a derived rule edge. Results from mode 2 come only from edges already derived by a VectorSimilar rule. Edge-traversal mode ignores `where` and `exact`. In both modes, the optional `mask` array limits visibility: hidden nodes never appear in results, and a hidden query key in edge mode behaves identically to a nonexistent key.",
"inputSchema": {
"type": "object",
"properties": {
"vector": {
"type": "array",
"items": { "type": "number" },
"description": "Query embedding vector for vector-similarity search. When present, vector-search mode is used and `key` is ignored."
},
"field": { "type": "string", "description": "Property field holding the embedding vectors (default: embedding). Used in vector-search mode." },
"label": { "type": "string", "description": "Restrict search to nodes with this label. Empty string means all labels. Used in vector-search mode." },
"k": { "type": "integer", "description": "Maximum results to return in vector-search mode (default: 10)." },
"min": { "type": "number", "description": "Minimum cosine similarity threshold in vector-search mode (default: 0.8 on every surface; before 0.7 HTTP and the Python binding defaulted to 0.0). Pass 0.0 to keep low-similarity hits." },
"mask": {
"type": "array",
"items": { "type": "string" },
"description": "Optional node key allow-list for vector-search mode. When present, only nodes whose key appears in this list are eligible for results. Hidden nodes are excluded before k-truncation. The beam widens until it has k visible hits, then falls back to an exhaustive masked scan at the same cap an exact VectorSimilar rule uses, so the result is not short while more visible hits exist. Unknown keys are silently ignored."
},
"where": {
"type": "object",
"description": "Optional property predicate for vector-search mode, same shape as visible_where: {\"field\": \"...\", \"eq\": value} or {\"field\": \"...\", \"in\": [values]}. Implies exact search (skips HNSW). Invalid predicates are a tool error. Edge-traversal mode ignores this."
},
"exact": {
"type": "boolean",
"description": "When true, vector-search mode uses exact GEMM brute force and does not consult HNSW. Default false. Edge-traversal mode ignores this."
},
"key": { "type": "string", "description": "Source node key for edge-traversal mode." },
"edge_type": { "type": "string", "description": "Edge type to filter by in edge-traversal mode (default: SIMILAR)." },
"limit": { "type": "integer", "description": "Maximum neighbors to return in edge-traversal mode (default: 10)." }
}
}
},
{
"name": "pairwise_similar",
"description": "Which of these are most like each other — exact cosine top-k among a caller key set. Self excluded. Never uses HNSW. Unknown keys, missing embeddings, zero-norm and wrong-dimension vectors are skipped. Duplicate keys collapse to first-seen order. Empty keys returns nothing. n above PAIRWISE_MAX_N is a tool error.",
"inputSchema": {
"type": "object",
"properties": {
"keys": {
"type": "array",
"items": { "type": "string" },
"description": "Node keys to score against each other."
},
"field": { "type": "string", "description": "Property field holding the embedding vectors." },
"k": { "type": "integer", "description": "Maximum neighbors per key (default: 10)." },
"min": { "type": "number", "description": "Minimum cosine similarity threshold (default: 0.0)." }
},
"required": ["keys", "field"]
}
},
{
"name": "hybrid_search",
"description": "What matches these words and this vector at once — Reciprocal Rank Fusion (RRF) over fulltext + vector results. Provide `query_text` and `text_field` for the fulltext leg. Optionally provide `vector` (embedding array) and `vector_field` (default: embedding) for the vector leg; omitting `vector` gives text-only ranking through the same RRF path. `label` restricts the vector search to nodes with that label (required for brute-force; omit to rely on HNSW rules). `k` controls result count (default: 10). RRF constant is fixed at 60; scores are 1/(60+rank) summed over lists a node appears in.",
"inputSchema": {
"type": "object",
"properties": {
"query_text": { "type": "string", "description": "Fulltext query string." },
"text_field": { "type": "string", "description": "Property field to search with fulltext." },
"vector": {
"type": "array",
"items": { "type": "number" },
"description": "Query embedding vector. Omit for text-only ranking."
},
"vector_field": { "type": "string", "description": "Property field holding embedding vectors (default: embedding)." },
"label": { "type": "string", "description": "Restrict vector search to nodes with this label. Required when relying on brute-force (no HNSW rule covers the field). If omitted, the vector leg always returns empty results (no rule-created HNSW index covers the unlabeled path); ranking is text-only in that case." },
"k": { "type": "integer", "description": "Maximum results to return (default: 10)." }
},
"required": ["query_text", "text_field"]
}
},
{
"name": "node_history",
"description": "What has happened to K — every recorded change to one node, newest last. Events include NodeInserted, PropSet, PropRemoved, EdgeAdded, EdgeRemoved, and NodeDeleted. The response includes `total_commits` (the horizon upper bound) and `horizon`, the oldest commit still retained; events before it are gone. History is WAL-scoped — pre-snapshot commits are not visible.",
"inputSchema": {
"type": "object",
"properties": {
"key": { "type": "string", "description": "Node key to look up." }
},
"required": ["key"]
}
},
{
"name": "edge_history",
"description": "When did A and B become linked, and when did it break — the full add/retract lifecycle for every edge between the two keys. Includes derived (rule-attributed) edges via DerivedEdgeAdded/DerivedEdgeRetracted WAL markers. The response includes `total_commits` (the horizon upper bound) and `horizon`, the oldest commit still retained; events before it are gone.",
"inputSchema": {
"type": "object",
"properties": {
"a": { "type": "string", "minLength": 1, "description": "First node key." },
"b": { "type": "string", "minLength": 1, "description": "Second node key." }
},
"required": ["a", "b"]
}
},
{
"name": "was_linked",
"description": "Were A and B linked at commit C — whether an edge of `edge_type` existed between the two keys (either direction) at that WAL commit. Returns an error when `at_commit` is outside the retained horizon (`horizon..total_commits`).",
"inputSchema": {
"type": "object",
"properties": {
"a": { "type": "string", "minLength": 1, "description": "First node key." },
"b": { "type": "string", "minLength": 1, "description": "Second node key." },
"edge_type": { "type": "string", "minLength": 1, "description": "Edge type to check." },
"at_commit": { "type": "integer", "minimum": 0, "description": "0-based WAL commit index to query." }
},
"required": ["a", "b", "edge_type", "at_commit"]
}
},
{
"name": "rename_node",
"description": "Rename K — the key changes and nothing else does. The dense id and all edges/properties remain stable. Returns 404 if `old_key` does not exist, 409 if `new_key` is already taken.",
"inputSchema": {
"type": "object",
"properties": {
"old_key": { "type": "string", "minLength": 1, "description": "Current node key." },
"new_key": { "type": "string", "minLength": 1, "description": "Desired new node key." }
},
"required": ["old_key", "new_key"]
}
}
]) else {
unreachable!("the literal above is an array")
};
tools
}
fn tool_ok(payload: Js) -> Js {
json!({
"content": [{ "type": "text", "text": payload.to_string() }]
})
}
/// A task tool's result: one text block, and no `structuredContent`.
///
/// The report used to ride along beside the digest, repeating it verbatim
/// under a `text` key. Nothing bound it — no task tool declares an
/// `outputSchema` — and it tripled the size of every reply, so a caller that
/// wants the numbers now asks for them with `json: true` and gets the report
/// *as* the text.
fn task_ok(text: &str) -> Js {
json!({
"content": [{ "type": "text", "text": text }]
})
}
fn tool_err(message: &str) -> Js {
json!({
"content": [{ "type": "text", "text": message }],
"isError": true
})
}
fn write_result(writer: &mut impl Write, id: Option<Js>, result: Js) -> io::Result<()> {
write_json(
writer,
&json!({
"jsonrpc": "2.0",
"id": id.unwrap_or(Js::Null),
"result": result
}),
)
}
fn write_error(
writer: &mut impl Write,
id: Option<Js>,
code: i64,
message: &str,
) -> io::Result<()> {
write_json(
writer,
&json!({
"jsonrpc": "2.0",
"id": id.unwrap_or(Js::Null),
"error": { "code": code, "message": message }
}),
)
}
fn write_json(writer: &mut impl Write, value: &Js) -> io::Result<()> {
let s = serde_json::to_string(value).map_err(io::Error::other)?;
writeln!(writer, "{s}")?;
writer.flush()
}
// ---------------------------------------------------------------------------
// Tests: MCP tool round-trips via stdio
// ---------------------------------------------------------------------------
#[cfg(test)]
mod tests {
use super::*;
use core_api::{AutoFk, IngestOptions, Predicate, RuleDef, Value};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
fn tmp_dir() -> PathBuf {
static SEQ: AtomicU64 = AtomicU64::new(0);
let n = SEQ.fetch_add(1, Ordering::Relaxed);
let d = std::env::temp_dir().join(format!("mcp-test-{}-{}", std::process::id(), n));
// These stores are never cleaned up, so a process id the OS hands out
// again lands on a previous run's data and every assertion about counts
// fails. `tests/mcp.rs::tmp` already clears its path for this reason.
let _ = std::fs::remove_dir_all(&d);
d
}
/// Open a SharedDb with two Person nodes and one derived SIMILAR edge.
fn demo_db() -> SharedDb {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
let opts = IngestOptions {
key_field: "id".into(),
auto_fk: AutoFk::Off,
};
// Two people with identical embeddings → will fire SIMILAR rule.
let people: Vec<BTreeMap<String, Value>> = vec![
[
("id", Value::Str("alice".into())),
("name", Value::Str("Alice".into())),
(
"emb",
Value::List(vec![Value::Float(1.0), Value::Float(0.0)]),
),
]
.into_iter()
.map(|(k, v)| (k.to_string(), v))
.collect(),
[
("id", Value::Str("bob".into())),
("name", Value::Str("Bob".into())),
(
"emb",
Value::List(vec![Value::Float(1.0), Value::Float(0.0)]),
),
]
.into_iter()
.map(|(k, v)| (k.to_string(), v))
.collect(),
];
g.ingest("Person", people, &opts).expect("ingest");
// Rule: VectorSimilar on emb → SIMILAR edge (cosine(ident,ident)=1.0 ≥ 0.9).
g.create_rule(RuleDef {
name: "sim_emb".into(),
src_label: "Person".into(),
dst_label: "Person".into(),
predicate: Predicate::VectorSimilar {
field: "emb".into(),
min: 0.9,
},
edge_type: "SIMILAR".into(),
weight_prop: Some("score".into()),
max_edges: None,
approximate: false,
via_label: None,
via_edge: None,
via_dir: None,
namespace: None,
})
.expect("rule");
}
db
}
fn roundtrip(db: &SharedDb, request: &str) -> Js {
roundtrip_with(db, false, request)
}
fn roundtrip_with(db: &SharedDb, all_tools: bool, request: &str) -> Js {
let input = format!("{request}\n");
let mut output = Vec::new();
run_mcp_stdio_with(db.clone(), None, all_tools, input.as_bytes(), &mut output)
.expect("mcp");
let s = std::str::from_utf8(&output).expect("utf8");
serde_json::from_str(s.trim()).expect("json response")
}
fn tool_call(db: &SharedDb, id: u64, tool: &str, args: Js) -> Js {
let req = json!({
"jsonrpc": "2.0",
"id": id,
"method": "tools/call",
"params": { "name": tool, "arguments": args }
});
roundtrip(db, &req.to_string())
}
/// Unwrap the `text` field from a successful tool response.
fn tool_text(resp: &Js) -> Js {
let text = resp["result"]["content"][0]["text"]
.as_str()
.expect("content[0].text");
serde_json::from_str(text).expect("tool text is json")
}
fn is_error(resp: &Js) -> bool {
resp["result"]["isError"].as_bool().unwrap_or(false)
}
fn tool_err_text(resp: &Js) -> String {
resp["result"]["content"][0]["text"]
.as_str()
.unwrap_or("")
.to_string()
}
// --- existing tools ---
#[test]
fn test_tools_list_includes_all_expected() {
let db = demo_db();
let resp = roundtrip_with(
&db,
true,
r#"{"jsonrpc":"2.0","id":1,"method":"tools/list"}"#,
);
let tools = resp["result"]["tools"].as_array().expect("tools array");
let names: Vec<&str> = tools
.iter()
.map(|t| t["name"].as_str().expect("name"))
.collect();
for expected in &[
// The task tools, first and in order.
"explain_association",
"node_edges",
"neighborhood",
"edges_at",
"what_if",
"recall",
"remember",
"schema",
"analyze",
"suggest_rules",
"forget",
// The fourteen graph tools.
"query",
"ingest_json",
"create_rule",
"explain",
"stats",
"node_info",
"upsert_entity",
"find_similar",
"pairwise_similar",
"hybrid_search",
"node_history",
"edge_history",
"was_linked",
"rename_node",
] {
assert!(names.contains(expected), "missing tool: {expected}");
}
assert_eq!(
names.len(),
25,
"expected exactly 25 tools, got {}",
names.len()
);
assert_eq!(
&names[..11],
[
"explain_association",
"node_edges",
"neighborhood",
"edges_at",
"what_if",
"recall",
"remember",
"schema",
"analyze",
"suggest_rules",
"forget",
],
"the task tools come first, in order"
);
assert_eq!(names[11], "query", "the graph tools follow them");
}
/// Binding: the default listing is the association tools, in [`ASSOCIATION_TOOLS`]
/// order, and nothing else.
#[test]
fn tools_list_defaults_to_the_association_tools_on_a_memory_store() {
let db = demo_db();
let resp = roundtrip(&db, r#"{"jsonrpc":"2.0","id":1,"method":"tools/list"}"#);
let names: Vec<&str> = resp["result"]["tools"]
.as_array()
.expect("tools array")
.iter()
.map(|t| t["name"].as_str().expect("name"))
.collect();
assert_eq!(names, ASSOCIATION_TOOLS.to_vec());
}
/// Binding: `pairwise_similar` is advertised on the memory surface
/// immediately after `find_similar`, and the listing is exactly `EXPECTED`.
#[test]
fn association_listing_includes_pairwise_similar_after_find_similar() {
let db = demo_db();
let resp = roundtrip(&db, r#"{"jsonrpc":"2.0","id":1,"method":"tools/list"}"#);
let names: Vec<&str> = resp["result"]["tools"]
.as_array()
.expect("tools array")
.iter()
.map(|t| t["name"].as_str().expect("name"))
.collect();
const EXPECTED: [&str; 23] = [
"query",
"explain_association",
"neighborhood",
"node_info",
"node_edges",
"was_linked",
"edges_at",
"what_if",
"node_history",
"edge_history",
"find_similar",
"pairwise_similar",
"hybrid_search",
"remember",
"recall",
"upsert_entity",
"ingest_json",
"create_rule",
"stats",
"schema",
"analyze",
"suggest_rules",
"forget",
];
assert_eq!(names, EXPECTED.to_vec());
assert_eq!(ASSOCIATION_TOOLS.as_slice(), EXPECTED.as_slice());
let find = names
.iter()
.position(|&n| n == "find_similar")
.expect("find_similar listed");
assert_eq!(names[find + 1], "pairwise_similar");
}
/// Binding: [`ASSOCIATION_TOOLS`] keeps the task tools an entity store can
/// answer with — the notes it wrote and the notes it kept — and every name
/// in it is served.
#[test]
fn the_association_surface_is_entity_tools_only() {
for kept in ["remember", "recall", "explain_association"] {
assert!(
ASSOCIATION_TOOLS.contains(&kept),
"{kept} answers on an entity graph and must be listed"
);
}
let served: Vec<String> = crate::mcp_tasks::task_tools()
.iter()
.chain(graph_tools().iter())
.filter_map(|t| t.get("name").and_then(Js::as_str))
.map(str::to_string)
.collect();
for name in ASSOCIATION_TOOLS {
assert!(
served.iter().any(|s| s == name),
"{name} is listed but not served"
);
}
}
#[test]
fn test_stats_returns_node_count() {
let db = demo_db();
let resp = tool_call(&db, 1, "stats", json!({}));
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["nodes_live"], 2);
}
/// Unscoped `stats` must not enumerate the store's namespaces. Asking with
/// neither `role` nor `namespace` gets the store-wide counts with the
/// roster key absent entirely — not an empty array, which would still
/// confirm the roster exists and invite a guess at its size.
#[test]
fn mcp_stats_unscoped_omits_namespace_roster() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
g.insert_node(
"Doc",
"a",
vec![("ns".into(), Value::Str("tenant-a".into()))],
)
.expect("insert a");
g.insert_node(
"Doc",
"b",
vec![("ns".into(), Value::Str("tenant-b".into()))],
)
.expect("insert b");
}
let unscoped = tool_text(&tool_call(&db, 1, "stats", json!({})));
assert!(
unscoped.get("namespaces").is_none(),
"unscoped stats must omit the roster entirely, not send an empty \
array: {unscoped}"
);
assert_eq!(
unscoped["nodes_live"], 2,
"the store-wide counts beside the roster are unchanged"
);
let scoped = tool_text(&tool_call(
&db,
2,
"stats",
json!({"namespace": "tenant-a"}),
));
let names: Vec<&str> = scoped["namespaces"]
.as_array()
.expect("a scoped call still carries the roster it may see")
.iter()
.map(|n| n["name"].as_str().expect("name"))
.collect();
assert_eq!(
names,
["tenant-a"],
"a call that names a namespace sees that one and no other"
);
}
/// A rule whose corpus is too large to index in one commit must not come
/// back as a bare "ok": the caller would go straight to querying edges that
/// do not exist yet.
#[test]
fn create_rule_reports_a_build_it_could_not_finish() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
for i in 0..300usize {
const D: usize = 32;
let axis = (i / 10) % D;
let mut xs = vec![0.0f64; D];
xs[axis] = 1.0;
xs[(axis + 1) % D] = (i % 10) as f64 * 0.001;
g.insert_node(
"V",
&format!("v{i}"),
vec![(
"emb".into(),
Value::List(xs.into_iter().map(Value::Float).collect()),
)],
)
.expect("insert");
}
g.set_hnsw_build_batch(Some(64));
}
let args = json!({
"name": "sim",
"src_label": "V",
"dst_label": "V",
"predicate": {"VectorSimilar": {"field": "emb", "min": 0.9}},
"edge_type": "SIM",
"weight_prop": null,
"max_edges": null,
"approximate": true
});
let resp = tool_call(&db, 1, "create_rule", args);
assert!(!is_error(&resp), "{resp}");
let result = tool_text(&resp);
assert_eq!(result["name"], json!("sim"));
assert_eq!(result["building"], json!({"indexed": 64, "total": 300}));
let note = result["note"].as_str().expect("a note explaining the wait");
assert!(
note.contains("derives no edges until it finishes") && note.contains("build-index"),
"the note must say the edges are not there yet and how to finish: {note}"
);
// `stats` carries the same progress.
let stats = tool_text(&tool_call(&db, 2, "stats", json!({})));
let rule = stats["rules"]
.as_array()
.expect("rules")
.iter()
.find(|r| r["name"] == "sim")
.expect("the rule is installed while it builds");
assert_eq!(rule["edges"], json!(0));
assert_eq!(
rule["building"],
json!({"rule": "sim", "indexed": 64, "total": 300})
);
// Finished, the report is the plain one again.
while !db.write().pump_index_build().expect("pump").is_empty() {}
let stats = tool_text(&tool_call(&db, 3, "stats", json!({})));
let rule = stats["rules"]
.as_array()
.expect("rules")
.iter()
.find(|r| r["name"] == "sim")
.expect("rule");
assert!(rule.get("building").is_none(), "{rule}");
assert!(rule["edges"].as_u64().expect("edges") > 0);
}
#[test]
fn test_query_runs_cypher() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"query",
json!({ "cypher": "MATCH (n:Person) RETURN n.name ORDER BY n.name" }),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
// columns + 2 rows
assert_eq!(result["columns"], json!(["n.name"]));
assert_eq!(result["rows"].as_array().map(|r| r.len()), Some(2));
}
#[test]
fn test_query_create_is_a_write() {
let db = SharedDb::open(&tmp_dir()).expect("open");
let resp = tool_call(
&db,
1,
"query",
json!({ "cypher": "CREATE (n:L {id: 'k'}) RETURN n" }),
);
assert!(
!is_error(&resp),
"CREATE via MCP query must succeed: {resp}"
);
let stats = tool_text(&tool_call(&db, 2, "stats", json!({})));
assert_eq!(stats["nodes_live"], 1);
}
#[test]
fn test_ingest_json_inserts_nodes() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"ingest_json",
json!({
"label": "Person",
"rows_json": r#"[{"id":"carol","name":"Carol"}]"#,
"key_field": "id"
}),
);
assert!(!is_error(&resp));
// Verify node visible via stats
let stats = tool_text(&tool_call(&db, 2, "stats", json!({})));
assert_eq!(stats["nodes_live"], 3);
}
#[test]
fn test_node_info_returns_props() {
let db = demo_db();
let resp = tool_call(&db, 1, "node_info", json!({ "key": "alice" }));
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["key"], "alice");
assert_eq!(result["label"], "Person");
assert_eq!(result["props"]["name"], "Alice");
}
/// Binding: `node_edges` groups by edge type and names the rule and score
/// behind each derived edge, in the report as in the digest.
#[test]
fn test_node_edges_returns_edges() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"node_edges",
json!({ "key": "alice", "json": true }),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["key"], "alice");
let types = result["types"].as_array().expect("types");
assert!(
!types.is_empty(),
"alice should have at least one edge type"
);
let similar = types
.iter()
.find(|t| t["edge_type"] == "SIMILAR")
.expect("the rule's edge type");
// A symmetric rule derives the edge both ways, and both are listed
// with the direction that tells them apart.
assert_eq!(similar["count"], json!(2));
let edges = similar["edges"].as_array().expect("edges");
let dirs: Vec<&str> = edges
.iter()
.map(|e| e["direction"].as_str().expect("direction"))
.collect();
assert!(dirs.contains(&"out") && dirs.contains(&"in"), "{similar}");
for edge in edges {
assert_eq!(edge["other"], json!("bob"));
assert_eq!(edge["derived"], json!(true));
assert_eq!(edge["rule"], json!("sim_emb"));
assert_eq!(edge["score"], json!(1.0));
assert!(
edge["predicate"]
.as_str()
.unwrap_or("")
.contains("vector_similar"),
"the predicate travels with the edge: {edge}"
);
}
}
/// Binding: a depth-1 `neighborhood` is the same relationship listing, and
/// anything deeper is still the traversal table.
#[test]
fn test_neighborhood_traverses_one_hop() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"neighborhood",
json!({ "key": "alice", "depth": 1, "json": true }),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["key"], "alice");
assert!(result["types"].as_array().is_some(), "{result}");
let deep = tool_call(
&db,
2,
"neighborhood",
json!({ "key": "alice", "depth": 2 }),
);
assert!(!is_error(&deep));
let table = tool_text(&deep);
assert_eq!(table["columns"], json!(["key", "label", "depth"]));
assert!(table["rows"].as_array().is_some());
}
#[test]
fn test_explain_returns_rule_info() {
let db = demo_db();
let resp = tool_call(&db, 1, "explain", json!({ "a": "alice", "b": "bob" }));
assert!(!is_error(&resp));
let result = tool_text(&resp);
let arr = result.as_array().expect("explain returns array");
assert!(!arr.is_empty(), "expected at least one explanation");
assert_eq!(arr[0]["rule"], "sim_emb");
}
#[test]
fn test_create_rule_backfills() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
let opts = IngestOptions {
key_field: "id".into(),
auto_fk: AutoFk::Off,
};
let rows: Vec<BTreeMap<String, Value>> = vec![
[
("id", Value::Str("x".into())),
("tag", Value::Str("a".into())),
]
.into_iter()
.map(|(k, v)| (k.to_string(), v))
.collect(),
[
("id", Value::Str("y".into())),
("tag", Value::Str("a".into())),
]
.into_iter()
.map(|(k, v)| (k.to_string(), v))
.collect(),
];
g.ingest("Item", rows, &opts).expect("ingest");
}
let resp = tool_call(
&db,
1,
"create_rule",
json!({
"name": "same_tag",
"src_label": "Item",
"dst_label": "Item",
"predicate": { "FieldEqual": { "field": "tag" } },
"edge_type": "SAME_TAG"
}),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["ok"], true);
// Derived edges should now exist.
let edges_resp = tool_call(&db, 2, "node_edges", json!({ "key": "x", "json": true }));
let edges_result = tool_text(&edges_resp);
let types = edges_result["types"].as_array().expect("types");
assert!(
types.iter().any(|t| t["edge_type"] == "SAME_TAG"),
"SAME_TAG edge not found after create_rule"
);
}
// --- new tools ---
#[test]
fn test_upsert_entity_creates_new_node() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"upsert_entity",
json!({
"key": "carol",
"label": "Person",
"props": { "name": "Carol", "age": 30 }
}),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["ok"], true);
assert_eq!(result["created"], true);
assert_eq!(result["key"], "carol");
// Verify node exists
let info = tool_text(&tool_call(&db, 2, "node_info", json!({ "key": "carol" })));
assert_eq!(info["props"]["name"], "Carol");
}
#[test]
fn test_upsert_entity_updates_existing_node() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"upsert_entity",
json!({
"key": "alice",
"props": { "name": "Alice Updated" }
}),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["ok"], true);
assert_eq!(result["created"], false);
assert_eq!(result["updated_fields"], 1);
// Verify prop changed
let info = tool_text(&tool_call(&db, 2, "node_info", json!({ "key": "alice" })));
assert_eq!(info["props"]["name"], "Alice Updated");
}
#[test]
fn test_upsert_entity_missing_label_on_create_is_error() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"upsert_entity",
json!({ "key": "new-node", "props": { "x": 1 } }),
);
assert!(is_error(&resp), "should error without label for new node");
}
#[test]
fn test_pairwise_similar_excludes_self() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"pairwise_similar",
json!({
"keys": ["alice", "bob"],
"field": "emb",
"k": 10,
"min": 0.0
}),
);
assert!(
!is_error(&resp),
"pairwise_similar must not error: {resp:?}"
);
let result = tool_text(&resp);
let results = result["results"].as_array().expect("results");
assert_eq!(results.len(), 2);
for row in results {
let key = row["key"].as_str().expect("key");
let neighbors = row["neighbors"].as_array().expect("neighbors");
assert!(
neighbors.iter().all(|n| n["key"].as_str() != Some(key)),
"self must be excluded: {row}"
);
assert!(!neighbors.is_empty(), "alice/bob are identical: {row}");
}
}
#[test]
fn test_find_similar_returns_similar_edges() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"find_similar",
json!({ "key": "alice", "edge_type": "SIMILAR" }),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
assert_eq!(result["key"], "alice");
assert_eq!(result["edge_type"], "SIMILAR");
let similar = result["similar"].as_array().expect("similar array");
assert!(!similar.is_empty(), "expected SIMILAR neighbors for alice");
assert_eq!(similar[0]["neighbor_key"], "bob");
}
#[test]
fn test_find_similar_limit_respected() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"find_similar",
json!({ "key": "alice", "edge_type": "SIMILAR", "limit": 0 }),
);
assert!(!is_error(&resp));
let result = tool_text(&resp);
let similar = result["similar"].as_array().expect("similar array");
assert_eq!(similar.len(), 0);
}
/// When `min` is omitted from a vector-mode find_similar call, the server
/// must apply the spec default of 0.8. A node whose cosine similarity to
/// the query is 0.0 (orthogonal) must not appear in the results.
#[test]
fn test_find_similar_vector_default_min_is_0_8() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
// close: [1,0] → cosine 1.0 with query [1,0] (above 0.8)
g.insert_node(
"Item",
"close",
vec![(
"emb".into(),
Value::List(vec![Value::Float(1.0), Value::Float(0.0)]),
)],
)
.unwrap();
// far: [0,1] → cosine 0.0 with query [1,0] (below 0.8, must be excluded)
g.insert_node(
"Item",
"far",
vec![(
"emb".into(),
Value::List(vec![Value::Float(0.0), Value::Float(1.0)]),
)],
)
.unwrap();
// mid: [0.6,0.8] → cosine 0.6 with query [1,0]: above the old HTTP
// and Python default (0.0), below this one (0.8).
g.insert_node(
"Item",
"mid",
vec![(
"emb".into(),
Value::List(vec![Value::Float(0.6), Value::Float(0.8)]),
)],
)
.unwrap();
}
// No `min` in the request — must default to 0.8.
let resp = tool_call(
&db,
1,
"find_similar",
json!({
"vector": [1.0, 0.0],
"field": "emb",
"label": "Item",
"k": 10
}),
);
assert!(!is_error(&resp), "vector search must not error");
let result = tool_text(&resp);
let results = result["results"].as_array().expect("results array");
let keys: Vec<&str> = results.iter().filter_map(|r| r["key"].as_str()).collect();
assert_eq!(
keys,
vec!["close"],
"only the hit at or above the default survives: mid (0.6) and far (0.0) do not"
);
assert_eq!(
result["min"],
json!(core_api::FIND_SIMILAR_DEFAULT_MIN),
"the reply echoes the floor it applied"
);
// Naming it still reaches everything, in score order: the argument is
// read, not replaced by the default.
let resp = tool_call(
&db,
2,
"find_similar",
json!({
"vector": [1.0, 0.0],
"field": "emb",
"label": "Item",
"k": 10,
"min": 0.0
}),
);
assert!(!is_error(&resp), "vector search must not error");
let result = tool_text(&resp);
let results = result["results"].as_array().expect("results array");
let keys: Vec<&str> = results.iter().filter_map(|r| r["key"].as_str()).collect();
assert_eq!(
keys,
vec!["close", "mid", "far"],
"an explicit min of 0.0 keeps the hits the default drops"
);
assert_eq!(
result["min"],
json!(0.0),
"the reply echoes the named floor"
);
}
/// `find_similar` with `mask` must exclude hidden node keys from results.
#[test]
fn test_find_similar_vector_mask_excludes_hidden() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
// visible: [1,0] — should appear in results.
g.insert_node(
"Item",
"visible",
vec![(
"emb".into(),
Value::List(vec![Value::Float(1.0), Value::Float(0.0)]),
)],
)
.unwrap();
// hidden: [1,0] — same direction as query but must not appear.
g.insert_node(
"Item",
"hidden",
vec![(
"emb".into(),
Value::List(vec![Value::Float(1.0), Value::Float(0.0)]),
)],
)
.unwrap();
}
let resp = tool_call(
&db,
1,
"find_similar",
json!({
"vector": [1.0, 0.0],
"field": "emb",
"label": "Item",
"k": 10,
"min": 0.0,
"mask": ["visible"]
}),
);
assert!(!is_error(&resp), "masked vector search must not error");
let result = tool_text(&resp);
let results = result["results"].as_array().expect("results array");
let keys: Vec<&str> = results.iter().filter_map(|r| r["key"].as_str()).collect();
assert!(
keys.contains(&"visible"),
"visible node must appear in masked results"
);
assert!(
!keys.contains(&"hidden"),
"hidden node must be excluded by mask"
);
}
/// `find_similar` with `mask` — bad mask value returns a tool error.
#[test]
fn test_find_similar_vector_mask_bad_type_is_error() {
let db = SharedDb::open(&tmp_dir()).expect("open");
let resp = tool_call(
&db,
1,
"find_similar",
json!({
"vector": [1.0, 0.0],
"field": "emb",
"k": 5,
"mask": [42]
}),
);
assert!(
is_error(&resp),
"non-string mask element must produce a tool error"
);
}
/// Vector-mode `where` eq filters to matching nodes. Default `min` stays 0.8.
#[test]
fn test_find_similar_vector_where_eq() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
g.insert_node(
"Document",
"in-scope",
vec![
(
"emb".into(),
Value::List(vec![Value::Float(1.0), Value::Float(0.0)]),
),
("resource_scope_id".into(), Value::Str("a".into())),
],
)
.unwrap();
g.insert_node(
"Document",
"out-scope",
vec![
(
"emb".into(),
Value::List(vec![Value::Float(1.0), Value::Float(0.0)]),
),
("resource_scope_id".into(), Value::Str("b".into())),
],
)
.unwrap();
}
let resp = tool_call(
&db,
1,
"find_similar",
json!({
"vector": [1.0, 0.0],
"field": "emb",
"label": "Document",
"k": 10,
"min": 0.0,
"where": { "field": "resource_scope_id", "eq": "a" }
}),
);
assert!(!is_error(&resp), "where eq must not error: {resp:?}");
let result = tool_text(&resp);
let keys: Vec<&str> = result["results"]
.as_array()
.expect("results")
.iter()
.filter_map(|r| r["key"].as_str())
.collect();
assert_eq!(keys, vec!["in-scope"]);
}
#[test]
fn test_find_similar_vector_where_invalid_is_error() {
let db = SharedDb::open(&tmp_dir()).expect("open");
let resp = tool_call(
&db,
1,
"find_similar",
json!({
"vector": [1.0, 0.0],
"field": "emb",
"where": { "field": "resource_scope_id", "eq": "a", "in": ["b"] }
}),
);
assert!(is_error(&resp), "invalid where must be a tool error");
let msg = format!("{resp:?}");
assert!(
msg.contains("where"),
"tool error must name where, got {msg}"
);
}
/// Edge-traversal mode ignores `where` and `exact`.
#[test]
fn test_find_similar_edge_ignores_where_and_exact() {
let db = demo_db();
let resp = tool_call(
&db,
1,
"find_similar",
json!({
"key": "alice",
"edge_type": "SIMILAR",
"where": { "field": "x", "eq": "y", "in": ["z"] },
"exact": true
}),
);
assert!(
!is_error(&resp),
"edge mode must ignore invalid where: {resp:?}"
);
}
/// Edge-traversal mode with `mask` must exclude hidden neighbors.
#[test]
fn test_find_similar_edge_mask_excludes_hidden_neighbor() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
g.insert_node("P", "alice", vec![]).unwrap();
g.insert_node("P", "bob", vec![]).unwrap(); // visible
g.insert_node("P", "carol", vec![]).unwrap(); // hidden
g.insert_edge("KNOWS", "alice", "bob").unwrap();
g.insert_edge("KNOWS", "alice", "carol").unwrap();
}
// Mask: alice and bob visible; carol hidden.
let resp = tool_call(
&db,
1,
"find_similar",
json!({
"key": "alice",
"edge_type": "KNOWS",
"mask": ["alice", "bob"]
}),
);
assert!(!is_error(&resp), "masked edge search must not error");
let result = tool_text(&resp);
let similar = result["similar"].as_array().expect("similar array");
let neighbors: Vec<&str> = similar
.iter()
.filter_map(|e| e["neighbor_key"].as_str())
.collect();
assert!(neighbors.contains(&"bob"), "bob (visible) must appear");
assert!(
!neighbors.contains(&"carol"),
"carol (hidden) must be excluded"
);
}
/// Edge-traversal mode with `mask`: a hidden query key must not reveal
/// its existence — response must be a tool error identical to a nonexistent key.
#[test]
fn test_find_similar_edge_mask_hidden_key_is_not_found() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
g.insert_node("P", "alice", vec![]).unwrap();
g.insert_node("P", "bob", vec![]).unwrap();
}
// alice exists but is not in the mask — must look like not-found.
let resp_masked = tool_call(
&db,
1,
"find_similar",
json!({ "key": "alice", "edge_type": "KNOWS", "mask": ["bob"] }),
);
// ghost never exists — use as the reference for "not found".
let resp_ghost = tool_call(
&db,
2,
"find_similar",
json!({ "key": "ghost", "edge_type": "KNOWS" }),
);
assert!(
is_error(&resp_masked),
"hidden query key must produce a tool error"
);
assert!(
is_error(&resp_ghost),
"nonexistent key must produce a tool error"
);
// Both errors must carry the same shape (both are key-not-found).
assert_eq!(
tool_err_text(&resp_masked).contains("alice"),
tool_err_text(&resp_ghost).contains("ghost"),
"error messages should follow same not-found template"
);
}
/// Binding: `explain_association` now answers in prose, and the report
/// behind it — what `json: true` returns — is still `explain`'s array,
/// with one `evidence` object added per relationship and every other
/// field unchanged.
#[test]
fn test_explain_association_same_as_explain() {
let db = demo_db();
let explain = tool_text(&tool_call(
&db,
1,
"explain",
json!({ "a": "alice", "b": "bob" }),
));
let assoc = tool_text(&tool_call(
&db,
2,
"explain_association",
json!({ "a": "alice", "b": "bob", "json": true }),
));
let explain: Vec<Js> = serde_json::from_value(explain).expect("explain array");
let mut assoc: Vec<Js> = serde_json::from_value(assoc).expect("assoc array");
for row in &mut assoc {
let ev = row
.as_object_mut()
.expect("object")
.remove("evidence")
.expect("every derived edge carries its evidence");
assert!(
ev["similarity"].is_number(),
"a vector_similar edge reports the cosine it scored: {ev}"
);
}
assert_eq!(explain, assoc, "evidence is the only addition");
let prose = tool_call(
&db,
3,
"explain_association",
json!({ "a": "alice", "b": "bob" }),
);
let text = prose["result"]["content"][0]["text"]
.as_str()
.expect("text content");
assert!(
text.contains("mushroomdb explain — alice ↔ bob:"),
"the default reply is the digest: {text}"
);
}
// ── history tools ──────────────────────────────────────────────────────────
/// `edge_history` must return the derived-edge lifecycle (Added event with
/// rule attribution) and include the `total_commits` horizon field.
#[test]
fn test_edge_history_returns_derived_lifecycle_with_rule() {
let db = demo_db(); // alice+bob + sim_emb rule → SIMILAR derived edge
let resp = tool_call(&db, 1, "edge_history", json!({ "a": "alice", "b": "bob" }));
assert!(!is_error(&resp), "edge_history must not error: {resp}");
let result = tool_text(&resp);
// Must carry horizon metadata.
let total = result["total_commits"].as_u64().expect("total_commits");
assert!(total > 0, "total_commits must be > 0 after ingest + rule");
// Must have at least one event (the SIMILAR derived-edge addition).
let events = result["events"].as_array().expect("events array");
assert!(!events.is_empty(), "expected at least one edge event");
// At least one event must be Added with a non-null rule (derived edge).
let derived_added = events
.iter()
.any(|ev| ev["event"].as_str() == Some("Added") && !ev["rule"].is_null());
assert!(
derived_added,
"expected a derived Added event with rule attribution: {events:?}"
);
}
/// `was_linked` must return `true` for an edge that was active at the given commit,
/// and the response must include the echo fields.
#[test]
fn test_was_linked_at_valid_commit() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
let opts = IngestOptions {
key_field: "id".into(),
auto_fk: AutoFk::Off,
};
let rows: Vec<BTreeMap<String, Value>> = vec![
[("id", Value::Str("x".into()))]
.into_iter()
.map(|(k, v)| (k.to_string(), v))
.collect(),
[("id", Value::Str("y".into()))]
.into_iter()
.map(|(k, v)| (k.to_string(), v))
.collect(),
];
g.ingest("N", rows, &opts).expect("ingest");
g.insert_edge("LINK", "x", "y").expect("edge");
}
// There are now at least 2 commits (ingest + edge). Check at the last one.
let g = db.read();
let total = g.wal_total_commits().expect("wal_total_commits");
drop(g);
let resp = tool_call(
&db,
1,
"was_linked",
json!({ "a": "x", "b": "y", "edge_type": "LINK", "at_commit": total - 1 }),
);
assert!(!is_error(&resp), "was_linked must not error: {resp}");
let result = tool_text(&resp);
assert_eq!(result["linked"], true);
assert_eq!(result["a"], "x");
assert_eq!(result["edge_type"], "LINK");
}
/// `was_linked` with an out-of-horizon commit must return a tool error (not
/// a protocol error), and the error message must mention the commit range.
#[test]
fn test_was_linked_out_of_horizon_returns_tool_error() {
let db = SharedDb::open(&tmp_dir()).expect("open");
{
let mut g = db.write();
g.insert_node("N", "a", vec![]).expect("node a");
g.insert_node("N", "b", vec![]).expect("node b");
}
// Commit 999 is well beyond the WAL.
let resp = tool_call(
&db,
1,
"was_linked",
json!({ "a": "a", "b": "b", "edge_type": "X", "at_commit": 999 }),
);
// isError true = tool-level error (not a JSON-RPC protocol error).
assert!(
is_error(&resp),
"out-of-range commit must be a tool error: {resp}"
);
let text = resp["result"]["content"][0]["text"].as_str().expect("text");
assert!(
text.contains("out of range") || text.contains("range"),
"error must mention range: {text}"
);
}
/// `node_history` tool must return the node's WAL history and the
/// `total_commits` horizon field.
#[test]
fn test_node_history_via_mcp() {
let db = demo_db(); // alice + bob, with a SIMILAR rule
let resp = tool_call(&db, 1, "node_history", json!({ "key": "alice" }));
assert!(!is_error(&resp), "node_history must not error: {resp}");
let result = tool_text(&resp);
assert_eq!(result["key"], "alice");
let total = result["total_commits"].as_u64().expect("total_commits");
assert!(total > 0, "total_commits must be > 0");
let history = result["history"].as_array().expect("history array");
assert!(
!history.is_empty(),
"alice should have at least one history entry"
);
// First event should be a NodeInserted.
let first_change = &history[0]["change"];
assert_eq!(first_change["type"], "NodeInserted");
assert_eq!(first_change["label"], "Person");
}
}