mirror of
https://github.com/allaunthefox/Research-Stack.git
synced 2026-08-11 20:10:35 +00:00
fix(rds): restore ENE storage observation probe
This commit is contained in:
parent
971f17034c
commit
e6f770324d
11 changed files with 690 additions and 243 deletions
|
|
@ -29,7 +29,9 @@ struct SearchQuery {
|
||||||
semantic: bool,
|
semantic: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
fn default_limit() -> i64 { 10 }
|
fn default_limit() -> i64 {
|
||||||
|
10
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Serialize)]
|
#[derive(Serialize)]
|
||||||
struct HealthResponse {
|
struct HealthResponse {
|
||||||
|
|
@ -58,7 +60,11 @@ async fn main() -> anyhow::Result<()> {
|
||||||
let ephemeral = EphemeralSurface::new(ephemeral_client);
|
let ephemeral = EphemeralSurface::new(ephemeral_client);
|
||||||
ephemeral.init_tables().await?;
|
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()
|
let app = Router::new()
|
||||||
.route("/health", get(health_handler))
|
.route("/health", get(health_handler))
|
||||||
|
|
@ -88,7 +94,10 @@ async fn list_sessions(
|
||||||
State(state): State<Arc<Mutex<AppState>>>,
|
State(state): State<Arc<Mutex<AppState>>>,
|
||||||
Query(params): Query<HashMap<String, String>>,
|
Query(params): Query<HashMap<String, String>>,
|
||||||
) -> Result<Json<serde_json::Value>, 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;
|
let guard = state.lock().await;
|
||||||
match guard.chat.list_sessions(limit).await {
|
match guard.chat.list_sessions(limit).await {
|
||||||
Ok(v) => Ok(Json(json!({"ok": true, "data": v}))),
|
Ok(v) => Ok(Json(json!({"ok": true, "data": v}))),
|
||||||
|
|
@ -128,7 +137,10 @@ async fn wiki_search(
|
||||||
Query(params): Query<HashMap<String, String>>,
|
Query(params): Query<HashMap<String, String>>,
|
||||||
) -> Result<Json<serde_json::Value>, String> {
|
) -> Result<Json<serde_json::Value>, String> {
|
||||||
let query = params.get("q").cloned().unwrap_or_default();
|
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;
|
let guard = state.lock().await;
|
||||||
match guard.wiki.search(&query, limit).await {
|
match guard.wiki.search(&query, limit).await {
|
||||||
Ok(v) => Ok(Json(json!({"ok": true, "data": v}))),
|
Ok(v) => Ok(Json(json!({"ok": true, "data": v}))),
|
||||||
|
|
@ -153,7 +165,10 @@ async fn list_ephemeral_nodes(
|
||||||
Query(params): Query<HashMap<String, String>>,
|
Query(params): Query<HashMap<String, String>>,
|
||||||
) -> Result<Json<serde_json::Value>, String> {
|
) -> Result<Json<serde_json::Value>, String> {
|
||||||
let zone = params.get("zone").cloned().unwrap_or_else(|| "cold".into());
|
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;
|
let guard = state.lock().await;
|
||||||
match guard.ephemeral.list_nodes_by_zone(&zone, limit).await {
|
match guard.ephemeral.list_nodes_by_zone(&zone, limit).await {
|
||||||
Ok(v) => Ok(Json(json!({"ok": true, "data": v}))),
|
Ok(v) => Ok(Json(json!({"ok": true, "data": v}))),
|
||||||
|
|
|
||||||
|
|
@ -57,7 +57,12 @@ impl GossipMessage {
|
||||||
pub fn new(sender: &str, msg_type: &str, payload: serde_json::Value) -> Self {
|
pub fn new(sender: &str, msg_type: &str, payload: serde_json::Value) -> Self {
|
||||||
let id = format!(
|
let id = format!(
|
||||||
"gossip_{}",
|
"gossip_{}",
|
||||||
&sha256_hex(&format!("{}:{}:{}", sender, msg_type, Utc::now().timestamp_millis()))[..16]
|
&sha256_hex(&format!(
|
||||||
|
"{}:{}:{}",
|
||||||
|
sender,
|
||||||
|
msg_type,
|
||||||
|
Utc::now().timestamp_millis()
|
||||||
|
))[..16]
|
||||||
);
|
);
|
||||||
Self {
|
Self {
|
||||||
message_id: id,
|
message_id: id,
|
||||||
|
|
@ -88,8 +93,8 @@ impl GossipMessage {
|
||||||
use hmac::{Hmac, Mac};
|
use hmac::{Hmac, Mac};
|
||||||
use sha2::Sha256;
|
use sha2::Sha256;
|
||||||
type HmacSha256 = Hmac<Sha256>;
|
type HmacSha256 = Hmac<Sha256>;
|
||||||
let mut mac = HmacSha256::new_from_slice(secret.as_bytes())
|
let mut mac =
|
||||||
.expect("HMAC can take key of any size");
|
HmacSha256::new_from_slice(secret.as_bytes()).expect("HMAC can take key of any size");
|
||||||
mac.update(&self.canonical_bytes());
|
mac.update(&self.canonical_bytes());
|
||||||
let result = mac.finalize();
|
let result = mac.finalize();
|
||||||
self.signature = Some(hex::encode(result.into_bytes()));
|
self.signature = Some(hex::encode(result.into_bytes()));
|
||||||
|
|
@ -97,7 +102,9 @@ impl GossipMessage {
|
||||||
|
|
||||||
/// Verify the HMAC signature against the cluster secret.
|
/// Verify the HMAC signature against the cluster secret.
|
||||||
pub fn verify(&self, secret: &str) -> bool {
|
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 hmac::{Hmac, Mac};
|
||||||
use sha2::Sha256;
|
use sha2::Sha256;
|
||||||
type HmacSha256 = Hmac<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) \
|
replication_version, capabilities, health_score_q16, is_active) \
|
||||||
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
|
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
|
||||||
rusqlite::params![
|
rusqlite::params![
|
||||||
&peer.node_id, &peer.public_key, &peer.ip_address, &peer.port,
|
&peer.node_id,
|
||||||
&peer.first_seen, &peer.last_seen, &peer.replication_version,
|
&peer.public_key,
|
||||||
|
&peer.ip_address,
|
||||||
|
&peer.port,
|
||||||
|
&peer.first_seen,
|
||||||
|
&peer.last_seen,
|
||||||
|
&peer.replication_version,
|
||||||
serde_json::to_string(&peer.capabilities)?,
|
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(())
|
Ok(())
|
||||||
|
|
@ -244,7 +257,7 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no
|
||||||
let mut stmt = self.conn.prepare(
|
let mut stmt = self.conn.prepare(
|
||||||
"SELECT node_id, public_key, ip_address, port, first_seen, last_seen, \
|
"SELECT node_id, public_key, ip_address, port, first_seen, last_seen, \
|
||||||
replication_version, capabilities, health_score_q16, is_active \
|
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 rows = stmt.query_map([], |row| {
|
||||||
let caps: String = row.get(7)?;
|
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();
|
let mut out = Vec::new();
|
||||||
for r in rows { out.push(r?); }
|
for r in rows {
|
||||||
|
out.push(r?);
|
||||||
|
}
|
||||||
Ok(out)
|
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) \
|
(message_id, sender_node, message_type, payload, timestamp) \
|
||||||
VALUES (?1, ?2, ?3, ?4, ?5)",
|
VALUES (?1, ?2, ?3, ?4, ?5)",
|
||||||
rusqlite::params![
|
rusqlite::params![
|
||||||
&msg.message_id, &msg.sender_node, &msg.message_type,
|
&msg.message_id,
|
||||||
&serde_json::to_string(&msg.payload)?, &msg.timestamp,
|
&msg.sender_node,
|
||||||
|
&msg.message_type,
|
||||||
|
&serde_json::to_string(&msg.payload)?,
|
||||||
|
&msg.timestamp,
|
||||||
],
|
],
|
||||||
)?;
|
)?;
|
||||||
Ok(())
|
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) \
|
(proposal_id, credential_id, proposer, timestamp, votes, resolved) \
|
||||||
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
|
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
|
||||||
rusqlite::params![
|
rusqlite::params![
|
||||||
&prop.proposal_id, &prop.credential_id, &prop.proposer,
|
&prop.proposal_id,
|
||||||
&prop.timestamp, &serde_json::to_string(&prop.votes)?,
|
&prop.credential_id,
|
||||||
|
&prop.proposer,
|
||||||
|
&prop.timestamp,
|
||||||
|
&serde_json::to_string(&prop.votes)?,
|
||||||
&prop.resolved as &dyn rusqlite::ToSql,
|
&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>> {
|
pub fn load_proposal(&self, proposal_id: &str) -> Result<Option<ConsensusProposal>> {
|
||||||
let mut stmt = self.conn.prepare(
|
let mut stmt = self.conn.prepare(
|
||||||
"SELECT proposal_id, credential_id, proposer, timestamp, votes, resolved \
|
"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 row = stmt
|
||||||
let votes_str: String = row.get(4)?;
|
.query_row([proposal_id], |row| {
|
||||||
Ok(ConsensusProposal {
|
let votes_str: String = row.get(4)?;
|
||||||
proposal_id: row.get(0)?,
|
Ok(ConsensusProposal {
|
||||||
credential_id: row.get(1)?,
|
proposal_id: row.get(0)?,
|
||||||
proposer: row.get(2)?,
|
credential_id: row.get(1)?,
|
||||||
timestamp: row.get(3)?,
|
proposer: row.get(2)?,
|
||||||
votes: serde_json::from_str(&votes_str).unwrap_or_default(),
|
timestamp: row.get(3)?,
|
||||||
resolved: row.get::<_, i32>(5)? != 0,
|
votes: serde_json::from_str(&votes_str).unwrap_or_default(),
|
||||||
|
resolved: row.get::<_, i32>(5)? != 0,
|
||||||
|
})
|
||||||
})
|
})
|
||||||
}).optional()?;
|
.optional()?;
|
||||||
Ok(row)
|
Ok(row)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -329,7 +352,9 @@ CREATE INDEX IF NOT EXISTS idx_replications_target ON ene_replications(target_no
|
||||||
})
|
})
|
||||||
})?;
|
})?;
|
||||||
let mut out = Vec::new();
|
let mut out = Vec::new();
|
||||||
for r in rows { out.push(r?); }
|
for r in rows {
|
||||||
|
out.push(r?);
|
||||||
|
}
|
||||||
Ok(out)
|
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) \
|
usage_count, last_rotated, health_score_q16, is_active) \
|
||||||
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
|
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
|
||||||
rusqlite::params![
|
rusqlite::params![
|
||||||
&frag.credential_id, &frag.provider, &frag.fragment,
|
&frag.credential_id,
|
||||||
&frag.access_level, &serde_json::to_string(&frag.node_assignments)?,
|
&frag.provider,
|
||||||
&frag.usage_count, &frag.last_rotated,
|
&frag.fragment,
|
||||||
&frag.health_score_q16, &frag.is_active as &dyn rusqlite::ToSql,
|
&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(())
|
Ok(())
|
||||||
|
|
@ -388,7 +418,10 @@ impl EneNode {
|
||||||
let loaded_peers = db.load_peers().unwrap_or_default();
|
let loaded_peers = db.load_peers().unwrap_or_default();
|
||||||
let mut identity = NodeIdentity::default();
|
let mut identity = NodeIdentity::default();
|
||||||
identity.node_id = node_id.unwrap_or_else(|| {
|
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();
|
identity.public_key = sha256_hex(&identity.node_id)[..32].to_string();
|
||||||
db.save_peer(&identity)?;
|
db.save_peer(&identity)?;
|
||||||
|
|
@ -439,8 +472,15 @@ impl EneNode {
|
||||||
drop(peers_guard);
|
drop(peers_guard);
|
||||||
|
|
||||||
self.with_db(|db| db.save_gossip(&msg))?;
|
self.with_db(|db| db.save_gossip(&msg))?;
|
||||||
self.seen_message_ids.write().await.insert(msg.message_id.clone());
|
self.seen_message_ids
|
||||||
info!("gossip {} -> {} peers", msg.message_type, self.peers.read().await.len());
|
.write()
|
||||||
|
.await
|
||||||
|
.insert(msg.message_id.clone());
|
||||||
|
info!(
|
||||||
|
"gossip {} -> {} peers",
|
||||||
|
msg.message_type,
|
||||||
|
self.peers.read().await.len()
|
||||||
|
);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -455,7 +495,10 @@ impl EneNode {
|
||||||
if self.seen_message_ids.read().await.contains(&msg.message_id) {
|
if self.seen_message_ids.read().await.contains(&msg.message_id) {
|
||||||
return Ok(());
|
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))?;
|
self.with_db(|db| db.save_gossip(&msg))?;
|
||||||
|
|
||||||
match msg.message_type.as_str() {
|
match msg.message_type.as_str() {
|
||||||
|
|
@ -471,9 +514,15 @@ impl EneNode {
|
||||||
|
|
||||||
async fn handle_discovery(&self, msg: &GossipMessage, from: SocketAddr) -> Result<()> {
|
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 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())
|
.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()]);
|
.unwrap_or_else(|| vec!["storage".into(), "compute".into()]);
|
||||||
|
|
||||||
if let Some(nid) = node_id {
|
if let Some(nid) = node_id {
|
||||||
|
|
@ -517,9 +566,18 @@ impl EneNode {
|
||||||
if let (Some(id), Some(b64)) = (cred_id, fragment_b64) {
|
if let (Some(id), Some(b64)) = (cred_id, fragment_b64) {
|
||||||
let frag = CredentialFragment {
|
let frag = CredentialFragment {
|
||||||
credential_id: id.into(),
|
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),
|
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()],
|
node_assignments: vec![msg.sender_node.clone()],
|
||||||
usage_count: 0,
|
usage_count: 0,
|
||||||
last_rotated: Utc::now().timestamp_millis(),
|
last_rotated: Utc::now().timestamp_millis(),
|
||||||
|
|
@ -569,7 +627,8 @@ impl EneNode {
|
||||||
|
|
||||||
pub async fn vote(&self, proposal_id: &str, approve: bool) -> Result<bool> {
|
pub async fn vote(&self, proposal_id: &str, approve: bool) -> Result<bool> {
|
||||||
let total = self.peers.read().await.len() + 1;
|
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")?;
|
.context("proposal not found")?;
|
||||||
prop.votes.insert(self.identity.node_id.clone(), approve);
|
prop.votes.insert(self.identity.node_id.clone(), approve);
|
||||||
self.with_db(|db| db.save_proposal(&prop))?;
|
self.with_db(|db| db.save_proposal(&prop))?;
|
||||||
|
|
@ -579,7 +638,10 @@ impl EneNode {
|
||||||
if approve_count >= threshold && !prop.resolved {
|
if approve_count >= threshold && !prop.resolved {
|
||||||
prop.resolved = true;
|
prop.resolved = true;
|
||||||
self.with_db(|db| db.save_proposal(&prop))?;
|
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);
|
return Ok(true);
|
||||||
}
|
}
|
||||||
Ok(false)
|
Ok(false)
|
||||||
|
|
@ -588,7 +650,11 @@ impl EneNode {
|
||||||
pub async fn propose_rotation(&self, credential_id: &str) -> Result<String> {
|
pub async fn propose_rotation(&self, credential_id: &str) -> Result<String> {
|
||||||
let proposal_id = format!(
|
let proposal_id = format!(
|
||||||
"prop_{}",
|
"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!({
|
let payload = serde_json::json!({
|
||||||
"proposal_id": &proposal_id,
|
"proposal_id": &proposal_id,
|
||||||
|
|
@ -596,7 +662,11 @@ impl EneNode {
|
||||||
"proposer": &self.identity.node_id,
|
"proposer": &self.identity.node_id,
|
||||||
"timestamp": Utc::now().timestamp_millis(),
|
"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?;
|
self.gossip(msg).await?;
|
||||||
Ok(proposal_id)
|
Ok(proposal_id)
|
||||||
}
|
}
|
||||||
|
|
@ -627,7 +697,10 @@ impl EneNode {
|
||||||
|
|
||||||
pub async fn get_mesh_health(&self) -> serde_json::Value {
|
pub async fn get_mesh_health(&self) -> serde_json::Value {
|
||||||
let peers = self.peers.read().await;
|
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;
|
let mesh_size = peers.len() + 1;
|
||||||
serde_json::json!({
|
serde_json::json!({
|
||||||
"mesh_size": mesh_size,
|
"mesh_size": mesh_size,
|
||||||
|
|
@ -639,7 +712,9 @@ impl EneNode {
|
||||||
|
|
||||||
fn base64_decode(s: &str) -> Vec<u8> {
|
fn base64_decode(s: &str) -> Vec<u8> {
|
||||||
use base64::Engine;
|
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)]
|
#[cfg(test)]
|
||||||
|
|
@ -648,7 +723,8 @@ mod tests {
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn gossip_sign_and_verify_roundtrip() {
|
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());
|
assert!(msg.signature.is_none());
|
||||||
msg.sign("secret-key");
|
msg.sign("secret-key");
|
||||||
assert!(msg.signature.is_some());
|
assert!(msg.signature.is_some());
|
||||||
|
|
@ -670,7 +746,8 @@ mod tests {
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn gossip_tamper_payload_invalidates_signature() {
|
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");
|
msg.sign("secret-key");
|
||||||
assert!(msg.verify("secret-key"));
|
assert!(msg.verify("secret-key"));
|
||||||
// Tamper with payload after signing
|
// Tamper with payload after signing
|
||||||
|
|
|
||||||
|
|
@ -35,11 +35,23 @@ async fn main() -> Result<()> {
|
||||||
let _ = std::fs::create_dir_all(parent);
|
let _ = std::fs::create_dir_all(parent);
|
||||||
}
|
}
|
||||||
|
|
||||||
let cluster_secret = cli.cluster_secret.or_else(|| std::env::var("ENE_CLUSTER_SECRET").ok());
|
let cluster_secret = cli
|
||||||
let node = EneNode::new(cli.node_id, &cli.db, cli.bind, cli.seed.clone(), cluster_secret).await?;
|
.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);
|
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 ────────────────────────────────────────────────
|
// ── UDP listener task ────────────────────────────────────────────────
|
||||||
let listen_node = node.clone();
|
let listen_node = node.clone();
|
||||||
|
|
@ -82,7 +94,8 @@ async fn main() -> Result<()> {
|
||||||
// ── Heartbeat loop ────────────────────────────────────────────────────
|
// ── Heartbeat loop ────────────────────────────────────────────────────
|
||||||
let beat_node = node.clone();
|
let beat_node = node.clone();
|
||||||
let heartbeat = tokio::spawn(async move {
|
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
|
interval.tick().await; // skip immediate first tick
|
||||||
loop {
|
loop {
|
||||||
interval.tick().await;
|
interval.tick().await;
|
||||||
|
|
@ -100,8 +113,14 @@ async fn main() -> Result<()> {
|
||||||
interval.tick().await;
|
interval.tick().await;
|
||||||
let status = report_node.get_status().await;
|
let status = report_node.get_status().await;
|
||||||
let health = report_node.get_mesh_health().await;
|
let health = report_node.get_mesh_health().await;
|
||||||
info!("status: {}", serde_json::to_string(&status).unwrap_or_default());
|
info!(
|
||||||
info!("health: {}", serde_json::to_string(&health).unwrap_or_default());
|
"status: {}",
|
||||||
|
serde_json::to_string(&status).unwrap_or_default()
|
||||||
|
);
|
||||||
|
info!(
|
||||||
|
"health: {}",
|
||||||
|
serde_json::to_string(&health).unwrap_or_default()
|
||||||
|
);
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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_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);
|
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");
|
info!("chat log schema initialized");
|
||||||
Ok(())
|
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> {
|
pub async fn delete_messages_for_session(&self, session_id: &str) -> Result<u64> {
|
||||||
let rows = self.client.inner()
|
let rows = self
|
||||||
.execute("DELETE FROM ene.chat_messages WHERE session_id = $1", &[&session_id])
|
.client
|
||||||
|
.inner()
|
||||||
|
.execute(
|
||||||
|
"DELETE FROM ene.chat_messages WHERE session_id = $1",
|
||||||
|
&[&session_id],
|
||||||
|
)
|
||||||
.await
|
.await
|
||||||
.context("delete messages")?;
|
.context("delete messages")?;
|
||||||
Ok(rows)
|
Ok(rows)
|
||||||
|
|
@ -181,19 +190,30 @@ CREATE INDEX IF NOT EXISTS idx_chat_messages_tool_search ON ene.chat_messages US
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.context("keyword search")?;
|
.context("keyword search")?;
|
||||||
Ok(rows.iter().map(|r| serde_json::json!({
|
Ok(rows
|
||||||
"session_id": r.get::<_, String>(0),
|
.iter()
|
||||||
"title": r.get::<_, String>(1),
|
.map(|r| {
|
||||||
"agent": r.get::<_, Option<String>>(2),
|
serde_json::json!({
|
||||||
"model": r.get::<_, Option<String>>(3),
|
"session_id": r.get::<_, String>(0),
|
||||||
"match_count": r.get::<_, i64>(4),
|
"title": r.get::<_, String>(1),
|
||||||
"rank": r.get::<_, f32>(5),
|
"agent": r.get::<_, Option<String>>(2),
|
||||||
})).collect())
|
"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 vec_str = vec_to_pgtext(embedding);
|
||||||
let rows = self.client.inner()
|
let rows = self
|
||||||
|
.client
|
||||||
|
.inner()
|
||||||
.query(
|
.query(
|
||||||
"SELECT session_id, title, agent, model, \
|
"SELECT session_id, title, agent, model, \
|
||||||
1 - (embedding <=> $1::vector) AS similarity \
|
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
|
.await
|
||||||
.context("similarity search")?;
|
.context("similarity search")?;
|
||||||
Ok(rows.iter().map(|r| serde_json::json!({
|
Ok(rows
|
||||||
"session_id": r.get::<_, String>(0),
|
.iter()
|
||||||
"title": r.get::<_, String>(1),
|
.map(|r| {
|
||||||
"agent": r.get::<_, Option<String>>(2),
|
serde_json::json!({
|
||||||
"model": r.get::<_, Option<String>>(3),
|
"session_id": r.get::<_, String>(0),
|
||||||
"similarity": r.get::<_, f32>(4),
|
"title": r.get::<_, String>(1),
|
||||||
})).collect())
|
"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>> {
|
pub async fn list_sessions(&self, limit: i64) -> Result<Vec<serde_json::Value>> {
|
||||||
let rows = self.client.inner()
|
let rows = self
|
||||||
|
.client
|
||||||
|
.inner()
|
||||||
.query(
|
.query(
|
||||||
"SELECT session_id, title, agent, model, message_count, \
|
"SELECT session_id, title, agent, model, message_count, \
|
||||||
token_input_total, token_output_total, created_at_ms, updated_at_ms \
|
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
|
.await
|
||||||
.context("list sessions")?;
|
.context("list sessions")?;
|
||||||
Ok(rows.iter().map(|r| serde_json::json!({
|
Ok(rows
|
||||||
"session_id": r.get::<_, String>(0),
|
.iter()
|
||||||
"title": r.get::<_, String>(1),
|
.map(|r| {
|
||||||
"agent": r.get::<_, Option<String>>(2),
|
serde_json::json!({
|
||||||
"model": r.get::<_, Option<String>>(3),
|
"session_id": r.get::<_, String>(0),
|
||||||
"message_count": r.get::<_, i32>(4),
|
"title": r.get::<_, String>(1),
|
||||||
"token_input_total": r.get::<_, i64>(5),
|
"agent": r.get::<_, Option<String>>(2),
|
||||||
"token_output_total": r.get::<_, i64>(6),
|
"model": r.get::<_, Option<String>>(3),
|
||||||
"created_at_ms": r.get::<_, i64>(7),
|
"message_count": r.get::<_, i32>(4),
|
||||||
"updated_at_ms": r.get::<_, i64>(8),
|
"token_input_total": r.get::<_, i64>(5),
|
||||||
})).collect())
|
"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>> {
|
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(
|
.query_opt(
|
||||||
"SELECT session_id, title, agent, model, message_count, \
|
"SELECT session_id, title, agent, model, message_count, \
|
||||||
token_input_total, token_output_total, created_at_ms, updated_at_ms, meta \
|
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
|
.await
|
||||||
.context("get session")?;
|
.context("get session")?;
|
||||||
let Some(sess) = sess else { return Ok(None) };
|
let Some(sess) = sess else { return Ok(None) };
|
||||||
let msgs = self.client.inner()
|
let msgs = self
|
||||||
|
.client
|
||||||
|
.inner()
|
||||||
.query(
|
.query(
|
||||||
"SELECT message_index, role, blocks, text_content, token_input, \
|
"SELECT message_index, role, blocks, text_content, token_input, \
|
||||||
token_output, tool_calls, created_at_ms \
|
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
|
.await
|
||||||
.context("get messages")?;
|
.context("get messages")?;
|
||||||
let messages: Vec<_> = msgs.iter().map(|r| serde_json::json!({
|
let messages: Vec<_> = msgs
|
||||||
"message_index": r.get::<_, i32>(0),
|
.iter()
|
||||||
"role": r.get::<_, String>(1),
|
.map(|r| {
|
||||||
"blocks": r.get::<_, serde_json::Value>(2),
|
serde_json::json!({
|
||||||
"text_content": r.get::<_, String>(3),
|
"message_index": r.get::<_, i32>(0),
|
||||||
"token_input": r.get::<_, i64>(4),
|
"role": r.get::<_, String>(1),
|
||||||
"token_output": r.get::<_, i64>(5),
|
"blocks": r.get::<_, serde_json::Value>(2),
|
||||||
"tool_calls": r.get::<_, serde_json::Value>(6),
|
"text_content": r.get::<_, String>(3),
|
||||||
"created_at_ms": r.get::<_, i64>(7),
|
"token_input": r.get::<_, i64>(4),
|
||||||
})).collect();
|
"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!({
|
Ok(Some(serde_json::json!({
|
||||||
"session_id": sess.get::<_, String>(0),
|
"session_id": sess.get::<_, String>(0),
|
||||||
"title": sess.get::<_, String>(1),
|
"title": sess.get::<_, String>(1),
|
||||||
|
|
|
||||||
|
|
@ -34,8 +34,9 @@ impl RdsClient {
|
||||||
if let Ok(dsn) = std::env::var("RDS_DSN") {
|
if let Ok(dsn) = std::env::var("RDS_DSN") {
|
||||||
return dsn;
|
return dsn;
|
||||||
}
|
}
|
||||||
let host = std::env::var("RDS_HOST")
|
let host = std::env::var("RDS_HOST").unwrap_or_else(|_| {
|
||||||
.unwrap_or_else(|_| "database-1.cluster-c9i0w8eu8fnv.us-east-2.rds.amazonaws.com".into());
|
"database-1.cluster-c9i0w8eu8fnv.us-east-2.rds.amazonaws.com".into()
|
||||||
|
});
|
||||||
let port = std::env::var("RDS_PORT").unwrap_or_else(|_| "5432".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 user = std::env::var("RDS_USER").unwrap_or_else(|_| "postgres".into());
|
||||||
let password = std::env::var("RDS_PASSWORD")
|
let password = std::env::var("RDS_PASSWORD")
|
||||||
|
|
@ -49,7 +50,11 @@ impl RdsClient {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Raw query helper.
|
/// 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?;
|
let rows = self.client.execute(sql, params).await?;
|
||||||
Ok(rows)
|
Ok(rows)
|
||||||
}
|
}
|
||||||
|
|
@ -119,7 +124,13 @@ pub fn sha256_text(text: &str) -> String {
|
||||||
|
|
||||||
/// Format a float vector as pgvector text: [0.1,0.2,...]
|
/// Format a float vector as pgvector text: [0.1,0.2,...]
|
||||||
pub fn vec_to_pgtext(v: &[f32]) -> String {
|
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)]
|
#[cfg(test)]
|
||||||
|
|
@ -128,7 +139,10 @@ mod tests {
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn dsn_from_env_uses_rds_dsn_when_set() {
|
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
|
// Clear other vars to ensure RDS_DSN takes precedence
|
||||||
for key in &["RDS_HOST", "RDS_PORT", "RDS_USER", "RDS_PASSWORD", "RDS_DB"] {
|
for key in &["RDS_HOST", "RDS_PORT", "RDS_USER", "RDS_PASSWORD", "RDS_DB"] {
|
||||||
let _ = std::env::remove_var(key);
|
let _ = std::env::remove_var(key);
|
||||||
|
|
@ -162,7 +176,14 @@ mod tests {
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn dsn_from_env_uses_defaults() {
|
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 _ = std::env::remove_var(key);
|
||||||
}
|
}
|
||||||
let dsn = RdsClient::dsn_from_env();
|
let dsn = RdsClient::dsn_from_env();
|
||||||
|
|
|
||||||
|
|
@ -20,10 +20,18 @@ pub struct ApiResponse {
|
||||||
|
|
||||||
impl ApiResponse {
|
impl ApiResponse {
|
||||||
pub fn success(data: serde_json::Value) -> Self {
|
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 {
|
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()),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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_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);
|
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");
|
info!("ephemeral node schema initialized");
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
@ -215,23 +219,27 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events(
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.context("list ephemeral nodes")?;
|
.context("list ephemeral nodes")?;
|
||||||
Ok(rows.iter().map(|r| EphemeralNode {
|
Ok(rows
|
||||||
node_id: r.get(0),
|
.iter()
|
||||||
thermal_zone: r.get(1),
|
.map(|r| EphemeralNode {
|
||||||
reliability_raw: r.get(2),
|
node_id: r.get(0),
|
||||||
latency_p95_ms: r.get(3),
|
thermal_zone: r.get(1),
|
||||||
scar_count: r.get(4),
|
reliability_raw: r.get(2),
|
||||||
last_seen_ms: r.get(5),
|
latency_p95_ms: r.get(3),
|
||||||
reputation_raw: r.get(6),
|
scar_count: r.get(4),
|
||||||
quarantine_until_ms: r.get(7),
|
last_seen_ms: r.get(5),
|
||||||
meta: r.get(8),
|
reputation_raw: r.get(6),
|
||||||
created_at_ms: r.get(9),
|
quarantine_until_ms: r.get(7),
|
||||||
updated_at_ms: r.get(10),
|
meta: r.get(8),
|
||||||
}).collect())
|
created_at_ms: r.get(9),
|
||||||
|
updated_at_ms: r.get(10),
|
||||||
|
})
|
||||||
|
.collect())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn insert_task(&self, task: &EphemeralTask) -> Result<()> {
|
pub async fn insert_task(&self, task: &EphemeralTask) -> Result<()> {
|
||||||
self.client.inner()
|
self.client
|
||||||
|
.inner()
|
||||||
.execute(
|
.execute(
|
||||||
"INSERT INTO ene.ephemeral_tasks \
|
"INSERT INTO ene.ephemeral_tasks \
|
||||||
(task_id, session_id, node_id, task_state, priority_raw, ttl_ms, \
|
(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, \
|
result_hash = EXCLUDED.result_hash, \
|
||||||
meta = EXCLUDED.meta",
|
meta = EXCLUDED.meta",
|
||||||
&[
|
&[
|
||||||
&task.task_id, &task.session_id, &task.node_id, &task.task_state,
|
&task.task_id,
|
||||||
&task.priority_raw, &task.ttl_ms, &task.dispatched_at_ms,
|
&task.session_id,
|
||||||
&task.completed_at_ms, &task.result_hash, &task.meta, &task.created_at_ms,
|
&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
|
.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<()> {
|
pub async fn insert_metric(&self, metric: &EphemeralMetric) -> Result<()> {
|
||||||
self.client.inner()
|
self.client
|
||||||
|
.inner()
|
||||||
.execute(
|
.execute(
|
||||||
"INSERT INTO ene.ephemeral_metrics \
|
"INSERT INTO ene.ephemeral_metrics \
|
||||||
(metric_id, node_id, metric_name, metric_value_raw, metric_scale, recorded_at_ms) \
|
(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, \
|
metric_scale = EXCLUDED.metric_scale, \
|
||||||
recorded_at_ms = EXCLUDED.recorded_at_ms",
|
recorded_at_ms = EXCLUDED.recorded_at_ms",
|
||||||
&[
|
&[
|
||||||
&metric.metric_id, &metric.node_id, &metric.metric_name,
|
&metric.metric_id,
|
||||||
&metric.metric_value_raw, &metric.metric_scale, &metric.recorded_at_ms,
|
&metric.node_id,
|
||||||
|
&metric.metric_name,
|
||||||
|
&metric.metric_value_raw,
|
||||||
|
&metric.metric_scale,
|
||||||
|
&metric.recorded_at_ms,
|
||||||
],
|
],
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
|
|
@ -294,7 +315,11 @@ CREATE INDEX IF NOT EXISTS idx_ephemeral_scar_node ON ene.ephemeral_scar_events(
|
||||||
Ok(())
|
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()
|
let rows = self.client.inner()
|
||||||
.query(
|
.query(
|
||||||
"SELECT scar_id, node_id, task_id, scar_pressure, failure_mode, coarsening_agent, created_at_ms \
|
"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
|
.await
|
||||||
.context("get node scars")?;
|
.context("get node scars")?;
|
||||||
Ok(rows.iter().map(|r| EphemeralScarEvent {
|
Ok(rows
|
||||||
scar_id: r.get(0),
|
.iter()
|
||||||
node_id: r.get(1),
|
.map(|r| EphemeralScarEvent {
|
||||||
task_id: r.get(2),
|
scar_id: r.get(0),
|
||||||
scar_pressure: r.get(3),
|
node_id: r.get(1),
|
||||||
failure_mode: r.get(4),
|
task_id: r.get(2),
|
||||||
coarsening_agent: r.get(5),
|
scar_pressure: r.get(3),
|
||||||
created_at_ms: r.get(6),
|
failure_mode: r.get(4),
|
||||||
}).collect())
|
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>> {
|
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(
|
.query_opt(
|
||||||
"SELECT receipt_id, task_id, node_id, cross_matrix, sidon_slack, \
|
"SELECT receipt_id, task_id, node_id, cross_matrix, sidon_slack, \
|
||||||
step_count, residual_series, write_timing_ms, scar_absent, created_at_ms \
|
step_count, residual_series, write_timing_ms, scar_absent, created_at_ms \
|
||||||
|
|
|
||||||
|
|
@ -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_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);
|
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");
|
info!("wiki schema initialized");
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn put_page(&self, title: &str, content: &str, editor: &str, summary: &str) -> Result<WikiPage> {
|
pub async fn put_page(
|
||||||
let slug = title.to_lowercase().replace(' ', "-").replace(|c: char| !c.is_alphanumeric() && c != '-', "");
|
&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()
|
let row = self.client.inner()
|
||||||
.query_one(
|
.query_one(
|
||||||
"INSERT INTO ene.wiki_pages (title, slug, content, updated_at) \
|
"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>> {
|
pub async fn get_page(&self, slug: &str) -> Result<Option<WikiPage>> {
|
||||||
let row = self.client.inner()
|
let row = self
|
||||||
|
.client
|
||||||
|
.inner()
|
||||||
.query_opt(
|
.query_opt(
|
||||||
"SELECT id, title, slug, content, concept_anchor, concept_vector, \
|
"SELECT id, title, slug, content, concept_anchor, concept_vector, \
|
||||||
created_at::text, updated_at::text FROM ene.wiki_pages WHERE slug = $1",
|
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
|
.await
|
||||||
.context("search wiki")?;
|
.context("search wiki")?;
|
||||||
Ok(rows.iter().map(|r| json!({
|
Ok(rows
|
||||||
"id": r.get::<_, i64>(0),
|
.iter()
|
||||||
"title": r.get::<_, String>(1),
|
.map(|r| {
|
||||||
"slug": r.get::<_, String>(2),
|
json!({
|
||||||
"rank": r.get::<_, f32>(3),
|
"id": r.get::<_, i64>(0),
|
||||||
})).collect())
|
"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>> {
|
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
|
.await
|
||||||
.context("recent wiki pages")?;
|
.context("recent wiki pages")?;
|
||||||
Ok(rows.iter().map(|r| WikiPage {
|
Ok(rows
|
||||||
id: r.get(0),
|
.iter()
|
||||||
title: r.get(1),
|
.map(|r| WikiPage {
|
||||||
slug: r.get(2),
|
id: r.get(0),
|
||||||
content: r.get(3),
|
title: r.get(1),
|
||||||
concept_anchor: r.get(4),
|
slug: r.get(2),
|
||||||
concept_vector: r.get(5),
|
content: r.get(3),
|
||||||
created_at: r.get(6),
|
concept_anchor: r.get(4),
|
||||||
updated_at: r.get(7),
|
concept_vector: r.get(5),
|
||||||
}).collect())
|
created_at: r.get(6),
|
||||||
|
updated_at: r.get(7),
|
||||||
|
})
|
||||||
|
.collect())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,7 @@ use serde::{Deserialize, Serialize};
|
||||||
use sha2::{Digest, Sha256};
|
use sha2::{Digest, Sha256};
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::path::PathBuf;
|
use std::path::PathBuf;
|
||||||
|
use std::time::Duration;
|
||||||
|
|
||||||
// ── constants ──────────────────────────────────────────────────────────────
|
// ── 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 ───────────────────────────────────────────────────────────────
|
// ── Decision ───────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
|
||||||
|
|
@ -145,19 +189,23 @@ pub fn decide(obs: &Observation) -> Decision {
|
||||||
if !obs.garage.up {
|
if !obs.garage.up {
|
||||||
d.alerts.push("ALERT: Garage S3 is unreachable".into());
|
d.alerts.push("ALERT: Garage S3 is unreachable".into());
|
||||||
d.trigger_garage_restart = true;
|
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 {
|
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.trigger_snap = true;
|
||||||
d.rationale.push("snapshot_count=0 → trigger_snap".into());
|
d.rationale.push("snapshot_count=0 → trigger_snap".into());
|
||||||
}
|
}
|
||||||
|
|
||||||
if !obs.backup_log.last_ok && obs.garage.up {
|
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.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
|
if obs.restic.dedup_ratio_q16 > 0
|
||||||
|
|
@ -180,14 +228,18 @@ pub fn decide(obs: &Observation) -> Decision {
|
||||||
}
|
}
|
||||||
|
|
||||||
if obs.cold_copy_needed {
|
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.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 {
|
if obs.garage.up {
|
||||||
d.trigger_offload = true;
|
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
|
d
|
||||||
|
|
@ -261,7 +313,10 @@ pub fn resume_chain() -> (i64, String) {
|
||||||
let mut last_hash = String::new();
|
let mut last_hash = String::new();
|
||||||
for line in text.lines() {
|
for line in text.lines() {
|
||||||
if let Ok(entry) = serde_json::from_str::<serde_json::Value>(line) {
|
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
|
last_hash = entry
|
||||||
.get("receipt_hash")
|
.get("receipt_hash")
|
||||||
.and_then(|v| v.as_str())
|
.and_then(|v| v.as_str())
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,8 @@
|
||||||
use anyhow::{Context, Result};
|
use anyhow::{Context, Result};
|
||||||
use chrono::Utc;
|
|
||||||
use clap::Parser;
|
use clap::Parser;
|
||||||
use ene_storage::{
|
use ene_storage::{
|
||||||
build_receipt, decide, load_garage_env, observe, resume_chain, ActionResult, Decision,
|
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 serde_json::json;
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
|
|
@ -25,13 +24,10 @@ async fn run_cmd(
|
||||||
for (k, v) in env_extra {
|
for (k, v) in env_extra {
|
||||||
cmd.env(k, v);
|
cmd.env(k, v);
|
||||||
}
|
}
|
||||||
let output = tokio::time::timeout(
|
let output = tokio::time::timeout(std::time::Duration::from_secs(timeout_secs), cmd.output())
|
||||||
std::time::Duration::from_secs(timeout_secs),
|
.await
|
||||||
cmd.output(),
|
.context("command timed out")?
|
||||||
)
|
.context("command failed to run")?;
|
||||||
.await
|
|
||||||
.context("command timed out")?
|
|
||||||
.context("command failed to run")?;
|
|
||||||
let rc = output.status.code().unwrap_or(-1);
|
let rc = output.status.code().unwrap_or(-1);
|
||||||
let stdout = String::from_utf8_lossy(&output.stdout).into_owned();
|
let stdout = String::from_utf8_lossy(&output.stdout).into_owned();
|
||||||
let stderr = String::from_utf8_lossy(&output.stderr).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 {
|
match run_cmd(args, env_extra, timeout_secs).await {
|
||||||
Ok((0, stdout, _)) => {
|
Ok((0, stdout, _)) => {
|
||||||
ar.actions_succeeded.push(label.to_string());
|
ar.actions_succeeded.push(label.to_string());
|
||||||
let tail = stdout.chars().rev().take(500).collect::<Vec<_>>().into_iter().rev().collect::<String>();
|
let tail = stdout
|
||||||
ar.details.insert(label.to_string(), json!({"rc": 0, "stdout_tail": tail.trim() }));
|
.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)) => {
|
Ok((rc, stdout, stderr)) => {
|
||||||
ar.actions_failed.push(label.to_string());
|
ar.actions_failed.push(label.to_string());
|
||||||
|
|
@ -67,7 +73,10 @@ async fn act_one(
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
ar.actions_failed.push(label.to_string());
|
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();
|
let mut ar = ActionResult::default();
|
||||||
|
|
||||||
if probe_only || dry_run {
|
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;
|
return ar;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -94,16 +106,35 @@ async fn act(
|
||||||
let consolidate_sh = storage_dir.join("garage/db-consolidate.sh");
|
let consolidate_sh = storage_dir.join("garage/db-consolidate.sh");
|
||||||
|
|
||||||
if d.trigger_garage_restart {
|
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()) {
|
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 {
|
if d.trigger_snap {
|
||||||
act_one(
|
act_one(
|
||||||
"restic_snap",
|
"restic_snap",
|
||||||
&["bash", &backup_sh.to_string_lossy(), "snap", "agent-triggered"],
|
&[
|
||||||
|
"bash",
|
||||||
|
&backup_sh.to_string_lossy(),
|
||||||
|
"snap",
|
||||||
|
"agent-triggered",
|
||||||
|
],
|
||||||
&mut ar,
|
&mut ar,
|
||||||
creds,
|
creds,
|
||||||
3600,
|
3600,
|
||||||
|
|
@ -177,7 +208,10 @@ fn emit_local(receipt: &serde_json::Value) -> Result<()> {
|
||||||
Ok(())
|
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"]
|
let ts = receipt["generated_at_utc"]
|
||||||
.as_str()
|
.as_str()
|
||||||
.unwrap_or("")
|
.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")?;
|
let tmp = tempfile::NamedTempFile::with_suffix(".json")?;
|
||||||
tokio::fs::write(tmp.path(), &payload).await?;
|
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(
|
let (rc, _, _) = run_cmd(
|
||||||
&[
|
&[
|
||||||
"aws",
|
"aws",
|
||||||
|
|
@ -251,7 +288,9 @@ async fn run_cycle(
|
||||||
let _ = emit_local(&receipt);
|
let _ = emit_local(&receipt);
|
||||||
|
|
||||||
let (s3_ok, s3_key) = if !no_s3 && obs.garage.up {
|
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 {
|
} else {
|
||||||
(false, String::new())
|
(false, String::new())
|
||||||
};
|
};
|
||||||
|
|
@ -286,10 +325,7 @@ async fn run_cycle(
|
||||||
warn!("! {}", alert);
|
warn!("! {}", alert);
|
||||||
}
|
}
|
||||||
|
|
||||||
receipt["receipt_hash"]
|
receipt["receipt_hash"].as_str().unwrap_or("").to_string()
|
||||||
.as_str()
|
|
||||||
.unwrap_or("")
|
|
||||||
.to_string()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[tokio::main]
|
#[tokio::main]
|
||||||
|
|
@ -305,7 +341,10 @@ async fn main() -> Result<()> {
|
||||||
tick += 1;
|
tick += 1;
|
||||||
|
|
||||||
if cli.loop_mode {
|
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 {
|
loop {
|
||||||
match tokio::time::timeout(
|
match tokio::time::timeout(
|
||||||
std::time::Duration::from_secs(cli.interval + 300),
|
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;
|
tokio::time::sleep(std::time::Duration::from_secs(cli.interval)).await;
|
||||||
}
|
}
|
||||||
} else {
|
} 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(())
|
Ok(())
|
||||||
|
|
|
||||||
|
|
@ -24,8 +24,14 @@ struct Cli {
|
||||||
|
|
||||||
#[derive(Subcommand, Debug)]
|
#[derive(Subcommand, Debug)]
|
||||||
enum Commands {
|
enum Commands {
|
||||||
Sync { #[arg(long)] since: Option<i64> },
|
Sync {
|
||||||
Watch { #[arg(long, default_value = "60")] interval: u64 },
|
#[arg(long)]
|
||||||
|
since: Option<i64>,
|
||||||
|
},
|
||||||
|
Watch {
|
||||||
|
#[arg(long, default_value = "60")]
|
||||||
|
interval: u64,
|
||||||
|
},
|
||||||
InitSchema,
|
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);
|
info!("opening opencode.db at {:?}", db_path);
|
||||||
let sqlite = Connection::open(db_path)?;
|
let sqlite = Connection::open(db_path)?;
|
||||||
sqlite.busy_timeout(Duration::from_secs(5))?;
|
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);
|
let mut chat_session = normalize_session(sess, &chat_msgs, None);
|
||||||
|
|
||||||
if let Some(ref emb) = embedder {
|
if let Some(ref emb) = embedder {
|
||||||
let text = format!("{} {} {}", sess.title,
|
let text = format!(
|
||||||
|
"{} {} {}",
|
||||||
|
sess.title,
|
||||||
sess.agent.as_deref().unwrap_or(""),
|
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 {
|
if let Ok(v) = emb.embed(&text).await {
|
||||||
chat_session.embedding = Some(v);
|
chat_session.embedding = Some(v);
|
||||||
}
|
}
|
||||||
|
|
@ -156,16 +170,25 @@ async fn cmd_init_schema(dsn: &str) -> Result<()> {
|
||||||
|
|
||||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
struct OpenCodeSession {
|
struct OpenCodeSession {
|
||||||
id: String, project_id: String, parent_id: Option<String>,
|
id: String,
|
||||||
slug: String, directory: String, title: String,
|
project_id: String,
|
||||||
agent: Option<String>, model: Option<String>,
|
parent_id: Option<String>,
|
||||||
time_created: i64, time_updated: i64,
|
slug: String,
|
||||||
tokens_input: i64, tokens_output: i64,
|
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)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
struct OpenCodeMessage {
|
struct OpenCodeMessage {
|
||||||
id: String, session_id: String, time_created: i64,
|
id: String,
|
||||||
|
session_id: String,
|
||||||
|
time_created: i64,
|
||||||
#[serde(rename = "role")]
|
#[serde(rename = "role")]
|
||||||
data_role: String,
|
data_role: String,
|
||||||
#[serde(default)]
|
#[serde(default)]
|
||||||
|
|
@ -194,15 +217,22 @@ fn load_sessions(conn: &Connection) -> Result<Vec<OpenCodeSession>> {
|
||||||
let mut stmt = conn.prepare(
|
let mut stmt = conn.prepare(
|
||||||
"SELECT id, project_id, parent_id, slug, directory, title, \
|
"SELECT id, project_id, parent_id, slug, directory, title, \
|
||||||
agent, model, time_created, time_updated, tokens_input, tokens_output \
|
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| {
|
let rows = stmt.query_map([], |row| {
|
||||||
Ok(OpenCodeSession {
|
Ok(OpenCodeSession {
|
||||||
id: row.get(0)?, project_id: row.get(1)?, parent_id: row.get(2)?,
|
id: row.get(0)?,
|
||||||
slug: row.get(3)?, directory: row.get(4)?, title: row.get(5)?,
|
project_id: row.get(1)?,
|
||||||
agent: row.get(6)?, model: row.get(7)?,
|
parent_id: row.get(2)?,
|
||||||
time_created: row.get(8)?, time_updated: row.get(9)?,
|
slug: row.get(3)?,
|
||||||
tokens_input: row.get(10)?, tokens_output: row.get(11)?,
|
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())
|
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(
|
let mut stmt = conn.prepare(
|
||||||
"SELECT id, project_id, parent_id, slug, directory, title, \
|
"SELECT id, project_id, parent_id, slug, directory, title, \
|
||||||
agent, model, time_created, time_updated, tokens_input, tokens_output \
|
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| {
|
let rows = stmt.query_map([since_ms], |row| {
|
||||||
Ok(OpenCodeSession {
|
Ok(OpenCodeSession {
|
||||||
id: row.get(0)?, project_id: row.get(1)?, parent_id: row.get(2)?,
|
id: row.get(0)?,
|
||||||
slug: row.get(3)?, directory: row.get(4)?, title: row.get(5)?,
|
project_id: row.get(1)?,
|
||||||
agent: row.get(6)?, model: row.get(7)?,
|
parent_id: row.get(2)?,
|
||||||
time_created: row.get(8)?, time_updated: row.get(9)?,
|
slug: row.get(3)?,
|
||||||
tokens_input: row.get(10)?, tokens_output: row.get(11)?,
|
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())
|
rows.collect::<Result<Vec<_>, _>>().map_err(|e| e.into())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn max_session_updated(conn: &Connection) -> Result<Option<i64>> {
|
fn max_session_updated(conn: &Connection) -> Result<Option<i64>> {
|
||||||
conn.query_row("SELECT MAX(time_updated) FROM session", [], |row| row.get(0))
|
conn.query_row("SELECT MAX(time_updated) FROM session", [], |row| {
|
||||||
.optional()
|
row.get(0)
|
||||||
.map_err(|e| e.into())
|
})
|
||||||
|
.optional()
|
||||||
|
.map_err(|e| e.into())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn messages_for_session(conn: &Connection, session_id: &str) -> Result<Vec<OpenCodeMessage>> {
|
fn messages_for_session(conn: &Connection, session_id: &str) -> Result<Vec<OpenCodeMessage>> {
|
||||||
let mut stmt = conn.prepare(
|
let mut stmt = conn.prepare(
|
||||||
"SELECT id, session_id, time_created, data FROM message \
|
"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 rows = stmt.query_map([session_id], |row| {
|
||||||
let data_str: String = row.get(3)?;
|
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 data: serde_json::Value =
|
||||||
let role = data.get("role").and_then(|v| v.as_str()).unwrap_or("unknown").to_string();
|
serde_json::from_str(&data_str).unwrap_or(serde_json::Value::Null);
|
||||||
Ok(OpenCodeMessage { id: row.get(0)?, session_id: row.get(1)?, time_created: row.get(2)?, data_role: role, data })
|
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())
|
rows.collect::<Result<Vec<_>, _>>().map_err(|e| e.into())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn parts_for_message(conn: &Connection, message_id: &str) -> Result<Vec<OpenCodePart>> {
|
fn parts_for_message(conn: &Connection, message_id: &str) -> Result<Vec<OpenCodePart>> {
|
||||||
let mut stmt = conn.prepare(
|
let mut stmt =
|
||||||
"SELECT data FROM part WHERE message_id = ?1 ORDER BY time_created, id"
|
conn.prepare("SELECT data FROM part WHERE message_id = ?1 ORDER BY time_created, id")?;
|
||||||
)?;
|
|
||||||
let rows = stmt.query_map([message_id], |row| {
|
let rows = stmt.query_map([message_id], |row| {
|
||||||
let data_str: String = row.get(0)?;
|
let data_str: String = row.get(0)?;
|
||||||
let part: OpenCodePart = serde_json::from_str(&data_str).unwrap_or(OpenCodePart {
|
let part: OpenCodePart = serde_json::from_str(&data_str).unwrap_or(OpenCodePart {
|
||||||
part_type: "unknown".into(), text: None, tool: None, call_id: None,
|
part_type: "unknown".into(),
|
||||||
input: None, output: None, is_error: None,
|
text: None,
|
||||||
|
tool: None,
|
||||||
|
call_id: None,
|
||||||
|
input: None,
|
||||||
|
output: None,
|
||||||
|
is_error: None,
|
||||||
});
|
});
|
||||||
Ok(part)
|
Ok(part)
|
||||||
})?;
|
})?;
|
||||||
rows.collect::<Result<Vec<_>, _>>().map_err(|e| e.into())
|
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 {
|
ChatSession {
|
||||||
session_id: sess.id.clone(),
|
session_id: sess.id.clone(),
|
||||||
workspace_fingerprint: Some(workspace_fingerprint(&sess.directory)),
|
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 blocks = Vec::new();
|
||||||
let mut text_parts = Vec::new();
|
let mut text_parts = Vec::new();
|
||||||
let mut tool_calls = Vec::new();
|
let mut tool_calls = Vec::new();
|
||||||
|
|
@ -295,23 +357,55 @@ fn normalize_message(msg: &OpenCodeMessage, parts: &[OpenCodePart], index: i32)
|
||||||
"text" => {
|
"text" => {
|
||||||
if let Some(ref t) = part.text {
|
if let Some(ref t) = part.text {
|
||||||
text_parts.push(t.clone());
|
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" => {
|
"reasoning" => {
|
||||||
if let Some(ref t) = part.text {
|
if let Some(ref t) = part.text {
|
||||||
text_parts.push(format!("[reasoning] {}", t));
|
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" => {
|
"tool" => {
|
||||||
let call_id = part.call_id.clone().unwrap_or_default();
|
let call_id = part.call_id.clone().unwrap_or_default();
|
||||||
let tool_name = part.tool.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 });
|
blocks.push(MessageBlock {
|
||||||
tool_calls.push(ToolCall { call_id: call_id.clone(), tool_name, input: part.input.clone().unwrap_or(serde_json::json!({})) });
|
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" => {
|
"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(),
|
role: msg.data_role.clone(),
|
||||||
blocks,
|
blocks,
|
||||||
text_content: text_parts.join("\n"),
|
text_content: text_parts.join("\n"),
|
||||||
token_input: 0, token_output: 0,
|
token_input: 0,
|
||||||
token_cache_creation: 0, token_cache_read: 0,
|
token_output: 0,
|
||||||
|
token_cache_creation: 0,
|
||||||
|
token_cache_read: 0,
|
||||||
tool_calls,
|
tool_calls,
|
||||||
embedding: None,
|
embedding: None,
|
||||||
receipt_hash: None,
|
receipt_hash: None,
|
||||||
|
|
@ -354,19 +450,34 @@ struct Embedder {
|
||||||
impl Embedder {
|
impl Embedder {
|
||||||
fn new() -> Self {
|
fn new() -> Self {
|
||||||
let base = std::env::var("OLLAMA_HOST").unwrap_or_else(|_| "http://localhost:11434".into());
|
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());
|
let model =
|
||||||
Self { client: reqwest::Client::new(), url: format!("{}/api/embeddings", base.trim_end_matches('/')), 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>> {
|
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}))
|
.json(&serde_json::json!({"model": self.model, "prompt": text}))
|
||||||
.send().await.context("embed POST")?;
|
.send()
|
||||||
|
.await
|
||||||
|
.context("embed POST")?;
|
||||||
if !resp.status().is_success() {
|
if !resp.status().is_success() {
|
||||||
anyhow::bail!("embed HTTP {}", resp.status());
|
anyhow::bail!("embed HTTP {}", resp.status());
|
||||||
}
|
}
|
||||||
let json: serde_json::Value = resp.json().await.context("embed JSON")?;
|
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")?;
|
let arr = json
|
||||||
Ok(arr.iter().map(|v| v.as_f64().unwrap_or(0.0) as f32).collect())
|
.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())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue