Phase B: identifier+key on the wire, JCS envelope signing (FRX/0.5), registry binding

This commit is contained in:
George Coles
2026-09-15 05:54:52 -04:00
parent 559c260639
commit d30e612be8
12 changed files with 159 additions and 89 deletions
+48 -17
View File
@@ -284,6 +284,14 @@ impl Node {
self.index.doc_count()
}
pub fn identifier(&self) -> String {
self.config
.node
.id
.clone()
.unwrap_or_else(|| self.key.public_hex())
}
pub fn sent(&self) -> u64 {
self.sent.load(Ordering::SeqCst)
}
@@ -359,7 +367,12 @@ impl Node {
let mut aggregates = self.aggregates.lock().expect("aggregates lock");
*aggregates.sent.entry(current_period()).or_default() += 1;
}
let envelope = Envelope::new(&self.key, TYPE_QUERY, serde_json::to_value(&query)?);
let envelope = Envelope::new(
&self.key,
&self.identifier(),
TYPE_QUERY,
serde_json::to_value(&query)?,
);
self.publish(&envelope).await;
let timeout = Duration::from_millis(timeout_ms.unwrap_or(self.config.query.timeout_ms));
let deadline = Instant::now() + timeout;
@@ -450,19 +463,19 @@ impl Node {
.as_ref()
.map(|signed| authorized_keys(signed, now_ts()));
match authorized {
Some(map) => {
let entry = map.get(&envelope.from);
(entry.map(|(_, class)| class.clone()), entry.is_some())
}
Some(map) => match map.get(&envelope.key) {
Some((id, class)) if id == &envelope.from => (Some(class.clone()), true),
_ => (None, false),
},
None => (None, false),
}
} else {
let members = self.members.read().expect("members lock");
let class = self.member_class(&members, &envelope.from);
let class = self.member_class(&members, &envelope.key);
let listed = if members.is_empty() {
self.config.node.dev_bootstrap
} else {
class.is_some()
class.is_some() && envelope.from == envelope.key
};
(class, listed)
};
@@ -471,7 +484,7 @@ impl Node {
}
match envelope.msg_type.as_str() {
TYPE_QUERY => {
if envelope.from == self.key.public_hex() || !self.config.node.responder {
if envelope.from == self.identifier() || !self.config.node.responder {
return;
}
let Ok(query) = envelope.parse_body::<QueryBody>() else {
@@ -491,7 +504,7 @@ impl Node {
}
let node = self.clone();
let relay = relay.to_string();
let querier = envelope.from.clone();
let querier = envelope.key.clone();
tokio::spawn(async move {
if let Err(error) = node.respond(&query, &querier, &relay).await {
eprintln!("responder failed: {error}");
@@ -535,10 +548,12 @@ impl Node {
}
let node = self.clone();
let period = period.to_string();
let requester = envelope.from.clone();
let requester_id = envelope.from.clone();
let requester_key = envelope.key.clone();
let relay = relay.to_string();
tokio::spawn(async move {
node.serve_aggregate(&period, &requester, &relay).await;
node.serve_aggregate(&period, &requester_id, &requester_key, &relay)
.await;
});
}
_ => {}
@@ -552,7 +567,12 @@ impl Node {
return Ok(());
}
let body = build_response(&query.qid, response_items(&hits), total, max);
let envelope = Envelope::new(&self.key, TYPE_RESPONSE, serde_json::to_value(&body)?);
let envelope = Envelope::new(
&self.key,
&self.identifier(),
TYPE_RESPONSE,
serde_json::to_value(&body)?,
);
self.send_unicast(querier, envelope, relay).await
}
@@ -572,7 +592,12 @@ impl Node {
.first()
.ok_or_else(|| anyhow!("no relays configured"))?
.clone();
let envelope = Envelope::new(&self.key, TYPE_AGGREGATE, json!({ "period": period }));
let envelope = Envelope::new(
&self.key,
&self.identifier(),
TYPE_AGGREGATE,
json!({ "period": period }),
);
self.send_unicast(to, envelope, &relay).await?;
let key = (to.to_string(), period.to_string());
let deadline = Instant::now()
@@ -594,13 +619,19 @@ impl Node {
}
}
async fn serve_aggregate(&self, period: &str, requester: &str, relay: &str) {
let body = self.aggregate_for(period, Some(requester));
async fn serve_aggregate(
&self,
period: &str,
requester_id: &str,
requester_key: &str,
relay: &str,
) {
let body = self.aggregate_for(period, Some(requester_id));
let Ok(value) = serde_json::to_value(&body) else {
return;
};
let envelope = Envelope::new(&self.key, TYPE_AGGREGATE, value);
if let Err(error) = self.send_unicast(requester, envelope, relay).await {
let envelope = Envelope::new(&self.key, &self.identifier(), TYPE_AGGREGATE, value);
if let Err(error) = self.send_unicast(requester_key, envelope, relay).await {
eprintln!("aggregate reply failed: {error}");
}
}