From 0c502bc58ad3f11621c24d7af62eddd4aaf4d937 Mon Sep 17 00:00:00 2001 From: Brandon Schneider Date: Wed, 20 May 2026 18:53:55 -0500 Subject: [PATCH] fix(rds): restore ENE storage observation probe --- .../infra/ene-rds/crates/ene-api/src/main.rs | 25 ++- .../infra/ene-rds/crates/ene-node/src/lib.rs | 167 ++++++++++---- .../infra/ene-rds/crates/ene-node/src/main.rs | 31 ++- .../ene-rds/crates/ene-rds-chat/src/lib.rs | 129 +++++++---- .../ene-rds/crates/ene-rds-core/src/lib.rs | 33 ++- .../ene-rds/crates/ene-rds-core/src/types.rs | 12 +- .../crates/ene-rds-ephemeral/src/lib.rs | 94 +++++--- .../ene-rds/crates/ene-rds-wiki/src/lib.rs | 63 ++++-- .../ene-rds/crates/ene-storage/src/lib.rs | 71 +++++- .../ene-rds/crates/ene-storage/src/main.rs | 97 +++++--- .../infra/ene-rds/crates/ene-sync/src/main.rs | 211 +++++++++++++----- 11 files changed, 690 insertions(+), 243 deletions(-) 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..4a891b9f 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 @@ -29,7 +29,9 @@ struct SearchQuery { semantic: bool, } -fn default_limit() -> i64 { 10 } +fn default_limit() -> i64 { + 10 +} #[derive(Serialize)] struct HealthResponse { @@ -58,7 +60,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)) @@ -88,7 +94,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 +137,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 +165,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-node/src/lib.rs b/4-Infrastructure/infra/ene-rds/crates/ene-node/src/lib.rs index 832da859..79ac386b 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-node/src/lib.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-node/src/lib.rs @@ -57,7 +57,12 @@ impl GossipMessage { pub fn new(sender: &str, msg_type: &str, payload: serde_json::Value) -> Self { let id = format!( "gossip_{}", - &sha256_hex(&format!("{}:{}:{}", sender, msg_type, Utc::now().timestamp_millis()))[..16] + &sha256_hex(&format!( + "{}:{}:{}", + sender, + msg_type, + Utc::now().timestamp_millis() + ))[..16] ); Self { message_id: id, @@ -88,8 +93,8 @@ impl GossipMessage { use hmac::{Hmac, Mac}; use sha2::Sha256; type HmacSha256 = Hmac; - let mut mac = HmacSha256::new_from_slice(secret.as_bytes()) - .expect("HMAC can take key of any size"); + let mut mac = + HmacSha256::new_from_slice(secret.as_bytes()).expect("HMAC can take key of any size"); mac.update(&self.canonical_bytes()); let result = mac.finalize(); self.signature = Some(hex::encode(result.into_bytes())); @@ -97,7 +102,9 @@ impl GossipMessage { /// Verify the HMAC signature against the cluster secret. pub fn verify(&self, secret: &str) -> bool { - let Some(ref sig) = self.signature else { return false }; + let Some(ref sig) = self.signature else { + return false; + }; use hmac::{Hmac, Mac}; use sha2::Sha256; type HmacSha256 = Hmac; @@ -231,10 +238,16 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no replication_version, capabilities, health_score_q16, is_active) \ VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", rusqlite::params![ - &peer.node_id, &peer.public_key, &peer.ip_address, &peer.port, - &peer.first_seen, &peer.last_seen, &peer.replication_version, + &peer.node_id, + &peer.public_key, + &peer.ip_address, + &peer.port, + &peer.first_seen, + &peer.last_seen, + &peer.replication_version, serde_json::to_string(&peer.capabilities)?, - &peer.health_score_q16, &peer.is_active as &dyn rusqlite::ToSql, + &peer.health_score_q16, + &peer.is_active as &dyn rusqlite::ToSql, ], )?; Ok(()) @@ -244,7 +257,7 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no let mut stmt = self.conn.prepare( "SELECT node_id, public_key, ip_address, port, first_seen, last_seen, \ replication_version, capabilities, health_score_q16, is_active \ - FROM ene_peers WHERE is_active = 1" + FROM ene_peers WHERE is_active = 1", )?; let rows = stmt.query_map([], |row| { let caps: String = row.get(7)?; @@ -263,7 +276,9 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no }) })?; let mut out = Vec::new(); - for r in rows { out.push(r?); } + for r in rows { + out.push(r?); + } Ok(out) } @@ -273,8 +288,11 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no (message_id, sender_node, message_type, payload, timestamp) \ VALUES (?1, ?2, ?3, ?4, ?5)", rusqlite::params![ - &msg.message_id, &msg.sender_node, &msg.message_type, - &serde_json::to_string(&msg.payload)?, &msg.timestamp, + &msg.message_id, + &msg.sender_node, + &msg.message_type, + &serde_json::to_string(&msg.payload)?, + &msg.timestamp, ], )?; Ok(()) @@ -286,8 +304,11 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no (proposal_id, credential_id, proposer, timestamp, votes, resolved) \ VALUES (?1, ?2, ?3, ?4, ?5, ?6)", rusqlite::params![ - &prop.proposal_id, &prop.credential_id, &prop.proposer, - &prop.timestamp, &serde_json::to_string(&prop.votes)?, + &prop.proposal_id, + &prop.credential_id, + &prop.proposer, + &prop.timestamp, + &serde_json::to_string(&prop.votes)?, &prop.resolved as &dyn rusqlite::ToSql, ], )?; @@ -297,19 +318,21 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no pub fn load_proposal(&self, proposal_id: &str) -> Result> { let mut stmt = self.conn.prepare( "SELECT proposal_id, credential_id, proposer, timestamp, votes, resolved \ - FROM ene_proposals WHERE proposal_id = ?1" + FROM ene_proposals WHERE proposal_id = ?1", )?; - let row = stmt.query_row([proposal_id], |row| { - let votes_str: String = row.get(4)?; - Ok(ConsensusProposal { - proposal_id: row.get(0)?, - credential_id: row.get(1)?, - proposer: row.get(2)?, - timestamp: row.get(3)?, - votes: serde_json::from_str(&votes_str).unwrap_or_default(), - resolved: row.get::<_, i32>(5)? != 0, + let row = stmt + .query_row([proposal_id], |row| { + let votes_str: String = row.get(4)?; + Ok(ConsensusProposal { + proposal_id: row.get(0)?, + credential_id: row.get(1)?, + proposer: row.get(2)?, + timestamp: row.get(3)?, + votes: serde_json::from_str(&votes_str).unwrap_or_default(), + resolved: row.get::<_, i32>(5)? != 0, + }) }) - }).optional()?; + .optional()?; Ok(row) } @@ -329,7 +352,9 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no }) })?; let mut out = Vec::new(); - for r in rows { out.push(r?); } + for r in rows { + out.push(r?); + } Ok(out) } @@ -353,10 +378,15 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no usage_count, last_rotated, health_score_q16, is_active) \ VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)", rusqlite::params![ - &frag.credential_id, &frag.provider, &frag.fragment, - &frag.access_level, &serde_json::to_string(&frag.node_assignments)?, - &frag.usage_count, &frag.last_rotated, - &frag.health_score_q16, &frag.is_active as &dyn rusqlite::ToSql, + &frag.credential_id, + &frag.provider, + &frag.fragment, + &frag.access_level, + &serde_json::to_string(&frag.node_assignments)?, + &frag.usage_count, + &frag.last_rotated, + &frag.health_score_q16, + &frag.is_active as &dyn rusqlite::ToSql, ], )?; Ok(()) @@ -388,7 +418,10 @@ impl EneNode { let loaded_peers = db.load_peers().unwrap_or_default(); let mut identity = NodeIdentity::default(); identity.node_id = node_id.unwrap_or_else(|| { - format!("ene_{}", &sha256_hex(&Utc::now().timestamp_millis().to_string())[..16]) + format!( + "ene_{}", + &sha256_hex(&Utc::now().timestamp_millis().to_string())[..16] + ) }); identity.public_key = sha256_hex(&identity.node_id)[..32].to_string(); db.save_peer(&identity)?; @@ -439,8 +472,15 @@ impl EneNode { drop(peers_guard); self.with_db(|db| db.save_gossip(&msg))?; - self.seen_message_ids.write().await.insert(msg.message_id.clone()); - info!("gossip {} -> {} peers", msg.message_type, self.peers.read().await.len()); + self.seen_message_ids + .write() + .await + .insert(msg.message_id.clone()); + info!( + "gossip {} -> {} peers", + msg.message_type, + self.peers.read().await.len() + ); Ok(()) } @@ -455,7 +495,10 @@ impl EneNode { if self.seen_message_ids.read().await.contains(&msg.message_id) { return Ok(()); } - self.seen_message_ids.write().await.insert(msg.message_id.clone()); + self.seen_message_ids + .write() + .await + .insert(msg.message_id.clone()); self.with_db(|db| db.save_gossip(&msg))?; match msg.message_type.as_str() { @@ -471,9 +514,15 @@ impl EneNode { async fn handle_discovery(&self, msg: &GossipMessage, from: SocketAddr) -> Result<()> { let node_id = msg.payload.get("node_id").and_then(|v| v.as_str()); - let caps = msg.payload.get("capabilities") + let caps = msg + .payload + .get("capabilities") .and_then(|v| v.as_array()) - .map(|arr| arr.iter().filter_map(|v| v.as_str().map(|s| s.to_string())).collect()) + .map(|arr| { + arr.iter() + .filter_map(|v| v.as_str().map(|s| s.to_string())) + .collect() + }) .unwrap_or_else(|| vec!["storage".into(), "compute".into()]); if let Some(nid) = node_id { @@ -517,9 +566,18 @@ impl EneNode { if let (Some(id), Some(b64)) = (cred_id, fragment_b64) { let frag = CredentialFragment { credential_id: id.into(), - provider: msg.payload.get("provider").and_then(|v| v.as_str()).unwrap_or("unknown").into(), + provider: msg + .payload + .get("provider") + .and_then(|v| v.as_str()) + .unwrap_or("unknown") + .into(), fragment: base64_decode(b64), - access_level: msg.payload.get("access_level").and_then(|v| v.as_i64()).unwrap_or(0) as i32, + access_level: msg + .payload + .get("access_level") + .and_then(|v| v.as_i64()) + .unwrap_or(0) as i32, node_assignments: vec![msg.sender_node.clone()], usage_count: 0, last_rotated: Utc::now().timestamp_millis(), @@ -569,7 +627,8 @@ impl EneNode { pub async fn vote(&self, proposal_id: &str, approve: bool) -> Result { let total = self.peers.read().await.len() + 1; - let mut prop = self.with_db(|db| db.load_proposal(proposal_id))? + let mut prop = self + .with_db(|db| db.load_proposal(proposal_id))? .context("proposal not found")?; prop.votes.insert(self.identity.node_id.clone(), approve); self.with_db(|db| db.save_proposal(&prop))?; @@ -579,7 +638,10 @@ impl EneNode { if approve_count >= threshold && !prop.resolved { prop.resolved = true; self.with_db(|db| db.save_proposal(&prop))?; - info!("consensus reached on {} ({} of {})", proposal_id, approve_count, total); + info!( + "consensus reached on {} ({} of {})", + proposal_id, approve_count, total + ); return Ok(true); } Ok(false) @@ -588,7 +650,11 @@ impl EneNode { pub async fn propose_rotation(&self, credential_id: &str) -> Result { let proposal_id = format!( "prop_{}", - &sha256_hex(&format!("{}:{}", credential_id, Utc::now().timestamp_millis()))[..12] + &sha256_hex(&format!( + "{}:{}", + credential_id, + Utc::now().timestamp_millis() + ))[..12] ); let payload = serde_json::json!({ "proposal_id": &proposal_id, @@ -596,7 +662,11 @@ impl EneNode { "proposer": &self.identity.node_id, "timestamp": Utc::now().timestamp_millis(), }); - let msg = GossipMessage::new(&self.identity.node_id, "credential_rotation_proposal", payload); + let msg = GossipMessage::new( + &self.identity.node_id, + "credential_rotation_proposal", + payload, + ); self.gossip(msg).await?; Ok(proposal_id) } @@ -627,7 +697,10 @@ impl EneNode { pub async fn get_mesh_health(&self) -> serde_json::Value { let peers = self.peers.read().await; - let healthy = peers.values().filter(|p| p.health_score_q16 > 32768).count(); + let healthy = peers + .values() + .filter(|p| p.health_score_q16 > 32768) + .count(); let mesh_size = peers.len() + 1; serde_json::json!({ "mesh_size": mesh_size, @@ -639,7 +712,9 @@ impl EneNode { fn base64_decode(s: &str) -> Vec { use base64::Engine; - base64::engine::general_purpose::STANDARD.decode(s.as_bytes()).unwrap_or_default() + base64::engine::general_purpose::STANDARD + .decode(s.as_bytes()) + .unwrap_or_default() } #[cfg(test)] @@ -648,7 +723,8 @@ mod tests { #[test] fn gossip_sign_and_verify_roundtrip() { - let mut msg = GossipMessage::new("ene_alpha", "heartbeat", serde_json::json!({"health": 1})); + let mut msg = + GossipMessage::new("ene_alpha", "heartbeat", serde_json::json!({"health": 1})); assert!(msg.signature.is_none()); msg.sign("secret-key"); assert!(msg.signature.is_some()); @@ -670,7 +746,8 @@ mod tests { #[test] fn gossip_tamper_payload_invalidates_signature() { - let mut msg = GossipMessage::new("ene_alpha", "heartbeat", serde_json::json!({"health": 1})); + let mut msg = + GossipMessage::new("ene_alpha", "heartbeat", serde_json::json!({"health": 1})); msg.sign("secret-key"); assert!(msg.verify("secret-key")); // Tamper with payload after signing diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-node/src/main.rs b/4-Infrastructure/infra/ene-rds/crates/ene-node/src/main.rs index ab41d526..ac0b3e9c 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-node/src/main.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-node/src/main.rs @@ -35,11 +35,23 @@ async fn main() -> Result<()> { let _ = std::fs::create_dir_all(parent); } - let cluster_secret = cli.cluster_secret.or_else(|| std::env::var("ENE_CLUSTER_SECRET").ok()); - let node = EneNode::new(cli.node_id, &cli.db, cli.bind, cli.seed.clone(), cluster_secret).await?; + let cluster_secret = cli + .cluster_secret + .or_else(|| std::env::var("ENE_CLUSTER_SECRET").ok()); + let node = EneNode::new( + cli.node_id, + &cli.db, + cli.bind, + cli.seed.clone(), + cluster_secret, + ) + .await?; let node = std::sync::Arc::new(node); - info!("ENE node {} starting on {}", node.identity.node_id, cli.bind); + info!( + "ENE node {} starting on {}", + node.identity.node_id, cli.bind + ); // ── UDP listener task ──────────────────────────────────────────────── let listen_node = node.clone(); @@ -82,7 +94,8 @@ async fn main() -> Result<()> { // ── Heartbeat loop ──────────────────────────────────────────────────── let beat_node = node.clone(); let heartbeat = tokio::spawn(async move { - let mut interval = tokio::time::interval(std::time::Duration::from_secs(cli.heartbeat_interval)); + let mut interval = + tokio::time::interval(std::time::Duration::from_secs(cli.heartbeat_interval)); interval.tick().await; // skip immediate first tick loop { interval.tick().await; @@ -100,8 +113,14 @@ async fn main() -> Result<()> { interval.tick().await; let status = report_node.get_status().await; let health = report_node.get_mesh_health().await; - info!("status: {}", serde_json::to_string(&status).unwrap_or_default()); - info!("health: {}", serde_json::to_string(&health).unwrap_or_default()); + info!( + "status: {}", + serde_json::to_string(&status).unwrap_or_default() + ); + info!( + "health: {}", + serde_json::to_string(&health).unwrap_or_default() + ); } }); 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 12f05740..0d34d382 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 @@ -63,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(()) } @@ -159,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) @@ -181,19 +190,30 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US ) .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, \ 1 - (embedding <=> $1::vector) AS similarity \ @@ -203,17 +223,24 @@ 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, \ token_input_total, token_output_total, created_at_ms, updated_at_ms \ @@ -222,21 +249,28 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US ) .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, \ token_input_total, token_output_total, created_at_ms, updated_at_ms, meta \ @@ -246,7 +280,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 \ @@ -255,16 +291,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/src/lib.rs b/4-Infrastructure/infra/ene-rds/crates/ene-rds-core/src/lib.rs index c8408c09..0f3fba58 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 @@ -34,8 +34,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") @@ -49,7 +50,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) } @@ -119,7 +124,13 @@ 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(",") + ) } #[cfg(test)] @@ -128,7 +139,10 @@ mod tests { #[test] fn dsn_from_env_uses_rds_dsn_when_set() { - std::env::set_var("RDS_DSN", "host=test port=5432 dbname=test user=test password=test sslmode=require"); + std::env::set_var( + "RDS_DSN", + "host=test port=5432 dbname=test user=test password=test sslmode=require", + ); // Clear other vars to ensure RDS_DSN takes precedence for key in &["RDS_HOST", "RDS_PORT", "RDS_USER", "RDS_PASSWORD", "RDS_DB"] { let _ = std::env::remove_var(key); @@ -162,7 +176,14 @@ mod tests { #[test] fn dsn_from_env_uses_defaults() { - for key in &["RDS_DSN", "RDS_HOST", "RDS_PORT", "RDS_USER", "RDS_PASSWORD", "RDS_DB"] { + for key in &[ + "RDS_DSN", + "RDS_HOST", + "RDS_PORT", + "RDS_USER", + "RDS_PASSWORD", + "RDS_DB", + ] { let _ = std::env::remove_var(key); } let dsn = RdsClient::dsn_from_env(); 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 269e98de..cf438718 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, &metric.recorded_at_ms, + &metric.metric_id, + &metric.node_id, + &metric.metric_name, + &metric.metric_value_raw, + &metric.metric_scale, + &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-storage/src/lib.rs b/4-Infrastructure/infra/ene-rds/crates/ene-storage/src/lib.rs index 7fa402ec..5dfa2cd6 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-storage/src/lib.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-storage/src/lib.rs @@ -3,6 +3,7 @@ use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; use std::collections::HashMap; use std::path::PathBuf; +use std::time::Duration; // ── constants ────────────────────────────────────────────────────────────── @@ -125,6 +126,49 @@ impl Observation { } } +pub async fn observe(creds: &HashMap) -> Observation { + let mut obs = Observation::new(); + let endpoint = creds + .get("AWS_ENDPOINT_URL") + .cloned() + .unwrap_or_else(|| GARAGE_ENDPOINT.to_string()); + + match tokio::time::timeout(Duration::from_secs(2), reqwest::get(&endpoint)).await { + Ok(Ok(response)) => { + obs.garage.up = true; + obs.garage.buckets.push(GARAGE_BUCKET.to_string()); + if !response.status().is_success() { + obs.errors.push(format!( + "garage endpoint responded with non-success status {}", + response.status() + )); + } + } + Ok(Err(err)) => { + obs.errors + .push(format!("garage endpoint probe failed: {err}")); + } + Err(_) => { + obs.errors + .push("garage endpoint probe timed out after 2s".to_string()); + } + } + + let backup_log = dirs::cache_dir() + .unwrap_or_else(|| PathBuf::from("/tmp")) + .join("garage-post-commit.log"); + if let Ok(text) = std::fs::read_to_string(&backup_log) { + obs.backup_log.last_ts = std::fs::metadata(&backup_log) + .and_then(|metadata| metadata.modified()) + .ok() + .map(|modified| chrono::DateTime::::from(modified).to_rfc3339()); + obs.backup_log.last_ok = + text.contains("completed") || text.contains("success") || text.contains("succeeded"); + } + + obs +} + // ── Decision ─────────────────────────────────────────────────────────────── #[derive(Debug, Default, Clone, Serialize, Deserialize)] @@ -145,19 +189,23 @@ pub fn decide(obs: &Observation) -> Decision { if !obs.garage.up { d.alerts.push("ALERT: Garage S3 is unreachable".into()); d.trigger_garage_restart = true; - d.rationale.push("garage_up=false → trigger_garage_restart".into()); + d.rationale + .push("garage_up=false → trigger_garage_restart".into()); } if obs.restic.snapshot_count == 0 && obs.garage.up { - d.alerts.push("ALERT: restic repo has zero snapshots — initial backup needed".into()); + d.alerts + .push("ALERT: restic repo has zero snapshots — initial backup needed".into()); d.trigger_snap = true; d.rationale.push("snapshot_count=0 → trigger_snap".into()); } if !obs.backup_log.last_ok && obs.garage.up { - d.alerts.push("WARN: No successful restic snapshot found in backup log".into()); + d.alerts + .push("WARN: No successful restic snapshot found in backup log".into()); d.trigger_snap = true; - d.rationale.push("backup_log_last_ok=false → trigger_snap".into()); + d.rationale + .push("backup_log_last_ok=false → trigger_snap".into()); } if obs.restic.dedup_ratio_q16 > 0 @@ -180,14 +228,18 @@ pub fn decide(obs: &Observation) -> Decision { } if obs.cold_copy_needed { - d.alerts.push("WARN: Newest restic snapshot is >26 h old — cold copy to gdrive appears stale".into()); + d.alerts.push( + "WARN: Newest restic snapshot is >26 h old — cold copy to gdrive appears stale".into(), + ); d.trigger_cold_copy = true; - d.rationale.push("cold_copy_needed=true → trigger_cold_copy".into()); + d.rationale + .push("cold_copy_needed=true → trigger_cold_copy".into()); } if obs.garage.up { d.trigger_offload = true; - d.rationale.push("garage_up=true → trigger_offload (idempotent)".into()); + d.rationale + .push("garage_up=true → trigger_offload (idempotent)".into()); } d @@ -261,7 +313,10 @@ pub fn resume_chain() -> (i64, String) { let mut last_hash = String::new(); for line in text.lines() { if let Ok(entry) = serde_json::from_str::(line) { - last_tick = entry.get("tick").and_then(|v| v.as_i64()).unwrap_or(last_tick); + last_tick = entry + .get("tick") + .and_then(|v| v.as_i64()) + .unwrap_or(last_tick); last_hash = entry .get("receipt_hash") .and_then(|v| v.as_str()) diff --git a/4-Infrastructure/infra/ene-rds/crates/ene-storage/src/main.rs b/4-Infrastructure/infra/ene-rds/crates/ene-storage/src/main.rs index a9c77cf9..c29d5f60 100644 --- a/4-Infrastructure/infra/ene-rds/crates/ene-storage/src/main.rs +++ b/4-Infrastructure/infra/ene-rds/crates/ene-storage/src/main.rs @@ -1,9 +1,8 @@ use anyhow::{Context, Result}; -use chrono::Utc; use clap::Parser; use ene_storage::{ build_receipt, decide, load_garage_env, observe, resume_chain, ActionResult, Decision, - Observation, GARAGE_BUCKET, GARAGE_ENDPOINT, RECEIPT_PREFIX, + GARAGE_BUCKET, GARAGE_ENDPOINT, RECEIPT_PREFIX, }; use serde_json::json; use std::collections::HashMap; @@ -25,13 +24,10 @@ async fn run_cmd( for (k, v) in env_extra { cmd.env(k, v); } - let output = tokio::time::timeout( - std::time::Duration::from_secs(timeout_secs), - cmd.output(), - ) - .await - .context("command timed out")? - .context("command failed to run")?; + let output = tokio::time::timeout(std::time::Duration::from_secs(timeout_secs), cmd.output()) + .await + .context("command timed out")? + .context("command failed to run")?; let rc = output.status.code().unwrap_or(-1); let stdout = String::from_utf8_lossy(&output.stdout).into_owned(); let stderr = String::from_utf8_lossy(&output.stderr).into_owned(); @@ -51,8 +47,18 @@ async fn act_one( match run_cmd(args, env_extra, timeout_secs).await { Ok((0, stdout, _)) => { ar.actions_succeeded.push(label.to_string()); - let tail = stdout.chars().rev().take(500).collect::>().into_iter().rev().collect::(); - ar.details.insert(label.to_string(), json!({"rc": 0, "stdout_tail": tail.trim() })); + let tail = stdout + .chars() + .rev() + .take(500) + .collect::>() + .into_iter() + .rev() + .collect::(); + ar.details.insert( + label.to_string(), + json!({"rc": 0, "stdout_tail": tail.trim() }), + ); } Ok((rc, stdout, stderr)) => { ar.actions_failed.push(label.to_string()); @@ -67,7 +73,10 @@ async fn act_one( } Err(e) => { ar.actions_failed.push(label.to_string()); - ar.details.insert(label.to_string(), json!({"rc": -1, "error": e.to_string() })); + ar.details.insert( + label.to_string(), + json!({"rc": -1, "error": e.to_string() }), + ); } } } @@ -81,7 +90,10 @@ async fn act( let mut ar = ActionResult::default(); if probe_only || dry_run { - ar.details.insert("mode".into(), json!(if probe_only { "probe_only" } else { "dry_run" })); + ar.details.insert( + "mode".into(), + json!(if probe_only { "probe_only" } else { "dry_run" }), + ); return ar; } @@ -94,16 +106,35 @@ async fn act( let consolidate_sh = storage_dir.join("garage/db-consolidate.sh"); if d.trigger_garage_restart { - act_one("garage_restart", &["systemctl", "--user", "restart", "garage.service"], &mut ar, creds, 30).await; + act_one( + "garage_restart", + &["systemctl", "--user", "restart", "garage.service"], + &mut ar, + creds, + 30, + ) + .await; if ar.actions_failed.contains(&"garage_restart".to_string()) { - act_one("garage_restart_system", &["sudo", "systemctl", "restart", "garage.service"], &mut ar, creds, 30).await; + act_one( + "garage_restart_system", + &["sudo", "systemctl", "restart", "garage.service"], + &mut ar, + creds, + 30, + ) + .await; } } if d.trigger_snap { act_one( "restic_snap", - &["bash", &backup_sh.to_string_lossy(), "snap", "agent-triggered"], + &[ + "bash", + &backup_sh.to_string_lossy(), + "snap", + "agent-triggered", + ], &mut ar, creds, 3600, @@ -177,7 +208,10 @@ fn emit_local(receipt: &serde_json::Value) -> Result<()> { Ok(()) } -async fn emit_s3(receipt: &serde_json::Value, creds: &HashMap) -> Result<(bool, String)> { +async fn emit_s3( + receipt: &serde_json::Value, + creds: &HashMap, +) -> Result<(bool, String)> { let ts = receipt["generated_at_utc"] .as_str() .unwrap_or("") @@ -196,7 +230,10 @@ async fn emit_s3(receipt: &serde_json::Value, creds: &HashMap) - let tmp = tempfile::NamedTempFile::with_suffix(".json")?; tokio::fs::write(tmp.path(), &payload).await?; - let endpoint = creds.get("AWS_ENDPOINT_URL").cloned().unwrap_or_else(|| GARAGE_ENDPOINT.into()); + let endpoint = creds + .get("AWS_ENDPOINT_URL") + .cloned() + .unwrap_or_else(|| GARAGE_ENDPOINT.into()); let (rc, _, _) = run_cmd( &[ "aws", @@ -251,7 +288,9 @@ async fn run_cycle( let _ = emit_local(&receipt); let (s3_ok, s3_key) = if !no_s3 && obs.garage.up { - emit_s3(&receipt, creds).await.unwrap_or((false, String::new())) + emit_s3(&receipt, creds) + .await + .unwrap_or((false, String::new())) } else { (false, String::new()) }; @@ -286,10 +325,7 @@ async fn run_cycle( warn!("! {}", alert); } - receipt["receipt_hash"] - .as_str() - .unwrap_or("") - .to_string() + receipt["receipt_hash"].as_str().unwrap_or("").to_string() } #[tokio::main] @@ -305,7 +341,10 @@ async fn main() -> Result<()> { tick += 1; if cli.loop_mode { - info!("loop mode, interval={}s, resuming tick={}", cli.interval, tick); + info!( + "loop mode, interval={}s, resuming tick={}", + cli.interval, tick + ); loop { match tokio::time::timeout( std::time::Duration::from_secs(cli.interval + 300), @@ -327,7 +366,15 @@ async fn main() -> Result<()> { tokio::time::sleep(std::time::Duration::from_secs(cli.interval)).await; } } else { - run_cycle(tick, &parent_hash, &creds, cli.probe_only, cli.dry_run, cli.no_s3).await; + run_cycle( + tick, + &parent_hash, + &creds, + cli.probe_only, + cli.dry_run, + cli.no_s3, + ) + .await; } Ok(()) 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()) } }