Compare commits
38 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| c21074b547 | |||
| 4e24d7d3f1 | |||
| 7557ea6e19 | |||
| 7351fb0a8d | |||
| d545b69614 | |||
| 598915cc61 | |||
| dab14f9100 | |||
| 40eb48e92f | |||
| c9f75e2592 | |||
| 09dade0613 | |||
| 38a8e32d6c | |||
| 15b66c8b1a | |||
| 9883aa7564 | |||
| d45a0882f3 | |||
| 1c9de03fdb | |||
| 3718bf0380 | |||
| 9d40f87926 | |||
| 90d3f0bc76 | |||
| ee39aa5f17 | |||
| 2d0aef4ef8 | |||
| 979e820f68 | |||
| e29fe4fd0b | |||
| 0e924f7df9 | |||
| 1bb1edc851 | |||
| 1010185978 | |||
| ff37835ae5 | |||
| e5c80359a8 | |||
| b53b5b4e8a | |||
| 20bd9ed00b | |||
| 70982498e0 | |||
| 373265c05d | |||
| ed722b9e2e | |||
| b0a78c5737 | |||
| 447d042022 | |||
| d4e82d3d56 | |||
| 40bb6ff579 | |||
| d5411fb58a | |||
| b2aac4bf89 |
+17
-5
@@ -41,17 +41,29 @@ fn strip_query(path: String) -> String {
|
||||
str_slice(path, 0, q)
|
||||
}
|
||||
|
||||
// query_param — extract one query-string value, URL-DECODED.
|
||||
//
|
||||
// The decode step was missing (found 2026-08-15): a claim sent as
|
||||
// "test%20claim" arrived at engram_assert_json still percent-encoded and was
|
||||
// stored/compared that way, so any value containing a space, &, =, or non-ASCII
|
||||
// character silently became a different string than the caller sent. Affects
|
||||
// every GET route that reads params this way, not just /api/assert.
|
||||
fn query_param(path: String, key: String) -> String {
|
||||
let q: Int = str_index_of(path, "?")
|
||||
if q < 0 { return "" }
|
||||
let qs: String = str_slice(path, q + 1, str_len(path))
|
||||
let needle: String = key + "="
|
||||
let pos: Int = str_index_of(qs, needle)
|
||||
// Anchor the match to a real key boundary: prefixing "&" and searching for
|
||||
// "&key=" means "q" can never match inside "faq=". (Found 2026-08-15:
|
||||
// "?faq=X&q=Y" returned X for key "q" — a silently wrong value, not an
|
||||
// error.) The leading "&" makes the first parameter match the same way.
|
||||
let hay: String = "&" + qs
|
||||
let needle: String = "&" + key + "="
|
||||
let pos: Int = str_index_of(hay, needle)
|
||||
if pos < 0 { return "" }
|
||||
let after: String = str_slice(qs, pos + str_len(needle), str_len(qs))
|
||||
let after: String = str_slice(hay, pos + str_len(needle), str_len(hay))
|
||||
let amp: Int = str_index_of(after, "&")
|
||||
if amp < 0 { return after }
|
||||
str_slice(after, 0, amp)
|
||||
let raw: String = if amp < 0 { after } else { str_slice(after, 0, amp) }
|
||||
return __url_decode(raw)
|
||||
}
|
||||
|
||||
fn query_int(path: String, key: String, default_val: Int) -> Int {
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
build/
|
||||
+215
-262
@@ -1,17 +1,28 @@
|
||||
// ingest.el — the native EL AFFERENT INGEST ORGAN
|
||||
//
|
||||
// The source-polymorphic ingest(source) primitive: point it at a directory,
|
||||
// file, url, llm-query, structured-primitive set, or stream; it EXTRACTS the
|
||||
// real content faithfully (no invention), TRANSDUCES it into a DISCRETE
|
||||
// MANIFOLD (multiple nodes + internal edges — meaning-structure, never a
|
||||
// single blob; the conversion from extracted surface content into geometry
|
||||
// is automatic and invisible to the caller, the way digestion is invisible
|
||||
// to the one who chose to eat — ingest is the conscious act, transduce is
|
||||
// the mechanism underneath it, and it is no less real for being unseen),
|
||||
// and MERGES that manifold into the engram geometry: shared
|
||||
// meanings DEDUP onto existing nodes (search + exact/cosine match), genuinely
|
||||
// new meanings add nodes, relations add edges. Every node enters with
|
||||
// PROVENANCE + grounding-level + stewardship class from the moment of entry.
|
||||
// file, url, llm-query, or stream; it EXTRACTS the real content faithfully
|
||||
// (no invention), TRANSDUCES it into a DISCRETE MANIFOLD (multiple nodes +
|
||||
// internal edges — meaning-structure, never a single blob; the conversion
|
||||
// from extracted surface content into geometry is automatic and invisible
|
||||
// to the caller, the way digestion is invisible to the one who chose to
|
||||
// eat — ingest is the conscious act, transduce is the mechanism underneath
|
||||
// it, and it is no less real for being unseen), and MERGES that manifold
|
||||
// into the engram geometry: shared meanings DEDUP onto existing nodes
|
||||
// (search + exact/cosine match), genuinely new meanings add nodes,
|
||||
// relations add edges. Every node enters with PROVENANCE + grounding-level
|
||||
// + stewardship class from the moment of entry.
|
||||
//
|
||||
// transduce() is THE single mechanism — one function, polymorphic, with no
|
||||
// content-type branch inside it. It does not ask whether a payload is
|
||||
// prose, structured data, or raw/opaque bytes (audio, or anything else);
|
||||
// it runs one boundary-scan-with-fixed-window-fallback chunking algorithm
|
||||
// and one dedup mechanism on whatever bytes it's handed, unconditionally.
|
||||
// Any deeper structure a payload might have (shared fields, relationships,
|
||||
// what a chunk of audio "means") is NOT interpreted here — that's left
|
||||
// entirely to the engram's own mechanisms (embedding, spreading activation,
|
||||
// dedup) acting on this real geometry over time. This organ claims zero
|
||||
// semantic understanding of any payload it transduces.
|
||||
//
|
||||
// It is a pure HTTP CLIENT of the engram server — it links only el_runtime.c
|
||||
// via fs/http/json/string builtins; it never links el_seed.c or the engram
|
||||
@@ -46,60 +57,6 @@ fn j_q(s: String) -> String {
|
||||
return "\"" + j_esc(s) + "\""
|
||||
}
|
||||
|
||||
// Extract the top-level keys of a JSON object string. A thin, self-contained
|
||||
// scanner (FLAGGED: the one non-trivial parser in this organ — everything else
|
||||
// is faithful text handling). Tracks string state + brace/bracket depth; a key
|
||||
// is a string at object-interior depth 1 immediately followed by ':'.
|
||||
fn json_object_keys(obj: String) -> [String] {
|
||||
let keys: [String] = el_list_empty()
|
||||
let n: Int = str_len(obj)
|
||||
let i: Int = 0
|
||||
let depth: Int = 0
|
||||
let in_str: Bool = false
|
||||
let esc: Bool = false
|
||||
let str_start: Int = -1
|
||||
let cur: String = ""
|
||||
let have_key: Bool = false
|
||||
while i < n {
|
||||
let c: String = str_char_at(obj, i)
|
||||
if in_str {
|
||||
if esc {
|
||||
esc = false
|
||||
} else {
|
||||
if str_eq(c, "\\") {
|
||||
esc = true
|
||||
} else {
|
||||
if str_eq(c, "\"") {
|
||||
in_str = false
|
||||
cur = str_slice(obj, str_start + 1, i)
|
||||
have_key = true
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if str_eq(c, "\"") {
|
||||
in_str = true
|
||||
str_start = i
|
||||
}
|
||||
if str_eq(c, "{") { depth = depth + 1 }
|
||||
if str_eq(c, "}") { depth = depth - 1 }
|
||||
if str_eq(c, "[") { depth = depth + 1 }
|
||||
if str_eq(c, "]") { depth = depth - 1 }
|
||||
if str_eq(c, ":") {
|
||||
if have_key {
|
||||
if depth == 1 {
|
||||
keys = el_list_append(keys, cur)
|
||||
}
|
||||
}
|
||||
have_key = false
|
||||
}
|
||||
if str_eq(c, ",") { have_key = false }
|
||||
}
|
||||
i = i + 1
|
||||
}
|
||||
return keys
|
||||
}
|
||||
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
// SECTION B — engram HTTP client (provenance-carrying afferent LOAD)
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
@@ -402,168 +359,125 @@ fn head80(s: String) -> String {
|
||||
// We accumulate into module-level lists carried by the caller.
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
// PROSE: chunk text into a discrete manifold. Split on blank lines into
|
||||
// paragraphs; every non-empty paragraph is its own node (NEVER one blob).
|
||||
// Edges: doc-root -contains-> chunk; chunk -precedes-> next chunk;
|
||||
// most-recent-heading -section_of-> chunk. Content is verbatim (substring of
|
||||
// the source) — pure extraction of ground truth.
|
||||
fn transduce_prose(nodes: [String], edges: [String], text: String,
|
||||
prov: String, ground: String, steward: String,
|
||||
root_lid: String, root_title: String) -> [String] {
|
||||
// returns [nodes_json_list_encoded, edges_json_list_encoded] is awkward in
|
||||
// EL; instead we mutate by returning a 2-list. We package results as a
|
||||
// single JSON array string carrying {nodes:[...],edges:[...]} additions.
|
||||
// (Kept simple: caller passes empty lists and receives the packaged pair.)
|
||||
// TRANSDUCE — the single mechanism. Takes ANY payload (prose, structured
|
||||
// data, raw/opaque bytes — audio, whatever) as one opaque string and turns
|
||||
// it into a discrete manifold: nodes + internal edges. There is no
|
||||
// content-type branch anywhere in this function. It never asks "is this
|
||||
// text," "is this JSON," "is this audio" — it runs ONE algorithm on the
|
||||
// bytes it is given, unconditionally:
|
||||
//
|
||||
// 1. BOUNDARY SCAN — split on "\n\n". This is a property of the bytes
|
||||
// (does a blank-line-style marker occur in them, yes or no), not a
|
||||
// classification of what the content IS. Prose paragraphs split on it
|
||||
// because that's how prose is typically written; that's a fact about
|
||||
// the bytes, not a rule this function knows about prose. Anything else
|
||||
// that happens to contain the same marker splits on it too, and
|
||||
// anything that doesn't, doesn't — same code path either way.
|
||||
// 2. FIXED-WINDOW FALLBACK — if step 1 found no boundary (0 or 1 non-empty
|
||||
// piece), the payload is cut into fixed-size windows instead. Same
|
||||
// chunk-per-node, edge-per-adjacency structure as step 1 produces; only
|
||||
// the source of the cut point differs.
|
||||
//
|
||||
// Every resulting chunk becomes its own node (never one blob), wired with
|
||||
// the same edges regardless of what's inside a chunk: root -contains->
|
||||
// chunk, chunk -precedes-> next chunk, and — if a chunk happens to start
|
||||
// with "#" — most-recent-heading -section_of-> chunk. That "#" check is a
|
||||
// structural marker (a fact about a chunk's first byte), not a decision
|
||||
// about whether this run is "the text case": chunks from any payload that
|
||||
// never happen to start with "#" simply never trigger it.
|
||||
//
|
||||
// Dedup is the existing, fully generic mechanism (find_existing_by_content,
|
||||
// via merge_manifold downstream of merge_packed) applied uniformly to every
|
||||
// chunk from every payload — there is no separate "structured" dedup path.
|
||||
// Any deeper structure that might exist inside a payload (shared fields,
|
||||
// repeated records, relationships) is NOT pre-computed here; that's left to
|
||||
// the engram's own mechanisms (embedding, spreading activation, dedup)
|
||||
// acting on this geometry over time, which is the whole point of handing it
|
||||
// raw bytes instead of a hand-coded interpretation of them.
|
||||
//
|
||||
// Byte-safety note: `source` must already be a string this function can
|
||||
// safely str_split/str_slice. Protecting it from silent truncation (El
|
||||
// strings are NUL-unsafe under strlen-based ops; fs_read()'s result
|
||||
// truncates at the first embedded NUL, which is routine in real binary
|
||||
// bytes) is a MECHANICAL fidelity concern that belongs to whatever produced
|
||||
// `source` (see ingest_file's file_source_string below) — not a
|
||||
// content-type judgment made in here. transduce() never learns whether a
|
||||
// chunk is plain text or a base64-encoded raw-byte window; every chunk is
|
||||
// handled identically either way.
|
||||
fn transduce(nodes: [String], edges: [String], source: String,
|
||||
prov: String, ground: String, steward: String,
|
||||
root_lid: String, root_title: String) -> [String] {
|
||||
let tagbase: String = "prov:" + prov + " ground:" + ground + " steward:" + steward
|
||||
// root node
|
||||
nodes = el_list_append(nodes, mk_node(root_lid, "document: " + root_title,
|
||||
"Concept", "Semantic", "0.6", "0.6", "0.9", tagbase + " kind:document"))
|
||||
nodes = el_list_append(nodes, mk_node(root_lid, "source: " + root_title,
|
||||
"Concept", "Semantic", "0.6", "0.6", "0.9", tagbase + " kind:source"))
|
||||
|
||||
let paras: [String] = str_split(text, "\n\n")
|
||||
let np: Int = el_list_len(paras)
|
||||
let idx: Int = 0
|
||||
// step 1: universal boundary scan
|
||||
let boundary_parts: [String] = str_split(source, "\n\n")
|
||||
let chunks: [String] = el_list_empty()
|
||||
let bp_n: Int = el_list_len(boundary_parts)
|
||||
let bp_i: Int = 0
|
||||
while bp_i < bp_n {
|
||||
let piece: String = str_trim(el_list_get(boundary_parts, bp_i))
|
||||
if !str_eq(piece, "") { chunks = el_list_append(chunks, piece) }
|
||||
bp_i = bp_i + 1
|
||||
}
|
||||
|
||||
// step 2: no boundary found -> fixed-size windows over the whole
|
||||
// payload. 4096 chars/node: low kilobytes — big enough to keep node
|
||||
// count sane on a large unbroken payload, small enough that each node
|
||||
// stays a legible, individually embeddable/dedupable unit rather than
|
||||
// one giant blob.
|
||||
if el_list_len(chunks) <= 1 {
|
||||
chunks = el_list_empty()
|
||||
let total: Int = str_len(source)
|
||||
let win: Int = 4096
|
||||
let off: Int = 0
|
||||
while off < total {
|
||||
let endp: Int = if off + win < total { off + win } else { total }
|
||||
let piece: String = str_slice(source, off, endp)
|
||||
if !str_eq(piece, "") { chunks = el_list_append(chunks, piece) }
|
||||
off = off + win
|
||||
}
|
||||
}
|
||||
|
||||
let nc: Int = el_list_len(chunks)
|
||||
let ci: Int = 0
|
||||
let last_chunk: String = ""
|
||||
let last_heading: String = ""
|
||||
let ci: Int = 0
|
||||
while idx < np {
|
||||
let raw: String = str_trim(el_list_get(paras, idx))
|
||||
if !str_eq(raw, "") {
|
||||
let lid: String = root_lid + ":c" + int_to_str(ci)
|
||||
let is_heading: Bool = str_starts_with(raw, "#")
|
||||
let kind: String = if is_heading { "kind:heading" } else { "kind:doc-chunk" }
|
||||
nodes = el_list_append(nodes, mk_node(lid, raw,
|
||||
"Knowledge", "Semantic", "0.55", "0.55", "0.9", tagbase + " " + kind))
|
||||
// containment: document root -contains-> chunk
|
||||
edges = el_list_append(edges, mk_edge(root_lid, "contains", lid))
|
||||
// sequence: previous chunk -precedes-> this chunk
|
||||
if !str_eq(last_chunk, "") {
|
||||
edges = el_list_append(edges, mk_edge(last_chunk, "precedes", lid))
|
||||
}
|
||||
// sectioning: most-recent heading -section_of-> this chunk
|
||||
if is_heading {
|
||||
last_heading = lid
|
||||
} else {
|
||||
if !str_eq(last_heading, "") {
|
||||
edges = el_list_append(edges, mk_edge(last_heading, "section_of", lid))
|
||||
}
|
||||
}
|
||||
last_chunk = lid
|
||||
ci = ci + 1
|
||||
while ci < nc {
|
||||
let raw: String = el_list_get(chunks, ci)
|
||||
let lid: String = root_lid + ":c" + int_to_str(ci)
|
||||
let is_heading: Bool = str_starts_with(raw, "#")
|
||||
let kind: String = if is_heading { "kind:heading" } else { "kind:chunk" }
|
||||
nodes = el_list_append(nodes, mk_node(lid, raw,
|
||||
"Knowledge", "Semantic", "0.55", "0.55", "0.9", tagbase + " " + kind))
|
||||
// containment: root -contains-> chunk
|
||||
edges = el_list_append(edges, mk_edge(root_lid, "contains", lid))
|
||||
// sequence: previous chunk -precedes-> this chunk
|
||||
if !str_eq(last_chunk, "") {
|
||||
edges = el_list_append(edges, mk_edge(last_chunk, "precedes", lid))
|
||||
}
|
||||
idx = idx + 1
|
||||
}
|
||||
// package: we return the two lists concatenated via a sentinel; but EL
|
||||
// lists can't nest heterogeneously here, so we instead return nodes and
|
||||
// rely on the caller holding edges by reference is not possible — so we
|
||||
// encode both into one list: [ "N" + nodejson ... , "E" + edgejson ... ].
|
||||
let packed: [String] = el_list_empty()
|
||||
let a: Int = 0
|
||||
let an: Int = el_list_len(nodes)
|
||||
while a < an { packed = el_list_append(packed, "N" + el_list_get(nodes, a)) a = a + 1 }
|
||||
let b: Int = 0
|
||||
let bn: Int = el_list_len(edges)
|
||||
while b < bn { packed = el_list_append(packed, "E" + el_list_get(edges, b)) b = b + 1 }
|
||||
return packed
|
||||
}
|
||||
|
||||
// STRUCTURED / RAW-GEOMETRY: ingest structured primitives (phonetics/formants,
|
||||
// instrument signatures, scene primitives) as GEOMETRY, faithfully. Normalized
|
||||
// input shape:
|
||||
// {"dataset":"<name>","primitive_type":"<t>",
|
||||
// "records":[{"key":"<id>","features":{...categorical...},"attributes":{...}}]}
|
||||
// Each record -> a primitive node; each categorical feature -> a SHARED feature
|
||||
// node (deduped across records: many primitives -> one feature node = real
|
||||
// connective geometry, meaning saturates); numeric attributes fold into the
|
||||
// primitive's content (unique values, no dedup benefit). This is knowledge
|
||||
// represented as geometry, not prose — the path speech/music/image ingest on.
|
||||
fn transduce_structured(nodes: [String], edges: [String], js: String,
|
||||
prov: String, ground: String, steward: String,
|
||||
root_lid: String) -> [String] {
|
||||
// grounding integrity: the SOURCE may declare its own epistemic grounding
|
||||
// (measured / derived / convention / ...) via a top-level "grounding" field;
|
||||
// honor it faithfully over the ingest-time default. This keeps the per-node
|
||||
// ground: facet consistent with the source's honest self-description.
|
||||
let src_ground: String = json_get_string(js, "grounding")
|
||||
let use_ground: String = if str_eq(src_ground, "") { ground } else { src_ground }
|
||||
let tagbase: String = "prov:" + prov + " ground:" + use_ground + " steward:" + steward
|
||||
let dsname: String = json_get_string(js, "dataset")
|
||||
let ptype: String = json_get_string(js, "primitive_type")
|
||||
// capture the source's own scholarly provenance citation (verbatim) onto
|
||||
// the dataset root — faithful attribution, retrievable, reachable from every
|
||||
// primitive via its -contains- edge back to the root.
|
||||
let src_cite: String = json_get_string(js, "provenance")
|
||||
let root_content: String = "dataset: " + dsname + " (" + ptype + ")"
|
||||
if !str_eq(src_cite, "") { root_content = root_content + " | provenance: " + src_cite }
|
||||
nodes = el_list_append(nodes, mk_node(root_lid, root_content,
|
||||
"Concept", "Semantic", "0.6", "0.6", "0.9", tagbase + " kind:dataset"))
|
||||
|
||||
let recs: String = json_get_raw(js, "records")
|
||||
let nr: Int = json_array_len(recs)
|
||||
let r: Int = 0
|
||||
while r < nr {
|
||||
let rec: String = json_array_get(recs, r)
|
||||
let rkey: String = json_get_string(rec, "key")
|
||||
let attrs: String = json_get_raw(rec, "attributes")
|
||||
// faithful compact serialization of the primitive's numeric signature
|
||||
let attr_str: String = flatten_pairs(attrs)
|
||||
let content: String = ptype + " " + rkey
|
||||
if !str_eq(attr_str, "") { content = content + " | " + attr_str }
|
||||
let plid: String = root_lid + ":" + rkey
|
||||
nodes = el_list_append(nodes, mk_node(plid, content,
|
||||
"Concept", "Semantic", "0.6", "0.6", "0.92",
|
||||
tagbase + " kind:primitive primitive:" + ptype + " key:" + rkey))
|
||||
edges = el_list_append(edges, mk_edge(root_lid, "contains", plid))
|
||||
|
||||
// categorical features -> SHARED (deduped) feature nodes + labelled edges
|
||||
let feats: String = json_get_raw(rec, "features")
|
||||
let fkeys: [String] = json_object_keys(feats)
|
||||
let fk: Int = el_list_len(fkeys)
|
||||
let k: Int = 0
|
||||
while k < fk {
|
||||
let fname: String = el_list_get(fkeys, k)
|
||||
let fval: String = json_get_string(feats, fname)
|
||||
// shared feature node: content is the feature=value pair; identical
|
||||
// pairs across records dedup onto ONE node (the geometry).
|
||||
let flid: String = "feat:" + fname + "=" + fval
|
||||
let fcontent: String = fname + "=" + fval
|
||||
nodes = el_list_append(nodes, mk_node(flid, fcontent,
|
||||
"Concept", "Semantic", "0.5", "0.5", "0.9",
|
||||
tagbase + " kind:feature feature:" + fname))
|
||||
edges = el_list_append(edges, mk_edge(plid, fname, flid))
|
||||
k = k + 1
|
||||
// sectioning: most-recent heading -section_of-> this chunk
|
||||
if is_heading {
|
||||
last_heading = lid
|
||||
} else {
|
||||
if !str_eq(last_heading, "") {
|
||||
edges = el_list_append(edges, mk_edge(last_heading, "section_of", lid))
|
||||
}
|
||||
}
|
||||
r = r + 1
|
||||
last_chunk = lid
|
||||
ci = ci + 1
|
||||
}
|
||||
let packed: [String] = el_list_empty()
|
||||
let a: Int = 0
|
||||
let an: Int = el_list_len(nodes)
|
||||
while a < an { packed = el_list_append(packed, "N" + el_list_get(nodes, a)) a = a + 1 }
|
||||
let b: Int = 0
|
||||
let bn: Int = el_list_len(edges)
|
||||
while b < bn { packed = el_list_append(packed, "E" + el_list_get(edges, b)) b = b + 1 }
|
||||
return packed
|
||||
}
|
||||
|
||||
// flatten a flat JSON object of scalar fields into "k=v k=v" (faithful; values
|
||||
// verbatim). Used for numeric attribute signatures.
|
||||
fn flatten_pairs(obj: String) -> String {
|
||||
if str_eq(obj, "") { return "" }
|
||||
let keys: [String] = json_object_keys(obj)
|
||||
let n: Int = el_list_len(keys)
|
||||
let out: String = ""
|
||||
let i: Int = 0
|
||||
while i < n {
|
||||
let k: String = el_list_get(keys, i)
|
||||
// json_get_raw returns the raw token — works for NUMBERS (bare, e.g.
|
||||
// "270") where json_get_string yields "" for non-string values. Strip
|
||||
// surrounding quotes if the value happens to be a string token.
|
||||
let raw: String = json_get_raw(obj, k)
|
||||
let v: String = str_replace(raw, "\"", "")
|
||||
let sep: String = if i == 0 { "" } else { " " }
|
||||
out = out + sep + k + "=" + v
|
||||
i = i + 1
|
||||
}
|
||||
return out
|
||||
// package both lists into one, "N"/"E"-prefixed (see merge_packed).
|
||||
let packed: [String] = el_list_empty()
|
||||
let pn_i: Int = 0
|
||||
let pn_n: Int = el_list_len(nodes)
|
||||
while pn_i < pn_n { packed = el_list_append(packed, "N" + el_list_get(nodes, pn_i)) pn_i = pn_i + 1 }
|
||||
let pe_i: Int = 0
|
||||
let pe_n: Int = el_list_len(edges)
|
||||
while pe_i < pe_n { packed = el_list_append(packed, "E" + el_list_get(edges, pe_i)) pe_i = pe_i + 1 }
|
||||
return packed
|
||||
}
|
||||
|
||||
// unpack the "N"/"E"-prefixed packed list back into two lists, then merge
|
||||
@@ -594,18 +508,7 @@ fn basename(path: String) -> String {
|
||||
return el_list_get(parts, n - 1)
|
||||
}
|
||||
|
||||
fn ends_with_ci(s: String, suf: String) -> Bool {
|
||||
return str_ends_with(str_to_lower(s), suf)
|
||||
}
|
||||
|
||||
fn is_text_file(path: String) -> Bool {
|
||||
return ends_with_ci(path, ".md") || ends_with_ci(path, ".txt")
|
||||
|| ends_with_ci(path, ".markdown") || ends_with_ci(path, ".text")
|
||||
}
|
||||
|
||||
// default ingestion grounding; overridable per-invocation via INGEST_GROUND.
|
||||
// Note: a source's OWN top-level "grounding" field (structured) takes precedence
|
||||
// over this — the author's honest self-description wins.
|
||||
fn default_ground() -> String {
|
||||
let g: String = env("INGEST_GROUND")
|
||||
if str_eq(g, "") { return "extracted" }
|
||||
@@ -618,25 +521,67 @@ fn default_steward() -> String {
|
||||
return s
|
||||
}
|
||||
|
||||
// ingest one file -> report JSON
|
||||
// Mechanical fidelity guard — NOT a content-type test. fs_read()'s el_val_t
|
||||
// result truncates at the first embedded NUL byte under El's strlen-based
|
||||
// string ops (see fs_size's doc comment in runtime/el_runtime.h); comparing
|
||||
// its length against fs_size() (a real stat()-based byte count) is a
|
||||
// technical fact about whether the string channel captured the file intact
|
||||
// — computed the same way for a poem, a JSON file, or a WAV, and saying
|
||||
// nothing about what the file IS. When the counts agree, `text` is
|
||||
// trustworthy verbatim. When they don't (silent truncation happened),
|
||||
// rebuild the payload as base64-encoded fixed-size windows read directly
|
||||
// off disk (fs_read_b64_chunk — binary-safe in C), joined with the same
|
||||
// "\n\n" boundary marker transduce()'s generic scan already looks for, so
|
||||
// transduce() sees one ordinary boundary-delimited payload and runs its one
|
||||
// algorithm on it exactly as it would on prose — it never learns that a
|
||||
// fidelity problem occurred upstream, let alone why.
|
||||
fn file_source_string(path: String, text: String, real_size: Int) -> String {
|
||||
if real_size <= 0 { return text }
|
||||
if str_len(text) == real_size { return text }
|
||||
// 3072 raw bytes -> 4096 base64 chars (3 divides evenly into base64's
|
||||
// 3-byte/4-char ratio); keeps each resulting node's content a clean,
|
||||
// bounded, low-kilobytes unit, same order of magnitude as the fixed
|
||||
// fallback window in transduce() itself.
|
||||
let win: Int = 3072
|
||||
let out: String = ""
|
||||
let off: Int = 0
|
||||
let first: Bool = true
|
||||
while off < real_size {
|
||||
let chunk_b64: String = fs_read_b64_chunk(path, off, win)
|
||||
if str_eq(chunk_b64, "") {
|
||||
off = real_size
|
||||
} else {
|
||||
let sep: String = if first { "" } else { "\n\n" }
|
||||
out = out + sep + chunk_b64
|
||||
first = false
|
||||
off = off + win
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// ingest one file -> report JSON. Uniform for every file regardless of
|
||||
// extension or content — transduce() decides nothing about content-type, so
|
||||
// neither does this function; it only decides whether the raw bytes made it
|
||||
// through the read intact (file_source_string), which is a fidelity
|
||||
// question, not a format one.
|
||||
fn ingest_file(path: String) -> String {
|
||||
let real_size: Int = fs_size(path)
|
||||
let text: String = fs_read(path)
|
||||
if str_eq(text, "") {
|
||||
let source: String = file_source_string(path, text, real_size)
|
||||
if str_eq(source, "") {
|
||||
return "{\"error\":\"empty or unreadable\",\"path\":" + j_q(path) + "}"
|
||||
}
|
||||
let prov: String = "file:" + path
|
||||
if ends_with_ci(path, ".json") {
|
||||
let packed: [String] = transduce_structured(el_list_empty(), el_list_empty(),
|
||||
text, prov, default_ground(), default_steward(), "ds:" + basename(path))
|
||||
return merge_packed(packed)
|
||||
}
|
||||
let packed: [String] = transduce_prose(el_list_empty(), el_list_empty(),
|
||||
text, prov, default_ground(), default_steward(),
|
||||
let packed: [String] = transduce(el_list_empty(), el_list_empty(),
|
||||
source, prov, default_ground(), default_steward(),
|
||||
"doc:" + basename(path), basename(path))
|
||||
return merge_packed(packed)
|
||||
}
|
||||
|
||||
// ingest a directory: walk one level, ingest each supported file, aggregate
|
||||
// ingest a directory: walk one level, ingest every file found, aggregate.
|
||||
// No extension filter — transduce() handles any payload uniformly now, so
|
||||
// there is no content-type gate at the directory boundary either.
|
||||
fn ingest_dir(path: String) -> String {
|
||||
let entries: [String] = fs_list(path)
|
||||
let n: Int = el_list_len(entries)
|
||||
@@ -649,14 +594,12 @@ fn ingest_dir(path: String) -> String {
|
||||
let name: String = str_trim(el_list_get(entries, i))
|
||||
if !str_eq(name, "") {
|
||||
let full: String = path + "/" + name
|
||||
if is_text_file(full) || ends_with_ci(full, ".json") {
|
||||
println("FILE " + full)
|
||||
let rep: String = ingest_file(full)
|
||||
tot_created = tot_created + json_get_int(rep, "nodes_created")
|
||||
tot_deduped = tot_deduped + json_get_int(rep, "nodes_deduped")
|
||||
tot_edges = tot_edges + json_get_int(rep, "edges_added")
|
||||
files = files + 1
|
||||
}
|
||||
println("FILE " + full)
|
||||
let rep: String = ingest_file(full)
|
||||
tot_created = tot_created + json_get_int(rep, "nodes_created")
|
||||
tot_deduped = tot_deduped + json_get_int(rep, "nodes_deduped")
|
||||
tot_edges = tot_edges + json_get_int(rep, "edges_added")
|
||||
files = files + 1
|
||||
}
|
||||
i = i + 1
|
||||
}
|
||||
@@ -667,11 +610,12 @@ fn ingest_dir(path: String) -> String {
|
||||
",\"edges_accepted\":" + int_to_str(tot_edges) + "}"
|
||||
}
|
||||
|
||||
// ingest a url: fetch, treat body as prose (faithful extraction of what's there)
|
||||
// ingest a url: fetch, hand the body straight to transduce (faithful
|
||||
// extraction of what's there — no interpretation of what it is)
|
||||
fn ingest_url(url: String) -> String {
|
||||
let body: String = http_get(url)
|
||||
if str_eq(body, "") { return "{\"error\":\"empty fetch\",\"url\":" + j_q(url) + "}" }
|
||||
let packed: [String] = transduce_prose(el_list_empty(), el_list_empty(),
|
||||
let packed: [String] = transduce(el_list_empty(), el_list_empty(),
|
||||
body, "url:" + url, "extracted", "public-web",
|
||||
"url:" + url, url)
|
||||
return merge_packed(packed)
|
||||
@@ -686,7 +630,7 @@ fn ingest_llm(query: String) -> String {
|
||||
let resp: String = http_post_json("http://127.0.0.1:11434/api/generate", body)
|
||||
let answer: String = json_get_string(resp, "response")
|
||||
if str_eq(answer, "") { return "{\"error\":\"no model response\"}" }
|
||||
let packed: [String] = transduce_prose(el_list_empty(), el_list_empty(),
|
||||
let packed: [String] = transduce(el_list_empty(), el_list_empty(),
|
||||
answer, "llm:" + model + ":" + query, "candidate-provisional", "guide-provisional",
|
||||
"llm:" + query, "guide answer: " + query)
|
||||
return merge_packed(packed)
|
||||
@@ -728,6 +672,19 @@ fn ingest_stream(path: String) -> String {
|
||||
// SECTION G — ENTRY
|
||||
// ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
// INGEST_KIND selects an ACQUISITION mechanism only — dir/file/url/llm/
|
||||
// stream — i.e. which RPC shape to use to go get the bytes (walk a
|
||||
// directory, open a file, fetch a URL, query an LLM, read a turn-stream).
|
||||
// That is a genuinely unavoidable choice at the process-entry boundary
|
||||
// (nothing about the string "/tmp/x" tells you whether it's a file to read
|
||||
// or a stream to read line-by-line, or distinguishes an LLM query from a
|
||||
// path), so it cannot be dropped the way content-type dispatch was.
|
||||
// It is NOT a content-type flag: it says nothing about what's inside the
|
||||
// bytes once fetched, and none of the five ingest_* functions it selects
|
||||
// among interpret their payload differently by content shape anymore —
|
||||
// they all hand off to the single, format-agnostic transduce(). The old
|
||||
// "structured" value (a caller-declared alias for "file", used only to hint
|
||||
// the now-removed JSON-vs-prose branch) is gone along with that branch.
|
||||
let kind: String = env("INGEST_KIND")
|
||||
let arg: String = env("INGEST_ARG")
|
||||
|
||||
@@ -741,20 +698,16 @@ if str_eq(kind, "dir") {
|
||||
if str_eq(kind, "file") {
|
||||
report = ingest_file(arg)
|
||||
} else {
|
||||
if str_eq(kind, "structured") {
|
||||
report = ingest_file(arg)
|
||||
if str_eq(kind, "url") {
|
||||
report = ingest_url(arg)
|
||||
} else {
|
||||
if str_eq(kind, "url") {
|
||||
report = ingest_url(arg)
|
||||
if str_eq(kind, "llm") {
|
||||
report = ingest_llm(arg)
|
||||
} else {
|
||||
if str_eq(kind, "llm") {
|
||||
report = ingest_llm(arg)
|
||||
if str_eq(kind, "stream") {
|
||||
report = ingest_stream(arg)
|
||||
} else {
|
||||
if str_eq(kind, "stream") {
|
||||
report = ingest_stream(arg)
|
||||
} else {
|
||||
report = "{\"error\":\"unknown INGEST_KIND: " + kind + "\"}"
|
||||
}
|
||||
report = "{\"error\":\"unknown INGEST_KIND: " + kind + "\"}"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+23
-12
@@ -62,34 +62,45 @@ This is where almost all work belongs. El programs are source files that get com
|
||||
|
||||
This is the self-contained C OS-boundary layer. It provides the `__`-prefixed primitives that compiled El programs call: libcurl HTTP, pthreads, filesystem I/O, arena allocation, etc. It is **not generated** — it is maintained by hand.
|
||||
|
||||
The old `el_runtime.c` has been archived to `runtime/legacy/`. The runtime is now native El (`runtime/*.el`). `el_seed.c` replaces `el_runtime.c` as the sole C compilation dependency.
|
||||
The runtime is native El (`runtime/*.el`) over a C OS-boundary. **Status (verified 2026-08-15):** the migration to a seed-only boundary is *in progress, not done*. Two files exist:
|
||||
- `runtime/el_runtime.c` (~860 KB) — **LIVE**. Holds the engram store (`EngramStore engram_global`) plus the `http_*`/`json_*`/`state_*`/`engram_*` impls. It is the authoritative single-file link target for the compiler, and `tools/install.sh` compiles it into `libel.a`. This is where a new C builtin's *implementation* must currently live to be linkable.
|
||||
- `runtime/el_seed.c` — the intended hand-maintained `__`-prefixed seed (thin wrappers over the above). It is compiled alongside `el_runtime.c` by `tools/install.sh`, but does **not** compile standalone yet (see the build-path caveat under "Rebuilding the Compiler").
|
||||
|
||||
**Only edit `el_seed.c` when you genuinely need OS-level access** (raw sockets, GPU calls, new libcurl features). For everything else, write El.
|
||||
**Only edit these when you genuinely need OS-level access** (raw sockets, GPU calls, new libcurl features, a new engram store op). For everything else, write El.
|
||||
|
||||
When you do add a C builtin:
|
||||
1. Add the C function to `el_seed.c`
|
||||
2. Declare it in `el_seed.h`
|
||||
3. Add it to the `builtin_arity` table in `el-compiler/src/codegen.el` (so the compiler knows the arg count)
|
||||
4. Rebuild the elc binary (see below)
|
||||
When you add a C builtin (verbatim-emit recipe — the El name is emitted as the exact C symbol; `builtin_arity` is an arity guard only, not a dispatch table):
|
||||
1. Implement the C function in `el_runtime.c` (and declare it in `el_runtime.h`).
|
||||
2. Add a `__`-prefixed thin wrapper in `el_seed.c` and declare it in `el_seed.h`.
|
||||
3. Add the name to `builtin_arity` in `el-compiler/src/codegen.el` — add **both** the plain and `__`-prefixed spellings.
|
||||
4. Rebuild the elc binary (see below) and confirm the self-host fixpoint is byte-identical.
|
||||
|
||||
Worked example: the `engram_assert_json` (op_assert seam) and `engram_node_full_in`/`engram_connect_in` (purview write-side) primitives added 2026-08-15 follow exactly this recipe.
|
||||
|
||||
---
|
||||
|
||||
## Rebuilding the Compiler
|
||||
|
||||
After changing any `.el` source in `el-compiler/src/`:
|
||||
After changing any `.el` source in `el-compiler/src/` (run from the `lang/` dir):
|
||||
|
||||
```bash
|
||||
cd /Users/will/Development/neuron-technologies/foundation/el
|
||||
# 1. Stage2: current elc compiles the (modified) compiler to C
|
||||
./dist/platform/elc elc-cli.el > elc-new.c
|
||||
# 2. Build the new compiler. The C link target is el_runtime.c — it holds the
|
||||
# engram store + http/json/state impls the compiler output calls. el_runtime.c
|
||||
# self-hosts elc on its own; el_seed.c is the (aspirational) seed layer and does
|
||||
# NOT compile standalone under clang (missing prototypes for the el_runtime.c
|
||||
# symbols it wraps — see caveat below), so link el_runtime.c here.
|
||||
cc -std=c11 -I runtime -lcurl -lpthread \
|
||||
-o dist/platform/elc-new \
|
||||
elc-new.c runtime/el_seed.c
|
||||
# Verify self-hosting:
|
||||
elc-new.c runtime/el_runtime.c
|
||||
# 3. Verify self-hosting FIXPOINT (stage3 == stage2 output, byte-identical):
|
||||
./dist/platform/elc-new elc-cli.el > elc-verify.c
|
||||
diff elc-new.c elc-verify.c # should be identical
|
||||
diff elc-new.c elc-verify.c # must be identical
|
||||
mv dist/platform/elc-new dist/platform/elc
|
||||
```
|
||||
|
||||
> **Build-path caveat (verified 2026-08-15).** `el_seed.c` is the intended hand-maintained OS-boundary seed, but it does **not** compile standalone under modern clang: it wraps ~16 unprefixed `el_runtime.c` symbols (`http_serve`, `json_*`, `state_*`, `http_response`) without prototypes, and clang treats implicit declarations as errors (C99+). The productionised install (`tools/install.sh`) builds `libel.a` from **both** `el_seed.o` + `el_runtime.o` together, which is why linking succeeds there. To make `el_seed.c` build on its own, add prototypes for those symbols (or `#include "el_runtime.h"`, reconciling the `__http_serve` return-type mismatch first). Until then, `el_runtime.c` is the authoritative single-file link target for the compiler.
|
||||
|
||||
After changing `el_seed.c` only (no El source changes), rebuild downstream programs but do NOT need to rebuild the compiler binary itself — the seed is linked at the application level, not the compiler level.
|
||||
|
||||
---
|
||||
|
||||
@@ -2760,12 +2760,18 @@ fn builtin_arity(name: String) -> Int {
|
||||
if str_eq(name, "__engram_neighbors_filtered") { return 3 }
|
||||
if str_eq(name, "__engram_activate") { return 2 }
|
||||
if str_eq(name, "__engram_activate_json") { return 2 }
|
||||
if str_eq(name, "__engram_op_assert_json") { return 2 }
|
||||
if str_eq(name, "__engram_node_full_in") { return 9 }
|
||||
if str_eq(name, "__engram_connect_in") { return 5 }
|
||||
if str_eq(name, "__engram_scan_nodes_json") { return 2 }
|
||||
if str_eq(name, "__engram_edges_json") { return 2 }
|
||||
if str_eq(name, "__generate") { return 1 }
|
||||
// Filesystem
|
||||
if str_eq(name, "fs_read") { return 1 }
|
||||
if str_eq(name, "fs_write") { return 2 }
|
||||
if str_eq(name, "fs_list") { return 1 }
|
||||
if str_eq(name, "fs_size") { return 1 }
|
||||
if str_eq(name, "fs_read_b64_chunk") { return 3 }
|
||||
// JSON
|
||||
if str_eq(name, "json_get") { return 2 }
|
||||
if str_eq(name, "json_parse") { return 1 }
|
||||
@@ -2857,9 +2863,13 @@ fn builtin_arity(name: String) -> Int {
|
||||
if str_eq(name, "engram_get_node_by_label") { return 1 }
|
||||
if str_eq(name, "engram_search_json") { return 2 }
|
||||
if str_eq(name, "engram_scan_nodes_json") { return 2 }
|
||||
if str_eq(name, "engram_edges_json") { return 2 }
|
||||
if str_eq(name, "engram_neighbors_json") { return 3 }
|
||||
if str_eq(name, "engram_activate_json") { return 2 }
|
||||
if str_eq(name, "engram_stats_json") { return 0 }
|
||||
if str_eq(name, "engram_op_assert_json") { return 2 }
|
||||
if str_eq(name, "engram_node_full_in") { return 9 }
|
||||
if str_eq(name, "engram_connect_in") { return 5 }
|
||||
// LLM
|
||||
if str_eq(name, "llm_call") { return 2 }
|
||||
if str_eq(name, "llm_call_system") { return 3 }
|
||||
|
||||
+804
-37
@@ -105,23 +105,14 @@ static void el_arena_track(char* p) {
|
||||
_tl_arena.ptrs[_tl_arena.count++] = p;
|
||||
}
|
||||
|
||||
/* Called by http_worker before dispatching the El handler. */
|
||||
void el_request_start(void) {
|
||||
_tl_arena.count = 0;
|
||||
_tl_arena_active = 1;
|
||||
_tl_fs_read_len = 0; /* never let a previous request's file length */
|
||||
_tl_fs_read_buf = NULL; /* leak into this response's byte accounting */
|
||||
}
|
||||
|
||||
/* Called by http_worker after the El handler returns and the response is sent.
|
||||
* Frees every intermediate string allocated during the request. */
|
||||
void el_request_end(void) {
|
||||
_tl_arena_active = 0;
|
||||
for (size_t i = 0; i < _tl_arena.count; i++) {
|
||||
free(_tl_arena.ptrs[i]);
|
||||
}
|
||||
_tl_arena.count = 0;
|
||||
}
|
||||
/* el_request_start / el_request_end moved to el_seed.c (see its comment at the
|
||||
* definition: "formerly defined in el_runtime.c. Now self-contained in
|
||||
* el_seed.c, delegating to the seed arena."). The copies here were left behind
|
||||
* during that move and made el_seed.o + el_runtime.o fail to link together with
|
||||
* duplicate symbols — which is exactly the link the real product build does.
|
||||
* Declared (not defined) here: el_runtime.c's http_worker still calls them. */
|
||||
void el_request_start(void);
|
||||
void el_request_end(void);
|
||||
|
||||
/* ── Scoped arena for CLI use ─────────────────────────────────────────────── *
|
||||
* CLI programs never call el_request_start/end, so all strdup allocations are
|
||||
@@ -2226,6 +2217,22 @@ el_val_t fs_exists(el_val_t pathv) {
|
||||
return (el_val_t)(stat(path, &st) == 0 ? 1 : 0);
|
||||
}
|
||||
|
||||
/* fs_size — real on-disk byte count of a file, via stat() (not strlen).
|
||||
* Needed alongside fs_read_b64_chunk() below: fs_read()'s el_val_t string
|
||||
* result is NUL-terminated and length-unsafe for arbitrary binary content
|
||||
* (str_len/str_slice fall back to strlen when the buffer isn't a tagged
|
||||
* binary value — see el_input_len), so any caller that needs to walk a
|
||||
* binary file in fixed-size windows (e.g. transduce()'s raw/opaque chunking
|
||||
* of audio bytes) must learn the true length here instead of from the
|
||||
* decoded string. Returns -1 if the path doesn't exist or isn't stat-able. */
|
||||
el_val_t fs_size(el_val_t pathv) {
|
||||
const char* path = EL_CSTR(pathv);
|
||||
if (!path || !*path) return -1;
|
||||
struct stat st;
|
||||
if (stat(path, &st) != 0) return -1;
|
||||
return (el_val_t)st.st_size;
|
||||
}
|
||||
|
||||
/* fs_mkdir — create directory at path with mode 0755, mkdir -p semantics.
|
||||
* Returns 1 if path exists or was created (incl. all parents); 0 on failure.
|
||||
* Walks the path component-by-component so missing intermediate dirs are
|
||||
@@ -6847,10 +6854,16 @@ static float* eg_embed_fetch(const char* text, int32_t* out_dim) {
|
||||
else esc[w++] = (char)c;
|
||||
}
|
||||
esc[w] = '\0';
|
||||
size_t blen = w + strlen(eg_embed_model()) + 64;
|
||||
size_t blen = w + strlen(eg_embed_model()) + 96;
|
||||
char* body = malloc(blen);
|
||||
if (!body) { free(esc); return NULL; }
|
||||
snprintf(body, blen, "{\"model\":\"%s\",\"prompt\":\"%s\"}",
|
||||
/* keep_alive:-1 pins the embed model resident in Ollama indefinitely
|
||||
* (2026-08-15, PR #105 port). Without it the tiny embed model is evicted
|
||||
* whenever a larger generation model loads (unified-memory pressure), so
|
||||
* the NEXT search pays a cold model reload — measured cold reload up to
|
||||
* ~2.2s vs ~0.02-0.05s warm, well inside ENGRAM_EMBED_TIMEOUT_MS but a
|
||||
* real tax on every activate() call that lands cold. Pinning removes it. */
|
||||
snprintf(body, blen, "{\"model\":\"%s\",\"keep_alive\":-1,\"prompt\":\"%s\"}",
|
||||
eg_embed_model(), esc);
|
||||
free(esc);
|
||||
struct curl_slist* h = curl_slist_append(NULL, "Content-Type: application/json");
|
||||
@@ -9508,6 +9521,36 @@ static inline double eg_cosq_at(EngramStore* g, double* cosq, unsigned char* cos
|
||||
return cosq[i];
|
||||
}
|
||||
|
||||
/* ── Beam cap for engram_activate spreading activation (2026-08-15, PR #105
|
||||
* port) ──────────────────────────────────────────────────────────────────
|
||||
* #105 measured the OLD (pre-adjacency-index, pre-qgate, pre-fan-effect)
|
||||
* BFS reaching multi-second/crash territory at depth 2-3 from unbounded
|
||||
* hub-node fan-out. That specific failure mode is already substantially
|
||||
* mitigated here by mechanisms #105's branch predates: the adjacency index
|
||||
* (O(degree) not O(E) per hop), the query-aware qgate (prunes semantically
|
||||
* irrelevant branches), the ACT-R fan-effect correction (dampens popular-
|
||||
* hub over-connectivity), and the 0.02 firing threshold. A beam cap is still
|
||||
* a genuine additional, orthogonal bound: it caps WORST-CASE per-hop
|
||||
* expansion width regardless of how many targets happen to pass the soft
|
||||
* gates above, so it is kept as defense in depth rather than dropped as
|
||||
* redundant.
|
||||
*
|
||||
* Bounds the number of frontier nodes EXPANDED per hop-level (see the
|
||||
* level-batching in the BFS below). Every reached node still gets its
|
||||
* best_bg[]/reached[] recorded and appears in the returned/promoted set —
|
||||
* the cap bounds only how far ASSOCIATIVE SPREAD continues past a level,
|
||||
* never the direct seed matches or the reported result set. Tunable via
|
||||
* ENGRAM_ACTIVATE_BEAM (default 128, matching #105); set very high (e.g.
|
||||
* the node count) to recover the pre-cap unbounded-per-level behaviour. */
|
||||
static int64_t engram_activate_beam(void) {
|
||||
static int64_t v = -1;
|
||||
if (v >= 0) return v;
|
||||
const char* s = getenv("ENGRAM_ACTIVATE_BEAM");
|
||||
int64_t d = 128;
|
||||
if (s && *s) { char* e = NULL; long t = strtol(s, &e, 10); if (e != s && t > 0) d = (int64_t)t; }
|
||||
v = d; return v;
|
||||
}
|
||||
|
||||
el_val_t engram_activate(el_val_t query, el_val_t depth) {
|
||||
EngramStore* g = engram_get();
|
||||
const char* q = EL_CSTR(query);
|
||||
@@ -9552,25 +9595,43 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
|
||||
backfilled++;
|
||||
}
|
||||
}
|
||||
/* Query embedding, cached single-slot: the curiosity loop re-issues the
|
||||
* same 4 rotating phrases, so consecutive identical queries skip the
|
||||
* HTTP round-trip entirely. */
|
||||
static char* _eg_qcache_text = NULL;
|
||||
static float* _eg_qcache_emb = NULL;
|
||||
static int32_t _eg_qcache_dim = 0;
|
||||
/* Query embedding cache (2026-08-15, PR #105 port: single-slot -> direct-
|
||||
* mapped multi-slot). The single-slot cache below this comment's history
|
||||
* only remembered the LAST query, so "the curiosity loop re-issues the
|
||||
* same 4 rotating phrases" only hit when two CONSECUTIVE calls used the
|
||||
* SAME phrase — any rotation among >1 phrase evicted the slot before it
|
||||
* could be reused. #105 measured this as a real cost (a repeated query
|
||||
* costing a full Ollama round-trip whenever a different phrase intervened)
|
||||
* and fixed it with a direct-mapped, FNV-1a-keyed cache sized for the
|
||||
* rotation. Ported here on TOP of the existing cosq/e_eff semantic layer
|
||||
* (this cache only ever supplies q_emb/q_dim into that unchanged
|
||||
* pipeline) rather than replacing it — see the M8/#105 reconciliation
|
||||
* note above eg_cosq_at. ENGRAM_QCACHE_SIZE must be a power of two (mask
|
||||
* indexing below). Full strcmp on lookup rejects hash collisions; each
|
||||
* slot owns its `text`/`vec` and is freed on eviction, matching the old
|
||||
* single-slot free/replace contract — q_emb below still points at cache-
|
||||
* owned memory the caller must NOT free, just as before. */
|
||||
#define ENGRAM_QCACHE_SIZE 1024
|
||||
typedef struct { char* text; uint64_t hash; float* vec; int32_t dim; } EgQCacheEntry;
|
||||
static EgQCacheEntry _eg_qcache[ENGRAM_QCACHE_SIZE];
|
||||
float* q_emb = NULL;
|
||||
int32_t q_dim = 0;
|
||||
if (_eg_qcache_text && strcmp(_eg_qcache_text, q) == 0) {
|
||||
q_emb = _eg_qcache_emb; q_dim = _eg_qcache_dim;
|
||||
} else {
|
||||
int32_t d = 0;
|
||||
float* v = eg_embed_fetch(q, &d);
|
||||
if (v) {
|
||||
free(_eg_qcache_text); free(_eg_qcache_emb);
|
||||
_eg_qcache_text = strdup(q);
|
||||
_eg_qcache_emb = v;
|
||||
_eg_qcache_dim = d;
|
||||
q_emb = v; q_dim = d;
|
||||
{
|
||||
uint64_t qh = engram_id_hash(q);
|
||||
EgQCacheEntry* slot = &_eg_qcache[qh & (ENGRAM_QCACHE_SIZE - 1)];
|
||||
if (slot->vec && slot->hash == qh && slot->text && strcmp(slot->text, q) == 0) {
|
||||
q_emb = slot->vec; q_dim = slot->dim;
|
||||
} else {
|
||||
int32_t d = 0;
|
||||
float* v = eg_embed_fetch(q, &d);
|
||||
if (v) {
|
||||
free(slot->text); free(slot->vec); /* evict prior occupant */
|
||||
slot->text = strdup(q);
|
||||
slot->hash = qh;
|
||||
slot->vec = v;
|
||||
slot->dim = d;
|
||||
q_emb = v; q_dim = d;
|
||||
}
|
||||
}
|
||||
}
|
||||
/* ── Context centroid fold-in (2026-07-29) ──────────────────────────
|
||||
@@ -9997,8 +10058,45 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
|
||||
const double FAN_DREF = (g->adj_connected > 0)
|
||||
? (2.0 * (double)g->edge_count / (double)g->adj_connected) : 0.0;
|
||||
_eg_act_fan_dref = FAN_DREF;
|
||||
const int64_t activate_beam = engram_activate_beam();
|
||||
while (fhead < ftail) {
|
||||
Frontier f = fr[fhead++];
|
||||
/* Level-batch (2026-08-15, PR #105 port): entries sharing .hops are
|
||||
* always contiguous — hop k+1 entries are appended only while
|
||||
* processing hop k, strictly after the current ftail, so they form one
|
||||
* block right after hop k's block (see engram_activate_beam's comment
|
||||
* for why this holds even with the improve-and-re-enqueue behavior
|
||||
* below). Find this level's extent, then beam-select which of it
|
||||
* EXPANDS; every entry in the level still gets recorded via
|
||||
* reached[]/best_bg[] regardless (that happened when it was enqueued,
|
||||
* one level up) — the cap bounds propagation width only. */
|
||||
int64_t level_hops = fr[fhead].hops;
|
||||
int64_t level_start = fhead;
|
||||
int64_t level_end = fhead;
|
||||
while (level_end < ftail && fr[level_end].hops == level_hops) level_end++;
|
||||
int64_t level_n = level_end - level_start;
|
||||
unsigned char* expand = NULL;
|
||||
if (level_n > activate_beam) {
|
||||
expand = calloc((size_t)level_n, 1);
|
||||
if (expand) {
|
||||
/* Partial selection: mark the top-`activate_beam` entries by
|
||||
* .act. O(beam*level_n) — beam is the small tunable. */
|
||||
for (int64_t bsel = 0; bsel < activate_beam; bsel++) {
|
||||
int64_t best = -1;
|
||||
for (int64_t k = 0; k < level_n; k++) {
|
||||
if (expand[k]) continue;
|
||||
if (best < 0 || fr[level_start+k].act > fr[level_start+best].act)
|
||||
best = k;
|
||||
}
|
||||
if (best < 0) break;
|
||||
expand[best] = 1;
|
||||
}
|
||||
}
|
||||
/* OOM on the selection map: expand stays NULL -> this level runs
|
||||
* unbounded, same as if beam were disabled. Never silently wrong. */
|
||||
}
|
||||
for (int64_t lk = level_start; lk < level_end; lk++) {
|
||||
if (expand && !expand[lk - level_start]) continue;
|
||||
Frontier f = fr[lk];
|
||||
if (f.hops >= max_depth) continue;
|
||||
int64_t cur = f.idx;
|
||||
int64_t new_hops = f.hops + 1;
|
||||
@@ -10130,6 +10228,9 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
free(expand);
|
||||
fhead = level_end;
|
||||
}
|
||||
/* Persist layer-1 background_activation to node store. */
|
||||
for (int64_t i = 0; i < g->node_count; i++) {
|
||||
@@ -13979,6 +14080,59 @@ el_val_t engram_wm_top_json(el_val_t n_v) {
|
||||
return el_wrap_str(b.buf);
|
||||
}
|
||||
|
||||
/* op_assert seam (realizer promotion, bl-53/#57).
|
||||
* Gathers the grounded ASSERTION ENVELOPE for a subject node —
|
||||
* { "subject": <node|null>, "grounding": [ {node,edge,hops}... ] }
|
||||
* i.e. the self-geometry a realizer renders as faithful first-person text.
|
||||
* Read-only: realization (geometry->text) stays in the faculty/realizer;
|
||||
* this native primitive produces its structured input from proven paths
|
||||
* (engram_emit_node_json + engram_neighbors_json). arity 2 (node_id, depth). */
|
||||
el_val_t engram_op_assert_json(el_val_t node_id, el_val_t depth) {
|
||||
const char* sid = EL_CSTR(node_id);
|
||||
JsonBuf b; jb_init(&b);
|
||||
jb_puts(&b, "{\"subject\":");
|
||||
EngramNode* n = (sid && *sid) ? engram_find_node(sid) : NULL;
|
||||
if (n) engram_emit_node_json(&b, n, 0); else jb_puts(&b, "null");
|
||||
jb_puts(&b, ",\"grounding\":");
|
||||
el_val_t nb = engram_neighbors_json(node_id, depth, EL_STR("both"));
|
||||
const char* nbs = EL_CSTR(nb);
|
||||
jb_puts(&b, (nbs && *nbs) ? nbs : "[]");
|
||||
jb_putc(&b, '}');
|
||||
return el_wrap_str(b.buf);
|
||||
}
|
||||
|
||||
/* Parametric mutation (purview write-side bounding, keystone 56ecbec6).
|
||||
* The mutation verbs travel with a TARGET MANIFOLD (purview) instead of the
|
||||
* implicit global singleton. purview==0 (EL_NULL) is the DEGENERATE/DEFAULT
|
||||
* case: G = live, behaviour identical to the base op. A non-zero purview is a
|
||||
* bounded target that the engine cannot yet resolve (multi-manifold store is a
|
||||
* promotion item), so we REFUSE rather than silently mutate the live set —
|
||||
* write-side bounding must never leak into G=live. */
|
||||
el_val_t engram_node_full_in(el_val_t purview,
|
||||
el_val_t content, el_val_t node_type, el_val_t label,
|
||||
el_val_t salience, el_val_t importance, el_val_t confidence,
|
||||
el_val_t tier, el_val_t tags) {
|
||||
if (purview == 0) {
|
||||
return engram_node_full(content, node_type, label, salience, importance,
|
||||
confidence, tier, tags);
|
||||
}
|
||||
fprintf(stderr, "[engram] purview write-side not yet resolvable (G != live); "
|
||||
"refusing to append to live store (purview=%lld)\n",
|
||||
(long long)purview);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
void engram_connect_in(el_val_t purview,
|
||||
el_val_t from_id, el_val_t to_id, el_val_t weight, el_val_t relation) {
|
||||
if (purview == 0) {
|
||||
engram_connect(from_id, to_id, weight, relation);
|
||||
return;
|
||||
}
|
||||
fprintf(stderr, "[engram] purview write-side not yet resolvable (G != live); "
|
||||
"refusing to connect in live store (purview=%lld)\n",
|
||||
(long long)purview);
|
||||
}
|
||||
|
||||
el_val_t engram_stats_json(void) {
|
||||
EngramStore* g = engram_get();
|
||||
/* embedded_count: how far the lazy backfill has progressed. The single
|
||||
@@ -16375,6 +16529,41 @@ el_val_t el_base64_encode_n(const unsigned char* data, size_t len, int url_safe)
|
||||
return el_wrap_str(out);
|
||||
}
|
||||
|
||||
/* fs_read_b64_chunk — binary-safe windowed file read: read up to `length`
|
||||
* raw bytes starting at byte `offset` from `path` and return them base64-
|
||||
* encoded (RFC 4648, standard alphabet, padded). The raw bytes are read into
|
||||
* a local C buffer and base64-encoded directly here — they never pass
|
||||
* through an el_val_t string as raw bytes, so embedded NUL bytes (common in
|
||||
* real PCM audio) never hit a strlen()-based code path. This mirrors the
|
||||
* existing llm_vision() image-attachment path (read file -> base64 in C ->
|
||||
* hand back a plain-ASCII string) and the http_*_to_file() rationale above:
|
||||
* bypass the string wrapper entirely for the part that must stay binary.
|
||||
*
|
||||
* Returns "" if the path can't be opened, offset is negative or past EOF,
|
||||
* or length <= 0. A short final chunk (less than `length` bytes remaining)
|
||||
* is returned truncated to what's actually on disk — never padded/invented. */
|
||||
el_val_t fs_read_b64_chunk(el_val_t pathv, el_val_t offsetv, el_val_t lengthv) {
|
||||
const char* path = EL_CSTR(pathv);
|
||||
int64_t offset = (int64_t)offsetv;
|
||||
int64_t length = (int64_t)lengthv;
|
||||
if (!path || !*path || offset < 0 || length <= 0) return el_wrap_str(el_strdup(""));
|
||||
FILE* f = fopen(path, "rb");
|
||||
if (!f) return el_wrap_str(el_strdup(""));
|
||||
fseek(f, 0, SEEK_END);
|
||||
long sz = ftell(f);
|
||||
if (sz < 0 || offset >= sz) { fclose(f); return el_wrap_str(el_strdup("")); }
|
||||
fseek(f, (long)offset, SEEK_SET);
|
||||
long remain = sz - (long)offset;
|
||||
size_t want = (size_t)(((int64_t)remain < length) ? remain : length);
|
||||
unsigned char* buf = malloc(want > 0 ? want : 1);
|
||||
if (!buf) { fclose(f); return el_wrap_str(el_strdup("")); }
|
||||
size_t got = fread(buf, 1, want, f);
|
||||
fclose(f);
|
||||
el_val_t out = el_base64_encode_n(buf, got, /*url_safe=*/0);
|
||||
free(buf);
|
||||
return out;
|
||||
}
|
||||
|
||||
/* Decode either alphabet — accepts both '+/' and '-_' transparently, and
|
||||
* tolerates missing padding (which JWTs typically omit). Whitespace is
|
||||
* skipped for robustness. Invalid characters cause the decode to stop and
|
||||
@@ -17606,3 +17795,581 @@ el_val_t emit_event(el_val_t name_v, el_val_t duration_ms_v) {
|
||||
return trace_span_end(h);
|
||||
}
|
||||
|
||||
|
||||
/* ── DHARMA runtime additions ────────────────────────────────────────────────
|
||||
*
|
||||
* Functions required by the dharma registry service. Added here so the
|
||||
* released el_runtime.c includes them without requiring dharma to bundle
|
||||
* its own stubs.
|
||||
*
|
||||
* Functions added:
|
||||
* list_len — alias for el_list_len (used in handlers.el)
|
||||
* list_get — alias for el_list_get (used in handlers.el)
|
||||
* json_array_push — append a pre-encoded JSON element to a JSON array string
|
||||
* now_millis — milliseconds since Unix epoch (alias for time_now)
|
||||
* unix_timestamp_ms — same as now_millis (alias)
|
||||
* time_now_ms — same as now_millis (alias)
|
||||
* log_info — stderr structured log at INFO level
|
||||
* log_warn — stderr structured log at WARN level
|
||||
* config — reads a config value from the environment
|
||||
* http_patch — HTTP PATCH with JSON Content-Type
|
||||
* http_post_engram — HTTP POST with optional X-API-Key header
|
||||
* http_get_engram — HTTP GET with optional X-API-Key header
|
||||
* str_to_bytes — encode a string as a JSON array of byte values
|
||||
* bytes_to_str — decode a JSON array of byte values back to a string
|
||||
* hash_sha256 — SHA-256 hex digest of a string
|
||||
*/
|
||||
|
||||
/* list_len — return the number of elements in a list. */
|
||||
el_val_t list_len(el_val_t list) {
|
||||
return el_list_len(list);
|
||||
}
|
||||
|
||||
/* list_get — return the element at index i in a list. */
|
||||
el_val_t list_get(el_val_t list, el_val_t index) {
|
||||
return el_list_get(list, index);
|
||||
}
|
||||
|
||||
/* json_array_push — append element (a pre-encoded JSON fragment, e.g. "\"foo\""
|
||||
* or "42") to the JSON array string arr. Returns a new JSON array string.
|
||||
* Example: json_array_push("[]", "\"alice\"") -> "[\"alice\"]"
|
||||
* json_array_push("[\"alice\"]", "\"bob\"") -> "[\"alice\",\"bob\"]" */
|
||||
el_val_t json_array_push(el_val_t arr_v, el_val_t elem_v) {
|
||||
const char* arr = EL_CSTR(arr_v);
|
||||
const char* elem = EL_CSTR(elem_v);
|
||||
if (!arr || !*arr) arr = "[]";
|
||||
if (!elem || !*elem) elem = "null";
|
||||
|
||||
/* Trim whitespace, find the closing ']'. */
|
||||
const char* p = arr;
|
||||
while (*p == ' ' || *p == '\t' || *p == '\n' || *p == '\r') p++;
|
||||
if (*p != '[') {
|
||||
/* Not an array — return a single-element array. */
|
||||
size_t n = strlen(elem) + 4;
|
||||
char* out = el_strbuf(n);
|
||||
snprintf(out, n, "[%s]", elem);
|
||||
return el_wrap_str(out);
|
||||
}
|
||||
size_t arr_len = strlen(arr);
|
||||
size_t elem_len = strlen(elem);
|
||||
|
||||
/* Walk from the end to find the matching ']'. */
|
||||
const char* end = arr + arr_len - 1;
|
||||
while (end > p && (*end == ' ' || *end == '\t' || *end == '\n' || *end == '\r')) end--;
|
||||
if (*end != ']') {
|
||||
/* Malformed — wrap elem in a new array. */
|
||||
size_t n = elem_len + 4;
|
||||
char* out = el_strbuf(n);
|
||||
snprintf(out, n, "[%s]", elem);
|
||||
return el_wrap_str(out);
|
||||
}
|
||||
|
||||
/* Content between '[' and ']'. */
|
||||
const char* inner_start = p + 1;
|
||||
const char* inner_end = end; /* points AT ']' */
|
||||
/* Check if the array is empty (only whitespace between brackets). */
|
||||
const char* q = inner_start;
|
||||
while (q < inner_end && (*q == ' ' || *q == '\t' || *q == '\n' || *q == '\r')) q++;
|
||||
int empty = (q == inner_end);
|
||||
|
||||
/* Build: prefix + (comma if non-empty) + elem + "]" */
|
||||
size_t prefix_len = (size_t)(inner_end - arr); /* up to but not including ']' */
|
||||
size_t sep_len = empty ? 0 : 1; /* "," if non-empty */
|
||||
size_t out_len = prefix_len + sep_len + elem_len + 2; /* +"]" + NUL */
|
||||
char* out = el_strbuf(out_len);
|
||||
memcpy(out, arr, prefix_len);
|
||||
if (!empty) out[prefix_len] = ',';
|
||||
memcpy(out + prefix_len + sep_len, elem, elem_len);
|
||||
out[prefix_len + sep_len + elem_len] = ']';
|
||||
out[prefix_len + sep_len + elem_len + 1] = '\0';
|
||||
return el_wrap_str(out);
|
||||
}
|
||||
|
||||
/* now_millis — milliseconds since Unix epoch. */
|
||||
el_val_t now_millis(void) {
|
||||
return time_now();
|
||||
}
|
||||
|
||||
/* unix_timestamp_ms — same as now_millis. */
|
||||
el_val_t unix_timestamp_ms(void) {
|
||||
return time_now();
|
||||
}
|
||||
|
||||
/* time_now_ms — same as now_millis. */
|
||||
el_val_t time_now_ms(void) {
|
||||
return time_now();
|
||||
}
|
||||
|
||||
/* log_info — write a structured [INFO] line to stderr. */
|
||||
void log_info(el_val_t msg_v) {
|
||||
const char* msg = EL_CSTR(msg_v);
|
||||
fprintf(stderr, "[INFO] %s\n", msg ? msg : "");
|
||||
}
|
||||
|
||||
/* log_warn — write a structured [WARN] line to stderr. */
|
||||
void log_warn(el_val_t msg_v) {
|
||||
const char* msg = EL_CSTR(msg_v);
|
||||
fprintf(stderr, "[WARN] %s\n", msg ? msg : "");
|
||||
}
|
||||
|
||||
/* config — read a configuration value from the environment.
|
||||
* Returns "" if the variable is not set (same as __env_get). */
|
||||
el_val_t config(el_val_t key_v) {
|
||||
const char* key = EL_CSTR(key_v);
|
||||
if (!key || !*key) return EL_STR("");
|
||||
const char* val = getenv(key);
|
||||
if (!val) return EL_STR("");
|
||||
return el_wrap_str(el_strdup(val));
|
||||
}
|
||||
|
||||
#if !defined(_WIN32) || defined(HAVE_CURL)
|
||||
/* http_patch — HTTP PATCH request with Content-Type: application/json.
|
||||
* Returns the response body (same error convention as http_post_json). */
|
||||
el_val_t http_patch(el_val_t url_v, el_val_t body_v) {
|
||||
const char* url = EL_CSTR(url_v);
|
||||
const char* body = EL_CSTR(body_v);
|
||||
if (!url || !*url) return http_error_json("empty url");
|
||||
CURL* c = curl_easy_init();
|
||||
if (!c) return http_error_json("curl_easy_init failed");
|
||||
HttpBuf rb; httpbuf_init(&rb);
|
||||
char errbuf[CURL_ERROR_SIZE]; errbuf[0] = '\0';
|
||||
struct curl_slist* h = NULL;
|
||||
h = curl_slist_append(h, "Content-Type: application/json");
|
||||
curl_easy_setopt(c, CURLOPT_URL, url);
|
||||
curl_easy_setopt(c, CURLOPT_CUSTOMREQUEST, "PATCH");
|
||||
curl_easy_setopt(c, CURLOPT_POSTFIELDS, body ? body : "");
|
||||
curl_easy_setopt(c, CURLOPT_POSTFIELDSIZE, (long)(body ? strlen(body) : 0));
|
||||
curl_easy_setopt(c, CURLOPT_HTTPHEADER, h);
|
||||
curl_easy_setopt(c, CURLOPT_WRITEFUNCTION, http_write_cb);
|
||||
curl_easy_setopt(c, CURLOPT_WRITEDATA, &rb);
|
||||
curl_easy_setopt(c, CURLOPT_FOLLOWLOCATION, 1L);
|
||||
curl_easy_setopt(c, CURLOPT_TIMEOUT_MS, el_http_timeout_ms());
|
||||
curl_easy_setopt(c, CURLOPT_NOSIGNAL, 1L);
|
||||
curl_easy_setopt(c, CURLOPT_ERRORBUFFER, errbuf);
|
||||
curl_easy_setopt(c, CURLOPT_USERAGENT, "el-runtime/1.0");
|
||||
CURLcode rc = curl_easy_perform(c);
|
||||
curl_slist_free_all(h);
|
||||
curl_easy_cleanup(c);
|
||||
if (rc != CURLE_OK) {
|
||||
free(rb.data);
|
||||
const char* m = errbuf[0] ? errbuf : curl_easy_strerror(rc);
|
||||
return http_error_json(m);
|
||||
}
|
||||
return el_wrap_str(rb.data);
|
||||
}
|
||||
|
||||
/* http_post_engram — HTTP POST with optional X-API-Key header.
|
||||
* If key is "" no authentication header is sent. */
|
||||
el_val_t http_post_engram(el_val_t url_v, el_val_t key_v, el_val_t body_v) {
|
||||
const char* url = EL_CSTR(url_v);
|
||||
const char* key = EL_CSTR(key_v);
|
||||
const char* body = EL_CSTR(body_v);
|
||||
if (!url || !*url) return http_error_json("empty url");
|
||||
CURL* c = curl_easy_init();
|
||||
if (!c) return http_error_json("curl_easy_init failed");
|
||||
HttpBuf rb; httpbuf_init(&rb);
|
||||
char errbuf[CURL_ERROR_SIZE]; errbuf[0] = '\0';
|
||||
struct curl_slist* h = NULL;
|
||||
h = curl_slist_append(h, "Content-Type: application/json");
|
||||
if (key && *key) {
|
||||
size_t n = strlen(key) + 32;
|
||||
char* hdr = malloc(n);
|
||||
snprintf(hdr, n, "X-API-Key: %s", key);
|
||||
h = curl_slist_append(h, hdr);
|
||||
free(hdr);
|
||||
}
|
||||
curl_easy_setopt(c, CURLOPT_URL, url);
|
||||
curl_easy_setopt(c, CURLOPT_POST, 1L);
|
||||
curl_easy_setopt(c, CURLOPT_POSTFIELDS, body ? body : "");
|
||||
curl_easy_setopt(c, CURLOPT_POSTFIELDSIZE, (long)(body ? strlen(body) : 0));
|
||||
curl_easy_setopt(c, CURLOPT_HTTPHEADER, h);
|
||||
curl_easy_setopt(c, CURLOPT_WRITEFUNCTION, http_write_cb);
|
||||
curl_easy_setopt(c, CURLOPT_WRITEDATA, &rb);
|
||||
curl_easy_setopt(c, CURLOPT_FOLLOWLOCATION, 1L);
|
||||
curl_easy_setopt(c, CURLOPT_TIMEOUT_MS, el_http_timeout_ms());
|
||||
curl_easy_setopt(c, CURLOPT_NOSIGNAL, 1L);
|
||||
curl_easy_setopt(c, CURLOPT_ERRORBUFFER, errbuf);
|
||||
curl_easy_setopt(c, CURLOPT_USERAGENT, "el-runtime/1.0");
|
||||
CURLcode rc = curl_easy_perform(c);
|
||||
curl_slist_free_all(h);
|
||||
curl_easy_cleanup(c);
|
||||
if (rc != CURLE_OK) {
|
||||
free(rb.data);
|
||||
const char* m = errbuf[0] ? errbuf : curl_easy_strerror(rc);
|
||||
return http_error_json(m);
|
||||
}
|
||||
return el_wrap_str(rb.data);
|
||||
}
|
||||
|
||||
/* http_get_engram — HTTP GET with optional X-API-Key header. */
|
||||
el_val_t http_get_engram(el_val_t url_v, el_val_t key_v) {
|
||||
const char* url = EL_CSTR(url_v);
|
||||
const char* key = EL_CSTR(key_v);
|
||||
if (!url || !*url) return http_error_json("empty url");
|
||||
CURL* c = curl_easy_init();
|
||||
if (!c) return http_error_json("curl_easy_init failed");
|
||||
HttpBuf rb; httpbuf_init(&rb);
|
||||
char errbuf[CURL_ERROR_SIZE]; errbuf[0] = '\0';
|
||||
struct curl_slist* h = NULL;
|
||||
if (key && *key) {
|
||||
size_t n = strlen(key) + 32;
|
||||
char* hdr = malloc(n);
|
||||
snprintf(hdr, n, "X-API-Key: %s", key);
|
||||
h = curl_slist_append(h, hdr);
|
||||
free(hdr);
|
||||
}
|
||||
curl_easy_setopt(c, CURLOPT_URL, url);
|
||||
curl_easy_setopt(c, CURLOPT_HTTPGET, 1L);
|
||||
if (h) curl_easy_setopt(c, CURLOPT_HTTPHEADER, h);
|
||||
curl_easy_setopt(c, CURLOPT_WRITEFUNCTION, http_write_cb);
|
||||
curl_easy_setopt(c, CURLOPT_WRITEDATA, &rb);
|
||||
curl_easy_setopt(c, CURLOPT_FOLLOWLOCATION, 1L);
|
||||
curl_easy_setopt(c, CURLOPT_TIMEOUT_MS, el_http_timeout_ms());
|
||||
curl_easy_setopt(c, CURLOPT_NOSIGNAL, 1L);
|
||||
curl_easy_setopt(c, CURLOPT_ERRORBUFFER, errbuf);
|
||||
curl_easy_setopt(c, CURLOPT_USERAGENT, "el-runtime/1.0");
|
||||
CURLcode rc = curl_easy_perform(c);
|
||||
if (h) curl_slist_free_all(h);
|
||||
curl_easy_cleanup(c);
|
||||
if (rc != CURLE_OK) {
|
||||
free(rb.data);
|
||||
const char* m = errbuf[0] ? errbuf : curl_easy_strerror(rc);
|
||||
return http_error_json(m);
|
||||
}
|
||||
return el_wrap_str(rb.data);
|
||||
}
|
||||
#endif /* HAVE_CURL */
|
||||
|
||||
/* str_to_bytes — encode a string as a JSON array of unsigned byte values.
|
||||
* "hello" -> "[104,101,108,108,111]"
|
||||
* Used by db.el to store binary content in Engram JSON nodes. */
|
||||
el_val_t str_to_bytes(el_val_t sv) {
|
||||
const char* s = EL_CSTR(sv);
|
||||
if (!s || !*s) return el_wrap_str(el_strdup("[]"));
|
||||
size_t n = strlen(s);
|
||||
/* Worst case: each byte is 3 digits + comma = 4 chars, plus "[]" + NUL. */
|
||||
char* out = el_strbuf(n * 4 + 3);
|
||||
size_t pos = 0;
|
||||
out[pos++] = '[';
|
||||
for (size_t i = 0; i < n; i++) {
|
||||
unsigned char b = (unsigned char)s[i];
|
||||
if (i > 0) out[pos++] = ',';
|
||||
/* Write decimal representation of b. */
|
||||
if (b >= 100) {
|
||||
out[pos++] = (char)('0' + b / 100);
|
||||
out[pos++] = (char)('0' + (b / 10) % 10);
|
||||
out[pos++] = (char)('0' + b % 10);
|
||||
} else if (b >= 10) {
|
||||
out[pos++] = (char)('0' + b / 10);
|
||||
out[pos++] = (char)('0' + b % 10);
|
||||
} else {
|
||||
out[pos++] = (char)('0' + b);
|
||||
}
|
||||
}
|
||||
out[pos++] = ']';
|
||||
out[pos] = '\0';
|
||||
return el_wrap_str(out);
|
||||
}
|
||||
|
||||
/* bytes_to_str — decode a JSON array of integer byte values back to a string.
|
||||
* "[104,101,108,108,111]" -> "hello"
|
||||
* Inverse of str_to_bytes. */
|
||||
el_val_t bytes_to_str(el_val_t arr_v) {
|
||||
const char* s = EL_CSTR(arr_v);
|
||||
if (!s) return el_wrap_str(el_strdup(""));
|
||||
/* Skip whitespace, expect '['. */
|
||||
while (*s == ' ' || *s == '\t' || *s == '\n' || *s == '\r') s++;
|
||||
if (*s != '[') return el_wrap_str(el_strdup(""));
|
||||
s++;
|
||||
|
||||
/* Count elements to size the output buffer. */
|
||||
int64_t n = (int64_t)json_array_len(arr_v);
|
||||
if (n <= 0) return el_wrap_str(el_strdup(""));
|
||||
|
||||
char* out = el_strbuf((size_t)n + 1);
|
||||
size_t pos = 0;
|
||||
|
||||
/* Walk the array, parse each integer, store as a byte. */
|
||||
while (*s) {
|
||||
while (*s == ' ' || *s == '\t' || *s == '\n' || *s == '\r') s++;
|
||||
if (*s == ']' || *s == '\0') break;
|
||||
/* Parse decimal integer. */
|
||||
char* end_ptr;
|
||||
long v = strtol(s, &end_ptr, 10);
|
||||
if (end_ptr == s) break; /* parse failure */
|
||||
s = end_ptr;
|
||||
if (v >= 0 && v <= 255) out[pos++] = (char)(unsigned char)v;
|
||||
while (*s == ' ' || *s == '\t' || *s == '\n' || *s == '\r') s++;
|
||||
if (*s == ',') { s++; continue; }
|
||||
if (*s == ']' || *s == '\0') break;
|
||||
}
|
||||
out[pos] = '\0';
|
||||
return el_wrap_str(out);
|
||||
}
|
||||
|
||||
/* hash_sha256 — return the SHA-256 hex digest of a string.
|
||||
* Uses the built-in el_sha256_oneshot implementation (no OpenSSL required). */
|
||||
el_val_t hash_sha256(el_val_t sv) {
|
||||
const char* s = EL_CSTR(sv);
|
||||
if (!s) s = "";
|
||||
unsigned char digest[32];
|
||||
el_sha256_oneshot((const unsigned char*)s, strlen(s), digest);
|
||||
return el_hex_encode(digest, 32);
|
||||
}
|
||||
|
||||
|
||||
/* HTTP client aliases — require curl; defined inside #ifdef HAVE_CURL below
|
||||
* with a matching stub in the #ifndef HAVE_CURL block. */
|
||||
#if !defined(_WIN32) || defined(HAVE_CURL)
|
||||
/* __http_do also lives in el_seed.c; marked weak so el_seed.c's definition
|
||||
* wins when both translation units are linked together (the real product build). */
|
||||
__attribute__((weak)) el_val_t __http_do(el_val_t method, el_val_t url, el_val_t body,
|
||||
el_val_t headers_map, el_val_t timeout_ms) {
|
||||
/* timeout_ms is accepted for API compatibility but ignored here;
|
||||
* el_runtime's http_do uses the EL_HTTP_TIMEOUT_MS env var instead. */
|
||||
(void)timeout_ms;
|
||||
struct curl_slist* h = headers_from_map(headers_map);
|
||||
el_val_t r = http_do(EL_CSTR(method), EL_CSTR(url), EL_CSTR(body), h);
|
||||
if (h) curl_slist_free_all(h);
|
||||
return r;
|
||||
}
|
||||
|
||||
/* __http_do_map — same as __http_do but headers_map arg is a JSON-string
|
||||
* rather than an ElMap. Parse it first, then delegate. */
|
||||
el_val_t __http_do_map(el_val_t method, el_val_t url, el_val_t body,
|
||||
el_val_t headers_json, el_val_t timeout_ms) {
|
||||
(void)timeout_ms;
|
||||
/* Build a curl_slist from a JSON object {"Header":"value",...}. */
|
||||
const char* hj = EL_CSTR(headers_json);
|
||||
struct curl_slist* h = NULL;
|
||||
if (hj && *hj && *hj == '{') {
|
||||
/* Walk the JSON pairs with a simple parser reusing json_get_string logic. */
|
||||
/* For correctness we just call the existing json_get iteration path.
|
||||
* We duplicate the key-extraction loop from headers_from_map but driven
|
||||
* by JSON rather than ElMap. Use json_get_raw to iterate is not easy
|
||||
* without knowing keys, so accept the JSON string and build a tmp map. */
|
||||
el_val_t map = json_parse(EL_STR(hj));
|
||||
h = headers_from_map(map);
|
||||
}
|
||||
el_val_t r = http_do(EL_CSTR(method), EL_CSTR(url), EL_CSTR(body), h);
|
||||
if (h) curl_slist_free_all(h);
|
||||
return r;
|
||||
}
|
||||
|
||||
/* __http_do_map_to_file — same as __http_do_map but streams response body
|
||||
* to a local file path rather than returning it as a string. */
|
||||
el_val_t __http_do_map_to_file(el_val_t method, el_val_t url, el_val_t body,
|
||||
el_val_t headers_json, el_val_t output_path) {
|
||||
const char* hj = EL_CSTR(headers_json);
|
||||
struct curl_slist* h = NULL;
|
||||
if (hj && *hj && *hj == '{') {
|
||||
el_val_t map = json_parse(EL_STR(hj));
|
||||
h = headers_from_map(map);
|
||||
}
|
||||
el_val_t r = http_do_to_file(EL_CSTR(method), EL_CSTR(url), EL_CSTR(body),
|
||||
h, EL_CSTR(output_path));
|
||||
if (h) curl_slist_free_all(h);
|
||||
return r;
|
||||
}
|
||||
#endif /* HAVE_CURL */
|
||||
|
||||
#if defined(_WIN32) && !defined(HAVE_CURL)
|
||||
/* ── HAVE_CURL=0 stubs — compile without -lcurl for the elc CLI binary. ───── *
|
||||
* These return a JSON error string so El programs get a clear message if they
|
||||
* call HTTP/LLM functions in a curl-less build. */
|
||||
static el_val_t _no_curl_err(void) {
|
||||
return el_wrap_str(el_strdup("{\"error\":\"not built with HAVE_CURL\"}"));
|
||||
}
|
||||
el_val_t http_get(el_val_t url) { (void)url; return _no_curl_err(); }
|
||||
el_val_t http_post(el_val_t url, el_val_t body) { (void)url; (void)body; return _no_curl_err(); }
|
||||
el_val_t http_post_json(el_val_t url, el_val_t body) { (void)url; (void)body; return _no_curl_err(); }
|
||||
el_val_t http_get_with_headers(el_val_t url, el_val_t h) { (void)url; (void)h; return _no_curl_err(); }
|
||||
el_val_t http_post_with_headers(el_val_t url, el_val_t b, el_val_t h) { (void)url; (void)b; (void)h; return _no_curl_err(); }
|
||||
el_val_t http_post_json_with_headers(el_val_t url, el_val_t h, el_val_t b) { (void)url; (void)h; (void)b; return _no_curl_err(); }
|
||||
el_val_t http_post_form_auth(el_val_t url, el_val_t b, el_val_t a) { (void)url; (void)b; (void)a; return _no_curl_err(); }
|
||||
el_val_t http_delete(el_val_t url) { (void)url; return _no_curl_err(); }
|
||||
el_val_t http_patch(el_val_t url, el_val_t body) { (void)url; (void)body; return _no_curl_err(); }
|
||||
el_val_t http_get_to_file(el_val_t url, el_val_t h, el_val_t p) { (void)url; (void)h; (void)p; return _no_curl_err(); }
|
||||
el_val_t http_post_to_file(el_val_t url, el_val_t b, el_val_t h, el_val_t p) { (void)url; (void)b; (void)h; (void)p; return _no_curl_err(); }
|
||||
el_val_t http_post_engram(el_val_t url, el_val_t k, el_val_t b) { (void)url; (void)k; (void)b; return _no_curl_err(); }
|
||||
el_val_t http_get_engram(el_val_t url, el_val_t k) { (void)url; (void)k; return _no_curl_err(); }
|
||||
el_val_t llm_call(el_val_t m, el_val_t p) { (void)m; (void)p; return _no_curl_err(); }
|
||||
el_val_t llm_call_system(el_val_t m, el_val_t s, el_val_t u) { (void)m; (void)s; (void)u; return _no_curl_err(); }
|
||||
el_val_t llm_call_agentic(el_val_t m, el_val_t s, el_val_t u, el_val_t t) { (void)m; (void)s; (void)u; (void)t; return _no_curl_err(); }
|
||||
el_val_t llm_vision(el_val_t m, el_val_t s, el_val_t p, el_val_t i) { (void)m; (void)s; (void)p; (void)i; return _no_curl_err(); }
|
||||
el_val_t llm_models(void) { return el_list_empty(); }
|
||||
void llm_register_tool(el_val_t n, el_val_t f) { (void)n; (void)f; }
|
||||
/* __ HTTP stubs (no-curl build) */
|
||||
el_val_t __http_do(el_val_t m, el_val_t u, el_val_t b, el_val_t h, el_val_t t) { (void)m; (void)u; (void)b; (void)h; (void)t; return _no_curl_err(); }
|
||||
el_val_t __http_do_map(el_val_t m, el_val_t u, el_val_t b, el_val_t h, el_val_t t) { (void)m; (void)u; (void)b; (void)h; (void)t; return _no_curl_err(); }
|
||||
el_val_t __http_do_map_to_file(el_val_t m, el_val_t u, el_val_t b, el_val_t h, el_val_t p) { (void)m; (void)u; (void)b; (void)h; (void)p; return _no_curl_err(); }
|
||||
#endif /* !HAVE_CURL */
|
||||
|
||||
/* ── Compiler-support builtins ───────────────────────────────────────────────
|
||||
* stdout_to_file / stdout_restore / el_mem_check are called by the El compiler's
|
||||
* own source (compiler.el:472,479,574 and codegen.el:4248) and are registered in
|
||||
* codegen.el's builtin_arity table, but were missing from this runtime — so
|
||||
* rebuilding elc from source failed with three implicit-declaration errors and
|
||||
* the committed elc binary could never be refreshed. The definitions below are
|
||||
* ported verbatim from ui/examples/native-hello-ios/NativeHello/el_runtime.c,
|
||||
* a divergent private copy of this runtime that still carried them.
|
||||
* ──────────────────────────────────────────────────────────────────────────── */
|
||||
|
||||
#include <sys/resource.h>
|
||||
|
||||
static int _el_saved_stdout_fd = -1;
|
||||
|
||||
/* Redirect process stdout to a file; used by the compiler's JS post-processing
|
||||
* pipeline to capture codegen output before piping it onward. */
|
||||
el_val_t stdout_to_file(el_val_t pathv) {
|
||||
const char* path = EL_CSTR(pathv);
|
||||
if (!path) return (el_val_t)(int64_t)-1;
|
||||
fflush(stdout);
|
||||
_el_saved_stdout_fd = dup(STDOUT_FILENO);
|
||||
int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC, 0600);
|
||||
if (fd < 0) return (el_val_t)(int64_t)-1;
|
||||
dup2(fd, STDOUT_FILENO);
|
||||
close(fd);
|
||||
return (el_val_t)(int64_t)0;
|
||||
}
|
||||
|
||||
el_val_t stdout_restore(void) {
|
||||
if (_el_saved_stdout_fd >= 0) {
|
||||
fflush(stdout);
|
||||
dup2(_el_saved_stdout_fd, STDOUT_FILENO);
|
||||
close(_el_saved_stdout_fd);
|
||||
_el_saved_stdout_fd = -1;
|
||||
}
|
||||
return (el_val_t)(int64_t)0;
|
||||
}
|
||||
|
||||
/* el_mem_check — self-terminating memory guard for long-running compiler runs.
|
||||
* Called periodically by the compiler to catch runaway growth before the OS
|
||||
* OOM-killer fires. Limit comes from ELC_MAX_MEM_MB (default 512 MB).
|
||||
* macOS reports ru_maxrss in bytes, Linux in kilobytes; normalised to MB. */
|
||||
el_val_t el_mem_check(void) {
|
||||
long limit_mb = 512;
|
||||
const char* env_val = getenv("ELC_MAX_MEM_MB");
|
||||
if (env_val && *env_val) {
|
||||
long v = atol(env_val);
|
||||
if (v > 0) limit_mb = v;
|
||||
}
|
||||
|
||||
struct rusage ru;
|
||||
if (getrusage(RUSAGE_SELF, &ru) != 0) return 0; /* can't read — skip check */
|
||||
|
||||
long rss_mb;
|
||||
#if defined(__APPLE__) || defined(__MACH__)
|
||||
rss_mb = (long)(ru.ru_maxrss / (1024L * 1024L));
|
||||
#else
|
||||
rss_mb = (long)(ru.ru_maxrss / 1024L);
|
||||
#endif
|
||||
|
||||
if (rss_mb >= limit_mb) {
|
||||
fprintf(stderr, "elc: memory limit exceeded (%ldMB), aborting\n", limit_mb);
|
||||
exit(1);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
/* ── engram_recall_json / cgi_* accessors — restored 2026-08-15 ──────────────
|
||||
*
|
||||
* These existed in the runtime neuron vendored (v1.0.0-20260501) and were lost
|
||||
* when this runtime moved on, so a soul built against current el would fail to
|
||||
* link — and, worse, the naive "fix" of pointing recall at engram_search_json
|
||||
* would have SILENTLY DOWNGRADED the mind's whole retrieval surface from
|
||||
* semantic to lexical, with no error at any layer.
|
||||
*
|
||||
* The lexical/semantic split is a real safety boundary, not redundant naming
|
||||
* (neuron-api.el:613 documents it): engram_search_json stays LEXICAL because
|
||||
* ~40 internal call sites pass a KEY and seven of them DELETE every record
|
||||
* returned — making those semantic would delete fuzzy matches. recall is the
|
||||
* SEMANTIC surface, used by the retrieval routes.
|
||||
*
|
||||
* The old implementation was eg_search_json_impl(q, limit, with_legs=1): embed
|
||||
* the query, cosine over the corpus, then a graph leg from semantic seeds.
|
||||
* In this runtime that is exactly what engram_activate() already does (it
|
||||
* embeds via eg_embed_fetch, scores by cosine, then spreads activation), so
|
||||
* recall delegates to it rather than re-deriving a second semantic path.
|
||||
* Output shape matches engram_search_json — a flat array of node objects via
|
||||
* engram_emit_node_json — because existing callers (memory.el:80,
|
||||
* neuron-api.el:618) parse it as search's shape, not activate's envelope.
|
||||
* ──────────────────────────────────────────────────────────────────────────── */
|
||||
|
||||
el_val_t engram_recall_json(el_val_t query, el_val_t limit) {
|
||||
int64_t lim = (int64_t)limit;
|
||||
if (lim <= 0) lim = 100;
|
||||
|
||||
/* depth 1: the associative leg, one hop out from the semantic seeds. */
|
||||
el_val_t lst = engram_activate(query, (el_val_t)(int64_t)1);
|
||||
ElList* arr = (ElList*)(uintptr_t)lst;
|
||||
|
||||
JsonBuf b; jb_init(&b);
|
||||
jb_putc(&b, '[');
|
||||
int64_t emitted = 0;
|
||||
if (arr) {
|
||||
for (int64_t i = 0; i < arr->length && emitted < lim; i++) {
|
||||
if (!arr->elems[i]) continue;
|
||||
el_val_t node_map = el_map_get(arr->elems[i], EL_STR("node"));
|
||||
el_val_t id_v = el_map_get(node_map, EL_STR("id"));
|
||||
const char* id_s = EL_CSTR(id_v);
|
||||
EngramNode* n = id_s ? engram_find_node(id_s) : NULL;
|
||||
if (!n) continue;
|
||||
if (emitted > 0) jb_putc(&b, ',');
|
||||
engram_emit_node_json(&b, n, 0);
|
||||
emitted++;
|
||||
}
|
||||
}
|
||||
jb_putc(&b, ']');
|
||||
return el_wrap_str(b.buf);
|
||||
}
|
||||
|
||||
/* cgi_* — read-only identity accessors over the process-wide CGI registration
|
||||
* set by cgi_register(). Read-only by design: there is no setter (studio.el:66). */
|
||||
el_val_t cgi_principal(void) { return EL_STR(_el_cgi_principal ? _el_cgi_principal : ""); }
|
||||
el_val_t cgi_network(void) { return EL_STR(_el_cgi_network ? _el_cgi_network : ""); }
|
||||
el_val_t cgi_engram(void) { return EL_STR(_el_cgi_engram ? _el_cgi_engram : ""); }
|
||||
|
||||
/* engram_edges_json(limit, offset) — emit edges straight from the store.
|
||||
*
|
||||
* Replaces a serialize-and-reread round trip that took production down on
|
||||
* 2026-08-15: /api/graph/edges called engram_save() to write the ENTIRE graph
|
||||
* to disk (128 MB) and then fs_read it back, just to answer a read query for
|
||||
* edges. One debug request cost a full snapshot write, a 128 MB read, and the
|
||||
* peak memory to hold it — on top of being O(whole graph) for a bounded slice.
|
||||
* The route's own comment had already named the fix: "Future: add an
|
||||
* engram_edges_json() builtin and drop the file round trip entirely."
|
||||
*
|
||||
* limit <= 0 defaults to 1000 rather than unbounded: this is the endpoint that
|
||||
* fell over, and an unbounded default would preserve the failure mode under a
|
||||
* different name. Pass an explicit limit to page.
|
||||
*/
|
||||
el_val_t engram_edges_json(el_val_t limit, el_val_t offset) {
|
||||
EngramStore* g = engram_get();
|
||||
int64_t lim = (int64_t)limit; if (lim <= 0) lim = 1000;
|
||||
int64_t off = (int64_t)offset; if (off < 0) off = 0;
|
||||
|
||||
JsonBuf b; jb_init(&b);
|
||||
jb_putc(&b, '[');
|
||||
int64_t emitted = 0;
|
||||
char t[192];
|
||||
for (int64_t i = off; i < g->edge_count && emitted < lim; i++) {
|
||||
EngramEdge* e = &g->edges[i];
|
||||
if (emitted > 0) jb_putc(&b, ',');
|
||||
jb_puts(&b, "{\"id\":"); jb_emit_escaped(&b, e->id ? e->id : "");
|
||||
jb_puts(&b, ",\"from_id\":"); jb_emit_escaped(&b, e->from_id ? e->from_id : "");
|
||||
jb_puts(&b, ",\"to_id\":"); jb_emit_escaped(&b, e->to_id ? e->to_id : "");
|
||||
jb_puts(&b, ",\"relation\":"); jb_emit_escaped(&b, e->relation ? e->relation : "");
|
||||
snprintf(t, sizeof t,
|
||||
",\"weight\":%.6g,\"hebb\":%.6g,\"confidence\":%.6g,"
|
||||
"\"created_at\":%lld,\"updated_at\":%lld,\"last_fired\":%lld,"
|
||||
"\"inhibitory\":%d,\"layer_id\":%u}",
|
||||
e->weight, e->hebb, e->confidence,
|
||||
(long long)e->created_at, (long long)e->updated_at,
|
||||
(long long)e->last_fired, e->inhibitory, (unsigned)e->layer_id);
|
||||
jb_puts(&b, t);
|
||||
emitted++;
|
||||
}
|
||||
jb_putc(&b, ']');
|
||||
return el_wrap_str(b.buf);
|
||||
}
|
||||
|
||||
@@ -235,6 +235,20 @@ el_val_t fs_list(el_val_t path);
|
||||
el_val_t fs_exists(el_val_t path);
|
||||
el_val_t fs_mkdir(el_val_t path); /* mkdir -p, mode 0755 */
|
||||
|
||||
/* Real on-disk byte count via stat() — not strlen(). Use this (not
|
||||
* str_len(fs_read(path))) when a file may contain binary content, since
|
||||
* fs_read()'s result truncates at the first embedded NUL under strlen-based
|
||||
* string ops. Returns -1 if the path doesn't exist. */
|
||||
el_val_t fs_size(el_val_t path);
|
||||
|
||||
/* Binary-safe windowed read: read up to `length` bytes starting at byte
|
||||
* `offset` from `path` and return them base64-encoded. Bytes are read and
|
||||
* encoded in C without ever passing through an el_val_t string as raw
|
||||
* bytes, so embedded NULs (routine in PCM audio) can't truncate the result.
|
||||
* Returns "" on any failure or when offset is past EOF; a final short
|
||||
* window returns only the bytes that actually exist on disk. */
|
||||
el_val_t fs_read_b64_chunk(el_val_t path, el_val_t offset, el_val_t length);
|
||||
|
||||
/* Length-explicit binary write. `length` is an Int (el_val_t holding the
|
||||
* byte count). The caller knows the length from context — typically because
|
||||
* `bytes` came from base64_decode (which produces a magic-tagged binary
|
||||
@@ -261,6 +275,10 @@ el_val_t json_set(el_val_t json_str, el_val_t key, el_val_t value);
|
||||
el_val_t json_array_len(el_val_t json_str);
|
||||
el_val_t json_array_get(el_val_t json_str, el_val_t index);
|
||||
el_val_t json_array_get_string(el_val_t json_str, el_val_t index);
|
||||
el_val_t json_escape_string(el_val_t sv);
|
||||
el_val_t json_build_object(el_val_t kvs);
|
||||
el_val_t json_build_array(el_val_t items);
|
||||
el_val_t json_array_push(el_val_t arr_v, el_val_t elem_v); /* defined in el_runtime.c */
|
||||
|
||||
/* ── Time ────────────────────────────────────────────────────────────────── */
|
||||
|
||||
@@ -287,6 +305,8 @@ el_val_t time_diff(el_val_t ts1, el_val_t ts2, el_val_t unit);
|
||||
|
||||
el_val_t el_now_instant(void);
|
||||
el_val_t now(void);
|
||||
el_val_t now_millis(void); /* wall-clock milliseconds (defined in el_runtime.c) */
|
||||
el_val_t now_ns(void); /* wall-clock nanoseconds (defined in el_runtime.c) */
|
||||
el_val_t unix_seconds(el_val_t n);
|
||||
el_val_t unix_millis(el_val_t n);
|
||||
el_val_t instant_from_iso8601(el_val_t s);
|
||||
@@ -698,6 +718,14 @@ el_val_t engram_label_df(el_val_t term);
|
||||
el_val_t engram_salient_term(el_val_t node_id, el_val_t max_df,
|
||||
el_val_t min_df, el_val_t tabu);
|
||||
el_val_t engram_embed_backfill(el_val_t count);
|
||||
/* op_assert seam: grounded assertion envelope {subject,grounding} for the realizer. */
|
||||
el_val_t engram_op_assert_json(el_val_t node_id, el_val_t depth);
|
||||
/* Parametric mutation (purview write-side): purview==0 => G=live (default), else refuse. */
|
||||
el_val_t engram_node_full_in(el_val_t purview, el_val_t content, el_val_t node_type, el_val_t label,
|
||||
el_val_t salience, el_val_t importance, el_val_t confidence,
|
||||
el_val_t tier, el_val_t tags);
|
||||
void engram_connect_in(el_val_t purview, el_val_t from_id, el_val_t to_id,
|
||||
el_val_t weight, el_val_t relation);
|
||||
el_val_t engram_list_layers_json(void);
|
||||
/* Working memory introspection — count, mean weight, and top-N snapshot.
|
||||
* Ported from runtime on 2026-06-30 self-review. */
|
||||
@@ -878,6 +906,130 @@ el_val_t trace_span_start(el_val_t name);
|
||||
el_val_t trace_span_end(el_val_t span_handle);
|
||||
el_val_t emit_event(el_val_t name, el_val_t duration_ms);
|
||||
|
||||
el_val_t __thread_create(el_val_t fn_name_v, el_val_t arg_v);
|
||||
el_val_t __thread_join(el_val_t tid_v);
|
||||
|
||||
/* Mutex + channel seed primitives (defined in el_runtime.c). Declared here so
|
||||
* that compiled El programs which use runtime/thread.el's with_mutex helper or
|
||||
* runtime/channel.el's Go-style channels see real prototypes instead of an
|
||||
* implicit int-return declaration (which the C11 ABI mis-truncates el_val_t). */
|
||||
el_val_t __mutex_new(void);
|
||||
void __mutex_lock(el_val_t m_v);
|
||||
void __mutex_unlock(el_val_t m_v);
|
||||
el_val_t __channel_new(el_val_t capacity_v);
|
||||
el_val_t __channel_send(el_val_t ch_v, el_val_t msg_v);
|
||||
el_val_t __channel_recv(el_val_t ch_v);
|
||||
el_val_t __channel_try_recv(el_val_t ch_v);
|
||||
el_val_t __channel_close(el_val_t ch_v);
|
||||
|
||||
/* ── __ prefixed aliases (self-hosting compiler ABI) ─────────────────────────
|
||||
* The El self-hosting compiler emits calls to __-prefixed names. These are
|
||||
* forwarding wrappers around the existing el_runtime functions above. */
|
||||
|
||||
/* I/O */
|
||||
el_val_t __println(el_val_t s);
|
||||
el_val_t __print(el_val_t s);
|
||||
el_val_t __readline(void);
|
||||
|
||||
/* String */
|
||||
el_val_t __int_to_str(el_val_t n);
|
||||
el_val_t __str_to_int(el_val_t s);
|
||||
el_val_t __float_to_str(el_val_t f);
|
||||
el_val_t __str_to_float(el_val_t s);
|
||||
el_val_t __str_len(el_val_t s);
|
||||
el_val_t __str_char_at(el_val_t s, el_val_t i);
|
||||
el_val_t __str_cmp(el_val_t a, el_val_t b);
|
||||
el_val_t __str_ncmp(el_val_t a, el_val_t b, el_val_t n);
|
||||
el_val_t __str_concat_raw(el_val_t a, el_val_t b);
|
||||
el_val_t __str_slice_raw(el_val_t s, el_val_t start, el_val_t end);
|
||||
el_val_t __str_alloc(el_val_t n);
|
||||
el_val_t __str_set_char(el_val_t s, el_val_t i, el_val_t c);
|
||||
|
||||
/* URL encoding */
|
||||
el_val_t __url_encode(el_val_t s);
|
||||
el_val_t __url_decode(el_val_t s);
|
||||
|
||||
/* Environment */
|
||||
el_val_t __env_get(el_val_t key);
|
||||
|
||||
/* Subprocess */
|
||||
el_val_t __exec(el_val_t cmd);
|
||||
el_val_t __exec_bg(el_val_t cmd);
|
||||
|
||||
/* Process */
|
||||
el_val_t __exit_program(el_val_t code);
|
||||
|
||||
/* Filesystem */
|
||||
el_val_t __fs_exists(el_val_t path);
|
||||
el_val_t __fs_mkdir(el_val_t path);
|
||||
el_val_t __fs_read(el_val_t path);
|
||||
el_val_t __fs_write(el_val_t path, el_val_t content);
|
||||
el_val_t __fs_write_bytes(el_val_t path, el_val_t bytes, el_val_t n);
|
||||
el_val_t __fs_list_raw(el_val_t path);
|
||||
|
||||
/* HTTP server */
|
||||
el_val_t __http_response(el_val_t status, el_val_t headers_json, el_val_t body);
|
||||
el_val_t __http_serve(el_val_t port, el_val_t handler);
|
||||
el_val_t __http_serve_v2(el_val_t port, el_val_t handler);
|
||||
|
||||
/* HTTP conn fd / SSE (weak; overridden by el_seed.c when linked together) */
|
||||
el_val_t __http_conn_fd(void);
|
||||
el_val_t __http_sse_open(el_val_t conn_id);
|
||||
el_val_t __http_sse_send(el_val_t conn_id, el_val_t data);
|
||||
el_val_t __http_sse_close(el_val_t conn_id);
|
||||
|
||||
/* HTTP client (requires HAVE_CURL; stubs provided for no-curl builds) */
|
||||
el_val_t __http_do(el_val_t method, el_val_t url, el_val_t body,
|
||||
el_val_t headers_map, el_val_t timeout_ms);
|
||||
el_val_t __http_do_map(el_val_t method, el_val_t url, el_val_t body,
|
||||
el_val_t headers_json, el_val_t timeout_ms);
|
||||
el_val_t __http_do_map_to_file(el_val_t method, el_val_t url, el_val_t body,
|
||||
el_val_t headers_json, el_val_t output_path);
|
||||
|
||||
/* JSON */
|
||||
el_val_t __json_array_get(el_val_t json, el_val_t index);
|
||||
el_val_t __json_array_get_string(el_val_t json, el_val_t index);
|
||||
el_val_t __json_array_len(el_val_t json);
|
||||
el_val_t __json_get(el_val_t json, el_val_t key);
|
||||
el_val_t __json_get_raw(el_val_t json, el_val_t key);
|
||||
el_val_t __json_set(el_val_t json, el_val_t key, el_val_t value);
|
||||
el_val_t __json_parse_map(el_val_t json_str);
|
||||
el_val_t __json_stringify_val(el_val_t val);
|
||||
|
||||
/* Hashing */
|
||||
el_val_t __sha256_hex(el_val_t s);
|
||||
|
||||
/* State K/V */
|
||||
el_val_t __state_del(el_val_t key);
|
||||
el_val_t __state_get(el_val_t key);
|
||||
el_val_t __state_keys(void);
|
||||
el_val_t __state_set(el_val_t key, el_val_t val);
|
||||
|
||||
/* UUID */
|
||||
el_val_t __uuid_v4(void);
|
||||
|
||||
/* Args */
|
||||
el_val_t __args_json(void);
|
||||
|
||||
/* Compiler-support builtins — called by the El compiler's own source
|
||||
* (compiler.el, codegen.el) and registered in codegen.el's builtin_arity. */
|
||||
el_val_t stdout_to_file(el_val_t path);
|
||||
el_val_t stdout_restore(void);
|
||||
el_val_t el_mem_check(void);
|
||||
|
||||
/* Semantic retrieval surface. NOT interchangeable with engram_search_json,
|
||||
* which is lexical by design — see the note at the definition. */
|
||||
el_val_t engram_recall_json(el_val_t query, el_val_t limit);
|
||||
|
||||
/* Edges straight from the store — replaces the engram_save()+fs_read()
|
||||
* whole-graph round trip that /api/graph/edges used to do. */
|
||||
el_val_t engram_edges_json(el_val_t limit, el_val_t offset);
|
||||
|
||||
/* CGI identity accessors (read-only). */
|
||||
el_val_t cgi_principal(void);
|
||||
el_val_t cgi_network(void);
|
||||
el_val_t cgi_engram(void);
|
||||
|
||||
#ifdef __cplusplus
|
||||
}
|
||||
#endif
|
||||
|
||||
@@ -37,6 +37,82 @@
|
||||
#include <dlfcn.h>
|
||||
#include <curl/curl.h>
|
||||
|
||||
/* el_runtime.c bridge prototypes.
|
||||
*
|
||||
* A block of __-prefixed wrappers further down in this file (http serving,
|
||||
* JSON access, key-val state, URL/HTML escaping, and the whole engram_*
|
||||
* node/edge/layer/search surface -- 51 symbols in total) delegate to
|
||||
* unprefixed counterparts that are implemented in el_runtime.c, not here.
|
||||
* Porting them into native el_seed.c or El has not happened yet.
|
||||
* tools/install.sh compiles el_seed.c and el_runtime.c as separate objects
|
||||
* and archives both into libel.a, so the symbols are always present at link
|
||||
* time. el_seed.c alone was just missing the prototypes, which made even a
|
||||
* standalone -c compile of this one file fail on a toolchain that now treats
|
||||
* an implicit function declaration as a hard error under C11.
|
||||
*
|
||||
* A plain include of el_runtime.h was tried first and rejected: it redefines
|
||||
* el_to_float and el_from_float, which el_seed.h already provides. Narrow
|
||||
* prototypes, copied verbatim from el_runtime.h, avoid that collision without
|
||||
* pulling in the rest of the retiring runtime header.
|
||||
*/
|
||||
el_val_t http_response(el_val_t status, el_val_t headers_json, el_val_t body);
|
||||
void http_serve(el_val_t port, el_val_t handler);
|
||||
void http_serve_v2(el_val_t port, el_val_t handler);
|
||||
el_val_t json_get(el_val_t json, el_val_t key);
|
||||
el_val_t json_get_string(el_val_t json_str, el_val_t key);
|
||||
el_val_t json_get_int(el_val_t json_str, el_val_t key);
|
||||
el_val_t json_get_float(el_val_t json_str, el_val_t key);
|
||||
el_val_t json_get_bool(el_val_t json_str, el_val_t key);
|
||||
el_val_t json_get_raw(el_val_t json_str, el_val_t key);
|
||||
el_val_t json_parse(el_val_t s);
|
||||
el_val_t json_set(el_val_t json_str, el_val_t key, el_val_t value);
|
||||
el_val_t json_stringify(el_val_t v);
|
||||
el_val_t json_array_len(el_val_t json_str);
|
||||
el_val_t json_array_get(el_val_t json_str, el_val_t index);
|
||||
el_val_t json_array_get_string(el_val_t json_str, el_val_t index);
|
||||
el_val_t state_set(el_val_t key, el_val_t value);
|
||||
el_val_t state_get(el_val_t key);
|
||||
el_val_t state_del(el_val_t key);
|
||||
el_val_t state_keys(void);
|
||||
el_val_t url_encode(el_val_t s);
|
||||
el_val_t url_decode(el_val_t s);
|
||||
el_val_t el_html_sanitize(el_val_t input_html, el_val_t allowlist_json);
|
||||
el_val_t engram_node(el_val_t content, el_val_t node_type, el_val_t salience);
|
||||
el_val_t engram_node_full(el_val_t content, el_val_t node_type, el_val_t label,
|
||||
el_val_t salience, el_val_t importance, el_val_t confidence,
|
||||
el_val_t tier, el_val_t tags);
|
||||
el_val_t engram_node_layered(el_val_t content, el_val_t node_type, el_val_t label,
|
||||
el_val_t salience, el_val_t certainty, el_val_t confidence,
|
||||
el_val_t status, el_val_t tags, el_val_t layer_id);
|
||||
el_val_t engram_add_layer(el_val_t name, el_val_t priority, el_val_t suppressible,
|
||||
el_val_t transparent, el_val_t injectable);
|
||||
el_val_t engram_remove_layer(el_val_t layer_id);
|
||||
el_val_t engram_list_layers(void);
|
||||
el_val_t engram_list_layers_json(void);
|
||||
el_val_t engram_get_node(el_val_t id);
|
||||
el_val_t engram_get_node_json(el_val_t id);
|
||||
el_val_t engram_get_node_by_label(el_val_t label);
|
||||
void engram_strengthen(el_val_t node_id);
|
||||
void engram_forget(el_val_t node_id);
|
||||
el_val_t engram_node_count(void);
|
||||
el_val_t engram_edge_count(void);
|
||||
el_val_t engram_scan_nodes(el_val_t limit, el_val_t offset);
|
||||
el_val_t engram_scan_nodes_json(el_val_t limit, el_val_t offset);
|
||||
el_val_t engram_scan_nodes_by_type_json(el_val_t node_type, el_val_t limit, el_val_t offset);
|
||||
el_val_t engram_search(el_val_t query, el_val_t limit);
|
||||
el_val_t engram_search_json(el_val_t query, el_val_t limit);
|
||||
el_val_t engram_activate(el_val_t query, el_val_t depth);
|
||||
el_val_t engram_activate_json(el_val_t query, el_val_t depth);
|
||||
el_val_t engram_compile_layered_json(el_val_t intent, el_val_t depth);
|
||||
el_val_t engram_stats_json(void);
|
||||
void engram_connect(el_val_t from_id, el_val_t to_id, el_val_t weight, el_val_t relation);
|
||||
el_val_t engram_edge_between(el_val_t from_id, el_val_t to_id);
|
||||
el_val_t engram_neighbors(el_val_t node_id);
|
||||
el_val_t engram_neighbors_filtered(el_val_t node_id, el_val_t max_depth, el_val_t direction);
|
||||
el_val_t engram_neighbors_json(el_val_t node_id, el_val_t max_depth, el_val_t direction);
|
||||
el_val_t engram_load(el_val_t path);
|
||||
el_val_t engram_save(el_val_t path);
|
||||
|
||||
/* ── Private allocator ───────────────────────────────────────────────────── */
|
||||
/*
|
||||
* el_seed.c carries its own arena for per-request allocation tracking.
|
||||
@@ -831,6 +907,219 @@ void __mutex_unlock(el_val_t m) {
|
||||
pthread_mutex_unlock(&_el_mutexes[slot]);
|
||||
}
|
||||
|
||||
/* ── Channels ─────────────────────────────────────────────────────────────── *
|
||||
* Buffered MPMC channel backed by a mutex + condvar + circular buffer.
|
||||
* Ported from the pre-restructure el_runtime.c (b2aac4b) — runtime/channel.el
|
||||
* has always called these five primitives, but they were never carried
|
||||
* forward into el_seed.c when el_runtime.c was consolidated onto the
|
||||
* canonical release copy. Native channels were silently unlinkable on dev
|
||||
* until this port.
|
||||
*
|
||||
* __channel_new(capacity) -> Int (handle)
|
||||
* __channel_send(ch, msg) — blocks if full (capacity > 0) or never (unbounded)
|
||||
* __channel_recv(ch) -> String — blocks until a message is available
|
||||
* __channel_try_recv(ch) -> String — non-blocking, returns "" if empty
|
||||
* __channel_close(ch) — signal no more sends; recv drains remaining
|
||||
*
|
||||
* Bounded channels (cap > 0): circular buffer, sender blocks when full.
|
||||
* Unbounded channels (cap == 0): dynamic array, sender never blocks.
|
||||
*/
|
||||
#define EL_CHANNEL_MAX 64
|
||||
#define EL_CHANNEL_BUF 1024
|
||||
|
||||
typedef struct {
|
||||
char** buf;
|
||||
int cap; /* 0 = unbounded (grows dynamically) */
|
||||
int head, tail, count;
|
||||
int dyn_cap; /* allocated slots for unbounded mode */
|
||||
int closed;
|
||||
pthread_mutex_t mu;
|
||||
pthread_cond_t not_empty;
|
||||
pthread_cond_t not_full;
|
||||
} ElChannel;
|
||||
|
||||
static ElChannel _channels[EL_CHANNEL_MAX];
|
||||
static int _channel_count = 0;
|
||||
static pthread_mutex_t _channel_alloc_mu = PTHREAD_MUTEX_INITIALIZER;
|
||||
|
||||
el_val_t __channel_new(el_val_t capacity_v) {
|
||||
int cap = (int)(int64_t)capacity_v;
|
||||
if (cap < 0) cap = 0;
|
||||
|
||||
pthread_mutex_lock(&_channel_alloc_mu);
|
||||
if (_channel_count >= EL_CHANNEL_MAX) {
|
||||
pthread_mutex_unlock(&_channel_alloc_mu);
|
||||
fprintf(stderr, "[__channel_new] channel table full\n");
|
||||
return EL_INT(-1);
|
||||
}
|
||||
int slot = _channel_count++;
|
||||
pthread_mutex_unlock(&_channel_alloc_mu);
|
||||
|
||||
ElChannel* ch = &_channels[slot];
|
||||
memset(ch, 0, sizeof(*ch));
|
||||
ch->cap = cap;
|
||||
ch->closed = 0;
|
||||
ch->head = 0;
|
||||
ch->tail = 0;
|
||||
ch->count = 0;
|
||||
|
||||
if (cap > 0) {
|
||||
/* Bounded: fixed circular buffer. */
|
||||
ch->buf = (char**)malloc((size_t)cap * sizeof(char*));
|
||||
ch->dyn_cap = cap;
|
||||
} else {
|
||||
/* Unbounded: start with EL_CHANNEL_BUF slots, grow as needed. */
|
||||
ch->buf = (char**)malloc(EL_CHANNEL_BUF * sizeof(char*));
|
||||
ch->dyn_cap = EL_CHANNEL_BUF;
|
||||
}
|
||||
if (!ch->buf) {
|
||||
fprintf(stderr, "[__channel_new] out of memory\n");
|
||||
return EL_INT(-1);
|
||||
}
|
||||
|
||||
pthread_mutex_init(&ch->mu, NULL);
|
||||
pthread_cond_init(&ch->not_empty, NULL);
|
||||
pthread_cond_init(&ch->not_full, NULL);
|
||||
|
||||
return EL_INT(slot);
|
||||
}
|
||||
|
||||
el_val_t __channel_send(el_val_t ch_v, el_val_t msg_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
const char* msg = EL_CSTR(msg_v);
|
||||
if (!msg) msg = "";
|
||||
char* copy = strdup(msg); /* channel owns the string */
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
|
||||
if (ch->closed) {
|
||||
/* Send on closed channel is a no-op (drop the message). */
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
free(copy);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
if (ch->cap > 0) {
|
||||
/* Bounded: block while full. */
|
||||
while (ch->count >= ch->cap && !ch->closed) {
|
||||
pthread_cond_wait(&ch->not_full, &ch->mu);
|
||||
}
|
||||
if (ch->closed) {
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
free(copy);
|
||||
return EL_STR("");
|
||||
}
|
||||
ch->buf[ch->tail] = copy;
|
||||
ch->tail = (ch->tail + 1) % ch->cap;
|
||||
ch->count++;
|
||||
} else {
|
||||
/* Unbounded: grow the buffer if needed. */
|
||||
if (ch->count >= ch->dyn_cap) {
|
||||
int new_cap = ch->dyn_cap * 2;
|
||||
char** grown = (char**)realloc(ch->buf, (size_t)new_cap * sizeof(char*));
|
||||
if (!grown) {
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
free(copy);
|
||||
fprintf(stderr, "[__channel_send] out of memory growing channel\n");
|
||||
return EL_STR("");
|
||||
}
|
||||
/* The circular buffer may have wrapped. Linearise it first.
|
||||
* In unbounded mode head is always 0 (we append at tail, drain
|
||||
* from head), so a simple memmove isn't needed — but if the
|
||||
* buffer did wrap (tail < head after growth), we need to fix up.
|
||||
* Simplest safe path: if tail wrapped, move the head..old_cap
|
||||
* segment to new_cap..new_cap+(old_cap-head). */
|
||||
if (ch->tail < ch->head) {
|
||||
/* Wrapped: [head..old_cap) is the front, [0..tail) is the back. */
|
||||
int front = ch->dyn_cap - ch->head;
|
||||
memmove(grown + ch->dyn_cap, grown + ch->head, (size_t)front * sizeof(char*));
|
||||
ch->head = ch->dyn_cap;
|
||||
}
|
||||
ch->buf = grown;
|
||||
ch->dyn_cap = new_cap;
|
||||
}
|
||||
ch->buf[ch->tail] = copy;
|
||||
ch->tail = (ch->tail + 1) % ch->dyn_cap;
|
||||
ch->count++;
|
||||
}
|
||||
|
||||
pthread_cond_signal(&ch->not_empty);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
el_val_t __channel_recv(el_val_t ch_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
|
||||
/* Block until there is a message or the channel is closed and drained. */
|
||||
while (ch->count == 0 && !ch->closed) {
|
||||
pthread_cond_wait(&ch->not_empty, &ch->mu);
|
||||
}
|
||||
|
||||
if (ch->count == 0) {
|
||||
/* Closed and empty — signal EOF. */
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
int buf_cap = (ch->cap > 0) ? ch->cap : ch->dyn_cap;
|
||||
char* msg = ch->buf[ch->head];
|
||||
ch->head = (ch->head + 1) % buf_cap;
|
||||
ch->count--;
|
||||
|
||||
pthread_cond_signal(&ch->not_full);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
|
||||
/* Hand the string to the arena so it is freed after the request. */
|
||||
seed_arena_track(msg);
|
||||
return EL_STR(msg);
|
||||
}
|
||||
|
||||
el_val_t __channel_try_recv(el_val_t ch_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
|
||||
if (ch->count == 0) {
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
int buf_cap = (ch->cap > 0) ? ch->cap : ch->dyn_cap;
|
||||
char* msg = ch->buf[ch->head];
|
||||
ch->head = (ch->head + 1) % buf_cap;
|
||||
ch->count--;
|
||||
|
||||
pthread_cond_signal(&ch->not_full);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
|
||||
seed_arena_track(msg);
|
||||
return EL_STR(msg);
|
||||
}
|
||||
|
||||
el_val_t __channel_close(el_val_t ch_v) {
|
||||
int slot = (int)(int64_t)ch_v;
|
||||
if (slot < 0 || slot >= EL_CHANNEL_MAX) return EL_STR("");
|
||||
ElChannel* ch = &_channels[slot];
|
||||
|
||||
pthread_mutex_lock(&ch->mu);
|
||||
ch->closed = 1;
|
||||
/* Wake all blocked recvers and senders so they can observe the close. */
|
||||
pthread_cond_broadcast(&ch->not_empty);
|
||||
pthread_cond_broadcast(&ch->not_full);
|
||||
pthread_mutex_unlock(&ch->mu);
|
||||
return EL_STR("");
|
||||
}
|
||||
|
||||
/* ── Subprocess ──────────────────────────────────────────────────────────── */
|
||||
|
||||
el_val_t __exec(el_val_t cmd) {
|
||||
@@ -1082,6 +1371,11 @@ el_val_t __engram_scan_nodes_json(el_val_t limit, el_val_t offset) {
|
||||
return engram_scan_nodes_json(limit, offset);
|
||||
}
|
||||
|
||||
el_val_t engram_edges_json(el_val_t limit, el_val_t offset);
|
||||
el_val_t __engram_edges_json(el_val_t limit, el_val_t offset) {
|
||||
return engram_edges_json(limit, offset);
|
||||
}
|
||||
|
||||
el_val_t __engram_scan_nodes_by_type_json(el_val_t node_type, el_val_t limit, el_val_t offset) {
|
||||
return engram_scan_nodes_by_type_json(node_type, limit, offset);
|
||||
}
|
||||
@@ -1094,7 +1388,27 @@ el_val_t __engram_activate_json(el_val_t query, el_val_t depth) {
|
||||
return engram_activate_json(query, depth);
|
||||
}
|
||||
|
||||
/* Forward decls for el_runtime.c symbols this file wraps. el_seed.c does not
|
||||
* include el_runtime.h (documented in lang/AGENTS.md), so each wrapped symbol
|
||||
* needs a prototype here or clang treats it as an implicit declaration (error
|
||||
* under C99+) and the ABI mis-truncates the el_val_t return. */
|
||||
el_val_t engram_op_assert_json(el_val_t node_id, el_val_t depth);
|
||||
el_val_t engram_node_full_in(el_val_t purview, el_val_t content, el_val_t node_type, el_val_t label,
|
||||
el_val_t salience, el_val_t importance, el_val_t confidence,
|
||||
el_val_t tier, el_val_t tags);
|
||||
void engram_connect_in(el_val_t purview, el_val_t from_id, el_val_t to_id,
|
||||
el_val_t weight, el_val_t relation);
|
||||
|
||||
el_val_t __engram_stats_json(void) { return engram_stats_json(); }
|
||||
el_val_t __engram_op_assert_json(el_val_t node_id, el_val_t depth) { return engram_op_assert_json(node_id, depth); }
|
||||
el_val_t __engram_node_full_in(el_val_t purview, el_val_t content, el_val_t node_type, el_val_t label,
|
||||
el_val_t salience, el_val_t importance, el_val_t confidence,
|
||||
el_val_t tier, el_val_t tags) {
|
||||
return engram_node_full_in(purview, content, node_type, label, salience, importance, confidence, tier, tags);
|
||||
}
|
||||
void __engram_connect_in(el_val_t purview, el_val_t from_id, el_val_t to_id, el_val_t weight, el_val_t relation) {
|
||||
engram_connect_in(purview, from_id, to_id, weight, relation);
|
||||
}
|
||||
el_val_t __engram_list_layers_json(void) { return engram_list_layers_json(); }
|
||||
|
||||
el_val_t __engram_compile_layered_json(el_val_t intent, el_val_t depth) {
|
||||
|
||||
@@ -139,6 +139,13 @@ el_val_t __mutex_new(void);
|
||||
void __mutex_lock(el_val_t m);
|
||||
void __mutex_unlock(el_val_t m);
|
||||
|
||||
/* Buffered MPMC channel (runtime/channel.el). capacity=0 means unbounded. */
|
||||
el_val_t __channel_new(el_val_t capacity);
|
||||
el_val_t __channel_send(el_val_t ch, el_val_t msg); /* blocks if bounded+full */
|
||||
el_val_t __channel_recv(el_val_t ch); /* blocks until available */
|
||||
el_val_t __channel_try_recv(el_val_t ch); /* non-blocking, "" if empty */
|
||||
el_val_t __channel_close(el_val_t ch);
|
||||
|
||||
/* ── Subprocess ──────────────────────────────────────────────────────────── */
|
||||
|
||||
el_val_t __exec(el_val_t cmd); /* popen, capture all stdout, return String */
|
||||
@@ -233,6 +240,12 @@ el_val_t __engram_scan_nodes_by_type_json(el_val_t node_type, el_val_t limit, e
|
||||
el_val_t __engram_neighbors_json(el_val_t node_id, el_val_t max_depth, el_val_t direction);
|
||||
el_val_t __engram_activate_json(el_val_t query, el_val_t depth);
|
||||
el_val_t __engram_stats_json(void);
|
||||
el_val_t __engram_op_assert_json(el_val_t node_id, el_val_t depth);
|
||||
el_val_t __engram_node_full_in(el_val_t purview, el_val_t content, el_val_t node_type, el_val_t label,
|
||||
el_val_t salience, el_val_t importance, el_val_t confidence,
|
||||
el_val_t tier, el_val_t tags);
|
||||
void __engram_connect_in(el_val_t purview, el_val_t from_id, el_val_t to_id,
|
||||
el_val_t weight, el_val_t relation);
|
||||
el_val_t __engram_list_layers_json(void);
|
||||
el_val_t __engram_compile_layered_json(el_val_t intent, el_val_t depth);
|
||||
|
||||
|
||||
@@ -0,0 +1,184 @@
|
||||
# Swarm + CCR + Work-Tracking — Neuron's bounded parallel execution, in native El
|
||||
|
||||
Bounded parallel agent execution on El's **native** concurrency — no external
|
||||
orchestrator. Grounded directly in two of Will's frameworks:
|
||||
|
||||
- **Swarm Architecture** (*Bounded Parallel Agent Execution*, Mar 2026)
|
||||
- **Compiled Context Runtime / CCR** (*Process-Driven Agent Execution with
|
||||
Unbounded Local Memory*, Mar 2026)
|
||||
|
||||
A swarm is a **coordinator** (the main thread) that mints a correlation identity,
|
||||
compiles a **bounded per-worker context (CCR)**, dispatches workers as **native
|
||||
pthreads** (`thread.el` `spawn`/`join`), tracks every unit of work durably, and
|
||||
**converges** results before returning control to the parent step.
|
||||
|
||||
```
|
||||
Parent step
|
||||
└─ swarm_run(blueprint, knowledge_refs, inputs, config)
|
||||
fan-out ──▶ worker_1 (CCR ctx_1) ─┐ native
|
||||
worker_2 (CCR ctx_2) ─┤ pthreads,
|
||||
worker_k (CCR ctx_k) ─┘ bounded by `concurrency`
|
||||
converge ─▶ collect | merge | vote | reduce ──▶ merged result
|
||||
```
|
||||
|
||||
## Why it runs on El natively
|
||||
|
||||
El is natively agentic. This capability composes El's shipped primitives — it
|
||||
adds no bespoke runtime:
|
||||
|
||||
| Primitive | Source | Role in the swarm |
|
||||
|-----------|--------|-------------------|
|
||||
| `spawn(fn,arg)` / `join(tid)` | `runtime/thread.el` → `__thread_create` (pthread + dlsym) | fan-out / rejoin |
|
||||
| `parallel_map`, `with_mutex` | `runtime/thread.el` | reference concurrency patterns |
|
||||
| Go-style channels | `runtime/channel.el` → `__channel_*` | available for vertical event streams |
|
||||
| `engram_*`, `http_*`, `fs_*`, `json_*` | `el_runtime.c` builtins | retrieval, tracking, I/O |
|
||||
|
||||
Every El fn compiles to a global C symbol, so any top-level `(String)->String`
|
||||
fn is directly threadable — the worker entry is exactly such a fn.
|
||||
|
||||
## Modules
|
||||
|
||||
| File | Framework grounding | What it does |
|
||||
|------|--------------------|--------------|
|
||||
| `worktrack.el` | Swarm §6 (correlation IDs, audit) | Durable, single-writer **JSONL journal** keyed by correlation ID; reconstructable status report; opt-in engram mirror (`SWARM_MIRROR=1`). |
|
||||
| `containment.el` | Swarm §3 + the single-writer invariant | Scope tokens w/ capabilities; **Rule 1** (no join), **Rule 2** (no open), **Rule 3** (no lateral edge), **Rule 4** (engram-write is @manager-only, by capability) enforced as checks. |
|
||||
| `ccr.el` | CCR §5 + Swarm §9.3 | Per-worker **Compiled Context Routing**: retrieve → scope → compact into a **bounded, minimal** package. The compiled-context boundary *is* the security boundary. |
|
||||
| `primitives.el` | CCR §2 (Five Primitives) | `attend / think / intend / act / learn` seam the swarm composes over. Engram-backed; explicit binding point for the API-surface reshape. |
|
||||
| `swarm.el` | Swarm §2, §4, §5 | The coordinator: fan-out/converge on native threads, bounded concurrency, four convergence strategies, integer failure threshold, full tracking. |
|
||||
|
||||
## Invariant: only the orchestrator mutates global engram state
|
||||
|
||||
**Only the orchestrator (@manager) writes to the engram / mutates global state.
|
||||
Workers are read-only against the full engram and may write only their own local
|
||||
geometry (their returned result + the journal). A worker is STRUCTURALLY UNABLE
|
||||
to mutate global engram state.**
|
||||
|
||||
This is **Rule 4** — an **authority gate, not a health gate**. Scope tokens carry
|
||||
a capability set: the orchestrator's token holds `engram:write` + `dharma:emit`
|
||||
(@manager-only, the VBD rule that only the manager mutates global state); a
|
||||
worker's token holds **only** `engram:read`. Every engram mutation
|
||||
(`op_write`/`op_relate`/`op_supersede` → `POST /api/nodes`, `/api/edges`,
|
||||
`DELETE`) flows through `swarm_engram_write`, which checks the caller's capability
|
||||
via the **same scope-token mechanism as the live Rule-2 denial** and rejects any
|
||||
worker **before any HTTP is issued**. Capability is fixed at mint time and cannot
|
||||
be acquired at runtime — so the guarantee holds regardless of engram health
|
||||
(distinct from the `SWARM_WRITE_HEALTHY` *health* gate).
|
||||
|
||||
The **curated merge is the only write path**: workers return geometry; the
|
||||
orchestrator, and only the orchestrator, commits the approved/verified geometry
|
||||
back (`commit=1`). Workers keep full-engram **read** access (`op_think`/`op_read`).
|
||||
|
||||
Proven in `harness_real_cognition.el` (§G): a worker `swarm_engram_write` is
|
||||
DENIED by capability with no node created and the violation journalled; the
|
||||
orchestrator passes the gate as the sole authorized writer.
|
||||
|
||||
## Containment → distribution
|
||||
|
||||
The three containment rules make workers **location-independent** (Swarm §9): a
|
||||
worker reads only its compiled context, shares no state with siblings, and its
|
||||
only outward edge is the returned result. The same coordinator can run workers
|
||||
as local threads today or dispatch them across machines later — the mechanism is
|
||||
identical; only the topology changes. Enforced here:
|
||||
|
||||
- **Rule 2** — `swarm_run` rejects any swarm opened under a worker token.
|
||||
- **Rules 1 + 3** — each worker gets a *closed* worker token; the coordinator is
|
||||
the only journal writer, so workers share no mutable state.
|
||||
|
||||
## Usage
|
||||
|
||||
```el
|
||||
// one process step fans out; results converge before the next step
|
||||
let inputs: String = "[\"billing\",\"payments\",\"ledger\"]"
|
||||
let refs: String = "[\"Volatility-Based Decomposition\"]" // CCR knowledge refs
|
||||
let cfg: String = "{\"concurrency\":\"4\",\"strategy\":\"collect\",\"min_success_ratio\":\"1.0\"}"
|
||||
let result: String = swarm_run("analyze_item", refs, inputs, cfg)
|
||||
// result: { corr_id, status, merged, report }
|
||||
```
|
||||
|
||||
Build any program that uses the swarm:
|
||||
|
||||
```bash
|
||||
lang/swarm/build.sh myprog.el ./myprog # concat + elc + cc (el_runtime.c)
|
||||
```
|
||||
|
||||
Config keys: `concurrency` (max workers at once), `strategy`
|
||||
(`collect|merge|vote|reduce`), `min_success_ratio` (decimal string, e.g. `0.8`),
|
||||
`caller_token` (containment). Env: `SWARM_TRACK_DIR` (journal dir),
|
||||
`CCR_TOKEN_BUDGET`, `ENGRAM_URL`/`ENGRAM_API_KEY` (retrieval + mirror),
|
||||
`SWARM_MIRROR=1`.
|
||||
|
||||
## Tests
|
||||
|
||||
```bash
|
||||
lang/swarm/build.sh lang/swarm/tests/test_swarm.el /tmp/t && SWARM_TRACK_DIR=/tmp/trk /tmp/t # 12/12
|
||||
lang/swarm/build.sh lang/swarm/tests/test_convergence.el /tmp/c && SWARM_TRACK_DIR=/tmp/trk /tmp/c # 8/8
|
||||
# integration against an isolated engram clone (never live):
|
||||
source <sandbox>/.nsbx-env
|
||||
lang/swarm/build.sh lang/swarm/tests/integ_engram.el /tmp/i && /tmp/i
|
||||
```
|
||||
|
||||
## Local-swarm integration harness (the one flip)
|
||||
|
||||
`tests/harness_local_swarm.el` proves the **full local-swarm mechanics today** on
|
||||
the isolated clone with the primitive seam pointed at the hermetic stub — 17/17
|
||||
green: 8 native-thread workers at concurrency 4, reduce + vote convergence, CCR
|
||||
scoping + non-leak, all three containment rules (incl. live Rule-2 denial),
|
||||
durable work-tracking, and **afferent telemetry** observed by the @manager.
|
||||
|
||||
Binding to the reshape's decorated primitives is **one flip and a run**:
|
||||
|
||||
```
|
||||
# in primitive_binding.el — change one line each:
|
||||
fn bound_think(ctx, instruction) { return think(ctx, instruction) } # decorated, dharma bus
|
||||
# then:
|
||||
SWARM_PRIMITIVE_SEAM=decorated lang/swarm/build.sh tests/harness_local_swarm.el ./h && ./h
|
||||
```
|
||||
|
||||
Nothing else in the swarm changes. `primitive_seam.el` (`seam_think/attend/learn`)
|
||||
already routes every worker primitive call through this one switch, and the same
|
||||
harness runs the bound path. Today `SWARM_PRIMITIVE_SEAM=decorated` still runs
|
||||
green because the binding falls back to the stub — proving the flip path executes.
|
||||
|
||||
## Real cognition — the seam is BOUND
|
||||
|
||||
`primitive_binding.el` is bound to the api-reshape agent's proven primitives
|
||||
(`wt/api-reshape@d4f401d`): `bound_think -> op_think` (GET `/api/think`), real
|
||||
768-dim gradients over the engram geometry. `reshape_surface.el` composes those
|
||||
read/cognition primitives verbatim (`op_think/read/attend/learn`).
|
||||
|
||||
`tests/harness_real_cognition.el` runs the **local swarm on real cognition**,
|
||||
17/17 green with `SWARM_PRIMITIVE_SEAM=decorated` against the `:8901` clone: 8
|
||||
native-thread workers, each a real `think` over its CCR-scoped **node-id anchor**
|
||||
(free-text anchors return "geometry unavailable"), `@manager` reduce+vote, all
|
||||
three containment rules, afferent telemetry, durable tracking. Per-anchor support
|
||||
counts (e.g. 6 / 16 / 87) drive a genuine, cognition-derived vote.
|
||||
|
||||
> **Build note (load-bearing):** the swarm build **must** define `HAVE_CURL`
|
||||
> (`build.sh` does). Without it every `http_*` builtin is a
|
||||
> `{"error":"not built with HAVE_CURL"}` stub — real HTTP silently disappears.
|
||||
|
||||
Writes (`attend`/`learn`, `POST`) are gated behind `SWARM_WRITE_HEALTHY=1` and the
|
||||
api-reshape agent's gate-1 write-healthy clone; the proven run is read-cognition.
|
||||
|
||||
## Built vs stubbed (honest)
|
||||
|
||||
**Real, tested:**
|
||||
- Native-thread fan-out/converge, bounded concurrency, order-preserving rejoin.
|
||||
- All three containment rules enforced (scope tokens + lateral-edge check).
|
||||
- CCR per-worker context: retrieval → scoping → compaction, bounded, non-leaking
|
||||
(a worker never receives sibling inputs) — verified against the live isolated mind.
|
||||
- Full durable work-tracking (JSONL journal, reconstructable report).
|
||||
- Four convergence strategies + integer failure threshold / partial-abort.
|
||||
|
||||
**Seam / not yet bound:**
|
||||
- `primitives.el` `think` is a deterministic, hermetic transform (no model call).
|
||||
Binding point is marked `PRIMITIVE_BINDING`; wire to the API-surface reshape's
|
||||
`think/act/attend/intend/learn` when it lands.
|
||||
- Blueprints are dispatched by name in `swarm_run_blueprint` (default +
|
||||
`classify`/`faildemo` demos). A YAML process-definition loader (Swarm §5) is
|
||||
future work — the runtime contract is in place.
|
||||
- Distributed placement (cloud/edge/federated topologies, Swarm §9.2) is
|
||||
structurally enabled by containment but not yet wired to a placement layer;
|
||||
today all workers are local native threads.
|
||||
- Engram work-tracking mirror is opt-in; the durable substrate is the journal.
|
||||
|
||||
Executable
+61
@@ -0,0 +1,61 @@
|
||||
#!/usr/bin/env bash
|
||||
# build.sh — compile an El program that uses the swarm capability.
|
||||
#
|
||||
# Concatenates the El native-concurrency stdlib (thread.el, channel.el) and the
|
||||
# swarm capability modules in dependency order, then the user program, compiles
|
||||
# with the canonical elc, and links against the shared C runtime.
|
||||
#
|
||||
# Usage:
|
||||
# swarm/build.sh <program.el> <out-binary>
|
||||
#
|
||||
# The swarm modules use only el_runtime.c builtins plus thread.el/channel.el,
|
||||
# so nothing else needs concatenating (engram_*, json_*, str_*, fs_*, http_*,
|
||||
# uuid_v4, now_millis are all C builtins in el_runtime.c).
|
||||
|
||||
set -uo pipefail
|
||||
cd "$(dirname "$0")/.." # -> lang/
|
||||
LANG_DIR="$(pwd)"
|
||||
ELC="${ELC:-${LANG_DIR}/dist/platform/elc}"
|
||||
RT="${LANG_DIR}/el-compiler/runtime"
|
||||
|
||||
PROG="${1:?usage: build.sh <program.el> <out-binary>}"
|
||||
OUT="${2:?usage: build.sh <program.el> <out-binary>}"
|
||||
|
||||
# swarm module load order (each may depend on those before it):
|
||||
# worktrack — durable work-tracking journal (no swarm deps)
|
||||
# containment — the three containment rules (no swarm deps)
|
||||
# primitives — think/act/attend/intend/learn seam (no swarm deps)
|
||||
# ccr — per-worker compiled bounded context (depends: primitives)
|
||||
# swarm — orchestrator: fan-out/converge (depends: all above + thread)
|
||||
SWARM_MODULES="
|
||||
swarm/worktrack.el
|
||||
swarm/containment.el
|
||||
swarm/primitives.el
|
||||
swarm/reshape_surface.el
|
||||
swarm/primitive_binding.el
|
||||
swarm/primitive_seam.el
|
||||
swarm/ccr.el
|
||||
swarm/swarm.el
|
||||
"
|
||||
|
||||
TMP_C="$(mktemp -t swarm_build.XXXXXX).c"
|
||||
COMBINED="$(mktemp -t swarm_combined.XXXXXX).el"
|
||||
|
||||
cat runtime/thread.el runtime/channel.el $SWARM_MODULES "$PROG" > "$COMBINED"
|
||||
|
||||
if ! "$ELC" "$COMBINED" > "$TMP_C" 2>/tmp/swarm.elc.err; then
|
||||
echo "elc FAILED:" >&2
|
||||
sed 's/^/ /' /tmp/swarm.elc.err >&2
|
||||
rm -f "$TMP_C" "$COMBINED"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
if ! cc -O2 -DHAVE_CURL -I "$RT" "$TMP_C" "$RT/el_runtime.c" -lcurl -lpthread -lm -o "$OUT" 2>/tmp/swarm.cc.err; then
|
||||
echo "cc FAILED:" >&2
|
||||
sed 's/^/ /' /tmp/swarm.cc.err >&2
|
||||
rm -f "$TMP_C" "$COMBINED"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
rm -f "$TMP_C" "$COMBINED"
|
||||
echo "built: $OUT"
|
||||
@@ -0,0 +1,153 @@
|
||||
// ccr.el — Compiled Context Routing for work distribution.
|
||||
//
|
||||
// The same spine as the API's vantage-read, applied per worker. Instead of
|
||||
// handing every worker the coordinator's full memory, CCR compiles a MINIMAL,
|
||||
// BOUNDED context package scoped to exactly one worker's input (CCR §5, "Compiled
|
||||
// Context Injection"; Swarm §9.3, "The Compiled Context Boundary as Security
|
||||
// Boundary").
|
||||
//
|
||||
// The pipeline is CCR §5.1: Retrieval -> Scoping -> Compilation -> (Injection,
|
||||
// which here is placing the package into the worker's task envelope).
|
||||
//
|
||||
// 1. Retrieval — resolve the blueprint's knowledge refs + the input's salient
|
||||
// terms against the mind (primitive_attend).
|
||||
// 2. Scoping — keep only what THIS input needs; drop everything else. A
|
||||
// worker never receives sibling inputs or unrelated memory.
|
||||
// 3. Compilation— compact to a CTX string within a token budget (lossless of
|
||||
// meaning, smaller in tokens): collapse blank runs, dedupe
|
||||
// lines, then bound to the budget.
|
||||
//
|
||||
// The package a worker receives is therefore (a) sufficient for its task and
|
||||
// (b) incapable of leaking what it was never given — the containment boundary
|
||||
// and the security boundary are the same object.
|
||||
|
||||
// ── token budget helpers ─────────────────────────────────────────────────────
|
||||
|
||||
// ccr_est_tokens — cheap token estimate (~4 chars/token).
|
||||
fn ccr_est_tokens(s: String) -> Int {
|
||||
return str_len(s) / 4
|
||||
}
|
||||
|
||||
// ccr_default_budget — default per-worker context budget in tokens.
|
||||
// Override with CCR_TOKEN_BUDGET.
|
||||
fn ccr_default_budget() -> Int {
|
||||
let b: String = env("CCR_TOKEN_BUDGET")
|
||||
if str_eq(b, "") {
|
||||
return 1200
|
||||
}
|
||||
return str_to_int(b)
|
||||
}
|
||||
|
||||
// ── stage 3: compaction ──────────────────────────────────────────────────────
|
||||
|
||||
// ccr_compact — collapse blank-line runs and drop exact duplicate lines, then
|
||||
// bound the result to `budget` tokens (truncate on a line boundary). Meaning is
|
||||
// preserved; token count falls (CCR §5.2).
|
||||
fn ccr_compact(text: String, budget: Int) -> String {
|
||||
let lines: [String] = str_split_lines(text)
|
||||
let n: Int = el_list_len(lines)
|
||||
let seen: String = "\n"
|
||||
let out: String = ""
|
||||
let out_tokens = 0
|
||||
let i = 0
|
||||
while i < n {
|
||||
let ln: String = str_trim(el_list_get(lines, i))
|
||||
if str_eq(ln, "") {
|
||||
let i = i + 1
|
||||
} else {
|
||||
let marker: String = "\n" + ln + "\n"
|
||||
if str_contains(seen, marker) {
|
||||
// duplicate line — skip
|
||||
let i = i + 1
|
||||
} else {
|
||||
let seen = seen + ln + "\n"
|
||||
let line_tokens: Int = ccr_est_tokens(ln) + 1
|
||||
if out_tokens + line_tokens > budget {
|
||||
// budget exhausted — stop (bounded)
|
||||
let i = n
|
||||
} else {
|
||||
let out = out + ln + "\n"
|
||||
let out_tokens = out_tokens + line_tokens
|
||||
let i = i + 1
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// ── stages 1+2: retrieve + scope ─────────────────────────────────────────────
|
||||
|
||||
// ccr_retrieve_scoped — pull context relevant to this input and its blueprint
|
||||
// knowledge refs, scoped to a fraction of the budget so no single source floods
|
||||
// the package. Returns compacted retrieved text (may be empty if the mind is
|
||||
// unreachable — the input alone is still a valid minimal context).
|
||||
fn ccr_retrieve_scoped(blueprint: String, knowledge_refs: String, input_item: String, budget: Int) -> String {
|
||||
let acc: String = ""
|
||||
// knowledge_refs is a JSON array of query strings.
|
||||
let m: Int = json_array_len(knowledge_refs)
|
||||
let i = 0
|
||||
while i < m {
|
||||
let ref: String = json_array_get_string(knowledge_refs, i)
|
||||
let hit: String = primitive_attend(ref, 3)
|
||||
let acc = acc + "# ref:" + ref + "\n" + hit + "\n"
|
||||
let i = i + 1
|
||||
}
|
||||
// the input's own salient text also seeds retrieval
|
||||
let hit2: String = primitive_attend(input_item, 3)
|
||||
let acc = acc + "# input-context\n" + hit2 + "\n"
|
||||
// scope retrieval to ~60% of budget; the input itself gets the rest
|
||||
let retr_budget: Int = (budget * 6) / 10
|
||||
return ccr_compact(acc, retr_budget)
|
||||
}
|
||||
|
||||
// ── ccr_compile — assemble the bounded per-worker context package ─────────────
|
||||
//
|
||||
// blueprint : task blueprint name
|
||||
// knowledge_refs : JSON array of retrieval queries from the blueprint
|
||||
// input_item : THIS worker's single input (and nothing else)
|
||||
// corr_id : swarm correlation ID
|
||||
// worker_id : this worker's ID
|
||||
// scope_token : the worker's containment token (closed boundary)
|
||||
//
|
||||
// Returns a JSON package: { blueprint, corr_id, worker_id, scope_token,
|
||||
// input, knowledge, budget_tokens, compiled_tokens }. `knowledge` is compiled
|
||||
// and bounded; the package as a whole is bounded by budget.
|
||||
fn ccr_compile(blueprint: String, knowledge_refs: String, input_item: String,
|
||||
corr_id: String, worker_id: String, scope_token: String) -> String {
|
||||
let budget: Int = ccr_default_budget()
|
||||
let knowledge: String = ccr_retrieve_scoped(blueprint, knowledge_refs, input_item, budget)
|
||||
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "blueprint")
|
||||
let kv = el_list_append(kv, blueprint)
|
||||
let kv = el_list_append(kv, "corr_id")
|
||||
let kv = el_list_append(kv, corr_id)
|
||||
let kv = el_list_append(kv, "worker_id")
|
||||
let kv = el_list_append(kv, worker_id)
|
||||
let kv = el_list_append(kv, "input")
|
||||
let kv = el_list_append(kv, input_item)
|
||||
let kv = el_list_append(kv, "knowledge")
|
||||
let kv = el_list_append(kv, knowledge)
|
||||
let kv = el_list_append(kv, "budget_tokens")
|
||||
let kv = el_list_append(kv, int_to_str(budget))
|
||||
let pkg: String = json_build_object(kv)
|
||||
// stamp the scope token as a nested object, and the measured size
|
||||
let pkg2: String = json_set(pkg, "scope_token", scope_token)
|
||||
let compiled_tokens: Int = ccr_est_tokens(pkg2)
|
||||
let pkg3: String = json_set(pkg2, "compiled_tokens", int_to_str(compiled_tokens))
|
||||
return pkg3
|
||||
}
|
||||
|
||||
// ccr_within_budget — did the compiled package stay within its budget?
|
||||
// (Retrieval is bounded to 60% and the input is small; this asserts the whole
|
||||
// package is bounded — the property distribution relies on.)
|
||||
fn ccr_within_budget(pkg: String) -> Bool {
|
||||
let budget: Int = str_to_int(json_get_string(pkg, "budget_tokens"))
|
||||
let compiled: Int = str_to_int(json_get_string(pkg, "compiled_tokens"))
|
||||
// allow a small envelope for JSON framing overhead
|
||||
if compiled <= budget + 200 {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
@@ -0,0 +1,181 @@
|
||||
// containment.el — the Swarm Architecture containment rules, enforced.
|
||||
//
|
||||
// "These rules are not conventions. They are enforced by the runtime."
|
||||
// (Swarm Architecture §3.2). The three rules that make bounded parallelism —
|
||||
// and therefore location-independent distribution — safe:
|
||||
//
|
||||
// Rule 1: a worker may NOT join another swarm.
|
||||
// Rule 2: a worker may NOT initiate a new swarm.
|
||||
// Rule 3: a worker may NOT communicate laterally with sibling workers.
|
||||
//
|
||||
// Enforcement is by SCOPE TOKEN. When a swarm fans out, the coordinator mints a
|
||||
// swarm scope token and stamps a distinct worker scope token into each worker's
|
||||
// task envelope. Any attempt to create or join a swarm checks the caller's
|
||||
// token: if the caller already holds a WORKER token, the operation is rejected.
|
||||
// Rule 3 is enforced structurally elsewhere — workers share no mutable state and
|
||||
// the only channels they hold are the vertical result path — but this module
|
||||
// provides the explicit lateral-edge check for the execution tree.
|
||||
//
|
||||
// A scope token is a JSON object: {"kind":"coordinator|worker","swarm":"<corr>",
|
||||
// "worker":"<id-or-empty>","depth":"<n>"}.
|
||||
|
||||
// ── Token minting ────────────────────────────────────────────────────────────
|
||||
|
||||
// CAPABILITIES. A scope token carries a `caps` set — the authority it holds.
|
||||
// This is an AUTHORITY gate, not a health gate: capability is decided at mint
|
||||
// time and cannot be acquired at runtime. Engram-WRITE (op_write/op_relate/
|
||||
// op_supersede -> POST /api/nodes, /api/edges, DELETE) and dharma_emit are
|
||||
// @manager-ONLY capabilities — exactly the VBD rule that only the orchestrator
|
||||
// mutates global state. The orchestrator's token carries them; a worker's token
|
||||
// NEVER does. A worker is therefore STRUCTURALLY UNABLE to mutate global engram
|
||||
// state, regardless of engram health.
|
||||
fn cap_orchestrator() -> String { return "engram:read,engram:write,dharma:emit,state:write" }
|
||||
fn cap_worker() -> String { return "engram:read" }
|
||||
|
||||
// containment_coordinator_token — the token the orchestrator (@manager) holds.
|
||||
// Depth 0. Carries the engram-WRITE + dharma-emit capabilities (@manager-only).
|
||||
fn containment_coordinator_token(corr_id: String) -> String {
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "kind")
|
||||
let kv = el_list_append(kv, "coordinator")
|
||||
let kv = el_list_append(kv, "swarm")
|
||||
let kv = el_list_append(kv, corr_id)
|
||||
let kv = el_list_append(kv, "worker")
|
||||
let kv = el_list_append(kv, "")
|
||||
let kv = el_list_append(kv, "depth")
|
||||
let kv = el_list_append(kv, "0")
|
||||
let kv = el_list_append(kv, "caps")
|
||||
let kv = el_list_append(kv, cap_orchestrator())
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// containment_worker_token — the token stamped into a worker's envelope. Depth 1.
|
||||
// A closed boundary: forbids opening/joining swarms AND carries ONLY the
|
||||
// engram:READ capability — no engram:write, no dharma:emit. Read-only against the
|
||||
// full engram; may write only its own local geometry (its returned result).
|
||||
fn containment_worker_token(corr_id: String, worker_id: String) -> String {
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "kind")
|
||||
let kv = el_list_append(kv, "worker")
|
||||
let kv = el_list_append(kv, "swarm")
|
||||
let kv = el_list_append(kv, corr_id)
|
||||
let kv = el_list_append(kv, "worker")
|
||||
let kv = el_list_append(kv, worker_id)
|
||||
let kv = el_list_append(kv, "depth")
|
||||
let kv = el_list_append(kv, "1")
|
||||
let kv = el_list_append(kv, "caps")
|
||||
let kv = el_list_append(kv, cap_worker())
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// containment_has_cap — does this token carry capability `cap`?
|
||||
fn containment_has_cap(token: String, cap: String) -> Bool {
|
||||
return str_contains(json_get_string(token, "caps"), cap)
|
||||
}
|
||||
|
||||
// ── Rule checks (return "" on allow, or a rejection reason string) ───────────
|
||||
|
||||
// containment_check_open — may the holder of `token` OPEN a new swarm?
|
||||
// Enforces Rule 2 (a worker may not initiate a new swarm). Only a coordinator
|
||||
// token, or an absent token (top-level process), may open one.
|
||||
fn containment_check_open(token: String) -> String {
|
||||
if str_eq(token, "") {
|
||||
return ""
|
||||
}
|
||||
let kind: String = json_get_string(token, "kind")
|
||||
if str_eq(kind, "worker") {
|
||||
return "CONTAINMENT rule 2: a swarm worker may not initiate a new swarm (worker=" + json_get_string(token, "worker") + " swarm=" + json_get_string(token, "swarm") + ")"
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// containment_check_join — may the holder of `token` JOIN swarm `target_corr`?
|
||||
// Enforces Rule 1 (a worker may not join another swarm). A worker already bound
|
||||
// to swarm A may not register into swarm B; and a worker may not re-join at all.
|
||||
fn containment_check_join(token: String, target_corr: String) -> String {
|
||||
if str_eq(token, "") {
|
||||
return ""
|
||||
}
|
||||
let kind: String = json_get_string(token, "kind")
|
||||
if str_eq(kind, "worker") {
|
||||
return "CONTAINMENT rule 1: a swarm worker may not join another swarm (worker=" + json_get_string(token, "worker") + " bound-swarm=" + json_get_string(token, "swarm") + " attempted-swarm=" + target_corr + ")"
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// containment_check_lateral — may `from_token` open a communication edge to a
|
||||
// sibling worker `to_worker_id`? Enforces Rule 3 (no lateral communication).
|
||||
// The only permitted edges are vertical: worker->coordinator and
|
||||
// coordinator->worker. Any worker->worker edge is rejected.
|
||||
fn containment_check_lateral(from_token: String, to_worker_id: String) -> String {
|
||||
let kind: String = json_get_string(from_token, "kind")
|
||||
if str_eq(kind, "worker") {
|
||||
if str_eq(to_worker_id, "") {
|
||||
// empty target = the coordinator (vertical) — allowed
|
||||
return ""
|
||||
}
|
||||
return "CONTAINMENT rule 3: a swarm worker may not communicate laterally with sibling workers (from=" + json_get_string(from_token, "worker") + " to=" + to_worker_id + ")"
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// containment_check_engram_write — RULE 4: only a token carrying the
|
||||
// engram:write capability (the orchestrator's) may mutate global engram state.
|
||||
// A worker token (engram:read only) is REJECTED — the authority gate. Reuses the
|
||||
// exact scope-token mechanism as Rule 2's open-denial. Returns "" on allow, or a
|
||||
// rejection reason. This is an AUTHORITY gate: it does not consult engram health.
|
||||
fn containment_check_engram_write(token: String, op: String) -> String {
|
||||
if containment_has_cap(token, "engram:write") {
|
||||
return ""
|
||||
}
|
||||
return "CONTAINMENT rule 4: engram-write is @manager-only — a worker is read-only against the engram and may not mutate global state (op=" + op + " kind=" + json_get_string(token, "kind") + " worker=" + json_get_string(token, "worker") + " caps=" + json_get_string(token, "caps") + ")"
|
||||
}
|
||||
|
||||
// containment_check_dharma_emit — the same @manager-only rule for dharma_emit,
|
||||
// grounding Rule 4 in VBD: global-state mutations (engram-write, dharma-emit) are
|
||||
// orchestrator-only, checked by the one capability mechanism.
|
||||
fn containment_check_dharma_emit(token: String) -> String {
|
||||
if containment_has_cap(token, "dharma:emit") {
|
||||
return ""
|
||||
}
|
||||
return "CONTAINMENT rule 4: dharma_emit is @manager-only (kind=" + json_get_string(token, "kind") + ")"
|
||||
}
|
||||
|
||||
// ── Enforcement helpers ──────────────────────────────────────────────────────
|
||||
|
||||
// containment_allows_open — Bool convenience over containment_check_open.
|
||||
fn containment_allows_open(token: String) -> Bool {
|
||||
return str_eq(containment_check_open(token), "")
|
||||
}
|
||||
|
||||
// containment_is_worker — is this a worker-scoped (closed-boundary) token?
|
||||
fn containment_is_worker(token: String) -> Bool {
|
||||
return str_eq(json_get_string(token, "kind"), "worker")
|
||||
}
|
||||
|
||||
// containment_guard_open — assert a swarm may be opened under this token.
|
||||
// Returns "" if allowed, or records a CONTAINMENT violation to the work-tracking
|
||||
// journal and returns the reason. Callers must abort on a non-empty return.
|
||||
fn containment_guard_open(token: String, corr_id: String) -> String {
|
||||
let reason: String = containment_check_open(token)
|
||||
if str_eq(reason, "") {
|
||||
return ""
|
||||
}
|
||||
let p: String = json_set_str("{}", "reason", reason)
|
||||
worktrack_append("containment.violation", corr_id, "open", p)
|
||||
return reason
|
||||
}
|
||||
|
||||
// containment_guard_engram_write — assert a token may mutate global engram state
|
||||
// (Rule 4). Returns "" if allowed; otherwise journals a containment.violation and
|
||||
// returns the reason. The write path MUST abort on a non-empty return.
|
||||
fn containment_guard_engram_write(token: String, corr_id: String, op: String) -> String {
|
||||
let reason: String = containment_check_engram_write(token, op)
|
||||
if str_eq(reason, "") {
|
||||
return ""
|
||||
}
|
||||
let p0: String = json_set_str("{}", "reason", reason)
|
||||
let p1: String = json_set_str(p0, "op", op)
|
||||
worktrack_append("containment.violation", corr_id, "engram-write", p1)
|
||||
return reason
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
// primitive_binding.el — THE ONE FLIP POINT.
|
||||
//
|
||||
// This file is the single seam between the swarm and the real agentic
|
||||
// primitives. Binding the reshape's decorated primitives is a one-line change
|
||||
// HERE and nothing else changes anywhere in the swarm.
|
||||
//
|
||||
// The api-reshape agent (wt/api-reshape) is wiring the primitives as DECORATED
|
||||
// El on the dharma_* event bus over the engram — think/attend/learn/ground/assert
|
||||
// become decorated fns that emit afferent events onto the bus. The moment they
|
||||
// land, flip `bound_think` (and its siblings) to call them.
|
||||
//
|
||||
// TODAY (stub fallback, compiles + runs now against :8901):
|
||||
// fn bound_think(...) { return primitive_think(ctx, instruction) }
|
||||
//
|
||||
// THE FLIP (when reshape's decorated primitives land — one line each):
|
||||
// fn bound_think(...) { return think(ctx, instruction) } // decorated, on dharma bus
|
||||
//
|
||||
// Keep the stub as fallback: `bound_think` is only reached when the seam mode is
|
||||
// "decorated" (SWARM_PRIMITIVE_SEAM=decorated). Until you flip these bodies AND
|
||||
// set that env, the harness runs entirely on the hermetic stub.
|
||||
|
||||
// bound_think — BOUND to the reshape's proven decorated `think` (op_think),
|
||||
// real cognition over the engram geometry. The worker's CCR slice carries a
|
||||
// NODE-ID anchor in ctx.input (free-text anchors return "geometry unavailable");
|
||||
// think re-origins at that node's region under the faculty and returns a real
|
||||
// 768-dim gradient.
|
||||
fn bound_think(ctx: String, instruction: String) -> String {
|
||||
let anchor: String = json_get_string(ctx, "input")
|
||||
let faculty: String = json_get_string(ctx, "faculty")
|
||||
return op_think(anchor, faculty)
|
||||
}
|
||||
|
||||
// bound_attend — BOUND to the reshape's op_attend (POST /api/attend). Needs the
|
||||
// gate-1 write-healthy clone; falls back to the read-side attend otherwise.
|
||||
fn bound_attend(query: String, limit: Int) -> String {
|
||||
if str_eq(env("SWARM_WRITE_HEALTHY"), "1") {
|
||||
return op_attend(query, "self")
|
||||
}
|
||||
return primitive_attend(query, limit)
|
||||
}
|
||||
|
||||
// bound_learn — BOUND to the reshape's op_learn (correspondence-beat). Needs the
|
||||
// gate-1 write-healthy clone; falls back to the opt-in journal-only learn.
|
||||
fn bound_learn(corr_id: String, observation: String) -> String {
|
||||
if str_eq(env("SWARM_WRITE_HEALTHY"), "1") {
|
||||
return op_learn(observation, "induce")
|
||||
}
|
||||
return primitive_learn(corr_id, observation)
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
// primitive_seam.el — the configurable primitive seam + telemetry.
|
||||
//
|
||||
// One switch selects where a worker's primitive invocation goes:
|
||||
// SWARM_PRIMITIVE_SEAM=stub (default) — hermetic in-process think.
|
||||
// SWARM_PRIMITIVE_SEAM=decorated — the reshape's decorated
|
||||
// primitives on the dharma bus
|
||||
// (see primitive_binding.el).
|
||||
//
|
||||
// Every seam invocation is an AFFERENT signal — a primitive call travelling
|
||||
// toward the manager. The seam stamps telemetry onto each thought (seam_mode +
|
||||
// one afferent tick) so the coordinator can aggregate afferent counters across
|
||||
// the swarm without any shared mutable state (containment-safe: counts ride the
|
||||
// vertical result path, not a shared bus register).
|
||||
|
||||
// seam_mode — "stub" (default) or "decorated".
|
||||
fn seam_mode() -> String {
|
||||
let m: String = env("SWARM_PRIMITIVE_SEAM")
|
||||
if str_eq(m, "decorated") {
|
||||
return "decorated"
|
||||
}
|
||||
return "stub"
|
||||
}
|
||||
|
||||
// seam_think — route a worker's `think` through the configured seam and stamp
|
||||
// telemetry. Returns the thought JSON augmented with:
|
||||
// seam_mode : which side of the seam served this call
|
||||
// afferent : "1" — one afferent primitive signal was emitted
|
||||
fn seam_think(ctx: String, instruction: String) -> String {
|
||||
let mode: String = seam_mode()
|
||||
let thought: String = ""
|
||||
if str_eq(mode, "decorated") {
|
||||
let thought = bound_think(ctx, instruction)
|
||||
} else {
|
||||
let thought = primitive_think(ctx, instruction)
|
||||
}
|
||||
let t1: String = json_set_str(thought, "seam_mode", mode)
|
||||
let t2: String = json_set_str(t1, "afferent", "1")
|
||||
return t2
|
||||
}
|
||||
|
||||
// seam_attend / seam_learn — same seam for the other primitives (used when a
|
||||
// blueprint retrieves or writes through the bus).
|
||||
fn seam_attend(query: String, limit: Int) -> String {
|
||||
if str_eq(seam_mode(), "decorated") {
|
||||
return bound_attend(query, limit)
|
||||
}
|
||||
return primitive_attend(query, limit)
|
||||
}
|
||||
|
||||
fn seam_learn(corr_id: String, observation: String) -> String {
|
||||
if str_eq(seam_mode(), "decorated") {
|
||||
return bound_learn(corr_id, observation)
|
||||
}
|
||||
return primitive_learn(corr_id, observation)
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
// primitives.el — the agentic primitive SEAM the swarm composes over.
|
||||
//
|
||||
// The swarm is orchestration OVER the five CCR primitives, not a replacement for
|
||||
// them (CCR §2, "The Five Primitives / The Execution Cycle"): a worker executes
|
||||
// its task blueprint as attend -> think -> intend -> act -> learn against its
|
||||
// compiled, bounded context.
|
||||
//
|
||||
// This file is the SEAM. The parallel API-surface reshape exposes the canonical
|
||||
// primitive tools; when it lands, bind each primitive below to the reshaped
|
||||
// implementation (see PRIMITIVE_BINDING). Until then these are thin, engram-
|
||||
// backed fallbacks so the swarm — its fan-out, containment, CCR context
|
||||
// compilation, convergence, and work-tracking — is fully exercisable today.
|
||||
//
|
||||
// Contract: every primitive takes and returns String (JSON where structured), so
|
||||
// any primitive is directly threadable via thread.el's spawn (which runs
|
||||
// top-level (String)->String El fns).
|
||||
//
|
||||
// PRIMITIVE_BINDING: to bind the reshape's real tools, replace each fallback body
|
||||
// with a call to the reshaped El fn / API endpoint. Signatures here are the
|
||||
// stable contract the swarm depends on; keep them.
|
||||
|
||||
// ── attend — retrieve the minimal relevant context for a focus ───────────────
|
||||
// Vantage-read: pull only what this focus needs from the mind. Backed by the
|
||||
// engram's spreading-activation retrieval.
|
||||
fn primitive_attend(query: String, limit: Int) -> String {
|
||||
if str_eq(query, "") {
|
||||
return "[]"
|
||||
}
|
||||
// Location-independent worker model: when an engram daemon is configured,
|
||||
// retrieve over HTTP (the worker may run anywhere). POST /api/search
|
||||
// {query,limit,_auth}. Falls back to the in-process store otherwise.
|
||||
let url: String = env("ENGRAM_URL")
|
||||
if str_eq(url, "") {
|
||||
return engram_activate(query, limit)
|
||||
}
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "query")
|
||||
let kv = el_list_append(kv, query)
|
||||
let body0: String = json_build_object(kv)
|
||||
let body1: String = json_set(body0, "limit", int_to_str(limit))
|
||||
let body2: String = json_set_str(body1, "_auth", env("ENGRAM_API_KEY"))
|
||||
return http_post(url + "/api/search", body2)
|
||||
}
|
||||
|
||||
// ── think — reason over the compiled context ─────────────────────────────────
|
||||
// In production this routes to a model (CCR dynamic model selection). Here it is
|
||||
// a deterministic, hermetic transform so swarm behaviour is testable without an
|
||||
// external model: it echoes a structured verdict derived from the context. The
|
||||
// binding point for a real model is explicit.
|
||||
fn primitive_think(compiled_ctx: String, instruction: String) -> String {
|
||||
// PRIMITIVE_BINDING: replace with the reshape's think() (model inference).
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "instruction")
|
||||
let kv = el_list_append(kv, instruction)
|
||||
let kv = el_list_append(kv, "ctx_bytes")
|
||||
let kv = el_list_append(kv, int_to_str(str_len(compiled_ctx)))
|
||||
let kv = el_list_append(kv, "conclusion")
|
||||
let kv = el_list_append(kv, "reasoned:" + instruction)
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// ── intend — form a bounded plan/decision from a thought ─────────────────────
|
||||
fn primitive_intend(thought: String) -> String {
|
||||
let concl: String = json_get_string(thought, "conclusion")
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "intent")
|
||||
let kv = el_list_append(kv, concl)
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// ── act — execute a bounded effect and return its result ─────────────────────
|
||||
// Workers defer real side-effects to the coordinator (idempotency requirement,
|
||||
// Swarm §7.3). Here act produces an artifact-shaped result the coordinator
|
||||
// collects during convergence.
|
||||
fn primitive_act(intent: String, input_item: String) -> String {
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "acted_on")
|
||||
let kv = el_list_append(kv, input_item)
|
||||
let kv = el_list_append(kv, "via")
|
||||
let kv = el_list_append(kv, json_get_string(intent, "intent"))
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// ── learn — record an observation into the mind, tagged by correlation ID ────
|
||||
// Append-only, naturally idempotent (Swarm §7.3). Best-effort: a worker that
|
||||
// cannot reach the mind still returns its result.
|
||||
fn primitive_learn(corr_id: String, observation: String) -> String {
|
||||
let url: String = env("ENGRAM_URL")
|
||||
if str_eq(url, "") {
|
||||
return ""
|
||||
}
|
||||
let content: String = "swarm-worker-obs corr=" + corr_id + " :: " + observation
|
||||
return engram_node(content, "Memory", 0.4)
|
||||
}
|
||||
@@ -0,0 +1,103 @@
|
||||
// reshape_surface.el — the api-reshape agent's PROVEN decorated primitives,
|
||||
// composed into the swarm build to bind real cognition.
|
||||
//
|
||||
// PROVENANCE: these fns are the reshape's surface at wt/api-reshape @ d4f401d
|
||||
// ("reshape: decorator-as-seam — port @route codegen, prove decorate->serve,
|
||||
// rewrite surface as decorated El"), verified live against
|
||||
// engram.cognition-20260814. Copied verbatim (read/cognition ops only) so the
|
||||
// swarm binds the REAL primitives, not a reimplementation. The write ops
|
||||
// (op_write/op_relate/op_supersede/op_ground) are intentionally NOT composed
|
||||
// here — they exercise the persist_node write path that needs the gate-1
|
||||
// write-healthy clone; the swarm's proven run is read-cognition (think/read).
|
||||
//
|
||||
// Ops route to the ENGRAM over ENGRAM_URL — pinned by THIS worktree's .nsbx-env
|
||||
// to the :8901 swarm clone (never the reshape agent's :8900). Separate clones,
|
||||
// no collision.
|
||||
|
||||
fn engram_url() -> String {
|
||||
let u: String = env("ENGRAM_URL")
|
||||
if str_eq(u, "") { return "http://127.0.0.1:8900" }
|
||||
return u
|
||||
}
|
||||
fn engram_key() -> String {
|
||||
let k: String = env("ENGRAM_API_KEY")
|
||||
if str_eq(k, "") { return "sbx-dev-api-reshape" }
|
||||
return k
|
||||
}
|
||||
fn SELF_KEY() -> String { return "kn-efeb4a5b-5aff-4759-8a97-7233099be6ee" }
|
||||
fn VALUES_KEY() -> String { return "kn-5b606390-a52d-4ca2-8e0e-eba141d13440" }
|
||||
|
||||
// self/values name -> keystone id; anything else passes through unchanged.
|
||||
fn resolve_named(v: String) -> String {
|
||||
if str_eq(v, "self") { return SELF_KEY() }
|
||||
if str_eq(v, "neuron") { return SELF_KEY() }
|
||||
if str_eq(v, "values") { return VALUES_KEY() }
|
||||
if str_eq(v, "values_hub") { return VALUES_KEY() }
|
||||
return v
|
||||
}
|
||||
|
||||
// read — THE VANTAGE-READ. Re-origin at a point + aperture -> a BOUNDED slice.
|
||||
fn op_read(vantage: String, typ: String, k: Int) -> String {
|
||||
let vid: String = resolve_named(vantage)
|
||||
if str_eq(typ, "edges") {
|
||||
return http_get(engram_url() + "/api/neighbors/" + vid)
|
||||
}
|
||||
if str_starts_with(vid, "kn-") {
|
||||
return http_get(engram_url() + "/api/neighbors/" + vid)
|
||||
}
|
||||
return http_get(engram_url() + "/api/search?q=" + url_encode(vid) + "&limit=" + int_to_str(k))
|
||||
}
|
||||
|
||||
// think — THE ONE OPERATION. anchor (node ids) steered by faculty -> gradient.
|
||||
fn op_think(seeds: String, faculty: String) -> String {
|
||||
let s: String = resolve_named(seeds)
|
||||
let f: String = if str_eq(faculty, "") { "reason" } else { faculty }
|
||||
return http_get(engram_url() + "/api/think?seeds=" + url_encode(s) + "&faculty=" + f)
|
||||
}
|
||||
|
||||
// attend — aim attention at a region. (POST — needs a write-healthy clone.)
|
||||
fn op_attend(node: String, observer: String) -> String {
|
||||
let n: String = resolve_named(node)
|
||||
let o: String = if str_eq(observer, "") { SELF_KEY() } else { resolve_named(observer) }
|
||||
let body: String = "{\"_auth\":\"" + engram_key() + "\",\"node\":\"" + n
|
||||
+ "\",\"observer\":\"" + o + "\",\"salience\":\"0.6\"}"
|
||||
return http_post_json(engram_url() + "/api/attend", body)
|
||||
}
|
||||
|
||||
fn identity_typed(t: String) -> Bool {
|
||||
if str_eq(t, "self") { return true }
|
||||
if str_eq(t, "values") { return true }
|
||||
return false
|
||||
}
|
||||
fn type_to_node_type(t: String) -> String {
|
||||
if str_eq(t, "knowledge") { return "Knowledge" }
|
||||
if str_eq(t, "artifact") { return "Artifact" }
|
||||
if str_eq(t, "backlog") { return "WorkItem" }
|
||||
if str_eq(t, "process") { return "Process" }
|
||||
if str_eq(t, "state") { return "InternalStateEvent" }
|
||||
return "Memory"
|
||||
}
|
||||
|
||||
// write — add a node (POST /api/nodes). Identity types refused. This is a
|
||||
// global-engram MUTATION — @manager-only (Rule 4); never called on a worker path.
|
||||
// (Reshape's op_write, with json_escape -> the available json_escape_string.)
|
||||
fn op_write(content: String, typ: String, importance: Float) -> String {
|
||||
if str_eq(content, "") { return "{\"error\":\"write: content required\"}" }
|
||||
if identity_typed(typ) {
|
||||
return "{\"error\":\"write type=" + typ + " is write-protected -> intentional-cultivation\"}"
|
||||
}
|
||||
let body: String = "{\"_auth\":\"" + engram_key() + "\",\"content\":\"" + json_escape_string(content)
|
||||
+ "\",\"node_type\":\"" + type_to_node_type(typ) + "\",\"tier\":\"Working\",\"importance\":"
|
||||
+ float_to_str(importance) + "}"
|
||||
return http_post_json(engram_url() + "/api/nodes", body)
|
||||
}
|
||||
|
||||
// learn — the reflexive correspondence-beat: calibrate the steering-prior.
|
||||
// (POST — needs a write-healthy clone.)
|
||||
fn op_learn(seeds: String, faculty: String) -> String {
|
||||
let s: String = resolve_named(seeds)
|
||||
let f: String = if str_eq(faculty, "") { "induce" } else { faculty }
|
||||
let body: String = "{\"_auth\":\"" + engram_key() + "\",\"seeds\":\"" + s
|
||||
+ "\",\"faculty\":\"" + f + "\",\"keystone\":\"false\"}"
|
||||
return http_post_json(engram_url() + "/api/correspondence-beat", body)
|
||||
}
|
||||
@@ -0,0 +1,483 @@
|
||||
// swarm.el — the swarm orchestrator: bounded parallel agent execution.
|
||||
//
|
||||
// Implements Swarm Architecture's single pattern — fan out, execute independently,
|
||||
// converge — on El's NATIVE concurrency (thread.el spawn/join). No external
|
||||
// orchestrator: a swarm is a coordinator (this file, the main thread) that mints
|
||||
// a correlation identity, compiles a bounded CCR context per worker, dispatches
|
||||
// workers as native pthreads, tracks every unit of work, and converges the
|
||||
// results before returning control to the parent step.
|
||||
//
|
||||
// The five properties of every swarm (Swarm §2.1) are all present:
|
||||
// parent step -> swarm_run is called from one process step
|
||||
// task blueprint -> `blueprint` name + knowledge refs, run by every worker
|
||||
// input set -> `inputs_json`, one item per worker
|
||||
// convergence -> `strategy` in config (collect|merge|vote|reduce)
|
||||
// correlation ID -> minted here, threaded through tracking + every worker
|
||||
//
|
||||
// Containment (Swarm §3) is enforced: the caller must hold a coordinator/absent
|
||||
// token to open a swarm (Rule 2), each worker is stamped a closed worker token
|
||||
// (Rules 1+3), and workers share no mutable state (the coordinator is the only
|
||||
// journal writer).
|
||||
|
||||
// ── worker entry — the top-level (String)->String fn native threads run ──────
|
||||
//
|
||||
// Every El fn compiles to a global C symbol; spawn() resolves this by name via
|
||||
// dlsym and runs it in a pthread. The envelope carries everything the worker is
|
||||
// permitted to see — its compiled context and nothing else (§9.3).
|
||||
//
|
||||
// Returns a result JSON: {worker_id, status:"completed"|"failed", output|error}.
|
||||
fn swarm_worker_entry(envelope_json: String) -> String {
|
||||
let worker_id: String = json_get_string(envelope_json, "worker_id")
|
||||
let ctx: String = json_get_raw(envelope_json, "ctx")
|
||||
|
||||
// The worker holds a CLOSED worker token (Rules 1+3): it shares no state
|
||||
// with siblings and may not open/join a swarm. That boundary is enforced at
|
||||
// the point of attempt — swarm_run rejects any swarm opened under a worker
|
||||
// token (Rule 2). A worker simply executing its blueprint is not opening a
|
||||
// swarm, so it proceeds. Its only outward edge is this returned result
|
||||
// (the vertical worker->coordinator path).
|
||||
let out: String = swarm_run_blueprint(ctx)
|
||||
// A worker reports failed iff its blueprint signalled failure. This is the
|
||||
// vertical status edge the coordinator reads during convergence (§4.3, §7).
|
||||
let bstatus: String = json_get_string(out, "blueprint_status")
|
||||
let status: String = "completed"
|
||||
if str_eq(bstatus, "failed") {
|
||||
let status = "failed"
|
||||
}
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "worker_id")
|
||||
let kv = el_list_append(kv, worker_id)
|
||||
let kv = el_list_append(kv, "status")
|
||||
let kv = el_list_append(kv, status)
|
||||
let res: String = json_build_object(kv)
|
||||
return json_set(res, "output", out)
|
||||
}
|
||||
|
||||
// swarm_run_blueprint — execute the task blueprint over a compiled context.
|
||||
// The default blueprint is the CCR execution cycle: think -> intend -> act over
|
||||
// the worker's bounded context. Specialise by dispatching on
|
||||
// json_get_string(ctx,"blueprint"). Idempotent: reads ctx, writes only its
|
||||
// returned output (§7.3).
|
||||
fn swarm_run_blueprint(ctx: String) -> String {
|
||||
let blueprint: String = json_get_string(ctx, "blueprint")
|
||||
let input_item: String = json_get_string(ctx, "input")
|
||||
let knowledge: String = json_get_string(ctx, "knowledge")
|
||||
|
||||
// classify — deterministic verdict for the `vote` convergence strategy:
|
||||
// verdict is "long" if the input has >4 chars, else "short".
|
||||
if str_eq(blueprint, "classify") {
|
||||
let verdict: String = "short"
|
||||
if str_len(input_item) > 4 {
|
||||
let verdict = "long"
|
||||
}
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "verdict")
|
||||
let kv = el_list_append(kv, verdict)
|
||||
let kv = el_list_append(kv, "blueprint_status")
|
||||
let kv = el_list_append(kv, "ok")
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// faildemo — a worker that fails on inputs beginning with "x" (exercises the
|
||||
// failure threshold + partial convergence path). Idempotent, side-effect-free.
|
||||
if str_eq(blueprint, "faildemo") {
|
||||
let st: String = "ok"
|
||||
if str_starts_with(input_item, "x") {
|
||||
let st = "failed"
|
||||
}
|
||||
return json_set_str("{}", "blueprint_status", st)
|
||||
}
|
||||
|
||||
// cognize — REAL-COGNITION blueprint. Routes think through the seam (bound to
|
||||
// op_think in decorated mode) over the worker's NODE-ID anchor, then derives a
|
||||
// vote verdict from the gradient's confidence. In stub mode there is no
|
||||
// gradient, so the verdict falls back to a deterministic slice hash — the
|
||||
// same blueprint runs green on either side of the seam.
|
||||
if str_eq(blueprint, "cognize") {
|
||||
let thought: String = seam_think(ctx, "reason over " + input_item)
|
||||
// Derive the vote verdict from the REAL gradient's support count
|
||||
// (json_get_int, since n_support is numeric). Different anchors have
|
||||
// different support -> genuine, cognition-driven vote diversity. In stub
|
||||
// mode there is no gradient (n_support -> 0) -> "uncertain".
|
||||
let nsup: Int = json_get_int(thought, "n_support")
|
||||
let verdict: String = "uncertain"
|
||||
if nsup >= 10 {
|
||||
let verdict = "confident"
|
||||
}
|
||||
let ck: [String] = el_list_empty()
|
||||
let ck = el_list_append(ck, "verdict")
|
||||
let ck = el_list_append(ck, verdict)
|
||||
let ck = el_list_append(ck, "blueprint_status")
|
||||
let ck = el_list_append(ck, "ok")
|
||||
let cout0: String = json_build_object(ck)
|
||||
let cout1: String = json_set_str(cout0, "n_support", int_to_str(nsup))
|
||||
let cout2: String = json_set_str(cout1, "seam_mode", json_get_string(thought, "seam_mode"))
|
||||
return json_set_str(cout2, "afferent", json_get_string(thought, "afferent"))
|
||||
}
|
||||
|
||||
// default (analyze_item): the CCR execution cycle think -> intend -> act,
|
||||
// with `think` routed through the CONFIGURABLE PRIMITIVE SEAM. Telemetry
|
||||
// (seam_mode + afferent tick) rides the worker's returned output.
|
||||
let instruction: String = "process input: " + input_item
|
||||
let thought: String = seam_think(ctx, instruction)
|
||||
let intent: String = primitive_intend(thought)
|
||||
let effect: String = primitive_act(intent, input_item)
|
||||
let e1: String = json_set_str(effect, "blueprint_status", "ok")
|
||||
let e2: String = json_set_str(e1, "seam_mode", json_get_string(thought, "seam_mode"))
|
||||
let e3: String = json_set_str(e2, "afferent", json_get_string(thought, "afferent"))
|
||||
return e3
|
||||
}
|
||||
|
||||
// ── native-thread fan-out, bounded by concurrency, order-preserving ──────────
|
||||
//
|
||||
// parallel_map (thread.el) spawns ALL threads at once. The swarm honours the
|
||||
// blueprint's `concurrency` cap (§5.1: a resource constraint, not a parallelism
|
||||
// constraint — all items are processed, at most N at a time) by dispatching in
|
||||
// waves of N native threads, joining each wave before the next. Results are
|
||||
// returned in input order.
|
||||
fn swarm_fanout(worker_fn: String, envelopes: [String], concurrency: Int) -> [String] {
|
||||
let n: Int = el_list_len(envelopes)
|
||||
let cap: Int = concurrency
|
||||
if cap < 1 {
|
||||
let cap = 1
|
||||
}
|
||||
let results: [String] = el_list_empty()
|
||||
let base = 0
|
||||
while base < n {
|
||||
// spawn a wave of up to `cap` workers
|
||||
let tids: [String] = el_list_empty()
|
||||
let k = 0
|
||||
while k < cap {
|
||||
let idx: Int = base + k
|
||||
if idx < n {
|
||||
let env_item: String = el_list_get(envelopes, idx)
|
||||
let tid: Int = spawn(worker_fn, env_item)
|
||||
let tids = el_list_append(tids, int_to_str(tid))
|
||||
}
|
||||
let k = k + 1
|
||||
}
|
||||
// join the wave in order
|
||||
let j = 0
|
||||
let jn: Int = el_list_len(tids)
|
||||
while j < jn {
|
||||
let tid: Int = str_to_int(el_list_get(tids, j))
|
||||
let r: String = join(tid)
|
||||
let results = el_list_append(results, r)
|
||||
let j = j + 1
|
||||
}
|
||||
let base = base + cap
|
||||
}
|
||||
return results
|
||||
}
|
||||
|
||||
// ── convergence strategies (Swarm §4.2) ──────────────────────────────────────
|
||||
|
||||
// swarm_converge_collect — ordered list, no transformation.
|
||||
fn swarm_converge_collect(results: [String]) -> String {
|
||||
let n: Int = el_list_len(results)
|
||||
let arr: String = "[]"
|
||||
let i = 0
|
||||
while i < n {
|
||||
let arr = json_array_push(arr, el_list_get(results, i))
|
||||
let i = i + 1
|
||||
}
|
||||
return arr
|
||||
}
|
||||
|
||||
// swarm_converge_merge — combine worker outputs into a single joined string.
|
||||
fn swarm_converge_merge(results: [String]) -> String {
|
||||
let n: Int = el_list_len(results)
|
||||
let merged: String = ""
|
||||
let i = 0
|
||||
while i < n {
|
||||
let out: String = json_get_raw(el_list_get(results, i), "output")
|
||||
if i > 0 {
|
||||
let merged = merged + " | "
|
||||
}
|
||||
let merged = merged + out
|
||||
let i = i + 1
|
||||
}
|
||||
return json_set_str("{}", "merged", merged)
|
||||
}
|
||||
|
||||
// swarm_converge_vote — tally a field across worker outputs, pick the majority.
|
||||
// Each worker output is expected to carry a "verdict" string field.
|
||||
fn swarm_converge_vote(results: [String]) -> String {
|
||||
let n: Int = el_list_len(results)
|
||||
// Collect verdicts (no mutable tally: json_set can't update an existing key
|
||||
// and there is no el_list_set). Then count each verdict by rescanning.
|
||||
let verdicts: [String] = el_list_empty()
|
||||
let i = 0
|
||||
while i < n {
|
||||
let out: String = json_get_raw(el_list_get(results, i), "output")
|
||||
let v: String = json_get_string(out, "verdict")
|
||||
if str_eq(v, "") {
|
||||
let i = i + 1
|
||||
} else {
|
||||
let verdicts = el_list_append(verdicts, v)
|
||||
let i = i + 1
|
||||
}
|
||||
}
|
||||
// pick the verdict with the highest count (first-past-the-post)
|
||||
let vn: Int = el_list_len(verdicts)
|
||||
let best: String = ""
|
||||
let bestc = 0
|
||||
let a = 0
|
||||
while a < vn {
|
||||
let cand: String = el_list_get(verdicts, a)
|
||||
// count occurrences of cand
|
||||
let c = 0
|
||||
let b = 0
|
||||
while b < vn {
|
||||
if str_eq(el_list_get(verdicts, b), cand) {
|
||||
let c = c + 1
|
||||
}
|
||||
let b = b + 1
|
||||
}
|
||||
if c > bestc {
|
||||
let bestc = c
|
||||
let best = cand
|
||||
}
|
||||
let a = a + 1
|
||||
}
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "winner")
|
||||
let kv = el_list_append(kv, best)
|
||||
let kv = el_list_append(kv, "votes")
|
||||
let kv = el_list_append(kv, int_to_str(bestc))
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// swarm_converge_reduce — fold outputs into an accumulator (count + concat).
|
||||
fn swarm_converge_reduce(results: [String]) -> String {
|
||||
let n: Int = el_list_len(results)
|
||||
let acc: String = ""
|
||||
let i = 0
|
||||
while i < n {
|
||||
let out: String = json_get_raw(el_list_get(results, i), "output")
|
||||
let acc = acc + out
|
||||
let i = i + 1
|
||||
}
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "count")
|
||||
let kv = el_list_append(kv, int_to_str(n))
|
||||
let kv = el_list_append(kv, "accumulated")
|
||||
let kv = el_list_append(kv, acc)
|
||||
return json_build_object(kv)
|
||||
}
|
||||
|
||||
// ratio_to_permille — parse a decimal ratio string ("1.0", "0.8") into an
|
||||
// integer per-mille (1000, 800) so failure thresholds use exact integer math.
|
||||
// (El float division is unreliable in this runtime — int_to_float(n)/int_to_float(n)
|
||||
// does not equal 1.0 — so the swarm deliberately avoids floats.)
|
||||
fn ratio_to_permille(s: String) -> Int {
|
||||
if str_eq(s, "") {
|
||||
return 1000
|
||||
}
|
||||
let parts: [String] = str_split(s, ".")
|
||||
let whole: Int = str_to_int(el_list_get(parts, 0))
|
||||
let permille: Int = whole * 1000
|
||||
if el_list_len(parts) > 1 {
|
||||
let frac_raw: String = el_list_get(parts, 1)
|
||||
let frac3: String = str_slice(str_pad_right(frac_raw, 3, "0"), 0, 3)
|
||||
let permille = permille + str_to_int(frac3)
|
||||
}
|
||||
return permille
|
||||
}
|
||||
|
||||
// swarm_converge — dispatch on strategy name.
|
||||
fn swarm_converge(strategy: String, results: [String]) -> String {
|
||||
if str_eq(strategy, "merge") {
|
||||
return swarm_converge_merge(results)
|
||||
}
|
||||
if str_eq(strategy, "vote") {
|
||||
return swarm_converge_vote(results)
|
||||
}
|
||||
if str_eq(strategy, "reduce") {
|
||||
return swarm_converge_reduce(results)
|
||||
}
|
||||
// default: collect
|
||||
return swarm_converge_collect(results)
|
||||
}
|
||||
|
||||
// ── the ONLY global-engram write path (Rule 4, @manager-only) ────────────────
|
||||
//
|
||||
// Every engram mutation flows through here and is gated by the caller's token
|
||||
// capability. Only the orchestrator's token carries engram:write, so a worker
|
||||
// (engram:read only) calling this is DENIED by capability before any HTTP is
|
||||
// issued — structurally unable to mutate global engram state, regardless of
|
||||
// engram health. This is the curated-merge write: the orchestrator committing
|
||||
// the geometry it approved. Workers never reach a successful branch here.
|
||||
fn swarm_engram_write(token: String, corr_id: String, content: String, typ: String, importance: Float) -> String {
|
||||
let deny: String = containment_guard_engram_write(token, corr_id, "engram.write")
|
||||
if str_eq(deny, "") {
|
||||
// authorized (orchestrator) — perform the write
|
||||
let res: String = op_write(content, typ, importance)
|
||||
let new_id: String = json_get_string(res, "id")
|
||||
let cp: String = json_set_str("{}", "node_id", new_id)
|
||||
worktrack_append("swarm.committed", corr_id, "orchestrator", cp)
|
||||
return res
|
||||
}
|
||||
// denied by capability — return the rejection, no engram mutation performed
|
||||
return json_set_str("{}", "denied", deny)
|
||||
}
|
||||
|
||||
// ── the coordinator: fan out -> track -> converge ────────────────────────────
|
||||
//
|
||||
// blueprint : task blueprint name run by every worker
|
||||
// knowledge_refs : JSON array of retrieval queries for CCR compilation
|
||||
// inputs_json : JSON array of input items (one per worker)
|
||||
// config_json : { concurrency, strategy, min_success_ratio,
|
||||
// failure_action, caller_token }
|
||||
//
|
||||
// Returns: { corr_id, status:"completed"|"aborted", merged, report }.
|
||||
fn swarm_run(blueprint: String, knowledge_refs: String, inputs_json: String, config_json: String) -> String {
|
||||
let corr_id: String = "swarm-" + uuid_v4()
|
||||
let caller_token: String = json_get_raw(config_json, "caller_token")
|
||||
let concurrency: Int = str_to_int(json_get_string(config_json, "concurrency"))
|
||||
if concurrency < 1 {
|
||||
let concurrency = 4
|
||||
}
|
||||
let strategy: String = json_get_string(config_json, "strategy")
|
||||
|
||||
// ── Containment Rule 2: only a coordinator/absent token may open a swarm ──
|
||||
let deny: String = containment_guard_open(caller_token, corr_id)
|
||||
if str_eq(deny, "") {
|
||||
// allowed — proceed
|
||||
let n: Int = json_array_len(inputs_json)
|
||||
|
||||
// swarm.created
|
||||
let cp: String = json_set_str("{}", "blueprint", blueprint)
|
||||
let cp2: String = json_set(cp, "input_count", int_to_str(n))
|
||||
worktrack_append("swarm.created", corr_id, corr_id, cp2)
|
||||
|
||||
// build per-worker envelopes: worker token + CCR-compiled bounded context
|
||||
let envelopes: [String] = el_list_empty()
|
||||
let i = 0
|
||||
while i < n {
|
||||
let worker_id: String = corr_id + "/worker-" + int_to_str(i)
|
||||
let input_item: String = json_array_get_string(inputs_json, i)
|
||||
let wtoken: String = containment_worker_token(corr_id, worker_id)
|
||||
let ctx: String = ccr_compile(blueprint, knowledge_refs, input_item, corr_id, worker_id, wtoken)
|
||||
// envelope: only this worker's compiled context + its closed token
|
||||
let ekv: [String] = el_list_empty()
|
||||
let ekv = el_list_append(ekv, "worker_id")
|
||||
let ekv = el_list_append(ekv, worker_id)
|
||||
let ekv = el_list_append(ekv, "corr_id")
|
||||
let ekv = el_list_append(ekv, corr_id)
|
||||
let env0: String = json_build_object(ekv)
|
||||
let env1: String = json_set(env0, "scope_token", wtoken)
|
||||
let env2: String = json_set(env1, "ctx", ctx)
|
||||
let envelopes = el_list_append(envelopes, env2)
|
||||
|
||||
let sp: String = json_set_str("{}", "input", input_item)
|
||||
worktrack_append("worker.started", corr_id, worker_id, sp)
|
||||
let i = i + 1
|
||||
}
|
||||
|
||||
// ── native-thread fan-out (bounded) ──
|
||||
let results: [String] = swarm_fanout("swarm_worker_entry", envelopes, concurrency)
|
||||
|
||||
// record per-worker terminal status + aggregate AFFERENT telemetry.
|
||||
// Afferent counters (primitive signals travelling toward the @manager)
|
||||
// are summed from the vertical result path — no shared bus register,
|
||||
// so the aggregation is containment-safe.
|
||||
let succ = 0
|
||||
let afferent = 0
|
||||
let seam_mode_seen: String = "stub"
|
||||
let rn: Int = el_list_len(results)
|
||||
let r = 0
|
||||
while r < rn {
|
||||
let res: String = el_list_get(results, r)
|
||||
let wid: String = json_get_string(res, "worker_id")
|
||||
let st: String = json_get_string(res, "status")
|
||||
let out: String = json_get_raw(res, "output")
|
||||
let aff: Int = str_to_int(json_get_string(out, "afferent"))
|
||||
let afferent = afferent + aff
|
||||
let sm: String = json_get_string(out, "seam_mode")
|
||||
if str_eq(sm, "") {
|
||||
let seam_mode_seen = seam_mode_seen
|
||||
} else {
|
||||
let seam_mode_seen = sm
|
||||
}
|
||||
if str_eq(st, "completed") {
|
||||
let succ = succ + 1
|
||||
worktrack_append("worker.completed", corr_id, wid, json_set_str("{}", "status", "completed"))
|
||||
} else {
|
||||
worktrack_append("worker.failed", corr_id, wid, json_set_str("{}", "error", json_get_string(res, "error")))
|
||||
}
|
||||
let r = r + 1
|
||||
}
|
||||
|
||||
// swarm.converging
|
||||
let vg: String = json_set("{}", "success_count", int_to_str(succ))
|
||||
worktrack_append("swarm.converging", corr_id, corr_id, vg)
|
||||
|
||||
// swarm.telemetry — afferent counters observed by the @manager.
|
||||
let tkv: [String] = el_list_empty()
|
||||
let tkv = el_list_append(tkv, "seam_mode")
|
||||
let tkv = el_list_append(tkv, seam_mode_seen)
|
||||
let telem0: String = json_build_object(tkv)
|
||||
let telem1: String = json_set_str(telem0, "afferent_think", int_to_str(afferent))
|
||||
let telemetry: String = json_set_str(telem1, "results_received", int_to_str(rn))
|
||||
worktrack_append("swarm.telemetry", corr_id, corr_id, telemetry)
|
||||
|
||||
// ── failure threshold (Swarm §4.3), integer per-mille math ──
|
||||
// require succ/n >= min_success_ratio <=> succ*1000 >= permille*n
|
||||
let permille: Int = ratio_to_permille(json_get_string(config_json, "min_success_ratio"))
|
||||
let status: String = "completed"
|
||||
if succ * 1000 < permille * n {
|
||||
let status = "aborted"
|
||||
}
|
||||
|
||||
if str_eq(status, "aborted") {
|
||||
let ap: String = json_set_str("{}", "reason", "success ratio below min_success_ratio")
|
||||
worktrack_append("swarm.aborted", corr_id, corr_id, ap)
|
||||
let rep: String = worktrack_swarm_report(corr_id)
|
||||
let ok: [String] = el_list_empty()
|
||||
let ok = el_list_append(ok, "corr_id")
|
||||
let ok = el_list_append(ok, corr_id)
|
||||
let ok = el_list_append(ok, "status")
|
||||
let ok = el_list_append(ok, "aborted")
|
||||
let out0: String = json_build_object(ok)
|
||||
return json_set(out0, "report", rep)
|
||||
}
|
||||
|
||||
// ── converge ──
|
||||
let merged: String = swarm_converge(strategy, results)
|
||||
let dp: String = json_set_str("{}", "strategy", strategy)
|
||||
worktrack_append("swarm.completed", corr_id, corr_id, dp)
|
||||
|
||||
// ── curated merge = the ONLY engram write path (Rule 4) ──
|
||||
// With "commit":"1", the ORCHESTRATOR (its token carries engram:write)
|
||||
// commits the approved merged geometry back to the engram. This is the
|
||||
// single writer. Workers returned geometry; only the orchestrator writes.
|
||||
let commit_id: String = ""
|
||||
if str_eq(json_get_string(config_json, "commit"), "1") {
|
||||
let orch_token: String = containment_coordinator_token(corr_id)
|
||||
let cres: String = swarm_engram_write(orch_token, corr_id, "swarm-merge " + corr_id + " :: " + merged, "memory", 0.5)
|
||||
let commit_id = json_get_string(cres, "id")
|
||||
}
|
||||
|
||||
let rep2: String = worktrack_swarm_report(corr_id)
|
||||
let ok2: [String] = el_list_empty()
|
||||
let ok2 = el_list_append(ok2, "corr_id")
|
||||
let ok2 = el_list_append(ok2, corr_id)
|
||||
let ok2 = el_list_append(ok2, "status")
|
||||
let ok2 = el_list_append(ok2, "completed")
|
||||
let out1: String = json_build_object(ok2)
|
||||
let out2: String = json_set(out1, "report", rep2)
|
||||
let out3: String = json_set(out2, "merged", merged)
|
||||
let out4: String = json_set(out3, "telemetry", telemetry)
|
||||
return json_set_str(out4, "committed_node", commit_id)
|
||||
}
|
||||
// ── denied: caller was a worker trying to open a swarm (Rule 2) ──
|
||||
let dkv: [String] = el_list_empty()
|
||||
let dkv = el_list_append(dkv, "corr_id")
|
||||
let dkv = el_list_append(dkv, corr_id)
|
||||
let dkv = el_list_append(dkv, "status")
|
||||
let dkv = el_list_append(dkv, "denied")
|
||||
let dkv = el_list_append(dkv, "error")
|
||||
let dkv = el_list_append(dkv, deny)
|
||||
return json_build_object(dkv)
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
// harness_local_swarm.el — LOCAL-SWARM INTEGRATION HARNESS.
|
||||
//
|
||||
// Proves the FULL local-swarm mechanics end-to-end, TODAY, on the isolated
|
||||
// engram clone (:8901), with the primitive seam pointed at the hermetic stub.
|
||||
// The moment the api-reshape agent lands the decorated primitives on the
|
||||
// dharma bus, binding is ONE flip (primitive_binding.el) + SWARM_PRIMITIVE_SEAM=
|
||||
// decorated — this same harness then runs the bound path with no other change.
|
||||
//
|
||||
// The @manager (the coordinator) fans out N native El worker threads at real
|
||||
// concurrency, each given a CCR-scoped engram slice, each invoking the primitive
|
||||
// seam (think over its slice), enforces all three containment rules, converges
|
||||
// (vote AND reduce), work-tracks durably, and observes afferent telemetry.
|
||||
//
|
||||
// Run with the sandbox env sourced (ENGRAM_URL=:8901) to also exercise CCR
|
||||
// retrieval against the real (isolated) mind; runs fully without it too.
|
||||
|
||||
fn ok(label: String, cond: Bool, fails: Int) -> Int {
|
||||
if cond { print(" ok " + label); return fails }
|
||||
print(" FAIL " + label); return fails + 1
|
||||
}
|
||||
|
||||
fn main() -> Int {
|
||||
let fails = 0
|
||||
print("== LOCAL-SWARM INTEGRATION HARNESS (seam=" + seam_mode() + ") ==")
|
||||
|
||||
// 8 independent slices, real concurrency of 4 (2 waves of native pthreads).
|
||||
let inputs: String = "[\"billing\",\"payments\",\"ledger\",\"invoicing\",\"tax\",\"payroll\",\"audit\",\"fx\"]"
|
||||
let refs: String = "[\"Volatility-Based Decomposition\"]"
|
||||
|
||||
// ── A) fan-out / converge at real concurrency (reduce) ──
|
||||
let cfg_r: String = "{\"concurrency\":\"4\",\"strategy\":\"reduce\",\"min_success_ratio\":\"1.0\"}"
|
||||
let rr: String = swarm_run("analyze_item", refs, inputs, cfg_r)
|
||||
let fails = ok("swarm completed at concurrency=4 over 8 native-thread workers", str_eq(json_get_string(rr, "status"), "completed"), fails)
|
||||
let corr: String = json_get_string(rr, "corr_id")
|
||||
let merged_r: String = json_get_raw(rr, "merged")
|
||||
let fails = ok("reduce converged all 8 worker outputs", str_to_int(json_get_string(merged_r, "count")) == 8, fails)
|
||||
|
||||
// ── B) afferent telemetry observed by the @manager ──
|
||||
let telem: String = json_get_raw(rr, "telemetry")
|
||||
let aff: Int = str_to_int(json_get_string(telem, "afferent_think"))
|
||||
let seen_mode: String = json_get_string(telem, "seam_mode")
|
||||
let fails = ok("afferent think-signals counted = 8 (one per worker)", aff == 8, fails)
|
||||
let fails = ok("telemetry records the active seam mode", str_eq(seen_mode, seam_mode()), fails)
|
||||
let telem_recs: Int = worktrack_count_kind(corr, "swarm.telemetry")
|
||||
let fails = ok("telemetry durably journalled", telem_recs == 1, fails)
|
||||
|
||||
// ── C) CCR scoping + non-leak per worker ──
|
||||
let wt: String = containment_worker_token(corr, corr + "/worker-3")
|
||||
let ctx3: String = ccr_compile("analyze_item", refs, "invoicing", corr, corr + "/worker-3", wt)
|
||||
let fails = ok("CCR context bounded within token budget", ccr_within_budget(ctx3), fails)
|
||||
let fails = ok("CCR context carries THIS slice", str_eq(json_get_string(ctx3, "input"), "invoicing"), fails)
|
||||
let leaks: Bool = str_contains(ctx3, "payroll") || str_contains(ctx3, "audit")
|
||||
let fails = ok("CCR context does NOT leak sibling slices (security boundary)", !leaks, fails)
|
||||
|
||||
// ── D) all three containment rules ──
|
||||
let deny: String = containment_check_open(wt)
|
||||
let fails = ok("Rule 2: worker token may not OPEN a swarm", !str_eq(deny, ""), fails)
|
||||
let denyj: String = containment_check_join(wt, "other-swarm")
|
||||
let fails = ok("Rule 1: worker token may not JOIN another swarm", !str_eq(denyj, ""), fails)
|
||||
let lat: String = containment_check_lateral(wt, "sibling-9")
|
||||
let fails = ok("Rule 3: worker->worker lateral edge rejected", !str_eq(lat, ""), fails)
|
||||
let ver: String = containment_check_lateral(wt, "")
|
||||
let fails = ok("Rule 3: worker->manager vertical edge allowed", str_eq(ver, ""), fails)
|
||||
// enforced live: a worker-token caller is denied opening a real swarm
|
||||
let wcfg: String = json_set(cfg_r, "caller_token", wt)
|
||||
let denied: String = swarm_run("analyze_item", refs, inputs, wcfg)
|
||||
let fails = ok("Rule 2 enforced live: worker-caller swarm denied", str_eq(json_get_string(denied, "status"), "denied"), fails)
|
||||
|
||||
// ── E) vote convergence strategy at concurrency ──
|
||||
let cfg_v: String = "{\"concurrency\":\"8\",\"strategy\":\"vote\",\"min_success_ratio\":\"1.0\"}"
|
||||
let rv: String = swarm_run("classify", refs, inputs, cfg_v)
|
||||
let winner: String = json_get_string(json_get_raw(rv, "merged"), "winner")
|
||||
// billing/payments/ledger/invoicing/payroll/audit = long(>4); tax/fx = short -> long wins
|
||||
let fails = ok("vote converged (winner=long)", str_eq(winner, "long"), fails)
|
||||
|
||||
// ── F) durable, inspectable work-tracking ──
|
||||
let started: Int = worktrack_count_kind(corr, "worker.started")
|
||||
let completed: Int = worktrack_count_kind(corr, "worker.completed")
|
||||
let fails = ok("work-tracking journal: 8 started + 8 completed", (started == 8) && (completed == 8), fails)
|
||||
|
||||
print("")
|
||||
if fails == 0 {
|
||||
print("HARNESS GREEN — full local-swarm mechanics proven with seam=" + seam_mode())
|
||||
return 0
|
||||
}
|
||||
print("HARNESS FAIL (" + int_to_str(fails) + ")")
|
||||
return 1
|
||||
}
|
||||
@@ -0,0 +1,119 @@
|
||||
// harness_real_cognition.el — the LOCAL SWARM running REAL cognition.
|
||||
//
|
||||
// Run with: SWARM_PRIMITIVE_SEAM=decorated + the sandbox env sourced
|
||||
// (ENGRAM_URL=:8901). Each worker's `think` is BOUND to the reshape's proven
|
||||
// op_think (GET /api/think) over its NODE-ID anchor — real 768-dim gradients from
|
||||
// the live (isolated) geometry, not the stub. The @manager fans out N native-El
|
||||
// worker threads at real concurrency, converges (reduce + vote) over the real
|
||||
// cognition, enforces all three containment rules, observes afferent telemetry,
|
||||
// and work-tracks durably.
|
||||
//
|
||||
// Anchors are real self-neighbourhood node ids on the :8901 clone (free-text
|
||||
// anchors return "geometry unavailable", so these must be node ids).
|
||||
|
||||
fn ok(label: String, cond: Bool, fails: Int) -> Int {
|
||||
if cond { print(" ok " + label); return fails }
|
||||
print(" FAIL " + label); return fails + 1
|
||||
}
|
||||
|
||||
fn main() -> Int {
|
||||
let fails = 0
|
||||
print("== REAL-COGNITION LOCAL SWARM (seam=" + seam_mode() + ", engram=" + env("ENGRAM_URL") + ") ==")
|
||||
|
||||
// ── 0) direct proof the bound primitive returns REAL cognition ──
|
||||
let g: String = op_think("self", "plan")
|
||||
let dim: Int = json_get_int(g, "dim")
|
||||
let nsup: Int = json_get_int(g, "n_support")
|
||||
let fails = ok("bound op_think returns a real 768-dim gradient", dim == 768, fails)
|
||||
let fails = ok("real gradient has support (n_support>0)", nsup > 0, fails)
|
||||
let gfree: String = op_think("this-is-free-text-not-a-node", "reason")
|
||||
let fails = ok("free-text anchor correctly refused (geometry unavailable)", str_contains(gfree, "geometry unavailable"), fails)
|
||||
|
||||
// ── the input set: 8 real NODE-ID anchors from self's neighbourhood ──
|
||||
let anchors: String = "[\"a1000001-0000-0000-0000-000000000001\",\"5f011441-fa43-4fe7-a9c0-c78a584ef11d\",\"kn-5adecd7e-d6db-4576-87fe-6ef8a935cea6\",\"76d7fd0b-0672-4511-a2f5-a095cf9c60ae\",\"7027e302-593f-441d-8fd6-9c400c163108\",\"2a730b18-6566-46ee-a21e-4f4dd0380908\",\"46b0e4dd-2c19-48d2-bcbc-19f61d6c79ae\",\"9162cde8-8739-4f00-bfc9-2850ed612e50\"]"
|
||||
let refs: String = "[\"self\"]"
|
||||
|
||||
// ── A) fan-out real cognition at concurrency, converge with REDUCE ──
|
||||
let cfg_r: String = "{\"concurrency\":\"4\",\"strategy\":\"reduce\",\"min_success_ratio\":\"1.0\"}"
|
||||
let rr: String = swarm_run("cognize", refs, anchors, cfg_r)
|
||||
let fails = ok("swarm completed: 8 workers each a real think, concurrency=4", str_eq(json_get_string(rr, "status"), "completed"), fails)
|
||||
let corr: String = json_get_string(rr, "corr_id")
|
||||
let merged_r: String = json_get_raw(rr, "merged")
|
||||
let fails = ok("reduce converged all 8 real-cognition outputs", str_to_int(json_get_string(merged_r, "count")) == 8, fails)
|
||||
let acc: String = json_get_string(merged_r, "accumulated")
|
||||
let fails = ok("converged output carries real gradient support (n_support)", str_contains(acc, "n_support"), fails)
|
||||
|
||||
// ── B) afferent telemetry: 8 real think-signals, decorated seam ──
|
||||
let telem: String = json_get_raw(rr, "telemetry")
|
||||
let aff: Int = str_to_int(json_get_string(telem, "afferent_think"))
|
||||
let fails = ok("afferent counters = 8 real think invocations", aff == 8, fails)
|
||||
let fails = ok("telemetry records seam_mode=decorated", str_eq(json_get_string(telem, "seam_mode"), "decorated"), fails)
|
||||
let fails = ok("telemetry durably journalled", worktrack_count_kind(corr, "swarm.telemetry") == 1, fails)
|
||||
|
||||
// ── C) converge with VOTE over real cognition ──
|
||||
let cfg_v: String = "{\"concurrency\":\"8\",\"strategy\":\"vote\",\"min_success_ratio\":\"1.0\"}"
|
||||
let rv: String = swarm_run("cognize", refs, anchors, cfg_v)
|
||||
let winner: String = json_get_string(json_get_raw(rv, "merged"), "winner")
|
||||
let fails = ok("vote converged over real cognition (winner=" + winner + ")", !str_eq(winner, ""), fails)
|
||||
|
||||
// ── D) all three containment rules still enforced ──
|
||||
let wt: String = containment_worker_token(corr, corr + "/worker-2")
|
||||
let fails = ok("Rule 2: worker may not open a swarm", !str_eq(containment_check_open(wt), ""), fails)
|
||||
let fails = ok("Rule 1: worker may not join another swarm", !str_eq(containment_check_join(wt, "s2"), ""), fails)
|
||||
let fails = ok("Rule 3: worker->worker lateral edge rejected", !str_eq(containment_check_lateral(wt, "sib"), ""), fails)
|
||||
let wcfg: String = json_set(cfg_r, "caller_token", wt)
|
||||
let denied: String = swarm_run("cognize", refs, anchors, wcfg)
|
||||
let fails = ok("Rule 2 enforced LIVE: worker-caller swarm denied", str_eq(json_get_string(denied, "status"), "denied"), fails)
|
||||
|
||||
// ── E) CCR scoping + non-leak over node-id anchors ──
|
||||
let ctx: String = ccr_compile("cognize", refs, "a1000001-0000-0000-0000-000000000001", corr, corr + "/worker-0", wt)
|
||||
let fails = ok("CCR context bounded within budget", ccr_within_budget(ctx), fails)
|
||||
let leaks: Bool = str_contains(ctx, "9162cde8")
|
||||
let fails = ok("CCR context does NOT leak sibling anchors", !leaks, fails)
|
||||
|
||||
// ── F) durable work-tracking ──
|
||||
let started: Int = worktrack_count_kind(corr, "worker.started")
|
||||
let completed: Int = worktrack_count_kind(corr, "worker.completed")
|
||||
let fails = ok("work-tracking: 8 started + 8 completed", (started == 8) && (completed == 8), fails)
|
||||
|
||||
// ── G) RULE 4 — engram-write is @manager-ONLY (authority gate) ──
|
||||
// A worker token (engram:read only) is STRUCTURALLY denied any engram write.
|
||||
let worker_tok: String = containment_worker_token(corr, corr + "/worker-1")
|
||||
let orch_tok: String = containment_coordinator_token(corr)
|
||||
let fails = ok("worker token carries engram:read", containment_has_cap(worker_tok, "engram:read"), fails)
|
||||
let fails = ok("worker token does NOT carry engram:write", !containment_has_cap(worker_tok, "engram:write"), fails)
|
||||
let fails = ok("orchestrator token carries engram:write", containment_has_cap(orch_tok, "engram:write"), fails)
|
||||
// a worker attempting an engram write is DENIED BY CAPABILITY (no HTTP issued)
|
||||
let wdeny: String = swarm_engram_write(worker_tok, corr, "worker tries to mutate global state", "memory", 0.5)
|
||||
let denied_reason: String = json_get_string(wdeny, "denied")
|
||||
let fails = ok("worker engram-write DENIED by capability (Rule 4)", str_contains(denied_reason, "rule 4"), fails)
|
||||
let fails = ok("denied worker write performed NO engram mutation (no node id)", str_eq(json_get_string(wdeny, "id"), ""), fails)
|
||||
let fails = ok("Rule-4 violation journalled", worktrack_count_kind(corr, "containment.violation") >= 1, fails)
|
||||
// the orchestrator passes the capability gate (sole authorized writer)
|
||||
let odeny: String = containment_check_engram_write(orch_tok, "engram.write")
|
||||
let fails = ok("orchestrator PASSES the engram-write capability gate (sole writer)", str_eq(odeny, ""), fails)
|
||||
|
||||
// ── H) curated merge = the only write path (orchestrator commits) ──
|
||||
// The AUTHORITY gate above is already proven (worker denied, orchestrator
|
||||
// authorized) WITHOUT issuing a write. The actual persisting commit exercises
|
||||
// the engram write path, which needs the gate-1 write-healthy clone — so it
|
||||
// runs only under SWARM_WRITE_HEALTHY=1 (else it would hit the known daemon
|
||||
// write-crash). Authority != health: the gate holds either way.
|
||||
if str_eq(env("SWARM_WRITE_HEALTHY"), "1") {
|
||||
let cfg_commit: String = "{\"concurrency\":\"4\",\"strategy\":\"reduce\",\"min_success_ratio\":\"1.0\",\"commit\":\"1\"}"
|
||||
let rc: String = swarm_run("cognize", refs, anchors, cfg_commit)
|
||||
let committed: String = json_get_string(rc, "committed_node")
|
||||
let fails2: Int = ok("orchestrator (sole writer) committed the merge to the engram", !str_eq(committed, ""), fails)
|
||||
let fails = fails2
|
||||
} else {
|
||||
print(" note curated-merge commit deferred to the gate-1 write-healthy clone (set SWARM_WRITE_HEALTHY=1); authority gate already proven above")
|
||||
}
|
||||
|
||||
print("")
|
||||
if fails == 0 {
|
||||
print("REAL-COGNITION SWARM GREEN — Neuron thinking in parallel over its own geometry.")
|
||||
return 0
|
||||
}
|
||||
print("REAL-COGNITION SWARM FAIL (" + int_to_str(fails) + ")")
|
||||
return 1
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
// integ_engram.el — integration proof against a LIVE (isolated) engram.
|
||||
//
|
||||
// Run with the sandbox env sourced (ENGRAM_URL=http://127.0.0.1:8901,
|
||||
// ENGRAM_API_KEY=sbx-dev-swarm-ccr). Proves:
|
||||
// (a) CCR retrieval pulls REAL content from the mind over HTTP;
|
||||
// (b) a full swarm runs and converges against the live mind;
|
||||
// (c) work-tracking mirrors records into the engram as SwarmTrack nodes.
|
||||
|
||||
fn main() -> Int {
|
||||
let url: String = env("ENGRAM_URL")
|
||||
if str_eq(url, "") {
|
||||
print("SKIP integ_engram (ENGRAM_URL not set)")
|
||||
return 0
|
||||
}
|
||||
|
||||
// (a) CCR compiles a bounded context whose retrieval hit the real mind.
|
||||
let refs: String = "[\"Volatility-Based Decomposition\",\"Swarm Architecture containment\"]"
|
||||
let wt: String = containment_worker_token("integ", "integ/w0")
|
||||
let ctx: String = ccr_compile("analyze_item", refs, "decompose the billing module", "integ", "integ/w0", wt)
|
||||
let knowledge: String = json_get_string(ctx, "knowledge")
|
||||
let pulled_real: Bool = str_contains(knowledge, "olatility") || str_contains(knowledge, "Anderson") || str_contains(knowledge, "VBD")
|
||||
if pulled_real {
|
||||
print(" ok CCR retrieval pulled real mind content (" + int_to_str(str_len(knowledge)) + " bytes, bounded)")
|
||||
} else {
|
||||
print(" FAIL CCR retrieval returned no mind content")
|
||||
}
|
||||
let bounded: Bool = ccr_within_budget(ctx)
|
||||
if bounded { print(" ok compiled context stayed within budget") } else { print(" FAIL context over budget") }
|
||||
|
||||
// (b) a real swarm over the live mind.
|
||||
let inputs: String = "[\"billing\",\"payments\",\"ledger\"]"
|
||||
let cfg: String = "{\"concurrency\":\"3\",\"strategy\":\"collect\",\"min_success_ratio\":\"1.0\"}"
|
||||
let res: String = swarm_run("analyze_item", refs, inputs, cfg)
|
||||
let status: String = json_get_string(res, "status")
|
||||
if str_eq(status, "completed") { print(" ok swarm completed against live engram") } else { print(" FAIL swarm status=" + status) }
|
||||
let corr: String = json_get_string(res, "corr_id")
|
||||
|
||||
// (c) work-tracking mirrored into the mind: search for this swarm's records.
|
||||
let hits: String = primitive_attend(corr, 5)
|
||||
let mirrored: Bool = str_contains(hits, "swarm-track") || str_contains(hits, corr)
|
||||
if mirrored { print(" ok work-tracking mirrored into the engram (queryable)") } else { print(" note mirror not yet visible to search (async index)") }
|
||||
|
||||
print("DONE integ_engram corr=" + corr)
|
||||
return 0
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
// test_convergence.el — convergence strategies + failure threshold / abort.
|
||||
|
||||
fn assert_true(label: String, cond: Bool, fails: Int) -> Int {
|
||||
if cond { print(" ok " + label); return fails }
|
||||
print(" FAIL " + label); return fails + 1
|
||||
}
|
||||
|
||||
fn main() -> Int {
|
||||
let fails = 0
|
||||
let refs: String = "[]"
|
||||
|
||||
// ── vote: classify 5 inputs; 3 "long" (>4 chars) vs 2 "short" -> winner long ──
|
||||
let inputs: String = "[\"alpha\",\"bravo\",\"hi\",\"charlie\",\"ok\"]"
|
||||
let cfg_v: String = "{\"concurrency\":\"3\",\"strategy\":\"vote\",\"min_success_ratio\":\"1.0\"}"
|
||||
let rv: String = swarm_run("classify", refs, inputs, cfg_v)
|
||||
let merged_v: String = json_get_raw(rv, "merged")
|
||||
let winner: String = json_get_string(merged_v, "winner")
|
||||
let votes: Int = str_to_int(json_get_string(merged_v, "votes"))
|
||||
let fails = assert_true("vote winner = long", str_eq(winner, "long"), fails)
|
||||
let fails = assert_true("vote count = 3", votes == 3, fails)
|
||||
|
||||
// ── merge: outputs joined ──
|
||||
let cfg_m: String = "{\"concurrency\":\"2\",\"strategy\":\"merge\",\"min_success_ratio\":\"1.0\"}"
|
||||
let rm: String = swarm_run("analyze_item", refs, "[\"a\",\"b\",\"c\"]", cfg_m)
|
||||
let merged_m: String = json_get_raw(rm, "merged")
|
||||
let joined: String = json_get_string(merged_m, "merged")
|
||||
let fails = assert_true("merge produced a joined string", str_contains(joined, "|"), fails)
|
||||
|
||||
// ── reduce: count accumulates ──
|
||||
let cfg_r: String = "{\"concurrency\":\"4\",\"strategy\":\"reduce\",\"min_success_ratio\":\"1.0\"}"
|
||||
let rr: String = swarm_run("analyze_item", refs, "[\"a\",\"b\",\"c\",\"d\"]", cfg_r)
|
||||
let merged_r: String = json_get_raw(rr, "merged")
|
||||
let rcount: Int = str_to_int(json_get_string(merged_r, "count"))
|
||||
let fails = assert_true("reduce count = 4", rcount == 4, fails)
|
||||
|
||||
// ── failure threshold: 2 of 5 fail (x-prefixed); ratio 3/5=0.6 < 0.8 -> aborted ──
|
||||
let fin: String = "[\"a\",\"xb\",\"c\",\"xd\",\"e\"]"
|
||||
let cfg_f: String = "{\"concurrency\":\"5\",\"strategy\":\"collect\",\"min_success_ratio\":\"0.8\"}"
|
||||
let rf: String = swarm_run("faildemo", refs, fin, cfg_f)
|
||||
let fstatus: String = json_get_string(rf, "status")
|
||||
let fails = assert_true("swarm aborted below min_success_ratio (0.6<0.8)", str_eq(fstatus, "aborted"), fails)
|
||||
let corr_f: String = json_get_string(rf, "corr_id")
|
||||
let failed_n: Int = worktrack_count_kind(corr_f, "worker.failed")
|
||||
let aborted_n: Int = worktrack_count_kind(corr_f, "swarm.aborted")
|
||||
let fails = assert_true("tracked 2 worker.failed", failed_n == 2, fails)
|
||||
let fails = assert_true("tracked swarm.aborted", aborted_n == 1, fails)
|
||||
|
||||
// ── same failures tolerated when min_success_ratio=0.5 (0.6>=0.5) -> completed ──
|
||||
let cfg_ok: String = "{\"concurrency\":\"5\",\"strategy\":\"collect\",\"min_success_ratio\":\"0.5\"}"
|
||||
let rok: String = swarm_run("faildemo", refs, fin, cfg_ok)
|
||||
let fails = assert_true("swarm completes when failures within tolerance", str_eq(json_get_string(rok, "status"), "completed"), fails)
|
||||
|
||||
if fails == 0 { print("PASS test_convergence"); return 0 }
|
||||
print("FAIL test_convergence (" + int_to_str(fails) + ")"); return 1
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
// test_swarm.el — end-to-end proof of the swarm capability on native El threads.
|
||||
//
|
||||
// Proves: native-thread fan-out/converge, bounded concurrency, per-worker CCR
|
||||
// bounded context (with the security-boundary property), containment Rule 2
|
||||
// enforcement, and durable work-tracking.
|
||||
|
||||
fn assert_true(label: String, cond: Bool, fails: Int) -> Int {
|
||||
if cond {
|
||||
print(" ok " + label)
|
||||
return fails
|
||||
}
|
||||
print(" FAIL " + label)
|
||||
return fails + 1
|
||||
}
|
||||
|
||||
fn main() -> Int {
|
||||
let fails = 0
|
||||
|
||||
// ── 1) fan-out / converge (collect) over native threads ──
|
||||
let inputs: String = "[\"alpha\",\"bravo\",\"charlie\",\"delta\",\"echo\"]"
|
||||
let refs: String = "[]"
|
||||
let cfg: String = "{\"concurrency\":\"2\",\"strategy\":\"collect\",\"min_success_ratio\":\"1.0\"}"
|
||||
let res: String = swarm_run("analyze_item", refs, inputs, cfg)
|
||||
let status: String = json_get_string(res, "status")
|
||||
let fails = assert_true("swarm completed", str_eq(status, "completed"), fails)
|
||||
|
||||
let merged: String = json_get_raw(res, "merged")
|
||||
let count: Int = json_array_len(merged)
|
||||
let fails = assert_true("collect returned 5 results (bounded concurrency=2)", count == 5, fails)
|
||||
|
||||
// ── 2) work-tracking is durable + complete ──
|
||||
let corr: String = json_get_string(res, "corr_id")
|
||||
let started: Int = worktrack_count_kind(corr, "worker.started")
|
||||
let completed: Int = worktrack_count_kind(corr, "worker.completed")
|
||||
let created: Int = worktrack_count_kind(corr, "swarm.created")
|
||||
let done: Int = worktrack_count_kind(corr, "swarm.completed")
|
||||
let fails = assert_true("tracked 5 worker.started", started == 5, fails)
|
||||
let fails = assert_true("tracked 5 worker.completed", completed == 5, fails)
|
||||
let fails = assert_true("tracked swarm.created + swarm.completed", (created == 1) && (done == 1), fails)
|
||||
|
||||
// ── 3) CCR: bounded, minimal, non-leaking per-worker context ──
|
||||
let wtoken: String = containment_worker_token(corr, corr + "/worker-0")
|
||||
let ctx: String = ccr_compile("analyze_item", refs, "alpha", corr, corr + "/worker-0", wtoken)
|
||||
let in_budget: Bool = ccr_within_budget(ctx)
|
||||
let fails = assert_true("CCR context within token budget", in_budget, fails)
|
||||
let this_input: String = json_get_string(ctx, "input")
|
||||
let fails = assert_true("CCR context contains THIS worker's input", str_eq(this_input, "alpha"), fails)
|
||||
// security boundary: a worker's compiled context must not carry a sibling input
|
||||
let leaks_sibling: Bool = str_contains(ctx, "charlie")
|
||||
let fails = assert_true("CCR context does NOT leak sibling inputs", !leaks_sibling, fails)
|
||||
|
||||
// ── 4) containment Rule 2: a worker may not open a swarm ──
|
||||
let worker_caller_cfg: String = json_set(cfg, "caller_token", wtoken)
|
||||
let denied: String = swarm_run("analyze_item", refs, inputs, worker_caller_cfg)
|
||||
let dstatus: String = json_get_string(denied, "status")
|
||||
let fails = assert_true("worker-token caller denied opening a swarm (Rule 2)", str_eq(dstatus, "denied"), fails)
|
||||
|
||||
// coordinator token IS allowed
|
||||
let coord: String = containment_coordinator_token("some-corr")
|
||||
let allow_reason: String = containment_check_open(coord)
|
||||
let fails = assert_true("coordinator token allowed to open a swarm", str_eq(allow_reason, ""), fails)
|
||||
|
||||
// ── 5) containment Rule 3: no lateral worker->worker edge ──
|
||||
let lateral: String = containment_check_lateral(wtoken, "some-sibling")
|
||||
let fails = assert_true("lateral worker->worker edge rejected (Rule 3)", !str_eq(lateral, ""), fails)
|
||||
let vertical: String = containment_check_lateral(wtoken, "")
|
||||
let fails = assert_true("vertical worker->coordinator edge allowed", str_eq(vertical, ""), fails)
|
||||
|
||||
if fails == 0 {
|
||||
print("PASS test_swarm")
|
||||
return 0
|
||||
}
|
||||
print("FAIL test_swarm (" + int_to_str(fails) + " failures)")
|
||||
return 1
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
// test_worktrack.el — durability + inspectability of the work-tracking journal.
|
||||
|
||||
fn main() -> Int {
|
||||
let corr: String = "test-" + uuid_v4()
|
||||
|
||||
// record a swarm lifecycle
|
||||
let p1: String = json_set("{}", "input_count", "3")
|
||||
worktrack_append("swarm.created", corr, "swarm-1", p1)
|
||||
worktrack_append("worker.started", corr, "worker-001", "{}")
|
||||
worktrack_append("worker.started", corr, "worker-002", "{}")
|
||||
worktrack_append("worker.completed", corr, "worker-001", "{}")
|
||||
worktrack_append("worker.failed", corr, "worker-002", "{}")
|
||||
worktrack_append("swarm.completed", corr, "swarm-1", "{}")
|
||||
|
||||
// inspect: reconstruct the report from the durable journal
|
||||
let report: String = worktrack_swarm_report(corr)
|
||||
print("report=" + report)
|
||||
|
||||
let recs_n: Int = el_list_len(worktrack_records(corr))
|
||||
print("records=" + int_to_str(recs_n))
|
||||
|
||||
let state: String = json_get_string(report, "state")
|
||||
let completed: Int = str_to_int(json_get_string(report, "workers_completed"))
|
||||
let failed: Int = str_to_int(json_get_string(report, "workers_failed"))
|
||||
|
||||
if str_eq(state, "completed") {
|
||||
if completed == 1 {
|
||||
if failed == 1 {
|
||||
if recs_n == 6 {
|
||||
print("PASS worktrack")
|
||||
return 0
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
print("FAIL worktrack")
|
||||
return 1
|
||||
}
|
||||
@@ -0,0 +1,217 @@
|
||||
// worktrack.el — full work-tracking for the swarm.
|
||||
//
|
||||
// "Intent all the way up, orchestrator at the top." Every unit of parallel
|
||||
// work a swarm fans out is recorded here: the swarm itself, each worker, its
|
||||
// status, its result summary, the convergence, and the final merged output —
|
||||
// all threaded by a single correlation ID so the entire execution graph can be
|
||||
// reconstructed and audited (Swarm Architecture §6.1).
|
||||
//
|
||||
// DURABILITY. Records are appended to a JSON-lines journal on disk. The journal
|
||||
// is append-only and single-writer: only the coordinator (the main thread, before
|
||||
// and after each fan-out and during convergence) writes to it. Workers never
|
||||
// touch it — they return structured results and the coordinator records them.
|
||||
// This is deliberate: it makes the tracking store race-free and, not
|
||||
// coincidentally, enforces Swarm containment rule 3 (no lateral worker state).
|
||||
//
|
||||
// INSPECTABILITY. The journal is plain JSONL — greppable, tailable, replayable.
|
||||
// worktrack_read() loads it back; worktrack_swarm_report() reconstructs a
|
||||
// swarm's full record from its correlation ID.
|
||||
//
|
||||
// ENGRAM MIRROR (optional). When ENGRAM_URL is set, each record is also mirrored
|
||||
// into the engram as a node (POST /api/node) tagged with the correlation ID, so
|
||||
// the swarm's execution becomes part of the durable mind, queryable by memory.
|
||||
//
|
||||
// Depends on: el_runtime.c builtins (fs_*, http_post, env, json_*, uuid_v4,
|
||||
// now_millis, str_*). No El-module concat dependencies of its own.
|
||||
|
||||
// ── JSON helper ──────────────────────────────────────────────────────────────
|
||||
// json_set inserts its value as a RAW JSON fragment (objects/arrays/numbers).
|
||||
// json_set_str sets a plain STRING value, correctly quoted and escaped. Use
|
||||
// json_set for nested JSON, json_set_str for strings.
|
||||
fn json_set_str(j: String, key: String, val: String) -> String {
|
||||
return json_set(j, key, "\"" + json_escape_string(val) + "\"")
|
||||
}
|
||||
|
||||
// ── Journal location ─────────────────────────────────────────────────────────
|
||||
|
||||
// worktrack_dir — directory holding the swarm journals.
|
||||
// Override with SWARM_TRACK_DIR; defaults to ./.swarm-track (relative to CWD).
|
||||
fn worktrack_dir() -> String {
|
||||
let d: String = env("SWARM_TRACK_DIR")
|
||||
if str_eq(d, "") {
|
||||
return ".swarm-track"
|
||||
}
|
||||
return d
|
||||
}
|
||||
|
||||
// worktrack_journal_path — the JSONL journal file for one correlation ID.
|
||||
fn worktrack_journal_path(corr_id: String) -> String {
|
||||
return worktrack_dir() + "/" + corr_id + ".jsonl"
|
||||
}
|
||||
|
||||
// worktrack_init — ensure the journal directory exists. Idempotent.
|
||||
fn worktrack_init() -> Bool {
|
||||
let d: String = worktrack_dir()
|
||||
if fs_exists(d) {
|
||||
return true
|
||||
}
|
||||
return fs_mkdir(d)
|
||||
}
|
||||
|
||||
// ── Record construction ──────────────────────────────────────────────────────
|
||||
|
||||
// worktrack_record — build one journal record as a JSON object string.
|
||||
// kind: the record kind (swarm.created, worker.started, ...)
|
||||
// corr_id: the swarm correlation ID (links every record)
|
||||
// subject: the entity the record is about (swarm id, worker id, "")
|
||||
// payload: a JSON object string with kind-specific fields
|
||||
fn worktrack_record(kind: String, corr_id: String, subject: String, payload: String) -> String {
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "kind")
|
||||
let kv = el_list_append(kv, kind)
|
||||
let kv = el_list_append(kv, "corr_id")
|
||||
let kv = el_list_append(kv, corr_id)
|
||||
let kv = el_list_append(kv, "subject")
|
||||
let kv = el_list_append(kv, subject)
|
||||
let kv = el_list_append(kv, "ts_ms")
|
||||
let kv = el_list_append(kv, int_to_str(now_millis()))
|
||||
let rec: String = json_build_object(kv)
|
||||
// Attach the payload as a nested raw JSON field.
|
||||
let rec2: String = json_set(rec, "data", payload)
|
||||
return rec2
|
||||
}
|
||||
|
||||
// ── Journal append (single-writer, durable) ──────────────────────────────────
|
||||
|
||||
// worktrack_append — append one record to the correlation journal (durable),
|
||||
// and mirror it to the engram if ENGRAM_URL is configured. Returns the record.
|
||||
//
|
||||
// fs_write here is used in append semantics: we read-modify-write the file. The
|
||||
// coordinator is the only writer, so this is safe and race-free.
|
||||
fn worktrack_append(kind: String, corr_id: String, subject: String, payload: String) -> String {
|
||||
worktrack_init()
|
||||
let rec: String = worktrack_record(kind, corr_id, subject, payload)
|
||||
let path: String = worktrack_journal_path(corr_id)
|
||||
let prior: String = ""
|
||||
if fs_exists(path) {
|
||||
let prior = fs_read(path)
|
||||
}
|
||||
let next: String = prior + rec + "\n"
|
||||
fs_write(path, next)
|
||||
worktrack_mirror_engram(rec, corr_id, kind, subject)
|
||||
return rec
|
||||
}
|
||||
|
||||
// worktrack_mirror_engram — best-effort mirror of a record into the engram.
|
||||
// No-op unless ENGRAM_URL is set. Failures are swallowed (tracking must not
|
||||
// depend on the mind being reachable).
|
||||
fn worktrack_mirror_engram(rec: String, corr_id: String, kind: String, subject: String) -> Bool {
|
||||
// Opt-in: the durable substrate is the JSONL journal (always written). The
|
||||
// engram mirror is an additional convenience, enabled with SWARM_MIRROR=1,
|
||||
// so a swarm never depends on — or loads — the mind just to track its work.
|
||||
if str_eq(env("SWARM_MIRROR"), "1") {
|
||||
// enabled — fall through to the mirror POST
|
||||
let _go: Int = 1
|
||||
} else {
|
||||
return false
|
||||
}
|
||||
let url: String = env("ENGRAM_URL")
|
||||
if str_eq(url, "") {
|
||||
return false
|
||||
}
|
||||
let content: String = "swarm-track " + kind + " " + subject + " :: " + rec
|
||||
let body_kv: [String] = el_list_empty()
|
||||
let body_kv = el_list_append(body_kv, "content")
|
||||
let body_kv = el_list_append(body_kv, content)
|
||||
let body_kv = el_list_append(body_kv, "node_type")
|
||||
let body_kv = el_list_append(body_kv, "SwarmTrack")
|
||||
let body_kv = el_list_append(body_kv, "salience")
|
||||
let body_kv = el_list_append(body_kv, "0.5")
|
||||
let body: String = json_build_object(body_kv)
|
||||
let key: String = env("ENGRAM_API_KEY")
|
||||
let body2: String = json_set_str(body, "_auth", key)
|
||||
let resp: String = http_post(url + "/api/nodes", body2)
|
||||
return true
|
||||
}
|
||||
|
||||
// ── Read / inspect ───────────────────────────────────────────────────────────
|
||||
|
||||
// worktrack_read — read the raw JSONL journal for a correlation ID.
|
||||
fn worktrack_read(corr_id: String) -> String {
|
||||
let path: String = worktrack_journal_path(corr_id)
|
||||
if fs_exists(path) {
|
||||
return fs_read(path)
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// worktrack_records — the journal as a [String] of record JSON objects, in order.
|
||||
fn worktrack_records(corr_id: String) -> [String] {
|
||||
let raw: String = worktrack_read(corr_id)
|
||||
let out: [String] = el_list_empty()
|
||||
if str_eq(raw, "") {
|
||||
return out
|
||||
}
|
||||
let lines: [String] = str_split_lines(raw)
|
||||
let n: Int = el_list_len(lines)
|
||||
let i = 0
|
||||
while i < n {
|
||||
let ln: String = el_list_get(lines, i)
|
||||
if str_eq(ln, "") {
|
||||
let i = i + 1
|
||||
} else {
|
||||
let out = el_list_append(out, ln)
|
||||
let i = i + 1
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// worktrack_count_kind — how many records of a given kind exist for a swarm.
|
||||
// Powers assertions and live status ("how many workers completed").
|
||||
fn worktrack_count_kind(corr_id: String, kind: String) -> Int {
|
||||
let recs: [String] = worktrack_records(corr_id)
|
||||
let n: Int = el_list_len(recs)
|
||||
let c = 0
|
||||
let i = 0
|
||||
while i < n {
|
||||
let r: String = el_list_get(recs, i)
|
||||
let k: String = json_get_string(r, "kind")
|
||||
if str_eq(k, kind) {
|
||||
let c = c + 1
|
||||
}
|
||||
let i = i + 1
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// worktrack_swarm_report — reconstruct a compact status report for a swarm from
|
||||
// its journal: counts of started/completed/failed workers and terminal state.
|
||||
// Inspectable, durable, derived purely from the append-only record.
|
||||
fn worktrack_swarm_report(corr_id: String) -> String {
|
||||
let started: Int = worktrack_count_kind(corr_id, "worker.started")
|
||||
let completed: Int = worktrack_count_kind(corr_id, "worker.completed")
|
||||
let failed: Int = worktrack_count_kind(corr_id, "worker.failed")
|
||||
let done: Int = worktrack_count_kind(corr_id, "swarm.completed")
|
||||
let aborted: Int = worktrack_count_kind(corr_id, "swarm.aborted")
|
||||
let state: String = "running"
|
||||
if aborted > 0 {
|
||||
let state = "aborted"
|
||||
} else {
|
||||
if done > 0 {
|
||||
let state = "completed"
|
||||
}
|
||||
}
|
||||
let kv: [String] = el_list_empty()
|
||||
let kv = el_list_append(kv, "corr_id")
|
||||
let kv = el_list_append(kv, corr_id)
|
||||
let kv = el_list_append(kv, "state")
|
||||
let kv = el_list_append(kv, state)
|
||||
let kv = el_list_append(kv, "workers_started")
|
||||
let kv = el_list_append(kv, int_to_str(started))
|
||||
let kv = el_list_append(kv, "workers_completed")
|
||||
let kv = el_list_append(kv, int_to_str(completed))
|
||||
let kv = el_list_append(kv, "workers_failed")
|
||||
let kv = el_list_append(kv, int_to_str(failed))
|
||||
return json_build_object(kv)
|
||||
}
|
||||
@@ -173,4 +173,7 @@ Each ad-hoc harness becomes `nsbx run <name> …` (or `--source` build) against
|
||||
## Env knobs
|
||||
|
||||
`NSBX_ROOT`, `NSBX_PORT_BASE`, `NSBX_RSS_BOUND_MB`, `NSBX_REMERGE_THRESHOLD`,
|
||||
`NSBX_READY_TIMEOUT_SECS` (default 15 — how long `up`/`create`/`build` wait for a
|
||||
daemon to answer `/api/stats` before reporting failure; raise it if a boot is
|
||||
legitimately slow under concurrent sandbox/CPU load rather than actually broken),
|
||||
`EL_REPO` (for `elc` + runtime sources), `ENGRAM_LIVE_DATA_DIR`, `ENGRAM_LIVE_PLIST`.
|
||||
|
||||
+134
-29
@@ -32,6 +32,12 @@ EL_REPO="${EL_REPO:-$HOME/Development/neuron-technologies/foundation/el}"
|
||||
PORT_BASE="${NSBX_PORT_BASE:-8900}"
|
||||
RSS_BOUND_MB="${NSBX_RSS_BOUND_MB:-550}" # from store-fix reboot-proof (aaf13f88)
|
||||
REMERGE_THRESHOLD="${NSBX_REMERGE_THRESHOLD:-40000}"
|
||||
# readiness-poll window for start_daemon (0.5s ticks). Default unchanged (15s) —
|
||||
# but a cold boot against the full live store, under concurrent CPU contention
|
||||
# from other running sandboxes, has been observed live to take well past that.
|
||||
# Bump per-invocation with NSBX_READY_TIMEOUT_SECS if `up`/`create` reports a
|
||||
# not-ready failure but the daemon looks otherwise fine (see its logs/daemon.log).
|
||||
READY_TICKS=$(( ${NSBX_READY_TIMEOUT_SECS:-15} * 2 ))
|
||||
KEYSTONES=( "kn-efeb4a5b-5aff-4759-8a97-7233099be6ee" "kn-5b606390-a52d-4ca2-8e0e-eba141d13440" )
|
||||
# fixed probe set for retrieval-parity (stable, identity-anchored)
|
||||
PARITY_QUERIES=( "who am I" "self identity core" "engram store durability" "keystone self anchor" "grounding honesty" )
|
||||
@@ -77,7 +83,7 @@ _port_claimed(){ # is another sandbox already assigned this port?
|
||||
live_stats(){ curl -s -m5 "$LIVE_URL/api/stats" 2>/dev/null; }
|
||||
api(){ # api <name> <path> [json-body]
|
||||
local name="$1" path="$2" body="${3:-}"
|
||||
local port; port="$(mget "$name" "['port']")"; [ -n "$port" ] || die "unknown sandbox: $name"
|
||||
local port; port="$(mget "$name" "['port']")"; [ -n "$port" ] || die "unknown sandbox: $name (run: nsbx list)"
|
||||
local url="http://127.0.0.1:${port}${path}"
|
||||
if [ -n "$body" ]; then curl -s -m30 -X POST -H 'Content-Type: application/json' -d "$body" "$url"
|
||||
else curl -s -m30 "$url"; fi
|
||||
@@ -88,6 +94,19 @@ stat_field(){ printf '%s' "$1" | sed -n "s/.*\"$2\":\([0-9]*\).*/\1/p"; }
|
||||
daemon_pid(){ local f; f="$(sdir "$1")/daemon.pid"; [ -f "$f" ] && cat "$f" || true; }
|
||||
daemon_alive(){ local p; p="$(daemon_pid "$1")"; [ -n "$p" ] && kill -0 "$p" 2>/dev/null; }
|
||||
|
||||
# daemon_health <name> : prints "stopped" | "running" | "unresponsive" (to stdout).
|
||||
# "running" means the pid is alive AND /api/stats actually answered — process
|
||||
# liveness alone (daemon_alive) is not proof the HTTP server is serving; a pegged
|
||||
# or hung process still passes kill -0. Short timeout (2s) since this runs per-row
|
||||
# in `nsbx list`.
|
||||
daemon_health(){
|
||||
local name="$1"
|
||||
daemon_alive "$name" || { echo "stopped"; return 0; }
|
||||
local port; port="$(mget "$name" "['port']")"
|
||||
local s; s="$(curl -s -m2 "http://127.0.0.1:${port}/api/stats" 2>/dev/null)"
|
||||
[ -n "$s" ] && echo "running" || echo "unresponsive"
|
||||
}
|
||||
|
||||
# ---------------------------------------------------------------- elc/build ----
|
||||
find_elc(){
|
||||
command -v elc 2>/dev/null && return 0
|
||||
@@ -123,6 +142,49 @@ _build_binary(){
|
||||
ok "built: $out ($(ls -lh "$out" | awk '{print $5}'), sha $(sha "$out" | cut -c1-12))"
|
||||
}
|
||||
|
||||
# bin_built_at <path> : human-readable build timestamp. `cp -p` preserves mtime,
|
||||
# so this is the ORIGINAL build time even for binaries copied stock-prod into a
|
||||
# sandbox — not the copy time.
|
||||
bin_built_at(){ stat -f '%Sm' -t '%Y-%m-%d %H:%M:%S' "$1" 2>/dev/null || echo "unknown"; }
|
||||
|
||||
# _binary_freshness <name> : best-effort staleness note (empty string if fresh/
|
||||
# unknown — never guesses). Two cases:
|
||||
# - stock-prod: compares the sha recorded at create time against the CURRENTLY
|
||||
# configured live real binary's sha (recomputed now, not cached) — catches
|
||||
# "live prod moved on since this sandbox was cloned".
|
||||
# - branch/source/rebuilt: compares the recorded source_commit against the
|
||||
# LOCAL origin/dev ref (no fetch — reads whatever the repo already has) —
|
||||
# catches "built from a commit that predates current dev tip".
|
||||
_binary_freshness(){
|
||||
local name="$1" src; src="$(mget "$name" "['source']")"
|
||||
case "$src" in
|
||||
stock-prod:*)
|
||||
local live_bin cur_sha rec_sha
|
||||
live_bin="$(_live_real_bin)"
|
||||
[ -n "$live_bin" ] && [ -x "$live_bin" ] || return 0
|
||||
cur_sha="$(sha "$live_bin")"; rec_sha="$(mget "$name" "['binary_sha256']")"
|
||||
[ -n "$cur_sha" ] && [ -n "$rec_sha" ] && [ "$cur_sha" != "$rec_sha" ] \
|
||||
&& printf 'stale: live prod binary has moved on since this sandbox was cloned (live is now %s, sha %s) — nsbx build %s --binary %s to catch up' \
|
||||
"$(basename "$live_bin")" "${cur_sha:0:12}" "$name" "$live_bin"
|
||||
;;
|
||||
branch:*|source:*|rebuilt:*)
|
||||
local commit cur behind
|
||||
commit="$(mget "$name" "['source_commit']")"
|
||||
[ -n "$commit" ] || return 0
|
||||
cur="$(git -C "$EL_REPO" rev-parse origin/dev 2>/dev/null)" || return 0
|
||||
[ -n "$cur" ] && [ "$commit" != "$cur" ] || return 0
|
||||
git -C "$EL_REPO" merge-base --is-ancestor "$commit" "$cur" 2>/dev/null || return 0
|
||||
behind="$(git -C "$EL_REPO" rev-list --count "$commit..$cur" 2>/dev/null)"
|
||||
printf 'stale: built from %s, %s commit(s) behind local origin/dev (%s) — nsbx build %s --branch origin/dev' \
|
||||
"${commit:0:12}" "${behind:-?}" "${cur:0:12}" "$name"
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
# _source_commit <dir> : best-effort git HEAD of a source tree used to build a
|
||||
# sandbox binary, empty if not a git repo (e.g. a prebuilt --binary path has none).
|
||||
_source_commit(){ git -C "$1" rev-parse HEAD 2>/dev/null || true; }
|
||||
|
||||
# ---------------------------------------------------------------- daemon -------
|
||||
# start_daemon <name> : boots the sandbox's real engram binary on its isolated
|
||||
# port against its cloned data dir, with the SAME auto-remerge net the live soul
|
||||
@@ -133,9 +195,9 @@ start_daemon(){
|
||||
local port bin data export key
|
||||
port="$(mget "$name" "['port']")"; bin="$d/bin/engram"; data="$d/data"
|
||||
key="sbx-$name"; export="$data/.scan-export.reseed-clean.json"
|
||||
[ -x "$bin" ] || die "sandbox binary missing: $bin"
|
||||
[ -x "$bin" ] || die "sandbox binary missing: $bin (run: nsbx build $name --source DIR | --branch REF, or nsbx destroy $name && nsbx create $name to reclone stock-prod)"
|
||||
[ "$port" != "$LIVE_BIND_PORT" ] && [ "$port" != "$SOUL_PORT" ] || die "refusing forbidden port $port"
|
||||
[ -f "$data/neuron.egm" ] || die "sandbox has no cloned store: $data/neuron.egm"
|
||||
[ -f "$data/neuron.egm" ] || die "sandbox has no cloned store: $data/neuron.egm (data dir is corrupt/incomplete — run: nsbx destroy $name && nsbx create $name)"
|
||||
# HARD guard: never point a sandbox daemon at the live data dir.
|
||||
[ "$(cd "$data" && pwd -P)" != "$(cd "$LIVE_DATA_DIR" && pwd -P)" ] || die "refusing: sandbox data dir resolves to LIVE store"
|
||||
|
||||
@@ -150,11 +212,11 @@ start_daemon(){
|
||||
echo "$pid" > "$d/daemon.pid"
|
||||
# readiness poll
|
||||
local url="http://127.0.0.1:$port" i s
|
||||
for i in $(seq 1 30); do
|
||||
for i in $(seq 1 "$READY_TICKS"); do
|
||||
s="$(curl -s -m3 "$url/api/stats" 2>/dev/null)"
|
||||
[ -n "$s" ] && break; sleep 0.5
|
||||
done
|
||||
[ -n "$s" ] || { warn "daemon did not become ready (see $d/logs/daemon.log)"; return 1; }
|
||||
[ -n "$s" ] || { warn "daemon did not become ready within ${NSBX_READY_TIMEOUT_SECS:-15}s (see $d/logs/daemon.log). pid $pid may still be alive and slow to boot under load — check: lsof -iTCP:$port -P, or retry with NSBX_READY_TIMEOUT_SECS=45"; return 1; }
|
||||
ok "ready pid=$pid boot-stats: $s"
|
||||
# auto-remerge net (idempotent): match live edge population if the export is present
|
||||
if [ -f "$export" ]; then
|
||||
@@ -231,33 +293,35 @@ cmd_create(){
|
||||
info "live baseline stats: ${lstats:-<unavailable>}"
|
||||
|
||||
# ---- determine + place the runtime binary (versioned into the snapshot) ----
|
||||
local source_desc live_bin
|
||||
local source_desc live_bin source_commit=""
|
||||
live_bin="$(_live_real_bin)"
|
||||
if [ -n "$binpath" ]; then
|
||||
[ -x "$binpath" ] || die "not an executable binary: $binpath"
|
||||
cp -p "$binpath" "$d/bin/engram"; source_desc="prebuilt:$binpath"
|
||||
elif [ -n "$src" ]; then
|
||||
_build_binary "$src" "$d/bin/engram" "$d/build"; source_desc="source:$src"
|
||||
source_commit="$(_source_commit "$src")"
|
||||
elif [ -n "$branch" ]; then
|
||||
log "worktree: $repo @ $branch -> $d/build/worktree"
|
||||
git -C "$repo" worktree add --detach "$d/build/worktree" "$branch" >/dev/null 2>&1 \
|
||||
|| die "git worktree add failed ($repo @ $branch)"
|
||||
_build_binary "$d/build/worktree" "$d/bin/engram" "$d/build"; source_desc="branch:$branch@$repo"
|
||||
source_commit="$(_source_commit "$d/build/worktree")"
|
||||
else
|
||||
[ -x "$live_bin" ] || die "cannot resolve live ENGRAM_REAL_BIN: $live_bin"
|
||||
cp -p "$live_bin" "$d/bin/engram"; source_desc="stock-prod:$live_bin"
|
||||
fi
|
||||
local bin_sha; bin_sha="$(sha "$d/bin/engram")"
|
||||
info "runtime: $source_desc (sha ${bin_sha:0:12})"
|
||||
info "runtime: $source_desc (sha ${bin_sha:0:12}, built $(bin_built_at "$d/bin/engram"))"
|
||||
|
||||
# ---- write manifest ----
|
||||
python3 - "$name" "$port" "$source_desc" "$bin_sha" "$egm_sha" "$base_nodes" "$base_edges" "$(sha "$live_bin" 2>/dev/null)" <<'PY' > "$(manifest "$name")"
|
||||
python3 - "$name" "$port" "$source_desc" "$bin_sha" "$egm_sha" "$base_nodes" "$base_edges" "$(sha "$live_bin" 2>/dev/null)" "$source_commit" <<'PY' > "$(manifest "$name")"
|
||||
import json,sys,datetime
|
||||
name,port,src,binsha,egmsha,bn,be,livebinsha=sys.argv[1:9]
|
||||
name,port,src,binsha,egmsha,bn,be,livebinsha,source_commit=sys.argv[1:10]
|
||||
json.dump({
|
||||
"name":name,"port":int(port),"created_at":datetime.datetime.now(datetime.timezone.utc).isoformat(),
|
||||
"source":src,"binary_sha256":binsha,"clone_egm_sha256":egmsha,
|
||||
"live_binary_sha256":livebinsha,
|
||||
"live_binary_sha256":livebinsha,"source_commit":source_commit,
|
||||
"live_baseline":{"node_count":int(bn or 0),"edge_count":int(be or 0)},
|
||||
"keystones":["kn-efeb4a5b-5aff-4759-8a97-7233099be6ee","kn-5b606390-a52d-4ca2-8e0e-eba141d13440"]
|
||||
}, sys.stdout, indent=2)
|
||||
@@ -268,6 +332,17 @@ PY
|
||||
start_daemon "$name" || die "daemon failed to start"
|
||||
local sstats; sstats="$(sbx_stats "$name")"
|
||||
local sbn sbe; sbn="$(stat_field "$sstats" node_count)"; sbe="$(stat_field "$sstats" edge_count)"
|
||||
# a boot immediately followed by an auto-remerge can leave the daemon briefly
|
||||
# busy — retry rather than silently folding a 0/0 baseline into the manifest.
|
||||
# `validate`'s zero-loss/reboot-prove checks compare current counts >= baseline,
|
||||
# so a 0/0 baseline would make them trivially PASS regardless of real data loss.
|
||||
local _bi
|
||||
for _bi in 1 2 3 4 5; do
|
||||
[ -n "$sbn" ] && [ "$sbn" != "0" ] && break
|
||||
sleep 1
|
||||
sstats="$(sbx_stats "$name")"; sbn="$(stat_field "$sstats" node_count)"; sbe="$(stat_field "$sstats" edge_count)"
|
||||
done
|
||||
[ -z "$sbn" ] || [ "$sbn" = "0" ] && warn "sandbox stats still empty/zero after retries — recording sbx_baseline 0/0. This makes 'nsbx validate $name' zero-loss checks trivially pass; investigate before trusting a validate PASS: nsbx status $name"
|
||||
_capture_retrieval "$name" "$d/baseline/retrieval.json"
|
||||
# fold sandbox baseline into manifest
|
||||
python3 - "$(manifest "$name")" "$sbn" "$sbe" <<'PY'
|
||||
@@ -320,7 +395,12 @@ except Exception: print("[]")' 2>/dev/null)"
|
||||
# thereafter. Prod on :$LIVE_BIND_PORT/:$SOUL_PORT is unreachable from here by design.
|
||||
cmd_up(){
|
||||
local name; if [ $# -gt 0 ] && [ "${1#-}" = "$1" ]; then name="$1"; shift; else name="${USER:-dev}-dev"; fi
|
||||
if mexists "$name"; then daemon_alive "$name" || start_daemon "$name"; else cmd_create "$name" "$@"; fi
|
||||
if mexists "$name"; then
|
||||
daemon_alive "$name" || start_daemon "$name" \
|
||||
|| die "daemon did not become ready — see $(sdir "$name")/logs/daemon.log (try: nsbx up $name again once you've checked the log)"
|
||||
else
|
||||
cmd_create "$name" "$@"
|
||||
fi
|
||||
local port; port="$(mget "$name" "['port']")"
|
||||
echo >&2
|
||||
ok "your sandbox '$name' is ready at http://127.0.0.1:$port (a private copy of the mind — prod is untouchable)"
|
||||
@@ -334,34 +414,37 @@ cmd_up(){
|
||||
# it on the SAME clone + port (the code-change dev loop, in place).
|
||||
cmd_build(){
|
||||
local name="$1"; shift || true
|
||||
mexists "$name" || die "no such sandbox: $name"
|
||||
mexists "$name" || die "no such sandbox: $name (run: nsbx list — or nsbx create $name to make it)"
|
||||
local src="" branch="" repo="$EL_REPO"
|
||||
while [ $# -gt 0 ]; do case "$1" in
|
||||
--source) src="$2"; shift 2;; --branch) branch="$2"; shift 2;; --repo) repo="$2"; shift 2;;
|
||||
*) die "unknown flag: $1";; esac; done
|
||||
local d; d="$(sdir "$name")"
|
||||
stop_daemon "$name"
|
||||
if [ -n "$src" ]; then _build_binary "$src" "$d/bin/engram" "$d/build"
|
||||
local source_commit=""
|
||||
if [ -n "$src" ]; then _build_binary "$src" "$d/bin/engram" "$d/build"; source_commit="$(_source_commit "$src")"
|
||||
elif [ -n "$branch" ]; then
|
||||
rm -rf "$d/build/worktree" 2>/dev/null; git -C "$repo" worktree prune 2>/dev/null
|
||||
git -C "$repo" worktree add --detach "$d/build/worktree" "$branch" >/dev/null 2>&1 || die "worktree add failed"
|
||||
_build_binary "$d/build/worktree" "$d/bin/engram" "$d/build"
|
||||
source_commit="$(_source_commit "$d/build/worktree")"
|
||||
else die "usage: nsbx build <name> --source DIR | --branch REF [--repo R]"; fi
|
||||
# record new binary sha
|
||||
python3 - "$(manifest "$name")" "$(sha "$d/bin/engram")" "${src:-branch:$branch}" <<'PY'
|
||||
import json,sys; mf,s,src=sys.argv[1:4]
|
||||
d=json.load(open(mf)); d["binary_sha256"]=s; d["source"]="rebuilt:"+src
|
||||
python3 - "$(manifest "$name")" "$(sha "$d/bin/engram")" "${src:-branch:$branch}" "$source_commit" <<'PY'
|
||||
import json,sys; mf,s,src,source_commit=sys.argv[1:5]
|
||||
d=json.load(open(mf)); d["binary_sha256"]=s; d["source"]="rebuilt:"+src; d["source_commit"]=source_commit
|
||||
json.dump(d,open(mf,'w'),indent=2)
|
||||
PY
|
||||
start_daemon "$name"
|
||||
start_daemon "$name" \
|
||||
|| die "rebuilt binary did not become ready — see $d/logs/daemon.log (the old binary is gone; fix the code and re-run nsbx build $name ...)"
|
||||
ok "rebuilt + restarted on :$(mget "$name" "['port']")"
|
||||
}
|
||||
|
||||
# ================================================================ run ==========
|
||||
cmd_run(){
|
||||
local name="$1"; shift || true
|
||||
mexists "$name" || die "no such sandbox: $name"
|
||||
daemon_alive "$name" || start_daemon "$name"
|
||||
mexists "$name" || die "no such sandbox: $name (run: nsbx list — or nsbx create $name to make it)"
|
||||
daemon_alive "$name" || start_daemon "$name" || die "daemon not running and failed to start — see $(sdir "$name")/logs/daemon.log"
|
||||
local d port; d="$(sdir "$name")"; port="$(mget "$name" "['port']")"
|
||||
# direct API form: nsbx run <name> api <path> [json]
|
||||
if [ "${1:-}" = "api" ]; then
|
||||
@@ -395,8 +478,8 @@ cmd_run(){
|
||||
# RSS bound; retrieval parity; keystone integrity.
|
||||
cmd_validate(){
|
||||
local name="$1"; shift || true
|
||||
mexists "$name" || die "no such sandbox: $name"
|
||||
daemon_alive "$name" || start_daemon "$name"
|
||||
mexists "$name" || die "no such sandbox: $name (run: nsbx list — or nsbx create $name to make it)"
|
||||
daemon_alive "$name" || start_daemon "$name" || die "daemon not running and failed to start — see $(sdir "$name")/logs/daemon.log"
|
||||
local d port key; d="$(sdir "$name")"; port="$(mget "$name" "['port']")"; key="sbx-$name"
|
||||
local url="http://127.0.0.1:$port"
|
||||
local bn be; bn="$(mget "$name" "['sbx_baseline']['node_count']")"; be="$(mget "$name" "['sbx_baseline']['edge_count']")"
|
||||
@@ -489,7 +572,7 @@ PY
|
||||
# Default is a DRY-RUN plan; requires --i-approve-prod-cutover to actually cut over.
|
||||
cmd_promote(){
|
||||
local name="$1"; shift || true
|
||||
mexists "$name" || die "no such sandbox: $name"
|
||||
mexists "$name" || die "no such sandbox: $name (run: nsbx list — or nsbx create $name to make it)"
|
||||
local approve=0 do_data=0
|
||||
while [ $# -gt 0 ]; do case "$1" in
|
||||
--i-approve-prod-cutover) approve=1; shift;;
|
||||
@@ -588,7 +671,7 @@ PY
|
||||
# ================================================================ destroy ======
|
||||
cmd_destroy(){
|
||||
local name="$1"; shift || true
|
||||
mexists "$name" || die "no such sandbox: $name"
|
||||
mexists "$name" || die "no such sandbox: $name (run: nsbx list — or nsbx create $name to make it)"
|
||||
local d; d="$(sdir "$name")"
|
||||
stop_daemon "$name"
|
||||
if [ -d "$d/build/worktree" ]; then
|
||||
@@ -604,22 +687,44 @@ cmd_destroy(){
|
||||
# ================================================================ list/status ==
|
||||
cmd_list(){
|
||||
[ -d "$SBX_ROOT" ] || { echo "no sandboxes"; return 0; }
|
||||
printf '%-16s %-6s %-8s %-9s %s\n' NAME PORT STATE PID SOURCE
|
||||
printf '%-16s %-6s %-13s %-9s %-19s %s\n' NAME PORT STATE PID "BUILT" SOURCE
|
||||
local m
|
||||
for m in "$SBX_ROOT"/*/manifest.json; do
|
||||
[ -f "$m" ] || continue
|
||||
local n p src pid state
|
||||
local n p src pid state bpath built fresh
|
||||
n="$(python3 -c "import json;print(json.load(open('$m'))['name'])")"
|
||||
p="$(python3 -c "import json;print(json.load(open('$m'))['port'])")"
|
||||
src="$(python3 -c "import json;print(json.load(open('$m'))['source'])")"
|
||||
pid="$(daemon_pid "$n")"; state="stopped"; daemon_alive "$n" && state="running"
|
||||
printf '%-16s %-6s %-8s %-9s %s\n' "$n" "$p" "$state" "${pid:-–}" "$src"
|
||||
pid="$(daemon_pid "$n")"
|
||||
case "$(daemon_health "$n")" in
|
||||
running) state="running";;
|
||||
unresponsive) state="running(!resp)";;
|
||||
*) state="stopped";;
|
||||
esac
|
||||
bpath="$(sdir "$n")/bin/engram"; built="$([ -f "$bpath" ] && bin_built_at "$bpath" || echo unknown)"
|
||||
fresh="$(_binary_freshness "$n")"; [ -n "$fresh" ] && src="[STALE] $src"
|
||||
printf '%-16s %-6s %-13s %-9s %-19s %s\n' "$n" "$p" "$state" "${pid:-–}" "$built" "$src"
|
||||
done
|
||||
info "state 'running(!resp)' = process alive but /api/stats didn't answer — see: nsbx status <name>"
|
||||
}
|
||||
cmd_status(){
|
||||
local name="$1"; mexists "$name" || die "no such sandbox: $name"
|
||||
local name="$1"; mexists "$name" || die "no such sandbox: $name (run: nsbx list to see what exists)"
|
||||
python3 -m json.tool "$(manifest "$name")"
|
||||
daemon_alive "$name" && echo "state: running (pid $(daemon_pid "$name")) stats: $(sbx_stats "$name")" || echo "state: stopped"
|
||||
local bpath; bpath="$(sdir "$name")/bin/engram"
|
||||
if [ -f "$bpath" ]; then
|
||||
echo "binary: sha=$(sha "$bpath" | cut -c1-12) built=$(bin_built_at "$bpath")"
|
||||
local fresh; fresh="$(_binary_freshness "$name")"
|
||||
[ -n "$fresh" ] && printf '%s%s%s\n' "$C_YEL" "$fresh" "$C_0"
|
||||
fi
|
||||
case "$(daemon_health "$name")" in
|
||||
running)
|
||||
echo "state: running (pid $(daemon_pid "$name")) stats: $(sbx_stats "$name")";;
|
||||
unresponsive)
|
||||
printf '%sstate: running but NOT RESPONDING%s (pid %s) — process alive, /api/stats returned nothing.\n' "$C_RED" "$C_0" "$(daemon_pid "$name")"
|
||||
info "check: tail -50 $(sdir "$name")/logs/daemon.log | next: kill -9 $(daemon_pid "$name") && nsbx up $name"
|
||||
;;
|
||||
*) echo "state: stopped";;
|
||||
esac
|
||||
[ -f "$(sdir "$name")/validate.json" ] && { echo "--- last validation ---"; python3 -m json.tool "$(sdir "$name")/validate.json"; }
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user