use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use super::networks::NetHub;
use crate::worlds::SharedCell;
const PASS_EVERY: Duration = Duration::from_secs(5);
#[derive(Clone, Default)]
pub struct Relay {
pub sent: Arc<AtomicU64>,
pub pending: Arc<AtomicU64>,
}
fn relay_file() -> std::path::PathBuf {
let home = std::env::var("HOME").unwrap_or_else(|_| ".".into());
std::path::Path::new(&home).join("cyb").join("relayed")
}
fn toggle_file() -> std::path::PathBuf {
let home = std::env::var("HOME").unwrap_or_else(|_| ".".into());
std::path::Path::new(&home).join("cyb").join("relay")
}
fn wanted() -> bool {
std::fs::read_to_string(toggle_file())
.map(|s| s.trim() != "off")
.unwrap_or(true)
}
fn load_marks() -> HashMap<String, u64> {
std::fs::read_to_string(relay_file())
.map(|text| {
text.lines()
.filter_map(|l| {
let (n, s) = l.split_once(' ')?;
Some((n.to_string(), s.trim().parse().ok()?))
})
.collect()
})
.unwrap_or_default()
}
fn save_marks(marks: &HashMap<String, u64>) {
let mut text = String::new();
for (n, s) in marks {
text.push_str(&format!("{n} {s}\n"));
}
let _ = std::fs::write(relay_file(), text);
}
fn hex32(b: &[u8; 32]) -> String {
b.iter().map(|x| format!("{x:02x}")).collect()
}
impl Relay {
pub fn start(shared: SharedCell, hub: NetHub) -> Self {
let relay = Relay::default();
let (sent, pending) = (relay.sent.clone(), relay.pending.clone());
std::thread::Builder::new()
.name("body-relay".into())
.spawn(move || {
let mut marks = load_marks();
let first_run = !relay_file().exists();
let agent = super::networks::agent();
loop {
std::thread::sleep(PASS_EVERY);
if !wanted() {
continue;
}
let Some((net_name, url)) = hub
.states
.lock()
.ok()
.and_then(|v| v.first().map(|n| (n.name.clone(), n.url.clone())))
else {
continue;
};
let fresh: Vec<(String, u64, serde_json::Value)> = {
let cell = shared.cell.lock().expect("shared cell poisoned");
let mut out = Vec::new();
for (neuron, chain) in cell.graph.chains.iter() {
let key = hex32(neuron);
let mark = marks.get(&key).copied().unwrap_or_else(|| {
if first_run {
chain.entries.keys().next_back().copied().unwrap_or(0)
} else {
0
}
});
marks.entry(key.clone()).or_insert(mark);
for (step, sig) in chain.entries.range(mark + 1..) {
for l in &sig.links {
out.push((
key.clone(),
*step,
serde_json::json!({
"neuron": hex32(&l.neuron),
"from": hex32(&l.from),
"to": hex32(&l.to),
"amount": l.amount,
"valence": l.valence,
}),
));
}
}
}
out
};
if first_run && !fresh.is_empty() {
}
if fresh.is_empty() {
if marks_dirty(&marks) {
save_marks(&marks);
}
pending.store(0, Ordering::Relaxed);
continue;
}
let mut failed = 0u64;
for (neuron_key, step, body) in fresh {
let posted = agent
.post(&format!("{url}/v1/link"))
.send_json(&body)
.ok()
.and_then(|mut r| {
r.body_mut().read_json::<serde_json::Value>().ok()
});
match posted {
Some(resp) => {
sent.fetch_add(1, Ordering::Relaxed);
if let (Some(h), Some(root)) = (
resp.get("height").and_then(|v| v.as_u64()),
resp.get("root").and_then(|v| v.as_str()),
) {
hub.note_block(&net_name, h, root);
}
let m = marks.entry(neuron_key).or_insert(0);
if step > *m {
*m = step;
}
}
None => {
failed += 1;
break;
}
}
}
pending.store(failed, Ordering::Relaxed);
save_marks(&marks);
}
})
.expect("spawn body-relay");
relay
}
}
fn marks_dirty(marks: &HashMap<String, u64>) -> bool {
!marks.is_empty() && !relay_file().exists()
}