Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 2865d6ad26 |
@@ -213,6 +213,11 @@ fn hist_append(hist: String, role: String, content: String) -> String {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn hist_trim(hist: String) -> String {
|
fn hist_trim(hist: String) -> String {
|
||||||
|
// Issue #9 (fragile parser): uses manual str_index_of scan rather than a real
|
||||||
|
// JSON parser. If the history JSON does not contain the expected marker pattern
|
||||||
|
// (e.g. corrupted or truncated), returns the unmodified hist silently — silent
|
||||||
|
// data corruption that causes LLM context-length errors on the next turn.
|
||||||
|
// TODO: replace with json_array_slice() once available in the EL runtime.
|
||||||
let inner: String = str_slice(hist, 1, str_len(hist) - 1)
|
let inner: String = str_slice(hist, 1, str_len(hist) - 1)
|
||||||
let marker: String = "{\"role\":"
|
let marker: String = "{\"role\":"
|
||||||
let i1: Int = str_index_of(inner, marker)
|
let i1: Int = str_index_of(inner, marker)
|
||||||
@@ -271,10 +276,20 @@ fn conv_history_load() -> String {
|
|||||||
fn handle_chat(body: String) -> String {
|
fn handle_chat(body: String) -> String {
|
||||||
let message: String = json_get(body, "message")
|
let message: String = json_get(body, "message")
|
||||||
if str_eq(message, "") {
|
if str_eq(message, "") {
|
||||||
return "{\"error\":\"message is required\",\"response\":\"\"}"
|
// Issue #5: missing required param — HTTP 400.
|
||||||
|
return "{\"__status__\":400,\"error\":\"message is required\",\"response\":\"\"}"
|
||||||
}
|
}
|
||||||
|
|
||||||
// Load history BEFORE compiling context so we can anchor activation to the thread.
|
// Load history BEFORE compiling context so we can anchor activation to the thread.
|
||||||
|
//
|
||||||
|
// TODO(reliability #3 — conv_history global race): "conv_history" is a process-global
|
||||||
|
// state key. Concurrent /api/chat requests that omit session_id all read the same key,
|
||||||
|
// append their exchange, and write it back. Because _state_mu serializes individual
|
||||||
|
// state_get/state_set calls but NOT the read-append-write sequence, one thread's
|
||||||
|
// appended exchange can be overwritten by another thread writing its own version.
|
||||||
|
// The fix is to require callers to supply a session_id (routing them through
|
||||||
|
// session_hist_<id>) and deprecate the global "conv_history" path. Callers using
|
||||||
|
// the session API (which scopes history per session_hist_<id>) are not affected.
|
||||||
let state_hist: String = state_get("conv_history")
|
let state_hist: String = state_get("conv_history")
|
||||||
let stored_hist: String = if str_eq(state_hist, "") { conv_history_load() } else { state_hist }
|
let stored_hist: String = if str_eq(state_hist, "") { conv_history_load() } else { state_hist }
|
||||||
let hist_len: Int = if str_eq(stored_hist, "") { 0 } else { json_array_len(stored_hist) }
|
let hist_len: Int = if str_eq(stored_hist, "") { 0 } else { json_array_len(stored_hist) }
|
||||||
@@ -380,7 +395,8 @@ fn handle_chat(body: String) -> String {
|
|||||||
|| str_starts_with(raw_response, "{\"type\":\"error\"")
|
|| str_starts_with(raw_response, "{\"type\":\"error\"")
|
||||||
|| str_contains(raw_response, "authentication_error")
|
|| str_contains(raw_response, "authentication_error")
|
||||||
if is_error {
|
if is_error {
|
||||||
return "{\"error\":\"llm unavailable\",\"response\":\"\"}"
|
// Issue #6: LLM failure — HTTP 503 (service unavailable).
|
||||||
|
return "{\"__status__\":503,\"error\":\"llm unavailable\",\"response\":\"\"}"
|
||||||
}
|
}
|
||||||
|
|
||||||
let clean_response: String = clean_llm_response(raw_response)
|
let clean_response: String = clean_llm_response(raw_response)
|
||||||
@@ -527,7 +543,15 @@ fn agentic_tools_all() -> String {
|
|||||||
fn call_mcp_bridge(tool_name: String, tool_input: String) -> String {
|
fn call_mcp_bridge(tool_name: String, tool_input: String) -> String {
|
||||||
let eff_input: String = if str_eq(tool_input, "") { "{}" } else { tool_input }
|
let eff_input: String = if str_eq(tool_input, "") { "{}" } else { tool_input }
|
||||||
let body: String = "{\"name\":\"" + tool_name + "\",\"input\":" + eff_input + "}"
|
let body: String = "{\"name\":\"" + tool_name + "\",\"input\":" + eff_input + "}"
|
||||||
let tmp: String = "/tmp/neuron-mcp-call.json"
|
// Issue #12: previously used a fixed path /tmp/neuron-mcp-call.json.
|
||||||
|
// Under concurrent load (64 worker threads), two simultaneous MCP tool calls
|
||||||
|
// race on this file — one call sends the other's input to the bridge.
|
||||||
|
// Fix: monotonic sequence counter makes the path unique per call.
|
||||||
|
let mcp_seq_s: String = state_get("mcp_call_seq")
|
||||||
|
let mcp_seq_n: Int = if str_eq(mcp_seq_s, "") { 0 } else { str_to_int(mcp_seq_s) }
|
||||||
|
let mcp_seq_next: Int = mcp_seq_n + 1
|
||||||
|
state_set("mcp_call_seq", int_to_str(mcp_seq_next))
|
||||||
|
let tmp: String = "/tmp/neuron-mcp-call-" + int_to_str(time_now()) + "-" + int_to_str(mcp_seq_next) + ".json"
|
||||||
fs_write(tmp, body)
|
fs_write(tmp, body)
|
||||||
return exec_capture("curl -s --max-time 30 -X POST http://127.0.0.1:7771/mcp/call -H 'Content-Type: application/json' -d @" + tmp)
|
return exec_capture("curl -s --max-time 30 -X POST http://127.0.0.1:7771/mcp/call -H 'Content-Type: application/json' -d @" + tmp)
|
||||||
}
|
}
|
||||||
@@ -802,15 +826,25 @@ fn is_builtin_tool(tool_name: String) -> Bool {
|
|||||||
|| str_starts_with(tool_name, "neuron_")
|
|| str_starts_with(tool_name, "neuron_")
|
||||||
}
|
}
|
||||||
|
|
||||||
// next_bridge_id — monotonic correlation id for a suspended agentic turn.
|
// next_bridge_id — unique correlation id for a suspended agentic turn.
|
||||||
// Combines boot-relative time with a per-process counter so two unknown-tool
|
// Uses uuid_v4() as the primary uniqueness guarantee so concurrent calls
|
||||||
// suspensions in the same second still get distinct ids.
|
// (even in the same millisecond) cannot collide. The "mcp_bridge_seq"
|
||||||
|
// counter is kept for human readability in logs/debugging but is no longer
|
||||||
|
// relied on for uniqueness.
|
||||||
|
//
|
||||||
|
// TODO(reliability #6): state_get/state_set on "mcp_bridge_seq" is a
|
||||||
|
// non-atomic read-modify-write — two concurrent calls can read the same
|
||||||
|
// counter and produce the same counter suffix. This is now benign because
|
||||||
|
// uuid_v4() provides collision-free uniqueness. A true counter fix would
|
||||||
|
// require an atomic_increment() builtin in el_runtime.c.
|
||||||
fn next_bridge_id() -> String {
|
fn next_bridge_id() -> String {
|
||||||
let prev: String = state_get("mcp_bridge_seq")
|
let prev: String = state_get("mcp_bridge_seq")
|
||||||
let n: Int = if str_eq(prev, "") { 0 } else { str_to_int(prev) }
|
let n: Int = if str_eq(prev, "") { 0 } else { str_to_int(prev) }
|
||||||
let next: Int = n + 1
|
let next: Int = n + 1
|
||||||
state_set("mcp_bridge_seq", int_to_str(next))
|
state_set("mcp_bridge_seq", int_to_str(next))
|
||||||
return "br-" + int_to_str(time_now()) + "-" + int_to_str(next)
|
// uuid_v4() provides collision-free uniqueness; counter is decorative.
|
||||||
|
let uid: String = uuid_v4()
|
||||||
|
return "br-" + uid
|
||||||
}
|
}
|
||||||
|
|
||||||
fn handle_chat_agentic(body: String) -> String {
|
fn handle_chat_agentic(body: String) -> String {
|
||||||
|
|||||||
@@ -7,65 +7,6 @@ import "neuron-api.el"
|
|||||||
import "sessions.el"
|
import "sessions.el"
|
||||||
import "soul.elh"
|
import "soul.elh"
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
|
||||||
// Rate limiting — simple in-memory per-IP sliding window counter.
|
|
||||||
//
|
|
||||||
// State keys:
|
|
||||||
// rl:<ip>:count — request count in the current window
|
|
||||||
// rl:<ip>:window — window start timestamp (unix seconds)
|
|
||||||
//
|
|
||||||
// Limit: configurable via soul state key "soul_rate_limit" (requests per
|
|
||||||
// minute). Falls back to 60 req/min if not set. The /health endpoint is
|
|
||||||
// exempt so monitoring does not consume quota.
|
|
||||||
//
|
|
||||||
// State growth: each unique source IP accumulates exactly 2 state keys
|
|
||||||
// (count + window) for the lifetime of the process. Per-IP storage is
|
|
||||||
// bounded and constant; values reset on window expiry. In aggregate, state
|
|
||||||
// grows linearly with distinct IPs — typical for a trusted-client service.
|
|
||||||
// EL has no state_delete builtin, so keys from inactive IPs persist.
|
|
||||||
// TODO: add state_delete sweep when the EL runtime exposes that primitive.
|
|
||||||
//
|
|
||||||
// Returns "" when the request is allowed, or a 429 JSON body when rejected.
|
|
||||||
// ---------------------------------------------------------------------------
|
|
||||||
fn rate_limit_check(ip: String, path: String) -> String {
|
|
||||||
// Health checks are exempt — they must never be blocked.
|
|
||||||
if str_eq(path, "/health") {
|
|
||||||
return ""
|
|
||||||
}
|
|
||||||
|
|
||||||
let limit_str: String = state_get("soul_rate_limit")
|
|
||||||
let limit: Int = if str_eq(limit_str, "") { 60 } else { str_to_int(limit_str) }
|
|
||||||
|
|
||||||
let now: Int = time_now()
|
|
||||||
let window_key: String = "rl:" + ip + ":window"
|
|
||||||
let count_key: String = "rl:" + ip + ":count"
|
|
||||||
|
|
||||||
let win_str: String = state_get(window_key)
|
|
||||||
let win_start: Int = if str_eq(win_str, "") { now } else { str_to_int(win_str) }
|
|
||||||
|
|
||||||
// New window every 60 seconds.
|
|
||||||
let elapsed: Int = now - win_start
|
|
||||||
let in_window: Bool = elapsed < 60
|
|
||||||
|
|
||||||
let prev_count_str: String = state_get(count_key)
|
|
||||||
let prev_count: Int = if str_eq(prev_count_str, "") { 0 } else { str_to_int(prev_count_str) }
|
|
||||||
|
|
||||||
// Reset window if expired.
|
|
||||||
let eff_count: Int = if in_window { prev_count } else { 0 }
|
|
||||||
let eff_win: Int = if in_window { win_start } else { now }
|
|
||||||
|
|
||||||
let new_count: Int = eff_count + 1
|
|
||||||
state_set(count_key, int_to_str(new_count))
|
|
||||||
state_set(window_key, int_to_str(eff_win))
|
|
||||||
|
|
||||||
if new_count > limit {
|
|
||||||
let retry_after: Int = 60 - (now - eff_win)
|
|
||||||
let eff_retry: Int = if retry_after < 0 { 0 } else { retry_after }
|
|
||||||
return "{\"__status__\":429,\"error\":\"rate limit exceeded\",\"code\":\"rate_limited\",\"retry_after_secs\":" + int_to_str(eff_retry) + "}"
|
|
||||||
}
|
|
||||||
return ""
|
|
||||||
}
|
|
||||||
|
|
||||||
fn strip_query(path: String) -> String {
|
fn strip_query(path: String) -> String {
|
||||||
let q: Int = str_index_of(path, "?")
|
let q: Int = str_index_of(path, "?")
|
||||||
if q < 0 {
|
if q < 0 {
|
||||||
@@ -75,14 +16,24 @@ fn strip_query(path: String) -> String {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn err_404(path: String) -> String {
|
fn err_404(path: String) -> String {
|
||||||
return "{\"error\":\"not found\",\"code\":\"not_found\",\"path\":\"" + path + "\"}"
|
// __status__ envelope — el_runtime reads the first key and emits HTTP 404.
|
||||||
|
// Issue #3: previously returned HTTP 200 with JSON error body.
|
||||||
|
return "{\"__status__\":404,\"error\":\"not found\",\"path\":\"" + path + "\"}"
|
||||||
}
|
}
|
||||||
|
|
||||||
fn err_405(method: String, path: String) -> String {
|
fn err_405(method: String, path: String) -> String {
|
||||||
return "{\"error\":\"method not allowed\",\"code\":\"method_not_allowed\",\"method\":\"" + method + "\",\"path\":\"" + path + "\"}"
|
// __status__ envelope — emits HTTP 405.
|
||||||
|
// Issue #3: previously returned HTTP 200 with JSON error body.
|
||||||
|
return "{\"__status__\":405,\"error\":\"method not allowed\",\"method\":\"" + method + "\",\"path\":\"" + path + "\"}"
|
||||||
}
|
}
|
||||||
|
|
||||||
fn route_health() -> String {
|
fn route_health() -> String {
|
||||||
|
// NOTE (issue #8): This endpoint performs live engram graph queries on every call
|
||||||
|
// (engram_node_count, engram_edge_count) and reads imprint state. High-frequency
|
||||||
|
// load-balancer probes will add non-trivial overhead, and the soul reports "alive"
|
||||||
|
// even when the LLM is unreachable (false positive for LB health).
|
||||||
|
// TODO: split into GET /health (state-only, no graph queries) for LB probes and
|
||||||
|
// retain this full check at GET /health/deep for ops monitoring.
|
||||||
let cgi_id: String = state_get("soul_cgi_id")
|
let cgi_id: String = state_get("soul_cgi_id")
|
||||||
let boot: String = state_get("soul_boot_count")
|
let boot: String = state_get("soul_boot_count")
|
||||||
let boot_num: String = if str_eq(boot, "") { "0" } else { boot }
|
let boot_num: String = if str_eq(boot, "") { "0" } else { boot }
|
||||||
@@ -90,35 +41,12 @@ fn route_health() -> String {
|
|||||||
let edge_ct: Int = engram_edge_count()
|
let edge_ct: Int = engram_edge_count()
|
||||||
let pulse: String = state_get("soul.pulse")
|
let pulse: String = state_get("soul.pulse")
|
||||||
let pulse_num: String = if str_eq(pulse, "") { "0" } else { pulse }
|
let pulse_num: String = if str_eq(pulse, "") { "0" } else { pulse }
|
||||||
|
|
||||||
// Uptime: soul records boot timestamp in state at startup via soul_boot_ts.
|
|
||||||
// Compute elapsed seconds; fall back to -1 if not yet set.
|
|
||||||
let boot_ts_str: String = state_get("soul_boot_ts")
|
|
||||||
let uptime_secs: Int = if str_eq(boot_ts_str, "") {
|
|
||||||
-1
|
|
||||||
} else {
|
|
||||||
time_now() - str_to_int(boot_ts_str)
|
|
||||||
}
|
|
||||||
|
|
||||||
// LLM connectivity: probe with a minimal call. Any non-error reply = ok.
|
|
||||||
// Use a short, fixed prompt so this never counts against conversation history.
|
|
||||||
let model: String = state_get("soul_model")
|
|
||||||
let eff_model: String = if str_eq(model, "") { "claude-sonnet-4-5" } else { model }
|
|
||||||
let llm_probe: String = llm_call_system(eff_model, "You are a health probe. Reply with the single word: ok", "ping")
|
|
||||||
let llm_ok: Bool = !str_eq(llm_probe, "")
|
|
||||||
&& !str_starts_with(llm_probe, "{\"error\"")
|
|
||||||
&& !str_starts_with(llm_probe, "{\"type\":\"error\"")
|
|
||||||
&& !str_contains(llm_probe, "authentication_error")
|
|
||||||
let llm_status: String = if llm_ok { "ok" } else { "unreachable" }
|
|
||||||
|
|
||||||
return "{\"status\":\"alive\""
|
return "{\"status\":\"alive\""
|
||||||
+ ",\"cgi_id\":\"" + cgi_id + "\""
|
+ ",\"cgi_id\":\"" + cgi_id + "\""
|
||||||
+ ",\"boot\":" + boot_num
|
+ ",\"boot\":" + boot_num
|
||||||
+ ",\"uptime_secs\":" + int_to_str(uptime_secs)
|
|
||||||
+ ",\"node_count\":" + int_to_str(node_ct)
|
+ ",\"node_count\":" + int_to_str(node_ct)
|
||||||
+ ",\"edge_count\":" + int_to_str(edge_ct)
|
+ ",\"edge_count\":" + int_to_str(edge_ct)
|
||||||
+ ",\"pulse\":" + pulse_num
|
+ ",\"pulse\":" + pulse_num
|
||||||
+ ",\"llm\":\"" + llm_status + "\""
|
|
||||||
+ ",\"layers\":{\"l0\":\"core\",\"l1\":\"safety\",\"l2\":\"stewardship\",\"l3\":\"" + imprint_current() + "\"}}"
|
+ ",\"layers\":{\"l0\":\"core\",\"l1\":\"safety\",\"l2\":\"stewardship\",\"l3\":\"" + imprint_current() + "\"}}"
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -141,7 +69,8 @@ fn route_lineage() -> String {
|
|||||||
|
|
||||||
fn route_imprint_contextual(body: String) -> String {
|
fn route_imprint_contextual(body: String) -> String {
|
||||||
if str_eq(body, "") {
|
if str_eq(body, "") {
|
||||||
return "{\"ok\":false,\"error\":\"empty body\"}"
|
// Issue #5: empty body is a client error — HTTP 400.
|
||||||
|
return "{\"__status__\":400,\"ok\":false,\"error\":\"empty body\"}"
|
||||||
}
|
}
|
||||||
let tags: String = "[\"imprint\",\"contextual\"]"
|
let tags: String = "[\"imprint\",\"contextual\"]"
|
||||||
let id: String = engram_node_full(
|
let id: String = engram_node_full(
|
||||||
@@ -163,7 +92,8 @@ fn route_imprint_contextual(body: String) -> String {
|
|||||||
|
|
||||||
fn route_imprint_user(body: String) -> String {
|
fn route_imprint_user(body: String) -> String {
|
||||||
if str_eq(body, "") {
|
if str_eq(body, "") {
|
||||||
return "{\"ok\":false,\"error\":\"empty body\"}"
|
// Issue #5: empty body is a client error — HTTP 400.
|
||||||
|
return "{\"__status__\":400,\"ok\":false,\"error\":\"empty body\"}"
|
||||||
}
|
}
|
||||||
let tags: String = "[\"imprint\",\"user\"]"
|
let tags: String = "[\"imprint\",\"user\"]"
|
||||||
let id: String = engram_node_full(
|
let id: String = engram_node_full(
|
||||||
@@ -185,15 +115,15 @@ fn route_imprint_user(body: String) -> String {
|
|||||||
|
|
||||||
fn route_synthesize(body: String) -> String {
|
fn route_synthesize(body: String) -> String {
|
||||||
if str_eq(body, "") {
|
if str_eq(body, "") {
|
||||||
return "{\"error\":\"body is required\",\"code\":\"missing_param\"}"
|
return "{\"mechanism\":\"did not engage\"}"
|
||||||
}
|
}
|
||||||
let parent_a: String = json_get(body, "parent_a")
|
let parent_a: String = json_get(body, "parent_a")
|
||||||
let parent_b: String = json_get(body, "parent_b")
|
let parent_b: String = json_get(body, "parent_b")
|
||||||
if str_eq(parent_a, "") {
|
if str_eq(parent_a, "") {
|
||||||
return "{\"error\":\"parent_a is required\",\"code\":\"missing_param\"}"
|
return "{\"mechanism\":\"did not engage\"}"
|
||||||
}
|
}
|
||||||
if str_eq(parent_b, "") {
|
if str_eq(parent_b, "") {
|
||||||
return "{\"error\":\"parent_b is required\",\"code\":\"missing_param\"}"
|
return "{\"mechanism\":\"did not engage\"}"
|
||||||
}
|
}
|
||||||
let req: String = "synthesize " + parent_a + " " + parent_b
|
let req: String = "synthesize " + parent_a + " " + parent_b
|
||||||
let tags: String = "[\"soul-inbox-pending\",\"synthesis-request\"]"
|
let tags: String = "[\"soul-inbox-pending\",\"synthesis-request\"]"
|
||||||
@@ -301,9 +231,13 @@ fn connectd_get(suffix: String) -> String {
|
|||||||
// so arbitrary JSON cannot reach the shell as a command-line argument.
|
// so arbitrary JSON cannot reach the shell as a command-line argument.
|
||||||
fn connectd_post(suffix: String, body: String) -> String {
|
fn connectd_post(suffix: String, body: String) -> String {
|
||||||
let eff: String = if str_eq(body, "") { "{}" } else { body }
|
let eff: String = if str_eq(body, "") { "{}" } else { body }
|
||||||
// Unique temp path per call — prevents collision if concurrency is ever added
|
// Issue #11: time_now() has second-granularity; two concurrent requests in the same
|
||||||
// or if two soul instances run on the same machine (latent correctness hazard).
|
// second collide on the same temp path. Added a monotonic per-process sequence counter.
|
||||||
let tmp: String = "/tmp/neuron-connectors-req-" + int_to_str(time_now()) + ".json"
|
let connectd_seq_s: String = state_get("connectd_post_seq")
|
||||||
|
let connectd_seq_n: Int = if str_eq(connectd_seq_s, "") { 0 } else { str_to_int(connectd_seq_s) }
|
||||||
|
let connectd_seq_next: Int = connectd_seq_n + 1
|
||||||
|
state_set("connectd_post_seq", int_to_str(connectd_seq_next))
|
||||||
|
let tmp: String = "/tmp/neuron-connectors-req-" + int_to_str(time_now()) + "-" + int_to_str(connectd_seq_next) + ".json"
|
||||||
fs_write(tmp, eff)
|
fs_write(tmp, eff)
|
||||||
let out: String = exec_capture("curl -s --max-time 20 -X POST http://127.0.0.1:7771" + suffix + " -H 'Content-Type: application/json' -d @" + tmp)
|
let out: String = exec_capture("curl -s --max-time 20 -X POST http://127.0.0.1:7771" + suffix + " -H 'Content-Type: application/json' -d @" + tmp)
|
||||||
if str_eq(out, "") {
|
if str_eq(out, "") {
|
||||||
@@ -338,20 +272,45 @@ fn handle_connectors(method: String, clean: String, body: String) -> String {
|
|||||||
return "{\"ok\":false,\"error\":\"unknown connectors route\"}"
|
return "{\"ok\":false,\"error\":\"unknown connectors route\"}"
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
// auth_check — validate NEURON_TOKEN bearer auth on every request.
|
||||||
|
// Returns "" when authorized, or a JSON 401 error string when not.
|
||||||
|
// /health and /lineage are public routes — always exempted.
|
||||||
|
// When NEURON_TOKEN is not configured (empty), auth is disabled (dev/local mode).
|
||||||
|
// Issue #4: previously no auth layer existed anywhere in the router.
|
||||||
|
// Clients pass the token in the JSON body as "__auth".
|
||||||
|
// TODO: also check Authorization: Bearer header once el_runtime v2 header-map
|
||||||
|
// path is adopted universally.
|
||||||
|
fn auth_check(clean: String, body: String) -> String {
|
||||||
|
if str_eq(clean, "/health") { return "" }
|
||||||
|
if str_eq(clean, "/lineage") { return "" }
|
||||||
|
let token: String = state_get("soul_token")
|
||||||
|
if str_eq(token, "") { return "" }
|
||||||
|
let auth_field: String = json_get(body, "__auth")
|
||||||
|
if str_eq(auth_field, token) { return "" }
|
||||||
|
return "{\"__status__\":401,\"error\":\"unauthorized\"}"
|
||||||
|
}
|
||||||
|
|
||||||
fn handle_request(method: String, path: String, body: String) -> String {
|
fn handle_request(method: String, path: String, body: String) -> String {
|
||||||
let clean: String = strip_query(path)
|
let clean: String = strip_query(path)
|
||||||
|
|
||||||
// Rate limit check. Extract caller IP from REMOTE_ADDR env var (set by the
|
// Issue #1/#2: EL has no exception/try-catch mechanism. A C-level crash inside
|
||||||
// EL HTTP runtime for each request). Skip enforcement when empty so
|
// an http_worker pthread drops the TCP connection (client gets RST) rather than
|
||||||
// loopback/internal callers are never blocked.
|
// returning HTTP 500. TODO: register a SIGSEGV/SIGBUS handler in el_runtime.c
|
||||||
let ip: String = env("REMOTE_ADDR")
|
// that writes a 500 JSON response to the current worker fd before aborting.
|
||||||
if !str_eq(ip, "") {
|
|
||||||
let rl_result: String = rate_limit_check(ip, clean)
|
// Issue #10: Rate limiting is not implemented.
|
||||||
if !str_eq(rl_result, "") {
|
// TODO: add a per-IP token-bucket counter returning HTTP 429 when exceeded.
|
||||||
return rl_result
|
// Requires a C-level counter in el_runtime.c or a sidecar reverse proxy.
|
||||||
}
|
|
||||||
|
// Auth — enforced on all routes except /health and /lineage.
|
||||||
|
// Issue #4: previously no auth check existed anywhere in the router.
|
||||||
|
let auth_err: String = auth_check(clean, body)
|
||||||
|
if !str_eq(auth_err, "") {
|
||||||
|
return auth_err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
if str_eq(method, "POST") && str_eq(clean, "/dharma/recv") {
|
if str_eq(method, "POST") && str_eq(clean, "/dharma/recv") {
|
||||||
return handle_dharma_recv(body)
|
return handle_dharma_recv(body)
|
||||||
}
|
}
|
||||||
@@ -379,7 +338,8 @@ fn handle_request(method: String, path: String, body: String) -> String {
|
|||||||
let raw_msg: String = json_get(body, "message")
|
let raw_msg: String = json_get(body, "message")
|
||||||
let eff_msg: String = if str_eq(raw_msg, "") { body } else { raw_msg }
|
let eff_msg: String = if str_eq(raw_msg, "") { body } else { raw_msg }
|
||||||
if str_eq(eff_msg, "") {
|
if str_eq(eff_msg, "") {
|
||||||
return "{\"error\":\"message is required\",\"code\":\"missing_param\"}"
|
// Issue #5: missing required param — HTTP 400.
|
||||||
|
return "{\"__status__\":400,\"error\":\"message required\"}"
|
||||||
}
|
}
|
||||||
let agentic_flag: Bool = json_get_bool(body, "agentic")
|
let agentic_flag: Bool = json_get_bool(body, "agentic")
|
||||||
let reply: String = if agentic_flag {
|
let reply: String = if agentic_flag {
|
||||||
@@ -519,13 +479,15 @@ fn handle_request(method: String, path: String, body: String) -> String {
|
|||||||
return handle_elp_chat(body)
|
return handle_elp_chat(body)
|
||||||
}
|
}
|
||||||
if str_eq(clean, "/api/chat") {
|
if str_eq(clean, "/api/chat") {
|
||||||
// NOTE: streaming (SSE / chunked transfer) is not implemented. All chat
|
// Issue #5: validate required params — return HTTP 400 when missing.
|
||||||
// responses are buffered and returned as a single JSON object. Streaming
|
|
||||||
// would require runtime-level SSE support in el_runtime.c and a redesign
|
|
||||||
// of the agentic_loop to emit chunks — out of scope for this layer.
|
|
||||||
let raw_msg: String = json_get(body, "message")
|
let raw_msg: String = json_get(body, "message")
|
||||||
if str_eq(raw_msg, "") {
|
if str_eq(raw_msg, "") {
|
||||||
return "{\"error\":\"message is required\",\"code\":\"missing_param\"}"
|
return "{\"__status__\":400,\"error\":\"message is required\",\"response\":\"\"}"
|
||||||
|
}
|
||||||
|
// Issue #7: reject oversized messages before engram_compile and the LLM.
|
||||||
|
// Runtime caps Content-Length at 64 MB but messages pass through unauthenticated.
|
||||||
|
if str_len(raw_msg) > 32768 {
|
||||||
|
return "{\"__status__\":400,\"error\":\"message too large (max 32768 chars)\",\"response\":\"\"}"
|
||||||
}
|
}
|
||||||
let agentic_flag: Bool = json_get_bool(body, "agentic")
|
let agentic_flag: Bool = json_get_bool(body, "agentic")
|
||||||
let reply: String = if agentic_flag {
|
let reply: String = if agentic_flag {
|
||||||
|
|||||||
@@ -369,7 +369,6 @@ load_identity_context()
|
|||||||
seed_persona_from_env()
|
seed_persona_from_env()
|
||||||
let boot_num: Int = mem_boot_count_inc()
|
let boot_num: Int = mem_boot_count_inc()
|
||||||
state_set("soul_boot_count", int_to_str(boot_num))
|
state_set("soul_boot_count", int_to_str(boot_num))
|
||||||
state_set("soul_boot_ts", int_to_str(time_now()))
|
|
||||||
println("[soul] boot #" + int_to_str(boot_num))
|
println("[soul] boot #" + int_to_str(boot_num))
|
||||||
emit_session_start_event()
|
emit_session_start_event()
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user