diff --git a/.mcp.json b/.mcp.json index 84863fb8..b9cd6485 100644 --- a/.mcp.json +++ b/.mcp.json @@ -121,7 +121,7 @@ "_comment": "Local ENE replacement for ContextStream-like memory/search/session recall. Writes local receipt-bearing SQLite memory and reads ene-api/session-sync when running.", "command": "python3", "args": [ - "4-Infrastructure/infra/ene_contextstream_mcp.py" + "/home/allaun/Research Stack/4-Infrastructure/infra/ene_contextstream_mcp.py" ], "env": { "ENE_API_URL": "http://127.0.0.1:3000", diff --git a/4-Infrastructure/infra/ene-rds/Cargo.lock b/4-Infrastructure/infra/ene-rds/Cargo.lock index 05e8d08f..8e5fc4e8 100644 --- a/4-Infrastructure/infra/ene-rds/Cargo.lock +++ b/4-Infrastructure/infra/ene-rds/Cargo.lock @@ -178,6 +178,15 @@ version = "2.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c4512299f36f043ab09a583e57bceb5a5aab7a73db1805848e8fef3c9e8c78b3" +[[package]] +name = "block-buffer" +version = "0.10.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3078c7629b62d3f0439517fa394996acacc5cbc91c5a20d8c658e77abd503a71" +dependencies = [ + "generic-array", +] + [[package]] name = "block-buffer" version = "0.12.0" @@ -228,8 +237,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6f8d983286843e49675a4b7a2d174efe136dc93a18d69130dd18198a6c167601" dependencies = [ "cfg-if", - "cpufeatures", - "rand_core", + "cpufeatures 0.3.0", + "rand_core 0.10.1", ] [[package]] @@ -330,6 +339,15 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" +[[package]] +name = "cpufeatures" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "59ed5838eebb26a2bb2e58f6d5b5316989ae9d08bab10e0e6d103e656d1b0280" +dependencies = [ + "libc", +] + [[package]] name = "cpufeatures" version = "0.3.0" @@ -339,6 +357,16 @@ dependencies = [ "libc", ] +[[package]] +name = "crypto-common" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "78c8292055d1c1df0cce5d180393dc8cce0abec0a7102adb6c7b1eef6016d60a" +dependencies = [ + "generic-array", + "typenum", +] + [[package]] name = "crypto-common" version = "0.2.1" @@ -357,15 +385,26 @@ dependencies = [ "cmov", ] +[[package]] +name = "digest" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ed9a281f7bc9b7576e61468ba615a66a5c8cfdff42420a70aa82701a3b1e292" +dependencies = [ + "block-buffer 0.10.4", + "crypto-common 0.1.7", + "subtle", +] + [[package]] name = "digest" version = "0.11.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" dependencies = [ - "block-buffer", + "block-buffer 0.12.0", "const-oid", - "crypto-common", + "crypto-common 0.2.1", "ctutils", ] @@ -428,6 +467,26 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "ene-node" +version = "0.1.0" +dependencies = [ + "anyhow", + "base64", + "chrono", + "clap", + "hex", + "hmac 0.12.1", + "rand 0.8.6", + "rusqlite", + "serde", + "serde_json", + "sha2 0.10.9", + "tokio", + "tracing", + "tracing-subscriber", +] + [[package]] name = "ene-rds-chat" version = "0.1.0" @@ -449,6 +508,8 @@ version = "0.1.0" dependencies = [ "anyhow", "chrono", + "native-tls", + "postgres-native-tls", "serde", "serde_json", "tokio", @@ -485,6 +546,25 @@ dependencies = [ "tracing", ] +[[package]] +name = "ene-storage" +version = "0.1.0" +dependencies = [ + "anyhow", + "chrono", + "clap", + "dirs", + "hex", + "reqwest", + "serde", + "serde_json", + "sha2 0.10.9", + "tempfile", + "tokio", + "tracing", + "tracing-subscriber", +] + [[package]] name = "ene-sync" version = "0.1.0" @@ -627,6 +707,16 @@ dependencies = [ "slab", ] +[[package]] +name = "generic-array" +version = "0.14.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85649ca51fd72272d7821adaf274ad91c288277713d9c18820d8499a7ff69e9a" +dependencies = [ + "typenum", + "version_check", +] + [[package]] name = "getrandom" version = "0.2.17" @@ -647,7 +737,7 @@ dependencies = [ "cfg-if", "libc", "r-efi", - "rand_core", + "rand_core 0.10.1", "wasip2", "wasip3", ] @@ -710,13 +800,28 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hex" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" + +[[package]] +name = "hmac" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c49c37c09c17a53d937dfbb742eb3a961d65a994e6bcdcf37e7399d0cc8ab5e" +dependencies = [ + "digest 0.10.7", +] + [[package]] name = "hmac" version = "0.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6303bc9732ae41b04cb554b844a762b4115a61bfaa81e3e83050991eeb56863f" dependencies = [ - "digest", + "digest 0.11.3", ] [[package]] @@ -1113,7 +1218,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69b6441f590336821bb897fb28fc622898ccceb1d6cea3fde5ea86b090c4de98" dependencies = [ "cfg-if", - "digest", + "digest 0.11.3", ] [[package]] @@ -1313,6 +1418,18 @@ version = "0.3.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e" +[[package]] +name = "postgres-native-tls" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fef4de47bb81477e0c3deaf153a1b10ae176484713ff1640969f4cb96b653ebc" +dependencies = [ + "native-tls", + "tokio", + "tokio-native-tls", + "tokio-postgres", +] + [[package]] name = "postgres-protocol" version = "0.6.11" @@ -1323,11 +1440,11 @@ dependencies = [ "byteorder", "bytes", "fallible-iterator 0.2.0", - "hmac", + "hmac 0.13.0", "md-5", "memchr", - "rand", - "sha2", + "rand 0.10.1", + "sha2 0.11.0", "stringprep", ] @@ -1354,6 +1471,15 @@ dependencies = [ "zerovec", ] +[[package]] +name = "ppv-lite86" +version = "0.2.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" +dependencies = [ + "zerocopy", +] + [[package]] name = "prettyplease" version = "0.2.37" @@ -1388,6 +1514,17 @@ version = "6.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" +[[package]] +name = "rand" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5ca0ecfa931c29007047d1bc58e623ab12e5590e8c7cc53200d5202b69266d8a" +dependencies = [ + "libc", + "rand_chacha", + "rand_core 0.6.4", +] + [[package]] name = "rand" version = "0.10.1" @@ -1396,7 +1533,26 @@ checksum = "d2e8e8bcc7961af1fdac401278c6a831614941f6164ee3bf4ce61b7edb162207" dependencies = [ "chacha20", "getrandom 0.4.2", - "rand_core", + "rand_core 0.10.1", +] + +[[package]] +name = "rand_chacha" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" +dependencies = [ + "ppv-lite86", + "rand_core 0.6.4", +] + +[[package]] +name = "rand_core" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" +dependencies = [ + "getrandom 0.2.17", ] [[package]] @@ -1679,6 +1835,17 @@ dependencies = [ "serde", ] +[[package]] +name = "sha2" +version = "0.10.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" +dependencies = [ + "cfg-if", + "cpufeatures 0.2.17", + "digest 0.10.7", +] + [[package]] name = "sha2" version = "0.11.0" @@ -1686,8 +1853,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "446ba717509524cb3f22f17ecc096f10f4822d76ab5c0b9822c5f9c284e825f4" dependencies = [ "cfg-if", - "cpufeatures", - "digest", + "cpufeatures 0.3.0", + "digest 0.11.3", ] [[package]] @@ -1948,7 +2115,7 @@ dependencies = [ "pin-project-lite", "postgres-protocol", "postgres-types", - "rand", + "rand 0.10.1", "socket2", "tokio", "tokio-util", diff --git a/4-Infrastructure/infra/ene-rds/Cargo.toml b/4-Infrastructure/infra/ene-rds/Cargo.toml index 110904e6..2f682691 100644 --- a/4-Infrastructure/infra/ene-rds/Cargo.toml +++ b/4-Infrastructure/infra/ene-rds/Cargo.toml @@ -15,6 +15,8 @@ clap = { version = "4", features = ["derive"] } reqwest = { version = "0.12", features = ["json"] } serde = { version = "1", features = ["derive"] } serde_json = "1" +native-tls = "0.2" +postgres-native-tls = "0.5" tokio = { version = "1", features = ["full"] } tokio-postgres = { version = "0.7", features = ["with-chrono-0_4", "with-serde_json-1"] } tracing = "0.1" diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-api/src/main.rs b/4-Infrastructure/infra/ene-rds/crates/ene-api/src/main.rs index 41f0b5d6..325555c4 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-api/src/main.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-api/src/main.rs @@ -11,6 +11,7 @@ use ene_rds_wiki::WikiSurface; use serde::{Deserialize, Serialize}; use serde_json::json; use std::collections::HashMap; +use std::net::SocketAddr; use std::sync::Arc; use tokio::sync::Mutex; @@ -29,7 +30,9 @@ struct SearchQuery { semantic: bool, } -fn default_limit() -> i64 { 10 } +fn default_limit() -> i64 { + 10 +} #[derive(Serialize)] struct HealthResponse { @@ -58,7 +61,11 @@ async fn main() -> anyhow::Result<()> { let ephemeral = EphemeralSurface::new(ephemeral_client); ephemeral.init_tables().await?; - let state = Arc::new(Mutex::new(AppState { chat, wiki, ephemeral })); + let state = Arc::new(Mutex::new(AppState { + chat, + wiki, + ephemeral, + })); let app = Router::new() .route("/health", get(health_handler)) @@ -71,8 +78,11 @@ async fn main() -> anyhow::Result<()> { .route("/ephemeral/nodes/:id", get(get_ephemeral_node)) .with_state(state); - let listener = tokio::net::TcpListener::bind("0.0.0.0:3000").await?; - tracing::info!("ENE API listening on http://0.0.0.0:3000"); + let bind_addr: SocketAddr = std::env::var("ENE_API_BIND") + .unwrap_or_else(|_| "0.0.0.0:3000".into()) + .parse()?; + let listener = tokio::net::TcpListener::bind(bind_addr).await?; + tracing::info!("ENE API listening on http://{bind_addr}"); axum::serve(listener, app).await?; Ok(()) } @@ -88,7 +98,10 @@ async fn list_sessions( State(state): State>>, Query(params): Query>, ) -> Result, String> { - let limit = params.get("limit").and_then(|s| s.parse().ok()).unwrap_or(10i64); + let limit = params + .get("limit") + .and_then(|s| s.parse().ok()) + .unwrap_or(10i64); let guard = state.lock().await; match guard.chat.list_sessions(limit).await { Ok(v) => Ok(Json(json!({"ok": true, "data": v}))), @@ -128,7 +141,10 @@ async fn wiki_search( Query(params): Query>, ) -> Result, String> { let query = params.get("q").cloned().unwrap_or_default(); - let limit = params.get("limit").and_then(|s| s.parse().ok()).unwrap_or(10i64); + let limit = params + .get("limit") + .and_then(|s| s.parse().ok()) + .unwrap_or(10i64); let guard = state.lock().await; match guard.wiki.search(&query, limit).await { Ok(v) => Ok(Json(json!({"ok": true, "data": v}))), @@ -153,7 +169,10 @@ async fn list_ephemeral_nodes( Query(params): Query>, ) -> Result, String> { let zone = params.get("zone").cloned().unwrap_or_else(|| "cold".into()); - let limit = params.get("limit").and_then(|s| s.parse().ok()).unwrap_or(100i64); + let limit = params + .get("limit") + .and_then(|s| s.parse().ok()) + .unwrap_or(100i64); let guard = state.lock().await; match guard.ephemeral.list_nodes_by_zone(&zone, limit).await { Ok(v) => Ok(Json(json!({"ok": true, "data": v}))), diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-rds-chat/src/lib.rs b/4-Infrastructure/infra/ene-rds/crates/ene-rds-chat/src/lib.rs index 79e3261d..72e4b910 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-rds-chat/src/lib.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-rds-chat/src/lib.rs @@ -1,6 +1,5 @@ use anyhow::{Context, Result}; use ene_rds_core::{vec_to_pgtext, RdsClient}; -use serde::Serialize; use tracing::info; pub mod models; @@ -64,7 +63,11 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_receipt ON ene.chat_messages(receip CREATE INDEX IF NOT EXISTS idx_chat_messages_text_search ON ene.chat_messages USING GIN(to_tsvector('english', text_content)); CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages USING GIN(tool_calls jsonb_path_ops); "#; - self.client.inner().batch_execute(ddl).await.context("init chat DDL")?; + self.client + .inner() + .batch_execute(ddl) + .await + .context("init chat DDL")?; info!("chat log schema initialized"); Ok(()) } @@ -160,8 +163,13 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US } pub async fn delete_messages_for_session(&self, session_id: &str) -> Result { - let rows = self.client.inner() - .execute("DELETE FROM ene.chat_messages WHERE session_id = $1", &[&session_id]) + let rows = self + .client + .inner() + .execute( + "DELETE FROM ene.chat_messages WHERE session_id = $1", + &[&session_id], + ) .await .context("delete messages")?; Ok(rows) @@ -170,33 +178,46 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US pub async fn search_keyword(&self, query: &str, limit: i64) -> Result> { let rows = self.client.inner() .query( - "SELECT s.session_id, s.title, s.agent, s.model, \ + "SELECT s.session_id, COALESCE(s.compaction_summary, s.session_id) AS title, \ + s.meta->>'agent' AS agent, s.meta->>'model' AS model, \ COUNT(m.id) AS match_count, \ - MAX(ts_rank(to_tsvector('english', m.text_content), plainto_tsquery('english', $1))) AS rank \ + MAX(ts_rank(to_tsvector('english', COALESCE(m.text_content, '')), plainto_tsquery('english', $1))) AS rank \ FROM ene.chat_sessions s \ JOIN ene.chat_messages m ON m.session_id = s.session_id \ - WHERE to_tsvector('english', m.text_content) @@ plainto_tsquery('english', $1) \ - GROUP BY s.session_id, s.title, s.agent, s.model \ + WHERE to_tsvector('english', COALESCE(m.text_content, '')) @@ plainto_tsquery('english', $1) \ + GROUP BY s.session_id, 2, 3, 4 \ ORDER BY rank DESC LIMIT $2", &[&query, &limit], ) .await .context("keyword search")?; - Ok(rows.iter().map(|r| serde_json::json!({ - "session_id": r.get::<_, String>(0), - "title": r.get::<_, String>(1), - "agent": r.get::<_, Option>(2), - "model": r.get::<_, Option>(3), - "match_count": r.get::<_, i64>(4), - "rank": r.get::<_, f32>(5), - })).collect()) + Ok(rows + .iter() + .map(|r| { + serde_json::json!({ + "session_id": r.get::<_, String>(0), + "title": r.get::<_, String>(1), + "agent": r.get::<_, Option>(2), + "model": r.get::<_, Option>(3), + "match_count": r.get::<_, i64>(4), + "rank": r.get::<_, f32>(5), + }) + }) + .collect()) } - pub async fn search_similar(&self, embedding: &[f32], limit: i64) -> Result> { + pub async fn search_similar( + &self, + embedding: &[f32], + limit: i64, + ) -> Result> { let vec_str = vec_to_pgtext(embedding); - let rows = self.client.inner() + let rows = self + .client + .inner() .query( - "SELECT session_id, title, agent, model, \ + "SELECT session_id, COALESCE(compaction_summary, session_id) AS title, \ + meta->>'agent' AS agent, meta->>'model' AS model, \ 1 - (embedding <=> $1::vector) AS similarity \ FROM ene.chat_sessions WHERE embedding IS NOT NULL \ ORDER BY embedding <=> $1::vector LIMIT $2", @@ -204,42 +225,58 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US ) .await .context("similarity search")?; - Ok(rows.iter().map(|r| serde_json::json!({ - "session_id": r.get::<_, String>(0), - "title": r.get::<_, String>(1), - "agent": r.get::<_, Option>(2), - "model": r.get::<_, Option>(3), - "similarity": r.get::<_, f32>(4), - })).collect()) + Ok(rows + .iter() + .map(|r| { + serde_json::json!({ + "session_id": r.get::<_, String>(0), + "title": r.get::<_, String>(1), + "agent": r.get::<_, Option>(2), + "model": r.get::<_, Option>(3), + "similarity": r.get::<_, f32>(4), + }) + }) + .collect()) } pub async fn list_sessions(&self, limit: i64) -> Result> { - let rows = self.client.inner() + let rows = self + .client + .inner() .query( - "SELECT session_id, title, agent, model, message_count, \ + "SELECT session_id, COALESCE(compaction_summary, session_id) AS title, \ + meta->>'agent' AS agent, meta->>'model' AS model, message_count, \ token_input_total, token_output_total, created_at_ms, updated_at_ms \ FROM ene.chat_sessions ORDER BY updated_at_ms DESC LIMIT $1", &[&limit], ) .await .context("list sessions")?; - Ok(rows.iter().map(|r| serde_json::json!({ - "session_id": r.get::<_, String>(0), - "title": r.get::<_, String>(1), - "agent": r.get::<_, Option>(2), - "model": r.get::<_, Option>(3), - "message_count": r.get::<_, i32>(4), - "token_input_total": r.get::<_, i64>(5), - "token_output_total": r.get::<_, i64>(6), - "created_at_ms": r.get::<_, i64>(7), - "updated_at_ms": r.get::<_, i64>(8), - })).collect()) + Ok(rows + .iter() + .map(|r| { + serde_json::json!({ + "session_id": r.get::<_, String>(0), + "title": r.get::<_, String>(1), + "agent": r.get::<_, Option>(2), + "model": r.get::<_, Option>(3), + "message_count": r.get::<_, i32>(4), + "token_input_total": r.get::<_, i64>(5), + "token_output_total": r.get::<_, i64>(6), + "created_at_ms": r.get::<_, i64>(7), + "updated_at_ms": r.get::<_, i64>(8), + }) + }) + .collect()) } pub async fn get_session(&self, session_id: &str) -> Result> { - let sess = self.client.inner() + let sess = self + .client + .inner() .query_opt( - "SELECT session_id, title, agent, model, message_count, \ + "SELECT session_id, COALESCE(compaction_summary, session_id) AS title, \ + meta->>'agent' AS agent, meta->>'model' AS model, message_count, \ token_input_total, token_output_total, created_at_ms, updated_at_ms, meta \ FROM ene.chat_sessions WHERE session_id = $1", &[&session_id], @@ -247,7 +284,9 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US .await .context("get session")?; let Some(sess) = sess else { return Ok(None) }; - let msgs = self.client.inner() + let msgs = self + .client + .inner() .query( "SELECT message_index, role, blocks, text_content, token_input, \ token_output, tool_calls, created_at_ms \ @@ -256,16 +295,21 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US ) .await .context("get messages")?; - let messages: Vec<_> = msgs.iter().map(|r| serde_json::json!({ - "message_index": r.get::<_, i32>(0), - "role": r.get::<_, String>(1), - "blocks": r.get::<_, serde_json::Value>(2), - "text_content": r.get::<_, String>(3), - "token_input": r.get::<_, i64>(4), - "token_output": r.get::<_, i64>(5), - "tool_calls": r.get::<_, serde_json::Value>(6), - "created_at_ms": r.get::<_, i64>(7), - })).collect(); + let messages: Vec<_> = msgs + .iter() + .map(|r| { + serde_json::json!({ + "message_index": r.get::<_, i32>(0), + "role": r.get::<_, String>(1), + "blocks": r.get::<_, serde_json::Value>(2), + "text_content": r.get::<_, String>(3), + "token_input": r.get::<_, i64>(4), + "token_output": r.get::<_, i64>(5), + "tool_calls": r.get::<_, serde_json::Value>(6), + "created_at_ms": r.get::<_, i64>(7), + }) + }) + .collect(); Ok(Some(serde_json::json!({ "session_id": sess.get::<_, String>(0), "title": sess.get::<_, String>(1), diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/Cargo.toml b/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/Cargo.toml index a5c0e979..2bfbc1c0 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/Cargo.toml +++ b/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/Cargo.toml @@ -10,6 +10,8 @@ anyhow = { workspace = true } chrono = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } +native-tls = { workspace = true } +postgres-native-tls = { workspace = true } tokio = { workspace = true } tokio-postgres = { workspace = true } tracing = { workspace = true } diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/lib.rs b/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/lib.rs index b47e2744..918abf1c 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/lib.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/lib.rs @@ -1,5 +1,7 @@ use anyhow::{Context, Result}; -use tokio_postgres::{Client, Config, NoTls}; +use native_tls::TlsConnector; +use postgres_native_tls::MakeTlsConnector; +use tokio_postgres::{Client, Config}; use tracing::{info, warn}; pub mod types; @@ -13,7 +15,11 @@ impl RdsClient { /// Connect from a libpq key=value DSN string. pub async fn connect(dsn: &str) -> Result { let config: Config = dsn.parse().context("parse PostgreSQL DSN")?; - let (client, connection) = config.connect(NoTls).await.context("connect to RDS")?; + let connector = TlsConnector::builder() + .build() + .context("build native TLS connector")?; + let connector = MakeTlsConnector::new(connector); + let (client, connection) = config.connect(connector).await.context("connect to RDS")?; tokio::spawn(async move { if let Err(e) = connection.await { warn!("PostgreSQL connection error: {}", e); @@ -27,8 +33,9 @@ impl RdsClient { if let Ok(dsn) = std::env::var("RDS_DSN") { return dsn; } - let host = std::env::var("RDS_HOST") - .unwrap_or_else(|_| "database-1.cluster-c9i0w8eu8fnv.us-east-2.rds.amazonaws.com".into()); + let host = std::env::var("RDS_HOST").unwrap_or_else(|_| { + "database-1.cluster-c9i0w8eu8fnv.us-east-2.rds.amazonaws.com".into() + }); let port = std::env::var("RDS_PORT").unwrap_or_else(|_| "5432".into()); let user = std::env::var("RDS_USER").unwrap_or_else(|_| "postgres".into()); let password = std::env::var("RDS_PASSWORD") @@ -42,7 +49,11 @@ impl RdsClient { } /// Raw query helper. - pub async fn execute(&self, sql: &str, params: &[&(dyn tokio_postgres::types::ToSql + Sync)]) -> Result { + pub async fn execute( + &self, + sql: &str, + params: &[&(dyn tokio_postgres::types::ToSql + Sync)], + ) -> Result { let rows = self.client.execute(sql, params).await?; Ok(rows) } @@ -112,5 +123,11 @@ pub fn sha256_text(text: &str) -> String { /// Format a float vector as pgvector text: [0.1,0.2,...] pub fn vec_to_pgtext(v: &[f32]) -> String { - format!("[{}]", v.iter().map(|f| f.to_string()).collect::>().join(",")) + format!( + "[{}]", + v.iter() + .map(|f| f.to_string()) + .collect::>() + .join(",") + ) } diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/types.rs b/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/types.rs index 45831cfb..83b3884c 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/types.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/types.rs @@ -20,10 +20,18 @@ pub struct ApiResponse { impl ApiResponse { pub fn success(data: serde_json::Value) -> Self { - Self { ok: true, data: Some(data), error: None } + Self { + ok: true, + data: Some(data), + error: None, + } } pub fn fail(msg: impl Into) -> Self { - Self { ok: false, data: None, error: Some(msg.into()) } + Self { + ok: false, + data: None, + error: Some(msg.into()), + } } } diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-rds-ephemeral/src/lib.rs b/4-Infrastructure/infra/ene-rds/crates/ene-rds-ephemeral/src/lib.rs index 58a622f9..df525620 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-rds-ephemeral/src/lib.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-rds-ephemeral/src/lib.rs @@ -146,7 +146,11 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_tasks_state ON ene.ephemeral_tasks(task CREATE INDEX IF NOT EXISTS idx_ephemeral_tasks_session ON ene.ephemeral_tasks(session_id); CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events(node_id, created_at_ms DESC); "#; - self.client.inner().batch_execute(ddl).await.context("init ephemeral DDL")?; + self.client + .inner() + .batch_execute(ddl) + .await + .context("init ephemeral DDL")?; info!("ephemeral node schema initialized"); Ok(()) } @@ -215,23 +219,27 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events( ) .await .context("list ephemeral nodes")?; - Ok(rows.iter().map(|r| EphemeralNode { - node_id: r.get(0), - thermal_zone: r.get(1), - reliability_raw: r.get(2), - latency_p95_ms: r.get(3), - scar_count: r.get(4), - last_seen_ms: r.get(5), - reputation_raw: r.get(6), - quarantine_until_ms: r.get(7), - meta: r.get(8), - created_at_ms: r.get(9), - updated_at_ms: r.get(10), - }).collect()) + Ok(rows + .iter() + .map(|r| EphemeralNode { + node_id: r.get(0), + thermal_zone: r.get(1), + reliability_raw: r.get(2), + latency_p95_ms: r.get(3), + scar_count: r.get(4), + last_seen_ms: r.get(5), + reputation_raw: r.get(6), + quarantine_until_ms: r.get(7), + meta: r.get(8), + created_at_ms: r.get(9), + updated_at_ms: r.get(10), + }) + .collect()) } pub async fn insert_task(&self, task: &EphemeralTask) -> Result<()> { - self.client.inner() + self.client + .inner() .execute( "INSERT INTO ene.ephemeral_tasks \ (task_id, session_id, node_id, task_state, priority_raw, ttl_ms, \ @@ -244,9 +252,17 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events( result_hash = EXCLUDED.result_hash, \ meta = EXCLUDED.meta", &[ - &task.task_id, &task.session_id, &task.node_id, &task.task_state, - &task.priority_raw, &task.ttl_ms, &task.dispatched_at_ms, - &task.completed_at_ms, &task.result_hash, &task.meta, &task.created_at_ms, + &task.task_id, + &task.session_id, + &task.node_id, + &task.task_state, + &task.priority_raw, + &task.ttl_ms, + &task.dispatched_at_ms, + &task.completed_at_ms, + &task.result_hash, + &task.meta, + &task.created_at_ms, ], ) .await @@ -275,7 +291,8 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events( } pub async fn insert_metric(&self, metric: &EphemeralMetric) -> Result<()> { - self.client.inner() + self.client + .inner() .execute( "INSERT INTO ene.ephemeral_metrics \ (metric_id, node_id, metric_name, metric_value_raw, metric_scale, recorded_at_ms) \ @@ -285,8 +302,12 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events( metric_scale = EXCLUDED.metric_scale, \ recorded_at_ms = EXCLUDED.recorded_at_ms", &[ - &metric.metric_id, &metric.node_id, &metric.metric_name, - &metric.metric_value_raw, &(metric.metric_scale as i32), &metric.recorded_at_ms, + &metric.metric_id, + &metric.node_id, + &metric.metric_name, + &metric.metric_value_raw, + &(metric.metric_scale as i32), + &metric.recorded_at_ms, ], ) .await @@ -294,7 +315,11 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events( Ok(()) } - pub async fn get_node_scars(&self, node_id: &str, limit: i64) -> Result> { + pub async fn get_node_scars( + &self, + node_id: &str, + limit: i64, + ) -> Result> { let rows = self.client.inner() .query( "SELECT scar_id, node_id, task_id, scar_pressure, failure_mode, coarsening_agent, created_at_ms \ @@ -303,19 +328,24 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events( ) .await .context("get node scars")?; - Ok(rows.iter().map(|r| EphemeralScarEvent { - scar_id: r.get(0), - node_id: r.get(1), - task_id: r.get(2), - scar_pressure: r.get(3), - failure_mode: r.get(4), - coarsening_agent: r.get(5), - created_at_ms: r.get(6), - }).collect()) + Ok(rows + .iter() + .map(|r| EphemeralScarEvent { + scar_id: r.get(0), + node_id: r.get(1), + task_id: r.get(2), + scar_pressure: r.get(3), + failure_mode: r.get(4), + coarsening_agent: r.get(5), + created_at_ms: r.get(6), + }) + .collect()) } pub async fn get_task_receipt(&self, task_id: &str) -> Result> { - let row = self.client.inner() + let row = self + .client + .inner() .query_opt( "SELECT receipt_id, task_id, node_id, cross_matrix, sidon_slack, \ step_count, residual_series, write_timing_ms, scar_absent, created_at_ms \ diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-rds-wiki/src/lib.rs b/4-Infrastructure/infra/ene-rds/crates/ene-rds-wiki/src/lib.rs index e4e6cde3..e5743d23 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-rds-wiki/src/lib.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-rds-wiki/src/lib.rs @@ -81,13 +81,26 @@ CREATE INDEX IF NOT EXISTS idx_wiki_pages_slug ON ene.wiki_pages(slug); CREATE INDEX IF NOT EXISTS idx_wiki_pages_title ON ene.wiki_pages USING gin(to_tsvector('english', title)); CREATE INDEX IF NOT EXISTS idx_wiki_revisions_page ON ene.wiki_revisions(page_id, created_at DESC); "#; - self.client.inner().batch_execute(ddl).await.context("init wiki DDL")?; + self.client + .inner() + .batch_execute(ddl) + .await + .context("init wiki DDL")?; info!("wiki schema initialized"); Ok(()) } - pub async fn put_page(&self, title: &str, content: &str, editor: &str, summary: &str) -> Result { - let slug = title.to_lowercase().replace(' ', "-").replace(|c: char| !c.is_alphanumeric() && c != '-', ""); + pub async fn put_page( + &self, + title: &str, + content: &str, + editor: &str, + summary: &str, + ) -> Result { + let slug = title + .to_lowercase() + .replace(' ', "-") + .replace(|c: char| !c.is_alphanumeric() && c != '-', ""); let row = self.client.inner() .query_one( "INSERT INTO ene.wiki_pages (title, slug, content, updated_at) \ @@ -119,7 +132,9 @@ CREATE INDEX IF NOT EXISTS idx_wiki_revisions_page ON ene.wiki_revisions(page_id } pub async fn get_page(&self, slug: &str) -> Result> { - let row = self.client.inner() + let row = self + .client + .inner() .query_opt( "SELECT id, title, slug, content, concept_anchor, concept_vector, \ created_at::text, updated_at::text FROM ene.wiki_pages WHERE slug = $1", @@ -150,12 +165,17 @@ CREATE INDEX IF NOT EXISTS idx_wiki_revisions_page ON ene.wiki_revisions(page_id ) .await .context("search wiki")?; - Ok(rows.iter().map(|r| json!({ - "id": r.get::<_, i64>(0), - "title": r.get::<_, String>(1), - "slug": r.get::<_, String>(2), - "rank": r.get::<_, f32>(3), - })).collect()) + Ok(rows + .iter() + .map(|r| { + json!({ + "id": r.get::<_, i64>(0), + "title": r.get::<_, String>(1), + "slug": r.get::<_, String>(2), + "rank": r.get::<_, f32>(3), + }) + }) + .collect()) } pub async fn recent(&self, limit: i64) -> Result> { @@ -167,15 +187,18 @@ CREATE INDEX IF NOT EXISTS idx_wiki_revisions_page ON ene.wiki_revisions(page_id ) .await .context("recent wiki pages")?; - Ok(rows.iter().map(|r| WikiPage { - id: r.get(0), - title: r.get(1), - slug: r.get(2), - content: r.get(3), - concept_anchor: r.get(4), - concept_vector: r.get(5), - created_at: r.get(6), - updated_at: r.get(7), - }).collect()) + Ok(rows + .iter() + .map(|r| WikiPage { + id: r.get(0), + title: r.get(1), + slug: r.get(2), + content: r.get(3), + concept_anchor: r.get(4), + concept_vector: r.get(5), + created_at: r.get(6), + updated_at: r.get(7), + }) + .collect()) } } diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-sync/src/main.rs b/4-Infrastructure/infra/ene-rds/crates/ene-sync/src/main.rs index 4f4e35b3..26dc090e 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-sync/src/main.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-sync/src/main.rs @@ -24,8 +24,14 @@ struct Cli { #[derive(Subcommand, Debug)] enum Commands { - Sync { #[arg(long)] since: Option }, - Watch { #[arg(long, default_value = "60")] interval: u64 }, + Sync { + #[arg(long)] + since: Option, + }, + Watch { + #[arg(long, default_value = "60")] + interval: u64, + }, InitSchema, } @@ -52,7 +58,12 @@ async fn main() -> Result<()> { } } -async fn cmd_sync(db_path: &PathBuf, dsn: &str, enable_embed: bool, since: Option) -> Result<()> { +async fn cmd_sync( + db_path: &PathBuf, + dsn: &str, + enable_embed: bool, + since: Option, +) -> Result<()> { info!("opening opencode.db at {:?}", db_path); let sqlite = Connection::open(db_path)?; sqlite.busy_timeout(Duration::from_secs(5))?; @@ -91,9 +102,12 @@ async fn cmd_sync(db_path: &PathBuf, dsn: &str, enable_embed: bool, since: Optio let mut chat_session = normalize_session(sess, &chat_msgs, None); if let Some(ref emb) = embedder { - let text = format!("{} {} {}", sess.title, + let text = format!( + "{} {} {}", + sess.title, sess.agent.as_deref().unwrap_or(""), - sess.model.as_deref().unwrap_or("")); + sess.model.as_deref().unwrap_or("") + ); if let Ok(v) = emb.embed(&text).await { chat_session.embedding = Some(v); } @@ -156,16 +170,25 @@ async fn cmd_init_schema(dsn: &str) -> Result<()> { #[derive(Debug, Clone, Serialize, Deserialize)] struct OpenCodeSession { - id: String, project_id: String, parent_id: Option, - slug: String, directory: String, title: String, - agent: Option, model: Option, - time_created: i64, time_updated: i64, - tokens_input: i64, tokens_output: i64, + id: String, + project_id: String, + parent_id: Option, + slug: String, + directory: String, + title: String, + agent: Option, + model: Option, + time_created: i64, + time_updated: i64, + tokens_input: i64, + tokens_output: i64, } #[derive(Debug, Clone, Serialize, Deserialize)] struct OpenCodeMessage { - id: String, session_id: String, time_created: i64, + id: String, + session_id: String, + time_created: i64, #[serde(rename = "role")] data_role: String, #[serde(default)] @@ -194,15 +217,22 @@ fn load_sessions(conn: &Connection) -> Result> { let mut stmt = conn.prepare( "SELECT id, project_id, parent_id, slug, directory, title, \ agent, model, time_created, time_updated, tokens_input, tokens_output \ - FROM session ORDER BY time_created" + FROM session ORDER BY time_created", )?; let rows = stmt.query_map([], |row| { Ok(OpenCodeSession { - id: row.get(0)?, project_id: row.get(1)?, parent_id: row.get(2)?, - slug: row.get(3)?, directory: row.get(4)?, title: row.get(5)?, - agent: row.get(6)?, model: row.get(7)?, - time_created: row.get(8)?, time_updated: row.get(9)?, - tokens_input: row.get(10)?, tokens_output: row.get(11)?, + id: row.get(0)?, + project_id: row.get(1)?, + parent_id: row.get(2)?, + slug: row.get(3)?, + directory: row.get(4)?, + title: row.get(5)?, + agent: row.get(6)?, + model: row.get(7)?, + time_created: row.get(8)?, + time_updated: row.get(9)?, + tokens_input: row.get(10)?, + tokens_output: row.get(11)?, }) })?; rows.collect::, _>>().map_err(|e| e.into()) @@ -212,56 +242,84 @@ fn sessions_since(conn: &Connection, since_ms: i64) -> Result ?1 ORDER BY time_created" + FROM session WHERE time_updated > ?1 ORDER BY time_created", )?; let rows = stmt.query_map([since_ms], |row| { Ok(OpenCodeSession { - id: row.get(0)?, project_id: row.get(1)?, parent_id: row.get(2)?, - slug: row.get(3)?, directory: row.get(4)?, title: row.get(5)?, - agent: row.get(6)?, model: row.get(7)?, - time_created: row.get(8)?, time_updated: row.get(9)?, - tokens_input: row.get(10)?, tokens_output: row.get(11)?, + id: row.get(0)?, + project_id: row.get(1)?, + parent_id: row.get(2)?, + slug: row.get(3)?, + directory: row.get(4)?, + title: row.get(5)?, + agent: row.get(6)?, + model: row.get(7)?, + time_created: row.get(8)?, + time_updated: row.get(9)?, + tokens_input: row.get(10)?, + tokens_output: row.get(11)?, }) })?; rows.collect::, _>>().map_err(|e| e.into()) } fn max_session_updated(conn: &Connection) -> Result> { - conn.query_row("SELECT MAX(time_updated) FROM session", [], |row| row.get(0)) - .optional() - .map_err(|e| e.into()) + conn.query_row("SELECT MAX(time_updated) FROM session", [], |row| { + row.get(0) + }) + .optional() + .map_err(|e| e.into()) } fn messages_for_session(conn: &Connection, session_id: &str) -> Result> { let mut stmt = conn.prepare( "SELECT id, session_id, time_created, data FROM message \ - WHERE session_id = ?1 ORDER BY time_created, id" + WHERE session_id = ?1 ORDER BY time_created, id", )?; let rows = stmt.query_map([session_id], |row| { let data_str: String = row.get(3)?; - let data: serde_json::Value = serde_json::from_str(&data_str).unwrap_or(serde_json::Value::Null); - let role = data.get("role").and_then(|v| v.as_str()).unwrap_or("unknown").to_string(); - Ok(OpenCodeMessage { id: row.get(0)?, session_id: row.get(1)?, time_created: row.get(2)?, data_role: role, data }) + let data: serde_json::Value = + serde_json::from_str(&data_str).unwrap_or(serde_json::Value::Null); + let role = data + .get("role") + .and_then(|v| v.as_str()) + .unwrap_or("unknown") + .to_string(); + Ok(OpenCodeMessage { + id: row.get(0)?, + session_id: row.get(1)?, + time_created: row.get(2)?, + data_role: role, + data, + }) })?; rows.collect::, _>>().map_err(|e| e.into()) } fn parts_for_message(conn: &Connection, message_id: &str) -> Result> { - let mut stmt = conn.prepare( - "SELECT data FROM part WHERE message_id = ?1 ORDER BY time_created, id" - )?; + let mut stmt = + conn.prepare("SELECT data FROM part WHERE message_id = ?1 ORDER BY time_created, id")?; let rows = stmt.query_map([message_id], |row| { let data_str: String = row.get(0)?; let part: OpenCodePart = serde_json::from_str(&data_str).unwrap_or(OpenCodePart { - part_type: "unknown".into(), text: None, tool: None, call_id: None, - input: None, output: None, is_error: None, + part_type: "unknown".into(), + text: None, + tool: None, + call_id: None, + input: None, + output: None, + is_error: None, }); Ok(part) })?; rows.collect::, _>>().map_err(|e| e.into()) } -fn normalize_session(sess: &OpenCodeSession, msgs: &[ChatMessage], _compaction: Option) -> ChatSession { +fn normalize_session( + sess: &OpenCodeSession, + msgs: &[ChatMessage], + _compaction: Option, +) -> ChatSession { ChatSession { session_id: sess.id.clone(), workspace_fingerprint: Some(workspace_fingerprint(&sess.directory)), @@ -285,7 +343,11 @@ fn normalize_session(sess: &OpenCodeSession, msgs: &[ChatMessage], _compaction: } } -fn normalize_message(msg: &OpenCodeMessage, parts: &[OpenCodePart], index: i32) -> Result { +fn normalize_message( + msg: &OpenCodeMessage, + parts: &[OpenCodePart], + index: i32, +) -> Result { let mut blocks = Vec::new(); let mut text_parts = Vec::new(); let mut tool_calls = Vec::new(); @@ -295,23 +357,55 @@ fn normalize_message(msg: &OpenCodeMessage, parts: &[OpenCodePart], index: i32) "text" => { if let Some(ref t) = part.text { text_parts.push(t.clone()); - blocks.push(MessageBlock { block_type: "text".into(), text: Some(t.clone()), tool_name: None, tool_input: None, tool_output: None, is_error: None }); + blocks.push(MessageBlock { + block_type: "text".into(), + text: Some(t.clone()), + tool_name: None, + tool_input: None, + tool_output: None, + is_error: None, + }); } } "reasoning" => { if let Some(ref t) = part.text { text_parts.push(format!("[reasoning] {}", t)); - blocks.push(MessageBlock { block_type: "reasoning".into(), text: Some(t.clone()), tool_name: None, tool_input: None, tool_output: None, is_error: None }); + blocks.push(MessageBlock { + block_type: "reasoning".into(), + text: Some(t.clone()), + tool_name: None, + tool_input: None, + tool_output: None, + is_error: None, + }); } } "tool" => { let call_id = part.call_id.clone().unwrap_or_default(); let tool_name = part.tool.clone().unwrap_or_default(); - blocks.push(MessageBlock { block_type: "tool_use".into(), text: None, tool_name: Some(tool_name.clone()), tool_input: part.input.clone(), tool_output: part.output.clone(), is_error: part.is_error }); - tool_calls.push(ToolCall { call_id: call_id.clone(), tool_name, input: part.input.clone().unwrap_or(serde_json::json!({})) }); + blocks.push(MessageBlock { + block_type: "tool_use".into(), + text: None, + tool_name: Some(tool_name.clone()), + tool_input: part.input.clone(), + tool_output: part.output.clone(), + is_error: part.is_error, + }); + tool_calls.push(ToolCall { + call_id: call_id.clone(), + tool_name, + input: part.input.clone().unwrap_or(serde_json::json!({})), + }); } "tool-result" => { - blocks.push(MessageBlock { block_type: "tool_result".into(), text: part.text.clone(), tool_name: part.tool.clone(), tool_input: None, tool_output: part.output.clone(), is_error: part.is_error }); + blocks.push(MessageBlock { + block_type: "tool_result".into(), + text: part.text.clone(), + tool_name: part.tool.clone(), + tool_input: None, + tool_output: part.output.clone(), + is_error: part.is_error, + }); } _ => {} } @@ -323,8 +417,10 @@ fn normalize_message(msg: &OpenCodeMessage, parts: &[OpenCodePart], index: i32) role: msg.data_role.clone(), blocks, text_content: text_parts.join("\n"), - token_input: 0, token_output: 0, - token_cache_creation: 0, token_cache_read: 0, + token_input: 0, + token_output: 0, + token_cache_creation: 0, + token_cache_read: 0, tool_calls, embedding: None, receipt_hash: None, @@ -354,19 +450,34 @@ struct Embedder { impl Embedder { fn new() -> Self { let base = std::env::var("OLLAMA_HOST").unwrap_or_else(|_| "http://localhost:11434".into()); - let model = std::env::var("OLLAMA_EMBED_MODEL").unwrap_or_else(|_| "nomic-embed-text".into()); - Self { client: reqwest::Client::new(), url: format!("{}/api/embeddings", base.trim_end_matches('/')), model } + let model = + std::env::var("OLLAMA_EMBED_MODEL").unwrap_or_else(|_| "nomic-embed-text".into()); + Self { + client: reqwest::Client::new(), + url: format!("{}/api/embeddings", base.trim_end_matches('/')), + model, + } } async fn embed(&self, text: &str) -> Result> { - let resp = self.client.post(&self.url) + let resp = self + .client + .post(&self.url) .json(&serde_json::json!({"model": self.model, "prompt": text})) - .send().await.context("embed POST")?; + .send() + .await + .context("embed POST")?; if !resp.status().is_success() { anyhow::bail!("embed HTTP {}", resp.status()); } let json: serde_json::Value = resp.json().await.context("embed JSON")?; - let arr = json.get("embedding").and_then(|v| v.as_array()).context("missing embedding")?; - Ok(arr.iter().map(|v| v.as_f64().unwrap_or(0.0) as f32).collect()) + let arr = json + .get("embedding") + .and_then(|v| v.as_array()) + .context("missing embedding")?; + Ok(arr + .iter() + .map(|v| v.as_f64().unwrap_or(0.0) as f32) + .collect()) } } diff --git a/4-Infrastructure/infra/ene-rds/systemd/ene-api-wrapper.sh b/4-Infrastructure/infra/ene-rds/systemd/ene-api-wrapper.sh new file mode 100755 index 00000000..bdb76363 --- /dev/null +++ b/4-Infrastructure/infra/ene-rds/systemd/ene-api-wrapper.sh @@ -0,0 +1,34 @@ +#!/usr/bin/env bash +set -euo pipefail + +ROOT="/home/allaun/Research Stack/4-Infrastructure/infra/ene-rds" +BIN="${ENE_API_BIN:-$ROOT/target/release/ene-api}" + +if [[ ! -x "$BIN" ]]; then + BIN="$ROOT/target/debug/ene-api" +fi + +if [[ ! -x "$BIN" ]]; then + echo "ene-api binary not found; run: cargo build --release -p ene-api" >&2 + exit 127 +fi + +export AWS_REGION="${AWS_REGION:-us-east-1}" +export RDS_HOST="${RDS_HOST:-database-1-instance-1.cghu8yqogqwo.us-east-1.rds.amazonaws.com}" +export RDS_PORT="${RDS_PORT:-5432}" +export RDS_USER="${RDS_USER:-postgres}" +export RDS_DB="${RDS_DB:-postgres}" +export ENE_API_BIND="${ENE_API_BIND:-0.0.0.0:3000}" + +if [[ -z "${RDS_IAM_TOKEN:-}" && -z "${RDS_PASSWORD:-}" ]]; then + export RDS_IAM_TOKEN="$( + aws rds generate-db-auth-token \ + --hostname "$RDS_HOST" \ + --port "$RDS_PORT" \ + --username "$RDS_USER" \ + --region "$AWS_REGION" + )" +fi + +cd "$ROOT" +exec "$BIN" diff --git a/4-Infrastructure/infra/ene-rds/systemd/ene-api.service b/4-Infrastructure/infra/ene-rds/systemd/ene-api.service new file mode 100644 index 00000000..a5af482b --- /dev/null +++ b/4-Infrastructure/infra/ene-rds/systemd/ene-api.service @@ -0,0 +1,15 @@ +[Unit] +Description=ENE API Server +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +ExecStart=/usr/bin/env bash "/home/allaun/Research Stack/4-Infrastructure/infra/ene-rds/systemd/ene-api-wrapper.sh" +Restart=on-failure +RestartSec=30 +Environment=AWS_REGION=us-east-1 +Environment=ENE_API_BIND=0.0.0.0:3000 + +[Install] +WantedBy=default.target diff --git a/4-Infrastructure/infra/ene_contextstream_mcp.py b/4-Infrastructure/infra/ene_contextstream_mcp.py index 8f85dde2..b684cb48 100755 --- a/4-Infrastructure/infra/ene_contextstream_mcp.py +++ b/4-Infrastructure/infra/ene_contextstream_mcp.py @@ -551,9 +551,23 @@ def handle(req: dict[str, Any]) -> dict[str, Any] | None: def main() -> int: + if len(sys.argv) > 1: + command = sys.argv[1] + if command in {"--status", "status"}: + print(json.dumps(tool_status({}), indent=2, sort_keys=True)) + return 0 + if command in {"--context", "context"}: + query = " ".join(sys.argv[2:]) + print(json.dumps(tool_context({"user_message": query}), indent=2, sort_keys=True)) + return 0 + print(f"unknown command: {command}", file=sys.stderr) + return 2 + + handled = 0 for line in sys.stdin: if not line.strip(): continue + handled += 1 try: response = handle(json.loads(line)) except json.JSONDecodeError as exc: @@ -565,6 +579,8 @@ def main() -> int: if response is not None: sys.stdout.write(json.dumps(response, separators=(",", ":")) + "\n") sys.stdout.flush() + if handled == 0: + print(json.dumps(tool_status({}), indent=2, sort_keys=True)) return 0 diff --git a/AGENTS.md b/AGENTS.md index 8190a7e2..491a3a39 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -217,7 +217,7 @@ STOP → Call ene-contextstream first: If the client cannot call the ENE MCP server, use the local fallback: ```bash -python3 4-Infrastructure/infra/ene_contextstream_mcp.py +python3 4-Infrastructure/infra/ene_contextstream_mcp.py --status ``` ## Remote Proof Agent