dd952c0e46
The soul obeys half of its own ownership rule. soul.el:571-573 says "when
ENGRAM_URL is set the HTTP Engram owns persistence — the soul must NEVER write
to the local snapshot", and it doesn't. But nothing was ever built to hand the
soul's writes TO that owner: sync is pull-only (/api/sync -> engram_load_merge),
so every node created inside the soul lived in process RAM and was shed on
restart. Measured live 2026-08-07: soul node_count=102184, engram 79197.
SCOPE CORRECTION vs the earlier internal spec: engram provisional claim 17's
"pull-then-push" is a PEER-ENGRAM to PEER-ENGRAM protocol (claims 15-18 say so
explicitly). The soul is a CALLER of the database API, not a peer. Claim 17 is
NOT authority for a soul<->engram contract and is no longer cited as such. The
design here follows from the ownership rule alone.
Mechanism: a new Accessor, persist.el, is the single boundary. Writes stage a
delta to a filesystem spool and are pushed to the owner via POST /api/load-merge
— NOT POST /api/nodes, which mints a new server-side id (breaking dedup and
edges) and drops label/tier/tags/importance/confidence (verified in a sandbox:
a tier "Canonical" probe came back "Working"). load-merge preserves the id and
every field, dedups nodes by id and edges by (from,to,relation) so retries are
no-ops, and calls persist_canonical() so THE OWNER writes its own file — the
ownership rule is honoured rather than worked around.
Spool-and-drain rather than push-per-write: measured ~0.38s per load-merge at
live scale (79k nodes/176MB), and a chat turn writes 5-7 nodes. The spool is on
disk, not in process state, because the soul serves each connection on its own
pthread and a shared buffer would lose entries to a read-modify-write race. That
also buys crash recovery: writes orphaned by kill -9 are drained on next boot.
Honesty: api_persisted (the gate all 10 MCP write handlers pass through) and
mem_store now assert AT THE OWNER instead of reading back the soul's own RAM.
With the owner down a write returns {"ok":false,"error":"write_not_persisted"}
and the delta is queued — where main returns {"ok":true} for a write that dies.
Coverage: 35 node sites + 9 edge sites routed through the boundary. Deliberately
excluded, with reasons in persist.el: 4 InternalStateEvent sites (Will's own
telemetry carve-out), the boot counter and the persona (both already have
bespoke owner-side write-backs), and soul.el's 54 genesis identity edges
(file-mode only). engram_strengthen and engram_forget are NOT propagated —
load-merge cannot update or delete, and hard-deleting at the owner would fail
verify-soul-contract.sh section B.
Also fixed here:
- routes.el GET /api/graph/edges engram_save()'d straight over the owner's
canonical snapshot.json — a read route, in a non-owner process, clobbering the
canonical on every call. Same defect class Will removed from the engram in el
dc39a61. Now exports to a scratch path. With this gone the soul writes nothing
at all in HTTP mode.
- persist.el must clear the runtime's _tl_fs_read_len hint after every fs_read.
In vendored runtime v1.0.0-20260501 that hint becomes the NEXT response's
Content-Length, so reading a spool file mid-request made an 86-byte reply go
out as 497 bytes with 411 bytes of adjacent heap trailing it. Caught and fixed
at our boundary; the runtime class was fixed upstream in el 43636ae, which is
not the pinned runtime here.
Rung: E2E-VERIFIED, discriminating. Same harness, same engram binary:
write-through: LEG 1 PRESENT at owner, LEG 2 SURVIVED kill -9 + restart
main: LEG 1 ABSENT at owner, LEG 2 LOST
verify-soul-contract.sh: GATE PASS on both builds (27/27 routes, immutability).
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
427 lines
22 KiB
EmacsLisp
427 lines
22 KiB
EmacsLisp
// persist.el — the soul→engram WRITE-THROUGH boundary (neuron#117).
|
|
//
|
|
// WHY THIS FILE EXISTS
|
|
// soul.el:571-573 states the ownership rule: "when ENGRAM_URL is set the HTTP
|
|
// Engram owns persistence — the soul must NEVER write to the local snapshot
|
|
// (not the persistence owner)." The soul obeys the NEGATIVE half. The POSITIVE
|
|
// half — how a write made inside the soul actually REACHES the owner — was
|
|
// never built. Sync is pull-only (awareness.el `/api/sync` -> engram_load_merge),
|
|
// so every node the soul creates lives in its process RAM and is shed on
|
|
// restart. Measured live 2026-08-07: soul node_count=102184, engram
|
|
// node_count=79197 — ~23k nodes existing nowhere but RAM.
|
|
//
|
|
// SCOPE NOTE ON THE PATENT (corrects an earlier internal reading)
|
|
// Engram provisional claims 15-18 describe a delta-sync protocol "with peer
|
|
// Engram instances"; claim 17's pull-then-push sequence is PEER-ENGRAM to
|
|
// PEER-ENGRAM. The soul is NOT a peer Engram — it is a CALLER of the database
|
|
// system API (cf. claim 27, "invoked explicitly by a caller of the database
|
|
// system API"). So claim 17 does not specify a soul↔engram contract and is not
|
|
// cited as authority here. This design is derived from the ownership rule
|
|
// alone: the owner owns the writes, therefore the soul must HAND writes to the
|
|
// owner and must never write the owner's file itself.
|
|
//
|
|
// THE MECHANISM, AND WHY NOT `POST /api/nodes`
|
|
// The obvious route is the one the persona/boot-counter write-backs already
|
|
// use, POST /api/nodes. It is the wrong instrument here, verified against the
|
|
// live engram binary in a sandbox:
|
|
// - it mints a NEW server-side id (engram_node), so the soul's id and the
|
|
// owner's id diverge — the next /api/sync pull re-imports the node as a
|
|
// DUPLICATE, and any edge referencing the soul's id never resolves;
|
|
// - it accepts only {content, node_type, salience} and drops label, tier,
|
|
// tags, importance, confidence, metadata. A probe posted with tier
|
|
// "Canonical" came back tier "Working", importance 0.5.
|
|
// POST /api/load-merge (Will's own route, el `dc39a61`) is the right one:
|
|
// - engram_load_merge PRESERVES the id and every field;
|
|
// - it dedups nodes by id and edges by (from_id,to_id,relation), so a
|
|
// re-submitted delta is a NO-OP — retry safety is free, and it is the same
|
|
// local-wins semantics the graph already uses;
|
|
// - it calls persist_canonical() — THE OWNER writes its own canonical file.
|
|
// The soul never touches it. The ownership rule is honoured in its
|
|
// strongest form rather than worked around;
|
|
// - it returns real counts {ok, nodes_added, edges_added, node_count},
|
|
// so a receipt can be a MEASUREMENT instead of a fixed success shape.
|
|
//
|
|
// SPOOL-AND-DRAIN, AND WHY IT IS NOT JUST A DIRECT POST
|
|
// Measured in a sandbox against a 79k-node / 176MB graph (live scale): one
|
|
// load-merge costs ~0.38s, essentially all of it the owner's persist_canonical.
|
|
// A chat turn writes 5-7 nodes; pushing each separately would add ~2.7s per
|
|
// turn. So writes are STAGED and pushed in one coalesced batch.
|
|
// The staging buffer is the FILESYSTEM, not process state, because the soul
|
|
// serves each HTTP connection on its own pthread (el_runtime http_serve_async)
|
|
// and a shared in-process buffer would lose entries to a read-modify-write
|
|
// race — silently, which is the one failure mode this file exists to end.
|
|
// One file per write, named with uuid_v4, is race-free by construction and
|
|
// buys a property a memory buffer cannot: writes that could not be pushed
|
|
// SURVIVE A SOUL CRASH and are drained on the next boot.
|
|
//
|
|
// WHAT IS DELIBERATELY NOT PUSHED
|
|
// - InternalStateEvent / heartbeat telemetry. Will's own carve-out, stated in
|
|
// engram server.el 8f8ccc9: "48h-pruned, loss-tolerant, ~2/min; snapshotting
|
|
// 28MB per heartbeat is waste."
|
|
// NOTE (ours, flagged for Will): we do NOT additionally exclude Working-tier
|
|
// nodes. That exclusion exists in `fb0bb55` to stop the boot counter leaking
|
|
// through the /api/sync PULL; it is about sync backflow, not durability.
|
|
// Applying it here would exclude mem_store — which writes tier "Working" — and
|
|
// mem_store is the single most important durable write path in the soul. Boot
|
|
// seeding reads the canonical file wholesale, so a pushed Working-tier node
|
|
// does survive restart. This is the one classification call this file makes
|
|
// that Will has not ruled on.
|
|
//
|
|
// WHAT THIS BOUNDARY CANNOT EXPRESS (by construction, not by omission)
|
|
// - engram_strengthen (salience/activation drift): load-merge SKIPS ids that
|
|
// already exist, so it cannot update an existing node. There is no owner-side
|
|
// update/upsert route. Not pushable through any current route; left as a
|
|
// follow-up that needs a change in the engram repo.
|
|
// - engram_forget (hard delete): load-merge is additive and has no delete verb.
|
|
// Propagating deletes would mean DELETE /api/nodes/<id>, a HARD delete at the
|
|
// owner — which scripts/verify-soul-contract.sh section B explicitly fails the
|
|
// build for ("to delete is to supersede/tombstone, never hard-remove"). Local
|
|
// deletes therefore stay local; the TOMBSTONE NODE and its "tombstones" edge
|
|
// (mem_tombstone) are pushed, and that is the sanctioned representation of a
|
|
// deletion in this graph.
|
|
|
|
// ── Configuration ─────────────────────────────────────────────────────────────
|
|
|
|
// wt_engram_url — same resolution order as ise_post: env, then the state key
|
|
// stashed at boot. NO hardcoded localhost fallback: unlike telemetry, inventing
|
|
// a destination for durable data would risk pushing a user's memories at whatever
|
|
// happens to be listening on 8742. Empty means "no HTTP owner" -> file mode.
|
|
fn wt_engram_url() -> String {
|
|
let env_url: String = env("ENGRAM_URL")
|
|
if !str_eq(env_url, "") { return env_url }
|
|
return state_get("soul_engram_url")
|
|
}
|
|
|
|
fn wt_api_key() -> String {
|
|
let env_key: String = env("ENGRAM_API_KEY")
|
|
if !str_eq(env_key, "") { return env_key }
|
|
return state_get("soul_engram_api_key")
|
|
}
|
|
|
|
// wt_enabled — true only in HTTP-engram mode. In file mode the soul IS the
|
|
// persistence owner and every path below is a no-op, so this whole feature is
|
|
// inert for genesis/local deployments. That is also what makes it reversible.
|
|
fn wt_enabled() -> Bool {
|
|
return !str_eq(wt_engram_url(), "")
|
|
}
|
|
|
|
// wt_spool_dir — where staged deltas live. MUST be readable by the engram
|
|
// process: /api/load-merge takes a PATH and the owner opens it itself. Both
|
|
// processes are same-host by construction (dev-stack LaunchAgents; the GKE
|
|
// image starts engram and soul in one container per entrypoint.sh).
|
|
fn wt_spool_dir() -> String {
|
|
let raw: String = env("SOUL_OUTBOX_DIR")
|
|
let dir: String = if str_eq(raw, "") { env("HOME") + "/.neuron/soul-outbox" } else { raw }
|
|
fs_mkdir(dir)
|
|
return dir
|
|
}
|
|
|
|
// ── Helpers ───────────────────────────────────────────────────────────────────
|
|
|
|
// wt_esc — minimal JSON string escape. Deliberately local rather than reusing
|
|
// chat.el's json_safe: persist.el is imported BY memory.el, which is imported by
|
|
// chat.el, so depending on chat.el here would be an import cycle.
|
|
fn wt_esc(s: String) -> String {
|
|
let s1: String = str_replace(s, "\\", "\\\\")
|
|
let s2: String = str_replace(s1, "\"", "\\\"")
|
|
let s3: String = str_replace(s2, "\n", "\\n")
|
|
let s4: String = str_replace(s3, "\r", "\\r")
|
|
let s5: String = str_replace(s4, "\t", "\\t")
|
|
return s5
|
|
}
|
|
|
|
// wt_durable_class — Will's telemetry carve-out, by node_type. See header.
|
|
fn wt_durable_class(node_type: String) -> Bool {
|
|
if str_eq(node_type, "InternalStateEvent") { return false }
|
|
return true
|
|
}
|
|
|
|
// wt_inner — strip the surrounding brackets off a JSON array so several arrays
|
|
// can be concatenated into one. Returns "" for "[]" / "" / anything too short.
|
|
fn wt_inner(arr: String) -> String {
|
|
let n: Int = str_len(arr)
|
|
if n < 3 { return "" }
|
|
if !str_starts_with(arr, "[") { return "" }
|
|
return str_slice(arr, 1, n - 1)
|
|
}
|
|
|
|
// wt_read — fs_read, plus a MANDATORY reset of the runtime's binary-length hint.
|
|
//
|
|
// THIS IS NOT OPTIONAL AND MUST NOT BE "SIMPLIFIED" BACK TO A BARE fs_read.
|
|
// The pinned runtime (vendor/el-runtime/v1.0.0-20260501) keeps a thread-local
|
|
// `_tl_fs_read_len` that fs_read SETS to the file's byte count (so binary files
|
|
// can be served with a correct Content-Length) and that http_send_response
|
|
// CONSUMES as the Content-Length of the next reply. Nothing else clears it
|
|
// except json_get_raw. So any fs_read during request handling that is not
|
|
// followed by a json_get_raw makes the NEXT HTTP response advertise the FILE's
|
|
// length instead of the body's — and the runtime then sends that many bytes,
|
|
// appending whatever adjacent heap memory follows the reply.
|
|
//
|
|
// Caught here, measured: a /api/neuron/memory reply that should be 86 bytes went
|
|
// out as 497, with 411 bytes of this module's own spool paths and log strings
|
|
// trailing the JSON. The drain reads spool files mid-request, so this boundary
|
|
// is exactly where the landmine gets stepped on.
|
|
//
|
|
// Upstream el fixed the class in `43636ae` ("pair fs_read length hint with its
|
|
// buffer"); that runtime is NOT the one vendored here, and re-pinning the
|
|
// runtime is deliberately out of scope for this change. Clearing the hint at
|
|
// our own boundary fixes our exposure without touching the pinned C.
|
|
// json_get_raw is used as the reset because it is the only builtin in this
|
|
// runtime that zeroes the hint, and it does so before any early return.
|
|
fn wt_clear_binlen() -> Void {
|
|
let discard: String = json_get_raw("{}", "_wt_reset")
|
|
}
|
|
|
|
fn wt_read(path: String) -> String {
|
|
let data: String = fs_read(path)
|
|
wt_clear_binlen()
|
|
return data
|
|
}
|
|
|
|
// wt_sweep — best-effort removal of the zero-byte husks left by truncation.
|
|
// The runtime exposes no unlink builtin, so a drained delta is emptied rather
|
|
// than deleted; this reclaims the directory entries.
|
|
//
|
|
// `-empty` is the safety property, not an optimisation: the command is
|
|
// STRUCTURALLY INCAPABLE of removing a delta that still has content, so it can
|
|
// never destroy a pending write even if it runs concurrently with a stage.
|
|
// Only the directory path is interpolated (never a filename), and it is quoted.
|
|
// The exit code is ignored — an un-swept husk costs one directory entry.
|
|
fn wt_sweep(dir: String) -> Void {
|
|
if str_eq(dir, "") { return }
|
|
if str_contains(dir, "'") { return }
|
|
exec_command("find '" + dir + "' -maxdepth 1 -name 'wt*.json' -empty -delete 2>/dev/null")
|
|
}
|
|
|
|
// ── Staging ───────────────────────────────────────────────────────────────────
|
|
|
|
// wt_stage — write ONE delta file. uuid_v4 in the name makes concurrent stagers
|
|
// collision-free without any lock. Returns true if the delta is on disk.
|
|
fn wt_stage(nodes_json: String, edges_json: String) -> Bool {
|
|
let dir: String = wt_spool_dir()
|
|
if str_eq(dir, "") { return false }
|
|
let payload: String = "{\"nodes\":" + nodes_json + ",\"edges\":" + edges_json + "}"
|
|
let path: String = dir + "/wt-" + uuid_v4() + ".json"
|
|
fs_write(path, payload)
|
|
// Read-back-verify the stage itself. A stage that did not land is a write we
|
|
// would otherwise believe was queued — exactly the hallucinated-save class.
|
|
if str_eq(wt_read(path), "") {
|
|
println("[persist] wt_stage: FAILED to write spool file " + path + " — delta not queued")
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
// ── The write boundary ────────────────────────────────────────────────────────
|
|
|
|
// wt_node — create a node locally AND queue it for the persistence owner.
|
|
// Same signature and same return contract as engram_node_full ("" on failure),
|
|
// so converting a call site is a rename and nothing else.
|
|
fn wt_node(content: String, node_type: String, label: String,
|
|
salience: Float, importance: Float, confidence: Float,
|
|
tier: String, tags: String) -> String {
|
|
let id: String = engram_node_full(content, node_type, label,
|
|
salience, importance, confidence,
|
|
tier, tags)
|
|
if str_eq(id, "") { return "" }
|
|
// engram_get_node_json emits the SAME record shape engram_save writes (minus
|
|
// the embedding vector, which the owner backfills lazily), so the read-back
|
|
// doubles as the delta payload — no second serialization to drift.
|
|
let rec: String = engram_get_node_json(id)
|
|
if str_eq(rec, "") || str_eq(rec, "{}") {
|
|
println("[persist] wt_node: local write did not read back, id=" + id + " label=" + label)
|
|
return ""
|
|
}
|
|
if wt_enabled() && wt_durable_class(node_type) {
|
|
wt_stage("[" + rec + "]", "[]")
|
|
}
|
|
return id
|
|
}
|
|
|
|
// wt_edge — create an edge locally AND queue it. Mirrors engram_connect.
|
|
//
|
|
// The edge id is freshly generated rather than read back: the runtime exposes no
|
|
// "id of the edge I just created" accessor, and the owner dedups edges by
|
|
// (from_id,to_id,relation), never by id — so the id is not load-bearing. The
|
|
// consequence, stated plainly: the soul's copy and the owner's copy of the same
|
|
// edge carry different edge ids. Nothing in either codebase looks an edge up by
|
|
// id (neighbors traversal scans from_id/to_id), so this is cosmetic.
|
|
fn wt_edge(from_id: String, to_id: String, weight: Float, relation: String) -> Void {
|
|
engram_connect(from_id, to_id, weight, relation)
|
|
if !wt_enabled() { return }
|
|
if str_eq(from_id, "") || str_eq(to_id, "") { return }
|
|
let ts: Int = time_now()
|
|
let rec: String = "{\"id\":\"" + uuid_v4() + "\""
|
|
+ ",\"from_id\":\"" + wt_esc(from_id) + "\""
|
|
+ ",\"to_id\":\"" + wt_esc(to_id) + "\""
|
|
+ ",\"relation\":\"" + wt_esc(relation) + "\""
|
|
+ ",\"metadata\":\"{}\""
|
|
+ ",\"weight\":" + float_to_str(weight)
|
|
+ ",\"confidence\":1"
|
|
+ ",\"created_at\":" + int_to_str(ts)
|
|
+ ",\"updated_at\":" + int_to_str(ts)
|
|
+ ",\"last_fired\":0,\"inhibitory\":0,\"layer_id\":1}"
|
|
wt_stage("[]", "[" + rec + "]")
|
|
}
|
|
|
|
// ── The drain ─────────────────────────────────────────────────────────────────
|
|
|
|
// wt_drain — coalesce every staged delta into ONE load-merge against the owner.
|
|
//
|
|
// Returns: nodes_added on success (>= 0), 0 when there was nothing to do, and
|
|
// -1 when the push FAILED. -1 is load-bearing: on failure the spool files are
|
|
// left untouched, so nothing is lost and the next drain retries them. A caller
|
|
// must never read a non-negative return as "my particular node is durable" —
|
|
// use wt_durable(id) for that.
|
|
//
|
|
// Concurrency: several threads may drain at once. Each builds its own batch file
|
|
// (uuid-named), and overlapping batches are harmless because load-merge dedups.
|
|
// Files are truncated ONLY after a confirmed ok:true, so a lost race costs a
|
|
// redundant push, never a dropped write.
|
|
fn wt_drain() -> Int {
|
|
if !wt_enabled() { return 0 }
|
|
let dir: String = wt_spool_dir()
|
|
if str_eq(dir, "") { return 0 }
|
|
|
|
// el_list_len/el_list_get, NOT json_stringify(fs_list(...)): fs_list builds
|
|
// a native list via el_list_append, and json_stringify does not serialize
|
|
// that type — it renders the raw pointer value. (Verified in isolation; the
|
|
// same latent defect is live in studio.el's /api/tools/file/list route,
|
|
// which returns e.g. {"entries":4386409744}. Noted, not fixed here.)
|
|
let listing = fs_list(dir)
|
|
let count: Int = el_list_len(listing)
|
|
if count == 0 { return 0 }
|
|
|
|
let nodes_acc: String = ""
|
|
let edges_acc: String = ""
|
|
let drained: String = ""
|
|
let found: Int = 0
|
|
let i: Int = 0
|
|
// No `continue` / `break`: elc lists them as keywords but not one line of
|
|
// the shipped soul uses either, so they are unexercised on this build path.
|
|
// Guard conditions are expressed as nested ifs instead, and every rebind is
|
|
// at the loop-body top level where `let x = ...` is assignment (the idiom
|
|
// memory.el's boot-counter loop relies on) — never inside a nested block,
|
|
// where it would shadow instead.
|
|
while i < count {
|
|
let name: String = el_list_get(listing, i)
|
|
// A delta is only usable when it ends with the closing "]}" that
|
|
// wt_stage writes last. fs_write is not atomic, so a file being written
|
|
// right now can be observed half-formed; requiring the terminator means
|
|
// it is picked up whole on the next drain instead of merged as garbage.
|
|
// An empty read means "already drained and truncated" — not an error.
|
|
let p: String = if str_starts_with(name, "wt-") { dir + "/" + name } else { "" }
|
|
let raw: String = if str_eq(p, "") { "" } else { wt_read(p) }
|
|
let usable: Bool = !str_eq(raw, "") && str_ends_with(raw, "]}")
|
|
let nj: String = if usable { wt_inner(json_get_raw(raw, "nodes")) } else { "" }
|
|
let ej: String = if usable { wt_inner(json_get_raw(raw, "edges")) } else { "" }
|
|
let nodes_acc = if str_eq(nj, "") { nodes_acc } else if str_eq(nodes_acc, "") { nj } else { nodes_acc + "," + nj }
|
|
let edges_acc = if str_eq(ej, "") { edges_acc } else if str_eq(edges_acc, "") { ej } else { edges_acc + "," + ej }
|
|
let drained = if !usable { drained } else if str_eq(drained, "") { p } else { drained + "\n" + p }
|
|
let found = if usable { found + 1 } else { found }
|
|
let i = i + 1
|
|
}
|
|
|
|
if found == 0 { return 0 }
|
|
|
|
let combined: String = "{\"nodes\":[" + nodes_acc + "],\"edges\":[" + edges_acc + "]}"
|
|
let batch: String = dir + "/wtb-" + uuid_v4() + ".json"
|
|
fs_write(batch, combined)
|
|
if str_eq(wt_read(batch), "") {
|
|
println("[persist] wt_drain: could not write batch file " + batch + " — " + int_to_str(found) + " deltas stay queued")
|
|
return -1
|
|
}
|
|
|
|
let url: String = wt_engram_url()
|
|
let key: String = wt_api_key()
|
|
let body: String = "{\"path\":\"" + wt_esc(batch) + "\",\"_auth\":\"" + wt_esc(key) + "\"}"
|
|
let resp: String = http_post_json(url + "/api/load-merge", body)
|
|
|
|
// The batch file is pure scratch — the retry is rebuilt from the SPOOL, not
|
|
// from it. Truncate it unconditionally, before branching on the outcome, so
|
|
// a persistently unreachable owner cannot accumulate one husk per attempt.
|
|
fs_write(batch, "")
|
|
|
|
// Distinguish the two failures rather than collapsing them: "cannot reach
|
|
// the owner" and "the owner refused this delta" need different human
|
|
// responses, and a log line that says the wrong one costs a debugging hour.
|
|
// curl surfaces transport errors as a JSON body, so an empty response is not
|
|
// the only unreachable signal.
|
|
// (str_contains rather than a strict parse on purpose — the engram's HTTP
|
|
// responses have been observed carrying trailing bytes past the JSON.)
|
|
let unreachable: Bool = str_eq(resp, "")
|
|
|| str_contains(resp, "Couldn't connect")
|
|
|| str_contains(resp, "Failed to connect")
|
|
|| str_contains(resp, "Could not resolve")
|
|
|| str_contains(resp, "timed out")
|
|
if unreachable {
|
|
wt_sweep(dir)
|
|
println("[persist] wt_drain: owner UNREACHABLE at " + url + " — " + int_to_str(found)
|
|
+ " deltas stay queued in " + dir + " (will retry): " + resp)
|
|
return -1
|
|
}
|
|
if !str_contains(resp, "\"ok\":true") {
|
|
wt_sweep(dir)
|
|
println("[persist] wt_drain: owner REJECTED the delta — " + int_to_str(found)
|
|
+ " stay queued in " + dir + ": " + resp)
|
|
return -1
|
|
}
|
|
|
|
let added: Int = json_get_int(resp, "nodes_added")
|
|
let added_e: Int = json_get_int(resp, "edges_added")
|
|
|
|
// Confirmed. Truncate the drained spool files so they are not re-pushed.
|
|
// Truncation (not deletion) because the runtime exposes no unlink builtin;
|
|
// an emptied file is inert to the loop above. The zero-byte husks are then
|
|
// swept below.
|
|
let paths = str_split(drained, "\n")
|
|
let pn: Int = el_list_len(paths)
|
|
let k: Int = 0
|
|
while k < pn {
|
|
let one: String = el_list_get(paths, k)
|
|
if !str_eq(one, "") { fs_write(one, "") }
|
|
let k = k + 1
|
|
}
|
|
wt_sweep(dir)
|
|
|
|
println("[persist] wt_drain: pushed " + int_to_str(found) + " deltas -> owner added "
|
|
+ int_to_str(added) + " nodes, " + int_to_str(added_e) + " edges")
|
|
return added
|
|
}
|
|
|
|
// wt_durable — is this id present AT THE OWNER? The only honest answer to
|
|
// "did my write persist" in HTTP mode.
|
|
//
|
|
// In file mode the soul IS the owner, so the local read-back is the owner-side
|
|
// read-back and this collapses to the pre-existing check.
|
|
//
|
|
// nodes_added from wt_drain is NOT a substitute: a concurrent drain may have
|
|
// already pushed this node, making our own added count 0 while the node is
|
|
// perfectly durable. Presence at the owner is the fact; counts are telemetry.
|
|
fn wt_durable(id: String) -> Bool {
|
|
if str_eq(id, "") { return false }
|
|
if !wt_enabled() {
|
|
let local: String = engram_get_node_json(id)
|
|
return !str_eq(local, "") && !str_eq(local, "null") && !str_eq(local, "{}")
|
|
}
|
|
let url: String = wt_engram_url()
|
|
let resp: String = http_get(url + "/api/nodes/" + id)
|
|
if str_eq(resp, "") { return false }
|
|
if str_eq(resp, "{}") { return false }
|
|
return str_contains(resp, "\"id\"")
|
|
}
|
|
|
|
// wt_commit — flush, then assert at the owner. The receipt callers should use.
|
|
// Deliberately NOT a fixed success shape: it can and does return false while the
|
|
// local write is perfectly fine in RAM, which is the true state of affairs when
|
|
// the owner is unreachable.
|
|
fn wt_commit(id: String) -> Bool {
|
|
if str_eq(id, "") { return false }
|
|
if !wt_enabled() {
|
|
let local: String = engram_get_node_json(id)
|
|
return !str_eq(local, "") && !str_eq(local, "null") && !str_eq(local, "{}")
|
|
}
|
|
let pushed: Int = wt_drain()
|
|
return wt_durable(id)
|
|
}
|