swarm: orchestrator, CCR context compilation, containment rules, primitive seam
- swarm.el: coordinator running fan-out/converge on El NATIVE threads (thread.el spawn/join) in bounded concurrency waves, order-preserving; convergence strategies collect/merge/vote/reduce; integer per-mille failure threshold (El float division is unreliable — avoided deliberately). - ccr.el: per-worker Compiled Context Routing — retrieval/scoping/compaction into a bounded, minimal package; the compiled-context boundary is the security boundary (a worker cannot receive or leak sibling inputs). - containment.el: the three Swarm containment rules enforced via scope tokens (Rule 1 no join, Rule 2 no open, Rule 3 no lateral edge) + execution-tree lateral-edge check. - primitives.el: attend/think/intend/act/learn seam the swarm composes over, with engram-backed fallbacks and an explicit binding point for the reshape. - prototype json_array_push in el_runtime.h (defined but unprototyped). test_swarm: 12/12 — native fan-out/converge, bounded concurrency, durable tracking, CCR bounding + non-leak, and all three containment rules.
This commit is contained in:
@@ -0,0 +1,360 @@
|
||||
// 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)
|
||||
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, "completed")
|
||||
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 input_item: String = json_get_string(ctx, "input")
|
||||
let knowledge: String = json_get_string(ctx, "knowledge")
|
||||
let instruction: String = "process input: " + input_item
|
||||
let thought: String = primitive_think(knowledge, instruction)
|
||||
let intent: String = primitive_intend(thought)
|
||||
let effect: String = primitive_act(intent, input_item)
|
||||
return effect
|
||||
}
|
||||
|
||||
// ── 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("{}", "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)
|
||||
// count occurrences by scanning; first-past-the-post
|
||||
let tally: String = "{}"
|
||||
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 cur: String = json_get_string(tally, v)
|
||||
let c: Int = 0
|
||||
if str_eq(cur, "") {
|
||||
let c = 1
|
||||
} else {
|
||||
let c = str_to_int(cur) + 1
|
||||
}
|
||||
let tally = json_set(tally, v, int_to_str(c))
|
||||
let i = i + 1
|
||||
}
|
||||
}
|
||||
// pick the max
|
||||
let best: String = ""
|
||||
let bestc = 0
|
||||
let j = 0
|
||||
while j < n {
|
||||
let out: String = json_get_raw(el_list_get(results, j), "output")
|
||||
let v: String = json_get_string(out, "verdict")
|
||||
if str_eq(v, "") {
|
||||
let j = j + 1
|
||||
} else {
|
||||
let c: Int = str_to_int(json_get_string(tally, v))
|
||||
if c > bestc {
|
||||
let bestc = c
|
||||
let best = v
|
||||
}
|
||||
let j = j + 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 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("{}", "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(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("{}", "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
|
||||
let succ = 0
|
||||
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")
|
||||
if str_eq(st, "completed") {
|
||||
let succ = succ + 1
|
||||
worktrack_append("worker.completed", corr_id, wid, json_set("{}", "status", "completed"))
|
||||
} else {
|
||||
worktrack_append("worker.failed", corr_id, wid, json_set("{}", "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)
|
||||
|
||||
// ── 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("{}", "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("{}", "strategy", strategy)
|
||||
worktrack_append("swarm.completed", corr_id, corr_id, dp)
|
||||
|
||||
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)
|
||||
return json_set(out2, "merged", merged)
|
||||
}
|
||||
// ── 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)
|
||||
}
|
||||
Reference in New Issue
Block a user