mod cli;
mod config;
mod workspace;
use std::collections::BTreeMap;
use std::io::{self, BufRead, Write};
use std::path::{Path, PathBuf};
use std::process::ExitCode;
use std::time::{SystemTime, UNIX_EPOCH};
use clap::Parser;
use plugmem_host::{
Database, ExportedFact, FactId, HostError, LinkInput, MaintenanceMode, MaintenanceOptions,
ReadOnlyDatabase, RecallQuery, RecallResult, RememberInput, RememberOutcome, Settings, Stats,
UnlinkInput, VALID_TO_OPEN,
};
use serde_json::json;
use crate::cli::{Cli, Command, HelpTopic, MaintainMode};
use crate::config::read_batch_size;
pub(crate) const ENV_DB: &str = "PLUGMEM_DB";
pub(crate) const DEFAULT_DB: &str = "plugmem.db";
#[derive(Debug)]
pub(crate) enum CliError {
Host(HostError),
Usage(String),
}
impl From<HostError> for CliError {
fn from(e: HostError) -> Self {
CliError::Host(e)
}
}
pub(crate) fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
pub fn run() -> ExitCode {
let stdout = io::stdout();
ExitCode::from(run_parsed(Cli::parse(), &mut stdout.lock()))
}
fn run_parsed(cli: Cli, out: &mut impl Write) -> u8 {
if let Command::Help { topic } = &cli.command {
return execute_help(topic, cli.json, out);
}
let table = match plugmem_host::read_config(cli.config.as_deref()) {
Ok(t) => t,
Err(e) => return report_err(&e.into()),
};
let cfg_batch_size = read_batch_size(table.as_ref());
let mut settings = match Settings::from_table(table.as_ref()) {
Ok(s) => s,
Err(e) => return report_err(&e.into()),
};
let root = workspace::resolve_root(cli.workspace.as_deref(), &settings);
if let Command::Workspace { command } = &cli.command {
return match workspace::execute(command, root, settings, cli.json, out) {
Ok(code) => code,
Err(e) => {
let _ = out.flush();
report_err(&e)
}
};
}
let path = resolve_db_path(
cli.db.as_deref(),
settings.database_path.as_deref(),
root.as_ref(),
);
if let Err(e) = workspace::ensure_dir(&path, root.as_ref()) {
return report_err(&e);
}
match &cli.command {
Command::Recover { dst } => return do_recover(&path, dst, &settings, cli.json, out),
Command::Scrub => return do_scrub(&path, &settings, cli.json, out),
Command::Repl { read_only: true } => {
return run_repl_ro(&path, settings, cli.json, io::stdin().lock(), out);
}
Command::Repl { read_only: false } => {
return run_repl(&path, settings, cli.json, io::stdin().lock(), out);
}
_ => {}
}
let readonly_ok = matches!(
&cli.command,
Command::Show { .. }
| Command::Stats
| Command::Export
| Command::Verify
| Command::Recall { .. }
);
if readonly_ok {
let recall_vector = match embed_recall_query(&mut settings, &cli.command) {
Ok(v) => v,
Err(e) => return report_err(&e),
};
match Database::open_readonly(&path, settings.config.clone()) {
Ok(ro) => {
return execute_ro(&ro, &cli.command, recall_vector.as_deref(), cli.json, out);
}
Err(HostError::Locked { path }) => return report_locked(&path),
Err(_) => {}
}
}
let db = match settings.open(&path) {
Ok(db) => db,
Err(HostError::Locked { path }) => return report_locked(&path),
Err(e) => return report_err(&CliError::Host(e)),
};
if let Command::Import { file, batch } = &cli.command {
let batch_size = batch
.or(cfg_batch_size.map(|n| n as usize))
.unwrap_or(DEFAULT_IMPORT_BATCH)
.max(1);
return match do_import(&db, now_ms(), file, batch_size, out) {
Ok(n) => {
if cli.json {
writeln!(out, "{}", json!({ "imported": n })).ok();
} else {
writeln!(out, "imported {n} facts").ok();
}
0
}
Err(e) => {
let _ = out.flush();
report_err(&e)
}
};
}
match execute(&db, &cli.command, cli.json, now_ms(), out) {
Ok(code) => code,
Err(e) => {
let _ = out.flush();
report_err(&e)
}
}
}
const DEFAULT_IMPORT_BATCH: usize = 128;
fn write_err(out: &mut impl Write, e: &CliError) {
let _ = match e {
CliError::Usage(msg) => writeln!(out, "plugmem: {msg}"),
CliError::Host(err) => {
writeln!(out, "plugmem: {err}").and_then(|()| match err.capacity_hint() {
Some(hint) => writeln!(out, "plugmem: {hint}"),
None => Ok(()),
})
}
};
}
fn report_err(e: &CliError) -> u8 {
write_err(&mut std::io::stderr(), e);
2
}
fn report_locked(path: &std::path::Path) -> u8 {
eprintln!(
"plugmem: database is locked by another process: {}",
path.display()
);
1
}
fn resolve_db_path(
flag: Option<&str>,
config_path: Option<&std::path::Path>,
root: Option<&PathBuf>,
) -> PathBuf {
flag.map(|value| workspace::resolve_target(value, root))
.or_else(|| {
std::env::var_os(ENV_DB).map(|v| workspace::resolve_target(&v.to_string_lossy(), root))
})
.or_else(|| config_path.map(PathBuf::from))
.or_else(plugmem_host::default_database_path)
.unwrap_or_else(|| PathBuf::from(DEFAULT_DB))
}
fn execute_help(topic: &HelpTopic, json_output: bool, out: &mut impl Write) -> u8 {
match topic {
HelpTopic::Settings => {
if json_output {
let help = plugmem_host::settings_help();
let settings: Vec<_> = help
.docs()
.iter()
.map(|doc| {
json!({
"section": doc.section,
"key": doc.key,
"type": doc.value_type,
"default": doc.default,
"description": doc.description,
"scope": doc.scope.as_str(),
})
})
.collect();
let value = json!({
"topic": "settings",
"config_path_precedence": help.config_path_precedence(),
"default_config_path": plugmem_host::default_config_path()
.map(|path| path.display().to_string()),
"settings": settings,
});
writeln!(out, "{value}").ok();
} else {
write!(out, "{}", plugmem_host::settings_help().render_human()).ok();
}
0
}
}
}
fn execute_ro(
ro: &ReadOnlyDatabase,
cmd: &Command,
recall_vector: Option<&[f32]>,
json: bool,
out: &mut impl Write,
) -> u8 {
match cmd {
Command::Recall { .. } => {
match with_recall_query(cmd, now_ms(), recall_vector, |q| ro.recall(q)) {
Ok(res) => {
render_recall(&res, json, out);
0
}
Err(e) => report_err(&CliError::Host(e)),
}
}
Command::Show { id } => render_show(ro.get(FactId(*id)), *id, json, out),
Command::Stats => {
render_stats(&ro.stats(), json, out);
0
}
Command::Export => {
ro.export_each(|f| write_export_line(out, &f));
0
}
Command::Verify => match ro.verify() {
Ok(()) => {
if json {
writeln!(out, "{}", json!({ "ok": true })).ok();
} else {
writeln!(out, "integrity ok").ok();
}
0
}
Err(e) => report_err(&CliError::Host(e)),
},
_ => unreachable!("execute_ro only receives read-only commands"),
}
}
fn execute(
db: &Database,
cmd: &Command,
json: bool,
now: u64,
out: &mut impl Write,
) -> Result<u8, CliError> {
match cmd {
Command::Remember {
text,
entity,
tags,
links,
meta,
valid_from,
vector,
} => {
let outcome = do_remember(
db,
now,
text,
entity,
tags,
links,
meta,
*valid_from,
vector,
None,
)?;
render_remember(&outcome, json, out);
Ok(0)
}
Command::Revise {
id,
text,
entity,
tags,
links,
meta,
valid_from,
vector,
} => {
let outcome = do_remember(
db,
now,
text,
entity,
tags,
links,
meta,
*valid_from,
vector,
Some(FactId(*id)),
)?;
render_remember(&outcome, json, out);
Ok(0)
}
Command::Recall { .. } => {
let res = with_recall_query(cmd, now, None, |q| db.recall(q))?;
render_recall(&res, json, out);
Ok(0)
}
Command::Forget { id } => {
let fresh = db.forget(now, FactId(*id))?;
if json {
writeln!(out, "{}", json!({ "id": id, "forgotten": fresh })).ok();
} else if fresh {
writeln!(out, "forgot fact {id}").ok();
} else {
writeln!(out, "fact {id} was already gone").ok();
}
Ok(0)
}
Command::Link { src, rel, dst } => {
db.link(LinkInput {
now,
src,
rel,
dst,
provenance: None,
})?;
if json {
writeln!(out, "{}", json!({ "src": src, "rel": rel, "dst": dst })).ok();
} else {
writeln!(out, "linked {src} -{rel}-> {dst}").ok();
}
Ok(0)
}
Command::Unlink { src, rel, dst } => {
let fresh = db.unlink(UnlinkInput { now, src, rel, dst })?;
if json {
writeln!(
out,
"{}",
json!({ "src": src, "rel": rel, "dst": dst, "unlinked": fresh })
)
.ok();
} else if fresh {
writeln!(out, "unlinked {src} -{rel}-> {dst}").ok();
} else {
writeln!(out, "edge {src} -{rel}-> {dst} was already absent").ok();
}
Ok(0)
}
Command::Show { id } => Ok(render_show(db.get(FactId(*id)), *id, json, out)),
Command::Stats => {
render_stats(&db.stats(), json, out);
Ok(0)
}
Command::Export => {
db.export_each(|f| write_export_line(out, &f));
Ok(0)
}
Command::Maintain { mode } => {
let report = db.maintain_with_options(now, maintenance_options(*mode))?;
if json {
writeln!(
out,
"{}",
json!({
"purged": report.purged,
"bytes_before": report.bytes_before,
"bytes_after": report.bytes_after,
"no_op": report.no_op,
"tombstones_before": report.tombstones_before,
"facts_before": report.facts_before,
"facts_after": report.facts_after,
"vectors_before": report.vectors_before,
"vectors_after": report.vectors_after,
"hnsw_indexed_before": report.hnsw_indexed_before,
"hnsw_indexed_after": report.hnsw_indexed_after,
"structural_compacted": report.structural_compacted,
"bm25_compacted": report.bm25_compacted,
"bm25_reindexed": report.bm25_reindexed,
"hnsw_rebuilt": report.hnsw_rebuilt,
"hnsw_remapped": report.hnsw_remapped,
"hnsw_inserted": report.hnsw_inserted,
"edges_compacted": report.edges_compacted,
"edges_before": report.edges_before,
"edge_versions_before": report.edge_versions_before,
})
)
.ok();
} else {
writeln!(
out,
"maintained: purged {}, {} -> {} bytes, hnsw +{}, bm25 {}{}{}",
report.purged,
report.bytes_before,
report.bytes_after,
report.hnsw_inserted,
if report.bm25_reindexed {
"reindexed"
} else if report.bm25_compacted {
"compacted"
} else {
"unchanged"
},
if report.edges_compacted {
", edges repacked"
} else {
""
},
if report.no_op { " (no-op)" } else { "" }
)
.ok();
}
Ok(0)
}
Command::Checkpoint => {
db.checkpoint(now)?;
if json {
writeln!(out, "{}", json!({ "ok": true })).ok();
} else {
writeln!(out, "checkpointed: journal flushed to snapshot").ok();
}
Ok(0)
}
Command::Verify => {
db.verify()?;
if json {
writeln!(out, "{}", json!({ "ok": true })).ok();
} else {
writeln!(out, "integrity ok").ok();
}
Ok(0)
}
Command::Scrub
| Command::Recover { .. }
| Command::Repl { .. }
| Command::Import { .. }
| Command::Workspace { .. }
| Command::Help { .. } => {
unreachable!("this command is dispatched before execute")
}
}
}
fn do_recover(src: &Path, dst: &Path, settings: &Settings, json: bool, out: &mut impl Write) -> u8 {
match Database::recover(src, dst, settings.config.clone(), now_ms()) {
Ok(r) => {
if json {
writeln!(
out,
"{}",
json!({
"kept": r.kept,
"dropped_text": r.dropped_text,
"dropped_vector": r.dropped_vector,
"dropped_metadata": r.dropped_metadata,
"dst": dst.display().to_string(),
})
)
.ok();
} else {
writeln!(
out,
"recovered to {}: kept {}, dropped {} text + {} vector + {} metadata",
dst.display(),
r.kept,
r.dropped_text,
r.dropped_vector,
r.dropped_metadata
)
.ok();
}
0
}
Err(HostError::Locked { path }) => report_locked(&path),
Err(e) => report_err(&CliError::Host(e)),
}
}
fn do_scrub(path: &Path, settings: &Settings, json: bool, out: &mut impl Write) -> u8 {
let ro = match Database::open_readonly(path, settings.config.clone()) {
Ok(ro) => ro,
Err(HostError::Locked { path }) => return report_locked(&path),
Err(e) => return report_err(&CliError::Host(e)),
};
let scrub = match ro.scrub() {
Ok(s) => s,
Err(e) => return report_err(&CliError::Host(e)),
};
let mut done = 0u64;
let mut total = 0u64;
for step in scrub {
match step {
Ok(p) => {
done = p.done_bytes;
total = p.total_bytes;
}
Err(e) => return report_err(&CliError::Host(e)),
}
}
if json {
writeln!(out, "{}", json!({ "ok": true, "bytes": done })).ok();
} else {
writeln!(out, "scrub ok: {done}/{total} bytes verified").ok();
}
0
}
#[derive(Parser)]
#[command(
no_binary_name = true,
name = "plugmem",
disable_help_subcommand = true
)]
struct ReplLine {
#[command(subcommand)]
command: Command,
}
fn split_line(line: &str) -> Vec<String> {
let mut tokens = Vec::new();
let mut cur = String::new();
let mut quote: Option<char> = None;
let mut has = false;
for c in line.chars() {
match quote {
Some(q) => {
if c == q {
quote = None;
} else {
cur.push(c);
}
}
None if c == '"' || c == '\'' => {
quote = Some(c);
has = true;
}
None if c.is_whitespace() => {
if has {
tokens.push(std::mem::take(&mut cur));
has = false;
}
}
None => {
cur.push(c);
has = true;
}
}
}
if has {
tokens.push(cur);
}
tokens
}
fn run_repl(
path: &Path,
settings: Settings,
json: bool,
input: impl BufRead,
out: &mut impl Write,
) -> u8 {
let db = match settings.open(path) {
Ok(db) => db,
Err(HostError::Locked { path }) => return report_locked(&path),
Err(e) => return report_err(&CliError::Host(e)),
};
eprintln!("plugmem repl — one open handle, host speed. `help` for verbs, `exit` to quit.");
eprint!("plugmem> ");
for line in input.lines() {
let Ok(line) = line else { break };
let line = line.trim();
if line.is_empty() {
eprint!("plugmem> ");
continue;
}
if line == "exit" || line == "quit" {
break;
} else if line == "help" {
writeln!(
out,
"verbs: remember recall revise forget link unlink show stats maintain checkpoint \
verify export import (scrub/recover stay one-shot) exit"
)
.ok();
} else {
run_repl_line(&db, line, json, out);
}
eprint!("plugmem> ");
}
eprintln!();
match db.checkpoint(now_ms()) {
Ok(()) => 0,
Err(e) => report_err(&CliError::Host(e)),
}
}
fn run_repl_line(db: &Database, line: &str, json: bool, out: &mut impl Write) {
let cmd = match ReplLine::try_parse_from(split_line(line)) {
Ok(r) => r.command,
Err(e) => {
let _ = writeln!(out, "{e}");
return;
}
};
match &cmd {
Command::Repl { .. } => {
let _ = writeln!(out, "already in a repl session");
}
Command::Scrub | Command::Recover { .. } => {
let _ = writeln!(
out,
"scrub/recover are one-shot commands; run them outside the repl"
);
}
_ => {
if let Err(e) = execute(db, &cmd, json, now_ms(), out) {
write_err(out, &e);
}
}
}
}
fn run_repl_ro(
path: &Path,
mut settings: Settings,
json: bool,
input: impl BufRead,
out: &mut impl Write,
) -> u8 {
let mut ro = match Database::open_readonly(path, settings.config.clone()) {
Ok(ro) => ro,
Err(HostError::Locked { path }) => return report_locked(&path),
Err(e) => return report_err(&CliError::Host(e)),
};
eprintln!(
"plugmem repl --read-only — observing generation {} of another process's writer. \
`help` for verbs, `refresh`/`generation` for cross-process freshness, `exit` to quit.",
ro.generation()
);
eprint!("plugmem(ro)> ");
for line in input.lines() {
let Ok(line) = line else { break };
let line = line.trim();
if line.is_empty() {
eprint!("plugmem(ro)> ");
continue;
}
match line {
"exit" | "quit" => break,
"help" => {
writeln!(
out,
"read verbs: recall show stats export verify \
freshness: generation refresh exit \
(writes and scrub/recover are refused in a read-only session)"
)
.ok();
}
"generation" => {
let g = ro.generation();
if json {
writeln!(out, "{}", json!({ "generation": g })).ok();
} else {
writeln!(out, "generation {g}").ok();
}
}
"refresh" => match ro.refresh() {
Ok(advanced) => {
let g = ro.generation();
if json {
writeln!(out, "{}", json!({ "advanced": advanced, "generation": g })).ok();
} else if advanced {
writeln!(out, "refreshed → generation {g}").ok();
} else {
writeln!(out, "already current → generation {g}").ok();
}
}
Err(e) => write_err(out, &CliError::Host(e)),
},
_ => run_repl_ro_line(&ro, &mut settings, line, json, out),
}
eprint!("plugmem(ro)> ");
}
eprintln!();
0
}
fn run_repl_ro_line(
ro: &ReadOnlyDatabase,
settings: &mut Settings,
line: &str,
json: bool,
out: &mut impl Write,
) {
let cmd = match ReplLine::try_parse_from(split_line(line)) {
Ok(r) => r.command,
Err(e) => {
let _ = writeln!(out, "{e}");
return;
}
};
let readable = matches!(
&cmd,
Command::Show { .. }
| Command::Stats
| Command::Export
| Command::Verify
| Command::Recall { .. }
);
if !readable {
let _ = writeln!(
out,
"read-only session: only recall/show/stats/export/verify run \
(plus refresh/generation); writes and one-shot commands need a writer handle"
);
return;
}
let recall_vector = match embed_recall_query(settings, &cmd) {
Ok(v) => v,
Err(e) => {
write_err(out, &e);
return;
}
};
let _ = execute_ro(ro, &cmd, recall_vector.as_deref(), json, out);
}
fn embed_recall_query(
settings: &mut Settings,
cmd: &Command,
) -> Result<Option<Vec<f32>>, CliError> {
let Command::Recall {
query: Some(text),
vector,
..
} = cmd
else {
return Ok(None);
};
if !vector.is_empty() {
return Ok(None);
}
let Some(embedder) = settings.embedder.as_mut() else {
return Ok(None);
};
let mut vectors = embedder.embed(&[text.as_str()]).map_err(CliError::Host)?;
Ok(vectors.pop())
}
fn with_recall_query<R>(
cmd: &Command,
now: u64,
override_vector: Option<&[f32]>,
f: impl FnOnce(RecallQuery<'_>) -> R,
) -> R {
let Command::Recall {
query,
tags,
entities,
as_of,
range,
k,
closed,
vector,
} = cmd
else {
unreachable!("with_recall_query called on a non-recall command");
};
let tag_refs: Vec<&str> = tags.iter().map(String::as_str).collect();
let ent_refs: Vec<&str> = entities.iter().map(String::as_str).collect();
let range_pair = range.as_ref().map(|v| (v[0], v[1]));
let explicit = (!vector.is_empty()).then_some(vector.as_slice());
let q = RecallQuery {
now,
text: query.as_deref(),
vector: explicit.or(override_vector),
tags: &tag_refs,
entities: &ent_refs,
as_of: *as_of,
range: range_pair,
k: *k,
token_budget: None,
include_closed: *closed,
ef: None,
};
f(q)
}
fn render_recall(res: &RecallResult, json: bool, out: &mut impl Write) {
if json {
let facts: Vec<_> = res
.facts
.iter()
.map(|f| {
json!({
"id": f.id.0,
"score": f.score,
"sources": f.sources,
"recorded_at": f.recorded_at,
"valid_from": f.valid_from,
"valid_to": open_or(f.valid_to),
})
})
.collect();
writeln!(
out,
"{}",
json!({ "facts": facts, "rendered": res.rendered, "truncated": res.truncated })
)
.ok();
} else if res.rendered.is_empty() {
writeln!(out, "(nothing recalled)").ok();
} else {
writeln!(out, "{}", res.rendered).ok();
}
}
fn render_show(
fact: Option<plugmem_host::FactSnapshot>,
id: u32,
json: bool,
out: &mut impl Write,
) -> u8 {
let Some(fact) = fact else {
if json {
writeln!(out, "{}", json!({ "id": id, "found": false })).ok();
} else {
writeln!(out, "fact {id} not found").ok();
}
return 1;
};
let r = &fact.record;
if json {
writeln!(
out,
"{}",
json!({
"id": r.id.0,
"text": fact.text,
"recorded_at": r.recorded_at,
"valid_from": r.valid_from,
"valid_to": open_or(r.valid_to),
"closed": r.is_closed(),
"tombstone": r.is_tombstone(),
"revises": (r.revises != FactId::NONE).then_some(r.revises.0),
"metadata": fact.metadata,
})
)
.ok();
} else {
writeln!(out, "fact {}", r.id.0).ok();
writeln!(out, " text {}", fact.text).ok();
writeln!(out, " recorded_at {}", r.recorded_at).ok();
write!(out, " valid [{}, ", r.valid_from).ok();
match r.valid_to {
VALID_TO_OPEN => writeln!(out, "open)").ok(),
to => writeln!(out, "{to})").ok(),
};
if r.revises != FactId::NONE {
writeln!(out, " revises fact {}", r.revises.0).ok();
}
if !fact.metadata.is_empty() {
let rendered = fact
.metadata
.iter()
.map(|(k, v)| format!("{k}={v}"))
.collect::<Vec<_>>()
.join(", ");
writeln!(out, " metadata {rendered}").ok();
}
if r.is_tombstone() {
writeln!(out, " state tombstoned").ok();
}
}
0
}
fn maintenance_options(mode: MaintainMode) -> MaintenanceOptions {
match MaintenanceMode::from(mode) {
MaintenanceMode::Auto => MaintenanceOptions::auto(),
MaintenanceMode::Full => MaintenanceOptions::full(),
mode => MaintenanceOptions {
mode,
..MaintenanceOptions::auto()
},
}
}
fn render_stats(s: &Stats, json: bool, out: &mut impl Write) {
if json {
writeln!(
out,
"{}",
json!({
"facts": s.facts,
"entities": s.entities,
"terms": s.terms,
"edges": s.edges,
"edge_versions": s.edge_versions,
"vectors": s.vectors,
"hnsw_indexed": s.hnsw_indexed,
"next_fact": s.next_fact,
"next_entity": s.next_entity,
"next_edge": s.next_edge,
"pool_bytes": s.pool_bytes,
"shards": {
"facts": s.shards.facts,
"entities": s.shards.entities,
"edges": s.shards.edges,
"temporal": s.shards.temporal,
"postings": s.shards.postings,
},
})
)
.ok();
} else {
writeln!(out, "facts {}", s.facts).ok();
writeln!(out, "entities {}", s.entities).ok();
writeln!(out, "terms {}", s.terms).ok();
writeln!(out, "edges {}", s.edges).ok();
writeln!(out, "edge_vers {}", s.edge_versions).ok();
writeln!(out, "vectors {}", s.vectors).ok();
writeln!(out, "hnsw_idx {}", s.hnsw_indexed).ok();
writeln!(out, "next_fact {}", s.next_fact).ok();
writeln!(out, "next_edge {}", s.next_edge).ok();
writeln!(out, "pool_bytes {}", s.pool_bytes).ok();
writeln!(
out,
"shards facts {} entities {} edges {} temporal {} postings {}",
s.shards.facts, s.shards.entities, s.shards.edges, s.shards.temporal, s.shards.postings,
)
.ok();
}
}
fn write_export_line(out: &mut impl Write, f: &ExportedFact) {
writeln!(
out,
"{}",
json!({
"text": f.text,
"entity": f.entity,
"tags": f.tags,
"metadata": f.metadata,
"recorded_at": f.recorded_at,
"valid_from": f.valid_from,
})
)
.ok();
}
#[cfg(test)]
fn render_export(facts: &[ExportedFact], _json: bool, out: &mut impl Write) {
for f in facts {
write_export_line(out, f);
}
}
fn do_import(
db: &Database,
now: u64,
file: &std::path::Path,
batch_size: usize,
_out: &mut impl Write,
) -> Result<usize, CliError> {
let f = std::fs::File::open(file)
.map_err(|e| CliError::Usage(format!("reading {}: {e}", file.display())))?;
let reader = io::BufReader::new(f);
let mut count = 0usize;
let mut batch: Vec<ParsedFact> = Vec::with_capacity(batch_size);
for (i, line) in reader.lines().enumerate() {
let line = line.map_err(|e| CliError::Usage(format!("line {}: {e}", i + 1)))?;
let line = line.trim();
if line.is_empty() {
continue;
}
batch.push(parse_import_line(line, i + 1)?);
if batch.len() >= batch_size {
count += flush_import_batch(db, now, &batch)?;
batch.clear();
}
}
count += flush_import_batch(db, now, &batch)?;
Ok(count)
}
struct ParsedFact {
text: String,
entity: Option<String>,
tags: Vec<String>,
metadata: Vec<(String, String)>,
valid_from: Option<u64>,
}
fn parse_import_line(line: &str, lineno: usize) -> Result<ParsedFact, CliError> {
let v: serde_json::Value =
serde_json::from_str(line).map_err(|e| CliError::Usage(format!("line {lineno}: {e}")))?;
let text = v["text"]
.as_str()
.ok_or_else(|| CliError::Usage(format!("line {lineno}: missing string \"text\"")))?
.to_string();
let entity = v["entity"].as_str().map(String::from);
let tags = v["tags"]
.as_array()
.map(|a| {
a.iter()
.filter_map(|t| t.as_str().map(String::from))
.collect()
})
.unwrap_or_default();
let metadata = v["metadata"]
.as_object()
.map(|m| {
m.iter()
.filter_map(|(k, val)| val.as_str().map(|s| (k.clone(), s.to_string())))
.collect::<BTreeMap<_, _>>()
.into_iter()
.collect()
})
.unwrap_or_default();
let valid_from = v["valid_from"].as_u64();
Ok(ParsedFact {
text,
entity,
tags,
metadata,
valid_from,
})
}
fn flush_import_batch(db: &Database, now: u64, batch: &[ParsedFact]) -> Result<usize, CliError> {
if batch.is_empty() {
return Ok(0);
}
let tag_refs: Vec<Vec<&str>> = batch
.iter()
.map(|p| p.tags.iter().map(String::as_str).collect())
.collect();
let meta_refs: Vec<Vec<(&str, &str)>> = batch
.iter()
.map(|p| {
p.metadata
.iter()
.map(|(k, v)| (k.as_str(), v.as_str()))
.collect()
})
.collect();
let inputs: Vec<RememberInput> = batch
.iter()
.zip(&tag_refs)
.zip(&meta_refs)
.map(|((p, tags), meta)| RememberInput {
entity: p.entity.as_deref(),
tags,
metadata: (!meta.is_empty()).then_some(meta.as_slice()),
valid_from: p.valid_from,
..RememberInput::text(now, &p.text)
})
.collect();
db.remember_many(inputs)?;
Ok(batch.len())
}
#[allow(clippy::too_many_arguments)]
fn do_remember(
db: &Database,
now: u64,
text: &str,
entity: &Option<String>,
tags: &[String],
links: &[String],
meta: &[String],
valid_from: Option<u64>,
vector: &[f32],
revise: Option<FactId>,
) -> Result<RememberOutcome, CliError> {
let tag_refs: Vec<&str> = tags.iter().map(String::as_str).collect();
let link_pairs = parse_links(links)?;
let link_refs: Vec<(&str, &str)> = link_pairs
.iter()
.map(|(r, e)| (r.as_str(), e.as_str()))
.collect();
let meta_map = parse_meta(meta)?;
let meta_refs: Vec<(&str, &str)> = meta_map
.iter()
.map(|(k, v)| (k.as_str(), v.as_str()))
.collect();
let input = RememberInput {
entity: entity.as_deref(),
tags: &tag_refs,
links: &link_refs,
metadata: (!meta_refs.is_empty()).then_some(meta_refs.as_slice()),
valid_from,
vector: (!vector.is_empty()).then_some(vector),
..RememberInput::text(now, text)
};
match revise {
Some(target) => Ok(db.revise(target, input)?),
None => Ok(db.remember(input)?),
}
}
fn parse_meta(meta: &[String]) -> Result<BTreeMap<String, String>, CliError> {
let mut map = BTreeMap::new();
for s in meta {
let (k, v) = s
.split_once('=')
.filter(|(k, _)| !k.is_empty())
.ok_or_else(|| CliError::Usage(format!("bad --meta `{s}` — expected KEY=VALUE")))?;
map.insert(k.to_string(), v.to_string());
}
Ok(map)
}
fn parse_links(links: &[String]) -> Result<Vec<(String, String)>, CliError> {
links
.iter()
.map(|s| {
s.split_once(':')
.filter(|(r, e)| !r.is_empty() && !e.is_empty())
.map(|(r, e)| (r.to_string(), e.to_string()))
.ok_or_else(|| CliError::Usage(format!("bad --link `{s}` — expected REL:ENTITY")))
})
.collect()
}
fn render_remember(outcome: &RememberOutcome, json: bool, out: &mut impl Write) {
if json {
let similar: Vec<_> = outcome
.similar
.iter()
.map(|s| json!({ "id": s.id.0, "score": s.score, "reason": format!("{:?}", s.reason) }))
.collect();
writeln!(
out,
"{}",
json!({
"id": outcome.id.0,
"entity": outcome.entity.map(|e| e.0),
"similar": similar,
})
)
.ok();
} else {
writeln!(out, "remembered fact {}", outcome.id.0).ok();
for s in &outcome.similar {
writeln!(
out,
" ~ similar to fact {} ({:?}, {:.2})",
s.id.0, s.reason, s.score
)
.ok();
}
}
}
fn open_or(valid_to: u64) -> Option<u64> {
(valid_to != VALID_TO_OPEN).then_some(valid_to)
}
#[cfg(test)]
mod tests {
use plugmem_host::Config;
use super::*;
struct StubEmbedder;
impl plugmem_host::Embedder for StubEmbedder {
fn dim(&self) -> usize {
3
}
fn embed(&mut self, texts: &[&str]) -> Result<Vec<Vec<f32>>, HostError> {
Ok(texts.iter().map(|_| vec![0.1, 0.2, 0.3]).collect())
}
}
fn recall_cmd(query: Option<&str>) -> Command {
Command::Recall {
query: query.map(str::to_owned),
tags: vec![],
entities: vec![],
as_of: None,
range: None,
k: 0,
closed: false,
vector: Vec::new(),
}
}
fn settings_with(embedder: Option<Box<dyn plugmem_host::Embedder>>) -> Settings {
Settings {
database_path: None,
config: Config::default(),
embedder,
snapshot_every_ops: None,
snapshot_journal_bytes: None,
maintain_every_forgets: None,
workspace: plugmem_host::WorkspaceSettings {
dir: None,
limits: plugmem_host::WorkspaceLimits::default(),
},
}
}
#[test]
fn embed_recall_query_embeds_recall_text_only_when_an_embedder_is_set() {
let mut with = settings_with(Some(Box::new(StubEmbedder)));
assert_eq!(
embed_recall_query(&mut with, &recall_cmd(Some("tokio"))).unwrap(),
Some(vec![0.1, 0.2, 0.3])
);
let mut without = settings_with(None);
assert_eq!(
embed_recall_query(&mut without, &recall_cmd(Some("tokio"))).unwrap(),
None
);
let mut with_empty = settings_with(Some(Box::new(StubEmbedder)));
assert_eq!(
embed_recall_query(&mut with_empty, &recall_cmd(None)).unwrap(),
None
);
let mut with_stats = settings_with(Some(Box::new(StubEmbedder)));
assert_eq!(
embed_recall_query(&mut with_stats, &Command::Stats).unwrap(),
None
);
}
#[test]
fn split_line_honors_quotes_and_whitespace() {
assert_eq!(split_line("remember hello"), ["remember", "hello"]);
assert_eq!(
split_line(r#"remember "two words" --tag x"#),
["remember", "two words", "--tag", "x"]
);
assert_eq!(split_line(" recall 'a b' "), ["recall", "a b"]);
assert_eq!(split_line(""), Vec::<String>::new());
assert_eq!(split_line(r#"remember """#), ["remember", ""]);
}
#[test]
fn repl_runs_over_one_handle_and_checkpoints_on_exit() {
let (db, tmp) = TempDb::open();
let path = tmp.0.join("m.plugmem");
drop(db);
let settings = settings_with(None);
let script = b"remember \"hello tokio world\"\nrecall tokio\nrevise 0 \"goodbye tokio\"\nbadcmd\nexit\n";
let mut out = Vec::new();
let code = run_repl(&path, settings, false, &script[..], &mut out);
let text = String::from_utf8(out).unwrap();
assert_eq!(code, 0);
assert!(text.contains("remembered fact 0"), "{text}");
assert!(text.contains("tokio"), "{text}");
assert!(text.contains("unrecognized subcommand"), "{text}");
let ro = Database::open_readonly(&path, Config::default()).unwrap();
assert_eq!(ro.stats().facts, 2, "original + successor after the revise");
}
#[test]
fn read_only_repl_observes_a_writer_reports_freshness_and_refuses_writes() {
let (db, tmp) = TempDb::open();
let path = tmp.0.join("m.plugmem");
let mut sink = Vec::new();
execute(
&db,
&remember("seed fact tokio", None, &[]),
false,
1_000,
&mut sink,
)
.unwrap();
db.checkpoint(1_001).unwrap();
let settings = settings_with(None);
let script = b"generation\nstats\nrefresh\nremember \"nope\"\nexit\n";
let mut out = Vec::new();
let code = run_repl_ro(&path, settings, false, &script[..], &mut out);
let text = String::from_utf8(out).unwrap();
assert_eq!(code, 0);
assert!(text.contains("generation 1"), "generation verb: {text}");
assert!(text.contains("fact"), "stats ran: {text}");
assert!(
text.contains("already current → generation 1"),
"refresh no-op: {text}"
);
assert!(text.contains("read-only session"), "write refused: {text}");
assert_eq!(db.stats().facts, 1);
}
#[test]
fn read_only_repl_refresh_advances_after_the_writer_checkpoints() {
let (db, tmp) = TempDb::open();
let path = tmp.0.join("m.plugmem");
let mut sink = Vec::new();
execute(&db, &remember("first", None, &[]), false, 1_000, &mut sink).unwrap();
db.checkpoint(1_001).unwrap();
struct HookOnFirstRead<'a> {
script: std::io::Cursor<&'a [u8]>,
db: &'a Database,
fired: bool,
}
impl std::io::Read for HookOnFirstRead<'_> {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
if !self.fired {
self.fired = true;
let mut s = Vec::new();
execute(
self.db,
&remember("second", None, &[]),
false,
2_000,
&mut s,
)
.unwrap();
self.db.checkpoint(2_001).unwrap();
}
self.script.read(buf)
}
}
let reader = std::io::BufReader::new(HookOnFirstRead {
script: std::io::Cursor::new(b"refresh\nstats\nexit\n" as &[u8]),
db: &db,
fired: false,
});
let mut out = Vec::new();
let code = run_repl_ro(&path, settings_with(None), false, reader, &mut out);
let text = String::from_utf8(out).unwrap();
assert_eq!(code, 0);
assert!(text.contains("refreshed → generation 2"), "advance: {text}");
assert!(text.contains("fact"), "stats after refresh: {text}");
assert_eq!(db.stats().facts, 2);
}
#[test]
fn read_only_repl_freshness_verbs_emit_json() {
let (db, tmp) = TempDb::open();
let path = tmp.0.join("m.plugmem");
let mut sink = Vec::new();
execute(&db, &remember("j", None, &[]), false, 1_000, &mut sink).unwrap();
db.checkpoint(1_001).unwrap();
let script = b"generation\nrefresh\nexit\n";
let mut out = Vec::new();
let code = run_repl_ro(&path, settings_with(None), true, &script[..], &mut out);
let text = String::from_utf8(out).unwrap();
assert_eq!(code, 0);
assert!(
text.contains(r#""generation":1"#),
"generation json: {text}"
);
assert!(text.contains(r#""advanced":false"#), "refresh json: {text}");
}
struct TempDb(PathBuf);
impl TempDb {
fn open() -> (Database, Self) {
let dir = std::env::temp_dir().join(format!(
"plugmem-cli-{}-{}",
std::process::id(),
now_ms_unique()
));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("m.plugmem");
let (db, _) = Database::open(&path, Config::default()).unwrap();
(db, TempDb(dir))
}
}
impl Drop for TempDb {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
fn now_ms_unique() -> String {
use std::sync::atomic::{AtomicU64, Ordering};
static N: AtomicU64 = AtomicU64::new(0);
format!("{}-{}", now_ms(), N.fetch_add(1, Ordering::Relaxed))
}
fn run_cmd(db: &Database, cmd: &Command, json: bool, now: u64) -> (u8, String) {
let mut buf = Vec::new();
let code = execute(db, cmd, json, now, &mut buf).expect("execute");
(code, String::from_utf8(buf).unwrap())
}
fn remember(text: &str, entity: Option<&str>, tags: &[&str]) -> Command {
Command::Remember {
text: text.into(),
entity: entity.map(Into::into),
tags: tags.iter().map(|t| (*t).into()).collect(),
links: Vec::new(),
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
}
}
fn remember_with_meta(text: &str, meta: &[&str]) -> Command {
Command::Remember {
text: text.into(),
entity: None,
tags: Vec::new(),
links: Vec::new(),
meta: meta.iter().map(|m| (*m).into()).collect(),
valid_from: None,
vector: Vec::new(),
}
}
#[test]
fn meta_flag_renders_sorted_in_show_and_export_and_rejects_bad_input() {
let (db, _t) = TempDb::open();
let cmd = remember_with_meta("a scan", &["uri=s3://b/x", "page=2", "page=3"]);
assert_eq!(run_cmd(&db, &cmd, false, 1_000).0, 0);
let (_, human) = run_cmd(&db, &Command::Show { id: 0 }, false, 2_000);
assert!(
human.contains("metadata page=3, uri=s3://b/x"),
"{human}"
);
let (_, jshow) = run_cmd(&db, &Command::Show { id: 0 }, true, 2_000);
let v: serde_json::Value = serde_json::from_str(&jshow).unwrap();
assert_eq!(v["metadata"]["page"], "3");
assert_eq!(v["metadata"]["uri"], "s3://b/x");
let (_, exp) = run_cmd(&db, &Command::Export, false, 2_000);
let line: serde_json::Value = serde_json::from_str(exp.lines().next().unwrap()).unwrap();
assert_eq!(line["metadata"]["uri"], "s3://b/x");
assert!(matches!(
parse_meta(&["noequals".to_string()]),
Err(CliError::Usage(_))
));
assert!(parse_meta(&["=noKey".to_string()]).is_err());
}
#[test]
fn remember_then_recall_human_and_json() {
let (db, _t) = TempDb::open();
let (code, out) = run_cmd(
&db,
&remember("prefers tokio", Some("user"), &["pref"]),
false,
1_000,
);
assert_eq!(code, 0);
assert!(out.starts_with("remembered fact 0"), "{out}");
let recall = Command::Recall {
query: Some("tokio".into()),
tags: Vec::new(),
entities: Vec::new(),
as_of: None,
range: None,
k: 0,
closed: false,
vector: Vec::new(),
};
let (code, out) = run_cmd(&db, &recall, false, 2_000);
assert_eq!(code, 0);
assert!(out.contains("tokio"), "{out}");
let (code, out) = run_cmd(&db, &recall, true, 2_000);
assert_eq!(code, 0);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert!(!v["facts"].as_array().unwrap().is_empty(), "{out}");
}
#[test]
fn recall_empty_is_ok_with_a_note() {
let (db, _t) = TempDb::open();
let recall = Command::Recall {
query: Some("nothing here".into()),
tags: Vec::new(),
entities: Vec::new(),
as_of: None,
range: None,
k: 0,
closed: false,
vector: Vec::new(),
};
let (code, out) = run_cmd(&db, &recall, false, 1_000);
assert_eq!(code, 0);
assert!(out.contains("nothing recalled"), "{out}");
}
#[test]
fn revise_closes_the_predecessor_and_conflict_is_surfaced() {
let (db, _t) = TempDb::open();
run_cmd(
&db,
&remember("lives in Moscow", Some("user"), &[]),
false,
1_000,
);
let (_, out) = run_cmd(
&db,
&remember("lives in Moscow now", Some("user"), &[]),
false,
1_500,
);
assert!(out.contains("similar to fact"), "{out}");
let revise = Command::Revise {
id: 0,
text: "lives in Berlin".into(),
entity: Some("user".into()),
tags: Vec::new(),
links: Vec::new(),
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
};
let (code, out) = run_cmd(&db, &revise, false, 2_000);
assert_eq!(code, 0);
assert!(out.starts_with("remembered fact"), "{out}");
}
#[test]
fn show_found_and_missing() {
let (db, _t) = TempDb::open();
run_cmd(&db, &remember("a note", None, &[]), false, 1_000);
let (code, out) = run_cmd(&db, &Command::Show { id: 0 }, false, 2_000);
assert_eq!(code, 0);
assert!(
out.contains("a note") && out.contains("recorded_at 1000"),
"{out}"
);
let (code, out) = run_cmd(&db, &Command::Show { id: 999 }, false, 2_000);
assert_eq!(code, 1, "missing id is a soft miss");
assert!(out.contains("not found"), "{out}");
let (_, out) = run_cmd(&db, &Command::Show { id: 0 }, true, 2_000);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["text"], "a note");
assert_eq!(v["valid_to"], serde_json::Value::Null); }
#[test]
fn forget_then_maintain_purges() {
let (db, _t) = TempDb::open();
run_cmd(&db, &remember("temp", None, &[]), false, 1_000);
let (code, out) = run_cmd(&db, &Command::Forget { id: 0 }, false, 2_000);
assert_eq!(code, 0);
assert!(out.contains("forgot fact 0"), "{out}");
let (_, out) = run_cmd(&db, &Command::Forget { id: 0 }, false, 2_100);
assert!(out.contains("already gone"), "{out}");
let (code, out) = run_cmd(
&db,
&Command::Maintain {
mode: MaintainMode::Auto,
},
false,
3_000,
);
assert_eq!(code, 0);
assert!(out.contains("purged 1"), "{out}");
}
#[test]
fn link_and_stats_and_json() {
let (db, _t) = TempDb::open();
run_cmd(
&db,
&remember("uses tokio", Some("plugmem"), &[]),
false,
1_000,
);
let link = Command::Link {
src: "plugmem".into(),
rel: "depends_on".into(),
dst: "tokio".into(),
};
let (code, out) = run_cmd(&db, &link, false, 2_000);
assert_eq!(code, 0);
assert!(out.contains("plugmem -depends_on-> tokio"), "{out}");
let unlink = Command::Unlink {
src: "plugmem".into(),
rel: "depends_on".into(),
dst: "tokio".into(),
};
let (code, out) = run_cmd(&db, &unlink, false, 2_500);
assert_eq!(code, 0);
assert!(
out.contains("unlinked plugmem -depends_on-> tokio"),
"{out}"
);
let (code, out) = run_cmd(&db, &Command::Stats, true, 3_000);
assert_eq!(code, 0);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["facts"], 1);
assert_eq!(v["edges"], 0);
assert_eq!(v["edge_versions"], 1);
}
#[test]
fn bad_link_is_a_usage_error() {
let (db, _t) = TempDb::open();
let cmd = Command::Remember {
text: "x".into(),
entity: Some("user".into()),
tags: Vec::new(),
links: vec!["not-a-pair".into()],
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
};
let mut buf = Vec::new();
let err = execute(&db, &cmd, false, 1_000, &mut buf).unwrap_err();
assert!(matches!(err, CliError::Usage(_)));
}
#[test]
fn as_of_time_travel_via_recall() {
let (db, _t) = TempDb::open();
run_cmd(
&db,
&remember("lives in Moscow", Some("user"), &[]),
false,
1_000,
);
let revise = Command::Revise {
id: 0,
text: "lives in Berlin".into(),
entity: Some("user".into()),
tags: Vec::new(),
links: Vec::new(),
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
};
run_cmd(&db, &revise, false, 2_000);
let as_of = Command::Recall {
query: Some("lives".into()),
tags: Vec::new(),
entities: vec!["user".into()],
as_of: Some(1_500),
range: None,
k: 0,
closed: false,
vector: Vec::new(),
};
let (_, out) = run_cmd(&db, &as_of, false, 3_000);
assert!(out.contains("Moscow"), "as-of 1500 → Moscow: {out}");
}
#[test]
fn every_command_has_a_json_shape() {
let (db, _t) = TempDb::open();
let (_, out) = run_cmd(
&db,
&remember("uses tokio", Some("plugmem"), &["pref"]),
true,
1_000,
);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["id"], 0);
assert!(v["similar"].is_array());
let revise = Command::Revise {
id: 0,
text: "uses tokio now".into(),
entity: Some("plugmem".into()),
tags: Vec::new(),
links: Vec::new(),
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
};
let (_, out) = run_cmd(&db, &revise, true, 1_500);
assert!(serde_json::from_str::<serde_json::Value>(out.trim()).is_ok());
let link = Command::Link {
src: "plugmem".into(),
rel: "depends_on".into(),
dst: "tokio".into(),
};
let (_, out) = run_cmd(&db, &link, true, 2_000);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["rel"], "depends_on");
let unlink = Command::Unlink {
src: "plugmem".into(),
rel: "depends_on".into(),
dst: "tokio".into(),
};
let (_, out) = run_cmd(&db, &unlink, true, 2_100);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["unlinked"], true);
let (_, out) = run_cmd(&db, &Command::Forget { id: 1 }, true, 2_500);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["forgotten"], true);
let (_, out) = run_cmd(
&db,
&Command::Maintain {
mode: MaintainMode::Auto,
},
true,
3_000,
);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert!(v["purged"].as_u64().unwrap() >= 1);
let (code, out) = run_cmd(&db, &Command::Show { id: 999 }, true, 3_500);
assert_eq!(code, 1);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["found"], false);
let recall = Command::Recall {
query: None,
tags: Vec::new(),
entities: vec!["plugmem".into()],
as_of: None,
range: Some(vec![0, 10_000]),
k: 4,
closed: true,
vector: Vec::new(),
};
let (_, out) = run_cmd(&db, &recall, true, 4_000);
assert!(serde_json::from_str::<serde_json::Value>(out.trim()).is_ok());
}
#[test]
fn stats_human_lists_the_counters() {
let (db, _t) = TempDb::open();
run_cmd(&db, &remember("a", None, &[]), false, 1_000);
let (code, out) = run_cmd(&db, &Command::Stats, false, 2_000);
assert_eq!(code, 0);
assert!(out.contains("facts") && out.contains("pool_bytes"), "{out}");
}
#[test]
fn verify_command_renders_human_and_json() {
let (db, _t) = TempDb::open();
run_cmd(&db, &remember("clean", None, &[]), false, 1_000);
let (code, out) = run_cmd(&db, &Command::Verify, false, 2_000);
assert_eq!(code, 0);
assert_eq!(out.trim(), "integrity ok");
let (code, out) = run_cmd(&db, &Command::Verify, true, 2_100);
assert_eq!(code, 0);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["ok"], true);
}
#[test]
fn show_json_of_a_revised_predecessor_is_closed() {
let (db, _t) = TempDb::open();
run_cmd(&db, &remember("v1", Some("e"), &[]), false, 1_000);
let revise = Command::Revise {
id: 0,
text: "v2".into(),
entity: Some("e".into()),
tags: Vec::new(),
links: Vec::new(),
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
};
run_cmd(&db, &revise, false, 2_000);
let (_, out) = run_cmd(&db, &Command::Show { id: 1 }, false, 3_000);
assert!(out.contains("revises fact 0"), "{out}");
}
#[test]
fn resolve_db_path_prefers_the_flag() {
let p = "/tmp/explicit.plugmem";
assert_eq!(resolve_db_path(Some(p), None, None), PathBuf::from(p));
let configured = std::path::Path::new("/tmp/configured.plugmem");
assert_eq!(
resolve_db_path(None, Some(configured), None),
PathBuf::from(configured)
);
let _ = resolve_db_path(None, None, None);
}
#[test]
fn a_bare_name_is_a_memory_only_when_a_workspace_is_configured() {
let root = PathBuf::from("/srv/bot");
assert_eq!(
resolve_db_path(Some("work"), None, None),
PathBuf::from("work")
);
assert_eq!(
resolve_db_path(Some("work"), None, Some(&root)),
PathBuf::from("/srv/bot/db/work.plugmem")
);
for path in ["./work", "work.plugmem", "/srv/other.plugmem", "../up"] {
assert_eq!(
resolve_db_path(Some(path), None, Some(&root)),
PathBuf::from(path),
"{path}"
);
}
let configured = std::path::Path::new("work");
assert_eq!(
resolve_db_path(None, Some(configured), Some(&root)),
PathBuf::from("work")
);
}
#[test]
fn settings_help_runs_without_opening_a_database() {
let cli = Cli::try_parse_from(["plugmem-cli", "help", "settings"]).unwrap();
let mut output = Vec::new();
assert_eq!(run_parsed(cli, &mut output), 0);
let output = String::from_utf8(output).unwrap();
assert!(output.contains("plugmem settings"));
assert!(output.contains("[database]"));
assert!(output.contains("path (path string"));
let cli = Cli::try_parse_from(["plugmem-cli", "--json", "help", "settings"]).unwrap();
let mut output = Vec::new();
assert_eq!(run_parsed(cli, &mut output), 0);
let output: serde_json::Value = serde_json::from_slice(&output).unwrap();
assert_eq!(output["topic"], "settings");
assert!(output["config_path_precedence"].is_array());
assert!(output["settings"].as_array().unwrap().len() > 10);
}
#[test]
fn run_parsed_opens_runs_and_reports() {
let dir = std::env::temp_dir().join(format!(
"plugmem-run-{}-{}",
std::process::id(),
now_ms_unique()
));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("m.plugmem");
let cli = Cli {
db: Some(path.display().to_string()),
workspace: None,
config: None,
json: false,
command: Command::Stats,
};
let mut buf = Vec::new();
let code = run_parsed(cli, &mut buf);
assert_eq!(code, 0);
assert!(String::from_utf8(buf).unwrap().contains("facts"));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn recover_and_scrub_render_json_and_human_shapes() {
let (db, tmp) = TempDb::open();
let path = tmp.0.join("m.plugmem");
run_cmd(&db, &remember("recoverable fact", None, &[]), false, 1_000);
run_cmd(&db, &Command::Checkpoint, false, 2_000);
drop(db);
let settings = settings_with(None);
let mut out = Vec::new();
assert_eq!(do_scrub(&path, &settings, true, &mut out), 0);
let scrub: serde_json::Value = serde_json::from_slice(&out).unwrap();
assert_eq!(scrub["ok"], true);
assert!(scrub["bytes"].as_u64().unwrap() > 0);
let mut out = Vec::new();
assert_eq!(do_scrub(&path, &settings, false, &mut out), 0);
let out = String::from_utf8(out).unwrap();
assert!(out.contains("scrub ok:"), "{out}");
let json_dst = tmp.0.join("copy-json.plugmem");
let mut out = Vec::new();
assert_eq!(do_recover(&path, &json_dst, &settings, true, &mut out), 0);
let recover: serde_json::Value = serde_json::from_slice(&out).unwrap();
assert_eq!(recover["kept"], 1);
assert_eq!(recover["dropped_text"], 0);
assert_eq!(recover["dst"], json_dst.display().to_string());
let human_dst = tmp.0.join("copy-human.plugmem");
let mut out = Vec::new();
assert_eq!(do_recover(&path, &human_dst, &settings, false, &mut out), 0);
let out = String::from_utf8(out).unwrap();
assert!(out.contains("recovered to"), "{out}");
assert!(out.contains("kept 1"), "{out}");
}
#[test]
fn readonly_dispatcher_renders_every_read_shape() {
let (db, tmp) = TempDb::open();
let path = tmp.0.join("m.plugmem");
run_cmd(
&db,
&remember("readonly tokio fact", Some("plugmem"), &["pref"]),
false,
1_000,
);
run_cmd(&db, &Command::Checkpoint, false, 2_000);
let ro = Database::open_readonly(&path, Config::default()).unwrap();
let mut out = Vec::new();
assert_eq!(execute_ro(&ro, &Command::Stats, None, true, &mut out), 0);
let stats: serde_json::Value = serde_json::from_slice(&out).unwrap();
assert_eq!(stats["facts"], 1);
let mut out = Vec::new();
assert_eq!(
execute_ro(&ro, &Command::Show { id: 0 }, None, false, &mut out),
0
);
let text = String::from_utf8(out).unwrap();
assert!(text.contains("readonly tokio fact"), "{text}");
let mut out = Vec::new();
assert_eq!(execute_ro(&ro, &Command::Export, None, false, &mut out), 0);
let exported: serde_json::Value =
serde_json::from_str(String::from_utf8(out).unwrap().lines().next().unwrap()).unwrap();
assert_eq!(exported["text"], "readonly tokio fact");
let mut out = Vec::new();
let recall = Command::Recall {
query: Some("tokio".into()),
tags: vec!["pref".into()],
entities: vec!["plugmem".into()],
as_of: None,
range: None,
k: 1,
closed: false,
vector: Vec::new(),
};
assert_eq!(execute_ro(&ro, &recall, None, false, &mut out), 0);
let text = String::from_utf8(out).unwrap();
assert!(text.contains("tokio"), "{text}");
let mut out = Vec::new();
assert_eq!(execute_ro(&ro, &Command::Verify, None, true, &mut out), 0);
let verify: serde_json::Value = serde_json::from_slice(&out).unwrap();
assert_eq!(verify["ok"], true);
}
#[test]
fn run_parsed_on_a_locked_database_returns_one() {
let (_held, dir) = {
let dir = std::env::temp_dir().join(format!(
"plugmem-lock-{}-{}",
std::process::id(),
now_ms_unique()
));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("m.plugmem");
(Database::open(&path, Config::default()).unwrap(), dir)
};
let cli = Cli {
db: Some(dir.join("m.plugmem").display().to_string()),
workspace: None,
config: None,
json: false,
command: Command::Stats,
};
let mut buf = Vec::new();
assert_eq!(run_parsed(cli, &mut buf), 1);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn run_parsed_propagates_a_usage_error_as_two() {
let dir = std::env::temp_dir().join(format!(
"plugmem-usage-{}-{}",
std::process::id(),
now_ms_unique()
));
std::fs::create_dir_all(&dir).unwrap();
let cli = Cli {
db: Some(dir.join("m.plugmem").display().to_string()),
workspace: None,
config: None,
json: false,
command: Command::Remember {
text: "x".into(),
entity: None,
tags: Vec::new(),
links: vec!["bad".into()],
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
},
};
let mut buf = Vec::new();
assert_eq!(run_parsed(cli, &mut buf), 2);
let _ = std::fs::remove_dir_all(&dir);
}
struct Scratch(PathBuf);
impl Scratch {
fn new(tag: &str) -> Self {
let dir = std::env::temp_dir().join(format!(
"plugmem-cli-{tag}-{}-{}",
std::process::id(),
now_ms_unique()
));
std::fs::create_dir_all(&dir).unwrap();
Scratch(dir)
}
}
impl Drop for Scratch {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
#[test]
fn export_import_roundtrip_preserves_open_facts() {
let (a, _ta) = TempDb::open();
run_cmd(
&a,
&Command::Remember {
text: "prefers tokio".into(),
entity: Some("user".into()),
tags: vec!["pref".into(), "lang".into()],
links: Vec::new(),
meta: vec!["uri=s3://b/x".into(), "src=chat".into()],
valid_from: Some(500),
vector: Vec::new(),
},
false,
1_000,
);
run_cmd(
&a,
&remember("lives in Moscow", Some("user"), &[]),
false,
1_100,
); run_cmd(
&a,
&Command::Revise {
id: 1,
text: "lives in Berlin".into(),
entity: Some("user".into()),
tags: vec!["geo".into()],
links: Vec::new(),
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
},
false,
1_200,
); run_cmd(&a, &remember("junk", None, &[]), false, 1_300); run_cmd(&a, &Command::Forget { id: 3 }, false, 1_400); run_cmd(
&a,
&remember("uses rust", Some("plugmem"), &["lang"]),
false,
1_500,
);
let mut dump = Vec::new();
render_export(&a.export(), false, &mut dump);
let scratch = Scratch::new("roundtrip");
let file = scratch.0.join("dump.jsonl");
std::fs::write(&file, &dump).unwrap();
let (b, _tb) = TempDb::open();
let n = do_import(&b, 9_000, &file, 128, &mut Vec::new()).unwrap();
let key = |f: &ExportedFact| {
let mut tags = f.tags.clone();
tags.sort();
(f.text.clone(), f.entity.clone(), tags, f.valid_from)
};
let mut ak: Vec<_> = a.export().iter().map(key).collect();
let mut bk: Vec<_> = b.export().iter().map(key).collect();
ak.sort();
bk.sort();
assert_eq!(n, ak.len());
assert_eq!(
ak, bk,
"roundtrip must preserve text/entity/tags/valid_from"
);
let b_open = b.export();
assert!(b_open.iter().any(|f| f.text == "prefers tokio"
&& f.valid_from == 500
&& f.entity.as_deref() == Some("user")
&& f.tags == vec!["pref".to_string(), "lang".to_string()]
&& f.metadata.get("uri").map(String::as_str) == Some("s3://b/x")
&& f.metadata.get("src").map(String::as_str) == Some("chat")));
assert!(b_open.iter().any(|f| f.text == "lives in Berlin"));
assert!(b_open.iter().any(|f| f.text == "uses rust"));
assert!(!b_open.iter().any(|f| f.text.contains("Moscow")));
assert!(!b_open.iter().any(|f| f.text == "junk"));
}
#[test]
fn export_command_emits_jsonl_regardless_of_json_flag() {
let (db, _t) = TempDb::open();
run_cmd(&db, &remember("a fact", Some("e"), &["t"]), false, 1_000);
for json in [false, true] {
let (code, out) = run_cmd(&db, &Command::Export, json, 2_000);
assert_eq!(code, 0);
let v: serde_json::Value = serde_json::from_str(out.trim()).unwrap();
assert_eq!(v["text"], "a fact");
assert_eq!(v["entity"], "e");
assert_eq!(v["tags"][0], "t");
}
}
#[test]
fn import_command_counts_and_rejects_bad_lines() {
let (db, _t) = TempDb::open();
let scratch = Scratch::new("import");
let good = scratch.0.join("in.jsonl");
std::fs::write(
&good,
"{\"text\":\"from jsonl\",\"entity\":\"user\",\"tags\":[\"x\"],\"valid_from\":42}\n\n{\"text\":\"second\"}\n",
)
.unwrap();
let n = do_import(&db, 9_000, &good, 1, &mut Vec::new()).unwrap();
assert_eq!(n, 2, "both facts imported, blank line skipped");
let bad = scratch.0.join("bad.jsonl");
std::fs::write(&bad, "not json at all\n").unwrap();
let err = do_import(&db, 9_000, &bad, 128, &mut Vec::new()).unwrap_err();
assert!(matches!(err, CliError::Usage(_)));
}
#[test]
fn import_batch_size_does_not_change_the_result() {
let scratch = Scratch::new("import-batch");
let file = scratch.0.join("facts.jsonl");
let mut jsonl = String::new();
for i in 0..5 {
jsonl.push_str(&format!("{{\"text\":\"fact number {i}\"}}\n"));
}
std::fs::write(&file, &jsonl).unwrap();
let (a, _ta) = TempDb::open();
let (b, _tb) = TempDb::open();
let na = do_import(&a, 9_000, &file, 1, &mut Vec::new()).unwrap();
let nb = do_import(&b, 9_000, &file, 100, &mut Vec::new()).unwrap();
assert_eq!(na, 5);
assert_eq!(nb, 5);
let texts = |db: &Database| {
let mut t: Vec<_> = db.export().into_iter().map(|f| f.text).collect();
t.sort();
t
};
assert_eq!(texts(&a), texts(&b), "batch size must not change the facts");
}
#[test]
fn config_table_feeds_settings_and_the_cli_batch_size() {
let scratch = Scratch::new("settings");
let cfgfile = scratch.0.join("config.toml");
std::fs::write(
&cfgfile,
"[engine]\ndim = 512\n[embedder]\nkind = \"none\"\n\
[maintenance]\nsnapshot_every_ops = 64\nbatch_size = 200\n",
)
.unwrap();
let table = plugmem_host::read_config(Some(&cfgfile)).unwrap();
let s = Settings::from_table(table.as_ref()).unwrap();
assert_eq!(s.config.dim, 512);
assert!(s.embedder.is_none());
assert_eq!(s.snapshot_every_ops, Some(64));
assert_eq!(read_batch_size(table.as_ref()), Some(200));
assert!(plugmem_host::read_config(Some(&scratch.0.join("nope.toml"))).is_err());
}
#[test]
fn checkpoint_command_flushes_the_journal_and_enables_the_readonly_path() {
let scratch = Scratch::new("checkpoint-cmd");
let path = scratch.0.join("m.plugmem");
let remember = Cli {
db: Some(path.display().to_string()),
workspace: None,
config: None,
json: false,
command: Command::Remember {
text: "hello tokio".into(),
entity: None,
tags: Vec::new(),
links: Vec::new(),
meta: Vec::new(),
valid_from: None,
vector: Vec::new(),
},
};
assert_eq!(run_parsed(remember, &mut Vec::new()), 0);
let checkpoint = |json| Cli {
db: Some(path.display().to_string()),
workspace: None,
config: None,
json,
command: Command::Checkpoint,
};
let mut buf = Vec::new();
assert_eq!(run_parsed(checkpoint(false), &mut buf), 0);
assert!(String::from_utf8(buf).unwrap().contains("checkpointed"));
let mut buf = Vec::new();
assert_eq!(run_parsed(checkpoint(true), &mut buf), 0);
let v: serde_json::Value =
serde_json::from_str(String::from_utf8(buf).unwrap().trim()).unwrap();
assert_eq!(v["ok"], true);
let scrub = Cli {
db: Some(path.display().to_string()),
workspace: None,
config: None,
json: false,
command: Command::Scrub,
};
let mut buf = Vec::new();
assert_eq!(run_parsed(scrub, &mut buf), 0);
assert!(String::from_utf8(buf).unwrap().contains("scrub ok"));
}
#[test]
fn run_parsed_uses_the_readonly_path_after_a_checkpoint() {
let scratch = Scratch::new("ro-route");
let path = scratch.0.join("m.plugmem");
{
let (db, _) = Database::open(&path, Config::default()).unwrap();
db.remember(RememberInput::text(1_000, "hello tokio"))
.unwrap();
db.checkpoint(2_000).unwrap(); }
let cli = Cli {
db: Some(path.display().to_string()),
workspace: None,
config: None,
json: false,
command: Command::Stats,
};
let mut buf = Vec::new();
assert_eq!(run_parsed(cli, &mut buf), 0);
assert!(String::from_utf8(buf).unwrap().contains("facts"));
let cli = Cli {
db: Some(path.display().to_string()),
workspace: None,
config: None,
json: false,
command: Command::Recall {
query: Some("tokio".into()),
tags: Vec::new(),
entities: Vec::new(),
as_of: None,
range: None,
k: 0,
closed: false,
vector: Vec::new(),
},
};
let mut buf = Vec::new();
assert_eq!(run_parsed(cli, &mut buf), 0);
assert!(String::from_utf8(buf).unwrap().contains("tokio"));
}
}