fix(rds): restore ENE storage observation probe

This commit is contained in:
Brandon Schneider 2026-05-20 18:53:55 -05:00
parent 7dd9abc837
commit 0c502bc58a
11 changed files with 690 additions and 243 deletions

View file

@ -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<Arc<Mutex<AppState>>>,
Query(params): Query<HashMap<String, String>>,
) -> Result<Json<serde_json::Value>, 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<HashMap<String, String>>,
) -> Result<Json<serde_json::Value>, 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<HashMap<String, String>>,
) -> Result<Json<serde_json::Value>, 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}))),

View file

@ -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<Sha256>;
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<Sha256>;
@ -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<Option<ConsensusProposal>> {
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<bool> {
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<String> {
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<u8> {
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

View file

@ -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()
);
}
});

View file

@ -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<u64> {
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<String>>(2),
"model": r.get::<_, Option<String>>(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<String>>(2),
"model": r.get::<_, Option<String>>(3),
"match_count": r.get::<_, i64>(4),
"rank": r.get::<_, f32>(5),
})
})
.collect())
}
pub async fn search_similar(&self, embedding: &[f32], limit: i64) -> Result<Vec<serde_json::Value>> {
pub async fn search_similar(
&self,
embedding: &[f32],
limit: i64,
) -> Result<Vec<serde_json::Value>> {
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<String>>(2),
"model": r.get::<_, Option<String>>(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<String>>(2),
"model": r.get::<_, Option<String>>(3),
"similarity": r.get::<_, f32>(4),
})
})
.collect())
}
pub async fn list_sessions(&self, limit: i64) -> Result<Vec<serde_json::Value>> {
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<String>>(2),
"model": r.get::<_, Option<String>>(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<String>>(2),
"model": r.get::<_, Option<String>>(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<Option<serde_json::Value>> {
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),

View file

@ -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<u64> {
pub async fn execute(
&self,
sql: &str,
params: &[&(dyn tokio_postgres::types::ToSql + Sync)],
) -> Result<u64> {
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::<Vec<_>>().join(","))
format!(
"[{}]",
v.iter()
.map(|f| f.to_string())
.collect::<Vec<_>>()
.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();

View file

@ -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<String>) -> Self {
Self { ok: false, data: None, error: Some(msg.into()) }
Self {
ok: false,
data: None,
error: Some(msg.into()),
}
}
}

View file

@ -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<Vec<EphemeralScarEvent>> {
pub async fn get_node_scars(
&self,
node_id: &str,
limit: i64,
) -> Result<Vec<EphemeralScarEvent>> {
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<Option<EphemeralReceipt>> {
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 \

View file

@ -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<WikiPage> {
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<WikiPage> {
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<Option<WikiPage>> {
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<Vec<WikiPage>> {
@ -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())
}
}

View file

@ -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<String, String>) -> 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::<Utc>::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::<serde_json::Value>(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())

View file

@ -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::<Vec<_>>().into_iter().rev().collect::<String>();
ar.details.insert(label.to_string(), json!({"rc": 0, "stdout_tail": tail.trim() }));
let tail = stdout
.chars()
.rev()
.take(500)
.collect::<Vec<_>>()
.into_iter()
.rev()
.collect::<String>();
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<String, String>) -> Result<(bool, String)> {
async fn emit_s3(
receipt: &serde_json::Value,
creds: &HashMap<String, String>,
) -> 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<String, String>) -
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(())

View file

@ -24,8 +24,14 @@ struct Cli {
#[derive(Subcommand, Debug)]
enum Commands {
Sync { #[arg(long)] since: Option<i64> },
Watch { #[arg(long, default_value = "60")] interval: u64 },
Sync {
#[arg(long)]
since: Option<i64>,
},
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<i64>) -> Result<()> {
async fn cmd_sync(
db_path: &PathBuf,
dsn: &str,
enable_embed: bool,
since: Option<i64>,
) -> 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<String>,
slug: String, directory: String, title: String,
agent: Option<String>, model: Option<String>,
time_created: i64, time_updated: i64,
tokens_input: i64, tokens_output: i64,
id: String,
project_id: String,
parent_id: Option<String>,
slug: String,
directory: String,
title: String,
agent: Option<String>,
model: Option<String>,
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<Vec<OpenCodeSession>> {
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::<Result<Vec<_>, _>>().map_err(|e| e.into())
@ -212,56 +242,84 @@ fn sessions_since(conn: &Connection, since_ms: i64) -> Result<Vec<OpenCodeSessio
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 WHERE time_updated > ?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::<Result<Vec<_>, _>>().map_err(|e| e.into())
}
fn max_session_updated(conn: &Connection) -> Result<Option<i64>> {
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<Vec<OpenCodeMessage>> {
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::<Result<Vec<_>, _>>().map_err(|e| e.into())
}
fn parts_for_message(conn: &Connection, message_id: &str) -> Result<Vec<OpenCodePart>> {
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::<Result<Vec<_>, _>>().map_err(|e| e.into())
}
fn normalize_session(sess: &OpenCodeSession, msgs: &[ChatMessage], _compaction: Option<String>) -> ChatSession {
fn normalize_session(
sess: &OpenCodeSession,
msgs: &[ChatMessage],
_compaction: Option<String>,
) -> 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<ChatMessage> {
fn normalize_message(
msg: &OpenCodeMessage,
parts: &[OpenCodePart],
index: i32,
) -> Result<ChatMessage> {
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<Vec<f32>> {
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())
}
}