mod corpus;
use areev::pack;
mod run_stack;
mod tune;
mod trigger_cli;
use std::collections::HashMap;
use std::process::ExitCode;
use areev_cal::{CalExecutor, CalExecutorConfig, AreevFacade};
use areev_core::error::Hash;
use areev_core::types::{ContentBlock, Event, Fact, Grain, TokenUsage, Tool};
use areev_store::{is_pg_dsn, Axis, Areev, Direction};
use areev_loop_adapter::{now_ms, AreevSubstrate};
use areev_loop::{Decision, Engine, ObserverType, Policy, RecStatus, RunOptions, ScopeSet, Severity};
const USAGE: &str = "\
areev — embedded memory engine for AI agents (OMS + CAL on Turso)
USAGE:
areev <command> [--db <memory.db>] [options] (-d = --db)
areev --version | -V print the version and exit
areev help | --help | -h show this help
COMMANDS:
init [--template blank|demo|coding-agent] [--ns NS] seed a backend +
print the Claude Code hook snippet (never writes your settings)
loop <run|reflect|list|show|approve|reject|apply|rollback|analyzers|policy|replay>
the governed self-improvement loop (deterministic core; optional verified LLM):
run [--min-new N --min-new-errors N --if-stale 6h --format json --quiet]
[--model provider:name | --llm-cmd 'CMD'] optional LLM reflection
(--model reads the key from $ANTHROPIC_API_KEY/$OPENAI_API_KEY/etc.)
[--ground-model provider:name | --ground-cmd 'CMD'] separate
grounding backend (defaults to the reflection model)
[--analyzer-cmd 'CMD'] register an external analyzer
(advisory only — never auto-applies)
reflect like run, but re-analyzes the whole memory (ignores the
incremental watermark) — a full sweep; same flags as run
list [--status pending|all|applied|...] [--fail-on high] (exit 2 on match)
show <hash> | approve/reject/apply/rollback <hash> --because \"...\" [--actor A]
outcomes the Verify gate: did applied advice hold, or regress?
replay --config FILE [--window 90d | --since MS] [--step per-pass|1d]
score a candidate config against the recorded past, beside the
incumbent — zero writes, prefix-only reads, LLM/command not replayed
[--policy FILE] grants auto-apply (else $AREEV_LOOP_POLICY); `policy` prints it
add <subject> <relation> <object> store a fact (positional)
[--subject S --relation R --object O] [--ns NS] [--confidence C]
[--idempotent] no-op if this exact value is already the head
record-tool-call --name NAME [--input JSON] --result TEXT [--is-error]
[--thread ID] [--call-id ID] [--ns NS] append one tool invocation
[--failure-cause CAUSE] [--failure-detail TEXT] why it failed
run-manifest --run-id ID --config JSON persist a reproducible harness
configuration and its run link in agent:harness
recall <subject> | --subject S [--relation R] [--ns NS] [-k N]
[--render sml|toon|markdown|plain|json] [--budget TOKENS]
--ns accepts a prefix scope 'org.*' (= org + its .-descendants);
writes and destruction never accept patterns
cal <QUERY> [--ns NS] execute a CAL statement
corpus --select '<READ CAL>' [--out FILE] [--recipient ID]
stream governed OpenAI chat
JSONL and record its provenance registry
tune --cmd 'TRAINER' --evalset HASH
(--select '<READ CAL>' --out FILE | --corpus FILE --manifest HASH)
[--recipient ID] [--timeout-secs N]
hand the corpus to YOUR trainer (JSON on
stdio; Areev never trains) and register
the returned adapter with full lineage;
promotion then rides the loop:
`loop run` proposes, `eval run --model`
gates, `loop apply --gating-run` admits
search --query TEXT [--subject S] [-k N] hybrid recall (BM25 + structural, RRF);
each JSON row carries `score` (fusion, top = 1.0, or the reranker's)
history --subject S --relation R [--ns NS]
provenance <source-hash> grains distilled from a source (reverse)
forks open forks (>1 head for a subject+relation)
merge --subject S --relation R --object O close a fork with a resolved value
subject-report <subject> [--ns NS] [--text-mentions] [--out FILE] [--bundle FILE]
DSAR read (GDPR Art. 15/20): everything
forget-subject WOULD erase — exact +
partition keys, history — as JSONL
(stdout or --out), and optionally a
portable .mgb bundle (--bundle)
forget-subject <subject> [--ns NS] [--text-mentions] [--because \"why\"]
[--override-hold] --yes
erase EVERY grain referencing an
identity (--because rides the audit
record) — exact +
partition keys (pat, pat#visit1), history
included, + its dictionary entries;
--text-mentions also erases grains whose
indexed text mentions the identity;
replicates as tombstones.
A namespace under a legal hold refuses
(STO-E009) and the refusal is itself
recorded; --override-hold --because
destroys anyway, names the hold in the
audit record, and needs admin on the
namespace as well as erase
purge-older-than <days> [--ns NS] [--type event] --yes retention sweep:
erase grains older than N days
(--ns \"\" sweeps every namespace)
retention <set|list|clear|sweep|floor|floor-clear|floors>
[--days N] [--min-days N] [--ns NS] [--type event]
[--because \"why\"] [--yes] declarative storage limitation: the
policy is a file-truth that travels
with the memory; `set` declares,
`sweep --yes` enforces (audited).
`floor --min-days N` declares a
MINIMUM any sweep must respect —
destruction younger than it refuses
hold <set|release|list> [--ns NS] --because \"why\" [--by PRINCIPAL]
list also takes [--format json] (adds at_ms, the placement time)
legal hold: while one is live on a
namespace, EVERY destruction path there
refuses (STO-E009) with the hold on
record — the age sweeps, FORGET by hash,
FORGET SUBJECT, the memory tool, the
loop rollback, and DROP SCHEMA.
A file-truth; `set` and `release` both
demand a reason and both land in
`areev audit export`
trigger add --type KIND --workflow HASH --because \"why\"
[--context-query SPEC] a saved query the evaluator runs at
fire time; its result rides into the
run input as `context` (how a
trigger-started run reads its own
memory on the embedded backend).
SPEC is `name` or `name($p = /ptr,
...)` — pointers resolve against the
firing item's payload (fail-closed)
[--interval SECS | --cron EXPR | --at MS] [--observer NAME]
[--connector-tool HASH] the connector's CODE as a grain: a Tool
Definition whose executor_uri carries the
blob. --observer/--connector still names
it (that name is half the run-id dedup
identity). Nothing runs unless the
evaluating host pinned the address with
--allow-executor
[--scope S] [--dedup-key PTR] [--catchup last|none|all]
[--where EXPR] [--members ALIAS=HASH,...] [--correlate PTR]
[--window 10m]
declare a trigger. Like retention, the
declaration is a file-truth that
travels with the memory and is enforced
separately. --where is CAL WHERE syntax:
it selects grains for a `memory` trigger
and gates member aliases for a
`composite` one
trigger run [--id T] [--connector-cmd CMD] [--tool-cmd CMD] [--dry-run]
[--lease SECS] [--max-items N] [--credential NAME=ENV_VAR|cmd:CMD|vault:P#F]
[--credential-ttl SECS] [--resolver-env VAR,...]
[--allow-executor HEX,...] [--sandbox-cmd CMD] [--executor-cache DIR]
[--executor-timeout SECS] [--tool-env VAR,...]
[--model SPEC] [--base-url URL] [--key-env VAR]
[--max-tokens N] [--max-usd USD] [--max-wall-ms MS] [--ask-ttl SECS]
[--max-effects N] [--llm-tool-result-chars N] [--llm-context-tokens N]
evaluate once and exit — the cadence is
data in the memory, so the heartbeat can
be coarse. Safe to invoke concurrently.
A firing gets the SAME runner `run start`
builds: the executor pin, the sandbox and
the model all reach it, so a plan that
runs by hand runs on a heartbeat — and
the same pin decides whether a trigger's
own --connector-tool code may run. Every
one of those also reads its $AREEV_RUN_*
variable, because a heartbeat is a cron
line, not an interactive command.
With none of --tool-cmd/--allow-executor/
--model it ingests without executing;
--dry-run touches nothing.
--credential names an env var whose value
the egress broker attaches on the way out,
so the connector never holds the token.
cmd:/vault: MINT one per call instead —
what an unattended heartbeat wants, since
a token minted at boot expires by morning.
The budget flags are the ones `run start`
takes, and bound the runs a firing starts
trigger render --target cron|launchd|systemd|k8s-cronjob
emit heartbeat config for infrastructure
you already run; never creates anything
trigger deliver --id T [< payload.json] [--tool-cmd CMD] [budget flags]
hand a webhook/manual payload to Areev.
The host owns the listener — Areev never
opens a port. Idempotent on the dedup key.
Takes the same runner and budget flags as
`trigger run`
trigger list | show <T> | status what is declared, and what has actually
fired (a trigger that never fired is
reported, not silent). `show` also
surfaces the cursor and, for a
superseded declaration, the state key
its evaluation actually lives under
trigger retarget <T> [--workflow HASH] --because \"why\"
point a trigger at the CURRENT version of
its plan (editing a plan mints a new hash
and triggers do not follow heads). Without
--workflow it follows the plan's own chain
to the live head; the destination is
validated before anything is written, and
the cursor/dedup fence carry over
trigger pause <T> --because \"why\" | resume <T> --because \"why\"
operational brake; unlike the
declaration's own enabled flag it stays
on this host and does not replicate
blob put <FILE> | put --stdin store bytes in the content-addressed
blob store; prints the cas:// URI
(idempotent — the address IS the content)
get <cas-uri> write those bytes to stdout, hash-verified.
Reads the .blobs sidecar WITHOUT opening
the memory, so a --tool-cmd subprocess can
fetch an attachment while its run holds
the single writer (an encrypted memory
still needs --passphrase-env)
blobs encrypt encrypt CAS attachments written before
sidecar encryption landed (needs the
memory's key; new blobs are sealed on
write). Idempotent
audit export [--since MS] [--until MS] [--out FILE] [--with-outcomes]
accountability evidence as JSONL: every
destructive op (who/what/why/how many)
+ the loop lifecycle with its gating
edge. BOTH trails hash-chained and
verified — a missing predecessor is a
chain_break, one merely outside --since
is previous_outside_window, and a
pre-chain record is unchained, never a
break. --with-outcomes adds the Verify
gate's measured checkpoints, marked
chained: false. It contains NO read
records: access logging is the host's
novelty --text T [--subject S] [--relation R] [-k N] [--ns NS]
nearest existing grains
(paraphrase check; needs --embed-cmd)
--ns accepts a prefix scope 'org.*' like recall does
log [--since OP] [--limit N] [--ns a,b]
op-log (change feed). --ns narrows it to
those namespaces and attributes every
row, TOMBSTONES INCLUDED — resolving a
forget's hash cannot, because the grain
is gone. op_seq stays the memory-wide
sequence, so a scoped cursor is still
comparable with an unscoped one
telemetry scrub --ns NS --yes drop every recall-telemetry row for one
namespace, plus buffered events. Reaches
the row no other scrub can: a
zero-result free-text query names no
grain hash
bundle --out FILE [--since OP] incremental backup (git-shaped)
import --bundle FILE apply a bundle (fast-forward)
pack validate DIR check a pack without opening a memory:
it builds every grain, addresses it, and
compares against the manifest's
expected_hash. Prints the pins its code
needs
pack install DIR [--dry-run] seed a pack into --db: blobs into the
CAS, `blob:`/`grain:` references
rewritten to the addresses they turn out
to have, saved queries restored. An
expected_hash that does not match is
REFUSED with nothing written.
[--expected-hash H] refuse unless the pack's plan builds to H
[--pin TOOL=ADDR,...] check this host's executor pins against
the pack's code (tool_name, grain id or
blob name → address); a mismatch or a
pin naming no code-carrying tool is
PCK-E005, nothing written. Pins are
checked, never stored
pack export --out DIR [--pack-format source|bundle] [--name N]
turn this memory's namespace into an
installable pack. `source` writes
reviewable grain JSON (and refuses to
write anything it cannot rebuild to the
same address); `bundle` writes a bundle
plus a manifest of what it must contain.
See docs/pack.md
migrate --from SRC --file PATH [--history PATH] import another system's
export: mem0 | mem0-history | langgraph | letta | letta-archival |
zep | basic-memory (PATH = vault dir) | jsonl (generic
{subject,relation,object}|{content} lines) | tool-log (OpenAI-style
tool-call JSONL → Tool grains, feeds tool-failure clustering).
Re-runs skip what is already imported; see docs/migrate.md.
reindex backfill + rebuild the BM25 text index
(e.g. after --index-text true on a file
written with indexing off)
vector-index <status|build|drop|check> [--m N] [--ef-construction N]
[--ef-search N] [--queries FILE.json --ns SCOPE] [--k N]
the ANN (pgvector HNSW) index over the
stored vectors. Postgres only: build
answers STO-E007 on a file memory, which
scans exactly. `check` grades the index
against the exact scan with YOUR query
vectors (recall@k) over the WIDE --ns
scope you actually query; --ef-search
retunes the session first (no rebuild).
`status` says whether reads are exact.
stream --to DIR [--interval-ms N] [--once] [--checkpoint] [--retain 30d]
continuous op-log shipping; --checkpoint
opens a new generation with a full
snapshot, --retain drops generations
older than the window (so erasure
reaches archives — see docs/gdpr.md)
restore --from DIR [--until-hlc T] rebuild from stream dir (PITR)
follow --from DIR [--interval-ms N] [--once] subscribe: apply new segments
(org/category distribution)
related --start TERM --relations R1,R2 [--direction out|in|both]
[--depth N] [--limit N] [--ns a,b]
walk the entity graph (bounded k-hop).
in/both only see relations the file
declares entity-valued. --ns with a
comma list walks the SET: a walk is not
composable from per-namespace calls, and
every named namespace is read-checked as
itself
entity-at --subject S --relation R --at MS [--axis world|knowledge] [--ns a,b]
as-of read: what was true at T (world)
or what was known at T (knowledge).
The world axis prefers the window that
took effect most recently, so an
out-of-order backfill reads correctly.
--ns with a comma list answers each
namespace independently
step-actions --workflow HASH [--node ID] [--limit N]
execution records for a workflow —
which grains ran which of its nodes
eval <create|run> the §7.4 gating edge: create stores an
evalset (--name N --cases FILE); run --evalset HASH executes it
against --tool-cmd CMD or --model provider:name ([--base-url URL]
[--key-env VAR] [--llm-max-tokens N] — how a tuned adapter served
by vLLM/SGLang/Ollama is graded; --baseline RUN [--tolerance N]
re-accepts against a recorded run within N percentage points,
the Verify gate's own regression rule), journals each case under an
eval- run id, and records the summary `areev loop apply
--gating-run` loads.
--case-ns NS puts the evalset and every case's input and output in
NS instead of agent:harness, so confidential cases are governed by
the namespace they came from; the summary stays in the harness.
A grader may report field metrics and usage out of band by writing
a JSON object with `metrics` and `usage` members to
$AREEV_EVAL_REPORT — each metric's MEAN lands in the summary as
<name> beside <name>_n, which the Verify gate reads like any
host-defined field, and usage sums into input_tokens /
output_tokens / usd_micros
tool provenance <hash> [--depth N] one-command code forensics: the
recommendations targeting this code, each transition's approver +
BECAUSE + gating edge, the runs that touched it, and the executor
blob its executor_uri points at (present? how many bytes? — read
lock-free, so it answers while the run is still holding the file)
run <start|resume|respond|input|pause|cancel|list|inspect|verify|fork|
shadow|oversight-report|demo> the governed
workflow runtime: journaled, checkpointed, HITL-pausable runs of
Workflow grains. list [--last N] [--offset N] prints the newest
runs with outcome + spend (default 20; stderr says when the page
truncated). An explicit --ns scopes it to that session namespace;
without one, every namespace is listed.
start --workflow HASH --run-id ID [--input JSON]
[--tool-cmd CMD] [--model provider:name] [--base-url URL]
[--key-env VAR] [--llm-max-tokens N] [--max-effects N]
[--llm-tool-result-chars N] [--llm-context-tokens N]
[--events] [--otel-endpoint http://HOST:4318]
[--as PRINCIPAL] [--max-tokens N --max-usd F ...]
[--allow-executor ADDR,...] [--executor-cache DIR]
[--sandbox-cmd 'areev-sandbox'] [--executor-timeout SECS]
[--tool-env VAR,...]
[--credential NAME=ENV_VAR[@PRINCIPAL],...] [--allow-host URL,...]
[--tool-egress TOOL:CRED[@HOST]+...:METHOD+METHOD,...]
[--credential-ttl SECS] [--resolver-env VAR,...]
[--initiator WHO] [--allow-confirmation-asks]
[--input-placement manifest|run-ns] [--harness-ns NS]
[--max-run-effects N] [--max-tool-calls N]
[--max-concurrent N] [--max-concurrent-per-principal N]
[--lease SECS] [--node ID]
[--llm-token-field max_tokens|max_completion_tokens]
[--llm-no-temperature] [--llm-extra-body JSON];
--initiator names who a run is started ON BEHALF OF: frozen in the
manifest, and an approval ask refuses them as well as the run's
principal. A resume refuses a model or scheduler-epoch mismatch
(RUN-E025/E026) before taking the lease — `run fork` is the way
through. --input-placement run-ns and --harness-ns move the run's
content-bearing records into the run's own namespace, so they are
governed by it rather than by the memory-wide agent:harness.
Every one of these has an $AREEV_RUN_* twin, because a run started
from a cron line is configured out of band.
--credential/--allow-host/--tool-egress broker a tool's outbound
calls: it gets the broker's address and a capability token, never
the secret. A tool with no grant gets nothing, and a grant naming
no method may only read. NAME=ENV_VAR@PRINCIPAL binds the
credential to a run principal, so a run executing as anyone else is
refused it — for a process holding several principals' secrets.
CRED@HOST pairs a credential with the bare hostname it may be sent
to, so a tool holding two services' secrets cannot send one to the
other. A credential can also be MINTED per call instead of read
from the environment: NAME=cmd:COMMAND takes the command's stdout
(`cmd:gcloud auth print-access-token`, `cmd:vault kv get -field=t
secret/x`) and NAME=vault:PATH#FIELD reads Vault/OpenBao directly
over $VAULT_ADDR/$VAULT_TOKEN — both cached for --credential-ttl
(default 300s), re-minted after it, and refused rather than sent
unauthenticated if the resolver fails. Bind a principal to a minted
credential on the NAME side (NAME@PRINCIPAL=cmd:...), because a
command may itself contain '@'. --resolver-env names the variables
a resolver needs for its OWN authentication ($VAULT_TOKEN,
$AWS_PROFILE): they are withheld from every other subprocess and
re-admitted only for resolvers. All five also read their
$AREEV_RUN_CREDENTIAL / _ALLOW_HOST / _TOOL_EGRESS /
_CREDENTIAL_TTL / _RESOLVER_ENV variable (the flag wins), which
is how `areev serve` and a heartbeat configure the broker;
--allow-executor pins the content address of a code-carrying tool
(a Definition whose executor_uri names a cas:// blob). Nothing
code-carrying runs unpinned, because the blob travels with the
memory and the authorization must not;
--executor-timeout overrides the fixed 300s wall-clock ceiling a
host-executed tool (--tool-cmd or a pinned --allow-executor blob)
otherwise runs under — a document-analysis leg making a dozen
model calls needs longer than that default was sized for. 0 means
wait forever;
--events streams the run's §6.10 events to stderr as JSON lines
(stdout stays the machine surface); --otel-endpoint also exports
the run as ONE OTLP/HTTP JSON trace at completion — http:// only,
terminate TLS at a local collector. Its spans carry the
OpenTelemetry GenAI attributes (gen_ai.request.model,
gen_ai.usage.*, gen_ai.tool.call.id, …) beside Areev's own
superstep/task_path/attempt provenance, so a GenAI-aware backend
reads them with no Areev-specific configuration;
--tool-env inverts how a tool's environment is decided: without it
a tool inherits this process's environment minus the variables
named to --passphrase-env/--token-env/--credential, with it the
environment is cleared and only the named variables (plus PATH and
the few a command needs to start) get through. A host that keeps
its own secrets in the environment should name what a tool sees
rather than name what it must not. A bare --tool-env passes nothing
but that minimal set. Naming a variable already registered as
holding a secret does NOT re-admit it: it is dropped and reported;
input --run-id ID --message TEXT [--as PRINCIPAL] queues a
steering message: the next superstep hands it to its nodes under
`$inbox`, so a person redirects a running run in band instead of
through a human-gate ask;
pause --run-id ID [--because TEXT] asks a live run to stop at its
next superstep boundary and park, resumably: the open superstep
finishes and checkpoints, nothing past it dispatches, and
`run resume` continues it under the same run id, manifest and
pins (no fork). Needs run.execute, the grant resume takes;
idempotent; RUN-E029 on a finished run or a pending cancel —
cancel wins, and cancel on a paused run finalizes it;
fork --run-id BASE --as-run NEW [--at N] [--plan HASH]
time-travels or migrates a run;
shadow [--runs a,b,c | --last N] [--plan HASH | --plan-file F]
[--reexecute pure] rehearses journaled runs. Bare, it re-drives
each under its OWN manifest and reports consistency; with a
candidate plan it re-drives them under that, answering every
effect from the journal — zero dispatches, zero writes — and
reports outcome, spend and out-of-support per run.
--reexecute pure additionally RE-RUNS the candidate's pure
wasm32-areev modules in the sandbox on the replayed input, so a
change to a tool's bytes stops rehearsing as `same`; the report
then names which terminal state keys moved (changed_keys /
added_keys / removed_keys — key paths, never values) and counts
sandbox_executions. Native, wasm32-areev-io, client and abstract
nodes still answer from the journal and are listed under
not_reexecuted with the reason, and a module still runs only if
this host --allow-executor pinned it and configured --sandbox-cmd.
`areev run demo` seeds the 10-minute proof
run-trace --run-id ID [--limit N] what a run recorded, and what it
produced downstream (facts/lessons)
runs-touching --hash H [--depth N] which runs produced or refined a grain —
and, for a tool definition, which ran it
(walks provenance both ways)
verify [--attestations] integrity + content-address recheck
(--attestations also checks every
attestation against --trusted-authors)
attest <hash> | --all [--ns PREFIX] sign a stored grain with the author
key (--signing-key-env); --all retro-
fills a memory that predates its key
stats store counters
serve --mcp [--ns NS] [--mount alias=path|DSN,...] [--no-destructive-ops] [--lock-ns NS] [--profile memory|full] MCP server on stdio
(--mount adds read-only memories for
cross-file ASSEMBLE; ns \"alias.inner\".
A target is a memory FILE or a postgres
DSN (postgres://…?schema=NAME) — either
backend, always opened read-only, so a
missing file is refused rather than
created and a SELECT-only pg role is
enough. Commas separate mounts only
before another alias=, so a multi-host
DSN survives. A mount's vector leg uses
ITS OWN embedder, and --embed-cmd
installs one on the primary only, so
ABOUT/NOVELTY do not cross a mount;
--profile memory drops the run/loop
tool family, default full)
repl [--ns NS] interactive CAL console in the terminal
remember --content TEXT [--facts JSON] [--observer ID]
[--session-id ID] [--run-id ID] [--role user|assistant|system|tool]
[--model SPEC | --llm-cmd CMD] [--extract-hint TEXT]
[--ground-model SPEC | --ground-cmd CMD]
[--min-confidence F] [--dry-run]
store free text as an Event, then
attach the facts distilled from it.
--facts is host-supplied and skips the
model; --model/--llm-cmd extract with an
LLM (stamped verification_status
\"unverified\" unless --ground-* checks
them). --dry-run prints, stores nothing.
hook claude-code print settings.json hook snippet
(auto recall-before-prompt, capture-on-stop,
and capture-before-compaction). Prints only —
areev never edits your settings
capture-stop [--policy-version V] the Stop AND PreCompact handler: reads
Claude Code hook JSON on stdin and stores
the transcript's turns as thread-indexed
Events. Idempotent by content address, so
both events can point at it. A compaction
summary is tagged compact_summary; a
missing transcript exits 0 silently — a
hook must never cost you a compaction
anonymize scan [--text T] [--policy-file F]
detect sensitive spans (Tier-0 chain:
structural, regex+checksum, secrets,
keyword cues, dictionary). Reads stdin
when --text is absent; prints JSON.
Pure text — needs no memory file.
anonymize test --fixtures F [--policy-file F]
Assert a policy against a fixture file
of must-redact / must-not-redact strings.
Exits non-zero on any miss or false
positive — run it in CI.
anonymize set --ns NS (--policy-file F | --policy JSON)
declare a per-namespace egress/audit
policy (travels with the file; stamps
min_reader_version). `list` shows the
declared policies, `clear --ns NS`
removes one, `mappings` prints this
process's live pseudonym mappings, and
`reveal --ns NS --token T` reverse-looks
one up (admin verb; Tier-2 audited by
fingerprint). Add --anonymize-egress to
any verb to force egress on as a host
floor, and --anonymize-cmd 'CMD' to
install a Tier-1 NER detector (JSON
probe/detect over stdin/stdout).
decide --state <TEXT|@FILE> (--questions <JSON|@FILE>
| --noul \"INSTRUCTIONS\"
| --choice \"INSTRUCTIONS\" --option KEY=DESC [--option ...]
| --score \"INSTRUCTIONS\" --level DESC [--level ...])
ask the decision chain (--decide /
--decide-cmd, below) typed questions
about a state; prints the wire response
plus provider, calibrated, latency_ms.
The shorthands ask one question, id
\"q\"; levels go lowest first. --state
is used as JSON when it parses. No --db:
it names no memory. A 429 with
Retry-After is retried once (<= 60s)
auth mint|list|revoke --auth FILE
manage the credential map (no --db: it
names no memory). mint --id NAME
--principal NAME [--label TEXT]
[--expires 90d|<ISO-8601>]
[--memories A,B] prints a 256-bit
areev_pat_ token ONCE on stdout and
stores only its SHA-256; revoke --id
NAME drops one credential without
touching the principal's others
(restart `areev ui` to apply)
provision [--check [--format json]] --db DSN [--schema NAME]
[--meta-schema NAME] [--telemetry off|aggregate|aggregate-hashed|full]
create/migrate a postgres memory's
schema AHEAD of use, so no request ever
pays for bootstrap DDL. Postgres only
(a file memory bootstraps itself at
open). Run it with a role that owns the
schema; the runtime role then needs no
CREATE, and can pin that with
`?provision=never` on its own DSN, which
refuses (STO-E008) instead of issuing
DDL. --schema supplies or overrides the
DSN's ?schema=; --meta-schema supplies
or overrides ?meta_schema=, the PAIRED
layout that keeps the engine's metadata
(meta, counters, ns_reg, telem_*) in a
second, physically separate schema —
every later DSN must then name both
(a mismatched layout is refused,
STO-E011). It also provisions the
telemetry tables, since --telemetry
defaults to aggregate on every verb —
pass --telemetry off to skip them.
--check is the READ-ONLY probe: SELECTs
only, no lock, no DDL, no meta write, so
it runs under a least-privilege role. It
reports every stamp's found-vs-wanted
value, what is pending, and
rolling_deploy: safe |
drain_writers_first | unknown. Exit 0
current, 2 pending or absent
ui [--addr HOST:PORT] [--allow-remote] [--allow-origin ORIGIN[,ORIGIN...]]
[--token-env VAR] [--no-destructive-ops] [--tls-cert PATH --tls-key PATH]
[--sso-header NAME --sso-secret-env VAR [--sso-secret-env-next VAR]]
[--sso-approvals deny|allow] [--sso-groups-header NAME]
[--sso-principal-prefix PREFIX]
[--oidc-issuer URL --oidc-client-id ID --oidc-client-secret-env VAR
--oidc-redirect-uri URL [--oidc-scopes S] [--oidc-principal-claim C]]
web console (default 127.0.0.1:7437).
--allow-origin accepts a cross-origin
POST from an EXACT origin (comma-
separated for more than one; no
wildcards) — the Origin check is not
lifted by --allow-remote, so a console
behind a reverse proxy needs its public
origin named here to be usable.
--sso-secret-env-next opens a rotation
window: both secrets prove the proxy, so
the fleet moves over one node at a time
(docs/runbooks/sso-secret-rotation.md)
--sso-approvals (default deny) decides
whether a proxy-asserted identity may
answer a HITL approval: its only proof
is the shared proxy secret, so 'allow'
makes every approval only as strong as
that one value. Approvers should hold a
per-principal credential (--auth).
--sso-groups-header maps IdP groups to
principals via --auth's \"groups\" table
(needs --auth); an identity with its own
grants outranks its groups, and a
group-derived principal is a ROLE — it
may never approve, under any setting.
--sso-principal-prefix stamps every
proxy-asserted principal so IdP names
stay distinct from local ones.
--oidc-* (build with --features oidc)
runs the auth-code+PKCE flow in-process
instead: the console verifies the IdP's
signature against its published key set
and issues an HttpOnly SameSite=Strict
session cookie. Login at /auth/login,
logout at /auth/logout. An OIDC identity
MAY approve (no --sso-approvals needed):
a verified signature is a stronger claim
than a shared proxy secret. The redirect
URI must match the provider's exactly.
Namespace defaults to \"shared\". Exit code 0 on success.
--db is optional for one-shot commands: it falls back to $AREEV_DB, then
~/.areev/default.db. `serve` and `ui` require an explicit --db (or $AREEV_DB).
Files carry their own settings (text index, entity relations, embedding
provenance) in an internal meta table; a bare open honors them.
--index-text true|false explicitly re-stamps the file's declaration.
--assembly-manifest-sample-rate 0..1 samples ASSEMBLE provenance into
agent:harness (host-only, default 0/off).
Principals: add --as <principal> to cal/repl/serve/subject-report/forget-subject/
purge-older-than to run the session under the file's grants for that
principal (fail closed; see GRANT in docs/cal-reference.md §9). --auth FILE
validates the principal against a credential map, and `areev ui --auth FILE`
serves the console in multi-principal mode (tokens → principals,
unauthenticated = read-only \"anonymous\").
Console auth: prefer --auth <map> — per-principal credentials, so a write or
an approval names WHO did it, one credential can be revoked without disturbing
the principal's others, and `expires_at` bounds a leak. `areev auth mint`
creates them. --token-env <VAR> remains the single-user shortcut: one shared
secret, implied admin, unattributable (it can never answer a HITL approval).
On that path the token's entropy is the ONLY control — use at least 32 hex
characters, or let `areev auth mint` choose; comparison against the presented
credential is constant-time. After 10 consecutive failed authentications from
one source address, further credential-bearing requests from it are refused
with 429 until the streak goes idle (15 min). Browsers authenticate via HTTP Basic, which the
browser caches per origin and RE-ATTACHES AUTOMATICALLY, including on
cross-site requests — so the Origin check (--allow-origin, above) is
load-bearing, not defence in depth. There is no logout: closing the tab does
not clear a cached Basic credential, so a browser used to demo the console
stays authenticated until it is restarted or its saved credentials are
cleared. A wrong token is counted per source IP and logged to stderr as
\"console auth FAILED from <ip> (<n> consecutive)\" (greppable for a
fail2ban-style rule; the token itself is never logged) — but it is NOT
delayed or locked out in-process: `ui` serves one connection at a time, so a
sleep before responding would let an unauthenticated caller stall the
console for everyone, and behind a reverse proxy every request logs the
PROXY's IP anyway, so an IP lockout here would lock out every user at once.
Rate limiting belongs at that proxy (the documented default deployment
shape) — write the fail2ban-style rule against its access log, not this
one. A SameSite/HttpOnly/Secure session cookie is the intended future fix
for the caching/logout problem above and is tracked separately — it does
not exist yet.
Encryption at rest: add --passphrase-env <VAR> to any command to derive an
AES-256 key (Argon2id) from the passphrase in environment variable VAR. The
non-secret salt is kept in a <memory.db>.kdf sidecar — back it up with the db.
Anonymization key: add --anon-key-env <VAR> to any command to supply the
32-byte root (64 hex characters in VAR) the pseudonym session/memory/vault
subkeys are derived from. Independent of the page cipher, so the mapping vault
and value-derived tokens also work on postgres and on plaintext files. Never
persisted; rotating it is a crypto-erasure of the mapping table.
Grain attestation: add --signing-key-env <VAR> (a 32-byte Ed25519 seed, 64
hex characters in VAR) to any command and every grain it writes is followed by
an attestation — a detached signature over the grain's content hash, stored as
an Observation in agent:attest. Add --trusted-authors FILE (a JSON map of
key_id -> public key, with an optional policy of off | verify | require) to
check attestations at import/follow and in `verify --attestations`; add
--require-attested to refuse any bundle carrying a grain without a valid
attestation from a trusted key. Keys are host config, never written to the
memory; the attested grain's hash never changes. docs/grain-attestation-plan.md
Vector recall: add --embed-cmd 'CMD' [--embed-model NAME] to any command to
install a command embedder — CMD gets the text on stdin and must print a JSON
array of numbers. Turns on the vector leg for search/serve, and embeds grains
written by add/remember/migrate.
Decision backends (optional, off unless configured): add --decide <CHAIN>
to any command — a comma-separated, ordered list of providers
(typesafe:MODEL, openrouter:MODEL, vercel:MODEL, openjev:MODEL,
cloudflare:MODEL, systemone:URL[#MODEL], llm:LLM-SPEC), tried in order;
only the first colon splits, so URLs and nested LLM specs pass through.
--decide-cmd 'CMD' appends a command backend as the LAST entry (wire request
JSON on stdin, wire response JSON on stdout; no shell). --decide-timeout-ms N
bounds the whole chain (default 2000); llm: entries generate text and usually
need a larger value. Each falls back to $AREEV_DECIDE, $AREEV_DECIDE_CMD,
$AREEV_DECIDE_TIMEOUT_MS. A bad spec or missing provider key fails at
startup (DEC-E001). A decision backend may score and order, never omit;
llm: entries are uncalibrated. docs/decision-model-proposal.md
Reranking and recall deadline: with --decide set, the chain is installed as
the recall reranker (each candidate scored for relevance, then reordered).
--rerank-cmd 'CMD' [--rerank-model NAME] installs a command reranker instead
(stdin {\"query\": \"...\", \"docs\": [...]}, stdout a JSON array of
docs.len() numbers) — an explicit command wins over --decide. search and
recall-hook rerank whenever one is installed; CAL does under WITH rerank.
--recall-deadline-ms N bounds each hybrid recall (0 = none, the default;
recall-hook defaults to 1500 when a reranker is installed); past it recall
fails open to what it gathered. Env: $AREEV_RERANK_CMD,
$AREEV_RECALL_DEADLINE_MS. Under an egress anonymization policy (or
--anonymize-egress) what a decision backend receives is pseudonymized.
Read-only: add --read-only to any command to open the memory refusing every
write (add/supersede/forget/reindex/blob-put and so on fail with a coded
STO-E004 error); reads and CAL SELECT work normally. On postgres this also
skips schema/table bootstrap, so a role holding only CONNECT + USAGE + SELECT
grants can open an existing, already-migrated memory — see
docs/deployment-profile.md for the grant recipe. --read-only combined with an
explicit --index-text is refused up front (--index-text always re-stamps the
file's declaration, which read-only mode cannot do).
Recall telemetry: add --telemetry off|aggregate|aggregate-hashed|full to any
command. Default aggregate (off under --read-only unless asked for
explicitly). It is HOST config, never a file-truth. aggregate-hashed keeps the
same rollups with the query key HMAC'd under a subkey derived from the
memory's own AEAD key, keeps no sample and writes no ring log — so what a
person typed is never retained and there is nothing to erase later. full adds
a per-recall ring log. `areev telemetry scrub --ns NS --yes` reaches the rows
a namespace erasure cannot.";
fn flag(args: &HashMap<String, String>, k: &str) -> Option<String> {
args.get(k).cloned()
}
fn assembly_manifest_sample_rate(flags: &HashMap<String, String>) -> Result<f64, String> {
let Some(raw) = flag(flags, "assembly-manifest-sample-rate") else {
return Ok(0.0);
};
let rate: f64 = raw
.parse()
.map_err(|_| "--assembly-manifest-sample-rate must be a number from 0 to 1".to_string())?;
if !rate.is_finite() || !(0.0..=1.0).contains(&rate) {
return Err("--assembly-manifest-sample-rate must be from 0 to 1".into());
}
Ok(rate)
}
fn short_flag(a: &str) -> Option<&'static str> {
match a {
"-d" => Some("db"),
"-k" => Some("k"),
_ => None,
}
}
fn need(args: &HashMap<String, String>, k: &str) -> Result<String, String> {
flag(args, k).ok_or_else(|| format!("missing required --{k}"))
}
fn parse_anon_key(hex: &str) -> Result<[u8; 32], String> {
if hex.len() != 64 {
return Err(format!(
"expected 64 hex characters (a 32-byte key), got {}",
hex.len()
));
}
let mut out = [0u8; 32];
for (i, byte) in out.iter_mut().enumerate() {
let pair = &hex[i * 2..i * 2 + 2];
*byte = u8::from_str_radix(pair, 16)
.map_err(|_| format!("not hex: {pair:?} at byte {i}"))?;
}
Ok(out)
}
#[cfg(feature = "postgres")]
fn open_postgres_store(
db: &str,
tel_mode: areev_store::TelemetryMode,
explicit_index: Option<&str>,
anon_key: Option<[u8; 32]>,
read_only: bool,
) -> Result<Areev, String> {
let (url, schema) = areev_store::pg::split_schema_url(db).map_err(|e| e.to_string())?;
if explicit_index.is_some() || anon_key.is_some() || read_only {
let o = areev_store::AreevOptions {
index_text: explicit_index
.map(|v| !matches!(v, "false" | "0" | "off" | "no"))
.unwrap_or(areev_store::AreevOptions::default().index_text),
anon_key,
telemetry: tel_mode,
read_only,
..Default::default()
};
Areev::open_postgres_with(&url, &schema, o)
} else if tel_mode != areev_store::TelemetryMode::Off {
Areev::open_postgres_with_telemetry(&url, &schema, tel_mode)
} else {
Areev::open_postgres(&url, &schema)
}
.map_err(|e| e.to_string())
}
#[cfg(not(feature = "postgres"))]
fn open_postgres_store(
_db: &str,
_tel_mode: areev_store::TelemetryMode,
_explicit_index: Option<&str>,
_anon_key: Option<[u8; 32]>,
_read_only: bool,
) -> Result<Areev, String> {
Err("this build lacks the postgres backend — reinstall with \
`cargo install areev --features postgres-tls` (or build with \
--features postgres-tls), or grab the `-postgres` asset from a \
GitHub release, which ships the backend already compiled in"
.into())
}
fn resolve_db(args: &HashMap<String, String>, require_explicit: bool) -> Result<String, String> {
if let Some(p) = flag(args, "db") {
return Ok(p);
}
if let Ok(p) = std::env::var("AREEV_DB") {
if !p.trim().is_empty() {
return Ok(p);
}
}
if require_explicit {
return Err(
"this command needs an explicit memory file — pass --db <file> (or -d), \
or set $AREEV_DB. It will not fall back to the personal default memory \
(~/.areev/default.db), to avoid serving the wrong file."
.to_string(),
);
}
let home = std::env::var("HOME")
.or_else(|_| std::env::var("USERPROFILE"))
.map_err(|_| "no --db given, and neither $AREEV_DB nor $HOME is set".to_string())?;
let path = format!("{home}/.areev/default.db");
if let Some(parent) = std::path::Path::new(&path).parent() {
std::fs::create_dir_all(parent)
.map_err(|e| format!("cannot create {}: {e}", parent.display()))?;
}
eprintln!("areev: using default memory {path} (override with -d/--db or $AREEV_DB)");
Ok(path)
}
fn resolve_decider(
flags: &HashMap<String, String>,
) -> Result<Option<std::sync::Arc<dyn areev_llm::DecisionBackend>>, String> {
let spec = run_stack::flag_or_env(flags, "decide", "AREEV_DECIDE");
let cmd = run_stack::flag_or_env(flags, "decide-cmd", "AREEV_DECIDE_CMD");
let deadline = match run_stack::flag_or_env(flags, "decide-timeout-ms", "AREEV_DECIDE_TIMEOUT_MS") {
None => areev_llm::decide::DEFAULT_DECIDE_TIMEOUT,
Some(ms) => match ms.parse::<u64>() {
Ok(n) if n > 0 => std::time::Duration::from_millis(n),
_ => {
return Err(areev_llm::DecideError::NotConfigured(format!(
"--decide-timeout-ms (or $AREEV_DECIDE_TIMEOUT_MS) must be a positive whole \
number of milliseconds, got {ms:?}"
))
.to_string())
}
},
};
for (key, v) in [("decide", &spec), ("decide-cmd", &cmd)] {
if flags.get(key).map(String::as_str) == Some("true") && v.is_some() {
return Err(areev_llm::DecideError::NotConfigured(format!("--{key} needs a value")).to_string());
}
}
areev_llm::resolve_chain(spec.as_deref(), cmd.as_deref(), Some(deadline)).map_err(|e| e.to_string())
}
fn decider_for_egress(
chain: std::sync::Arc<dyn areev_llm::DecisionBackend>,
egress: bool,
) -> Result<std::sync::Arc<dyn areev_llm::DecisionBackend>, String> {
if !egress {
return Ok(chain);
}
let policy = areev_core::anon::AnonPolicy { scope: "session".into(), ..Default::default() };
Ok(std::sync::Arc::new(
areev_llm::PseudonymizingDecider::new(chain, policy).map_err(|e| e.to_string())?,
))
}
fn store_egress_active(m: &Areev) -> bool {
m.anonymize_egress_floor() || m.anon_declared().iter().any(|(_, mode)| mode == "egress")
}
fn flag_values(argv: &[String], name: &str) -> Vec<String> {
let long = format!("--{name}");
let mut out = Vec::new();
let mut i = 0;
while i < argv.len() {
if argv[i] == long {
if let Some(v) = argv.get(i + 1).filter(|v| !v.starts_with('-')) {
out.push(v.clone());
i += 2;
continue;
}
}
i += 1;
}
out
}
fn text_or_file(flag_name: &str, v: &str) -> Result<String, String> {
match v.strip_prefix('@') {
Some(path) => std::fs::read_to_string(path).map_err(|e| format!("--{flag_name} @{path}: {e}")),
None => Ok(v.to_string()),
}
}
const DECIDE_RETRY_AFTER_CAP_SECS: u64 = 60;
fn run_decide(
decider: Option<std::sync::Arc<dyn areev_llm::DecisionBackend>>,
flags: &HashMap<String, String>,
argv: &[String],
) -> Result<(), String> {
use areev_llm::decide::{questions_from_wire, Question};
use areev_llm::{DecideError, DecideRequest};
use std::collections::BTreeMap;
let decider = decider.ok_or_else(|| {
DecideError::NotConfigured(
"areev decide needs a backend — pass --decide <chain> and/or --decide-cmd <cmd> \
(or set $AREEV_DECIDE / $AREEV_DECIDE_CMD)"
.into(),
)
.to_string()
})?;
let decider = decider_for_egress(decider, flags.contains_key("anonymize-egress"))?;
let raw_state = text_or_file("state", &need(flags, "state")?)?;
let state = match serde_json::from_str::<serde_json::Value>(&raw_state) {
Ok(v) if v.is_string() || v.is_object() || v.is_array() => v,
_ => serde_json::Value::String(raw_state),
};
let forms: Vec<&str> = ["questions", "noul", "choice", "score"]
.into_iter()
.filter(|k| flags.contains_key(*k))
.collect();
if forms.len() != 1 {
return Err(format!(
"areev decide needs exactly one of --questions, --noul, --choice or --score (got {})",
if forms.is_empty() { "none".to_string() } else { forms.iter().map(|f| format!("--{f}")).collect::<Vec<_>>().join(", ") }
));
}
let one = |q: Question| BTreeMap::from([("q".to_string(), q)]);
let questions: BTreeMap<String, Question> = match forms[0] {
"questions" => {
let raw = text_or_file("questions", &need(flags, "questions")?)?;
let v: serde_json::Value = serde_json::from_str(&raw)
.map_err(|e| DecideError::InvalidQuestion(format!("--questions is not JSON: {e}")).to_string())?;
questions_from_wire(&v).map_err(|e| e.to_string())?
}
"noul" => one(Question::noul(need(flags, "noul")?)),
"choice" => {
let mut options = Vec::new();
for o in flag_values(argv, "option") {
let (k, d) = o.split_once('=').ok_or_else(|| {
format!("--option {o:?}: expected KEY=DESCRIPTION (e.g. --option billing=\"money owed\")")
})?;
let k = k.trim().to_string();
if options.iter().any(|(seen, _)| *seen == k) {
return Err(DecideError::InvalidQuestion(format!("--option {k:?} is given twice")).to_string());
}
options.push((k, d.trim().to_string()));
}
one(Question::choice(need(flags, "choice")?, options))
}
"score" => one(Question::score(need(flags, "score")?, flag_values(argv, "level"))),
_ => unreachable!("forms holds only the four names above"),
};
let req = DecideRequest::new(state, questions);
let decision = match decider.decide(&req) {
Ok(d) => d,
Err(e) => {
let retry_after = match &e {
DecideError::ChainExhausted(errs) => errs.iter().filter_map(|(_, e)| e.retry_after_secs()).min(),
other => other.retry_after_secs(),
};
match retry_after {
Some(secs) => {
let secs = secs.min(DECIDE_RETRY_AFTER_CAP_SECS);
eprintln!("areev: {e} — retrying once in {secs}s");
std::thread::sleep(std::time::Duration::from_secs(secs));
decider.decide(&req).map_err(|e| e.to_string())?
}
None => return Err(e.to_string()),
}
}
};
println!(
"{}",
serde_json::to_string_pretty(&decision.to_json()).map_err(|e| e.to_string())?
);
Ok(())
}
fn apply_principal(
facade: areev_cal::AreevFacade,
flags: &std::collections::HashMap<String, String>,
) -> Result<areev_cal::AreevFacade, String> {
let Some(p) = flag(flags, "as") else {
return Ok(facade);
};
if let Some(path) = flag(flags, "auth") {
let text = std::fs::read_to_string(&path).map_err(|e| format!("--auth {path}: {e}"))?;
let map = areev_core::authz::CredentialMap::from_json(&text).map_err(|e| e.to_string())?;
map.knows_principal(&p).map_err(|e| e.to_string())?;
}
let bound = facade.with_principal(&p).map_err(|e| e.to_string())?;
eprintln!("areev: session bound to principal '{p}' (rights from the file's grants)");
Ok(bound)
}
fn write_erasure_audit(
m: &mut areev_store::Areev,
flags: &HashMap<String, String>,
verb: &str,
target: &str,
count: usize,
stale_exports: &[areev_store::CorpusExportRegistry],
overridden: Option<&areev_store::HoldRecord>,
) {
let principal = flag(flags, "as").unwrap_or_else(|| "local:owner".to_string());
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
let because = flag(flags, "because");
let mut obs = areev_core::authz::audit_observation(
&principal,
verb,
target,
because.as_deref(),
count,
now,
);
if !stale_exports.is_empty() {
let mut context = obs
.common
.context
.take()
.and_then(|value| value.as_object().cloned())
.unwrap_or_default();
context.insert(
"stale_corpora".into(),
serde_json::to_value(stale_exports).unwrap_or_else(|_| serde_json::json!([])),
);
let mut stale_adapters = Vec::new();
for export in stale_exports {
let Ok(h) = Hash::from_hex(&export.manifest_hash) else {
continue;
};
let kids = m.grains_derived_from(&h).unwrap_or_default();
for g in kids {
if g.get_str("relation") != Some("mg:adapter") {
continue;
}
stale_adapters.push(serde_json::json!({
"subject": g.get_str("subject"),
"grain_hash": g.hash.to_hex(),
"export_id": export.export_id,
}));
}
}
if !stale_adapters.is_empty() {
context.insert("stale_adapters".into(), serde_json::json!(stale_adapters));
}
obs.common.context = Some(serde_json::Value::Object(context));
for export in stale_exports {
let recipient = export
.recipient
.as_deref()
.map(|id| format!(" for recipient {id}"))
.unwrap_or_default();
eprintln!(
"areev: corpus {} at {}{} is stale and must be retired or re-derived",
export.export_id, export.destination, recipient
);
}
for a in &stale_adapters {
eprintln!(
"areev: adapter {} (grain {}) derives from stale corpus {} — retire or re-derive",
a["subject"].as_str().unwrap_or("?"),
a["grain_hash"].as_str().unwrap_or("?"),
a["export_id"].as_str().unwrap_or("?"),
);
}
}
if let Some(hold) = overridden {
let ctx = obs.common.context.get_or_insert_with(|| serde_json::json!({}));
if let Some(o) = ctx.as_object_mut() {
o.insert(
"hold_overridden".into(),
serde_json::json!({
"ns": hold.ns,
"placed_by": hold.placed_by,
"because": hold.because,
}),
);
}
}
if let Err(e) = m.append_audit(&mut obs) {
eprintln!("areev: warning: erasure succeeded but its audit record failed to write: {e}");
}
}
fn report_meta_warnings(facade: &areev_cal::AreevFacade) {
use areev_cal::CalStoreFacade;
let _ = facade.list_queries();
let _ = facade.list_templates();
for w in facade.meta_warnings() {
eprintln!("areev: warning: {w}");
}
}
#[cfg(feature = "postgres")]
fn mount_is_the_primary(primary: &str, mount: &str) -> bool {
if !(is_pg_dsn(primary) && is_pg_dsn(mount)) {
return false;
}
let norm = |d: &str| {
areev_store::pg::split_schema_url(&areev_store::pg::strip_provision(d)).ok()
};
match (norm(primary), norm(mount)) {
(Some(a), Some(b)) => a == b,
_ => false,
}
}
#[cfg(not(feature = "postgres"))]
fn mount_is_the_primary(_primary: &str, _mount: &str) -> bool {
false
}
fn redact_dsn(label: &str) -> String {
areev_server::redact_dsn(label)
}
fn starts_a_mount_entry(s: &str) -> bool {
let Some((name, _)) = s.split_once('=') else { return false };
!name.is_empty()
&& name
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '.' | '-'))
}
fn parse_mounts(spec: Option<&str>) -> Result<Vec<(String, String)>, String> {
let Some(spec) = spec else { return Ok(Vec::new()) };
let mut entries: Vec<&str> = Vec::new();
let mut start = 0usize;
for (i, c) in spec.char_indices() {
if c == ',' && starts_a_mount_entry(spec[i + 1..].trim_start()) {
entries.push(&spec[start..i]);
start = i + 1;
}
}
entries.push(&spec[start..]);
let mut out = Vec::new();
for entry in entries.into_iter().map(str::trim).filter(|e| !e.is_empty()) {
let (alias, target) = entry
.split_once('=')
.map(|(a, t)| (a.trim(), t.trim()))
.filter(|(a, t)| !a.is_empty() && !t.is_empty())
.ok_or_else(|| {
format!("--mount expects alias=<path|DSN>, got '{}'", redact_dsn(entry))
})?;
if alias.is_empty()
|| !alias.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-'))
{
return Err(format!(
"--mount alias {:?} is not a name: an alias is [A-Za-z0-9_-]+ and becomes a \
namespace prefix (\"alias.inner\"). Write --mount <alias>=<path|DSN>",
redact_dsn(alias)
));
}
out.push((alias.to_string(), target.to_string()));
}
Ok(out)
}
fn open_mount(target: &str) -> Result<Areev, String> {
if is_pg_dsn(target) {
return open_postgres_store(target, areev_store::TelemetryMode::Off, None, None, true);
}
let opts = areev_store::AreevOptions { read_only: true, ..Default::default() };
Areev::open_with(target, opts).map_err(|e| {
let mut msg = e.to_string();
if e.code() == "STO-E004" {
msg.push_str(
" — a mount is opened read-only, and a read-only open cannot re-stamp a \
declaration. Open this memory read-write once with the settings it should \
declare, then mount it",
);
}
msg
})
}
fn addr_is_loopback(addr: &str) -> bool {
let a = addr.trim();
if let Some(rest) = a.strip_prefix('[') {
let host = rest.split(']').next().unwrap_or("");
return host == "::1" || host.eq_ignore_ascii_case("localhost");
}
let host = a.rsplit_once(':').map(|(h, _)| h).unwrap_or(a);
if host.eq_ignore_ascii_case("localhost") {
return true;
}
match host.parse::<std::net::IpAddr>() {
Ok(ip) => ip.is_loopback(),
Err(_) => false,
}
}
fn parse_args(rest: &[String]) -> (HashMap<String, String>, Vec<String>) {
let mut flags = HashMap::new();
let mut positional = Vec::new();
let mut i = 0;
while i < rest.len() {
let a = &rest[i];
if let Some(name) = a.strip_prefix("--") {
if i + 1 < rest.len() && !rest[i + 1].starts_with('-') {
flags.insert(name.to_string(), rest[i + 1].clone());
i += 2;
} else {
flags.insert(name.to_string(), "true".to_string());
i += 1;
}
} else if let Some(long) = short_flag(a) {
if i + 1 < rest.len() && !rest[i + 1].starts_with('-') {
flags.insert(long.to_string(), rest[i + 1].clone());
i += 2;
} else {
flags.insert(long.to_string(), "true".to_string());
i += 1;
}
} else {
positional.push(a.clone());
i += 1;
}
}
(flags, positional)
}
fn transcript_content_blocks(
content: &serde_json::Value,
turn: usize,
) -> (String, Vec<ContentBlock>, Vec<serde_json::Value>) {
const TOOL_RESULT_CHARS: usize = 32 * 1024;
fn truncate(s: &str, max: usize) -> (String, bool) {
if s.chars().count() <= max {
(s.to_string(), false)
} else {
(s.chars().take(max).collect(), true)
}
}
fn tool_result_body(content: &serde_json::Value) -> String {
match content {
serde_json::Value::String(s) => s.clone(),
serde_json::Value::Array(blocks) => blocks
.iter()
.filter_map(|b| b["text"].as_str().or_else(|| b.as_str()))
.collect::<Vec<_>>()
.join("\n"),
other if other.is_object() => other["text"].as_str().unwrap_or("").to_string(),
_ => String::new(),
}
}
let input_blocks: Vec<&serde_json::Value> = match content {
serde_json::Value::String(_) => Vec::new(),
serde_json::Value::Array(blocks) => blocks.iter().collect(),
_ => Vec::new(),
};
if let serde_json::Value::String(s) = content {
return (
s.clone(),
vec![ContentBlock::Text { text: s.clone() }],
Vec::new(),
);
}
let mut rendered = Vec::new();
let mut typed = Vec::new();
let mut elisions = Vec::new();
for (block_index, b) in input_blocks.into_iter().enumerate() {
match b["type"].as_str() {
Some("text") => {
if let Some(text) = b["text"].as_str() {
rendered.push(text.to_string());
typed.push(ContentBlock::Text { text: text.to_string() });
}
}
Some("tool_use") => {
let name = b["name"].as_str().unwrap_or("tool").to_string();
let id = b["id"]
.as_str()
.map(str::to_string)
.unwrap_or_else(|| format!("synthetic:{turn}:{block_index}"));
let input = b.get("input").cloned().unwrap_or(serde_json::Value::Null);
rendered.push(format!("[tool_use {name}] {input}"));
typed.push(ContentBlock::ToolUse { id, name, input });
}
Some("tool_result") => {
let original = tool_result_body(&b["content"]);
let original_bytes = original.len();
let (body, elided) = truncate(&original, TOOL_RESULT_CHARS);
let is_error = b.get("is_error").and_then(serde_json::Value::as_bool);
let flag = if is_error == Some(true) { " ERROR" } else { "" };
rendered.push(format!("[tool_result{flag}] {body}"));
typed.push(ContentBlock::ToolResult {
tool_use_id: b["tool_use_id"]
.as_str()
.map(str::to_string)
.unwrap_or_else(|| format!("synthetic:{turn}:{block_index}")),
content: body,
is_error,
});
if elided {
elisions.push(serde_json::json!({
"block_index": block_index,
"kind": "tool_result",
"original_bytes": original_bytes,
"kept_chars": TOOL_RESULT_CHARS,
}));
}
}
_ => {
if let Some(text) = b["text"].as_str() {
rendered.push(text.to_string());
typed.push(ContentBlock::Text { text: text.to_string() });
}
}
}
}
(rendered.join("\n"), typed, elisions)
}
fn transcript_turn_ms(record: &serde_json::Value) -> Option<i64> {
let value = record.get("timestamp").or_else(|| record.get("created_at"))?;
if let Some(ms) = value.as_i64() {
return Some(ms);
}
let text = value.as_str()?;
if let Ok(ms) = text.parse::<i64>() {
return Some(ms);
}
chrono::DateTime::parse_from_rfc3339(text)
.ok()
.map(|dt| dt.timestamp_millis())
}
fn transcript_token_usage(message: &serde_json::Value) -> Option<TokenUsage> {
let usage = message.get("usage")?.as_object()?;
let u32_field = |primary: &str, alias: &str| {
usage
.get(primary)
.or_else(|| usage.get(alias))
.and_then(serde_json::Value::as_u64)
.map(|n| n.min(u32::MAX as u64) as u32)
};
Some(TokenUsage {
input_tokens: u32_field("input_tokens", "input").unwrap_or(0),
output_tokens: u32_field("output_tokens", "output").unwrap_or(0),
cache_read_tokens: u32_field("cache_read_tokens", "cache_read_input_tokens"),
cache_creation_tokens: u32_field(
"cache_creation_tokens",
"cache_creation_input_tokens",
),
})
}
fn run_migrate(
m: &mut Areev,
ns: &str,
from: &str,
file: &str,
history: Option<String>,
) -> Result<areev_store::migrate::MigrateReport, String> {
use areev_store::migrate as mig;
let read = |p: &str| std::fs::read_to_string(p).map_err(|e| format!("{p}: {e}"));
let parse = |p: &str, s: &str| {
serde_json::from_str::<serde_json::Value>(s).map_err(|e| format!("{p}: bad JSON: {e}"))
};
let rep = match from {
"mem0" => {
let export = parse(file, &read(file)?)?;
let history = match &history {
Some(h) => Some(parse(h, &read(h)?)?),
None => None,
};
mig::migrate_mem0(m, ns, Some(&export), history.as_ref())
}
"mem0-history" => {
let h = parse(file, &read(file)?)?;
mig::migrate_mem0(m, ns, None, Some(&h))
}
"langgraph" | "langmem" => mig::migrate_langgraph(m, ns, &read(file)?),
"letta" => {
let af = parse(file, &read(file)?)?;
mig::migrate_letta(m, ns, &af)
}
"letta-archival" => mig::migrate_letta_archival(m, ns, &read(file)?),
"zep" | "graphiti" => {
let v = parse(file, &read(file)?)?;
mig::migrate_zep(m, ns, &v)
}
"jsonl" => mig::migrate_jsonl(m, ns, &read(file)?),
"tool-log" | "openai-tools" => {
let mut rep = mig::MigrateReport::default();
for (i, line) in read(file)?.lines().enumerate() {
let line = line.trim();
if line.is_empty() {
continue;
}
let v: serde_json::Value = serde_json::from_str(line)
.map_err(|e| format!("{file}:{}: bad JSON: {e}", i + 1))?;
for (record_index, record) in extract_tool_records(&v).into_iter().enumerate() {
use sha2::{Digest, Sha256};
let synthetic_call_id = || {
let mut digest = Sha256::new();
digest.update(file.as_bytes());
digest.update([0]);
digest.update(line.as_bytes());
digest.update((i as u64).to_le_bytes());
digest.update((record_index as u64).to_le_bytes());
format!("import:{}", &hex::encode(digest.finalize())[..24])
};
let mut tool = Tool::new(&record.name)
.is_error(record.is_error)
.created_at(record.created_at.unwrap_or(0));
if let Some(input) = record.input {
tool = tool.input(input);
}
if !record.content.is_empty() {
tool = tool.content(&record.content);
}
let call_id = record.call_id.unwrap_or_else(synthetic_call_id);
tool = tool.tool_call_id(&call_id);
let tool = tool.namespace(ns);
let (_, hash) = areev_core::format::serialize::serialize_grain(&tool)
.map_err(|e| e.to_string())?;
if m.get(&hash).is_ok() {
rep.skipped += 1;
} else {
m.add(&tool).map_err(|e| e.to_string())?;
rep.added += 1;
}
}
}
Ok(rep)
}
"basic-memory" => {
let root = std::path::PathBuf::from(file);
if !root.is_dir() {
return Err(format!(
"--from basic-memory expects --file to be the vault directory (got {file})"
));
}
let mut notes: Vec<std::path::PathBuf> = Vec::new();
let mut stack = vec![root.clone()];
while let Some(d) = stack.pop() {
let entries = std::fs::read_dir(&d).map_err(|e| format!("{}: {e}", d.display()))?;
for entry in entries {
let path = entry.map_err(|e| e.to_string())?.path();
if path.is_dir() {
stack.push(path);
} else if path.extension().is_some_and(|x| x == "md") {
notes.push(path);
}
}
}
notes.sort();
let mut rep = mig::MigrateReport::default();
for path in notes {
let rel = path
.strip_prefix(&root)
.unwrap_or(&path)
.to_string_lossy()
.replace('\\', "/");
let md = read(&path.to_string_lossy())?;
let mtime_ms = std::fs::metadata(&path)
.and_then(|md| md.modified())
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_millis() as i64);
mig::migrate_basic_memory_note(m, ns, &rel, &md, mtime_ms, &mut rep)
.map_err(|e| e.to_string())?;
}
Ok(rep)
}
other => {
return Err(format!(
"unknown --from '{other}' — sources: mem0, mem0-history, langgraph, letta, \
letta-archival, zep, basic-memory, jsonl, tool-log"
))
}
};
rep.map_err(|e| e.to_string())
}
#[derive(Debug, Clone, PartialEq)]
struct ImportedToolRecord {
name: String,
input: Option<serde_json::Value>,
content: String,
is_error: bool,
call_id: Option<String>,
created_at: Option<i64>,
}
fn extract_tool_records(v: &serde_json::Value) -> Vec<ImportedToolRecord> {
let stringify = |x: &serde_json::Value| match x {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
};
let created_at = v
.get("created_at")
.or_else(|| v.get("timestamp"))
.and_then(|value| {
value
.as_i64()
.or_else(|| value.as_str().and_then(|s| s.parse::<i64>().ok()))
});
if let Some(calls) = v.get("tool_calls").and_then(|c| c.as_array()) {
return calls
.iter()
.filter_map(|c| {
let name = c
.get("function")
.and_then(|f| f.get("name"))
.and_then(|n| n.as_str())
.or_else(|| c.get("name").and_then(|n| n.as_str()))?;
let input = c
.get("function")
.and_then(|f| f.get("arguments"))
.or_else(|| c.get("input"))
.map(|value| match value {
serde_json::Value::String(raw) => serde_json::from_str(raw)
.unwrap_or_else(|_| serde_json::Value::String(raw.clone())),
other => other.clone(),
});
let is_error = c
.get("is_error")
.and_then(|e| e.as_bool())
.unwrap_or_else(|| c.get("error").is_some());
let call_id = c
.get("id")
.or_else(|| c.get("tool_call_id"))
.and_then(|id| id.as_str())
.map(str::to_string);
Some(ImportedToolRecord {
name: name.to_string(),
input,
content: String::new(),
is_error,
call_id,
created_at,
})
})
.collect();
}
let name = v
.get("tool_name")
.and_then(|n| n.as_str())
.or_else(|| v.get("name").and_then(|n| n.as_str()))
.or_else(|| v.get("function").and_then(|f| f.get("name")).and_then(|n| n.as_str()));
match name {
Some(name) => {
let content = v
.get("content")
.or_else(|| v.get("output"))
.or_else(|| v.get("result"))
.map(&stringify)
.unwrap_or_default();
let is_error = v
.get("is_error")
.and_then(|e| e.as_bool())
.unwrap_or_else(|| v.get("error").is_some());
let input = v.get("input").or_else(|| v.get("arguments")).map(|value| match value {
serde_json::Value::String(raw) => serde_json::from_str(raw)
.unwrap_or_else(|_| serde_json::Value::String(raw.clone())),
other => other.clone(),
});
let call_id = v
.get("tool_call_id")
.or_else(|| v.get("call_id"))
.or_else(|| v.get("id"))
.and_then(|id| id.as_str())
.map(str::to_string);
vec![ImportedToolRecord {
name: name.to_string(),
input,
content,
is_error,
call_id,
created_at,
}]
}
None => Vec::new(),
}
}
fn run() -> Result<(), String> {
let argv: Vec<String> = std::env::args().skip(1).collect();
let cmd = match argv.first() {
Some(c) => c.clone(),
None => {
println!("{USAGE}");
return Ok(());
}
};
if cmd == "help" || cmd == "--help" || cmd == "-h" {
println!("{USAGE}");
return Ok(());
}
if cmd == "version" || cmd == "--version" || cmd == "-V" {
println!("areev {}", env!("CARGO_PKG_VERSION"));
return Ok(());
}
let (flags, positional) = parse_args(&argv[1..]);
for var_flag in [
"passphrase-env",
"token-env",
"anon-key-env",
"signing-key-env",
"sso-secret-env",
"sso-secret-env-next",
"oidc-client-secret-env",
] {
if let Some(var) = flag(&flags, var_flag) {
areev_core::proc::deny_env_var(&var);
areev_loop::proc::deny_env_var(&var);
}
}
if let Some(list) = run_stack::flag_or_env(&flags, "credential", areev_run::egress_spec::ENV_CREDENTIAL) {
for pair in list.split(',') {
if let Some((_, spec)) = pair.split_once('=') {
let spec = spec.trim();
if spec.starts_with("cmd:") {
continue;
}
if spec.starts_with("vault:") {
for var in ["VAULT_TOKEN", "VAULT_ADDR", "VAULT_NAMESPACE"] {
areev_core::proc::deny_env_var(var);
areev_loop::proc::deny_env_var(var);
}
continue;
}
let var = spec.split_once('@').map(|(v, _)| v).unwrap_or(spec).trim();
if !var.is_empty() {
areev_core::proc::deny_env_var(var);
areev_loop::proc::deny_env_var(var);
}
}
}
}
if let Some(list) = run_stack::flag_or_env(&flags, "resolver-env", areev_run::egress_spec::ENV_RESOLVER_ENV) {
for var in list.split(',').map(str::trim).filter(|v| !v.is_empty()) {
areev_core::proc::deny_env_var(var);
areev_loop::proc::deny_env_var(var);
}
}
let decider = resolve_decider(&flags)?;
run_stack::recall_deadline(&flags)?;
if cmd == "decide" {
return run_decide(decider, &flags, &argv[1..]);
}
if cmd == "auth" {
return run_auth(&flags, &positional);
}
if cmd == "provision" {
return run_provision(&flags);
}
if cmd == "pack" && positional.first().map(|s| s.as_str()) == Some("validate") {
return pack::run_pack(None, "shared", &flags, &positional);
}
if cmd == "anonymize" && positional.first().map(String::as_str) == Some("scan") {
let text = match flag(&flags, "text") {
Some(t) => t,
None => {
let mut buf = String::new();
std::io::Read::read_to_string(&mut std::io::stdin(), &mut buf)
.map_err(|e| format!("reading stdin: {e}"))?;
buf
}
};
let policy = match flag(&flags, "policy-file") {
Some(path) => {
let json = std::fs::read_to_string(&path)
.map_err(|e| format!("--policy-file {path}: {e}"))?;
areev_core::anon::AnonPolicy::from_json(&json).map_err(|e| e.to_string())?
}
None => areev_core::anon::AnonPolicy::default(),
};
let out = areev_core::anon::scan(&text, &policy, &[]).map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&out).map_err(|e| e.to_string())?
);
return Ok(());
}
if cmd == "anonymize" && positional.first().map(String::as_str) == Some("test") {
let path = flag(&flags, "fixtures")
.ok_or("anonymize test requires --fixtures <FILE> (JSON)")?;
let raw = std::fs::read_to_string(&path)
.map_err(|e| format!("--fixtures {path}: {e}"))?;
let fx: serde_json::Value =
serde_json::from_str(&raw).map_err(|e| format!("--fixtures {path}: {e}"))?;
let must_redact = fixture_strings(&fx, "must_redact", &path)?;
let must_not_redact = fixture_strings(&fx, "must_not_redact", &path)?;
let policy = match flag(&flags, "policy-file") {
Some(pp) => {
let json = std::fs::read_to_string(&pp)
.map_err(|e| format!("--policy-file {pp}: {e}"))?;
areev_core::anon::AnonPolicy::from_json(&json).map_err(|e| e.to_string())?
}
None => match fx.get("policy") {
Some(v) => areev_core::anon::AnonPolicy::from_json(&v.to_string())
.map_err(|e| e.to_string())?,
None => return Err(format!("{path} has no \"policy\" and no --policy-file given")),
},
};
let mut misses: Vec<String> = Vec::new();
let mut false_positives: Vec<String> = Vec::new();
for want in &must_redact {
let got = areev_core::anon::anonymize(want, &policy, &[], None)
.map_err(|e| e.to_string())?;
if got.text == *want {
misses.push(want.clone());
}
}
for want in &must_not_redact {
let got = areev_core::anon::anonymize(want, &policy, &[], None)
.map_err(|e| e.to_string())?;
if got.text != *want {
false_positives.push(format!("{want} -> {}", got.text));
}
}
let total = must_redact.len() + must_not_redact.len();
let failed = misses.len() + false_positives.len();
for m in &misses {
eprintln!("MISS (should have been redacted): {m}");
}
for f in &false_positives {
eprintln!("FALSE POSITIVE (should have been left alone): {f}");
}
println!(
"{} fixtures: {} passed, {} failed ({} missed, {} false positive)",
total,
total - failed,
failed,
misses.len(),
false_positives.len()
);
if failed > 0 {
std::process::exit(1);
}
return Ok(());
}
let db = resolve_db(&flags, matches!(cmd.as_str(), "serve" | "ui"))?;
let ns = flag(&flags, "ns").unwrap_or_else(|| "shared".to_string());
if cmd == "hook" {
let target = positional.first().map(String::as_str).unwrap_or("claude-code");
if target != "claude-code" {
return Err(format!("unknown hook target '{target}'"));
}
let exe = std::env::current_exe()
.map(|p| p.display().to_string())
.unwrap_or_else(|_| "areev".into());
println!(
r#"Add to ~/.claude/settings.json (hooks section) to close the learning
loop automatically — inject relevant memory before each prompt, and capture
each exchange (with tool outcomes) when a turn ends:
{{
"hooks": {{
"UserPromptSubmit": [{{ "hooks": [{{
"type": "command",
"command": "{exe} recall-hook --db {db} --ns {ns} --with-loop"
}}] }}],
"Stop": [{{ "hooks": [{{
"type": "command",
"command": "{exe} capture-stop --db {db} --ns {ns}"
}}] }}],
"PreCompact": [{{ "hooks": [{{
"type": "command",
"command": "{exe} capture-stop --db {db} --ns {ns}"
}}] }}]
}}
}}
PreCompact is the same verb on the event that says the conversation is about
to be dropped: it captures whatever the last Stop did not, so a compaction
costs the model its context and costs the memory nothing. No matcher, so it
fires for manual /compact and automatic compaction alike, and it exits
silently when the host gives it no transcript — a hook must never be why a
compaction stalls. The summary a compaction produces is stored with
`compact_summary: true`, so a reader can tell a machine's recap from what the
person actually typed.
recall-hook reads the prompt and prints matching memories to stdout, which
Claude Code injects as context — so retrieval no longer depends on the model
choosing to call a tool. For on-demand reads/writes by the model itself, also
register the MCP server:
claude mcp add areev -- {exe} serve --mcp --db {db} --ns {ns}
Nothing was written — apply the snippet yourself (or rerun with your own paths)."#
);
return Ok(());
}
let read_only = flags.contains_key("read-only");
let is_pg_url = is_pg_dsn(&db);
let enc_key = match flag(&flags, "passphrase-env") {
Some(_) if is_pg_url => {
return Err(
"--passphrase-env applies to file-backed memories (page cipher + .kdf \
sidecar); on the postgres backend use TDE/pgcrypto at the deployment layer"
.into(),
);
}
Some(var) => {
let pass = zeroize::Zeroizing::new(std::env::var(&var).map_err(|_| {
format!("--passphrase-env {var}: environment variable is not set")
})?);
if pass.trim().is_empty() {
return Err(format!("--passphrase-env {var}: passphrase is empty"));
}
if !is_pg_url {
areev_store::read_only_requires_existing(&db, read_only)
.map_err(|e| e.to_string())?;
}
Some(Areev::derive_key_for(&db, pass.as_str()).map_err(|e| e.to_string())?)
}
None => None,
};
let anon_key = match flag(&flags, "anon-key-env") {
Some(var) => {
let raw = zeroize::Zeroizing::new(std::env::var(&var).map_err(|_| {
format!("--anon-key-env {var}: environment variable is not set")
})?);
Some(parse_anon_key(raw.trim()).map_err(|e| format!("--anon-key-env {var}: {e}"))?)
}
None => None,
};
let tel_mode = match flag(&flags, "telemetry") {
Some(v) => areev_store::TelemetryMode::parse(&v)
.ok_or_else(|| {
format!("--telemetry: unknown mode '{v}' (off|aggregate|aggregate-hashed|full)")
})?,
None if read_only => areev_store::TelemetryMode::Off,
None => areev_store::TelemetryMode::Aggregate,
};
if cmd == "blob" && positional.first().map(String::as_str) == Some("get") {
let uri = positional
.get(1)
.ok_or("usage: areev blob get <cas-uri> --db <file|dsn>")?;
if let Some(bytes) = areev_store::read_blob_offline(&db, uri).map_err(|e| e.to_string())? {
use std::io::Write;
std::io::stdout().write_all(&bytes).map_err(|e| e.to_string())?;
return Ok(());
}
}
let explicit_index = flag(&flags, "index-text");
if read_only && explicit_index.is_some() {
return Err(
"--read-only cannot be combined with an explicit --index-text: --index-text \
always re-stamps the file's declaration, and a read-only open never writes. \
Drop --index-text (a read-only open honors whatever the file already declares) \
or drop --read-only"
.into(),
);
}
let mut m = if is_pg_url {
open_postgres_store(&db, tel_mode, explicit_index.as_deref(), anon_key, read_only)?
} else {
if explicit_index.is_some() || enc_key.is_some() || anon_key.is_some() || read_only {
let mut o = areev_store::AreevOptions::default();
if let Some(v) = &explicit_index {
o.index_text = !matches!(v.as_str(), "false" | "0" | "off" | "no");
}
if let Some(key) = &enc_key {
o.encryption_key = Some(**key);
}
o.anon_key = anon_key;
o.telemetry = tel_mode;
o.read_only = read_only;
Areev::open_with(&db, o)
} else if tel_mode != areev_store::TelemetryMode::Off {
Areev::open_with_telemetry(&db, tel_mode)
} else {
Areev::open(&db)
}
.map_err(|e| {
let mut msg = e.to_string();
if enc_key.is_none() && std::path::Path::new(&format!("{db}.kdf")).exists() {
msg.push_str(&format!(
" — {db}.kdf exists, so this memory is encrypted; pass --passphrase-env <VAR>"
));
}
msg
})?
};
for w in m.open_warnings() {
eprintln!("areev: warning: {w}");
}
m.set_run_id(flag(&flags, "run-id").as_deref());
if flags.contains_key("anonymize-egress") {
m.set_anonymize_egress_floor(true);
}
if let Some(cmd_line) = flag(&flags, "anonymize-cmd") {
let backend =
areev_store::CommandAnonymize::new(&cmd_line).map_err(|e| e.to_string())?;
m.set_anonymizer(Box::new(backend));
}
if let Some(cmd_line) = flag(&flags, "anonymize-llm-cmd") {
let model = flag(&flags, "anonymize-llm-model");
let llm = areev_loop::CommandLlm::new(&cmd_line, model.as_deref())
.map_err(|e| e.to_string())?;
m.set_anonymizer(Box::new(areev_llm::LlmDetector::new(llm)));
}
if let Some(cmd_line) = flag(&flags, "embed-cmd") {
let model = flag(&flags, "embed-model");
let seen = m.open_warnings().len();
let ce = areev_store::CommandEmbed::new(&cmd_line, model.as_deref())
.map_err(|e| e.to_string())?;
m.set_embedder(Box::new(ce));
for w in &m.open_warnings()[seen..] {
eprintln!("areev: warning: {w}");
}
}
if let Some(cmd_line) = run_stack::flag_or_env(&flags, "rerank-cmd", "AREEV_RERANK_CMD") {
let program = cmd_line.split_whitespace().next().unwrap_or_default();
if !areev_mcp::program_resolves(program) {
return Err(format!(
"--rerank-cmd (or $AREEV_RERANK_CMD): '{program}' is not an executable file or a \
command on PATH"
));
}
let model = flag(&flags, "rerank-model");
let rr = areev_store::CommandRerank::new(&cmd_line, model.as_deref())
.map_err(|e| format!("--rerank-cmd: {e}"))?;
m.set_reranker(Box::new(rr));
} else if let Some(chain) = &decider {
let chain = decider_for_egress(chain.clone(), store_egress_active(&m))?;
m.set_reranker(Box::new(areev_store::DecisionRerank::new(chain)));
}
if let Some(var) = flag(&flags, "signing-key-env") {
let raw = zeroize::Zeroizing::new(std::env::var(&var).map_err(|_| {
format!("--signing-key-env {var}: environment variable is not set")
})?);
m.set_signing_key_hex(raw.trim())
.map_err(|e| format!("--signing-key-env {var}: {e}"))?;
}
if let Some(path) = flag(&flags, "trusted-authors") {
let doc = std::fs::read_to_string(&path)
.map_err(|e| format!("--trusted-authors {path}: {e}"))?;
m.set_trusted_authors(&doc).map_err(|e| format!("--trusted-authors {path}: {e}"))?;
}
if flags.contains_key("require-attested") {
if flag(&flags, "trusted-authors").is_none() {
return Err("--require-attested needs --trusted-authors FILE (nothing to require against)".into());
}
m.set_attest_policy(areev_store::AttestPolicy::Require);
}
match cmd.as_str() {
"add" => {
let (s, r, o) = match (
flag(&flags, "subject"),
flag(&flags, "relation"),
flag(&flags, "object"),
) {
(Some(s), Some(r), Some(o)) => (s, r, o),
_ if positional.len() >= 3 => (
positional[0].clone(),
positional[1].clone(),
positional[2].clone(),
),
_ => {
return Err("usage: areev add <subject> <relation> <object> \
(or --subject S --relation R --object O)"
.to_string())
}
};
for (name, v) in [("subject", &s), ("relation", &r), ("object", &o)] {
if v.trim().is_empty() {
return Err(format!("{name} must not be empty"));
}
}
let conf: f64 = flag(&flags, "confidence")
.map(|c| c.parse().map_err(|_| "bad --confidence".to_string()))
.transpose()?
.unwrap_or(0.9);
let mut f = Fact::new(&s, &r, &o).confidence(conf);
f.common.namespace = Some(ns);
if flags.contains_key("idempotent") {
let (h, inserted) = m.add_if_novel(&f).map_err(|e| e.to_string())?;
println!("{h}");
if !inserted {
eprintln!("(unchanged — value already current, no new grain)");
}
} else {
let h = m.add(&f).map_err(|e| e.to_string())?;
println!("{h}");
}
}
"record-tool-call" => {
let name = flag(&flags, "name")
.or_else(|| positional.first().cloned())
.ok_or_else(|| "usage: areev record-tool-call --name NAME --result TEXT [--input JSON] [--is-error] [--thread ID] [--call-id ID] [--run-id ID] [--workflow HASH --node ID] [--status pending|completed|failed] [--failure-cause CAUSE] [--failure-detail TEXT] [--executor-kind host|client] [--correlation-id ID]".to_string())?;
let result = flag(&flags, "result")
.or_else(|| positional.get(1).cloned())
.ok_or_else(|| "record-tool-call requires --result TEXT".to_string())?;
let input = flag(&flags, "input");
let is_error = flag(&flags, "is-error")
.map(|v| !matches!(v.as_str(), "false" | "0" | "off" | "no"))
.unwrap_or(false);
let facade = AreevFacade::with_session(m, Some(ns.clone()), None);
let facade = apply_principal(facade, &flags)?;
let hash = facade
.record_tool_call(
&ns,
&name,
input.as_deref(),
&result,
is_error,
flag(&flags, "thread").as_deref(),
flag(&flags, "call-id").as_deref(),
flag(&flags, "run-id").as_deref(),
flag(&flags, "workflow").as_deref(),
flag(&flags, "node").as_deref(),
flag(&flags, "status").as_deref(),
flag(&flags, "failure-cause").as_deref(),
flag(&flags, "failure-detail").as_deref(),
flag(&flags, "executor-kind").as_deref(),
flag(&flags, "correlation-id").as_deref(),
)
.map_err(|e| e.to_string())?;
println!("{hash}");
}
"run-manifest" => {
let run_id = flag(&flags, "run-id")
.or_else(|| positional.first().cloned())
.ok_or_else(|| "run-manifest requires --run-id ID".to_string())?;
let config = flag(&flags, "config")
.ok_or_else(|| "run-manifest requires --config JSON".to_string())?;
let facade = AreevFacade::with_session(m, Some(ns), None);
let facade = apply_principal(facade, &flags)?;
let (config_hash, link_hash) = facade
.record_run_manifest(&run_id, &config)
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"run_id": run_id,
"config_hash": config_hash.to_hex(),
"link_hash": link_hash.to_hex(),
})
);
}
"recall" => {
let s = flag(&flags, "subject")
.or_else(|| positional.first().cloned())
.ok_or_else(|| "usage: areev recall <subject> (or --subject S)".to_string())?;
let rel = flag(&flags, "relation");
let k: usize = flag(&flags, "k").and_then(|v| v.parse().ok()).unwrap_or(16);
let grains = m
.recall(&ns, &s, rel.as_deref(), k)
.map_err(|e| e.to_string())?;
match flag(&flags, "render") {
None => {
for g in grains {
let line = serde_json::json!({
"hash": g.hash.to_hex(),
"type": format!("{:?}", g.grain_type).to_lowercase(),
"fields": g.fields,
});
println!("{line}");
}
}
Some(render) => {
use areev_context::{ContextAssembler, FormatPolicy, OutputFormat};
let mut policy = match render.as_str() {
"sml" => FormatPolicy::claude(),
"markdown" => FormatPolicy::gpt4(),
"toon" => FormatPolicy::new(OutputFormat::Toon),
"plain" => FormatPolicy::new(OutputFormat::PlainText),
"json" => FormatPolicy::json_api(),
other => return Err(format!("unknown --render '{other}'")),
};
if let Some(b) = flag(&flags, "budget").and_then(|v| v.parse().ok()) {
policy.token_budget = Some(b);
}
let hits: Vec<areev_cal::store_types::SearchHit> = grains
.into_iter()
.map(|grain| {
let hash = grain.hash;
areev_cal::store_types::SearchHit {
grain,
score: 1.0,
hash,
score_breakdown: None,
explanation: None,
scope_depth: None,
source_namespace: None,
relative_time: None,
conflict_status: None,
supersession_status: None,
superseded_by_hash: None,
recall_source: None,
}
})
.collect();
let ctx = ContextAssembler::new().format(&hits, &policy);
println!("{}", ctx.text);
eprintln!(
"-- {} grains, ~{} tokens{}",
ctx.included_count,
ctx.estimated_tokens,
if ctx.truncated { " (truncated to budget)" } else { "" }
);
}
}
}
"search" => {
let q = need(&flags, "query")?;
let subject = flag(&flags, "subject");
let k: usize = flag(&flags, "k").and_then(|v| v.parse().ok()).unwrap_or(10);
let tuning = areev_store::RecallTuning { rerank: m.has_reranker(), ..Default::default() };
let hits = m
.recall_hybrid_scored(
&ns,
subject.as_deref(),
None,
Some(&q),
k,
run_stack::recall_deadline(&flags)?,
tuning,
)
.map_err(|e| e.to_string())?;
for (g, score) in hits {
println!("{}", serde_json::json!({
"hash": g.hash.to_hex(),
"type": format!("{:?}", g.grain_type).to_lowercase(),
"score": score,
"fields": g.fields,
}));
}
}
"cal" => {
let query = positional
.first()
.ok_or_else(|| "usage: areev cal '<QUERY>' --db <file>".to_string())?
.clone();
let facade = run_stack::host_facade(m, Some(ns), &flags)?;
report_meta_warnings(&facade);
let facade = apply_principal(facade, &flags)?;
let ex = CalExecutor::new(CalExecutorConfig {
allow_destructive_ops: !flags.contains_key("no-destructive-ops"),
assembly_manifest_sample_rate: assembly_manifest_sample_rate(&flags)?,
..CalExecutorConfig::default()
})
.with_governance(std::sync::Arc::new(areev_loop_adapter::LoopGovernance::new()));
let res = ex.execute(&query, &facade).map_err(|e| e.to_string())?;
let payload = serde_json::to_string_pretty(&res.result).map_err(|e| e.to_string())?;
println!("{payload}");
for w in res.warnings {
eprintln!("warning: {w}");
}
}
"corpus" => {
let selector = flag(&flags, "select")
.or_else(|| positional.first().cloned())
.ok_or_else(|| {
"usage: areev corpus --select '<READ CAL>' [--out FILE] [--recipient ID]"
.to_string()
})?;
let facade = run_stack::host_facade(m, Some(ns), &flags)?;
report_meta_warnings(&facade);
let facade = apply_principal(facade, &flags)?;
let destination = flag(&flags, "out").unwrap_or_else(|| "stdout".to_string());
let recipient = flag(&flags, "recipient");
let (summary, manifest) =
corpus::export(&facade, &selector, &destination, recipient.as_deref(), now_ms())?;
eprintln!(
"corpus export: {} row(s), {} source grain(s), manifest {}",
summary.rows,
summary.source_hashes.len(),
manifest
);
}
"tune" => {
let facade = run_stack::host_facade(m, Some(ns), &flags)?;
report_meta_warnings(&facade);
let facade = apply_principal(facade, &flags)?;
return tune::run_tune(&facade, &flags);
}
"history" => {
let s = need(&flags, "subject")?;
let r = need(&flags, "relation")?;
let versions = m.history(&ns, &s, &r).map_err(|e| e.to_string())?;
for v in versions {
println!(
"{}",
serde_json::json!({
"hash": v.hash.to_hex(),
"object": v.object,
"created_at": v.created_at,
"confidence": v.confidence,
"superseded_by": v.superseded_by.map(|h| h.to_hex()),
})
);
}
}
"log" => {
let since: i64 = flag(&flags, "since").and_then(|v| v.parse().ok()).unwrap_or(0);
let limit: usize = flag(&flags, "limit").and_then(|v| v.parse().ok()).unwrap_or(50);
let scope: Vec<String> = flag(&flags, "ns")
.map(|v| {
v.split(',')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect()
})
.unwrap_or_default();
let ops = m
.changes_since_scoped(since, &scope, limit)
.map_err(|e| e.to_string())?;
for op in ops {
let kind = match op.op {
areev_store::OP_ADD => "add",
areev_store::OP_SUPERSEDE => "supersede",
areev_store::OP_FORGET => "forget",
_ => "?",
};
println!(
"{:>6} {:>20} {:<9} {} {}",
op.op_seq,
op.hlc,
kind,
op.hash.to_hex(),
op.ns.as_deref().unwrap_or("-")
);
}
}
"telemetry" => {
let action = positional.first().map(String::as_str).unwrap_or("");
if action != "scrub" {
return Err(
"usage: areev telemetry scrub --ns NS --yes".into()
);
}
let ns = need(&flags, "ns")?;
if !flags.contains_key("yes") {
return Err(format!(
"refusing to scrub telemetry for {ns:?} without --yes — this removes \
recall evidence permanently"
));
}
m.telemetry_scrub_namespace(&ns).map_err(|e| e.to_string())?;
println!("scrubbed telemetry for namespace {ns:?}");
}
"stream" => {
let to = need(&flags, "to")?;
let interval: u64 = flag(&flags, "interval-ms").and_then(|v| v.parse().ok()).unwrap_or(500);
let once = flags.contains_key("once");
let checkpoint = flags.contains_key("checkpoint");
let retain_ms = match flag(&flags, "retain") {
Some(v) => Some(
parse_duration(&v)
.ok_or_else(|| format!("--retain takes a duration like 30d, got '{v}'"))?,
),
None => None,
};
std::fs::create_dir_all(&to).map_err(|e| e.to_string())?;
let cursor_path = format!("{to}/CURSOR");
let (mut gen_id, mut cursor, mut seg) = read_stream_cursor(&cursor_path, &to);
if gen_id.is_empty() || checkpoint {
let g = format!(
"{:016x}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos() as u64
);
std::fs::create_dir_all(format!("{to}/gen-{g}")).map_err(|e| e.to_string())?;
let snap = format!("{to}/gen-{g}/segment-{:08}.mgb", 0);
let st = m.bundle_since(0, &snap).map_err(|e| e.to_string())?;
gen_id = g;
cursor = st.last_op_seq;
seg = 1;
std::fs::write(&cursor_path, format!("{gen_id} {cursor} {seg}"))
.map_err(|e| e.to_string())?;
eprintln!("checkpoint: full snapshot of {} ops → {snap}", st.ops);
if let Some(window) = retain_ms {
let dropped = prune_generations(&to, &gen_id, window)?;
if dropped > 0 {
eprintln!(
"retention: dropped {dropped} generation(s) older than the window \
(their contents are in this checkpoint, minus anything erased since)"
);
}
}
}
loop {
let ops_now = m.changes_since(cursor, 1).map_err(|e| e.to_string())?;
if !ops_now.is_empty() {
let seg_path = format!("{to}/gen-{gen_id}/segment-{seg:08}.mgb");
let st = m.bundle_since(cursor, &seg_path).map_err(|e| e.to_string())?;
cursor = st.last_op_seq;
seg += 1;
std::fs::write(&cursor_path, format!("{gen_id} {cursor} {seg}"))
.map_err(|e| e.to_string())?;
eprintln!("shipped {} ops → {seg_path}", st.ops);
}
if once {
break;
}
std::thread::sleep(std::time::Duration::from_millis(interval));
}
}
"follow" => {
let from = need(&flags, "from")?;
let interval: u64 = flag(&flags, "interval-ms").and_then(|v| v.parse().ok()).unwrap_or(1000);
let once = flags.contains_key("once");
#[cfg(feature = "postgres")]
let fcur_path = if is_pg_url {
let (_, schema) =
areev_store::pg::split_schema_url(&db).map_err(|e| e.to_string())?;
format!("{from}/{schema}.follow")
} else {
format!("{db}.follow")
};
#[cfg(not(feature = "postgres"))]
let fcur_path = format!("{db}.follow");
loop {
let cursor = std::fs::read_to_string(format!("{from}/CURSOR")).unwrap_or_default();
let gen_id = cursor.trim().split(' ').next().unwrap_or("").to_string();
if !gen_id.is_empty() {
let (fgen, fseg) = match std::fs::read_to_string(&fcur_path) {
Ok(s) => {
let mut it = s.trim().splitn(2, ' ');
(it.next().unwrap_or("").to_string(),
it.next().and_then(|v| v.parse::<u32>().ok()).unwrap_or(0))
}
Err(_) => (String::new(), 0),
};
if !fgen.is_empty()
&& fgen != gen_id
&& !std::path::Path::new(&format!("{from}/gen-{fgen}")).exists()
{
eprintln!(
"follow: generation {fgen} has aged out of the stream dir — \
re-baselining from generation {gen_id}'s snapshot"
);
}
let mut fseg = if fgen != gen_id { 0 } else { fseg };
loop {
let seg_path = format!("{from}/gen-{gen_id}/segment-{fseg:08}.mgb");
if !std::path::Path::new(&seg_path).exists() {
break;
}
let st = m.import_bundle(&seg_path).map_err(|e| e.to_string())?;
eprintln!("applied segment {fseg} ({} ops, {} skipped)", st.applied, st.skipped);
fseg += 1;
std::fs::write(&fcur_path, format!("{gen_id} {fseg}")).map_err(|e| e.to_string())?;
}
}
if once {
break;
}
std::thread::sleep(std::time::Duration::from_millis(interval));
}
}
"restore" => {
let from = need(&flags, "from")?;
let until: Option<i64> = flag(&flags, "until-hlc").and_then(|v| v.parse().ok());
let cursor = std::fs::read_to_string(format!("{from}/CURSOR"))
.map_err(|_| "no CURSOR in stream dir".to_string())?;
let gen_id = cursor.trim().split(' ').next().unwrap_or("").to_string();
let dir = format!("{from}/gen-{gen_id}");
let mut segs: Vec<String> = std::fs::read_dir(&dir)
.map_err(|e| e.to_string())?
.filter_map(|e| e.ok().map(|e| e.path().display().to_string()))
.filter(|p| p.ends_with(".mgb"))
.collect();
segs.sort();
let mut applied = 0usize;
let mut meta_applied = 0usize;
let mut meta_skipped = 0usize;
for s in &segs {
let st = m.import_bundle_until(s, until).map_err(|e| e.to_string())?;
applied += st.applied;
meta_applied += st.meta_applied;
meta_skipped += st.meta_skipped;
}
let registry = if meta_applied > 0 {
format!(", {meta_applied} registry entries")
} else {
String::new()
};
println!(
"restored {} ops from {} segments (gen {gen_id}){registry}",
applied,
segs.len()
);
if until.is_some() && meta_skipped > 0 {
eprintln!(
"note: {meta_skipped} registry entries (saved queries/templates/retention) \
skipped — a point-in-time restore replays grain history only"
);
}
}
"bundle" => {
let out = need(&flags, "out")?;
let since: i64 = flag(&flags, "since").and_then(|v| v.parse().ok()).unwrap_or(0);
let st = m.bundle_since(since, &out).map_err(|e| e.to_string())?;
println!(
"bundled {} ops ({} bytes) → {} (next --since {})",
st.ops, st.bytes, out, st.last_op_seq
);
}
"import" => {
let bundle = need(&flags, "bundle")?;
let st = m.import_bundle(&bundle).map_err(|e| e.to_string())?;
let registry = if st.meta_applied > 0 {
format!(", {} registry entries", st.meta_applied)
} else {
String::new()
};
println!("applied {} ops, skipped {}{registry}", st.applied, st.skipped);
}
"migrate" => {
let from = need(&flags, "from")?;
let file = need(&flags, "file")?;
let deferred = m.defer_text_index().map_err(|e| e.to_string())?;
let rep = run_migrate(&mut m, &ns, &from, &file, flag(&flags, "history"));
if deferred {
m.rebuild_text_index().map_err(|e| e.to_string())?;
}
let rep = rep?;
for n in &rep.notes {
eprintln!("areev: migrate: {n}");
}
println!("{}", rep.to_json());
}
"reindex" => {
let n = m.rebuild_text_index().map_err(|e| e.to_string())?;
println!(
"text index rebuilt: {} grains indexed ({n} needed their text backfilled)",
m.indexed_documents()
);
let links = m.rebuild_link_indexes().map_err(|e| e.to_string())?;
println!("link indexes rebuilt: {links} rows (provenance, runs, cross-links)");
}
"vector-index" => {
let usage = "usage: areev vector-index <status|build|drop|check> --db DSN \
[build: --m N --ef-construction N --ef-search N] \
[check: --queries FILE.json --ns SCOPE --k N --ef-search N]";
let num = |key: &str, default: usize| -> Result<usize, String> {
flag(&flags, key)
.map_or(Ok(default), |v| v.parse::<usize>())
.map_err(|_| format!("--{key} must be a number"))
};
match positional.first().map(String::as_str).unwrap_or("status") {
"status" => {
let name = m.vector_index().map_err(|e| e.to_string())?;
println!("{}", serde_json::json!({"index": name}));
}
"build" => {
let (hm, efc, efs) =
(num("m", 16)?, num("ef-construction", 64)?, num("ef-search", 40)?);
m.ensure_vector_index(hm, efc, efs).map_err(|e| e.to_string())?;
let name = m.vector_index().map_err(|e| e.to_string())?;
println!("{}", serde_json::json!({"index": name, "m": hm, "ef_construction": efc, "ef_search": efs}));
}
"drop" => {
m.drop_vector_index().map_err(|e| e.to_string())?;
println!("{}", serde_json::json!({"index": serde_json::Value::Null}));
}
"check" => {
let path = flag(&flags, "queries")
.ok_or("vector-index check requires --queries FILE.json (a JSON array of vectors)")?;
if !flags.contains_key("ns") {
return Err("vector-index check requires --ns SCOPE — the WIDE scope you actually \
query (e.g. --ns 'firm.deals.*'); a narrow scope is exact regardless"
.into());
}
let text = std::fs::read_to_string(&path)
.map_err(|e| format!("cannot read --queries {path}: {e}"))?;
let queries: Vec<Vec<f32>> = serde_json::from_str(&text)
.map_err(|e| format!("--queries must be a JSON array of number arrays: {e}"))?;
let k = num("k", 10)?;
if let Some(ef) = flag(&flags, "ef-search") {
let ef: usize = ef.parse().map_err(|_| "--ef-search must be a number")?;
m.set_vector_ef_search(ef).map_err(|e| e.to_string())?;
}
let report = m.vector_recall_check(&ns, &queries, k).map_err(|e| e.to_string())?;
println!("{}", serde_json::to_string_pretty(&report).unwrap());
}
_ => return Err(usage.into()),
}
}
"related" => {
let start = flag(&flags, "start").ok_or("related requires --start")?;
let rels = areev_store::parse_relations(
&flag(&flags, "relations").ok_or("related requires --relations")?,
);
if rels.is_empty() {
return Err("--relations must name at least one relation".to_string());
}
let dir = Direction::parse(&flag(&flags, "direction").unwrap_or_default())
.ok_or("--direction must be one of: out, in, both")?;
let depth: usize = flag(&flags, "depth")
.map_or(Ok(2), |d| d.parse())
.map_err(|_| "--depth must be a number")?;
let cap: usize = flag(&flags, "limit")
.map_or(Ok(64), |l| l.parse())
.map_err(|_| "--limit must be a number")?;
let refs: Vec<&str> = rels.iter().map(String::as_str).collect();
let scope: Vec<String> = ns
.split(',')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect();
let reached = if scope.len() > 1 {
m.related_scoped(&scope, &start, &refs, dir, depth, cap)
} else {
m.related(&ns, &start, &refs, dir, depth, cap)
}
.map_err(|e| e.to_string())?;
if reached.is_empty() {
println!("(nothing reachable from {start})");
}
for r in reached {
println!("{r}");
}
}
"entity-at" => {
let subject = flag(&flags, "subject").ok_or("entity-at requires --subject")?;
let relation = flag(&flags, "relation").ok_or("entity-at requires --relation")?;
let at: i64 = flag(&flags, "at")
.ok_or("entity-at requires --at (epoch ms)")?
.parse()
.map_err(|_| "--at must be epoch milliseconds")?;
let axis = Axis::parse(&flag(&flags, "axis").unwrap_or_default())
.ok_or("--axis must be one of: world, knowledge")?;
let scope: Vec<String> = ns
.split(',')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect();
if scope.len() > 1 {
let rows = m
.entity_at_scoped(&scope, &subject, &relation, at, axis)
.map_err(|e| e.to_string())?;
if rows.is_empty() {
println!("(nothing known for {subject} {relation} at {at})");
} else {
let out: Vec<_> = rows
.iter()
.map(|(n, g)| serde_json::json!({"namespace": n, "grain": g}))
.collect();
println!("{}", serde_json::to_string_pretty(&out).unwrap_or_default());
}
return Ok(());
}
match m
.entity_at(&ns, &subject, &relation, at, axis)
.map_err(|e| e.to_string())?
{
Some(g) => println!("{}", serde_json::to_string_pretty(&g).unwrap_or_default()),
None => println!("(nothing known for {subject} {relation} at {at})"),
}
}
"step-actions" => {
let wf = flag(&flags, "workflow").ok_or("step-actions requires --workflow")?;
let wf = Hash::from_hex(&wf).map_err(|e| e.to_string())?;
let node = flag(&flags, "node");
let limit: usize = flag(&flags, "limit")
.map_or(Ok(64), |l| l.parse())
.map_err(|_| "--limit must be a number")?;
let rows = m
.step_actions(&ns, &wf, node.as_deref(), limit)
.map_err(|e| e.to_string())?;
if rows.is_empty() {
println!("(no execution records for {})", wf.to_hex());
}
for (n, h) in rows {
println!("{n}\t{}", h.to_hex());
}
}
"run-trace" => {
let run = flag(&flags, "run-id").ok_or("run-trace requires --run-id")?;
let limit: usize = flag(&flags, "limit")
.map_or(Ok(64), |l| l.parse())
.map_err(|_| "--limit must be a number")?;
let trace = m.run_trace(&ns, &run, limit).map_err(|e| e.to_string())?;
let produced = m.run_yield(&ns, &run, limit).map_err(|e| e.to_string())?;
if trace.is_empty() {
println!("(no grains recorded for run {run})");
}
println!("recorded during {run}: {} grain(s)", trace.len());
for g in &trace {
println!(" {} {:?}", g.hash.to_hex(), g.grain_type);
}
println!("produced by {run}: {} grain(s)", produced.len());
for g in &produced {
println!(" {} {:?}", g.hash.to_hex(), g.grain_type);
}
}
"runs-touching" => {
let h = flag(&flags, "hash").ok_or("runs-touching requires --hash")?;
let h = Hash::from_hex(&h).map_err(|e| e.to_string())?;
let depth: usize = flag(&flags, "depth")
.map_or(Ok(4), |d| d.parse())
.map_err(|_| "--depth must be a number")?;
let runs = m.runs_touching(&ns, &h, depth).map_err(|e| e.to_string())?;
if runs.is_empty() {
println!("(no run produced or refined {})", h.to_hex());
}
for r in runs {
println!("{r}");
}
}
"verify" => {
let rep = m.verify().map_err(|e| e.to_string())?;
println!(
"integrity: {} | grains: {} | hash mismatches: {} | undecodable: {}",
rep.integrity, rep.grains, rep.hash_mismatches, rep.undecodable
);
let mut failed = rep.integrity != "ok" || rep.hash_mismatches > 0 || rep.undecodable > 0;
if flags.contains_key("attestations") {
let a = m.verify_attestations().map_err(|e| e.to_string())?;
println!(
"attestations: {} | policy: {} | attested: {} | unattested: {} | invalid: {} | unknown key: {} | orphaned: {}",
a.attestations, a.policy, a.attested, a.unattested, a.attest_invalid,
a.attest_unknown_key, a.attest_orphaned
);
for line in &a.invalid {
eprintln!("areev: invalid attestation: {line}");
}
failed |= a.attest_invalid > 0
|| (a.policy == "require" && (a.unattested > 0 || a.attest_unknown_key > 0));
}
if failed {
return Err("verification FAILED".to_string());
}
}
"attest" => {
if m.signing_key().is_none() {
return Err("areev attest needs --signing-key-env VAR (the author key)".into());
}
if flags.contains_key("all") {
let prefix = flag(&flags, "ns");
let st = m.attest_all(prefix.as_deref()).map_err(|e| e.to_string())?;
println!("attested: {} | already attested: {}", st.attested, st.skipped);
} else {
let Some(h) = positional.first().cloned().or_else(|| flag(&flags, "hash")) else {
return Err("usage: areev attest <hash> | areev attest --all [--ns PREFIX]".into());
};
let hash = Hash::from_hex(&h).map_err(|e| e.to_string())?;
let att = m.attest(&hash).map_err(|e| e.to_string())?;
println!("{}", att.to_hex());
}
}
"stats" => {
let s = m.stats().map_err(|e| e.to_string())?;
println!(
"grains: {} ({} current) | triples: {} | terms: {} | ops: {} | thread-indexed events: {}",
s.grains, s.current, s.triples, s.terms, s.ops, s.events_indexed
);
}
"serve" => {
if !flags.contains_key("mcp") {
return Err("only --mcp transport is available (areev serve --mcp)".to_string());
}
let mut facade = run_stack::host_facade(m, Some(ns), &flags)?;
report_meta_warnings(&facade);
for (alias, target) in parse_mounts(flag(&flags, "mount").as_deref())? {
if mount_is_the_primary(&db, &target) {
return Err(format!(
"mount '{alias}' names the same postgres memory as --db — a mount is a \
SECOND store, and mounting the primary would make \"{alias}.<ns>\" and \
\"<ns>\" the same rows read twice. Point it at another schema, or drop \
the mount and query the namespace directly"
));
}
let store = open_mount(&target)
.map_err(|e| format!("mount '{alias}' ({}): {e}", redact_dsn(&target)))?;
eprintln!("areev: mounted '{alias}' (read-only) → {}", redact_dsn(&target));
facade.mount(&alias, store);
}
let facade = apply_principal(facade, &flags)?;
let mut server = areev_mcp::McpServer::host_configured(facade, None)
.with_memory_path(&db)
.assembly_manifest_sample_rate(assembly_manifest_sample_rate(&flags)?);
if flags.contains_key("no-destructive-ops") {
server = server.allow_destructive_ops(false);
eprintln!("areev: destructive operations disabled (read-only session)");
}
if let Some(lock) = flag(&flags, "lock-ns") {
server = server.lock_namespace(lock.clone());
eprintln!("areev: namespace locked to '{lock}' (per-call namespace ignored)");
}
if let Some(p) = flag(&flags, "profile") {
let profile = areev_mcp::ToolProfile::parse(&p)?;
server = server.tool_profile(profile);
if profile == areev_mcp::ToolProfile::Memory {
eprintln!("areev: MCP tool profile 'memory' (run/loop tools unavailable)");
}
}
if let Some(p) = load_policy(&flags)? {
server = server.with_loop_policy(p);
eprintln!("areev: loop host policy attached to areev_loop");
}
server.serve_stdio().map_err(|e| e.to_string())?;
}
"capture-stop" => {
use std::io::Read as IoRead;
let mut input = String::new();
std::io::stdin().read_to_string(&mut input).map_err(|e| e.to_string())?;
let hook: serde_json::Value =
serde_json::from_str(&input).map_err(|e| format!("bad hook json: {e}"))?;
let session = hook["session_id"].as_str().unwrap_or("unknown-session").to_string();
let tpath = match hook["transcript_path"].as_str().map(str::trim) {
Some(p) if !p.is_empty() => p,
_ => return Ok(()),
};
let Ok(transcript) = std::fs::read_to_string(tpath) else { return Ok(()) };
let compact_trigger = hook["hook_event_name"]
.as_str()
.filter(|e| e.eq_ignore_ascii_case("precompact"))
.map(|_| hook["trigger"].as_str().unwrap_or("unknown").to_string());
let policy_version = flag(&flags, "policy-version");
let mut parent: Option<String> = None;
let mut stored = 0;
let mut skipped = 0;
for (turn, line) in transcript.lines().enumerate() {
let Ok(v) = serde_json::from_str::<serde_json::Value>(line) else { continue };
let message = &v["message"];
let Some(role) = message["role"].as_str() else { continue };
let Some(role) = areev_core::types::Role::from_str(role) else {
continue;
};
let (text, blocks, elisions) =
transcript_content_blocks(&message["content"], turn);
if text.trim().is_empty() && blocks.is_empty() {
continue;
}
let mut event = Event::new(&text)
.namespace(&ns)
.session(session.clone())
.created_at(transcript_turn_ms(&v).unwrap_or(turn as i64))
.role(role)
.content_blocks(blocks);
if let Some(p) = parent.clone() {
event = event.parent_message(p);
}
if let Some(model) = message["model"].as_str().or_else(|| v["model"].as_str()) {
event = event.model(model.to_string());
}
if let Some(stop) = message["stop_reason"]
.as_str()
.or_else(|| v["stop_reason"].as_str())
{
event = event.stop_reason(stop.to_string());
}
if let Some(usage) = transcript_token_usage(message) {
event = event.token_usage(usage);
}
if !elisions.is_empty() {
event
.common
.extra_fields
.insert("elisions".into(), serde_json::Value::Array(elisions));
}
if v["isCompactSummary"].as_bool() == Some(true) {
event
.common
.extra_fields
.insert("compact_summary".into(), serde_json::Value::Bool(true));
}
if let Some(policy) = policy_version.as_deref() {
event.common.context = Some(serde_json::json!({
"policy_version": policy,
}));
}
let (_, hash) = areev_core::format::serialize::serialize_grain(&event)
.map_err(|e| e.to_string())?;
if m.get(&hash).is_ok() {
skipped += 1;
} else {
m.add(&event).map_err(|e| e.to_string())?;
stored += 1;
}
parent = Some(hash.to_hex());
}
if let Some(trigger) = &compact_trigger {
let mut obs = areev_core::types::Observation::new("agent:hook", "system")
.subject(&format!("session:{session}"))
.object(&format!(
"compaction (trigger={trigger}) — {stored} turns captured, \
{skipped} already stored"
))
.namespace(areev_core::authz::HARNESS_NS)
.created_at((stored + skipped) as i64);
let ex = &mut obs.common.extra_fields;
ex.insert("observation_kind".into(), serde_json::json!("compaction"));
ex.insert("session_id".into(), serde_json::json!(session));
ex.insert("compact_trigger".into(), serde_json::json!(trigger));
ex.insert("captured".into(), serde_json::json!(stored));
ex.insert("already_stored".into(), serde_json::json!(skipped));
let _ = m.add(&obs);
}
println!("captured {stored} events for session {session} ({skipped} already stored)");
}
"recall-hook" => {
use std::io::Read as IoRead;
let mut input = String::new();
std::io::stdin().read_to_string(&mut input).map_err(|e| e.to_string())?;
let hook: serde_json::Value =
serde_json::from_str(&input).unwrap_or(serde_json::Value::Null);
let query = hook["prompt"].as_str().unwrap_or("").trim().to_string();
if query.is_empty() {
return Ok(());
}
let k: usize = flag(&flags, "k").and_then(|v| v.parse().ok()).unwrap_or(5);
let with_loop = flag(&flags, "with-loop").is_some();
let deadline = match run_stack::recall_deadline(&flags)? {
Some(d) => Some(d),
None if m.has_reranker() => Some(std::time::Duration::from_millis(1500)),
None => None,
};
let tuning = areev_store::RecallTuning { rerank: m.has_reranker(), ..Default::default() };
let grains = m
.recall_hybrid_scored(&ns, None, None, Some(&query), k, deadline, tuning)
.map_err(|e| e.to_string())?;
if grains.is_empty() && !with_loop {
return Ok(());
}
if !grains.is_empty() {
use areev_context::{ContextAssembler, DecidePolicy, FormatPolicy};
let mut policy = FormatPolicy::claude().query_text(query.clone());
policy.token_budget =
Some(flag(&flags, "budget").and_then(|v| v.parse().ok()).unwrap_or(400));
let mut assembler = ContextAssembler::new();
if let Some(chain) = &decider {
let chain = decider_for_egress(chain.clone(), store_egress_active(&m))?;
let deadline_ms = deadline
.or(Some(std::time::Duration::from_millis(1500)))
.map(|d| d.as_millis() as u64);
policy = policy.decide(DecidePolicy { deadline_ms, ..Default::default() });
assembler = assembler.with_decider(chain);
}
let hits: Vec<areev_cal::store_types::SearchHit> = grains
.into_iter()
.map(|(grain, score)| {
let hash = grain.hash;
areev_cal::store_types::SearchHit {
grain,
score: f64::from(score),
hash,
score_breakdown: None,
explanation: None,
scope_depth: None,
source_namespace: None,
relative_time: None,
conflict_status: None,
supersession_status: None,
superseded_by_hash: None,
recall_source: None,
}
})
.collect();
let ctx = assembler.format(&hits, &policy);
if !ctx.text.trim().is_empty() {
println!("Relevant memory from Areev:\n{}", ctx.text);
}
}
if with_loop {
let sub = AreevSubstrate::new(m, Some(ns.to_string()));
let engine = Engine::with_builtins();
let mut pending = engine
.recommendations(&sub, Some(RecStatus::Pending))
.map_err(|e| e.to_string())?;
if !pending.is_empty() {
pending.sort_by(|a, b| b.severity.cmp(&a.severity).then(a.hash.cmp(&b.hash)));
const CAP: usize = 3;
println!(
"\nAreev Loop: {} pending recommendation(s) for this memory (review with `areev loop list`, act with approve/apply/reject --because):",
pending.len()
);
for r in pending.iter().take(CAP) {
let origin = match r.origin {
areev_loop::Origin::Llm { .. } => " [llm]",
areev_loop::Origin::Command { .. } => " [external]",
_ => "",
};
println!(" [{}] {}{} {}", r.severity.as_str(), short(&r.hash), origin, r.summary.render());
}
if pending.len() > CAP {
println!(" … and {} more", pending.len() - CAP);
}
}
}
}
"remember" => run_remember(m, &ns, &flags)?,
"memtool" => {
let cmd = positional
.first()
.ok_or_else(|| "usage: areev memtool '<json>' --db <file>".to_string())?;
let cmd: serde_json::Value =
serde_json::from_str(cmd).map_err(|e| format!("bad command json: {e}"))?;
let mut t = areev_store::memory_tool::MemoryTool::new(&mut m, &ns);
println!("{}", t.execute(&cmd).map_err(|e| e.to_string())?);
}
"repl" => {
use std::io::{BufRead, Write as IoWrite};
let facade = run_stack::host_facade(m, Some(ns.clone()), &flags)?;
report_meta_warnings(&facade);
let facade = apply_principal(facade, &flags)?;
let ex = CalExecutor::new(CalExecutorConfig {
allow_destructive_ops: !flags.contains_key("no-destructive-ops"),
assembly_manifest_sample_rate: assembly_manifest_sample_rate(&flags)?,
..CalExecutorConfig::default()
})
.with_governance(std::sync::Arc::new(areev_loop_adapter::LoopGovernance::new()));
eprintln!("areev repl — namespace '{ns}' · CAL statements, or .stats .log .help .quit");
let stdin = std::io::stdin();
let mut lines = stdin.lock().lines();
loop {
eprint!("cal> ");
std::io::stderr().flush().ok();
let line = match lines.next() {
Some(Ok(l)) => l,
_ => break,
};
let line = line.trim();
match line {
"" => continue,
".quit" | ".exit" => break,
".help" => eprintln!(
"CAL: RECALL / ASSEMBLE / EXISTS / HISTORY / ADD / SUPERSEDE / DESCRIBE / | COUNT\n\
dot: .stats .log .verify .quit"
),
".stats" => match facade.with_store(|st| st.stats()) {
Ok(s) => eprintln!(
"grains {} ({} current) · triples {} · ops {}",
s.grains, s.current, s.triples, s.ops
),
Err(e) => eprintln!("error: {e}"),
},
".log" => match facade.with_store(|st| st.changes_since(0, 20)) {
Ok(ops) => {
for o in ops {
eprintln!("{:>5} {:<9} {}", o.op_seq, o.op, o.hash.to_hex());
}
}
Err(e) => eprintln!("error: {e}"),
},
".verify" => match facade.with_store(|st| st.verify()) {
Ok(r) => eprintln!(
"integrity {} · {} grains · {} mismatches",
r.integrity, r.grains, r.hash_mismatches
),
Err(e) => eprintln!("error: {e}"),
},
q => match ex.execute(q, &facade) {
Ok(res) => {
println!(
"{}",
serde_json::to_string_pretty(&res.result).unwrap_or_default()
);
for w in res.warnings {
eprintln!("warning: {w}");
}
}
Err(e) => eprintln!("error: {e}"),
},
}
}
}
"ui" => {
let addr = flag(&flags, "addr").unwrap_or_else(|| "127.0.0.1:7437".to_string());
let allow_remote = flag(&flags, "allow-remote").is_some();
let auth_token = match flag(&flags, "token-env") {
Some(var) => {
let t = std::env::var(&var).map_err(|_| {
format!("--token-env {var}: environment variable is not set")
})?;
if t.trim().is_empty() {
return Err(format!("--token-env {var}: token is empty"));
}
if !areev_core::authz::token_is_minted(&t) {
eprintln!(
"areev: ⚠ --token-env {var} holds a token Areev did not mint, so \
its entropy is unknown. Prefer `areev auth mint` (256-bit) and a \
per-principal credential map (--auth), which is also the only \
way approvals get an attributable identity."
);
}
Some(t)
}
None => None,
};
let has_auth = auth_token.is_some();
if !addr_is_loopback(&addr) && !allow_remote {
let why = if has_auth {
"It is authenticated, but still plaintext HTTP — the token and all \
memory cross the wire in the clear. Terminate TLS in front of it, \
or pass --allow-remote to accept that risk."
} else {
"With no --token-env it is an UNAUTHENTICATED, writable console over \
plaintext HTTP — anyone who can reach it could read or modify your \
memory. Add --token-env <VAR>, keep it on loopback, put it behind a \
TLS-terminating proxy, or pass --allow-remote to override."
};
return Err(format!(
"refusing to bind the console to a non-loopback address ({addr}). {why}"
));
}
let facade = run_stack::host_facade(m, Some(ns), &flags)?;
report_meta_warnings(&facade);
let mut server = areev_server::UiServer::new(facade, db.clone());
if let Some(d) = &decider {
server = server.with_decider(d.describe(), d.calibrated());
}
if let Some(path) = flag(&flags, "auth") {
let text = std::fs::read_to_string(&path)
.map_err(|e| format!("--auth {path}: {e}"))?;
let map = areev_core::authz::CredentialMap::from_json(&text)
.map_err(|e| e.to_string())?;
eprintln!(
"areev: credential map loaded ({} token(s)); rights come from the \
file's grant grains; unauthenticated requests run as 'anonymous'",
map.tokens.len()
);
let now = areev_core::time::now_ms();
for (id, when) in map.expiring_within(now, 14 * 86_400_000) {
let entry = map.tokens.iter().find(|t| t.id() == id);
if entry.is_some_and(|t| t.is_expired_at(now)) {
eprintln!(
"areev: ⚠ credential {id} EXPIRED {when} — it authenticates nobody; \
mint a replacement (areev auth mint) and revoke it"
);
} else {
eprintln!("areev: credential {id} expires {when} (within 14 days)");
}
}
server = server.with_credentials(map);
}
if allow_remote {
server = server.allow_remote(true);
}
if let Some(raw) = flag(&flags, "allow-origin") {
fn validate_origin(o: &str) -> Result<String, String> {
let o = o.trim();
if o.is_empty() {
return Err("--allow-origin: empty origin".to_string());
}
if o == "*" {
return Err(
"--allow-origin: '*' is not accepted — name the exact origin(s) \
allowed to POST"
.to_string(),
);
}
let (scheme, rest) = if let Some(r) = o.strip_prefix("https://") {
("https", r)
} else if let Some(r) = o.strip_prefix("http://") {
("http", r)
} else {
return Err(format!("--allow-origin {o}: must start with http:// or https://"));
};
if rest.contains('@') {
return Err(format!("--allow-origin {o}: must not contain userinfo"));
}
if rest.contains('?') || rest.contains('#') {
return Err(format!("--allow-origin {o}: must not contain a query or fragment"));
}
let host = rest.strip_suffix('/').unwrap_or(rest);
if host.is_empty() {
return Err(format!("--allow-origin {o}: missing host"));
}
if host.contains('/') {
return Err(format!(
"--allow-origin {o}: must not contain a path beyond an optional \
trailing '/'"
));
}
Ok(format!("{scheme}://{host}"))
}
let mut origins = Vec::new();
for part in raw.split(',') {
let part = part.trim();
if part.is_empty() {
continue;
}
let normalized = validate_origin(part)?;
if normalized.starts_with("http://")
&& !addr_is_loopback(normalized.trim_start_matches("http://"))
{
eprintln!(
"areev: WARNING — --allow-origin {normalized} is http:// \
(non-TLS) and non-loopback; a Basic credential the browser \
caches for that origin would cross the wire in the clear."
);
}
server = server.allow_origin(normalized.clone());
origins.push(normalized);
}
if origins.is_empty() {
return Err("--allow-origin: no origins given".to_string());
}
eprintln!("areev: cross-origin POSTs accepted from: {}", origins.join(", "));
}
if flags.contains_key("no-destructive-ops") {
server = server.allow_destructive_ops(false);
eprintln!("areev: destructive operations disabled (read-only console)");
}
if let Some(tok) = auth_token {
server = server.with_auth(tok);
eprintln!(
"areev: console authentication ENABLED — HTTP Basic \
(any username, password = the token)"
);
} else {
eprintln!(
"areev: console is UNAUTHENTICATED — enable with --auth <map> \
(per-principal, attributable; `areev auth mint` creates one) or \
--token-env <VAR> (one shared secret, cannot approve)"
);
}
if let Some(p) = load_policy(&flags)? {
server = server.with_loop_policy(p);
eprintln!("areev: loop host policy attached to the console's loop routes");
}
let sso_approvals = match flag(&flags, "sso-approvals").as_deref() {
None | Some("deny") => false,
Some("allow") => true,
Some(other) => {
return Err(format!(
"--sso-approvals {other}: expected 'deny' (default) or 'allow'"
))
}
};
if flag(&flags, "sso-approvals").is_some() && flag(&flags, "sso-header").is_none() {
return Err(
"--sso-approvals applies to trusted-header SSO only — it has no effect \
without --sso-header/--sso-secret-env. Credential-map principals (--auth) \
may always approve; shared-token and anonymous callers never may."
.into(),
);
}
match (flag(&flags, "sso-header"), flag(&flags, "sso-secret-env")) {
(Some(header), Some(var)) => {
let secret = std::env::var(&var)
.map_err(|_| format!("--sso-secret-env {var}: environment variable is not set"))?;
if secret.trim().is_empty() {
return Err(format!("--sso-secret-env {var}: secret is empty"));
}
let next = match flag(&flags, "sso-secret-env-next") {
None => None,
Some(nvar) => {
let v = std::env::var(&nvar).map_err(|_| {
format!(
"--sso-secret-env-next {nvar}: environment variable is not set"
)
})?;
if v.trim().is_empty() {
return Err(format!(
"--sso-secret-env-next {nvar}: secret is empty"
));
}
if v == secret {
return Err(format!(
"--sso-secret-env-next {nvar} holds the same value as \
--sso-secret-env {var} — that is a rotation that did not \
happen, and it would read as 'both live' everywhere"
));
}
Some(v)
}
};
let rotating = next.is_some();
server = server.with_sso_rotating(&header, secret, next);
eprintln!(
"areev: trusted-header SSO enabled ({header} + x-areev-proxy-secret)"
);
server = server.with_sso_approvals(sso_approvals);
if let Some(gh) = flag(&flags, "sso-groups-header") {
if flag(&flags, "auth").is_none() {
return Err(
"--sso-groups-header needs --auth <map>: the group → principal \
table lives in the credential map, and without one every group \
would resolve to nothing"
.into(),
);
}
server = server.with_sso_groups_header(&gh);
eprintln!(
"areev: SSO groups honored via {gh} (group-derived principals are \
roles: they may never answer a HITL approval)"
);
}
if let Some(prefix) = flag(&flags, "sso-principal-prefix") {
if prefix.trim().is_empty() {
return Err("--sso-principal-prefix: empty prefix".into());
}
server = server.with_sso_principal_prefix(&prefix);
eprintln!("areev: proxy-asserted principals are prefixed {prefix:?}");
}
if sso_approvals {
eprintln!(
"areev: ⚠ --sso-approvals allow — a proxy-asserted identity may \
answer HITL approvals. Every approval's audit record is then only \
as strong as the proxy secret: anyone holding it can approve as \
anyone. Prefer per-principal credentials (--auth) for approvers."
);
} else {
eprintln!(
"areev: SSO identities may review but NOT approve (run.respond); \
approvers need a per-principal credential (--auth). Override with \
--sso-approvals allow."
);
}
if rotating {
eprintln!(
"areev: ⚠ SSO secret ROTATION WINDOW open — two proxy secrets are \
accepted. Retire the old one (drop --sso-secret-env-next and \
promote its value) as soon as every proxy presents the new one: \
docs/runbooks/sso-secret-rotation.md"
);
}
}
(None, None) => {}
_ => return Err("--sso-header and --sso-secret-env must be given together".into()),
}
if let Some(issuer) = flag(&flags, "oidc-issuer") {
#[cfg(feature = "oidc")]
{
let client_id = flag(&flags, "oidc-client-id").ok_or(
"--oidc-issuer needs --oidc-client-id".to_string(),
)?;
let secret_var = flag(&flags, "oidc-client-secret-env").ok_or(
"--oidc-issuer needs --oidc-client-secret-env VAR (the secret is \
named, never passed on the command line)"
.to_string(),
)?;
let client_secret = std::env::var(&secret_var).map_err(|_| {
format!(
"--oidc-client-secret-env {secret_var}: environment variable is \
not set"
)
})?;
if client_secret.trim().is_empty() {
return Err(format!(
"--oidc-client-secret-env {secret_var}: secret is empty"
));
}
let redirect_uri = flag(&flags, "oidc-redirect-uri").ok_or(
"--oidc-issuer needs --oidc-redirect-uri (registered VERBATIM with \
the provider — RFC 9700 requires exact matching)"
.to_string(),
)?;
if !issuer.starts_with("https://") {
return Err(format!("--oidc-issuer {issuer}: must be https"));
}
let cfg = areev_server::oidc::OidcConfig {
issuer,
client_id,
client_secret,
redirect_uri: redirect_uri.clone(),
scopes: flag(&flags, "oidc-scopes")
.unwrap_or_else(|| "openid email profile".to_string()),
principal_claim: flag(&flags, "oidc-principal-claim")
.unwrap_or_else(|| "email".to_string()),
principal_prefix: flag(&flags, "oidc-principal-prefix"),
};
let provider = areev_server::oidc::OidcProvider::discover(cfg)
.map_err(|e| e.to_string())?;
eprintln!("areev: native OIDC enabled — login at /auth/login");
if !redirect_uri.starts_with("https://") {
eprintln!(
"areev: ⚠ --oidc-redirect-uri is not https, so the session \
cookie cannot be marked Secure. Acceptable on loopback; on any \
real deployment terminate TLS first."
);
}
server = server.with_oidc(provider);
}
#[cfg(not(feature = "oidc"))]
{
let _ = issuer;
return Err(
"this build has no native OIDC — rebuild with `--features oidc`, or \
use trusted-header SSO behind an authenticating proxy \
(--sso-header), which is the documented default"
.into(),
);
}
}
match (flag(&flags, "tls-cert"), flag(&flags, "tls-key")) {
(Some(cert), Some(key)) => {
#[cfg(feature = "tls")]
{
server = server.with_tls(&cert, &key).map_err(|e| e.to_string())?;
eprintln!("areev: native TLS enabled (cert {cert})");
}
#[cfg(not(feature = "tls"))]
{
let _ = (&cert, &key);
return Err(
"this build has no native TLS — rebuild with `--features tls`, \
or use the documented TLS-terminating proxy"
.into(),
);
}
}
(None, None) => {}
_ => return Err("--tls-cert and --tls-key must be given together".into()),
}
let listener = areev_server::UiServer::bind(&addr).map_err(|e| e.to_string())?;
if !addr_is_loopback(&addr) {
eprintln!(
"areev: WARNING — bound to non-loopback {addr} over plaintext HTTP; {}",
if has_auth {
"the token crosses the wire in the clear — use a TLS-terminating proxy."
} else {
"UNAUTHENTICATED — anyone who can reach it can read/modify memory."
}
);
}
eprintln!(
"areev console → http://{} (Ctrl-C to stop)",
listener.local_addr().map_err(|e| e.to_string())?
);
server.serve(listener).map_err(|e| e.to_string())?;
}
"get" => {
let h = positional
.first()
.ok_or_else(|| "usage: areev get <hash> --db <file>".to_string())?;
let hash = Hash::from_hex(h).map_err(|e| e.to_string())?;
let g = m.get(&hash).map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"hash": g.hash.to_hex(),
"type": format!("{:?}", g.grain_type).to_lowercase(),
"fields": g.fields,
})
);
}
"provenance" => {
let h = flag(&flags, "of")
.or_else(|| positional.first().cloned())
.ok_or_else(|| {
"usage: areev provenance <source-hash> --db <file> \
(lists grains whose derived_from is that hash)"
.to_string()
})?;
let h = h.strip_prefix("sha256:").unwrap_or(&h);
let parent = Hash::from_hex(h).map_err(|e| e.to_string())?;
let kids = m.grains_derived_from(&parent).map_err(|e| e.to_string())?;
for g in kids {
println!(
"{}",
serde_json::json!({
"hash": g.hash.to_hex(),
"type": format!("{:?}", g.grain_type).to_lowercase(),
"subject": g.get_str("subject"),
"relation": g.get_str("relation"),
"object": g.get_str("object"),
})
);
}
}
"forks" => {
let forks = m.open_forks().map_err(|e| e.to_string())?;
if forks.is_empty() {
eprintln!("no open forks");
}
for f in forks {
println!(
"{}",
serde_json::json!({
"namespace": f.namespace,
"subject": f.subject,
"relation": f.relation,
"heads": f.heads.iter().map(|h| h.to_hex()).collect::<Vec<_>>(),
})
);
}
}
"subject-report" => {
let subject = positional.first().cloned().or_else(|| flag(&flags, "subject")).ok_or(
"usage: areev subject-report <subject> [--ns NS] [--text-mentions] \
[--out FILE] [--bundle FILE]"
.to_string(),
)?;
let opts = areev_store::ErasureOptions {
text_mentions: flag(&flags, "text-mentions")
.map(|v| !matches!(v.as_str(), "false" | "0" | "off" | "no"))
.unwrap_or(false),
};
if let Some(principal) = flag(&flags, "as") {
let grants = m.authz_grants(&principal).map_err(|e| e.to_string())?;
areev_core::authz::AuthzSet::restricted(&principal, grants)
.check(areev_core::authz::Verb::Read, &ns)
.map_err(|e| e.to_string())?;
}
let report = m.subject_report_with(&ns, &subject, opts).map_err(|e| e.to_string())?;
let mut lines = String::new();
for g in &report.grains {
let fields: serde_json::Map<String, serde_json::Value> =
g.fields.clone().into_iter().collect();
let row = serde_json::json!({
"hash": g.hash.to_hex(),
"type": g.grain_type.as_str(),
"fields": fields,
});
lines.push_str(&row.to_string());
lines.push('\n');
}
match flag(&flags, "out") {
Some(path) => {
std::fs::write(&path, &lines).map_err(|e| e.to_string())?;
println!("wrote {} grains to {path}", report.grains.len());
}
None => print!("{lines}"),
}
if let Some(bpath) = flag(&flags, "bundle") {
let stats = m
.subject_bundle_with(&ns, &subject, opts, &bpath)
.map_err(|e| e.to_string())?;
eprintln!("bundle: {} records, {} bytes -> {bpath}", stats.ops, stats.bytes);
}
eprintln!(
"subject-report '{subject}' ns '{ns}': {} grains; matched identities: {}",
report.grains.len(),
report.identity_names.join(", ")
);
}
"forget-subject" => {
let subject = positional.first().cloned().or_else(|| flag(&flags, "subject")).ok_or(
"usage: areev forget-subject <subject> [--ns NS] [--text-mentions] --yes"
.to_string(),
)?;
if flags.contains_key("no-destructive-ops") {
return Err(
"destructive operations are disabled for this invocation \
(--no-destructive-ops)"
.into(),
);
}
if !flags.contains_key("yes") {
return Err(format!(
"forget-subject erases EVERY grain referencing '{subject}' in namespace \
'{ns}' (history included) and replicates the erasure to peers. \
Re-run with --yes to proceed."
));
}
let opts = areev_store::ErasureOptions {
text_mentions: flag(&flags, "text-mentions")
.map(|v| !matches!(v.as_str(), "false" | "0" | "off" | "no"))
.unwrap_or(false),
};
if let Some(principal) = flag(&flags, "as") {
let grants = m.authz_grants(&principal).map_err(|e| e.to_string())?;
areev_core::authz::AuthzSet::restricted(&principal, grants)
.check(areev_core::authz::Verb::Erase, &ns)
.map_err(|e| e.to_string())?;
}
let override_hold = flags.contains_key("override-hold");
let audit_target = format!(
"subject:{} ns:{ns}",
areev_core::authz::subject_fingerprint(&subject)
);
let (rep, overridden) = if override_hold {
let because = flag(&flags, "because").ok_or_else(|| {
"--override-hold requires --because: destroying records someone \
placed a hold on is a decision, and it has a ground"
.to_string()
})?;
let who = flag(&flags, "as").unwrap_or_else(|| "local:owner".to_string());
let over = areev_store::HoldOverride::new(who, because)
.map_err(|e| e.to_string())?;
m.forget_subject_overriding(&ns, &subject, opts, &over)
.map_err(|e| e.to_string())?
} else {
match m.forget_subject_with(&ns, &subject, opts) {
Ok(r) => (r, None),
Err(e) => {
if e.code() == "STO-E009" {
let mut obs = areev_core::authz::audit_observation(
&flag(&flags, "as").unwrap_or_else(|| "local:owner".into()),
"erase.refused",
&audit_target,
Some(&e.to_string()),
0,
areev_core::time::now_ms(),
);
let _ = m.append_audit(&mut obs);
}
return Err(e.to_string());
}
}
};
if let Some(hold) = &overridden {
eprintln!(
"areev: overrode the legal hold on '{}' placed by {} ({})",
hold.ns, hold.placed_by, hold.because
);
}
let stale_exports = m
.corpus_exports_touching_subject(&subject)
.unwrap_or_else(|e| {
eprintln!("areev: warning: could not consult corpus registry: {e}");
Vec::new()
});
write_erasure_audit(
&mut m,
&flags,
"erase",
&audit_target,
rep.grains_erased,
&stale_exports,
overridden.as_ref(),
);
println!(
"erased {} grains ({} dictionary entries, {} vocabulary tokens, {} blobs reclaimed)",
rep.grains_erased, rep.terms_removed, rep.vocab_removed, rep.blobs_reclaimed
);
}
"purge-older-than" => {
let days: i64 = positional
.first()
.and_then(|v| v.parse().ok())
.or_else(|| flag(&flags, "days").and_then(|v| v.parse().ok()))
.ok_or("usage: areev purge-older-than <days> [--ns NS] [--type event] --yes".to_string())?;
let gt = match flag(&flags, "type") {
Some(t) => Some(
areev_core::types::GrainType::from_str(&t)
.ok_or_else(|| format!("unknown grain type '{t}'"))?,
),
None => None,
};
if flags.contains_key("no-destructive-ops") {
return Err(
"destructive operations are disabled for this invocation \
(--no-destructive-ops)"
.into(),
);
}
if !flags.contains_key("yes") {
let scope = if ns.is_empty() {
"across EVERY namespace".to_string()
} else {
format!("in namespace '{ns}'")
};
return Err(format!(
"purge-older-than erases every{} grain older than {days} days {scope} \
and replicates the erasure to peers. Re-run with --yes to proceed.",
gt.map(|g| format!(" {g:?}")).unwrap_or_default()
));
}
let cutoff = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0)
- days * 24 * 3600 * 1000;
let ns_opt = if ns.is_empty() { None } else { Some(ns.as_str()) };
if let Some(principal) = flag(&flags, "as") {
let grants = m.authz_grants(&principal).map_err(|e| e.to_string())?;
areev_core::authz::AuthzSet::restricted(&principal, grants)
.check(areev_core::authz::Verb::Erase, ns_opt.unwrap_or("*"))
.map_err(|e| e.to_string())?;
}
let exports = m.corpus_exports().unwrap_or_else(|e| {
eprintln!("areev: warning: could not consult corpus registry: {e}");
Vec::new()
});
let rep = m.forget_older_than(ns_opt, cutoff, gt).map_err(|e| e.to_string())?;
let stale_exports =
areev_store::exports_touching_hashes(&exports, &rep.erased_hashes);
write_erasure_audit(
&mut m,
&flags,
"erase",
&format!(
"older_than:{days}d ns:{}{}",
if ns.is_empty() { "*" } else { ns.as_str() },
flag(&flags, "type").map(|t| format!(" type:{t}")).unwrap_or_default()
),
rep.grains_erased,
&stale_exports,
None,
);
println!(
"erased {} grains ({} vocabulary tokens, {} blobs reclaimed)",
rep.grains_erased, rep.vocab_removed, rep.blobs_reclaimed
);
}
"merge" => {
let s = need(&flags, "subject")?;
let r = need(&flags, "relation")?;
let o = need(&flags, "object")?;
let conf: f64 = flag(&flags, "confidence")
.map(|c| c.parse().map_err(|_| "bad --confidence".to_string()))
.transpose()?
.unwrap_or(0.9);
let mut merged = Fact::new(&s, &r, &o).confidence(conf);
merged.common.namespace = Some(ns.clone());
let h = m
.merge_heads(&ns, &s, &r, &mut merged)
.map_err(|e| e.to_string())?;
println!("{h}");
}
"novelty" => {
let text = need(&flags, "text")?;
let subject = flag(&flags, "subject");
let relation = flag(&flags, "relation");
let k: usize = flag(&flags, "k").and_then(|v| v.parse().ok()).unwrap_or(5);
let matches = m
.nearest_semantic(&ns, subject.as_deref(), relation.as_deref(), &text, k)
.map_err(|e| e.to_string())?;
for (h, sim) in matches {
let object = m.get(&h).ok().and_then(|g| g.get_str("object").map(String::from));
println!(
"{}",
serde_json::json!({
"hash": h.to_hex(),
"similarity": (sim * 1000.0).round() / 1000.0,
"object": object,
})
);
}
}
"init" => {
run_init(m, &ns, &flags, &positional)?;
}
"loop" => {
run_loop(m, &ns, &flags, &positional)?;
}
"audit" => {
run_audit(&mut m, &flags, &positional)?;
}
"anonymize" => run_anonymize(m, &flags, &positional)?,
"pack" => {
return pack::run_pack(Some(m), &ns, &flags, &positional);
}
"trigger" => {
return trigger_cli::run_trigger(m, &ns, &db, &flags, &positional);
}
"retention" => {
run_retention(&mut m, &ns, &flags, &positional)?;
}
"hold" => {
run_hold(&mut m, &ns, &flags, &positional)?;
}
"run" => {
return run_run(m, &db, &ns, &flags, &positional);
}
"eval" => {
return run_eval(m, &flags, &positional);
}
"tool" => {
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("");
let hash_arg = positional.get(1).cloned();
if sub_cmd != "provenance" || hash_arg.is_none() {
return Err("usage: areev tool provenance <hash> [--depth N]".into());
}
return run_tool_provenance(m, &db, &ns, &hash_arg.unwrap(), &flags);
}
"blob" => {
match positional.first().map(String::as_str).unwrap_or("") {
"put" => {
let bytes = match positional.get(1) {
Some(path) => std::fs::read(path)
.map_err(|e| format!("blob put: reading {path}: {e}"))?,
None => {
if !flags.contains_key("stdin") {
return Err("usage: areev blob put <FILE> --db <file> (or --stdin)".into());
}
use std::io::Read;
let mut buf = Vec::new();
std::io::stdin()
.read_to_end(&mut buf)
.map_err(|e| format!("blob put: reading stdin: {e}"))?;
buf
}
};
println!("{}", m.put_blob(&bytes).map_err(|e| e.to_string())?);
}
"get" => {
let uri = positional
.get(1)
.ok_or("usage: areev blob get <cas-uri> --db <file>")?;
use std::io::Write;
let bytes = m.get_blob(uri).map_err(|e| e.to_string())?;
std::io::stdout().write_all(&bytes).map_err(|e| e.to_string())?;
}
other => {
return Err(format!(
"unknown blob subcommand '{other}' — try `areev blob put|get`"
))
}
}
}
"blobs" => {
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("");
if sub_cmd != "encrypt" {
return Err("usage: areev blobs encrypt --db <file> --passphrase-env VAR".into());
}
let (encrypted, already) = m.encrypt_blobs().map_err(|e| e.to_string())?;
println!(
"encrypted {encrypted} attachment(s); {already} were already encrypted"
);
}
other => return Err(format!("unknown command '{other}' — try `areev help`")),
}
Ok(())
}
fn fixture_strings(root: &serde_json::Value, key: &str, path: &str) -> Result<Vec<String>, String> {
match root.get(key) {
None => Ok(Vec::new()),
Some(serde_json::Value::Array(items)) => items
.iter()
.map(|v| {
v.as_str()
.map(str::to_string)
.ok_or_else(|| format!("{path}: every \"{key}\" entry must be a string"))
})
.collect(),
Some(_) => Err(format!("{path}: \"{key}\" must be an array of strings")),
}
}
fn run_anonymize(
mut m: Areev,
flags: &HashMap<String, String>,
positional: &[String],
) -> Result<(), String> {
let sub = positional.first().map(String::as_str).unwrap_or("");
match sub {
"set" => {
let ns = flag(flags, "ns").ok_or("anonymize set requires --ns")?;
let policy = match (flag(flags, "policy"), flag(flags, "policy-file")) {
(Some(j), _) => j,
(None, Some(path)) => std::fs::read_to_string(&path)
.map_err(|e| format!("--policy-file {path}: {e}"))?,
(None, None) => return Err("anonymize set requires --policy or --policy-file".into()),
};
m.set_anon_policy(&ns, &policy).map_err(|e| e.to_string())?;
println!("anonymization policy declared for namespace '{ns}' — egress reads are transformed from now on; older builds opening this file will warn");
Ok(())
}
"list" => {
let policies = m.anon_policies().map_err(|e| e.to_string())?;
if policies.is_empty() {
println!("no anonymization policies declared");
return Ok(());
}
for (ns, p) in policies {
println!(
"{}",
serde_json::json!({"ns": ns, "policy": p})
);
}
Ok(())
}
"clear" => {
let ns = flag(flags, "ns").ok_or("anonymize clear requires --ns")?;
m.clear_anon_policy(&ns).map_err(|e| e.to_string())?;
println!("anonymization policy cleared for namespace '{ns}'");
Ok(())
}
"reveal" => {
let ns = flag(flags, "ns").ok_or("anonymize reveal requires --ns")?;
let token = flag(flags, "token").ok_or("anonymize reveal requires --token")?;
let facade =
areev_cal::AreevFacade::with_session(m, Some(ns.clone()), None);
let out = facade
.reveal_tokens(&ns, &[token])
.map_err(|e| e.to_string())?;
println!("{out}");
Ok(())
}
"mappings" => {
let maps = m.anon_mappings().map_err(|e| e.to_string())?;
for (ns, id, mapping) in maps {
println!(
"{}",
serde_json::json!({"ns": ns, "mapping_id": id, "mapping": mapping})
);
}
Ok(())
}
other => Err(format!(
"unknown anonymize subcommand '{other}' (expected scan, set, list, clear, or mappings)"
)),
}
}
fn run_retention(
m: &mut Areev,
ns: &str,
flags: &HashMap<String, String>,
positional: &[String],
) -> Result<(), String> {
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("list");
match sub_cmd {
"set" => {
let days: f64 = flag(flags, "days")
.and_then(|v| v.parse().ok())
.or_else(|| positional.get(1).and_then(|v| v.parse().ok()))
.ok_or_else(|| {
"usage: areev retention set --days N [--ns NS] [--type event] \
[--because \"why\"]"
.to_string()
})?;
let policy = areev_store::RetentionPolicy {
days,
grain_type: flag(flags, "type"),
because: flag(flags, "because"),
};
m.set_retention_policy(ns, &policy).map_err(|e| e.to_string())?;
println!(
"retention policy for '{ns}': {days} days{}{}",
policy.grain_type.map(|t| format!(", type {t}")).unwrap_or_default(),
policy.because.map(|b| format!(" — {b}")).unwrap_or_default()
);
eprintln!(
"areev: declared, not enforced — run `areev retention sweep --yes` (cron), \
or let the loop's retention_sweep analyzer propose it for review"
);
}
"list" => {
let policies = m.retention_policies().map_err(|e| e.to_string())?;
if policies.is_empty() {
println!("no retention policies declared in this memory");
}
for (pns, p) in policies {
println!(
"{}",
serde_json::json!({
"namespace": pns,
"days": p.days,
"type": p.grain_type,
"because": p.because,
})
);
}
}
"clear" => {
m.clear_retention_policy(ns).map_err(|e| e.to_string())?;
println!("cleared the retention policy for '{ns}'");
}
"sweep" => {
if flags.contains_key("no-destructive-ops") {
return Err(
"destructive operations are disabled for this invocation \
(--no-destructive-ops)"
.into(),
);
}
let policies = m.retention_policies().map_err(|e| e.to_string())?;
if policies.is_empty() {
println!("no retention policies declared — nothing to sweep");
return Ok(());
}
if !flags.contains_key("yes") {
let plan: Vec<String> = policies
.iter()
.map(|(pns, p)| {
format!(
" {pns}: erase grains older than {} days{}",
p.days,
p.grain_type.as_deref().map(|t| format!(" (type {t})")).unwrap_or_default()
)
})
.collect();
return Err(format!(
"retention sweep would apply {} policy/policies:\n{}\nRe-run with --yes \
to proceed.",
policies.len(),
plan.join("\n")
));
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
let exports = m.corpus_exports().unwrap_or_else(|e| {
eprintln!("areev: warning: could not consult corpus registry: {e}");
Vec::new()
});
let results = m.sweep_retention(now).map_err(|e| e.to_string())?;
let mut total = 0usize;
for sweep in &results {
let (pns, days) = (&sweep.ns, sweep.days);
match &sweep.outcome {
areev_store::RetentionOutcome::Skipped { why } => {
println!("{pns}: SKIPPED — {why}");
continue;
}
areev_store::RetentionOutcome::Swept(rep) => {
total += rep.grains_erased;
println!(
"{pns}: erased {} grains older than {days} days ({} vocabulary tokens, \
{} blobs reclaimed)",
rep.grains_erased, rep.vocab_removed, rep.blobs_reclaimed
);
if rep.grains_erased > 0 {
let stale_exports = areev_store::exports_touching_hashes(
&exports,
&rep.erased_hashes,
);
write_erasure_audit(
m,
flags,
"erase",
&format!("retention:{days}d ns:{pns}"),
rep.grains_erased,
&stale_exports,
None,
);
}
}
}
}
println!("retention sweep: {total} grains erased across {} policies", results.len());
}
"floor" => {
let min_days: f64 = flag(flags, "min-days")
.and_then(|v| v.parse().ok())
.ok_or_else(|| {
"usage: areev retention floor --min-days N --because \"why\" [--ns NS]"
.to_string()
})?;
let because = flag(flags, "because").ok_or_else(|| {
"retention floor requires --because (a floor with no recorded \
rationale is not auditable)"
.to_string()
})?;
m.set_retention_floor(ns, min_days, &because)
.map_err(|e| e.to_string())?;
println!("retention floor for '{ns}': ≥{min_days} days — {because}");
}
"floor-clear" => {
m.clear_retention_floor(ns).map_err(|e| e.to_string())?;
println!("retention floor cleared for '{ns}'");
}
"floors" => {
let floors = m.retention_floors().map_err(|e| e.to_string())?;
if floors.is_empty() {
println!("no retention floors declared in this memory");
}
for (fns, days, because) in floors {
println!("{fns}: ≥{days} days — {because}");
}
}
other => {
return Err(format!(
"unknown retention subcommand '{other}' — usage: areev retention \
<set|list|clear|sweep|floor|floor-clear|floors>"
))
}
}
Ok(())
}
fn run_run(
m: Areev,
db: &str,
ns: &str,
flags: &HashMap<String, String>,
positional: &[String],
) -> Result<(), String> {
use areev_run::{RunSession, Runner, SystemClock};
use std::sync::Arc;
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("");
let principal = flag(flags, "as").unwrap_or_else(|| "user:local".to_string());
let mut broker_guard: Option<Arc<areev_run::Broker>> = None;
let egress: Option<areev_run::EgressHandle> = match run_stack::build_egress(flags)? {
None => None,
Some(broker) => {
broker.serve_blobs(db);
let broker = Arc::new(broker);
broker_guard = Some(Arc::clone(&broker));
Some(areev_run::EgressHandle::new(broker))
}
};
let executor = run_stack::tool_executor(flags, egress.as_ref());
let llm = run_stack::toolcall_llm(flags)?;
let observer = run_stack::observer(flags)?;
let decider = run_stack::decider(flags, &m)?;
let runner = Runner {
facade: Arc::new(run_stack::host_facade(m, Some(ns.to_string()), flags)?),
clock: Arc::new(SystemClock),
executor,
llm,
observer,
ns: ns.to_string(),
principal: principal.clone(),
};
let runner = match decider {
Some(d) => runner.with_decider(d),
None => runner,
};
let opts = run_stack::run_options(flags);
let need = |key: &str, usage: &str| -> Result<String, String> {
flag(flags, key).ok_or_else(|| format!("usage: {usage}"))
};
let print_session = |session: RunSession| match session {
RunSession::Finished { outcome, run_id } => {
println!(
"{}",
serde_json::json!({"run_id": run_id, "finished": format!("{outcome:?}")})
);
}
RunSession::Parked { envelope, run_id } => {
println!("{envelope}");
if envelope["kind"] == "paused" {
eprintln!(
"areev: run paused by its host — `areev run resume --run-id {run_id}` \
continues it under the same run id"
);
} else {
eprintln!(
"areev: run parked — answer with `areev run respond --run-id ID \
--ask TOOL_CALL_ID --result JSON --as PRINCIPAL`, then \
`areev run resume --run-id ID`"
);
}
}
};
let report_refusals = run_stack::report_refusals;
match sub_cmd {
"start" => {
let wf = need("workflow", "areev run start --workflow HASH --run-id ID [--input JSON] [--tool-cmd CMD] [--as PRINCIPAL]")?;
let run_id = need("run-id", "areev run start --workflow HASH --run-id ID")?;
let input: serde_json::Value = match flag(flags, "input") {
Some(raw) => serde_json::from_str(&raw)
.map_err(|e| format!("--input must be valid JSON: {e}"))?,
None => serde_json::json!({}),
};
let h = Hash::from_hex(&wf).map_err(|e| e.to_string())?;
if let Ok(plan_grain) = runner.facade.with_store(|m| m.get(&h)) {
if let Some(plan_ns) = plan_grain.get_str("namespace") {
if plan_ns != runner.ns {
eprintln!(
"areev: note: this plan lives in namespace '{plan_ns}', but the run \
reads and journals in '{}' — pass `--ns {plan_ns}` if its record \
should sit with the plan",
runner.ns
);
}
}
}
let session = runner.start(&h, &run_id, input, &opts).map_err(|e| e.to_string())?;
print_session(session);
}
"resume" => {
let run_id = need("run-id", "areev run resume --run-id ID [--tool-cmd CMD]")?;
let session = runner.resume(&run_id, &opts).map_err(|e| e.to_string())?;
print_session(session);
}
"respond" => {
let run_id = need("run-id", "areev run respond --run-id ID --ask TOOL_CALL_ID --result JSON --as PRINCIPAL")?;
let ask = need("ask", "areev run respond --run-id ID --ask TOOL_CALL_ID --result JSON --as PRINCIPAL")?;
let result: serde_json::Value = serde_json::from_str(
&need("result", "areev run respond … --result JSON")?,
)
.map_err(|e| format!("--result must be valid JSON: {e}"))?;
let is_error = flag(flags, "is-error")
.map(|v| !matches!(v.as_str(), "false" | "0" | "off" | "no"))
.unwrap_or(false);
runner
.respond(&run_id, &ask, result, is_error, &principal)
.map_err(|e| e.to_string())?;
println!("response recorded — `areev run resume --run-id {run_id}` to continue");
}
"input" => {
let run_id = need("run-id", "areev run input --run-id ID --message TEXT")?;
let message = need("message", "areev run input --run-id ID --message TEXT")?;
runner.input(&run_id, &message, &principal).map_err(|e| e.to_string())?;
println!(
"steering message queued for '{run_id}' — the next superstep hands it \
to its nodes under `$inbox`"
);
}
"cancel" => {
let run_id = need("run-id", "areev run cancel --run-id ID [--because \"why\"]")?;
let because = flag(flags, "because").unwrap_or_else(|| "canceled".into());
runner
.cancel(&run_id, &principal, &because)
.map_err(|e| e.to_string())?;
println!(
"cancel recorded for '{run_id}' — a live driver drains at its next \
superstep boundary; `areev run resume --run-id {run_id}` finalizes a \
parked one"
);
}
"pause" => {
let usage = "areev run pause --run-id ID [--because \"why\"]";
let run_id = need("run-id", usage)?;
let because = flag(flags, "because").unwrap_or_else(|| "paused".into());
let receipt = runner
.pause(&run_id, &principal, &because)
.map_err(|e| e.to_string())?;
println!("{}", serde_json::to_string(&receipt).unwrap());
eprintln!(
"areev: pause {} for '{run_id}' — a live driver parks at its next \
superstep boundary; `areev run resume --run-id {run_id}` continues it",
if receipt.already { "already standing" } else { "recorded" }
);
}
"verify" => {
let run_id = need("run-id", "areev run verify --run-id ID")?;
let report = runner.verify(&run_id).map_err(|e| e.to_string())?;
println!("{}", serde_json::to_string_pretty(&report).unwrap());
if !report.verified {
return Err("verification did not pass".into());
}
}
"list" => {
let n = flag(flags, "last").and_then(|v| v.parse().ok()).unwrap_or(20);
let offset = flag(flags, "offset").and_then(|v| v.parse().ok()).unwrap_or(0);
let scope = flags.contains_key("ns").then_some(ns);
let page = runner.list_runs(scope, offset, n).map_err(|e| e.to_string())?;
let rows: Vec<serde_json::Value> = page
.runs
.iter()
.map(|r| {
let has_manifest = runner
.facade
.with_store(|m| {
m.latest("agent:harness", &format!("run:{}", r.run_id), "mg:harness")
})
.ok()
.flatten()
.is_some();
serde_json::json!({
"run_id": r.run_id,
"ns": r.ns,
"outcome": r.outcome,
"usd_micros": r.spent_usd_micros,
"tokens": r.spent_tokens,
"has_manifest": has_manifest,
})
})
.collect();
println!("{}", serde_json::to_string_pretty(&rows).unwrap());
if page.truncated {
eprintln!(
"areev: showing {} of {} runs (newest first) — raise --last or page with --offset",
rows.len(),
page.total
);
}
if page.unattributed > 0 {
eprintln!(
"areev: {} more run(s) predate the session-namespace stamp and cannot be \
attributed to a namespace — list them without --ns",
page.unattributed
);
}
}
"inspect" => {
let run_id = need("run-id", "areev run inspect --run-id ID")?;
let report = runner.inspect(&run_id).map_err(|e| e.to_string())?;
println!("{}", serde_json::to_string_pretty(&report).unwrap());
}
"oversight-report" => {
let plan_hash = flag(flags, "plan")
.map(|p| {
areev_core::error::Hash::from_hex(&p)
.map_err(|_| format!("--plan is not a valid hash: {p}"))
})
.transpose()?;
let report = runner
.oversight_report(flag(flags, "run-id").as_deref(), plan_hash.as_ref())
.map_err(|e| e.to_string())?;
println!("{}", serde_json::to_string_pretty(&report).unwrap());
}
"shadow" => {
let ids: Vec<String> = match flag(flags, "runs") {
Some(csv) => csv.split(',').map(|s| s.trim().to_string()).filter(|s| !s.is_empty()).collect(),
None => {
let n = flag(flags, "last").and_then(|v| v.parse().ok()).unwrap_or(5);
runner.recent_runs(n).map_err(|e| e.to_string())?
}
};
if ids.is_empty() {
return Err("no runs to shadow (give --runs a,b or record some runs first)".into());
}
let candidate = match (flag(flags, "plan"), flag(flags, "plan-file")) {
(Some(_), Some(_)) => return Err("give --plan or --plan-file, not both".into()),
(Some(h), None) => Some(areev_run::PlanCandidate::Hash(
Hash::from_hex(&h).map_err(|e| e.to_string())?,
)),
(None, Some(path)) => {
let text = std::fs::read_to_string(&path).map_err(|e| format!("{path}: {e}"))?;
let body: serde_json::Value =
serde_json::from_str(&text).map_err(|e| format!("{path}: {e}"))?;
let fields = body
.as_object()
.cloned()
.ok_or_else(|| format!("{path}: a plan draft must be a JSON object"))?;
Some(areev_run::PlanCandidate::Body(fields))
}
(None, None) => None,
};
let reexecute = match flag(flags, "reexecute") {
None => areev_run::Reexecute::Off,
Some(mode) => areev_run::Reexecute::parse(&mode).ok_or_else(|| {
format!("--reexecute takes 'pure' or 'off', not {mode:?}")
})?,
};
if let Some(candidate) = candidate {
let opts = areev_run::ShadowOptions::reexecute(reexecute);
let report =
runner.shadow_plan_with(&ids, &candidate, &opts).map_err(|e| e.to_string())?;
println!("{}", serde_json::to_string_pretty(&report).unwrap());
return Ok(());
}
if reexecute != areev_run::Reexecute::Off {
return Err(
"--reexecute needs a candidate: add --plan HASH or --plan-file draft.json"
.into(),
);
}
let report = runner.shadow_eval(&ids);
println!("{}", serde_json::to_string_pretty(&report).unwrap());
if !report.all_consistent {
return Err("shadow evaluation found inconsistent runs".into());
}
}
"fork" => {
let usage = "areev run fork --run-id BASE --as-run NEW [--at SUPERSTEP] [--plan HASH]";
let base = need("run-id", usage)?;
let new_id = need("as-run", usage)?;
let at = flag(flags, "at").and_then(|v| v.parse().ok());
let plan = match flag(flags, "plan") {
Some(h) => Some(Hash::from_hex(&h).map_err(|e| e.to_string())?),
None => None,
};
let seed = runner
.fork(&base, at, &new_id, plan.as_ref(), &opts)
.map_err(|e| e.to_string())?;
println!(
"fork '{new_id}' created from '{base}' (seed checkpoint {seed}) — \
`areev run resume --run-id {new_id}` continues it"
);
}
"demo" => {
let (wf_hash, greet_hash, approve_hash) = runner.facade.with_store(|m| {
let greet = areev_core::types::Tool::new("greet")
.kind(areev_core::types::ToolKind::Definition)
.tool_description("Greets the subject (host-executed)")
.created_at(1_000)
.namespace(ns);
let gh = m.add(&greet)?;
let approve = areev_core::types::Tool::new("approve")
.kind(areev_core::types::ToolKind::Definition)
.tool_description("A human approves the greeting")
.executor_kind(areev_core::types::ExecutorKind::Client)
.created_at(1_001)
.namespace(ns);
let ah = m.add(&approve)?;
let wf = areev_core::types::Workflow::new(vec!["greet".into(), "approve".into()])
.edge("greet", "approve")
.bind("greet", &gh.to_hex())
.bind("approve", &ah.to_hex())
.created_at(1_002)
.namespace(ns);
let wh = m.add(&wf)?;
Ok::<_, areev_core::error::AreevError>((wh, gh, ah))
})
.map_err(|e| e.to_string())?;
println!("demo plan seeded (workflow {wf_hash})");
println!(" greet -> host tool {greet_hash}");
println!(" approve -> human approval {approve_hash}");
println!();
println!("the 10-minute proof (no LLM key required):");
println!(
" 1. areev run start --db <DB> --workflow {wf_hash} --run-id demo-1 \\\n\
\x20 --input '{{\"who\":\"world\"}}' --tool-cmd 'printf '\\''{{\"greeting\":\"hello\"}}'\\'''"
);
println!(" → parks with a requires_action envelope; copy its tool_call_id");
println!(
" 2. areev run respond --db <DB> --run-id demo-1 --ask <TOOL_CALL_ID> \\\n\
\x20 --result '{{\"approved\":true}}' --as user:officer"
);
println!(" (responding as the starter is REFUSED — separation of duties)");
println!(" 3. areev run resume --db <DB> --run-id demo-1 → Completed");
println!(" 4. areev run verify --db <DB> --run-id demo-1 → journal-consistent replay");
println!(" 5. areev run-trace --db <DB> --run-id demo-1 → the full journal");
}
other => {
return Err(format!(
"unknown run subcommand '{other}' — usage: areev run \
<start|resume|respond|input|pause|cancel|list|inspect|verify|fork|shadow|oversight-report|demo>"
))
}
}
report_refusals(&broker_guard);
Ok(())
}
fn run_tool_provenance(
mut m: Areev,
db_path: &str,
ns: &str,
hash_arg: &str,
flags: &HashMap<String, String>,
) -> Result<(), String> {
use areev_loop::SubstrateRead;
let depth = flag(flags, "depth").and_then(|v| v.parse().ok()).unwrap_or(2);
let mut executor: Option<serde_json::Value> = None;
let runs = match Hash::from_hex(hash_arg) {
Ok(h) => {
executor = m
.get(&h)
.ok()
.and_then(|g| g.to_tool().ok())
.and_then(|t| t.executor_uri.clone())
.map(|uri| executor_blob_report(db_path, &uri));
let mut all = m.runs_touching(ns, &h, depth).unwrap_or_default();
if ns != "agent:harness" {
all.extend(m.runs_touching("agent:harness", &h, depth).unwrap_or_default());
}
all.sort();
all.dedup();
all
}
Err(_) => Vec::new(),
};
let sub = AreevSubstrate::new(m, None);
let engine = Engine::with_builtins();
let target = format!("tool:{hash_arg}");
let recs: Vec<_> = engine
.recommendations(&sub, None)
.map_err(|e| e.to_string())?
.into_iter()
.filter(|r| r.target_ref == target)
.collect();
let audits = sub
.grains_of_type("observation", Some("areev-loop"), Default::default())
.unwrap_or_default();
let mut out_recs = Vec::new();
for r in &recs {
let mut chain: Vec<serde_json::Value> = audits
.iter()
.filter(|a| a.str_field("rec_hash") == Some(r.hash.as_str()))
.map(|a| {
let mut row = serde_json::json!({
"to": a.str_field("to_status"),
"actor": a.str_field("actor"),
"because": a.str_field("because"),
"at_ms": a.fields.get("at_ms"),
});
if let Some(es) = a.str_field("gating_evalset") {
row["gating"] = serde_json::json!({
"evalset": es,
"run_id": a.str_field("gating_run_id"),
"passed": a.fields.get("gating_passed"),
"failed": a.fields.get("gating_failed"),
});
}
row
})
.collect();
chain.sort_by_key(|c| c.get("at_ms").and_then(|v| v.as_i64()).unwrap_or(0));
out_recs.push(serde_json::json!({
"recommendation": r.hash,
"status": r.status.as_str(),
"summary": r.summary.render(),
"evalset_pin": r.evalset_hash,
"lifecycle": chain,
}));
}
let mut report = serde_json::json!({
"tool": hash_arg,
"recommendations": out_recs,
"runs_touching": runs,
"runs_touching_count": runs.len(),
});
if let Some(e) = executor {
report["executor"] = e;
}
println!("{}", serde_json::to_string_pretty(&report).unwrap());
if out_recs.is_empty() && runs.is_empty() {
eprintln!("note: nothing recorded for this hash yet (no recommendations target it; no runs touch it)");
}
Ok(())
}
fn executor_blob_report(db_path: &str, uri: &str) -> serde_json::Value {
let mut out = serde_json::json!({"executor_uri": uri});
if !uri.starts_with("cas://") {
out["blob_present"] = serde_json::json!(false);
out["note"] = serde_json::json!(
"not a cas:// address — the code is not content-addressed in this memory, \
so what ran cannot be verified from here"
);
return out;
}
match areev_store::read_blob_offline(db_path, uri) {
Ok(Some(bytes)) => {
out["blob_present"] = serde_json::json!(true);
out["blob_bytes"] = serde_json::json!(bytes.len());
}
Ok(None) => {
out["blob_present"] = serde_json::json!(true);
out["blob_sealed"] = serde_json::json!(true);
}
Err(e) => {
out["blob_present"] = serde_json::json!(false);
out["error"] = serde_json::json!(e.to_string());
}
}
out
}
fn run_eval(
m: Areev,
flags: &HashMap<String, String>,
positional: &[String],
) -> Result<(), String> {
const EVAL_NS: &str = "agent:harness";
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("");
let case_ns = flag(flags, "case-ns").unwrap_or_else(|| EVAL_NS.to_string());
let facade = run_stack::host_facade(m, Some(EVAL_NS.to_string()), flags)?;
let need = |key: &str, usage: &str| -> Result<String, String> {
flag(flags, key).ok_or_else(|| format!("usage: {usage}"))
};
match sub_cmd {
"create" => {
let usage = "areev eval create --name NAME --cases FILE";
let name = need("name", usage)?;
let path = need("cases", usage)?;
let text = std::fs::read_to_string(&path).map_err(|e| format!("{path}: {e}"))?;
let cases: Vec<serde_json::Value> =
serde_json::from_str(&text).map_err(|e| format!("cases must be a JSON list: {e}"))?;
if cases.is_empty() {
return Err("an evalset needs at least one case".into());
}
for (i, c) in cases.iter().enumerate() {
let ok = c.get("name").and_then(|v| v.as_str()).is_some()
&& c.get("input").is_some()
&& c.get("expect").is_some_and(|e| {
e.get("equals").is_some() || e.get("contains").and_then(|v| v.as_str()).is_some()
});
if !ok {
return Err(format!(
"case {i} must be {{\"name\", \"input\", \"expect\": {{\"equals\": …}} | {{\"contains\": \"…\"}}}}"
));
}
}
let payload = serde_json::json!({"name": name, "cases": cases});
let mut fact = areev_core::types::Fact::new(
&format!("evalset:{name}"),
"mg:evalset",
&payload.to_string(),
)
.namespace(&case_ns);
fact.common.extra_fields.insert("evalset_name".into(), serde_json::json!(name));
let h = facade.with_store(|m| m.add(&fact)).map_err(|e| e.to_string())?;
println!("evalset '{name}' stored: {h} ({} cases)", cases.len());
Ok(())
}
"run" => {
let usage = "areev eval run --evalset HASH (--tool-cmd CMD | --model provider:name \
[--base-url URL] [--key-env VAR] [--llm-max-tokens N]) \
[--baseline RUN_ID [--tolerance POINTS]]";
let evalset_hex = need("evalset", usage)?;
let cmd = flag(flags, "tool-cmd");
let model = flag(flags, "model");
let llm = match (&cmd, &model) {
(Some(_), Some(_)) => {
return Err("give either --tool-cmd or --model, not both".into());
}
(None, None) => return Err(usage.to_string()),
(Some(_), None) => None,
(None, Some(spec)) => Some(
areev_llm::resolve_toolcall(
spec,
flag(flags, "base-url").as_deref(),
flag(flags, "key-env").as_deref(),
)
.map_err(|e| e.to_string())?,
),
};
let llm_max_tokens: u32 = flag(flags, "llm-max-tokens")
.map(|v| v.parse().map_err(|_| "--llm-max-tokens must be a number".to_string()))
.transpose()?
.unwrap_or(1024);
let h = Hash::from_hex(&evalset_hex).map_err(|e| e.to_string())?;
let grain = facade.with_store(|m| m.get(&h)).map_err(|e| e.to_string())?;
let payload: serde_json::Value = grain
.get_str("object")
.and_then(|o| serde_json::from_str(o).ok())
.ok_or("grain is not an evalset (no JSON object payload)")?;
let cases = payload
.get("cases")
.and_then(|c| c.as_array())
.cloned()
.ok_or("evalset payload carries no cases")?;
let run_id = format!(
"eval-{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis())
.unwrap_or(0)
);
if llm.is_some() {
for case in &cases {
let cname = case.get("name").and_then(|v| v.as_str()).unwrap_or("case");
let input = case.get("input").cloned().unwrap_or(serde_json::json!({}));
eval_model_messages(&input).map_err(|e| format!("case '{cname}': {e}"))?;
}
}
let (mut passed, mut failed) = (0u64, 0u64);
let (mut effects, mut wall_ms) = (0u64, 0u64);
let (mut input_tokens, mut output_tokens) = (0u64, 0u64);
let mut metric_totals: std::collections::BTreeMap<String, (f64, u64)> =
std::collections::BTreeMap::new();
let mut usd_micros = 0u64;
let mut reported_usage = false;
let mut rows = Vec::new();
for case in &cases {
let cname = case.get("name").and_then(|v| v.as_str()).unwrap_or("case");
let input = case.get("input").cloned().unwrap_or(serde_json::json!({}));
let expect = case.get("expect").cloned().unwrap_or(serde_json::json!({}));
let started = std::time::Instant::now();
let mut case_report: Option<Result<EvalCaseReport, String>> = None;
let (ok, got) = match &llm {
Some(llm) => {
let (ok, got, usage) =
run_eval_case_model(llm.as_ref(), &input, &expect, llm_max_tokens);
if let Some((i, o)) = usage {
input_tokens += i;
output_tokens += o;
}
(ok, got)
}
None => {
let (ok, got, report) = run_eval_case(
cmd.as_deref().unwrap_or_default(),
&evalset_hex,
cname,
&input,
&expect,
);
case_report = report;
(ok, got)
}
};
let mut ok = ok;
let mut got = got;
if let Some(Err(why)) = &case_report {
ok = false;
got = format!("{got}\n[eval report rejected: {why}]");
}
effects += 1;
wall_ms += started.elapsed().as_millis() as u64;
if ok { passed += 1 } else { failed += 1 };
let mut case_metrics = serde_json::Map::new();
let mut case_usage = serde_json::Map::new();
if let Some(Ok(report)) = &case_report {
for (name, v) in &report.metrics {
let e = metric_totals.entry(name.clone()).or_insert((0.0, 0));
e.0 += v;
e.1 += 1;
case_metrics.insert(name.clone(), serde_json::json!(v));
}
if let Some(n) = report.input_tokens {
input_tokens += n;
reported_usage = true;
case_usage.insert("input_tokens".into(), serde_json::json!(n));
}
if let Some(n) = report.output_tokens {
output_tokens += n;
reported_usage = true;
case_usage.insert("output_tokens".into(), serde_json::json!(n));
}
if let Some(n) = report.usd_micros {
usd_micros += n;
reported_usage = true;
case_usage.insert("usd_micros".into(), serde_json::json!(n));
}
}
facade
.record_tool_call(
&case_ns,
&format!("eval:{cname}"),
Some(&input.to_string()),
&got,
!ok,
None,
None,
Some(&run_id),
None,
None,
Some(if ok { "completed" } else { "failed" }),
None,
None,
Some("host"),
None,
)
.map_err(|e| e.to_string())?;
let mut row = serde_json::json!({"case": cname, "ok": ok});
if !case_metrics.is_empty() {
row["metrics"] = serde_json::Value::Object(case_metrics);
}
if !case_usage.is_empty() {
row["usage"] = serde_json::Value::Object(case_usage);
}
rows.push(row);
}
let mut summary = serde_json::json!({
"run_id": run_id, "passed": passed, "failed": failed,
"effects": effects, "wall_ms": wall_ms,
});
if let Some(spec) = &model {
summary["model"] = serde_json::json!(spec);
summary["input_tokens"] = serde_json::json!(input_tokens);
summary["output_tokens"] = serde_json::json!(output_tokens);
}
if case_ns != EVAL_NS {
summary["case_ns"] = serde_json::json!(case_ns);
}
for (name, (sum, n)) in &metric_totals {
summary[name.as_str()] = serde_json::json!(sum / *n as f64);
summary[format!("{name}_n")] = serde_json::json!(n);
}
if reported_usage {
summary["input_tokens"] = serde_json::json!(input_tokens);
summary["output_tokens"] = serde_json::json!(output_tokens);
if usd_micros > 0 {
summary["usd_micros"] = serde_json::json!(usd_micros);
}
}
let mut fact = areev_core::types::Fact::new(
&format!("evalset:{evalset_hex}"),
"mg:eval_run",
&summary.to_string(),
)
.namespace(EVAL_NS);
fact.common.extra_fields.insert("run_id".into(), serde_json::json!(run_id));
facade.with_store(|m| m.add(&fact)).map_err(|e| e.to_string())?;
let mut reacceptance = serde_json::Value::Null;
if let Some(baseline_run) = flag(flags, "baseline") {
let tolerance: f64 = match flag(flags, "tolerance") {
None => 0.0,
Some(v) => match v.parse::<f64>() {
Ok(t) if t.is_finite() && t >= 0.0 => t,
_ => return Err("--tolerance must be a number of percentage points ≥ 0".into()),
},
};
let baseline = json_from_facade_recall(&facade, EVAL_NS, &format!("evalset:{evalset_hex}"))?
.into_iter()
.filter(|row| row["fields"]["relation"] == "mg:eval_run")
.filter_map(|row| {
serde_json::from_str::<serde_json::Value>(row["fields"]["object"].as_str()?)
.ok()
})
.find(|s| s["run_id"] == baseline_run.as_str())
.ok_or_else(|| {
format!("no recorded gate run '{baseline_run}' for this evalset")
})?;
let rate = |p: u64, f: u64| {
if p + f == 0 { 0.0 } else { p as f64 / (p + f) as f64 * 100.0 }
};
let base_rate = rate(
baseline["passed"].as_u64().unwrap_or(0),
baseline["failed"].as_u64().unwrap_or(0),
);
let this_rate = rate(passed, failed);
let within =
!areev_loop::recommendation::is_regression(base_rate, this_rate, true, tolerance);
reacceptance = serde_json::json!({
"baseline_run": baseline_run,
"baseline_pass_rate": base_rate,
"pass_rate": this_rate,
"tolerance_points": tolerance,
"accepted": within,
});
let mut cmp = areev_core::types::Fact::new(
&format!("evalset:{evalset_hex}"),
"mg:reacceptance",
&reacceptance.to_string(),
)
.namespace(EVAL_NS);
cmp.common.extra_fields.insert("run_id".into(), serde_json::json!(run_id));
facade.with_store(|m| m.add(&cmp)).map_err(|e| e.to_string())?;
if !within {
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"evalset": evalset_hex, "run_id": run_id,
"passed": passed, "failed": failed,
"reacceptance": reacceptance,
}))
.unwrap()
);
return Err(format!(
"re-acceptance FAILED: pass rate {this_rate:.1}% vs baseline \
{base_rate:.1}% (tolerance {tolerance} points)"
));
}
}
println!(
"{}",
serde_json::to_string_pretty(&{
let mut printed = summary.clone();
printed["evalset"] = serde_json::json!(evalset_hex);
printed["cases"] = serde_json::Value::Array(rows.clone());
printed["reacceptance"] = reacceptance.clone();
if !reported_usage && model.is_none() {
printed["input_tokens"] = serde_json::Value::Null;
printed["output_tokens"] = serde_json::Value::Null;
}
printed
})
.unwrap()
);
if failed > 0 {
return Err(format!("{failed} case(s) failed"));
}
Ok(())
}
other => Err(format!(
"unknown eval subcommand '{other}' — usage: areev eval <create|run>"
)),
}
}
fn json_from_facade_recall(
facade: &AreevFacade,
ns: &str,
subject: &str,
) -> Result<Vec<serde_json::Value>, String> {
let rows = facade
.with_store(|m| m.recall(ns, subject, None, 100_000))
.map_err(|e| e.to_string())?;
Ok(rows
.into_iter()
.map(|g| serde_json::json!({"fields": g.fields, "hash": g.hash.to_hex()}))
.collect())
}
#[derive(Debug, Default, Clone)]
pub(crate) struct EvalCaseReport {
pub(crate) metrics: std::collections::BTreeMap<String, f64>,
pub(crate) input_tokens: Option<u64>,
pub(crate) output_tokens: Option<u64>,
pub(crate) usd_micros: Option<u64>,
}
pub(crate) const RESERVED_EVAL_KEYS: &[&str] = &[
"run_id",
"passed",
"failed",
"effects",
"wall_ms",
"input_tokens",
"output_tokens",
"usd_micros",
"model",
"case_ns",
];
pub(crate) fn parse_eval_report(text: &str) -> Result<EvalCaseReport, String> {
let v: serde_json::Value = serde_json::from_str(text.trim())
.map_err(|e| format!("$AREEV_EVAL_REPORT is not valid JSON: {e}"))?;
let obj = v
.as_object()
.ok_or_else(|| "$AREEV_EVAL_REPORT must be a JSON object".to_string())?;
let mut out = EvalCaseReport::default();
if let Some(metrics) = obj.get("metrics") {
let metrics = metrics
.as_object()
.ok_or_else(|| "\"metrics\" must be an object".to_string())?;
for (name, value) in metrics {
if !name
.chars()
.next()
.is_some_and(|c| c.is_ascii_lowercase())
|| !name
.chars()
.all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '_')
{
return Err(format!(
"metric name {name:?} must match [a-z][a-z0-9_]*"
));
}
if RESERVED_EVAL_KEYS.contains(&name.as_str()) {
return Err(format!(
"metric name {name:?} collides with a summary key this verb owns \
(reserved: {})",
RESERVED_EVAL_KEYS.join(", ")
));
}
let n = value
.as_f64()
.filter(|n| n.is_finite())
.ok_or_else(|| format!("metric {name:?} must be a finite number"))?;
out.metrics.insert(name.clone(), n);
}
}
if let Some(usage) = obj.get("usage") {
let usage = usage
.as_object()
.ok_or_else(|| "\"usage\" must be an object".to_string())?;
for (key, slot) in [
("input_tokens", &mut out.input_tokens),
("output_tokens", &mut out.output_tokens),
("usd_micros", &mut out.usd_micros),
] {
if let Some(v) = usage.get(key) {
*slot = Some(
v.as_u64()
.ok_or_else(|| format!("usage.{key} must be a non-negative integer"))?,
);
}
}
}
Ok(out)
}
#[allow(clippy::type_complexity)]
fn run_eval_case(
cmd: &str,
evalset: &str,
case: &str,
input: &serde_json::Value,
expect: &serde_json::Value,
) -> (bool, String, Option<Result<EvalCaseReport, String>>) {
use areev_core::proc::{self, SpawnPolicy};
use std::process::Command;
#[cfg(not(windows))]
let mut shell = Command::new("/bin/sh");
#[cfg(not(windows))]
shell.arg("-c").arg(cmd);
#[cfg(windows)]
let mut shell = Command::new("cmd");
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
shell.raw_arg("/C").raw_arg(cmd);
}
let report_dir = std::env::temp_dir();
let report_path = report_dir.join(format!(
"areev-eval-{}-{}.json",
std::process::id(),
case.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
.collect::<String>()
));
let _ = std::fs::remove_file(&report_path);
let report_env = report_path.to_string_lossy().to_string();
let out = match proc::run(
shell,
Some(input.to_string().as_bytes()),
&[
("AREEV_EVALSET", evalset),
("AREEV_EVAL_CASE", case),
("AREEV_EVAL_REPORT", report_env.as_str()),
],
&SpawnPolicy::default(),
) {
Ok(o) => o,
Err(e) => return (false, format!("spawn failed: {e}"), None),
};
let report = match std::fs::read_to_string(&report_path) {
Ok(text) if !text.trim().is_empty() => Some(parse_eval_report(&text)),
_ => None,
};
let _ = std::fs::remove_file(&report_path);
let stdout = String::from_utf8_lossy(&out.stdout).trim().to_string();
if let Some(why) = out.failure("eval case") {
return (false, why, report);
}
(eval_expect_matches(expect, &stdout), stdout, report)
}
fn eval_expect_matches(expect: &serde_json::Value, got: &str) -> bool {
if let Some(eq) = expect.get("equals") {
serde_json::from_str::<serde_json::Value>(got)
.map(|parsed| &parsed == eq)
.unwrap_or(false)
} else if let Some(needle) = expect.get("contains").and_then(|v| v.as_str()) {
got.contains(needle)
} else {
false
}
}
fn eval_model_messages(
input: &serde_json::Value,
) -> Result<(Option<String>, Vec<areev_llm::ChatMessage>), String> {
use areev_llm::ChatMessage;
if let Some(prompt) = input.as_str() {
return Ok((None, vec![ChatMessage::User(prompt.to_string())]));
}
let Some(turns) = input.as_array() else {
return Err(
"--model cases need a string prompt or an array of {role, content} \
messages (tool-cmd cases take any JSON)"
.into(),
);
};
let mut system = None;
let mut messages = Vec::new();
for (i, turn) in turns.iter().enumerate() {
let role = turn.get("role").and_then(|r| r.as_str()).unwrap_or_default();
let content = turn
.get("content")
.and_then(|c| c.as_str())
.ok_or_else(|| format!("message {i} has no string content"))?;
match role {
"system" if i == 0 => system = Some(content.to_string()),
"system" => return Err(format!("message {i}: system must be the first message")),
"user" => messages.push(ChatMessage::User(content.to_string())),
"assistant" => messages.push(ChatMessage::Assistant {
text: Some(content.to_string()),
tool_calls: vec![],
provider_content: None,
}),
other => return Err(format!("message {i} has unsupported role {other:?}")),
}
}
if messages.is_empty() {
return Err("a --model case needs at least one user or assistant message".into());
}
Ok((system, messages))
}
fn run_eval_case_model(
llm: &dyn areev_llm::ToolCallLlm,
input: &serde_json::Value,
expect: &serde_json::Value,
max_tokens: u32,
) -> (bool, String, Option<(u64, u64)>) {
use areev_llm::{ToolCallRequest, ToolChoice};
let (system, messages) = match eval_model_messages(input) {
Ok(v) => v,
Err(e) => return (false, e, None),
};
let req = ToolCallRequest {
system,
messages,
tools: &[],
tool_choice: ToolChoice::None,
max_tokens,
temperature: 0.0,
};
match llm.call(&req) {
Ok(resp) => {
let usage = (resp.usage.input_tokens, resp.usage.output_tokens);
let text = resp.text.unwrap_or_default().trim().to_string();
(eval_expect_matches(expect, &text), text, Some(usage))
}
Err(e) => (false, format!("model call failed: {}", e.message), None),
}
}
fn audit_hold(
m: &mut Areev,
verb: &str,
ns: &str,
because: &str,
by: &str,
now: i64,
) -> Result<(), String> {
let mut obs = areev_core::authz::audit_observation(
by,
verb,
&format!("hold ns:{ns}"),
Some(because),
0,
now,
);
m.append_audit(&mut obs).map(|_| ()).map_err(|e| e.to_string())
}
fn run_hold(
m: &mut Areev,
ns: &str,
flags: &HashMap<String, String>,
positional: &[String],
) -> Result<(), String> {
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("list");
match sub_cmd {
"set" => {
let because = flag(flags, "because").ok_or_else(|| {
"usage: areev hold set --because \"litigation X\" [--ns NS] [--by PRINCIPAL]"
.to_string()
})?;
let by = flag(flags, "by").unwrap_or_else(|| "user:local".into());
let now = now_ms();
m.place_hold(ns, &because, &by, now).map_err(|e| e.to_string())?;
audit_hold(m, "hold.set", ns, &because, &by, now)?;
println!("hold placed on '{ns}' by {by}: {because}");
}
"release" => {
let because = flag(flags, "because").ok_or_else(|| {
"usage: areev hold release --because \"matter closed\" [--ns NS] \
[--by PRINCIPAL] — a release with no recorded rationale is the \
half of the trail an auditor actually asks for"
.to_string()
})?;
let by = flag(flags, "by").unwrap_or_else(|| "user:local".into());
let now = now_ms();
m.release_hold(ns).map_err(|e| e.to_string())?;
audit_hold(m, "hold.release", ns, &because, &by, now)?;
println!("hold released on '{ns}' by {by}: {because}");
}
"list" => {
let holds = m.hold_records().map_err(|e| e.to_string())?;
if flag(flags, "format").as_deref() == Some("json") {
let rows: Vec<serde_json::Value> = holds
.iter()
.map(|h| {
serde_json::json!({
"ns": h.ns,
"because": h.because,
"placed_by": h.placed_by,
"at_ms": h.at_ms,
})
})
.collect();
println!("{}", serde_json::to_string(&rows).map_err(|e| e.to_string())?);
return Ok(());
}
if holds.is_empty() {
println!("no legal holds in this memory");
}
for h in holds {
let when = if h.at_ms == 0 {
"placed at an unrecorded time".to_string()
} else {
format!("placed at {}", h.at_ms)
};
println!("{}: held by {} — {} ({when})", h.ns, h.placed_by, h.because);
}
}
other => {
return Err(format!(
"unknown hold subcommand '{other}' — usage: areev hold <set|release|list>"
))
}
}
Ok(())
}
fn read_stream_cursor(cursor_path: &str, to: &str) -> (String, i64, u32) {
let Ok(s) = std::fs::read_to_string(cursor_path) else {
return (String::new(), 0, 0);
};
let mut parts = s.trim().splitn(3, ' ');
let gen_id = parts.next().unwrap_or("").to_string();
let cursor: i64 = parts.next().and_then(|v| v.parse().ok()).unwrap_or(0);
let seg = match parts.next().and_then(|v| v.parse::<u32>().ok()) {
Some(n) => n,
None => std::fs::read_dir(format!("{to}/gen-{gen_id}"))
.map(|d| d.filter_map(|e| e.ok()).filter(|e| e.path().extension().is_some_and(|x| x == "mgb")).count() as u32)
.unwrap_or(0),
};
(gen_id, cursor, seg)
}
fn prune_generations(to: &str, current_gen: &str, window_ms: i64) -> Result<usize, String> {
let now = std::time::SystemTime::now();
let mut dropped = 0usize;
for entry in std::fs::read_dir(to).map_err(|e| e.to_string())?.flatten() {
let path = entry.path();
let name = entry.file_name().to_string_lossy().to_string();
if !path.is_dir() || !name.starts_with("gen-") || name == format!("gen-{current_gen}") {
continue;
}
let newest = std::fs::read_dir(&path)
.map_err(|e| e.to_string())?
.flatten()
.filter_map(|f| f.metadata().ok().and_then(|md| md.modified().ok()))
.max();
let Some(newest) = newest else { continue };
let age_ms = now.duration_since(newest).map(|d| d.as_millis() as i64).unwrap_or(0);
if age_ms > window_ms {
std::fs::remove_dir_all(&path).map_err(|e| e.to_string())?;
dropped += 1;
}
}
Ok(dropped)
}
fn parse_duration(s: &str) -> Option<i64> {
let s = s.trim();
let split = s.find(|c: char| !c.is_ascii_digit())?;
let n: i64 = s[..split].parse().ok()?;
let mult = match &s[split..] {
"s" => 1_000,
"m" => 60_000,
"h" => 3_600_000,
"d" => 86_400_000,
_ => return None,
};
Some(n * mult)
}
fn parse_severity(s: &str) -> Severity {
match s.to_ascii_lowercase().as_str() {
"high" => Severity::High,
"medium" => Severity::Medium,
"low" => Severity::Low,
_ => Severity::Info,
}
}
fn status_filter(flags: &HashMap<String, String>) -> Option<RecStatus> {
match flag(flags, "status").as_deref() {
Some("approved") => Some(RecStatus::Approved),
Some("rejected") => Some(RecStatus::Rejected),
Some("applied") => Some(RecStatus::Applied),
Some("rolled_back") => Some(RecStatus::RolledBack),
Some("expired") => Some(RecStatus::Expired),
Some("withdrawn") => Some(RecStatus::Withdrawn),
Some("all") => None,
_ => Some(RecStatus::Pending),
}
}
fn short(hash: &str) -> &str {
&hash[..hash.len().min(12)]
}
fn load_gating_evidence(
sub: &AreevSubstrate,
engine: &Engine,
rec_hash: &str,
run_id: &str,
) -> Result<areev_loop::GatingEvidence, String> {
engine.gating_evidence(sub, rec_hash, run_id).map_err(|e| e.to_string())
}
fn resolve_hash(engine: &Engine, sub: &AreevSubstrate, prefix: &str) -> Result<String, String> {
let recs = engine.recommendations(sub, None).map_err(|e| e.to_string())?;
Ok(areev_loop_adapter::find_recommendation(&recs, prefix)?.hash.clone())
}
fn resolve_llm(
flags: &HashMap<String, String>,
cmd_flag: &str,
model_flag: &str,
) -> Result<Option<Box<dyn areev_loop::LlmBackend>>, String> {
if let Some(cmd) = flag(flags, cmd_flag) {
let label = flag(flags, "llm-model");
let llm = areev_loop::CommandLlm::new(&cmd, label.as_deref()).map_err(|e| e.to_string())?;
return Ok(Some(Box::new(llm)));
}
if let Some(spec) = flag(flags, model_flag) {
let base = flag(flags, "llm-base-url");
let key_env = flag(flags, "llm-api-key-env");
let llm = areev_llm::resolve(&spec, base.as_deref(), key_env.as_deref())
.map_err(|e| e.to_string())?;
return Ok(Some(llm));
}
Ok(None)
}
fn run_remember(mut m: Areev, ns: &str, flags: &HashMap<String, String>) -> Result<(), String> {
let mut content = need(flags, "content")?;
if let Some(t) = m.ingress_transform_text(ns, &content).map_err(|e| e.to_string())? {
content = t;
}
let egress_wrap = matches!(
m.anon_active_mode(ns).map_err(|e| e.to_string())?.as_deref(),
Some("egress")
);
let observer = flag(flags, "observer").unwrap_or_else(|| "cli".to_string());
let dry_run = flag(flags, "dry-run").is_some();
let hint = flag(flags, "extract-hint");
let min_confidence = match flag(flags, "min-confidence") {
Some(s) => s
.parse::<f64>()
.map_err(|e| format!("bad --min-confidence: {e}"))?,
None => 0.0,
};
let role = flags.get("role").map(String::as_str);
if let Some(r) = role.filter(|r| areev_core::types::Role::from_str(r).is_none()) {
return Err(format!("--role {r:?}: expected user|assistant|system|tool"));
}
let explicit = match flag(flags, "facts") {
Some(j) => Some(areev_store::FactDraft::from_json_array(&j).map_err(|e| e.to_string())?),
None => None,
};
let llm = match explicit {
Some(_) => None,
None => resolve_llm(flags, "llm-cmd", "model")?,
};
let grounder = match llm {
Some(_) => resolve_llm(flags, "ground-cmd", "ground-model")?,
None => None,
};
let wrap = |b: Box<dyn areev_loop::LlmBackend>| -> Result<Box<dyn areev_loop::LlmBackend>, String> {
if egress_wrap {
let policy = areev_core::anon::AnonPolicy {
scope: "session".into(),
..Default::default()
};
Ok(Box::new(
areev_llm::PseudonymizingBackend::new(b, policy).map_err(|e| e.to_string())?,
))
} else {
Ok(b)
}
};
let llm = llm.map(wrap).transpose()?;
let grounder = grounder.map(wrap).transpose()?;
let model = llm.as_ref().map(|l| l.model().to_string());
if dry_run {
let (proposed, facts) = match &llm {
Some(l) => {
let (proposed, facts, _) = extract_drafts(
l.as_ref(),
grounder.as_deref(),
"dry-run",
&content,
hint.as_deref(),
min_confidence,
)?;
(proposed, facts)
}
None => {
let f = explicit.unwrap_or_default();
(f.len(), f)
}
};
println!(
"{}",
serde_json::json!({
"dry_run": true,
"model": model,
"proposed": proposed,
"facts": facts.iter().map(draft_json).collect::<Vec<_>>(),
})
);
return Ok(());
}
let capture = areev_store::Capture {
observer: Some(observer.as_str()),
session_id: flags.get("session-id").map(String::as_str),
role,
run_id: flags.get("run-id").map(String::as_str),
};
let event = m
.capture(ns, &content, &capture)
.map_err(|e| e.to_string())?;
let report = |facts: &[Hash], extra: serde_json::Value| {
let mut out = serde_json::json!({
"event": event.to_hex(),
"facts": facts.iter().map(|h| h.to_hex()).collect::<Vec<_>>(),
});
if let (Some(obj), Some(extra)) = (out.as_object_mut(), extra.as_object()) {
obj.extend(extra.clone());
}
println!("{out}");
};
let failed = |e: String| -> String {
report(&[], serde_json::json!({ "error": e }));
e
};
let (proposed, drafts, status) = match &llm {
None => {
let d = explicit.unwrap_or_default();
(d.len(), d, None)
}
Some(l) => match extract_drafts(
l.as_ref(),
grounder.as_deref(),
&event.to_hex(),
&content,
hint.as_deref(),
min_confidence,
) {
Ok((proposed, drafts, grounded)) => (
proposed,
drafts,
Some(if grounded { "verified" } else { "unverified" }),
),
Err(e) => return Err(failed(e)),
},
};
let attribution = areev_store::FactAttribution {
verification_status: status,
extractor_model: model.as_deref(),
};
let facts = m
.attach_facts(ns, &event, &drafts, &attribution)
.map_err(|e| e.to_string())?;
let extra = match &model {
Some(model) => serde_json::json!({
"model": model,
"verification_status": status,
"proposed": proposed,
"dropped": proposed.saturating_sub(facts.len()),
}),
None => serde_json::json!({}),
};
report(&facts, extra);
Ok(())
}
fn extract_drafts(
llm: &dyn areev_loop::LlmBackend,
grounder: Option<&dyn areev_loop::LlmBackend>,
source: &str,
content: &str,
hint: Option<&str>,
min_confidence: f64,
) -> Result<(usize, Vec<areev_store::FactDraft>, bool), String> {
let ex =
areev_llm::extract_pipeline(llm, grounder, source, content, hint, min_confidence)
.map_err(|e| e.to_string())?;
if ex.proposed >= areev_llm::extract::MAX_FACTS {
eprintln!(
"note: extraction hit the {}-fact cap; some facts may have been dropped",
areev_llm::extract::MAX_FACTS
);
}
let drafts = ex
.facts
.into_iter()
.map(|f| areev_store::FactDraft {
subject: f.subject,
relation: f.relation,
object: f.object,
confidence: f.confidence,
})
.collect();
Ok((ex.proposed, drafts, ex.grounded))
}
fn draft_json(d: &areev_store::FactDraft) -> serde_json::Value {
serde_json::json!({
"subject": d.subject,
"relation": d.relation,
"object": d.object,
"confidence": d.confidence,
})
}
fn run_provision(flags: &HashMap<String, String>) -> Result<(), String> {
let db = flag(flags, "db")
.or_else(|| std::env::var("AREEV_DB").ok().filter(|v| !v.trim().is_empty()))
.ok_or_else(|| {
"areev provision: --db <postgres DSN> names the memory to create".to_string()
})?;
if !is_pg_dsn(&db) {
return Err(format!(
"areev provision is postgres-only, and {db:?} is not a postgres DSN. A file-backed \
memory has nothing to provision: it creates and migrates itself at open, with no \
privilege system to fail against"
));
}
let tel_mode = match flag(flags, "telemetry") {
Some(v) => areev_store::TelemetryMode::parse(&v)
.ok_or_else(|| {
format!("--telemetry: unknown mode '{v}' (off|aggregate|aggregate-hashed|full)")
})?,
None => areev_store::TelemetryMode::Aggregate,
};
let db = match flag(flags, "meta-schema") {
Some(m) => with_meta_schema(&db, &m),
None => db,
};
if flags.contains_key("check") || flags.contains_key("dry-run") {
return check_provision_cmd(
&db,
flag(flags, "schema").as_deref(),
tel_mode,
flag(flags, "format").as_deref() == Some("json"),
);
}
provision_postgres(&db, flag(flags, "schema").as_deref(), tel_mode)
}
#[cfg(feature = "postgres")]
fn with_meta_schema(db: &str, m: &str) -> String {
let base = areev_store::pg::strip_meta_schema(db);
let sep = if base.contains('?') { '&' } else { '?' };
format!("{base}{sep}meta_schema={m}")
}
#[cfg(not(feature = "postgres"))]
fn with_meta_schema(db: &str, _m: &str) -> String {
db.to_string()
}
#[cfg(feature = "postgres")]
pub(crate) const EXIT_PENDING: i32 = 2;
#[cfg(feature = "postgres")]
fn check_provision_cmd(
db: &str,
schema_override: Option<&str>,
tel_mode: areev_store::TelemetryMode,
json_out: bool,
) -> Result<(), String> {
let dsn = areev_store::pg::strip_provision(db);
let (url, dsn_schema) = match schema_override {
Some(_) => (areev_store::pg::strip_schema(&dsn), String::new()),
None => areev_store::pg::split_schema_url(&dsn).map_err(|e| e.to_string())?,
};
let schema = schema_override.map(str::to_string).unwrap_or(dsn_schema);
if schema.is_empty() {
return Err("areev provision --check: name the schema with --schema NAME or ?schema=NAME".into());
}
let report = areev_store::pg::check_provision(&url, &schema, tel_mode)
.map_err(|e| e.to_string())?;
if json_out {
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
} else {
println!(
"schema {:?}{}: {}",
report.schema,
match &report.meta_schema {
Some(m) => format!(" + metadata schema {m:?}"),
None => String::new(),
},
if report.exists { "exists" } else { "ABSENT" }
);
for st in &report.stamps {
println!(
" {:14} found {:8} wanted {:8} {}",
st.name,
st.found.as_deref().unwrap_or("-"),
st.wanted,
if st.is_current() { "ok" } else { "PENDING" }
);
}
println!(" rolling_deploy: {}", report.rolling_deploy);
if !report.pending.is_empty() {
println!(" pending: {}", report.pending.join(", "));
}
}
if !report.is_current() {
std::process::exit(EXIT_PENDING);
}
Ok(())
}
#[cfg(not(feature = "postgres"))]
fn check_provision_cmd(
_db: &str,
_schema_override: Option<&str>,
_tel_mode: areev_store::TelemetryMode,
_json_out: bool,
) -> Result<(), String> {
Err("this build lacks the postgres backend — rebuild with --features postgres-tls".into())
}
#[cfg(feature = "postgres")]
fn provision_postgres(
db: &str,
schema_override: Option<&str>,
tel_mode: areev_store::TelemetryMode,
) -> Result<(), String> {
let dsn = areev_store::pg::strip_provision(db);
let (url, dsn_schema) = match schema_override {
Some(_) => (areev_store::pg::strip_schema(&dsn), String::new()),
None => areev_store::pg::split_schema_url(&dsn).map_err(|e| e.to_string())?,
};
let schema = schema_override.map(str::to_string).unwrap_or(dsn_schema);
if schema.is_empty() {
return Err("areev provision: name the schema with --schema NAME or ?schema=NAME".into());
}
let mut m = match tel_mode {
areev_store::TelemetryMode::Off => Areev::open_postgres(&url, &schema),
mode => Areev::open_postgres_with_telemetry(&url, &schema, mode),
}
.map_err(|e| e.to_string())?;
for w in m.open_warnings() {
eprintln!("areev: warning: {w}");
}
let _ = m.telemetry_flush();
drop(m);
println!(
"provisioned postgres schema {schema:?}{}{} — the next open runs no DDL",
match areev_store::pg::meta_schema(&url).map_err(|e| e.to_string())? {
Some(m) => format!(" + metadata schema {m:?}"),
None => String::new(),
},
match tel_mode {
areev_store::TelemetryMode::Off => " (telemetry tables NOT created)",
_ => " (including telemetry tables)",
}
);
Ok(())
}
#[cfg(not(feature = "postgres"))]
fn provision_postgres(
_db: &str,
_schema_override: Option<&str>,
_tel_mode: areev_store::TelemetryMode,
) -> Result<(), String> {
Err("this build lacks the postgres backend — rebuild with --features postgres-tls".into())
}
fn run_auth(flags: &HashMap<String, String>, positional: &[String]) -> Result<(), String> {
use areev_core::authz::{CredentialMap, TOKEN_PREFIX};
let sub = positional.first().map(String::as_str).unwrap_or("");
let path = flag(flags, "auth")
.ok_or_else(|| "areev auth: --auth <FILE> names the credential map".to_string())?;
let expires_at = match flag(flags, "expires") {
None => None,
Some(spec) => Some(parse_expiry(&spec)?),
};
match sub {
"mint" => {
let principal = flag(flags, "principal").ok_or_else(|| {
"areev auth mint: --principal <NAME> says who the token authenticates as".to_string()
})?;
let id = flag(flags, "id").ok_or_else(|| {
"areev auth mint: --id <NAME> names this credential so it can be revoked \
independently of the principal's other tokens"
.to_string()
})?;
let mut raw = [0u8; 32];
getrandom::fill(&mut raw)
.map_err(|e| format!("areev auth mint: no system randomness available: {e}"))?;
let token = format!("{TOKEN_PREFIX}{}", areev_core::authz::encode_token_body(&raw));
let digest = {
use sha2::Digest;
hex::encode(sha2::Sha256::digest(token.as_bytes()))
};
let mut map = read_map_or_empty(&path)?;
if map.tokens.iter().any(|t| t.id() == id) {
return Err(format!(
"areev auth mint: {path} already has a credential with id {id:?} — \
revoke it first, or choose another id"
));
}
map.tokens.push(areev_core::authz::CredentialEntry {
id: Some(id.clone()),
label: flag(flags, "label"),
sha256: Some(digest),
env: None,
principal: principal.clone(),
memories: flag(flags, "memories")
.map(|m| m.split(',').map(|s| s.trim().to_string()).collect()),
expires_at: expires_at.clone(),
});
let json = serde_json::to_string_pretty(&map)
.map_err(|e| format!("areev auth mint: {e}"))?;
CredentialMap::from_json(&json).map_err(|e| format!("areev auth mint: {e}"))?;
write_map(&path, &json)?;
println!("{token}");
eprintln!("areev: minted {id} for {principal} — recorded in {path}");
if let Some(exp) = &expires_at {
eprintln!("areev: expires {exp}");
}
eprintln!(
"areev: this is the ONLY time the token is shown; {path} stores its SHA-256, \
which cannot be reversed. Store it now."
);
}
"list" => {
let map = read_map_or_empty(&path)?;
let now = areev_core::time::now_ms();
if map.tokens.is_empty() {
eprintln!("areev: {path} holds no credentials");
return Ok(());
}
for t in &map.tokens {
let state = match t.expires_at.as_deref() {
None => "no expiry".to_string(),
Some(raw) if t.is_expired_at(now) => format!("EXPIRED {raw}"),
Some(raw) => format!("expires {raw}"),
};
let scope = match &t.memories {
Some(m) => format!(" memories={}", m.join(",")),
None => String::new(),
};
let label = t.label.as_deref().map(|l| format!(" — {l}")).unwrap_or_default();
println!("{}\t{}\t{state}{scope}{label}", t.id(), t.principal);
}
}
"revoke" => {
let id = flag(flags, "id")
.ok_or_else(|| "areev auth revoke: --id <NAME> (see `areev auth list`)".to_string())?;
let mut map = read_map_or_empty(&path)?;
let before = map.tokens.len();
map.tokens.retain(|t| t.id() != id);
if map.tokens.len() == before {
return Err(format!("areev auth revoke: {path} has no credential with id {id:?}"));
}
let json = serde_json::to_string_pretty(&map)
.map_err(|e| format!("areev auth revoke: {e}"))?;
write_map(&path, &json)?;
eprintln!(
"areev: revoked {id} from {path} — RESTART `areev ui` for it to take effect \
(the map is read once at startup)"
);
}
other => {
return Err(format!(
"unknown auth subcommand {other:?} — one of: mint, list, revoke\n\n\
areev auth mint --auth FILE --id NAME --principal NAME [--label TEXT]\n\
\x20 [--expires 90d|2026-12-31T23:59:59Z] [--memories A,B]\n\
areev auth list --auth FILE\n\
areev auth revoke --auth FILE --id NAME"
));
}
}
Ok(())
}
fn parse_expiry(spec: &str) -> Result<String, String> {
let spec = spec.trim();
if let Some(ms) = areev_core::time::iso8601_to_ms(spec) {
return Ok(format_epoch_ms(ms));
}
let (num, unit) = spec.split_at(spec.len().saturating_sub(1));
let n: i64 = num
.parse()
.map_err(|_| format!("--expires {spec:?}: expected <n>d, <n>h, or an ISO-8601 instant"))?;
let ms = match unit {
"d" => n * 86_400_000,
"h" => n * 3_600_000,
_ => {
return Err(format!(
"--expires {spec:?}: unit must be 'd' (days) or 'h' (hours), or give a full \
ISO-8601 instant"
))
}
};
if n <= 0 {
return Err(format!("--expires {spec:?}: must be in the future"));
}
Ok(format_epoch_ms(areev_core::time::now_ms() + ms))
}
fn format_epoch_ms(ms: i64) -> String {
let secs = ms.div_euclid(1000);
let days = secs.div_euclid(86_400);
let tod = secs.rem_euclid(86_400);
let z = days + 719_468;
let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
let doe = z - era * 146_097;
let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146_096) / 365;
let y = yoe + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let d = doy - (153 * mp + 2) / 5 + 1;
let m = if mp < 10 { mp + 3 } else { mp - 9 };
let y = if m <= 2 { y + 1 } else { y };
format!(
"{y:04}-{m:02}-{d:02}T{:02}:{:02}:{:02}Z",
tod / 3600,
(tod % 3600) / 60,
tod % 60
)
}
fn read_map_or_empty(path: &str) -> Result<areev_core::authz::CredentialMap, String> {
match std::fs::read_to_string(path) {
Ok(s) => areev_core::authz::CredentialMap::from_json(&s).map_err(|e| format!("{path}: {e}")),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
Ok(areev_core::authz::CredentialMap {
version: 1,
tokens: Vec::new(),
groups: None,
})
}
Err(e) => Err(format!("{path}: {e}")),
}
}
fn write_map(path: &str, json: &str) -> Result<(), String> {
std::fs::write(path, format!("{json}\n")).map_err(|e| format!("{path}: {e}"))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600));
}
Ok(())
}
fn load_policy(flags: &HashMap<String, String>) -> Result<Option<Policy>, String> {
let path = flag(flags, "policy").or_else(|| std::env::var("AREEV_LOOP_POLICY").ok());
match path {
Some(p) => {
let s = std::fs::read_to_string(&p).map_err(|e| format!("{p}: {e}"))?;
Ok(Some(Policy::from_json(&s).map_err(|e| e.to_string())?))
}
None => Ok(None),
}
}
fn run_audit(
m: &mut Areev,
flags: &HashMap<String, String>,
positional: &[String],
) -> Result<(), String> {
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("export");
if sub_cmd != "export" {
return Err(format!(
"unknown audit subcommand '{sub_cmd}' — usage: areev audit export \
[--since MS] [--until MS] [--out FILE]"
));
}
let parse_ms = |k: &str| -> Result<Option<i64>, String> {
match flag(flags, k) {
None => Ok(None),
Some(v) => v
.parse::<i64>()
.map(Some)
.map_err(|_| format!("--{k} takes epoch milliseconds, got '{v}'")),
}
};
let since = parse_ms("since")?;
let until = parse_ms("until")?;
let in_window = |at: i64| since.is_none_or(|s| at >= s) && until.is_none_or(|u| at <= u);
let mut rows: Vec<(i64, serde_json::Value)> = Vec::new();
let mut breaks = 0usize;
let mut tier2: Vec<(i64, serde_json::Value, serde_json::Value)> = Vec::new();
let mut seen_tier2: std::collections::HashSet<String> = std::collections::HashSet::new();
let authz = m
.recent(
areev_core::authz::AUTHZ_NS,
Some(areev_core::types::GrainType::Observation),
usize::MAX / 2,
)
.map_err(|e| e.to_string())?;
for g in authz {
let ctx = g.fields.get("context").cloned().unwrap_or(serde_json::Value::Null);
if ctx.get("audit").and_then(|v| v.as_str()) != Some("tier2") {
continue;
}
let at = g.fields.get("created_at").and_then(|v| v.as_i64()).unwrap_or(0);
if !in_window(at) {
continue;
}
let mut row = serde_json::json!({
"trail": "destruction",
"hash": g.hash.to_hex(),
"at_ms": at,
"principal": g.fields.get("observer_id"),
"observer_type": g.fields.get("observer_type"),
"verb": ctx.get("verb"),
"target": ctx.get("target"),
"subject_ref": ctx.get("subject_ref"),
"because": ctx.get("because"),
"grains_erased": ctx.get("grains_erased"),
"stale_corpora": ctx.get("stale_corpora"),
"stale_adapters": ctx.get("stale_adapters"),
});
if let Some(h) = ctx.get("hold_overridden") {
row["hold_overridden"] = h.clone();
}
match ctx.get("seq").and_then(|v| v.as_u64()) {
None => {
row["unchained"] = serde_json::json!(true);
}
Some(seq) => {
row["seq"] = serde_json::json!(seq);
if let Some(p) = ctx.get("previous_audit").and_then(|v| v.as_str()) {
row["previous_audit"] = serde_json::json!(p);
}
if ctx.get("chain_root").and_then(|v| v.as_bool()) == Some(true) {
row["chain_root"] = serde_json::json!(true);
}
}
}
seen_tier2.insert(g.hash.to_hex());
tier2.push((at, row, ctx.clone()));
}
let mut prev_seq: Option<u64> = None;
for (at, mut row, ctx) in tier2 {
if let Some(p) = ctx.get("previous_audit").and_then(|v| v.as_str()) {
if !seen_tier2.contains(p) {
let exists = areev_core::error::Hash::from_hex(p)
.ok()
.map(|h| m.has(&h).unwrap_or(false))
.unwrap_or(false);
if exists {
row["previous_outside_window"] = serde_json::json!(true);
} else {
breaks += 1;
row["chain_break"] = serde_json::json!(true);
}
}
}
if let (Some(prev), Some(seq)) = (prev_seq, row.get("seq").and_then(|v| v.as_u64())) {
if seq > prev + 1 {
breaks += 1;
row["chain_break"] = serde_json::json!(true);
row["missing_seq"] = serde_json::json!(seq - prev - 1);
}
}
if let Some(seq) = row.get("seq").and_then(|v| v.as_u64()) {
prev_seq = Some(seq);
}
rows.push((at, row));
}
let mut chain_prev: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut loop_rows: Vec<(i64, serde_json::Value, Option<String>)> = Vec::new();
let loop_grains = m
.recent("areev-loop", Some(areev_core::types::GrainType::Fact), usize::MAX / 2)
.map_err(|e| e.to_string())?;
for g in loop_grains {
if g.get_str("relation") != Some("loop_audit") {
continue;
}
let Some(body) = g.get_str("object").and_then(|o| {
serde_json::from_str::<serde_json::Value>(o).ok()
}) else {
continue;
};
let at = body.get("at_ms").and_then(|v| v.as_i64()).unwrap_or(0);
if !in_window(at) {
continue;
}
chain_prev.insert(g.hash.to_hex());
let prev = body
.get("derived_from")
.and_then(|v| v.as_array())
.and_then(|a| a.get(1))
.and_then(|v| v.as_str())
.map(|s| s.to_string());
loop_rows.push((
at,
serde_json::json!({
"trail": "loop",
"hash": g.hash.to_hex(),
"at_ms": at,
"rec_hash": body.get("rec_hash"),
"from_status": body.get("from_status"),
"to_status": body.get("to_status"),
"actor": body.get("actor"),
"observer_type": body.get("observer_type"),
"because": body.get("because"),
"previous_audit": prev.clone(),
"gating": body.get("gating_run_id").map(|run| serde_json::json!({
"evalset": body.get("gating_evalset"),
"run_id": run,
"passed": body.get("gating_passed"),
"failed": body.get("gating_failed"),
})),
"scope": body.get("scope"),
}),
prev,
));
}
for (at, mut row, prev) in loop_rows {
if let Some(p) = prev {
if !chain_prev.contains(&p) {
let exists = areev_core::error::Hash::from_hex(&p)
.ok()
.map(|h| m.has(&h).unwrap_or(false))
.unwrap_or(false);
if exists {
row["previous_outside_window"] = serde_json::json!(true);
} else {
breaks += 1;
row["chain_break"] = serde_json::json!(true);
}
}
}
rows.push((at, row));
}
if flags.contains_key("with-outcomes") {
let state = areev_loop_adapter::loop_state_of(m).map_err(|e| e.to_string())?;
let outcomes: Vec<areev_loop::OutcomeResult> = if state.is_null() {
Vec::new()
} else {
areev_loop::config::LoopPersisted::from_value(state)
.map(|p| p.outcomes.into_values().flatten().collect())
.unwrap_or_default()
};
{
for o in outcomes {
if !in_window(o.measured_at_ms) {
continue;
}
rows.push((
o.measured_at_ms,
serde_json::json!({
"trail": "loop_outcome",
"rec_hash": o.rec_hash,
"metric": o.metric,
"verdict": o.verdict,
"baseline": o.baseline,
"current": o.current,
"baseline_run_id": o.baseline_run_id,
"current_run_id": o.current_run_id,
"cost": o.cost,
"measured_at_ms": o.measured_at_ms,
"chained": false,
}),
));
}
}
}
rows.sort_by_key(|(at, _)| *at);
let mut out = String::new();
for (_, row) in &rows {
out.push_str(&row.to_string());
out.push('\n');
}
match flag(flags, "out") {
Some(path) => {
std::fs::write(&path, &out).map_err(|e| e.to_string())?;
println!("wrote {} audit records to {path}", rows.len());
}
None => print!("{out}"),
}
if breaks > 0 {
eprintln!(
"areev: warning: {breaks} loop audit record(s) name a predecessor absent from this \
export — the chain is truncated (widen --since, or the trail was tampered with)"
);
}
eprintln!("audit export: {} records", rows.len());
Ok(())
}
fn run_loop(
m: Areev,
ns: &str,
flags: &HashMap<String, String>,
positional: &[String],
) -> Result<(), String> {
let sub_cmd = positional.first().map(|s| s.as_str()).unwrap_or("status");
let egress = store_egress_active(&m);
let mut sub = AreevSubstrate::new(m, Some(ns.to_string()));
let policy = load_policy(flags)?;
let mut engine = match policy {
Some(p) => Engine::with_builtins().with_policy(p),
None => Engine::with_builtins(),
};
if let Some(cmd) = flag(flags, "llm-cmd") {
let model = flag(flags, "llm-model");
let llm = areev_loop::CommandLlm::new(&cmd, model.as_deref()).map_err(|e| e.to_string())?;
engine = engine.with_llm(Box::new(llm));
} else if let Some(spec) = flag(flags, "model") {
let base = flag(flags, "llm-base-url");
let key_env = flag(flags, "llm-api-key-env");
let llm = areev_llm::resolve(&spec, base.as_deref(), key_env.as_deref())
.map_err(|e| e.to_string())?;
engine = engine.with_llm(llm);
}
if let Some(cmd) = flag(flags, "ground-cmd") {
let g = areev_loop::CommandLlm::new(&cmd, None).map_err(|e| e.to_string())?;
engine = engine.with_ground_llm(Box::new(g));
} else if let Some(spec) = flag(flags, "ground-model") {
let base = flag(flags, "llm-base-url");
let key_env = flag(flags, "llm-api-key-env");
let g = areev_llm::resolve(&spec, base.as_deref(), key_env.as_deref())
.map_err(|e| e.to_string())?;
engine = engine.with_ground_llm(g);
}
if let Some(chain) = resolve_decider(flags)? {
let chain = decider_for_egress(chain, egress)?;
engine = engine.with_decider(Box::new(areev_loop_adapter::LoopDecider(chain)));
}
if let Some(cmd) = flag(flags, "analyzer-cmd") {
let a = areev_loop::CommandAnalyzer::new(&cmd).map_err(|e| e.to_string())?;
engine.register(Box::new(a));
}
let now = now_ms();
let actor = flag(flags, "actor").unwrap_or_else(|| "user:local".to_string());
let observer = ObserverType::Human;
let scopes = ScopeSet::all(); let json = flag(flags, "format").as_deref() == Some("json");
match sub_cmd {
"run" | "reflect" => {
let opts = RunOptions {
min_new: flag(flags, "min-new").and_then(|v| v.parse().ok()),
min_new_errors: flag(flags, "min-new-errors").and_then(|v| v.parse().ok()),
if_stale_ms: flag(flags, "if-stale").and_then(|v| parse_duration(&v)),
namespaces: Vec::new(),
full_sweep: sub_cmd == "reflect",
triggering_actor: Some(actor.clone()),
};
let res = engine.run(&mut sub, &opts, now).map_err(|e| e.to_string())?;
if json {
println!("{}", serde_json::to_string(&res).map_err(|e| e.to_string())?);
} else if res.ran() {
if !flags.contains_key("quiet") {
eprintln!(
"loop: ran — proposed {} ({} deduped, {} auto-applied) across {} analyzer(s)",
res.stored,
res.deduped,
res.auto_applied,
res.analyzers_run.len()
);
}
if res.stored > 0 {
eprintln!("loop: {} new — areev loop list", res.stored);
}
} else if !flags.contains_key("quiet") {
eprintln!("loop: skipped ({:?})", res.skip_reason);
}
}
"list" => {
let filter = status_filter(flags);
let recs = areev_loop_adapter::visible_recommendations(
&engine,
&sub,
&sub.authz(),
filter,
)
.map_err(|e| e.to_string())?;
if json {
let rows: Vec<_> = recs
.iter()
.map(|r| {
let mut row = serde_json::json!({
"hash": r.hash,
"status": r.status.as_str(),
"severity": r.severity.as_str(),
"analyzer": r.analyzer,
"destructive": r.destructive,
"summary": r.summary.render(),
});
if !r.near_duplicate_of.is_empty() {
row["near_duplicate_of"] = serde_json::json!(r.near_duplicate_of);
}
row
})
.collect();
println!("{}", serde_json::to_string(&rows).map_err(|e| e.to_string())?);
} else if recs.is_empty() {
eprintln!("no recommendations — run `areev loop run` first");
} else {
for r in &recs {
println!(
"{} {:<6} {:<28} {}",
short(&r.hash),
r.severity.as_str(),
r.analyzer,
r.summary.render()
);
}
}
if let Some(sev) = flag(flags, "fail-on") {
let threshold = parse_severity(&sev);
let hit = recs
.iter()
.any(|r| r.status == RecStatus::Pending && r.severity >= threshold);
if hit {
eprintln!("loop: pending recommendation(s) at or above severity '{sev}'");
std::process::exit(2);
}
}
}
"show" => {
let prefix = positional
.get(1)
.ok_or_else(|| "usage: areev loop show <hash>".to_string())?;
let recs = engine.recommendations(&sub, None).map_err(|e| e.to_string())?;
let r = areev_loop_adapter::find_recommendation(&recs, prefix)?;
let out = areev_loop_adapter::recommendation_detail(r);
println!("{}", serde_json::to_string_pretty(&out).map_err(|e| e.to_string())?);
}
"approve" | "reject" => {
let prefix = positional
.get(1)
.ok_or_else(|| format!("usage: areev loop {sub_cmd} <hash> --because \"...\""))?;
let because = need(flags, "because")?;
let hash = resolve_hash(&engine, &sub, prefix)?;
let decision = if sub_cmd == "approve" { Decision::Approve } else { Decision::Reject };
engine
.review(&mut sub, &hash, decision, &actor, observer, &scopes, &because, now)
.map_err(|e| e.to_string())?;
eprintln!("{sub_cmd}d {}", short(&hash));
}
"apply" => {
let prefix = positional
.get(1)
.ok_or_else(|| "usage: areev loop apply <hash> --because \"...\" [--gating-run RUN_ID]".to_string())?;
let because = need(flags, "because")?;
let hash = resolve_hash(&engine, &sub, prefix)?;
let allow_destructive = flags.contains_key("allow-destructive");
let applied = match flag(flags, "gating-run") {
Some(run_id) => {
let gating = load_gating_evidence(&sub, &engine, &hash, &run_id)?;
engine
.apply_gated(&mut sub, &hash, &actor, observer, &scopes, &because, allow_destructive, &gating, now)
.map_err(|e| e.to_string())?
}
None => engine
.apply(&mut sub, &hash, &actor, observer, &scopes, &because, allow_destructive, now)
.map_err(|e| e.to_string())?,
};
eprintln!(
"applied {} ({})",
short(&hash),
if applied.rollbackable { "rollbackable" } else { "non-rollbackable" }
);
}
"rollback" => {
let prefix = positional
.get(1)
.ok_or_else(|| "usage: areev loop rollback <hash> --because \"...\"".to_string())?;
let because = need(flags, "because")?;
let hash = resolve_hash(&engine, &sub, prefix)?;
let target_ref = engine
.recommendations(&sub, None)
.ok()
.and_then(|recs| recs.into_iter().find(|r| r.hash == hash))
.map(|r| r.target_ref);
if let Some(model) = target_ref.as_deref().and_then(|t| t.strip_prefix("model:")) {
eprintln!(
"adapter promotion for model:{model} will be retracted — the host \
must stop serving it; runs answered since promotion are NOT reverted"
);
}
let code_target = target_ref
.as_deref()
.and_then(|t| t.strip_prefix("tool:").map(str::to_string));
if let Some(tool_hash) = &code_target {
if let Ok(h) = Hash::from_hex(tool_hash) {
let touched = sub
.facade()
.with_store(|m| m.runs_touching("agent:harness", &h, 2))
.unwrap_or_default();
eprintln!(
"blast radius: {} run(s) touched tool {} since apply{}",
touched.len(),
short(tool_hash),
if touched.is_empty() { String::new() } else { format!(": {}", touched.join(", ")) }
);
eprintln!(
" the runs' outputs are NOT reverted — inspect each with \
`areev run-trace --run-id <id>`"
);
}
}
engine
.rollback(&mut sub, &hash, &actor, observer, &scopes, &because, now)
.map_err(|e| e.to_string())?;
eprintln!("rolled back {}", short(&hash));
}
"analyzers" => {
for a in engine.analyzers() {
let m = a.manifest();
println!(
"{:<28} {:?} {:<7} on={} {}",
m.id,
m.tier,
format!("{:?}", m.trust_class).to_lowercase(),
m.default_on,
m.title
);
}
}
"replay" => {
use areev_loop::replay::{ReplayOptions, ReplayRequest};
let path = flag(flags, "config")
.ok_or("usage: areev loop replay --config FILE [--window 90d | --since MS] [--step per-pass|1d] [--format json]")?;
let text = std::fs::read_to_string(&path).map_err(|e| format!("{path}: {e}"))?;
let req = ReplayRequest::from_json(&text).map_err(|e| e.to_string())?;
let (candidate, _) = req.resolve(now).map_err(|e| e.to_string())?;
let since: Option<i64> = flag(flags, "since")
.map(|v| v.parse().map_err(|_| "--since must be epoch milliseconds".to_string()))
.transpose()?;
let opts = ReplayOptions::from_args(
flag(flags, "window").as_deref(),
since,
flag(flags, "step").as_deref(),
Vec::new(),
now,
)
.map_err(|e| e.to_string())?;
let ops_before = sub.facade().with_store(|m| m.stats()).map_err(|e| e.to_string())?.ops;
let report = engine.replay(&sub, &candidate, &opts).map_err(|e| e.to_string())?;
let ops_after = sub.facade().with_store(|m| m.stats()).map_err(|e| e.to_string())?.ops;
if ops_after != ops_before {
return Err(format!(
"replay wrote to the memory (op-log {ops_before} → {ops_after}) — this is a bug; the report is discarded"
));
}
if json {
let mut v = serde_json::to_value(&report).map_err(|e| e.to_string())?;
v["oplog_len"] = serde_json::json!({"before": ops_before, "after": ops_after});
println!("{}", serde_json::to_string(&v).map_err(|e| e.to_string())?);
} else {
println!(
"replay: {} step(s) ({}) from {} to {}; op-log {ops_before} → {ops_after} (unchanged)",
report.steps.len(),
report.step,
report.since_ms,
report.until_ms
);
let row = |name: &str, t: &areev_loop::replay::ReplayTally| {
format!(
"{name:<28} {:>8} {:>9} {:>9} {:>9} {:>9} {:>10} {:>8} {:>6}",
t.findings, t.approved, t.rejected, t.never_reviewed, t.never_proposed, t.regressed, t.drifted, t.held
)
};
println!(
"{:<28} {:>8} {:>9} {:>9} {:>9} {:>9} {:>10} {:>8} {:>6}",
"", "findings", "approved", "rejected", "unrevwd", "never", "regressed", "drifted", "held"
);
println!("{}", row("incumbent", &report.incumbent.total));
println!("{}", row("candidate", &report.candidate.total));
let analyzers: std::collections::BTreeSet<&String> = report
.incumbent
.per_analyzer
.keys()
.chain(report.candidate.per_analyzer.keys())
.collect();
for a in analyzers {
let empty = areev_loop::replay::ReplayTally::default();
let i = report.incumbent.per_analyzer.get(a).unwrap_or(&empty);
let c = report.candidate.per_analyzer.get(a).unwrap_or(&empty);
println!("{}", row(&format!(" {a} (incumbent)"), i));
println!("{}", row(&format!(" {a} (candidate)"), c));
}
println!(
"queue per step — incumbent {:?}, candidate {:?}",
report.incumbent.queue_per_step, report.candidate.queue_per_step
);
for n in &report.not_replayed {
println!("not replayed: {} — {}", n.what, n.reason);
}
println!("{}", report.matching);
}
}
"outcomes" => {
let outcomes = engine.outcomes(&sub).map_err(|e| e.to_string())?;
if json {
println!("{}", serde_json::to_string(&outcomes).map_err(|e| e.to_string())?);
} else if outcomes.is_empty() {
eprintln!(
"no measured outcomes yet — outcome review runs after an applied \
recommendation's review window elapses"
);
} else {
for o in &outcomes {
let horizon = match o.checkpoint {
Some(cp) => cp.label(),
None if o.horizon_ms % 86_400_000 == 0 => format!("{}d", o.horizon_ms / 86_400_000),
None => format!("{}h", o.horizon_ms / 3_600_000),
};
let mut line = format!(
"{} {:<22} @{:<4} baseline {} → current {} [{}]",
short(&o.rec_hash),
o.metric,
horizon,
o.baseline,
o.current,
o.verdict
);
if let Some(run) = &o.baseline_run_id {
line.push_str(&format!(" baseline={} ({run})", o.baseline_kind));
}
if let Some(best) = o.best_before {
line.push_str(&format!(" best_before {best}"));
}
if let Some(c) = &o.cost {
match (c.baseline, c.current) {
(Some(b), Some(cur)) => {
let ratio = if b > 0.0 { format!(" ×{:.2}", cur / b) } else { String::new() };
line.push_str(&format!(
" cost {} {b} → {cur}{ratio} [{}, bound ×{}]",
c.field, c.status, c.max_increase_ratio
));
}
_ => line.push_str(&format!(" cost {} not measurable", c.field)),
}
}
println!("{line}");
}
}
}
"policy" => {
println!(
"{}",
serde_json::to_string_pretty(engine.policy()).map_err(|e| e.to_string())?
);
}
_ => {
let h = engine.health(&sub, now).map_err(|e| e.to_string())?;
if json {
println!("{}", serde_json::to_string(&h).map_err(|e| e.to_string())?);
} else {
println!(
"loop: {} recommendation(s) — {} pending, {} applied",
h.total, h.pending, h.applied
);
match h.last_run_ms {
None => println!(" never run — areev loop run"),
Some(last) => {
let days = (now - last) / 86_400_000;
println!(
" last run {days}d ago; {} new grain(s) ({} tool error(s)) since",
h.grains_since_run, h.error_events_since_run
);
}
}
let lm = engine.llm_metrics(&sub).map_err(|e| e.to_string())?;
if lm.proposed > 0 {
match lm.approval_rate {
Some(rate) => println!(
" LLM findings: {} surfaced, {:.0}% approved ({} approved / {} rejected, {} pending)",
lm.proposed, rate * 100.0, lm.approved, lm.rejected, lm.pending
),
None => println!(" LLM findings: {} surfaced, none decided yet", lm.proposed),
}
}
if h.stale {
eprintln!(" ⚠ the loop may be stale — run `areev loop run` or wire the SessionEnd hook");
} else if h.pending > 0 {
println!(" review with: areev loop list");
}
}
}
}
Ok(())
}
fn run_init(
m: Areev,
ns: &str,
flags: &HashMap<String, String>,
_positional: &[String],
) -> Result<(), String> {
let template = flag(flags, "template").unwrap_or_else(|| "blank".to_string());
let mut m = m;
let seeded = match template.as_str() {
"blank" => 0,
"demo" => seed_demo(&mut m, ns)?,
"coding-agent" | "support-agent" => {
m.add(&Fact::new("agent", "instruction", "review pending recommendations before acting").namespace(ns))
.map_err(|e| e.to_string())?;
1
}
other => return Err(format!("unknown --template '{other}' (blank|demo|coding-agent|support-agent)")),
};
println!("initialized backend (template: {template}, seeded {seeded} grain(s))");
if template == "demo" {
println!("next: areev loop run --db <file> # ~4 recommendations across analyzers");
}
let exe = std::env::current_exe()
.ok()
.and_then(|p| p.to_str().map(str::to_string))
.unwrap_or_else(|| "areev".to_string());
eprintln!("\nClaude Code hooks (paste into settings.json — absolute path baked in):");
eprintln!(" UserPromptSubmit → {exe} recall-hook --ns {ns} --with-loop");
eprintln!(" Stop → {exe} capture-stop --ns {ns}");
eprintln!(" PreCompact → {exe} capture-stop --ns {ns}");
eprintln!(" SessionEnd → {exe} loop run --min-new 20 --min-new-errors 3 --quiet --ns {ns}");
Ok(())
}
fn seed_demo(m: &mut Areev, ns: &str) -> Result<usize, String> {
let past = 1_000_000_000_000; let grains = [
Fact::new("acme", "tier", "Enterprise").namespace(ns),
Fact::new("acme", "tier", "Enterprise").namespace(ns),
Fact::new("acme", "deploy_target", "us-east-1").namespace(ns),
Fact::new("acme", "deploy_target", "eu-west-1").namespace(ns),
Fact::new("promo", "active", "true").namespace(ns).valid_to(past),
];
for g in &grains {
m.add(g).map_err(|e| e.to_string())?;
}
Ok(grains.len())
}
const CLI_STACK_BYTES: usize = 16 * 1024 * 1024;
fn main() -> ExitCode {
let worker = std::thread::Builder::new()
.stack_size(CLI_STACK_BYTES)
.spawn(run);
let result = match worker {
Ok(h) => h.join(),
Err(_) => Ok(run()),
};
match result {
Ok(Ok(())) => ExitCode::SUCCESS,
Ok(Err(e)) => {
eprintln!("areev: {e}");
ExitCode::FAILURE
}
Err(_) => ExitCode::FAILURE,
}
}
#[cfg(test)]
mod tests {
use super::*;
fn args(a: &[&str]) -> (HashMap<String, String>, Vec<String>) {
parse_args(&a.iter().map(|s| s.to_string()).collect::<Vec<_>>())
}
#[test]
fn mount_spec_splits_on_the_alias_not_on_every_comma() {
assert_eq!(
parse_mounts(Some("org=/srv/org.db,kb=/srv/kb.db")).unwrap(),
vec![
("org".to_string(), "/srv/org.db".to_string()),
("kb".to_string(), "/srv/kb.db".to_string()),
]
);
assert_eq!(
parse_mounts(Some("org=postgres://u:p@h1:5432,h2:5432/db?schema=org_kb")).unwrap(),
vec![(
"org".to_string(),
"postgres://u:p@h1:5432,h2:5432/db?schema=org_kb".to_string()
)]
);
assert_eq!(
parse_mounts(Some(
"org=postgres://u:p@h1:5432,h2:5432/db?schema=org_kb, kb=/srv/kb.db"
))
.unwrap(),
vec![
(
"org".to_string(),
"postgres://u:p@h1:5432,h2:5432/db?schema=org_kb".to_string()
),
("kb".to_string(), "/srv/kb.db".to_string()),
]
);
assert_eq!(
parse_mounts(Some(
"org=postgres://h/db?schema=s&options=-c%20a=1,-c%20b=2"
))
.unwrap()
.len(),
1
);
assert_eq!(parse_mounts(None).unwrap(), Vec::new());
assert!(parse_mounts(Some("no-equals-sign")).is_err());
}
#[test]
fn mount_errors_and_logs_never_print_a_dsn_password() {
let dsn = "postgres://areev:SUPERSECRET@pg:5432/db?schema=org_kb";
assert!(!redact_dsn(dsn).contains("SUPERSECRET"), "{}", redact_dsn(dsn));
assert!(redact_dsn(dsn).contains("areev:***@pg:5432"));
let err = parse_mounts(Some(dsn)).unwrap_err();
assert!(!err.contains("SUPERSECRET"), "{err}");
assert_eq!(redact_dsn("/srv/org.db"), "/srv/org.db");
}
#[cfg(feature = "postgres")]
#[test]
fn a_postgres_mount_of_the_primary_schema_is_recognized() {
let primary = "postgres://u:p@h/db?schema=main";
assert!(mount_is_the_primary(primary, "postgres://u:p@h/db?schema=main"));
assert!(!mount_is_the_primary(primary, "postgres://u:p@h/db?schema=org_kb"));
assert!(!mount_is_the_primary("/srv/main.db", "/srv/main.db"));
}
#[test]
fn expiry_formatting_round_trips_through_the_parser() {
for iso in [
"1970-01-01T00:00:00Z",
"2026-08-27T12:34:56Z",
"2024-02-28T23:59:59Z",
"2024-02-29T12:00:00Z",
"2024-03-01T00:00:00Z",
"2000-02-29T00:00:00Z",
"2100-02-28T00:00:00Z",
"2026-12-31T23:59:59Z",
"2027-01-01T00:00:00Z",
] {
let ms = areev_core::time::iso8601_to_ms(iso).expect(iso);
assert_eq!(format_epoch_ms(ms), iso, "round-trip failed for {iso}");
}
}
#[test]
fn expiry_parsing_accepts_windows_and_instants_and_refuses_the_rest() {
assert_eq!(
parse_expiry("2026-12-31T23:59:59Z").unwrap(),
"2026-12-31T23:59:59Z"
);
assert_eq!(parse_expiry("2026-12-31").unwrap(), "2026-12-31T00:00:00Z");
let now = areev_core::time::now_ms();
let d90 = areev_core::time::iso8601_to_ms(&parse_expiry("90d").unwrap()).unwrap();
let h12 = areev_core::time::iso8601_to_ms(&parse_expiry("12h").unwrap()).unwrap();
assert!(h12 > now && d90 > h12, "90d={d90} 12h={h12} now={now}");
assert!((d90 - now - 90 * 86_400_000).abs() < 2_000);
for bad in ["90x", "d", "", "-5d", "0d", "soon", "90 d"] {
assert!(parse_expiry(bad).is_err(), "{bad:?} must be refused");
}
}
#[test]
fn valueless_long_flag_does_not_swallow_short_flag() {
let (flags, pos) = args(&["--mcp", "-d", "mem.db", "--ns", "caller"]);
assert_eq!(flags.get("mcp").map(String::as_str), Some("true"));
assert_eq!(flags.get("db").map(String::as_str), Some("mem.db"));
assert_eq!(flags.get("ns").map(String::as_str), Some("caller"));
assert!(pos.is_empty(), "no stray positionals: {pos:?}");
}
#[test]
fn db_flag_equivalent_across_form_and_order() {
for a in [
&["--mcp", "-d", "mem.db"][..],
&["-d", "mem.db", "--mcp"][..],
&["--mcp", "--db", "mem.db"][..],
&["--db", "mem.db", "--mcp"][..],
] {
let (flags, _) = args(a);
assert_eq!(flags.get("db").map(String::as_str), Some("mem.db"), "args: {a:?}");
assert_eq!(flags.get("mcp").map(String::as_str), Some("true"), "args: {a:?}");
}
}
#[test]
fn adjacent_valueless_flags_stay_true() {
let (flags, _) = args(&["--mcp", "--no-destructive-ops"]);
assert_eq!(flags.get("mcp").map(String::as_str), Some("true"));
assert_eq!(flags.get("no-destructive-ops").map(String::as_str), Some("true"));
}
#[test]
fn openai_tool_import_preserves_arguments_error_and_call_id() {
let rows = extract_tool_records(&serde_json::json!({
"created_at": 42,
"tool_calls": [{
"id": "call-7",
"is_error": true,
"function": {
"name": "charge",
"arguments": "{\"amount\":42}"
}
}]
}));
assert_eq!(rows.len(), 1);
assert_eq!(rows[0].name, "charge");
assert_eq!(rows[0].input, Some(serde_json::json!({"amount": 42})));
assert!(rows[0].content.is_empty(), "arguments are not a result");
assert!(rows[0].is_error);
assert_eq!(rows[0].call_id.as_deref(), Some("call-7"));
assert_eq!(rows[0].created_at, Some(42));
}
}