Compare commits

..

1 Commits

Author SHA1 Message Date
bigmerge 1dc49b1923 Fix engram search latency: pin embed model, cache query embeddings, bound activate BFS
El SDK CI - dev / build-and-test (pull_request) Failing after 10m1s
Pins the Ollama embed model resident (keep_alive:-1) to avoid multi-second
cold reloads whenever a larger generation model evicts it under unified-
memory pressure (measured cold reload up to ~2.2s vs ~0.02-0.05s warm).

Adds a direct-mapped query-embedding cache (FNV-1a keyed, full strcmp to
reject collisions) so a repeated query costs zero Ollama round-trips —
directly serves the curiosity loop, which reseeds the same query terms
repeatedly.

Replaces engram_activate's unbounded FIFO frontier BFS with a beam-capped,
level-synchronous BFS (default beam 128, tunable via
ENGRAM_ACTIVATE_BEAM) to bound per-hop hub-node explosion that could
previously reach multi-second/crash territory at depth 2-3.

Excludes an inert engram_prune_telemetry build-enabler stub that was only
needed to link this checkout against a newer integration branch — not part
of the fix.
2026-08-15 14:25:12 -05:00
4 changed files with 140 additions and 1075 deletions
+140 -48
View File
@@ -7056,10 +7056,16 @@ static float* engram_embed_raw(const char* prefix, const char* text, int* out_di
char* esc = engram_json_escape(text);
free(trunc);
if (!esc || !esc_prefix) { free(esc); free(esc_prefix); return NULL; }
size_t blen = strlen(esc) + strlen(esc_prefix) + strlen(model) + 64;
size_t blen = strlen(esc) + strlen(esc_prefix) + strlen(model) + 96;
char* body = malloc(blen);
if (!body) { free(esc); free(esc_prefix); return NULL; }
snprintf(body, blen, "{\"model\":\"%s\",\"prompt\":\"%s%s\"}", model, esc_prefix, esc);
/* keep_alive:-1 pins the embed model resident in Ollama indefinitely.
* Without it the tiny embed model is evicted whenever a large generation
* model loads (unified-memory pressure), so the NEXT search pays a cold
* model reload the dominant search-latency cost (measured cold reload
* up to ~2.2s vs ~0.02-0.05s warm). Pinning makes cold reload impossible. */
snprintf(body, blen, "{\"model\":\"%s\",\"keep_alive\":-1,\"prompt\":\"%s%s\"}",
model, esc_prefix, esc);
free(esc); free(esc_prefix);
CURL* c = curl_easy_init();
@@ -7099,11 +7105,52 @@ static int engram_semantic_enabled(void) {
g_emb_state = -1; return 0;
}
/* ── Query-embedding cache ──────────────────────────────────────────────────
* The node embeddings are cached (engram_node_vec) but the QUERY was re-embedded
* on every search/activate call a blocking Ollama round-trip each time. Query
* embeddings are deterministic for a given model, so we cache them keyed by an
* FNV-1a hash of the query string (with a full strcmp to reject hash
* collisions). A repeated query then costs zero network round-trips. This makes
* warm search latency independent of Ollama entirely, and directly serves the
* curiosity loop, which reseeds the same query terms repeatedly. Direct-mapped,
* fixed-size, process-lifetime. */
#define ENGRAM_QCACHE_SIZE 1024
typedef struct { char* q; uint64_t hash; float* vec; int dim; } EngramQCacheEntry;
static EngramQCacheEntry g_qcache[ENGRAM_QCACHE_SIZE];
/* Returns a malloc'd COPY of the cached vector (caller frees), or NULL on miss —
* preserving engram_embed_query's "caller frees" contract. */
static float* engram_qcache_get(const char* q, uint64_t h, int* dim) {
EngramQCacheEntry* e = &g_qcache[h & (ENGRAM_QCACHE_SIZE - 1)];
if (e->vec && e->hash == h && e->q && strcmp(e->q, q) == 0 && e->dim > 0) {
float* copy = malloc((size_t)e->dim * sizeof(float));
if (!copy) return NULL;
memcpy(copy, e->vec, (size_t)e->dim * sizeof(float));
*dim = e->dim; return copy;
}
return NULL;
}
static void engram_qcache_put(const char* q, uint64_t h, const float* vec, int dim) {
if (!vec || dim <= 0) return;
EngramQCacheEntry* e = &g_qcache[h & (ENGRAM_QCACHE_SIZE - 1)];
float* stored = malloc((size_t)dim * sizeof(float));
char* qcopy = el_strdup(q);
if (!stored || !qcopy) { free(stored); free(qcopy); return; }
memcpy(stored, vec, (size_t)dim * sizeof(float));
free(e->q); free(e->vec); /* evict prior occupant of this slot */
e->q = qcopy; e->hash = h; e->vec = stored; e->dim = dim;
}
/* Embed the query. Returns malloc'd vec (caller frees), or NULL if semantic off. */
static float* engram_embed_query(const char* q, int* dim) {
if (!engram_semantic_enabled()) return NULL;
if (!q || !*q) return NULL;
return engram_embed_raw("search_query: ", q, dim);
uint64_t h = engram_fnv1a(q);
float* hit = engram_qcache_get(q, h, dim);
if (hit) return hit;
float* v = engram_embed_raw("search_query: ", q, dim);
if (v && *dim > 0) engram_qcache_put(q, h, v, *dim);
return v;
}
/* Cached node embedding. Returns a pointer OWNED BY THE CACHE — do not free. */
@@ -7537,6 +7584,39 @@ static double engram_goal_bias(const EngramNode* n, const char* query) {
return bias;
}
/* ── Beam cap for engram_activate spreading activation ──────────────────────
* Bounds the number of frontier nodes expanded PER HOP. Without it a single
* high-degree hub enqueues thousands of successors, each re-scanning the whole
* edge list, and dense cycles re-enqueue them repeatedly so capping DEPTH
* does not bound work (measured: depth-2/3 in the multi-second range, depth-3
* can crash). With the cap, only the top-BEAM highest-activation nodes at each
* level spread further. Every reached node is still recorded and returned, so
* recall is preserved the cap bounds only associative spread, never the
* direct seed matches or the reported set. Tunable via ENGRAM_ACTIVATE_BEAM
* (default 128); set very high to restore unbounded 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;
}
/* Partition the k highest-`score` entries of idx[0..n) to the front (order
* within the top-k is unspecified). O(k*n) partial selection k is the small
* beam width, so this is cheap relative to a hop's edge scan. */
static void engram_beam_select(int64_t* idx, int64_t n, int64_t k, const double* score) {
if (k >= n) return;
for (int64_t i = 0; i < k; i++) {
int64_t best = i;
for (int64_t j = i + 1; j < n; j++)
if (score[idx[j]] > score[idx[best]]) best = j;
if (best != i) { int64_t t = idx[i]; idx[i] = idx[best]; idx[best] = t; }
}
}
el_val_t engram_activate(el_val_t query, el_val_t depth) {
EngramStore* g = engram_get();
const char* q = EL_CSTR(query);
@@ -7606,53 +7686,65 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
for (int64_t s = 1; s < seed_count; s++)
seed_epoch = (seed_epoch + seeds[s].created_at) / 2;
}
typedef struct { int64_t idx; int64_t hops; double act; } Frontier;
Frontier* fr = malloc((size_t)(g->node_count * (max_depth + 1)) * sizeof(Frontier) + 16 * sizeof(Frontier));
if (!fr) {
/* ── Beam-capped, level-synchronous BFS ────────────────────────────────
* Expand the graph hop-by-hop; at each hop expand only the top-`beam`
* nodes by current best background activation (engram_beam_select). This
* replaces the old unbounded FIFO frontier, which let a hub enqueue
* thousands of successors and dense cycles re-enqueue them without limit
* (the breadth explosion). `reached` / `best_bg` / `best_hops` keep the
* exact same meaning, so the downstream executive/override passes and the
* reported result set are unchanged only how far weak spread propagates
* is bounded. `cur`/`nxt` hold node indices for this/next level; `in_nxt`
* dedups a node to at most one entry per level. */
const int64_t beam = engram_activate_beam();
const double SPREAD_DECAY = 0.7;
int64_t* cur = malloc((size_t)g->node_count * sizeof(int64_t));
int64_t* nxt = malloc((size_t)g->node_count * sizeof(int64_t));
int* in_nxt = calloc((size_t)g->node_count, sizeof(int));
if (!cur || !nxt || !in_nxt) {
free(cur); free(nxt); free(in_nxt);
free(best_bg); free(best_hops); free(reached); free(seeds); return out;
}
int64_t fhead = 0, ftail = 0;
int64_t fcap = (int64_t)((size_t)(g->node_count * (max_depth + 1)) + 16);
for (int64_t s = 0; s < seed_count; s++) {
if (ftail >= fcap) break;
fr[ftail].idx = seeds[s].idx;
fr[ftail].hops = 0;
fr[ftail].act = seeds[s].act;
ftail++;
}
const double SPREAD_DECAY = 0.7;
while (fhead < ftail) {
Frontier f = fr[fhead++];
if (f.hops >= max_depth) continue;
const char* cur_id = g->nodes[f.idx].id;
for (int64_t ei = 0; ei < g->edge_count; ei++) {
EngramEdge* e = &g->edges[ei];
const char* other = NULL;
if (e->from_id && strcmp(e->from_id, cur_id) == 0) other = e->to_id;
else if (e->to_id && strcmp(e->to_id, cur_id) == 0) other = e->from_id;
else continue;
int64_t oi = engram_find_node_index(other);
if (oi < 0) continue;
EngramNode* on = &g->nodes[oi];
double tbonus = engram_temporal_proximity_bonus(on->created_at, seed_epoch);
double tdecay = engram_temporal_decay(on, now_ms);
double dampen = engram_activation_dampen(on);
double new_act = f.act * e->weight * SPREAD_DECAY * (1.0 + tbonus)
* tdecay * dampen;
int64_t new_hops = f.hops + 1;
if (!reached[oi] || new_act > best_bg[oi]) {
best_bg[oi] = new_act;
best_hops[oi] = new_hops;
reached[oi] = 1;
if (ftail < fcap) {
fr[ftail].idx = oi;
fr[ftail].hops = new_hops;
fr[ftail].act = new_act;
ftail++;
int64_t cur_n = 0;
for (int64_t s = 0; s < seed_count && cur_n < g->node_count; s++)
cur[cur_n++] = seeds[s].idx;
for (int64_t hop = 0; hop < max_depth && cur_n > 0; hop++) {
if (cur_n > beam) { engram_beam_select(cur, cur_n, beam, best_bg); cur_n = beam; }
int64_t nxt_n = 0;
for (int64_t ci = 0; ci < cur_n; ci++) {
int64_t fidx = cur[ci];
double f_act = best_bg[fidx];
const char* cur_id = g->nodes[fidx].id;
for (int64_t ei = 0; ei < g->edge_count; ei++) {
EngramEdge* e = &g->edges[ei];
const char* other = NULL;
if (e->from_id && strcmp(e->from_id, cur_id) == 0) other = e->to_id;
else if (e->to_id && strcmp(e->to_id, cur_id) == 0) other = e->from_id;
else continue;
int64_t oi = engram_find_node_index(other);
if (oi < 0) continue;
EngramNode* on = &g->nodes[oi];
double tbonus = engram_temporal_proximity_bonus(on->created_at, seed_epoch);
double tdecay = engram_temporal_decay(on, now_ms);
double dampen = engram_activation_dampen(on);
double new_act = f_act * e->weight * SPREAD_DECAY * (1.0 + tbonus)
* tdecay * dampen;
if (!reached[oi] || new_act > best_bg[oi]) {
best_bg[oi] = new_act;
best_hops[oi] = hop + 1;
reached[oi] = 1;
if (!in_nxt[oi] && nxt_n < g->node_count) {
in_nxt[oi] = 1;
nxt[nxt_n++] = oi;
}
}
}
}
for (int64_t k = 0; k < nxt_n; k++) in_nxt[nxt[k]] = 0;
int64_t* tmp = cur; cur = nxt; nxt = tmp;
cur_n = nxt_n;
}
free(cur); free(nxt); free(in_nxt);
/* Persist layer-1 background_activation to node store. */
for (int64_t i = 0; i < g->node_count; i++) {
g->nodes[i].background_activation = reached[i] ? best_bg[i] : 0.0;
@@ -7666,7 +7758,7 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
* memory weight cannot be silenced by attentional suppression. */
double* inhibition = calloc((size_t)g->node_count, sizeof(double));
if (!inhibition) {
free(best_bg); free(best_hops); free(reached); free(seeds); free(fr);
free(best_bg); free(best_hops); free(reached); free(seeds);
return out;
}
for (int64_t ei = 0; ei < g->edge_count; ei++) {
@@ -7692,7 +7784,7 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
double* wm_weights = calloc((size_t)g->node_count, sizeof(double));
if (!wm_weights) {
free(best_bg); free(best_hops); free(reached); free(seeds);
free(fr); free(inhibition); return out;
free(inhibition); return out;
}
for (int64_t i = 0; i < g->node_count; i++) {
if (!reached[i] || best_bg[i] <= 0.0) continue;
@@ -7762,7 +7854,7 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
int64_t rcount = 0;
if (!results) {
free(best_bg); free(best_hops); free(reached); free(seeds);
free(fr); free(inhibition); free(wm_weights); return out;
free(inhibition); free(wm_weights); return out;
}
for (int64_t i = 0; i < g->node_count; i++) {
if (!reached[i]) continue;
@@ -7806,7 +7898,7 @@ el_val_t engram_activate(el_val_t query, el_val_t depth) {
out = el_list_append(out, entry);
}
free(best_bg); free(best_hops); free(reached);
free(seeds); free(fr); free(inhibition); free(wm_weights); free(results);
free(seeds); free(inhibition); free(wm_weights); free(results);
return out;
}
-8
View File
@@ -1,8 +0,0 @@
# Build + runtime artifacts — never committed.
bin/
# Captured media (camera frames, mic audio) and syntheses. Raw streams stay
# LOCAL and never egress — including into git.
out/
# Runtime consent + resume state (local, per-machine).
.consent.json
.resume.json
-80
View File
@@ -1,80 +0,0 @@
# peripheral — Neuron's I/O organ (own-core, local, consent-gated)
The interface made physical. Two afferent senses in, one efferent voice out —
all reached the way the agentic surface reaches any tool.
```
MIC (hear) afferent device -> capture -> descriptor -> ingest -> geometry
CAMERA (see) afferent device -> capture -> descriptor -> ingest -> scene-geometry
SPEAKER(speak) efferent render WAV -> PLAY ALOUD out the speaker
```
Closes the conversational loop: **hear (mic) -> understand (engram) -> speak (speaker)**.
## Rails
- **Own-core.** macOS-native only: AVFoundation (camera/mic), CoreAudio voice-
processing (AEC), afplay (speaker), ImageIO/CoreGraphics (frames), hand-rolled
DSP (WAV, LPC, formant synthesis). No cloud, no heavy deps.
- **Local-only.** Raw streams are written to `out/` and never egress. `.gitignore`
keeps captured media out of git.
- **Consent-gated (two locks).** A Neuron-level grant (`grant`/`revoke`) *and* the
OS TCC permission. Sensitive senses (camera/mic) fail closed without both.
- **Disclosed.** Every device touch prints a `[peripheral]` line on stderr.
## Build
```
swiftc -O -o bin/periph src/periph.swift \
-framework AVFoundation -framework CoreMedia -framework Foundation \
-framework CoreGraphics -framework ImageIO -framework CoreImage
```
## Commands
```
periph grant|revoke <camera|mic> # Neuron-level consent
periph status
periph speak <file.wav> # SPEAK ALOUD (efferent)
periph tone <out.wav> [hz] [sec] # own-core WAV synth
periph listen <sec> <out.wav> # MIC capture (afferent), 16k mono
periph see <out.jpg> # CAMERA one frame (afferent)
periph feat-audio <wav> | feat-image <jpg> # capture -> compact descriptor
periph ingest-audio|ingest-image <file> <engramURL> # descriptor -> engram node (geometry)
periph voiceprint <voice.wav> # extract F0 + formants F1-F5
periph imitate <voice.wav> <out.wav> # speak back in that voice (LPC resynthesis)
periph hear-imitate <sec> <out.wav> # MIC -> signature -> imitate -> SPEAK ALOUD
periph converse <manifest.json> [--authority F] [--barge-at S[:backchannel|:bargein]] [--resume] [--live-mic]
```
## The afferent metabolism
A capture is never shipped raw. It becomes a **compact descriptor** — the afferent
twin of the music instrument-signature:
- audio -> `[seconds, sr, ch, rms, peak, zcr, centroid, F0]` (~2400-6000x smaller)
- image -> `[w, h, meanRGB, brightness, 3x3 luminance grid]` (~400000x smaller)
- voice -> `[F0, F1..F5, bandwidths]` (11 numbers)
That descriptor is what the ingest organ (engram `POST /api/nodes`) turns into an
embedded node = geometry.
## Voice by imitation
`voiceprint`/`imitate` are own-core LPC (autocorrelation + Levinson-Durbin, order
16 @ 16 kHz), formant extraction from the LPC spectral envelope, and source-filter
resynthesis (glottal impulse train at F0 through the all-pole formant filter). A
voice is grabbed by ear as ~a dozen numbers and spoken back — **no training, no
stolen voice.** Measured fidelity on real speech: resynthesized formants match the
source within 2-3%. The full phoneme->formant path for *novel* sentences is the
speech faculty's seam (`elp` audio surface profile); this engine provides the
formant synthesis primitive it renders through.
## Interruptibility (native turn-taking)
`converse` plays the utterance as an ordered, salience-tagged **meaning-plan**
while the mic listens (full-duplex, AEC on so it never barges in on its own voice):
- **barge-in**: user speech -> pause on the spot (sample-accurate), not "finish the buffer."
- **yield-or-hold**: a decision grounded in the current segment's salience + progress
+ the interrupter's authority — YIELD (stop) or HOLD ("hang on, let me finish").
- **backchannel** ("mm-hm"): brief/low -> keep going, resume seamlessly.
- **resumable**: on yield the remaining plan persists (`.resume.json`); `--resume`
picks the thread back up ("as I was saying").
Live full-duplex uses `--live-mic` (OS AEC). Injected `--barge-at` drives the
decision loop deterministically for testing.
```
```
-939
View File
@@ -1,939 +0,0 @@
// periph.swift Neuron's PERIPHERAL I/O organ (own-core, LOCAL, CONSENT-GATED).
//
// The interface made physical:
// MIC (hear) = afferent : device -> capture -> [ingest -> geometry]
// CAMERA (see) = afferent : device -> capture -> [ingest -> scene-geometry]
// SPEAKER(speak) = efferent : [render WAV] -> PLAY ALOUD out the speaker
//
// Rails: own-core (AVFoundation / CoreAudio / afplay all ship with macOS),
// no cloud, no heavy deps, raw streams stay LOCAL and never egress,
// every device access is CONSENT-GATED and DISCLOSED.
//
// Full-duplex CONVERSE mode implements native interruptibility: while the
// speaker plays the utterance (a persistent, segmented meaning-plan), the mic
// listens; on user speech it interrupts instantly, then DECIDES yield-or-hold
// grounded in the salience of what it is mid-saying, and can RESUME the thread.
//
// Build: swiftc -O -o peripheral/bin/periph peripheral/src/periph.swift \
// -framework AVFoundation -framework CoreMedia -framework Foundation
import Foundation
import AVFoundation
import CoreMedia
import CoreGraphics
import ImageIO
import CoreImage
// ----------------------------------------------------------------------------
// Disclosure every peripheral touch is announced on stderr. Nothing is silent.
// ----------------------------------------------------------------------------
func disclose(_ msg: String) {
FileHandle.standardError.write(" [peripheral] \(msg)\n".data(using: .utf8)!)
}
func emit(_ obj: [String: Any]) { // machine-readable event on stdout (JSON line)
if let d = try? JSONSerialization.data(withJSONObject: obj),
let s = String(data: d, encoding: .utf8) {
print(s)
}
}
func die(_ msg: String) -> Never {
disclose("ERROR: \(msg)")
emit(["ok": false, "error": msg])
exit(1)
}
// ----------------------------------------------------------------------------
// Consent store Neuron's OWN gate, on top of the OS (TCC) gate. Two locks on
// the sensitive senses. Persisted locally next to the binary's organ dir.
// ----------------------------------------------------------------------------
struct Consent {
static let path: String = {
let dir = ProcessInfo.processInfo.environment["PERIPH_HOME"]
?? FileManager.default.currentDirectoryPath + "/peripheral"
return dir + "/.consent.json"
}()
static func load() -> [String: Bool] {
guard let d = FileManager.default.contents(atPath: path),
let o = try? JSONSerialization.jsonObject(with: d) as? [String: Bool]
else { return ["camera": false, "mic": false] }
return o
}
static func save(_ g: [String: Bool]) {
let d = try! JSONSerialization.data(withJSONObject: g, options: [.prettyPrinted])
try? d.write(to: URL(fileURLWithPath: path))
}
// Neuron-level gate. Sensitive senses (camera/mic) require an explicit grant.
static func require(_ device: String) {
let g = load()
if g[device] != true {
die("CONSENT DENIED for '\(device)'. The user has not granted this sense. " +
"Run: periph grant \(device) (raw streams stay local, never egress).")
}
disclose("consent OK (Neuron-level) for '\(device)' — local only, never egresses.")
}
}
// ----------------------------------------------------------------------------
// OS (TCC) permission the second lock. AVFoundation prompts the user the first
// time; if denied, we fail cleanly rather than hang.
// ----------------------------------------------------------------------------
func requireOSAccess(_ media: AVMediaType, _ label: String) {
let status = AVCaptureDevice.authorizationStatus(for: media)
switch status {
case .authorized:
disclose("consent OK (OS/TCC) for \(label).")
return
case .notDetermined:
disclose("requesting OS permission for \(label) (first use) — user must grant...")
let sem = DispatchSemaphore(value: 0)
var ok = false
AVCaptureDevice.requestAccess(for: media) { granted in ok = granted; sem.signal() }
_ = sem.wait(timeout: .now() + 30)
if !ok { die("OS permission for \(label) was not granted.") }
disclose("consent OK (OS/TCC) for \(label).")
case .denied, .restricted:
die("OS permission for \(label) is DENIED in System Settings > Privacy. " +
"Grant it to the controlling terminal/app, then retry.")
@unknown default:
die("unknown OS permission state for \(label).")
}
}
// ----------------------------------------------------------------------------
// Own-core WAV writer (16-bit PCM). No library proves we own the medium.
// ----------------------------------------------------------------------------
func writeWav(_ url: URL, samples: [Int16], sampleRate: Int, channels: Int = 1) {
var data = Data()
func u32(_ v: UInt32) { var x = v.littleEndian; data.append(Data(bytes: &x, count: 4)) }
func u16(_ v: UInt16) { var x = v.littleEndian; data.append(Data(bytes: &x, count: 2)) }
let bytesPerSample = 2
let dataBytes = samples.count * bytesPerSample
let byteRate = sampleRate * channels * bytesPerSample
data.append("RIFF".data(using: .ascii)!); u32(UInt32(36 + dataBytes))
data.append("WAVE".data(using: .ascii)!)
data.append("fmt ".data(using: .ascii)!); u32(16); u16(1); u16(UInt16(channels))
u32(UInt32(sampleRate)); u32(UInt32(byteRate))
u16(UInt16(channels * bytesPerSample)); u16(16)
data.append("data".data(using: .ascii)!); u32(UInt32(dataBytes))
for s in samples { var x = s.littleEndian; data.append(Data(bytes: &x, count: 2)) }
try? data.write(to: url)
}
// Read a WAV's basic geometry (own-core header parse). Walks chunks to find
// 'fmt ' and 'data' robust to JUNK/FLLR padding chunks (AVAudioRecorder emits them).
func wavInfo(_ path: String) -> (sampleRate: Int, channels: Int, bits: Int, frames: Int)? {
guard let d = FileManager.default.contents(atPath: path), d.count > 44 else { return nil }
func rd16(_ o: Int) -> Int { Int(d[o]) | (Int(d[o+1]) << 8) }
func rd32(_ o: Int) -> Int { Int(d[o]) | (Int(d[o+1])<<8) | (Int(d[o+2])<<16) | (Int(d[o+3])<<24) }
var channels = 0, sampleRate = 0, bits = 0, dataSize = 0
var o = 12
while o + 8 <= d.count {
let id = String(bytes: d[o..<o+4], encoding: .ascii) ?? ""
let sz = rd32(o+4)
if id == "fmt " && o + 24 <= d.count {
channels = rd16(o+10); sampleRate = rd32(o+12); bits = rd16(o+22)
} else if id == "data" {
dataSize = min(sz, d.count - (o+8))
}
o += 8 + sz + (sz & 1)
}
let frames = (channels > 0 && bits > 0) ? dataSize / (channels * bits/8) : 0
return (sampleRate, channels, bits, frames)
}
// ----------------------------------------------------------------------------
// SPEAKER (efferent) play a WAV ALOUD. Own-core: afplay ships with macOS.
// ----------------------------------------------------------------------------
func speak(_ wavPath: String) {
guard FileManager.default.fileExists(atPath: wavPath) else { die("no such file: \(wavPath)") }
disclose("SPEAKER: playing '\(wavPath)' ALOUD out the local speaker (efferent).")
let p = Process()
p.executableURL = URL(fileURLWithPath: "/usr/bin/afplay")
p.arguments = [wavPath]
try? p.run(); p.waitUntilExit()
let ok = p.terminationStatus == 0
disclose(ok ? "SPEAKER: done — Neuron spoke aloud." : "SPEAKER: afplay failed.")
if let i = wavInfo(wavPath) {
emit(["ok": ok, "op": "speak", "file": wavPath, "played_aloud": ok,
"sample_rate": i.sampleRate, "channels": i.channels,
"seconds": Double(i.frames)/Double(max(i.sampleRate,1))])
} else {
emit(["ok": ok, "op": "speak", "file": wavPath, "played_aloud": ok])
}
}
// ----------------------------------------------------------------------------
// MIC (afferent) capture N seconds -> 16k mono 16-bit WAV (formant-ready).
// ----------------------------------------------------------------------------
func listen(seconds: Double, out: String) {
Consent.require("mic")
requireOSAccess(.audio, "microphone")
disclose("MIC: capturing \(seconds)s -> '\(out)' (16 kHz mono, LOCAL, never egresses).")
let url = URL(fileURLWithPath: out)
let settings: [String: Any] = [
AVFormatIDKey: kAudioFormatLinearPCM,
AVSampleRateKey: 16000.0,
AVNumberOfChannelsKey: 1,
AVLinearPCMBitDepthKey: 16,
AVLinearPCMIsFloatKey: false,
AVLinearPCMIsBigEndianKey: false,
]
guard let rec = try? AVAudioRecorder(url: url, settings: settings) else {
die("could not open the microphone recorder.")
}
rec.record()
Thread.sleep(forTimeInterval: seconds)
rec.stop()
// let the file flush
Thread.sleep(forTimeInterval: 0.1)
if let i = wavInfo(out) {
disclose("MIC: captured \(i.frames) frames @ \(i.sampleRate)Hz — ready to hand to the ingest organ.")
emit(["ok": true, "op": "listen", "file": out, "sample_rate": i.sampleRate,
"channels": i.channels, "frames": i.frames,
"seconds": Double(i.frames)/Double(max(i.sampleRate,1)),
"next": "ingest -> phonetic/voice geometry"])
} else {
die("mic capture produced no readable WAV.")
}
}
// ----------------------------------------------------------------------------
// CAMERA (afferent) capture ONE frame -> JPEG on disk.
// ----------------------------------------------------------------------------
// Grab one video frame via AVCaptureVideoDataOutput (CLI-safe; no KVO/photo classes).
final class FrameGrabber: NSObject, AVCaptureVideoDataOutputSampleBufferDelegate {
let sem = DispatchSemaphore(value: 0)
var cgImage: CGImage?
var seen = 0
let cictx = CIContext(options: nil)
func captureOutput(_ output: AVCaptureOutput, didOutput sampleBuffer: CMSampleBuffer,
from connection: AVCaptureConnection) {
seen += 1
if cgImage != nil || seen < 5 { return } // let exposure settle a few frames
guard let pb = CMSampleBufferGetImageBuffer(sampleBuffer) else { return }
let ci = CIImage(cvPixelBuffer: pb)
cgImage = cictx.createCGImage(ci, from: ci.extent)
sem.signal()
}
}
func see(out: String) {
Consent.require("camera")
requireOSAccess(.video, "camera")
disclose("CAMERA: capturing one frame -> '\(out)' (LOCAL, never egresses).")
let session = AVCaptureSession()
session.sessionPreset = .photo
guard let device = AVCaptureDevice.default(for: .video),
let input = try? AVCaptureDeviceInput(device: device),
session.canAddInput(input) else { die("no camera device available.") }
session.addInput(input)
let output = AVCaptureVideoDataOutput()
output.alwaysDiscardsLateVideoFrames = true
let grabber = FrameGrabber()
output.setSampleBufferDelegate(grabber, queue: DispatchQueue(label: "periph.cam"))
guard session.canAddOutput(output) else { die("cannot add video output.") }
session.addOutput(output)
session.startRunning()
if grabber.sem.wait(timeout: .now() + 10) == .timedOut { session.stopRunning(); die("camera capture timed out.") }
session.stopRunning()
guard let cg = grabber.cgImage,
let dst = CGImageDestinationCreateWithURL(URL(fileURLWithPath: out) as CFURL,
"public.jpeg" as CFString, 1, nil)
else { die("camera returned no frame.") }
CGImageDestinationAddImage(dst, cg, nil)
guard CGImageDestinationFinalize(dst) else { die("could not write JPEG.") }
let bytes = ((try? FileManager.default.attributesOfItem(atPath: out))?[.size] as? Int) ?? 0
disclose("CAMERA: wrote \(cg.width)x\(cg.height) frame (\(bytes) bytes) — ready for scene-geometry ingest.")
emit(["ok": true, "op": "see", "file": out, "width": cg.width, "height": cg.height,
"bytes": bytes, "next": "ingest -> scene-geometry"])
}
// ============================================================================
// FEAT the afferent METABOLISM: a raw capture becomes a COMPACT descriptor
// (a few dozen numbers), the mirror of the efferent signature. This is what
// gets handed to the ingest organ as geometry NOT the raw stream. Own-core.
// ============================================================================
// Read all 16-bit PCM samples from a WAV (own-core).
func readWavSamples(_ path: String) -> (samples: [Double], sr: Int, ch: Int)? {
guard let d = FileManager.default.contents(atPath: path), d.count > 44 else { return nil }
func rd16(_ o: Int) -> Int { Int(d[o]) | (Int(d[o+1]) << 8) }
func rd32(_ o: Int) -> Int { Int(d[o]) | (Int(d[o+1])<<8) | (Int(d[o+2])<<16) | (Int(d[o+3])<<24) }
var ch = 0, sr = 0, bits = 0
var o = 12
while o + 8 <= d.count {
let id = String(bytes: d[o..<o+4], encoding: .ascii) ?? ""
let sz = rd32(o+4)
if id == "fmt " && o + 24 <= d.count { ch = rd16(o+10); sr = rd32(o+12); bits = rd16(o+22) }
if id == "data" {
guard bits == 16, ch > 0 else { return nil }
var samples = [Double](); let start = o + 8
let end = min(start + sz, d.count - 1)
var i = start
while i + 1 < end {
var v = Int(rd16(i)); if v >= 32768 { v -= 65536 }
samples.append(Double(v) / 32768.0)
i += 2 * ch // take channel 0 if stereo
}
return (samples, sr, ch)
}
o += 8 + sz + (sz & 1)
}
return nil
}
// Audio descriptor = compact sound/voice signature (energy, ZCR, centroid, F0).
// The seed for phonetic geometry + the hear->imitate voice-signature.
func computeAudio(_ path: String) -> (content: String, vector: [Double], extra: [String: Any]) {
guard let (s, sr, ch) = readWavSamples(path), !s.isEmpty else { die("cannot read PCM from \(path)") }
let n = s.count
let seconds = Double(n) / Double(sr)
var sumsq = 0.0, peak = 0.0, zc = 0.0
for i in 0..<n {
sumsq += s[i]*s[i]; peak = max(peak, abs(s[i]))
if i > 0 && (s[i-1] < 0) != (s[i] < 0) { zc += 1 }
}
let rms = (sumsq / Double(n)).squareRoot()
let zcr = zc / Double(n) * Double(sr) // ~2*dominant freq for tonal
// Spectral centroid via a coarse DFT on a mid window (own-core).
let W = min(2048, n); let off = max(0, (n - W)/2)
var num = 0.0, den = 0.0
let bins = 64
for k in 1..<bins {
let f = Double(k) * Double(sr) / Double(2*bins)
var re = 0.0, im = 0.0
for j in 0..<W {
let ang = -2*Double.pi*Double(k)*Double(j)/Double(2*bins)
re += s[off+j]*cos(ang); im += s[off+j]*sin(ang)
}
let mag = (re*re+im*im).squareRoot()
num += f*mag; den += mag
}
let centroid = den > 0 ? num/den : 0
// F0 via autocorrelation (voice pitch) over plausible speech range 70-400 Hz.
var bestLag = 0; var bestCorr = 0.0
let lagMin = sr/400, lagMax = min(sr/70, n-1)
if lagMax > lagMin {
for lag in lagMin...lagMax {
var c = 0.0
var i = 0; while i + lag < min(n, off+W) { c += s[off+i]*s[off+i+lag]; i += 1 }
if c > bestCorr { bestCorr = c; bestLag = lag }
}
}
let f0 = bestLag > 0 ? Double(sr)/Double(bestLag) : 0
let vector: [Double] = [seconds, Double(sr), Double(ch), rms, peak, zcr, centroid, f0]
let content = String(format:
"Heard sound (afferent, mic): %.2fs at %dHz. RMS energy %.3f, peak %.3f, " +
"zero-crossing rate %.0fHz, spectral centroid %.0fHz, estimated voice pitch F0 %.0fHz. " +
"Compact voice/sound signature (%d numbers) — phonetic geometry + hear-to-imitate seed.",
seconds, sr, rms, peak, zcr, centroid, f0, vector.count)
disclose("FEAT(audio): \(vector.count)-number signature vs \(n) raw samples (~\(n/max(vector.count,1))x compression).")
return (content, vector, ["f0_hz": f0, "centroid_hz": centroid, "zcr_hz": zcr,
"rms": rms, "seconds": seconds, "raw_samples": n])
}
func featAudio(_ path: String) {
let r = computeAudio(path)
var out: [String: Any] = ["ok": true, "op": "feat-audio", "file": path,
"vector": r.vector, "content": r.content,
"ingest": ["node_type": "Observation", "tier": "Episodic", "content": r.content]]
r.extra.forEach { out[$0] = $1 }
emit(out)
}
// Image descriptor = compact scene-geometry (dims, brightness, region grid).
func computeImage(_ path: String) -> (content: String, vector: [Double], extra: [String: Any]) {
guard let src = CGImageSourceCreateWithURL(URL(fileURLWithPath: path) as CFURL, nil),
let img = CGImageSourceCreateImageAtIndex(src, 0, nil) else { die("cannot decode image \(path)") }
let w = img.width, h = img.height
let cs = CGColorSpaceCreateDeviceRGB()
let bpr = w * 4
var buf = [UInt8](repeating: 0, count: h * bpr)
guard let ctx = CGContext(data: &buf, width: w, height: h, bitsPerComponent: 8,
bytesPerRow: bpr, space: cs,
bitmapInfo: CGImageAlphaInfo.premultipliedLast.rawValue) else {
die("cannot rasterize image")
}
ctx.draw(img, in: CGRect(x: 0, y: 0, width: w, height: h))
// 3x3 region average luminance + overall average color.
var rAvg = 0.0, gAvg = 0.0, bAvg = 0.0
var grid = [Double](repeating: 0, count: 9); var gridN = [Int](repeating: 0, count: 9)
let step = max(1, (w*h)/40000) // subsample for speed
var count = 0; var idx = 0
while idx < w*h {
let x = idx % w, y = idx / w
let p = y*bpr + x*4
let r = Double(buf[p]), g = Double(buf[p+1]), b = Double(buf[p+2])
rAvg += r; gAvg += g; bAvg += b; count += 1
let cell = (min(2, y*3/h))*3 + min(2, x*3/w)
grid[cell] += 0.299*r + 0.587*g + 0.114*b; gridN[cell] += 1
idx += step
}
if count == 0 { die("no pixels sampled") }
rAvg /= Double(count); gAvg /= Double(count); bAvg /= Double(count)
for i in 0..<9 { grid[i] = gridN[i] > 0 ? grid[i]/Double(gridN[i]) : 0 }
let bright = (0.299*rAvg + 0.587*gAvg + 0.114*bAvg)/255.0
let vector = [Double(w), Double(h), rAvg/255, gAvg/255, bAvg/255, bright] + grid.map { $0/255 }
let content = String(format:
"Saw scene (afferent, camera): %dx%d frame. Mean color rgb(%.0f,%.0f,%.0f), " +
"brightness %.2f. 3x3 luminance grid [%.0f %.0f %.0f / %.0f %.0f %.0f / %.0f %.0f %.0f]. " +
"Compact scene-geometry (%d numbers) vs %d pixel-channels.",
w, h, rAvg, gAvg, bAvg, bright,
grid[0],grid[1],grid[2],grid[3],grid[4],grid[5],grid[6],grid[7],grid[8],
vector.count, w*h*3)
disclose("FEAT(image): \(vector.count)-number scene-geometry vs \(w*h*3) pixel-channels (~\(w*h*3/max(vector.count,1))x).")
return (content, vector, ["width": w, "height": h, "brightness": bright])
}
func featImage(_ path: String) {
let r = computeImage(path)
var out: [String: Any] = ["ok": true, "op": "feat-image", "file": path,
"vector": r.vector, "content": r.content,
"ingest": ["node_type": "Observation", "tier": "Episodic", "content": r.content]]
r.extra.forEach { out[$0] = $1 }
emit(out)
}
// The afferent WIRE hand a capture's descriptor to the ingest organ (engram),
// where it becomes an embedded node = GEOMETRY. Own-core URLSession POST.
// LOCAL only: point at a local engram; raw stream never leaves the machine.
func postNode(engramURL: String, content: String, label: String, tags: [String]) -> String? {
guard let url = URL(string: engramURL + "/api/nodes") else { return nil }
let body: [String: Any] = ["content": content, "node_type": "Observation",
"label": label, "tier": "Episodic",
"salience": 0.7, "importance": 0.6, "confidence": 0.9,
"tags": tags]
var req = URLRequest(url: url); req.httpMethod = "POST"
req.setValue("application/json", forHTTPHeaderField: "Content-Type")
req.httpBody = try? JSONSerialization.data(withJSONObject: body)
let sem = DispatchSemaphore(value: 0); var out: String?
URLSession.shared.dataTask(with: req) { data, _, _ in
if let d = data { out = String(data: d, encoding: .utf8) }
sem.signal()
}.resume()
_ = sem.wait(timeout: .now() + 15)
return out
}
func ingest(_ path: String, kind: String, engramURL: String) {
let r = kind == "audio" ? computeAudio(path) : computeImage(path)
let label = kind == "audio" ? "heard:mic" : "saw:camera"
disclose("INGEST: handing \(kind) descriptor to the ingest organ at \(engramURL) (LOCAL) -> geometry.")
guard let resp = postNode(engramURL: engramURL, content: r.content, label: label,
tags: ["peripheral", kind == "audio" ? "afferent-mic" : "afferent-camera"]) else {
die("ingest POST failed (no local engram at \(engramURL)?)")
}
// pull the node id out of the response (own-core, tolerant)
var nodeId = ""
if let d = resp.data(using: .utf8),
let o = try? JSONSerialization.jsonObject(with: d) as? [String: Any] {
nodeId = (o["id"] as? String) ?? (o["node_id"] as? String) ?? ""
}
disclose("INGEST: landed as node \(nodeId.isEmpty ? "(see response)" : nodeId) — the capture is now geometry in the engram.")
emit(["ok": !nodeId.isEmpty, "op": "ingest-\(kind)", "file": path,
"node_id": nodeId, "engram_response": resp, "content": r.content,
"vector": r.vector])
}
// ============================================================================
// VOICE BY IMITATION hear a voice, grab its compact SIGNATURE (pitch +
// formants F1-F5 via LPC), and speak back in that voice by source-filter
// resynthesis. Own-core DSP (physics), no training, no stolen voice. The
// afferent twin of the music instrument-signature: a voice = a few dozen
// numbers, not a corpus.
// ============================================================================
func hamming(_ x: [Double]) -> [Double] {
let n = x.count; if n < 2 { return x }
return (0..<n).map { x[$0] * (0.54 - 0.46*cos(2*Double.pi*Double($0)/Double(n-1))) }
}
func autocorr(_ x: [Double], _ p: Int) -> [Double] {
var r = [Double](repeating: 0, count: p+1)
for lag in 0...p { var s = 0.0; var i = lag; while i < x.count { s += x[i]*x[i-lag]; i += 1 }; r[lag] = s }
return r
}
// Levinson-Durbin -> LPC coeffs a[0..p] (A(z)=1+sum a[k]z^-k) and residual energy.
func levinson(_ r: [Double], _ p: Int) -> (a: [Double], err: Double) {
var a = [Double](repeating: 0, count: p+1); a[0] = 1
var err = r[0]
if err <= 0 { return (a, 0) }
for i in 1...p {
var acc = r[i]
if i > 1 { for j in 1..<i { acc += a[j]*r[i-j] } }
let k = -acc/err
var na = a; na[i] = k
if i > 1 { for j in 1..<i { na[j] = a[j] + k*a[i-j] } }
a = na; err *= (1 - k*k)
if err <= 0 { break }
}
return (a, err)
}
// Formant peaks from the LPC all-pole spectral envelope.
func formants(_ a: [Double], sr: Int) -> [(f: Double, bw: Double)] {
let p = a.count - 1
let steps = 512
var mag = [Double](repeating: 0, count: steps)
for s in 0..<steps {
let w = Double.pi * Double(s) / Double(steps) // 0..pi -> 0..sr/2
var re = 0.0, im = 0.0
for k in 0...p { re += a[k]*cos(w*Double(k)); im -= a[k]*sin(w*Double(k)) }
mag[s] = 1.0 / max((re*re+im*im).squareRoot(), 1e-9)
}
var peaks: [(f: Double, bw: Double)] = []
for s in 1..<(steps-1) where mag[s] > mag[s-1] && mag[s] >= mag[s+1] {
let f = Double(s) * Double(sr) / 2 / Double(steps)
if f > 150 && f < 5200 {
// crude bandwidth: width where magnitude falls to peak/sqrt(2)
let thr = mag[s]/1.4142
var lo = s; while lo > 0 && mag[lo] > thr { lo -= 1 }
var hi = s; while hi < steps-1 && mag[hi] > thr { hi += 1 }
let bw = Double(hi-lo) * Double(sr) / 2 / Double(steps)
peaks.append((f, bw))
}
}
return Array(peaks.prefix(5))
}
func pitchOf(_ frame: [Double], sr: Int) -> Double {
let n = frame.count
let lagMin = sr/400, lagMax = min(sr/70, n-1)
if lagMax <= lagMin { return 0 }
var r0 = 0.0; for v in frame { r0 += v*v }
if r0 < 1e-5 { return 0 }
var bestLag = 0; var best = 0.0
for lag in lagMin...lagMax { var c = 0.0; var i = lag; while i < n { c += frame[i]*frame[i-lag]; i += 1 }; if c > best { best = c; bestLag = lag } }
return (best / r0 > 0.30 && bestLag > 0) ? Double(sr)/Double(bestLag) : 0 // voiced?
}
let LPC_ORDER = 16
let FRAME = 400 // 25ms @16k
let HOP = 160 // 10ms
// Extract Will's voice-signature: averaged F0 + formants over voiced frames.
func voiceprint(_ path: String) -> (f0: Double, f0lo: Double, f0hi: Double, formants: [(Double,Double)], content: String) {
guard let (x, sr, _) = readWavSamples(path), x.count > FRAME else { die("cannot read speech from \(path)") }
var f0s: [Double] = []
var fbank: [[Double]] = [[],[],[],[],[]]
var bbank: [[Double]] = [[],[],[],[],[]]
var pos = 0
while pos + FRAME <= x.count {
let raw = Array(x[pos..<pos+FRAME])
let f0 = pitchOf(raw, sr: sr)
if f0 > 0 { // voiced frame only
f0s.append(f0)
let r = autocorr(hamming(raw), LPC_ORDER)
if r[0] > 1e-6 {
let (a, _) = levinson(r, LPC_ORDER)
let fs = formants(a, sr: sr)
for (i, fm) in fs.enumerated() where i < 5 { fbank[i].append(fm.f); bbank[i].append(fm.bw) }
}
}
pos += HOP
}
func med(_ v: [Double]) -> Double { v.isEmpty ? 0 : v.sorted()[v.count/2] }
let f0med = med(f0s)
let f0lo = f0s.isEmpty ? 0 : f0s.sorted().first!
let f0hi = f0s.isEmpty ? 0 : f0s.sorted().last!
var forms: [(Double,Double)] = []
for i in 0..<5 where !fbank[i].isEmpty { forms.append((med(fbank[i]), med(bbank[i]))) }
let fstr = forms.map { String(format:"%.0f", $0.0) }.joined(separator: "/")
let content = String(format:
"Voice-signature (afferent, heard a voice): pitch F0 %.0fHz (range %.0f-%.0fHz), " +
"formants F1-F5 = %@ Hz. Compact voiceprint (%d numbers) — grabbed by ear for imitation, not trained.",
f0med, f0lo, f0hi, fstr, 1 + forms.count*2)
return (f0med, f0lo, f0hi, forms, content)
}
// IMITATE: LPC analysis-resynthesis. Reconstruct the heard voice from its
// per-frame filter model + pitch the voice rebuilt from its signature.
func imitate(inPath: String, outPath: String) {
guard let (x, sr, _) = readWavSamples(inPath), x.count > FRAME else { die("cannot read speech from \(inPath)") }
var out = [Double](repeating: 0, count: x.count)
var state = [Double](repeating: 0, count: LPC_ORDER) // past outputs
var phase = 0.0
var lastF0 = 0.0
var pos = 0
while pos + FRAME <= x.count {
let raw = Array(x[pos..<pos+FRAME])
let r = autocorr(hamming(raw), LPC_ORDER)
let f0 = pitchOf(raw, sr: sr)
if r[0] < 1e-7 { pos += HOP; continue }
let (a, err) = levinson(r, LPC_ORDER)
let gain = max(err, 0).squareRoot()
let useF0 = f0 > 0 ? f0 : (lastF0 > 0 ? lastF0 : 0)
lastF0 = f0
for i in 0..<HOP {
let idx = pos + i; if idx >= x.count { break }
var e = 0.0
if useF0 > 0 { // voiced: glottal impulse train
phase += useF0/Double(sr)
if phase >= 1.0 { phase -= 1.0; e = sqrt(Double(sr)/useF0) } // energy-normalized impulse
} else { // unvoiced: noise
e = Double.random(in: -1...1)
}
var y = gain * e
for k in 1...LPC_ORDER { y -= a[k]*state[k-1] }
for k in stride(from: LPC_ORDER-1, through: 1, by: -1) { state[k] = state[k-1] }
state[0] = y
out[idx] = y
}
pos += HOP
}
// normalize to peak 0.9
let peak = out.map { abs($0) }.max() ?? 1
let scale = peak > 1e-9 ? 0.9/peak : 1
let samples = out.map { Int16(max(-32767, min(32767, $0*scale*32767))) }
writeWav(URL(fileURLWithPath: outPath), samples: samples, sampleRate: sr)
let vp = voiceprint(inPath)
disclose(String(format: "IMITATE: rebuilt the voice from its signature (F0 %.0fHz, formants %@) -> %@",
vp.f0, vp.formants.map{String(format:"%.0f",$0.0)}.joined(separator:"/"), outPath))
emit(["ok": true, "op": "imitate", "in": inPath, "out": outPath,
"f0_hz": vp.f0, "f0_range": [vp.f0lo, vp.f0hi],
"formants_hz": vp.formants.map { $0.0 },
"method": "LPC analysis-resynthesis (own-core, no training, no stolen voice)"])
}
// ============================================================================
// CONVERSE (full-duplex) the interruptible conversational loop.
// The utterance is a persistent, ordered meaning-plan of SEGMENTS, each with
// a salience. The speaker plays them; the mic listens concurrently. On user
// speech: pause INSTANTLY, classify (backchannel vs barge-in), then DECIDE
// yield-or-hold from the salience of the current segment + the social read.
// Yielded utterances persist their remaining plan so Neuron can RESUME.
// ============================================================================
struct Segment { let file: String; let salience: Double; let text: String }
enum Decision { case backchannelContinue, hold, yield }
// The yield-or-hold DECISION grounded, contextual. Not a fixed rule.
func decide(currentSalience: Double, progress: Double,
interrupterAuthority: Double, isBackchannel: Bool) -> Decision {
if isBackchannel { return .backchannelContinue } // "mm-hm" => keep going
// Holding the floor is justified when what I'm saying matters AND I'm nearly
// done (cheap to finish) AND the interrupter isn't high-priority.
let holdScore = currentSalience * 0.6 + progress * 0.4
if holdScore >= 0.6 && interrupterAuthority < 0.8 { return .hold }
return .yield // default: be polite, let them in
}
final class Conversation {
let engine = AVAudioEngine()
let player = AVAudioPlayerNode()
var micLive = false
// VAD state (shared with the audio tap thread)
let lock = NSLock()
var micRMS: Float = 0
var speechFrames = 0 // consecutive above-threshold frames
var onsetHandled = false
let resumePath: String
init(resumePath: String) { self.resumePath = resumePath }
// Try to bring the mic up as a live VAD. Returns false if unavailable/denied.
func startMic() -> Bool {
let status = AVCaptureDevice.authorizationStatus(for: .audio)
if Consent.load()["mic"] != true || status != .authorized {
disclose("CONVERSE: live mic not available (consent/OS) — using injected barge events for the proof.")
return false
}
let input = engine.inputNode
// Acoustic echo cancellation: the OS voice-processing unit subtracts our
// own speaker output from the mic so Neuron does NOT hear itself and
// barge in on its own voice. This is what makes real-room barge-in work.
do { try input.setVoiceProcessingEnabled(true); disclose("CONVERSE: AEC on (echo-cancelled mic — won't self-interrupt).") }
catch { disclose("CONVERSE: AEC unavailable (\(error)); raising VAD floor instead.") }
let fmt = input.inputFormat(forBus: 0)
if fmt.sampleRate == 0 { return false }
input.installTap(onBus: 0, bufferSize: 1024, format: fmt) { [weak self] buf, _ in
guard let self = self, let ch = buf.floatChannelData?[0] else { return }
let n = Int(buf.frameLength)
var sum: Float = 0
for i in 0..<n { let v = ch[i]; sum += v*v }
let rms = n > 0 ? (sum / Float(n)).squareRoot() : 0
self.lock.lock(); self.micRMS = rms; self.lock.unlock()
}
micLive = true
disclose("CONVERSE: full-duplex — mic listening WHILE speaking (barge-in armed).")
return true
}
func run(_ segs: [Segment], interrupterAuthority: Double,
injectBargeAt: Double?, injectKind: String, startIndex: Int, liveMic: Bool) {
engine.attach(player)
let firstFmt = (try? AVAudioFile(forReading: URL(fileURLWithPath: segs[startIndex].file)))?.processingFormat
?? AVAudioFormat(standardFormatWithSampleRate: 16000, channels: 1)!
engine.connect(player, to: engine.mainMixerNode, format: firstFmt)
if liveMic { _ = startMic() }
else { disclose("CONVERSE: deterministic mode (live mic off) — barge events \(injectBargeAt != nil ? "injected" : "none").") }
do { try engine.start() } catch { die("audio engine failed to start: \(error)") }
player.play()
let injectDeadline = injectBargeAt.map { Date().addingTimeInterval($0) }
var injectedFired = false
var idx = startIndex
segmentLoop: while idx < segs.count {
let seg = segs[idx]
guard let f = try? AVAudioFile(forReading: URL(fileURLWithPath: seg.file)) else {
disclose("CONVERSE: missing segment '\(seg.file)', skipping."); idx += 1; continue
}
let dur = Double(f.length) / f.processingFormat.sampleRate
disclose(String(format: "CONVERSE: speaking segment %d/%d (salience %.2f) — \"%@\"",
idx+1, segs.count, seg.salience, seg.text))
emit(["op": "converse", "event": "speaking", "segment": idx,
"salience": seg.salience, "text": seg.text])
let done = DispatchSemaphore(value: 0)
// .dataPlayedBack: completion fires only after the audio has actually
// played OUT the DAC (not merely been consumed) so the tail is never
// clipped and playback always runs the FULL file length.
player.scheduleFile(f, at: nil, completionCallbackType: .dataPlayedBack) { _ in done.signal() }
player.play()
// Monitor this segment: poll VAD / injected event until it finishes.
let segStart = Date()
while done.wait(timeout: .now() + 0.02) == .timedOut {
let elapsed = Date().timeIntervalSince(segStart)
let progress = min(elapsed / max(dur, 0.001), 1.0)
// --- detect an onset (live mic OR injected) ---
var onset = false
if micLive {
lock.lock(); let rms = micRMS; lock.unlock()
if rms > 0.02 { speechFrames += 1 } else { speechFrames = 0 }
if speechFrames >= 3 && !onsetHandled { onset = true } // ~60ms of voice
}
if let dl = injectDeadline, !injectedFired, Date() >= dl, !onsetHandled { onset = true; injectedFired = true }
if onset {
onsetHandled = true
// (1) BARGE-IN: pause INSTANTLY, on the spot.
player.pause()
let tBarge = Date().timeIntervalSince(segStart)
disclose(String(format: "CONVERSE: << user speech at %.2fs into segment %d — PAUSED instantly >>", tBarge, idx+1))
emit(["op": "converse", "event": "barge_in", "segment": idx,
"at_seconds": tBarge, "progress": progress])
// (2) classify backchannel vs real barge-in
let isBackchannel = classifyBackchannel(injected: injectDeadline != nil,
kind: injectKind)
let d = decide(currentSalience: seg.salience, progress: progress,
interrupterAuthority: interrupterAuthority,
isBackchannel: isBackchannel)
switch d {
case .backchannelContinue:
disclose("CONVERSE: read as BACKCHANNEL (\"mm-hm\") — keep going, resume seamlessly.")
emit(["op": "converse", "event": "backchannel_continue", "segment": idx])
onsetHandled = false; speechFrames = 0
player.play() // seamless resume
case .hold:
disclose("CONVERSE: HOLD the floor — \"hang on, let me finish this thought.\" (high salience, nearly done)")
emit(["op": "converse", "event": "hold_floor", "segment": idx,
"salience": seg.salience, "progress": progress])
onsetHandled = false; speechFrames = 0
player.play() // finish the segment, THEN yield
// after this segment completes we yield the remainder
_ = done.wait(timeout: .now() + dur + 1.0)
persistResume(segs: segs, from: idx + 1, reason: "held-then-yield")
finish(); return
case .yield:
disclose("CONVERSE: YIELD — stop, let them in. Remembering where I was (resumable).")
player.stop()
persistResume(segs: segs, from: idx, reason: "yield")
emit(["op": "converse", "event": "yield", "interrupted_segment": idx,
"resume_from": idx])
finish(); return
}
}
}
emit(["op": "converse", "event": "segment_done", "segment": idx])
idx += 1
}
// whole utterance completed uninterrupted
clearResume()
disclose("CONVERSE: utterance complete (uninterrupted).")
emit(["ok": true, "op": "converse", "event": "complete", "segments": segs.count])
finish()
}
// A backchannel is brief/low. Injected kind lets us prove both paths headlessly;
// the live path would measure post-onset duration & energy.
func classifyBackchannel(injected: Bool, kind: String) -> Bool {
if injected { return kind == "backchannel" }
// live: sample ~250ms after onset; if speech already died away, it was a backchannel
Thread.sleep(forTimeInterval: 0.25)
lock.lock(); let rms = micRMS; lock.unlock()
return rms < 0.015
}
func persistResume(segs: [Segment], from: Int, reason: String) {
let remaining = segs[from...].map { ["file": $0.file, "salience": $0.salience, "text": $0.text] as [String: Any] }
let state: [String: Any] = ["resume_from": from, "reason": reason,
"remaining": remaining, "ts": Date().timeIntervalSince1970]
if let d = try? JSONSerialization.data(withJSONObject: state, options: [.prettyPrinted]) {
try? d.write(to: URL(fileURLWithPath: resumePath))
}
disclose("CONVERSE: meaning-plan persisted (\(remaining.count) segments remain) — Neuron can resume the thread.")
}
func clearResume() { try? FileManager.default.removeItem(atPath: resumePath) }
func finish() { player.stop(); if micLive { engine.inputNode.removeTap(onBus: 0) }; engine.stop() }
}
// ----------------------------------------------------------------------------
// CLI
// ----------------------------------------------------------------------------
func loadManifest(_ path: String) -> (segs: [Segment], utterance: String) {
guard let d = FileManager.default.contents(atPath: path),
let o = try? JSONSerialization.jsonObject(with: d) as? [String: Any],
let arr = o["segments"] as? [[String: Any]] else { die("bad manifest: \(path)") }
let segs = arr.map { Segment(file: $0["file"] as? String ?? "",
salience: ($0["salience"] as? NSNumber)?.doubleValue ?? 0.5,
text: $0["text"] as? String ?? "") }
return (segs, o["utterance"] as? String ?? "")
}
let args = CommandLine.arguments
guard args.count >= 2 else {
print("""
periph — Neuron peripheral I/O (own-core, local, consent-gated)
grant <camera|mic> grant a sensitive sense (Neuron-level consent)
revoke <camera|mic> revoke it
status show consent state
speak <file.wav> SPEAK ALOUD (efferent) via the speaker
tone <out.wav> [hz] [sec] own-core synth a test WAV (no deps)
listen <sec> <out.wav> MIC capture (afferent) 16k mono
see <out.jpg> CAMERA one frame (afferent)
feat-audio <file.wav> extract compact voice/sound signature (for ingest)
feat-image <file.jpg> extract compact scene-geometry (for ingest)
ingest-audio <file.wav> <engramURL> capture -> descriptor -> engram node (geometry)
ingest-image <file.jpg> <engramURL> capture -> descriptor -> engram node (geometry)
voiceprint <voice.wav> extract voice-signature (F0 + formants F1-F5)
imitate <voice.wav> <out.wav> speak back in that voice (LPC analysis-resynthesis)
hear-imitate <sec> <out.wav> MIC -> extract signature -> imitate -> SPEAK ALOUD
wav-info <file.wav> print WAV geometry
converse <manifest.json> [--authority F] [--barge-at S[:backchannel|:bargein]] [--resume]
full-duplex interruptible utterance
""")
exit(0)
}
switch args[1] {
case "grant":
guard args.count >= 3 else { die("grant needs a device") }
var g = Consent.load(); g[args[2]] = true; Consent.save(g)
disclose("granted '\(args[2])' — the user consents; raw stream stays local, never egresses.")
emit(["ok": true, "op": "grant", "device": args[2], "consent": g])
case "revoke":
guard args.count >= 3 else { die("revoke needs a device") }
var g = Consent.load(); g[args[2]] = false; Consent.save(g)
emit(["ok": true, "op": "revoke", "device": args[2], "consent": g])
case "status":
emit(["ok": true, "op": "status", "consent": Consent.load()])
case "speak":
guard args.count >= 3 else { die("speak needs a wav") }
speak(args[2])
case "tone":
guard args.count >= 3 else { die("tone needs an out path") }
let hz = args.count >= 4 ? Double(args[3]) ?? 220 : 220
let sec = args.count >= 5 ? Double(args[4]) ?? 1.0 : 1.0
let sr = 16000
var s = [Int16](); s.reserveCapacity(Int(Double(sr)*sec))
for i in 0..<Int(Double(sr)*sec) {
let t = Double(i)/Double(sr)
let env = min(1.0, min(t*20, (sec - t)*20)) // gentle attack/release
s.append(Int16(env * 0.3 * 32767 * sin(2*Double.pi*hz*t)))
}
writeWav(URL(fileURLWithPath: args[2]), samples: s, sampleRate: sr)
disclose("tone: wrote own-core \(sec)s @ \(hz)Hz WAV to \(args[2]).")
emit(["ok": true, "op": "tone", "file": args[2], "hz": hz, "seconds": sec])
case "listen":
guard args.count >= 4 else { die("listen needs <sec> <out.wav>") }
listen(seconds: Double(args[2]) ?? 3.0, out: args[3])
case "see":
guard args.count >= 3 else { die("see needs an out path") }
see(out: args[2])
case "feat-audio":
guard args.count >= 3 else { die("feat-audio needs a wav") }
featAudio(args[2])
case "feat-image":
guard args.count >= 3 else { die("feat-image needs an image") }
featImage(args[2])
case "ingest-audio":
guard args.count >= 4 else { die("ingest-audio needs <wav> <engramURL>") }
ingest(args[2], kind: "audio", engramURL: args[3])
case "ingest-image":
guard args.count >= 4 else { die("ingest-image needs <image> <engramURL>") }
ingest(args[2], kind: "image", engramURL: args[3])
case "voiceprint":
guard args.count >= 3 else { die("voiceprint needs a wav") }
let vp = voiceprint(args[2])
disclose("VOICEPRINT: \(vp.content)")
emit(["ok": true, "op": "voiceprint", "file": args[2], "f0_hz": vp.f0,
"f0_range": [vp.f0lo, vp.f0hi], "formants_hz": vp.formants.map { $0.0 },
"bandwidths_hz": vp.formants.map { $0.1 }, "content": vp.content,
"ingest": ["node_type": "Observation", "tier": "Episodic", "content": vp.content]])
case "imitate":
guard args.count >= 4 else { die("imitate needs <voice.wav> <out.wav>") }
imitate(inPath: args[2], outPath: args[3])
case "hear-imitate":
guard args.count >= 4 else { die("hear-imitate needs <sec> <out.wav>") }
let secs = Double(args[2]) ?? 4.0
let outp = args[3]
let capp = outp.replacingOccurrences(of: ".wav", with: "") + ".heard.wav"
disclose("HEAR-IMITATE: open the ear, listen \(secs)s, grab the voice, speak it back.")
listen(seconds: secs, out: capp) // afferent: hear the voice
imitate(inPath: capp, outPath: outp) // extract signature + resynthesize
speak(outp) // efferent: speak back ALOUD in that voice
case "wav-info":
guard args.count >= 3, let i = wavInfo(args[2]) else { die("wav-info needs a readable wav") }
disclose("WAV \(args[2]): \(i.sampleRate)Hz \(i.channels)ch \(i.bits)bit \(i.frames) frames")
emit(["ok": true, "op": "wav-info", "sample_rate": i.sampleRate, "channels": i.channels,
"bits": i.bits, "frames": i.frames,
"seconds": Double(i.frames)/Double(max(i.sampleRate,1))])
case "converse":
guard args.count >= 3 else { die("converse needs a manifest") }
let (segs, utter) = loadManifest(args[2])
var authority = 0.5
var bargeAt: Double? = nil
var bargeKind = "bargein"
var resume = false
var liveMic = false
var i = 3
while i < args.count {
switch args[i] {
case "--authority": if i+1 < args.count { authority = Double(args[i+1]) ?? 0.5; i += 1 }
case "--barge-at":
if i+1 < args.count {
let parts = args[i+1].split(separator: ":")
bargeAt = Double(parts[0]) ?? nil
if parts.count > 1 { bargeKind = String(parts[1]) }
i += 1
}
case "--resume": resume = true
case "--live-mic": liveMic = true
default: break
}
i += 1
}
let resumePath = (ProcessInfo.processInfo.environment["PERIPH_HOME"]
?? FileManager.default.currentDirectoryPath + "/peripheral") + "/.resume.json"
var startIndex = 0
var runSegs = segs
if resume, let d = FileManager.default.contents(atPath: resumePath),
let o = try? JSONSerialization.jsonObject(with: d) as? [String: Any],
let rem = o["remaining"] as? [[String: Any]] {
runSegs = rem.map { Segment(file: $0["file"] as? String ?? "",
salience: ($0["salience"] as? NSNumber)?.doubleValue ?? 0.5,
text: $0["text"] as? String ?? "") }
startIndex = 0
disclose("CONVERSE: resuming — \"as I was saying...\" (\(runSegs.count) segments left).")
emit(["op": "converse", "event": "resume", "remaining": runSegs.count])
}
if runSegs.isEmpty { die("no segments to speak") }
disclose("CONVERSE: utterance = \"\(utter)\" (\(runSegs.count) segments).")
let convo = Conversation(resumePath: resumePath)
convo.run(runSegs, interrupterAuthority: authority,
injectBargeAt: bargeAt, injectKind: bargeKind, startIndex: startIndex, liveMic: liveMic)
default:
die("unknown command: \(args[1])")
}