merge: reuse accepted AI provider setup in publishing worktree
This commit is contained in:
@@ -352,20 +352,19 @@ impl RpcHandler {
|
||||
(Some(u), Some(t)) if t > 0 => {
|
||||
serde_json::json!((u as f64 / t as f64 * 100.0).round())
|
||||
}
|
||||
_ => serde_json::json!(0),
|
||||
_ => serde_json::Value::Null,
|
||||
}
|
||||
};
|
||||
let apps = state.map(|s| s.apps.as_slice()).unwrap_or(&[]);
|
||||
let reported_at = state
|
||||
.map(|s| s.timestamp.clone())
|
||||
.or_else(|| n.last_seen.clone())
|
||||
.unwrap_or_else(|| n.added_at.clone());
|
||||
.or_else(|| n.last_seen.clone());
|
||||
|
||||
let mut report = serde_json::json!({
|
||||
"node_id": n.did,
|
||||
"node_name": state.and_then(|s| s.node_name.clone()).or_else(|| n.name.clone()),
|
||||
"uptime_secs": state.and_then(|s| s.uptime_secs).unwrap_or(0),
|
||||
"cpu_pct": state.and_then(|s| s.cpu_usage_percent).map(|v| v.round()).unwrap_or(0.0),
|
||||
"uptime_secs": state.and_then(|s| s.uptime_secs),
|
||||
"cpu_pct": state.and_then(|s| s.cpu_usage_percent).filter(|v| v.is_finite() && (0.0..=100.0).contains(v)).map(|v| v.round()),
|
||||
"mem_pct": pct(state.and_then(|s| s.mem_used_bytes), state.and_then(|s| s.mem_total_bytes)),
|
||||
"disk_pct": pct(state.and_then(|s| s.disk_used_bytes), state.and_then(|s| s.disk_total_bytes)),
|
||||
"container_count": apps.len(),
|
||||
@@ -561,7 +560,7 @@ fn annotate_fleet_report(report: &mut serde_json::Value) {
|
||||
let is_online = reported
|
||||
.map(|dt| {
|
||||
let age = chrono::Utc::now().signed_duration_since(dt);
|
||||
age.num_minutes() < 30
|
||||
age.num_seconds() >= -60 && age.num_seconds() < 1800
|
||||
})
|
||||
.unwrap_or(false);
|
||||
|
||||
@@ -569,7 +568,9 @@ fn annotate_fleet_report(report: &mut serde_json::Value) {
|
||||
.map(|dt| {
|
||||
let age = chrono::Utc::now().signed_duration_since(dt);
|
||||
let mins = age.num_minutes();
|
||||
if mins < 1 {
|
||||
if age.num_seconds() < -60 {
|
||||
"unknown (clock ahead)".to_string()
|
||||
} else if mins < 1 {
|
||||
"just now".to_string()
|
||||
} else if mins < 60 {
|
||||
format!("{}m ago", mins)
|
||||
|
||||
@@ -602,14 +602,39 @@ impl RpcHandler {
|
||||
None
|
||||
};
|
||||
|
||||
// Reuse the minute collector instead of running expensive probes for
|
||||
// every peer. An absent/stalled collector is unknown, never zero load.
|
||||
let now = chrono::Utc::now().timestamp();
|
||||
let latest = self
|
||||
.metrics_store
|
||||
.latest()
|
||||
.await
|
||||
.filter(|sample| (0..=180).contains(&now.saturating_sub(sample.timestamp)));
|
||||
let metrics = latest.as_ref().map(|sample| &sample.system);
|
||||
let uptime = tokio::fs::read_to_string("/proc/uptime")
|
||||
.await
|
||||
.ok()
|
||||
.and_then(|s| s.split_whitespace().next()?.parse::<f64>().ok())
|
||||
.filter(|v| v.is_finite() && *v >= 0.0)
|
||||
.map(|v| v as u64);
|
||||
let state = federation::build_local_state(
|
||||
apps,
|
||||
0.0,
|
||||
0,
|
||||
0,
|
||||
0,
|
||||
0,
|
||||
0,
|
||||
metrics
|
||||
.map(|m| m.cpu_percent)
|
||||
.filter(|v| v.is_finite() && (0.0..=100.0).contains(v)),
|
||||
metrics
|
||||
.filter(|m| m.mem_total_bytes > 0)
|
||||
.map(|m| m.mem_used_bytes),
|
||||
metrics
|
||||
.filter(|m| m.mem_total_bytes > 0)
|
||||
.map(|m| m.mem_total_bytes),
|
||||
metrics
|
||||
.filter(|m| m.disk_total_bytes > 0)
|
||||
.map(|m| m.disk_used_bytes),
|
||||
metrics
|
||||
.filter(|m| m.disk_total_bytes > 0)
|
||||
.map(|m| m.disk_total_bytes),
|
||||
uptime,
|
||||
tor_active,
|
||||
server_name,
|
||||
nostr_npub,
|
||||
@@ -1254,76 +1279,120 @@ impl RpcHandler {
|
||||
);
|
||||
}
|
||||
|
||||
let reply = self.prepare_peer_approval_reply(&req).await?;
|
||||
// Persist the operator decision before transport. A relay outage must
|
||||
// not require another approval or lose the already-authorized reply.
|
||||
pending::decide(&self.config.data_dir, id, pending::PendingState::Approved).await?;
|
||||
let delivered = self.deliver_peer_approval_reply(&reply).await?;
|
||||
Ok(serde_json::json!({ "approved": true, "id": id, "delivery_pending": !delivered }))
|
||||
}
|
||||
|
||||
async fn prepare_peer_approval_reply(
|
||||
&self,
|
||||
req: &pending::PendingPeerRequest,
|
||||
) -> Result<federation::handshake_delivery::ApprovalReply> {
|
||||
use federation::handshake_delivery::{self, ApprovalReply};
|
||||
if let Some(reply) = handshake_delivery::find(&self.config.data_dir, &req.id).await? {
|
||||
anyhow::ensure!(
|
||||
reply.recipient == req.from_nostr_pubkey && reply.expected_did == req.from_did,
|
||||
"Approval recipient changed"
|
||||
);
|
||||
return Ok(reply);
|
||||
}
|
||||
let (data, _) = self.state_manager.get_snapshot().await;
|
||||
let local_did = identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?;
|
||||
let local_onion = data
|
||||
.server_info
|
||||
.tor_address
|
||||
.clone()
|
||||
.as_deref()
|
||||
.ok_or_else(|| anyhow::anyhow!("Tor address not available"))?;
|
||||
let local_pubkey = data.server_info.pubkey.clone();
|
||||
|
||||
// Generate a one-shot federation invite. The code embeds OUR onion
|
||||
// and OUR pubkey, but it leaves this box only inside the NIP-44
|
||||
// ciphertext below.
|
||||
let identity_dir = self.config.data_dir.join("identity");
|
||||
let local_fips_npub = identity::fips_npub(&identity_dir).await.unwrap_or(None);
|
||||
// Discovery/connection-request approvals admit the requester as
|
||||
// Observer — the invite itself now carries that level, so both
|
||||
// sides converge on Observer without post-hoc demotion.
|
||||
let local_fips_npub = identity::fips_npub(&self.config.data_dir.join("identity"))
|
||||
.await
|
||||
.unwrap_or(None);
|
||||
let invite_code = federation::create_invite(
|
||||
&self.config.data_dir,
|
||||
&local_did,
|
||||
&local_onion,
|
||||
&local_pubkey,
|
||||
local_onion,
|
||||
&data.server_info.pubkey,
|
||||
local_fips_npub.as_deref(),
|
||||
TrustLevel::Observer,
|
||||
)
|
||||
.await?;
|
||||
handshake_delivery::stage(
|
||||
&self.config.data_dir,
|
||||
ApprovalReply {
|
||||
request_id: req.id.clone(),
|
||||
recipient: req.from_nostr_pubkey.clone(),
|
||||
expected_did: req.from_did.clone(),
|
||||
invite_code,
|
||||
attempts: 0,
|
||||
next_attempt: 0,
|
||||
},
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
// Pre-add the requester to OUR federation list as Observer so that
|
||||
// when their `federation.peer-joined` callback arrives over Tor we
|
||||
// already trust their pubkey enough to accept the join. Their DID
|
||||
// and pubkey come from the request — we'll cross-check the pubkey
|
||||
// against the eventual peer-joined signature in the existing
|
||||
// verification path (handlers.rs line ~365).
|
||||
if !req.from_did.is_empty() {
|
||||
// We don't know the requester's onion or ed25519 pubkey yet —
|
||||
// they'll send those in the federation.peer-joined callback
|
||||
// after they apply our invite. Until then we can't add a real
|
||||
// FederatedNode entry. We just store the pending row as
|
||||
// Approved so the UI shows progress, and trust the existing
|
||||
// peer-joined handler to admit them as Observer when they call.
|
||||
//
|
||||
// Caveat: peer-joined currently hardcodes TrustLevel::Trusted.
|
||||
// We override that below by demoting on success.
|
||||
debug!(
|
||||
requester_did = %req.from_did,
|
||||
"Approval pending — waiting for federation.peer-joined callback over Tor"
|
||||
);
|
||||
}
|
||||
|
||||
// Encrypt + send the invite over NIP-44 to the requester.
|
||||
let identity_dir = self.config.data_dir.join("identity");
|
||||
nostr_handshake::send_peer_invite(
|
||||
&identity_dir,
|
||||
&req.from_nostr_pubkey,
|
||||
&invite_code,
|
||||
&self.config.nostr_relays,
|
||||
async fn deliver_peer_approval_reply(
|
||||
&self,
|
||||
reply: &federation::handshake_delivery::ApprovalReply,
|
||||
) -> Result<bool> {
|
||||
let Some(claimed) = federation::handshake_delivery::claim(
|
||||
&self.config.data_dir,
|
||||
&reply.request_id,
|
||||
chrono::Utc::now().timestamp(),
|
||||
)
|
||||
.await?
|
||||
else {
|
||||
return Ok(false);
|
||||
};
|
||||
let result = nostr_handshake::send_peer_invite(
|
||||
&self.config.data_dir.join("identity"),
|
||||
&claimed.recipient,
|
||||
&claimed.invite_code,
|
||||
&self.handshake_relays().await,
|
||||
self.config.nostr_tor_proxy.as_deref(),
|
||||
)
|
||||
.await?;
|
||||
.await;
|
||||
if result.is_err() {
|
||||
warn!(request_id = %reply.request_id, "Peer approval delivery deferred; durable retry scheduled");
|
||||
}
|
||||
Ok(result.is_ok())
|
||||
}
|
||||
|
||||
pending::set_state(&self.config.data_dir, id, pending::PendingState::Approved).await?;
|
||||
info!(
|
||||
id = %id,
|
||||
from = %req.from_nostr_pubkey,
|
||||
"Approved peer request and shipped invite over NIP-44"
|
||||
);
|
||||
Ok(serde_json::json!({
|
||||
"approved": true,
|
||||
"id": id,
|
||||
}))
|
||||
/// Recover relay loss and legacy approvals without changing trust or
|
||||
/// resurrecting a node that the operator explicitly removed.
|
||||
pub(in crate::api::rpc) async fn retry_peer_approval_replies(&self) -> Result<()> {
|
||||
let requests = pending::load_pending(&self.config.data_dir).await?;
|
||||
let nodes = federation::load_nodes(&self.config.data_dir).await?;
|
||||
let removed = federation::load_removed_dids(&self.config.data_dir).await?;
|
||||
let cutoff = chrono::Utc::now() - chrono::Duration::days(30);
|
||||
let mut sent = 0;
|
||||
for req in requests {
|
||||
if req.outbound || req.state != pending::PendingState::Approved {
|
||||
federation::handshake_delivery::remove(&self.config.data_dir, &req.id).await?;
|
||||
continue;
|
||||
}
|
||||
let expired = chrono::DateTime::parse_from_rfc3339(&req.received_at)
|
||||
.map(|time| time < cutoff)
|
||||
.unwrap_or(true);
|
||||
if expired
|
||||
|| removed.contains(&req.from_did)
|
||||
|| nodes.iter().any(|node| node.did == req.from_did)
|
||||
{
|
||||
federation::handshake_delivery::remove(&self.config.data_dir, &req.id).await?;
|
||||
continue;
|
||||
}
|
||||
if req.from_did.is_empty() || sent >= 4 {
|
||||
continue;
|
||||
}
|
||||
let reply = self.prepare_peer_approval_reply(&req).await?;
|
||||
if reply.next_attempt > chrono::Utc::now().timestamp() {
|
||||
continue;
|
||||
}
|
||||
self.deliver_peer_approval_reply(&reply).await?;
|
||||
sent += 1;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// federation.reject-request — drop a pending request and, if requested,
|
||||
@@ -1353,19 +1422,19 @@ impl RpcHandler {
|
||||
);
|
||||
}
|
||||
|
||||
pending::decide(&self.config.data_dir, id, pending::PendingState::Rejected).await?;
|
||||
if notify {
|
||||
let identity_dir = self.config.data_dir.join("identity");
|
||||
let _ = nostr_handshake::send_peer_reject(
|
||||
&identity_dir,
|
||||
&req.from_nostr_pubkey,
|
||||
reason,
|
||||
&self.config.nostr_relays,
|
||||
&self.handshake_relays().await,
|
||||
self.config.nostr_tor_proxy.as_deref(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
pending::set_state(&self.config.data_dir, id, pending::PendingState::Rejected).await?;
|
||||
info!(id = %id, from = %req.from_nostr_pubkey, "Rejected peer request");
|
||||
Ok(serde_json::json!({ "rejected": true, "id": id }))
|
||||
}
|
||||
@@ -1410,7 +1479,7 @@ impl RpcHandler {
|
||||
&identity_dir,
|
||||
&req.from_nostr_pubkey,
|
||||
reason,
|
||||
&self.config.nostr_relays,
|
||||
&self.handshake_relays().await,
|
||||
self.config.nostr_tor_proxy.as_deref(),
|
||||
)
|
||||
.await
|
||||
|
||||
@@ -0,0 +1,276 @@
|
||||
//! Exercise encrypted replies through a relay configured in the UI only.
|
||||
use crate::federation::pending::{self, PendingState};
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use nostr_sdk::prelude::{nip44, Event, Keys};
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
#[tokio::test]
|
||||
async fn managed_relay_receives_approval_rejection_and_cancellation() {
|
||||
for (operation, accepted) in [
|
||||
("approve", true),
|
||||
("reject", true),
|
||||
("cancel", true),
|
||||
("approve", false),
|
||||
("retry", true),
|
||||
] {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
|
||||
let relay_url = format!("ws://{}", listener.local_addr().unwrap());
|
||||
let relay = tokio::spawn(async move {
|
||||
let (socket, _) = listener.accept().await.unwrap();
|
||||
let mut ws = tokio_tungstenite::accept_async(socket).await.unwrap();
|
||||
while let Some(Ok(message)) = ws.next().await {
|
||||
if !message.is_text() {
|
||||
continue;
|
||||
}
|
||||
let value: serde_json::Value =
|
||||
serde_json::from_str(message.to_text().unwrap()).unwrap();
|
||||
if value[0] != "EVENT" {
|
||||
continue;
|
||||
}
|
||||
let event: Event = serde_json::from_value(value[1].clone()).unwrap();
|
||||
event.verify().unwrap();
|
||||
ws.send(tokio_tungstenite::tungstenite::Message::Text(
|
||||
serde_json::json!([
|
||||
"OK",
|
||||
event.id.to_hex(),
|
||||
accepted,
|
||||
"blocked: fixture rejection"
|
||||
])
|
||||
.to_string(),
|
||||
))
|
||||
.await
|
||||
.unwrap();
|
||||
return event;
|
||||
}
|
||||
panic!("relay closed without a signed event");
|
||||
});
|
||||
let mut config = crate::config::Config::default();
|
||||
config.data_dir = tmp.path().to_path_buf();
|
||||
config.nostr_relays.clear();
|
||||
config.nostr_tor_proxy = None;
|
||||
crate::nostr_relays::save_relays(
|
||||
tmp.path(),
|
||||
&crate::nostr_relays::RelayStore {
|
||||
relays: vec![crate::nostr_relays::RelayConfig {
|
||||
url: relay_url,
|
||||
enabled: true,
|
||||
added_at: chrono::Utc::now().to_rfc3339(),
|
||||
}],
|
||||
},
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let sender = Keys::parse(&"11".repeat(32)).unwrap();
|
||||
let recipient = Keys::parse(&"22".repeat(32)).unwrap();
|
||||
let identity_dir = tmp.path().join("identity");
|
||||
tokio::fs::create_dir_all(&identity_dir).await.unwrap();
|
||||
tokio::fs::write(identity_dir.join("nostr_secret"), "11".repeat(32))
|
||||
.await
|
||||
.unwrap();
|
||||
tokio::fs::write(identity_dir.join("node_key"), [0x33; 32])
|
||||
.await
|
||||
.unwrap();
|
||||
let state = Arc::new(crate::state::StateManager::new());
|
||||
state
|
||||
.mutate_data(|data| {
|
||||
data.server_info.pubkey = "33".repeat(32);
|
||||
data.server_info.tor_address = Some(format!("{}.onion", "a".repeat(56)));
|
||||
})
|
||||
.await;
|
||||
let handler = crate::api::rpc::RpcHandler::new(
|
||||
config,
|
||||
state,
|
||||
Arc::new(crate::monitoring::MetricsStore::new()),
|
||||
crate::session::SessionStore::new_for_tests(tmp.path().join("sessions.json")),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let row = if operation == "cancel" {
|
||||
pending::insert_outbound(
|
||||
tmp.path(),
|
||||
recipient.public_key().to_hex(),
|
||||
String::new(),
|
||||
crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
} else {
|
||||
pending::insert_inbound(
|
||||
tmp.path(),
|
||||
recipient.public_key().to_hex(),
|
||||
String::new(),
|
||||
crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
};
|
||||
if operation == "retry" {
|
||||
pending::decide(tmp.path(), &row.id, PendingState::Approved)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
let params = Some(serde_json::json!({"id": row.id, "notify": true}));
|
||||
let action = async {
|
||||
match operation {
|
||||
"approve" => handler.handle_federation_approve_request(params).await,
|
||||
"reject" => handler.handle_federation_reject_request(params).await,
|
||||
"retry" => handler
|
||||
.retry_peer_approval_replies()
|
||||
.await
|
||||
.map(|_| serde_json::json!({"ok": true})),
|
||||
_ => handler.handle_federation_cancel_request(params).await,
|
||||
}
|
||||
};
|
||||
let outcome = tokio::time::timeout(Duration::from_secs(20), action)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(outcome.is_ok(), accepted || operation == "approve");
|
||||
let event = tokio::time::timeout(Duration::from_secs(5), relay)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(event.pubkey, sender.public_key());
|
||||
let plaintext =
|
||||
nip44::decrypt(recipient.secret_key(), &event.pubkey, &event.content).unwrap();
|
||||
let message: serde_json::Value = serde_json::from_str(&plaintext).unwrap();
|
||||
let expected = match operation {
|
||||
"approve" | "retry" => "peer-invite",
|
||||
"reject" => "peer-reject",
|
||||
_ => "peer-cancel",
|
||||
};
|
||||
assert_eq!(message["type"], expected);
|
||||
let saved = pending::find_by_id(tmp.path(), &row.id).await.unwrap();
|
||||
if !accepted {
|
||||
assert_eq!(
|
||||
saved.unwrap().state,
|
||||
if operation == "approve" {
|
||||
PendingState::Approved
|
||||
} else {
|
||||
PendingState::Pending
|
||||
}
|
||||
);
|
||||
if operation == "approve" {
|
||||
assert_eq!(outcome.unwrap()["delivery_pending"], true);
|
||||
let durable = crate::federation::handshake_delivery::find(tmp.path(), &row.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(durable.recipient, recipient.public_key().to_hex());
|
||||
assert_eq!(durable.attempts, 1);
|
||||
}
|
||||
assert!(crate::federation::load_nodes(tmp.path())
|
||||
.await
|
||||
.unwrap()
|
||||
.is_empty());
|
||||
continue;
|
||||
}
|
||||
match operation {
|
||||
"approve" | "retry" => {
|
||||
assert_eq!(saved.unwrap().state, PendingState::Approved);
|
||||
let first = crate::federation::handshake_delivery::find(tmp.path(), &row.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(first.attempts, 1);
|
||||
// Concurrent/background polling honors the persisted backoff.
|
||||
handler.retry_peer_approval_replies().await.unwrap();
|
||||
let second = crate::federation::handshake_delivery::find(tmp.path(), &row.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(second.attempts, 1);
|
||||
assert_eq!(first.invite_code, second.invite_code);
|
||||
let invite =
|
||||
crate::federation::parse_invite(message["invite_code"].as_str().unwrap())
|
||||
.unwrap();
|
||||
assert_eq!(invite.trust_level, crate::federation::TrustLevel::Observer);
|
||||
assert!(!event.content.contains(".onion"));
|
||||
}
|
||||
"reject" => assert_eq!(saved.unwrap().state, PendingState::Rejected),
|
||||
_ => assert!(saved.is_none()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn federation_metrics_are_collected_values_or_unknown_never_placeholders() {
|
||||
for age in [None, Some(0), Some(181), Some(-120)] {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let mut config = crate::config::Config::default();
|
||||
config.data_dir = tmp.path().to_path_buf();
|
||||
let metrics = Arc::new(crate::monitoring::MetricsStore::new());
|
||||
if let Some(age) = age {
|
||||
metrics
|
||||
.push(
|
||||
serde_json::from_value(serde_json::json!({
|
||||
"timestamp": chrono::Utc::now().timestamp() - age,
|
||||
"system": {"cpu_percent": 37.5, "mem_used_bytes": 200,
|
||||
"mem_total_bytes": 800, "disk_used_bytes": 600,
|
||||
"disk_total_bytes": 1000, "net_rx_bytes": 0, "net_tx_bytes": 0,
|
||||
"load_avg_1": 0.0, "load_avg_5": 0.0, "load_avg_15": 0.0},
|
||||
"containers": [], "rpc_latency_ms": 0.0, "ws_connections": 0
|
||||
}))
|
||||
.unwrap(),
|
||||
)
|
||||
.await;
|
||||
}
|
||||
let handler = crate::api::rpc::RpcHandler::new(
|
||||
config,
|
||||
Arc::new(crate::state::StateManager::new()),
|
||||
metrics,
|
||||
crate::session::SessionStore::new_for_tests(tmp.path().join("sessions.json")),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let snapshot = handler.handle_federation_get_state().await.unwrap();
|
||||
if age == Some(0) {
|
||||
assert_eq!(snapshot["cpu_usage_percent"], 37.5);
|
||||
assert_eq!(snapshot["mem_used_bytes"], 200);
|
||||
assert_eq!(snapshot["disk_total_bytes"], 1000);
|
||||
} else {
|
||||
for field in [
|
||||
"cpu_usage_percent",
|
||||
"mem_used_bytes",
|
||||
"mem_total_bytes",
|
||||
"disk_used_bytes",
|
||||
"disk_total_bytes",
|
||||
] {
|
||||
assert!(
|
||||
snapshot.get(field).is_none_or(|v| v.is_null()),
|
||||
"{age:?}: {field}"
|
||||
);
|
||||
}
|
||||
}
|
||||
let peer = serde_json::from_value(serde_json::json!({
|
||||
"did": "did:key:test", "pubkey": "11".repeat(32), "onion": "test.onion",
|
||||
"trust_level": "trusted", "added_at": chrono::Utc::now().to_rfc3339(),
|
||||
"last_state": snapshot
|
||||
}))
|
||||
.unwrap();
|
||||
crate::federation::save_nodes(tmp.path(), &[peer])
|
||||
.await
|
||||
.unwrap();
|
||||
let fleet = handler.handle_telemetry_fleet_status().await.unwrap();
|
||||
let report = &fleet["nodes"][0];
|
||||
if age == Some(0) {
|
||||
assert_eq!(report["cpu_pct"], 38.0);
|
||||
assert_eq!(report["mem_pct"], 25.0);
|
||||
assert_eq!(report["disk_pct"], 60.0);
|
||||
} else {
|
||||
for field in ["cpu_pct", "mem_pct", "disk_pct"] {
|
||||
assert!(report[field].is_null(), "{age:?}: {field}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,6 @@
|
||||
mod handlers;
|
||||
#[cfg(test)]
|
||||
mod handshake_tests;
|
||||
|
||||
use anyhow::Result;
|
||||
|
||||
|
||||
@@ -276,6 +276,11 @@ impl RpcHandler {
|
||||
}
|
||||
Err(e) => tracing::debug!("background handshake poll failed: {e:#}"),
|
||||
}
|
||||
if load_discovery_state(&self.config.data_dir).await.enabled {
|
||||
if let Err(error) = self.retry_peer_approval_replies().await {
|
||||
tracing::warn!("Peer approval retry could not complete: {error:#}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) async fn handle_handshake_poll(&self) -> Result<serde_json::Value> {
|
||||
@@ -356,6 +361,16 @@ impl RpcHandler {
|
||||
);
|
||||
continue;
|
||||
};
|
||||
let scoped_invite = match crate::federation::restrict_discovery_invite(
|
||||
invite_code,
|
||||
&row.from_did,
|
||||
) {
|
||||
Ok(code) => code,
|
||||
Err(_) => {
|
||||
tracing::warn!("Rejected peer invite with mismatched identity");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let row_id = row.id.clone();
|
||||
let (data, _) = self.state_manager.get_snapshot().await;
|
||||
let local_did =
|
||||
@@ -373,7 +388,7 @@ impl RpcHandler {
|
||||
let local_name = data.server_info.name.clone();
|
||||
match crate::federation::accept_invite(
|
||||
&self.config.data_dir,
|
||||
invite_code,
|
||||
&scoped_invite,
|
||||
&local_did,
|
||||
&local_onion,
|
||||
&local_pubkey,
|
||||
|
||||
@@ -1025,10 +1025,46 @@ impl RpcHandler {
|
||||
.ok_or_else(|| anyhow::anyhow!("Missing key"))?;
|
||||
|
||||
match key {
|
||||
"claude_api_key_set" => {
|
||||
let key_file = self.config.data_dir.join("secrets/claude-api-key");
|
||||
let has_key = tokio::fs::metadata(&key_file).await.is_ok();
|
||||
Ok(serde_json::json!({ "value": has_key }))
|
||||
"claude_api_key_set" | "openai_api_key_set" => {
|
||||
let provider = if key == "claude_api_key_set" {
|
||||
"claude"
|
||||
} else {
|
||||
"openai"
|
||||
};
|
||||
Ok(
|
||||
serde_json::json!({ "value": crate::settings::model_provider::has_key(&self.config.data_dir, provider).await }),
|
||||
)
|
||||
}
|
||||
"ai_provider" => {
|
||||
let settings =
|
||||
crate::settings::model_provider::ModelProvider::load(&self.config.data_dir)
|
||||
.await?;
|
||||
Ok(serde_json::json!({ "value": settings }))
|
||||
}
|
||||
"ai_provider_status" => {
|
||||
let settings =
|
||||
crate::settings::model_provider::ModelProvider::load(&self.config.data_dir)
|
||||
.await?;
|
||||
let local = tokio::time::timeout(std::time::Duration::from_secs(4), async {
|
||||
let (detected, _) = crate::api::rpc::mesh::assistant::detect_ollama().await;
|
||||
detected
|
||||
&& crate::assistant::backends::ollama::model_supports_tools(
|
||||
crate::assistant::backends::ollama::OLLAMA_BASE_URL,
|
||||
crate::assistant::backends::ollama::OLLAMA_DEFAULT_MODEL,
|
||||
)
|
||||
.await
|
||||
});
|
||||
let (claude, openai, local) = tokio::join!(
|
||||
crate::settings::model_provider::has_key(&self.config.data_dir, "claude"),
|
||||
crate::settings::model_provider::has_key(&self.config.data_dir, "openai"),
|
||||
local,
|
||||
);
|
||||
let budget = crate::assistant::AssistantBudget::load(&self.config.data_dir).await;
|
||||
Ok(serde_json::json!({ "value": {
|
||||
"schema": 1, "settings": settings, "claude_configured": claude,
|
||||
"openai_configured": openai, "local_ready": local.ok(),
|
||||
"routstr_remaining_sats": budget.remaining_sats(),
|
||||
}}))
|
||||
}
|
||||
_ => Ok(serde_json::json!({ "value": null })),
|
||||
}
|
||||
@@ -1210,38 +1246,21 @@ impl RpcHandler {
|
||||
let value = params.get("value").and_then(|v| v.as_str()).unwrap_or("");
|
||||
|
||||
match key {
|
||||
"claude_api_key" => {
|
||||
let secrets_dir = self.config.data_dir.join("secrets");
|
||||
tokio::fs::create_dir_all(&secrets_dir)
|
||||
.await
|
||||
.context("Failed to create secrets dir")?;
|
||||
let key_file = secrets_dir.join("claude-api-key");
|
||||
|
||||
if value.is_empty() {
|
||||
// Remove key
|
||||
tokio::fs::remove_file(&key_file).await.ok();
|
||||
info!("Claude API key removed");
|
||||
"claude_api_key" | "openai_api_key" => {
|
||||
let provider = if key == "claude_api_key" {
|
||||
"claude"
|
||||
} else {
|
||||
// Save key
|
||||
tokio::fs::write(&key_file, value)
|
||||
.await
|
||||
.context("Failed to write API key")?;
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
std::fs::set_permissions(&key_file, std::fs::Permissions::from_mode(0o600))
|
||||
.ok();
|
||||
}
|
||||
info!("Claude API key saved");
|
||||
}
|
||||
|
||||
// `secrets/claude-api-key` (above) is deliberately the ONLY
|
||||
// Claude key ledger on this node (13-02-PLAN.md). A second
|
||||
// copy used to be written alongside it for a standalone,
|
||||
// unauthenticated sidecar process on port 3142 — that
|
||||
// sidecar and its key copy are retired; the session-gated
|
||||
// Rust daemon reads this one file directly.
|
||||
|
||||
"openai"
|
||||
};
|
||||
crate::settings::model_provider::save_key(&self.config.data_dir, provider, value)
|
||||
.await?;
|
||||
info!(provider, "AI provider credential updated");
|
||||
Ok(serde_json::json!({ "saved": true }))
|
||||
}
|
||||
"ai_provider" => {
|
||||
let settings: crate::settings::model_provider::ModelProvider =
|
||||
serde_json::from_str(value).context("Invalid AI provider settings")?;
|
||||
settings.save(&self.config.data_dir).await?;
|
||||
Ok(serde_json::json!({ "saved": true }))
|
||||
}
|
||||
_ => anyhow::bail!("Unknown setting: {}", key),
|
||||
|
||||
@@ -11,6 +11,7 @@ use crate::api::rpc::RpcHandler;
|
||||
|
||||
pub mod claude;
|
||||
pub mod ollama;
|
||||
pub mod openai;
|
||||
pub mod routstr;
|
||||
#[cfg(test)]
|
||||
pub mod scripted;
|
||||
@@ -39,6 +40,8 @@ pub trait Backend: Send + Sync {
|
||||
pub enum BackendId {
|
||||
Ollama,
|
||||
Claude,
|
||||
Openai,
|
||||
Unavailable,
|
||||
/// 13-13: the third D-04 leg. Not currently returned as the "primary"
|
||||
/// id by `select_backend` (mirroring the existing convention that the
|
||||
/// returned id names the primary attempt, not necessarily which leg of
|
||||
@@ -52,11 +55,23 @@ impl std::fmt::Display for BackendId {
|
||||
match self {
|
||||
BackendId::Ollama => write!(f, "ollama"),
|
||||
BackendId::Claude => write!(f, "claude"),
|
||||
BackendId::Openai => write!(f, "openai"),
|
||||
BackendId::Unavailable => write!(f, "unavailable"),
|
||||
BackendId::Routstr => write!(f, "routstr"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
struct InvalidProviderSettings;
|
||||
#[async_trait]
|
||||
impl Backend for InvalidProviderSettings {
|
||||
async fn send(&self, _: &str, _: &[ToolDef], _: &[ChatMessage]) -> Result<BackendTurn> {
|
||||
anyhow::bail!(
|
||||
"AI connection settings could not be loaded. Review them before sending a message."
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// D-04's per-call fallback: try `primary`'s `send()`, and on a transport
|
||||
/// error fall through to `secondary` for that SAME call rather than
|
||||
/// failing the whole turn — a local model that answers earlier turns and
|
||||
@@ -115,6 +130,54 @@ fn ollama_is_selectable(detected: bool, tool_capable: bool) -> bool {
|
||||
/// tools-free degrade.
|
||||
pub async fn select_backend(handler: &RpcHandler) -> (Box<dyn Backend>, BackendId) {
|
||||
let data_dir = handler.data_dir();
|
||||
// An explicit provider is a privacy and billing choice. Never silently
|
||||
// fall through to another provider if its credentials or network fail.
|
||||
match crate::settings::model_provider::ModelProvider::load(data_dir).await {
|
||||
Ok(settings) => match settings.provider {
|
||||
crate::settings::model_provider::Provider::Openai => {
|
||||
return (
|
||||
Box::new(openai::OpenaiBackend::new(
|
||||
data_dir.to_path_buf(),
|
||||
settings.openai_model,
|
||||
)),
|
||||
BackendId::Openai,
|
||||
)
|
||||
}
|
||||
crate::settings::model_provider::Provider::Claude => {
|
||||
return (
|
||||
Box::new(claude::ClaudeBackend::new(data_dir.to_path_buf())),
|
||||
BackendId::Claude,
|
||||
)
|
||||
}
|
||||
crate::settings::model_provider::Provider::Local => {
|
||||
return (
|
||||
Box::new(ollama::OllamaBackend::new(
|
||||
ollama::OLLAMA_BASE_URL.to_string(),
|
||||
ollama::OLLAMA_DEFAULT_MODEL.to_string(),
|
||||
)),
|
||||
BackendId::Ollama,
|
||||
)
|
||||
}
|
||||
crate::settings::model_provider::Provider::Routstr => {
|
||||
let budget = crate::assistant::AssistantBudget::load(data_dir).await;
|
||||
let mints = crate::wallet::ecash::load_accepted_mints(data_dir)
|
||||
.await
|
||||
.map(|m| m.mints)
|
||||
.unwrap_or_default();
|
||||
return (
|
||||
Box::new(routstr::RoutstrBackend::new(
|
||||
data_dir.to_path_buf(),
|
||||
budget.payment_policy(),
|
||||
mints,
|
||||
handler.nostr_tor_proxy(),
|
||||
)),
|
||||
BackendId::Routstr,
|
||||
);
|
||||
}
|
||||
crate::settings::model_provider::Provider::Auto => {}
|
||||
},
|
||||
Err(_) => return (Box::new(InvalidProviderSettings), BackendId::Unavailable),
|
||||
}
|
||||
let (detected, _models) = crate::api::rpc::mesh::assistant::detect_ollama().await;
|
||||
let model = ollama::OLLAMA_DEFAULT_MODEL;
|
||||
let tool_capable = if detected {
|
||||
|
||||
@@ -0,0 +1,279 @@
|
||||
//! Explicit OpenAI API selection using the shared tool loop and egress policy.
|
||||
//! Keys stay node-side; no redirects, automatic retries, or provider fallback.
|
||||
use super::{Backend, BackendTurn};
|
||||
use crate::assistant::{
|
||||
egress::{self, EgressVerdict},
|
||||
tools::{ChatMessage, ToolCall, ToolDef},
|
||||
};
|
||||
use anyhow::{Context, Result};
|
||||
use async_trait::async_trait;
|
||||
use serde_json::{json, Value};
|
||||
use std::{path::PathBuf, time::Duration};
|
||||
|
||||
const URL: &str = "https://api.openai.com/v1/chat/completions";
|
||||
const RESPONSE_LIMIT: usize = 2 * 1024 * 1024;
|
||||
pub struct OpenaiBackend {
|
||||
data_dir: PathBuf,
|
||||
model: String,
|
||||
}
|
||||
impl OpenaiBackend {
|
||||
pub fn new(data_dir: PathBuf, model: String) -> Self {
|
||||
Self { data_dir, model }
|
||||
}
|
||||
async fn send_at(
|
||||
&self,
|
||||
url: &str,
|
||||
system: &str,
|
||||
tools: &[ToolDef],
|
||||
history: &[ChatMessage],
|
||||
) -> Result<BackendTurn> {
|
||||
let key = tokio::fs::read_to_string(self.data_dir.join("secrets/openai-api-key"))
|
||||
.await
|
||||
.map_err(|_| {
|
||||
anyhow::anyhow!("OpenAI API key is not configured. Open AI connection settings.")
|
||||
})?;
|
||||
anyhow::ensure!(!key.trim().is_empty(), "OpenAI API key is not configured");
|
||||
anyhow::ensure!(
|
||||
!self.model.is_empty(),
|
||||
"Choose an OpenAI model in AI connection settings"
|
||||
);
|
||||
let mut messages = vec![json!({"role": "system", "content": system})];
|
||||
messages.extend(history.iter().flat_map(super::routstr::message_to_wire));
|
||||
let mut body = json!({"model": self.model, "messages": messages, "stream": false,
|
||||
"store": false, "max_completion_tokens": 2048, "n": 1});
|
||||
if !tools.is_empty() {
|
||||
body["tools"] = json!(tools.iter().map(|tool| json!({"type": "function", "function": {
|
||||
"name": tool.name, "description": tool.description, "parameters": tool.parameters,
|
||||
}})).collect::<Vec<_>>());
|
||||
body["parallel_tool_calls"] = json!(false);
|
||||
}
|
||||
let context = egress::EgressContext::from_turn(
|
||||
history,
|
||||
&tools.iter().map(|tool| tool.name).collect::<Vec<_>>(),
|
||||
&self.data_dir.join("secrets"),
|
||||
)
|
||||
.await;
|
||||
match egress::screen_outbound(&body.to_string(), &context) {
|
||||
EgressVerdict::Allow => {}
|
||||
EgressVerdict::Truncate(value) => {
|
||||
body = serde_json::from_str(&value)
|
||||
.context("Could not apply outbound privacy filter")?;
|
||||
}
|
||||
EgressVerdict::BlockFallBackLocal => {
|
||||
crate::assistant::global_counters().note_blocked_egress();
|
||||
anyhow::bail!("This message contains private key or recovery material and was not sent to OpenAI");
|
||||
}
|
||||
}
|
||||
let client = reqwest::Client::builder()
|
||||
.timeout(Duration::from_secs(180))
|
||||
.connect_timeout(Duration::from_secs(15))
|
||||
.redirect(reqwest::redirect::Policy::none())
|
||||
.build()?;
|
||||
let mut response = client.post(url).bearer_auth(key.trim()).json(&body).send().await
|
||||
.map_err(|_| anyhow::anyhow!("OpenAI is temporarily unreachable. Your request was not retried automatically."))?;
|
||||
if !response.status().is_success() {
|
||||
anyhow::bail!("{}", error_message(response.status().as_u16()));
|
||||
}
|
||||
let mut bytes = Vec::new();
|
||||
while let Some(chunk) = response
|
||||
.chunk()
|
||||
.await
|
||||
.context("OpenAI response interrupted")?
|
||||
{
|
||||
anyhow::ensure!(
|
||||
bytes.len().saturating_add(chunk.len()) <= RESPONSE_LIMIT,
|
||||
"OpenAI response exceeded the size limit"
|
||||
);
|
||||
bytes.extend_from_slice(&chunk);
|
||||
}
|
||||
parse_response(
|
||||
&serde_json::from_slice(&bytes).context("OpenAI returned an invalid response")?,
|
||||
)
|
||||
}
|
||||
}
|
||||
#[async_trait]
|
||||
impl Backend for OpenaiBackend {
|
||||
async fn send(
|
||||
&self,
|
||||
system: &str,
|
||||
tools: &[ToolDef],
|
||||
history: &[ChatMessage],
|
||||
) -> Result<BackendTurn> {
|
||||
self.send_at(URL, system, tools, history).await
|
||||
}
|
||||
}
|
||||
fn error_message(status: u16) -> &'static str {
|
||||
match status {
|
||||
401 | 403 => "OpenAI rejected the API key or project access. Check AI connection settings.",
|
||||
404 => "This OpenAI model is unavailable for your account. Choose another model in AI connection settings.",
|
||||
429 => "OpenAI usage or rate limit reached. Check your API billing and retry later.",
|
||||
500..=599 => "OpenAI is temporarily unavailable. Retry later.",
|
||||
_ => "OpenAI rejected the request. Check the selected model and retry.",
|
||||
}
|
||||
}
|
||||
fn parse_response(value: &Value) -> Result<BackendTurn> {
|
||||
let choice = value["choices"]
|
||||
.as_array()
|
||||
.and_then(|items| items.first())
|
||||
.context("OpenAI returned no answer")?;
|
||||
anyhow::ensure!(
|
||||
choice["finish_reason"] != "length",
|
||||
"OpenAI reached the response limit. Try a shorter request."
|
||||
);
|
||||
let message = &choice["message"];
|
||||
if let Some(calls) = message["tool_calls"]
|
||||
.as_array()
|
||||
.filter(|calls| !calls.is_empty())
|
||||
{
|
||||
let mut parsed = Vec::new();
|
||||
for call in calls {
|
||||
anyhow::ensure!(
|
||||
call["type"] == "function",
|
||||
"Unsupported OpenAI tool response"
|
||||
);
|
||||
let id = call["id"]
|
||||
.as_str()
|
||||
.filter(|id| !id.is_empty())
|
||||
.context("Missing OpenAI tool call ID")?;
|
||||
let name = call["function"]["name"]
|
||||
.as_str()
|
||||
.filter(|name| !name.is_empty())
|
||||
.context("Missing OpenAI tool name")?;
|
||||
let arguments: Value = serde_json::from_str(
|
||||
call["function"]["arguments"]
|
||||
.as_str()
|
||||
.context("Invalid OpenAI tool arguments")?,
|
||||
)
|
||||
.context("Invalid OpenAI tool arguments")?;
|
||||
anyhow::ensure!(
|
||||
arguments.is_object()
|
||||
&& !parsed.iter().any(|previous: &ToolCall| previous.id == id),
|
||||
"Invalid OpenAI tool call"
|
||||
);
|
||||
parsed.push(ToolCall {
|
||||
id: id.into(),
|
||||
name: name.into(),
|
||||
arguments,
|
||||
});
|
||||
}
|
||||
return Ok(BackendTurn::ToolCalls(parsed));
|
||||
}
|
||||
let text = message["content"]
|
||||
.as_str()
|
||||
.or_else(|| message["refusal"].as_str())
|
||||
.filter(|text| !text.trim().is_empty())
|
||||
.context("OpenAI returned no text; check model compatibility")?;
|
||||
Ok(BackendTurn::Text(text.into()))
|
||||
}
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::assistant::tools::{Role, ToolResult};
|
||||
#[test]
|
||||
fn parses_text_and_rejects_incomplete_or_malformed_tool_calls() {
|
||||
assert!(
|
||||
matches!(parse_response(&json!({"choices":[{"message":{"content":"hello"}}]})).unwrap(), BackendTurn::Text(text) if text == "hello")
|
||||
);
|
||||
let valid = json!({"choices":[{"message":{"tool_calls":[{"id":"call_1","type":"function","function":{"name":"status","arguments":"{\"count\":1}"}}]}}]});
|
||||
assert!(
|
||||
matches!(parse_response(&valid).unwrap(), BackendTurn::ToolCalls(calls) if calls[0].arguments["count"] == 1)
|
||||
);
|
||||
for args in ["{", "null", "[]"] {
|
||||
let mut invalid = valid.clone();
|
||||
invalid["choices"][0]["message"]["tool_calls"][0]["function"]["arguments"] =
|
||||
json!(args);
|
||||
assert!(parse_response(&invalid).is_err());
|
||||
}
|
||||
assert!(parse_response(
|
||||
&json!({"choices":[{"finish_reason":"length","message":{"content":"partial"}}]})
|
||||
)
|
||||
.is_err());
|
||||
assert!(parse_response(&json!({"choices":[]})).is_err());
|
||||
}
|
||||
#[test]
|
||||
fn errors_distinguish_credentials_limits_and_outages_without_raw_provider_data() {
|
||||
assert!(error_message(401).contains("API key"));
|
||||
assert!(error_message(429).contains("limit"));
|
||||
assert!(!error_message(503).contains("key"));
|
||||
let wire = super::super::routstr::message_to_wire(&ChatMessage {
|
||||
role: Role::Tool,
|
||||
text: None,
|
||||
tool_calls: vec![],
|
||||
tool_results: vec![ToolResult {
|
||||
call_id: "call_1".into(),
|
||||
content: "result".into(),
|
||||
is_error: false,
|
||||
}],
|
||||
});
|
||||
assert_eq!(wire[0]["tool_call_id"], "call_1");
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn real_http_adapter_sends_private_key_only_in_header_and_never_follows_redirect() {
|
||||
use hyper::{
|
||||
service::{make_service_fn, service_fn},
|
||||
Body, Response, Server,
|
||||
};
|
||||
use std::sync::{Arc, Mutex};
|
||||
let captured = Arc::new(Mutex::new(Vec::new()));
|
||||
let capture = captured.clone();
|
||||
let server = Server::bind(&([127, 0, 0, 1], 0).into()).serve(make_service_fn(move |_| {
|
||||
let capture = capture.clone();
|
||||
async move {
|
||||
Ok::<_, hyper::Error>(service_fn(move |request: hyper::Request<Body>| {
|
||||
let capture = capture.clone();
|
||||
async move {
|
||||
let (parts, body) = request.into_parts();
|
||||
let body = hyper::body::to_bytes(body).await?;
|
||||
capture.lock().unwrap().push((
|
||||
parts.headers,
|
||||
serde_json::from_slice::<Value>(&body).unwrap(),
|
||||
));
|
||||
Ok::<_, hyper::Error>(
|
||||
Response::builder()
|
||||
.status(302)
|
||||
.header("Location", "/leak")
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
}
|
||||
}))
|
||||
}
|
||||
}));
|
||||
let url = format!("http://{}/v1/chat/completions", server.local_addr());
|
||||
let task = tokio::spawn(server);
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
crate::settings::model_provider::save_key(dir.path(), "openai", "fixture-private-key")
|
||||
.await
|
||||
.unwrap();
|
||||
let backend = OpenaiBackend::new(dir.path().into(), "test-model".into());
|
||||
let history = [ChatMessage {
|
||||
role: Role::User,
|
||||
text: Some("Hello".into()),
|
||||
tool_calls: vec![],
|
||||
tool_results: vec![],
|
||||
}];
|
||||
assert!(backend
|
||||
.send_at(&url, "Be helpful", &[], &history)
|
||||
.await
|
||||
.is_err());
|
||||
let requests = captured.lock().unwrap();
|
||||
assert_eq!(requests.len(), 1);
|
||||
assert_eq!(requests[0].0["authorization"], "Bearer fixture-private-key");
|
||||
assert_eq!(requests[0].1["store"], false);
|
||||
assert_eq!(requests[0].1["max_completion_tokens"], 2048);
|
||||
assert!(!requests[0].1.to_string().contains("fixture-private-key"));
|
||||
drop(requests);
|
||||
let private = [ChatMessage {
|
||||
role: Role::User,
|
||||
text: Some("fixture-private-key".into()),
|
||||
tool_calls: vec![],
|
||||
tool_results: vec![],
|
||||
}];
|
||||
assert!(backend
|
||||
.send_at(&url, "Be helpful", &[], &private)
|
||||
.await
|
||||
.is_err());
|
||||
assert_eq!(captured.lock().unwrap().len(), 1);
|
||||
task.abort();
|
||||
}
|
||||
}
|
||||
@@ -295,7 +295,7 @@ fn parse_openai_tool_calls(raw_calls: &[Value]) -> Vec<ToolCall> {
|
||||
/// (the wire-format inverse of `parse_openai_tool_calls`), and tool-result
|
||||
/// turns carry `tool_call_id` so each call's id is echoed back exactly —
|
||||
/// the OpenAI-shape contract this adapter's edge is responsible for.
|
||||
fn message_to_wire(msg: &ChatMessage) -> Vec<Value> {
|
||||
pub(super) fn message_to_wire(msg: &ChatMessage) -> Vec<Value> {
|
||||
match msg.role {
|
||||
Role::System => vec![],
|
||||
Role::User => vec![json!({
|
||||
|
||||
@@ -0,0 +1,199 @@
|
||||
//! Durable, node-encrypted approval replies. Relay acknowledgement is not peer
|
||||
//! acceptance: keep retrying the same invite until reciprocal membership exists.
|
||||
use anyhow::{Context, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::path::Path;
|
||||
use tokio::{fs, io::AsyncWriteExt};
|
||||
|
||||
const FILE: &str = "federation/handshake-delivery.enc";
|
||||
const DOMAIN: &[u8] = b"archipelago-handshake-delivery-v1";
|
||||
static LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
|
||||
|
||||
#[derive(Clone, Serialize, Deserialize)]
|
||||
pub(crate) struct ApprovalReply {
|
||||
pub request_id: String,
|
||||
pub recipient: String,
|
||||
pub expected_did: String,
|
||||
pub invite_code: String,
|
||||
pub attempts: u32,
|
||||
pub next_attempt: i64,
|
||||
}
|
||||
|
||||
async fn load(data_dir: &Path) -> Result<Vec<ApprovalReply>> {
|
||||
let bytes = match fs::read(data_dir.join(FILE)).await {
|
||||
Ok(bytes) => bytes,
|
||||
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
|
||||
Err(e) => return Err(e.into()),
|
||||
};
|
||||
let key = crate::storage_crypto::derive_key(data_dir, DOMAIN).await?;
|
||||
let plaintext = crate::storage_crypto::open(&bytes, &key)?;
|
||||
serde_json::from_slice(&plaintext)
|
||||
.context("Invalid handshake delivery store; preserved for recovery")
|
||||
}
|
||||
|
||||
async fn save(data_dir: &Path, entries: &[ApprovalReply]) -> Result<()> {
|
||||
let key = crate::storage_crypto::derive_key(data_dir, DOMAIN).await?;
|
||||
let bytes = crate::storage_crypto::seal(&serde_json::to_vec(entries)?, &key)?;
|
||||
let path = data_dir.join(FILE);
|
||||
let parent = path.parent().context("Delivery parent missing")?;
|
||||
fs::create_dir_all(parent).await?;
|
||||
let temporary = parent.join(format!(".delivery-{}.tmp", uuid::Uuid::new_v4()));
|
||||
let result = async {
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.create_new(true)
|
||||
.write(true)
|
||||
.mode(0o600)
|
||||
.open(&temporary)
|
||||
.await?;
|
||||
file.write_all(&bytes).await?;
|
||||
file.sync_all().await?;
|
||||
drop(file);
|
||||
fs::rename(&temporary, &path).await?;
|
||||
fs::File::open(parent).await?.sync_all().await?;
|
||||
Ok::<_, anyhow::Error>(())
|
||||
}
|
||||
.await;
|
||||
if result.is_err() {
|
||||
let _ = fs::remove_file(temporary).await;
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
pub(crate) async fn find(data_dir: &Path, request_id: &str) -> Result<Option<ApprovalReply>> {
|
||||
let _guard = LOCK.lock().await;
|
||||
Ok(load(data_dir)
|
||||
.await?
|
||||
.into_iter()
|
||||
.find(|entry| entry.request_id == request_id))
|
||||
}
|
||||
|
||||
pub(crate) async fn stage(data_dir: &Path, reply: ApprovalReply) -> Result<ApprovalReply> {
|
||||
let _guard = LOCK.lock().await;
|
||||
let mut entries = load(data_dir).await?;
|
||||
if let Some(existing) = entries
|
||||
.iter()
|
||||
.find(|entry| entry.request_id == reply.request_id)
|
||||
{
|
||||
anyhow::ensure!(
|
||||
existing.recipient == reply.recipient && existing.expected_did == reply.expected_did,
|
||||
"Approval recipient changed; refusing delivery"
|
||||
);
|
||||
return Ok(existing.clone());
|
||||
}
|
||||
anyhow::ensure!(entries.len() < 1024, "Handshake delivery queue is full");
|
||||
entries.push(reply.clone());
|
||||
save(data_dir, &entries).await?;
|
||||
Ok(reply)
|
||||
}
|
||||
|
||||
/// Claim before sending, including failed sends. Concurrent polls cannot create
|
||||
/// retry storms; a crash after this write delays but never loses the reply.
|
||||
pub(crate) async fn claim(
|
||||
data_dir: &Path,
|
||||
request_id: &str,
|
||||
now: i64,
|
||||
) -> Result<Option<ApprovalReply>> {
|
||||
let _guard = LOCK.lock().await;
|
||||
let mut entries = load(data_dir).await?;
|
||||
let Some(entry) = entries
|
||||
.iter_mut()
|
||||
.find(|entry| entry.request_id == request_id)
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
if entry.next_attempt > now {
|
||||
return Ok(None);
|
||||
}
|
||||
entry.attempts = entry.attempts.saturating_add(1);
|
||||
let delay = 30_i64
|
||||
.saturating_mul(1_i64 << entry.attempts.min(7))
|
||||
.min(3600);
|
||||
entry.next_attempt = now.saturating_add(delay);
|
||||
let claimed = entry.clone();
|
||||
save(data_dir, &entries).await?;
|
||||
Ok(Some(claimed))
|
||||
}
|
||||
|
||||
pub(crate) async fn remove(data_dir: &Path, request_id: &str) -> Result<()> {
|
||||
let _guard = LOCK.lock().await;
|
||||
let mut entries = load(data_dir).await?;
|
||||
let before = entries.len();
|
||||
entries.retain(|entry| entry.request_id != request_id);
|
||||
if entries.len() != before {
|
||||
save(data_dir, &entries).await?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
async fn fixture() -> tempfile::TempDir {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
fs::create_dir_all(dir.path().join("identity"))
|
||||
.await
|
||||
.unwrap();
|
||||
fs::write(dir.path().join("identity/node_key"), [7; 32])
|
||||
.await
|
||||
.unwrap();
|
||||
dir
|
||||
}
|
||||
fn reply() -> ApprovalReply {
|
||||
ApprovalReply {
|
||||
request_id: "request-1".into(),
|
||||
recipient: "recipient".into(),
|
||||
expected_did: "did:key:peer".into(),
|
||||
invite_code: "secret-invite".into(),
|
||||
attempts: 0,
|
||||
next_attempt: 0,
|
||||
}
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn encrypted_reply_survives_reload_and_retry_claim_is_exclusive() {
|
||||
let dir = fixture().await;
|
||||
stage(dir.path(), reply()).await.unwrap();
|
||||
let raw = fs::read(dir.path().join(FILE)).await.unwrap();
|
||||
assert!(!raw.windows(13).any(|bytes| bytes == b"secret-invite"));
|
||||
assert_eq!(
|
||||
find(dir.path(), "request-1")
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.invite_code,
|
||||
"secret-invite"
|
||||
);
|
||||
let (a, b) = tokio::join!(
|
||||
claim(dir.path(), "request-1", 100),
|
||||
claim(dir.path(), "request-1", 100)
|
||||
);
|
||||
assert_eq!(
|
||||
usize::from(a.unwrap().is_some()) + usize::from(b.unwrap().is_some()),
|
||||
1
|
||||
);
|
||||
assert!(claim(dir.path(), "request-1", 159).await.unwrap().is_none());
|
||||
assert_eq!(
|
||||
claim(dir.path(), "request-1", 160)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.attempts,
|
||||
2
|
||||
);
|
||||
remove(dir.path(), "request-1").await.unwrap();
|
||||
assert!(find(dir.path(), "request-1").await.unwrap().is_none());
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn corrupt_store_is_not_overwritten_and_recipient_cannot_change() {
|
||||
let dir = fixture().await;
|
||||
stage(dir.path(), reply()).await.unwrap();
|
||||
let mut other = reply();
|
||||
other.recipient = "different-recipient".into();
|
||||
assert!(stage(dir.path(), other).await.is_err());
|
||||
let path = dir.path().join(FILE);
|
||||
let mut raw = fs::read(&path).await.unwrap();
|
||||
raw[20] ^= 1;
|
||||
fs::write(&path, &raw).await.unwrap();
|
||||
assert!(stage(dir.path(), reply()).await.is_err());
|
||||
assert_eq!(fs::read(path).await.unwrap(), raw);
|
||||
}
|
||||
}
|
||||
@@ -134,6 +134,32 @@ pub fn parse_invite(code: &str) -> Result<ParsedInvite> {
|
||||
})
|
||||
}
|
||||
|
||||
/// Bind a Nostr-discovery reply to the node the operator requested, and cap
|
||||
/// its grant before any local node entry or callback is written. Legacy invites
|
||||
/// default to Trusted, which must never transiently authorize discovery peers.
|
||||
pub(crate) fn restrict_discovery_invite(code: &str, expected_did: &str) -> Result<String> {
|
||||
use base64::Engine;
|
||||
let parsed = parse_invite(code)?;
|
||||
anyhow::ensure!(
|
||||
!expected_did.is_empty() && parsed.did == expected_did,
|
||||
"Peer invite does not match the requested node"
|
||||
);
|
||||
anyhow::ensure!(
|
||||
crate::identity::did_key_from_pubkey_hex(&parsed.pubkey)? == parsed.did,
|
||||
"Peer invite DID does not match its identity key"
|
||||
);
|
||||
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD.decode(
|
||||
code.strip_prefix("fed1:")
|
||||
.context("Invalid invite prefix")?,
|
||||
)?;
|
||||
let mut payload: serde_json::Value = serde_json::from_slice(&bytes)?;
|
||||
payload["trust"] = serde_json::json!("observer");
|
||||
Ok(format!(
|
||||
"fed1:{}",
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(serde_json::to_vec(&payload)?)
|
||||
))
|
||||
}
|
||||
|
||||
/// Accept an invite: parse code, verify the remote node, add to federation.
|
||||
pub async fn accept_invite(
|
||||
data_dir: &Path,
|
||||
@@ -621,3 +647,35 @@ mod tests {
|
||||
assert_eq!(nodes.len(), 1, "re-accept should not duplicate");
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod discovery_invite_scope_tests {
|
||||
use super::*;
|
||||
use base64::Engine;
|
||||
#[test]
|
||||
fn discovery_reply_binds_identity_and_caps_legacy_trust_before_acceptance() {
|
||||
let key = "33".repeat(32);
|
||||
let did = crate::identity::did_key_from_pubkey_hex(&key).unwrap();
|
||||
let payload =
|
||||
serde_json::json!({"did":did,"pubkey":key,"onion":"test.onion","token":"test-token"});
|
||||
let code = format!(
|
||||
"fed1:{}",
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD
|
||||
.encode(serde_json::to_vec(&payload).unwrap())
|
||||
);
|
||||
let restricted = restrict_discovery_invite(&code, &did).unwrap();
|
||||
let parsed = parse_invite(&restricted).unwrap();
|
||||
assert_eq!(parsed.trust_level, TrustLevel::Observer);
|
||||
assert_eq!(parsed.token, "test-token");
|
||||
assert!(restrict_discovery_invite(&code, "did:key:someone-else").is_err());
|
||||
assert!(restrict_discovery_invite(&code, "").is_err());
|
||||
let mut forged = payload;
|
||||
forged["pubkey"] = serde_json::json!("44".repeat(32));
|
||||
let forged = format!(
|
||||
"fed1:{}",
|
||||
base64::engine::general_purpose::URL_SAFE_NO_PAD
|
||||
.encode(serde_json::to_vec(&forged).unwrap())
|
||||
);
|
||||
assert!(restrict_discovery_invite(&forged, &did).is_err());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,12 +6,14 @@
|
||||
|
||||
mod invites;
|
||||
pub mod pending;
|
||||
pub(crate) mod handshake_delivery;
|
||||
mod storage;
|
||||
mod sync;
|
||||
mod types;
|
||||
|
||||
// Re-export all public items so `crate::federation::*` continues to work.
|
||||
pub use invites::{accept_invite, create_invite, parse_invite};
|
||||
pub(crate) use invites::restrict_discovery_invite;
|
||||
// Crate-internal: used by the periodic federation auto-sync to re-assert
|
||||
// membership to peers that don't list us back (asymmetry self-heal).
|
||||
pub(crate) use invites::notify_join;
|
||||
|
||||
@@ -14,6 +14,9 @@ use anyhow::{Context, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::path::Path;
|
||||
use tokio::fs;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
static PENDING_STORE_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
|
||||
|
||||
const PENDING_FILE: &str = "federation/pending_requests.json";
|
||||
const MAX_PENDING_PER_PUBKEY: usize = 5;
|
||||
@@ -76,11 +79,12 @@ pub async fn load_pending(data_dir: &Path) -> Result<Vec<PendingPeerRequest>> {
|
||||
let content = fs::read_to_string(&path)
|
||||
.await
|
||||
.context("Failed to read pending requests file")?;
|
||||
let file: PendingRequestsFile = serde_json::from_str(&content).unwrap_or_default();
|
||||
let file: PendingRequestsFile = serde_json::from_str(&content)
|
||||
.context("Invalid pending requests file; preserving existing data")?;
|
||||
Ok(file.requests)
|
||||
}
|
||||
|
||||
pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> Result<()> {
|
||||
async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> Result<()> {
|
||||
let path = data_dir.join(PENDING_FILE);
|
||||
if let Some(parent) = path.parent() {
|
||||
fs::create_dir_all(parent)
|
||||
@@ -92,9 +96,27 @@ pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> R
|
||||
};
|
||||
let content =
|
||||
serde_json::to_string_pretty(&file).context("Failed to serialize pending requests")?;
|
||||
fs::write(&path, content)
|
||||
.await
|
||||
.context("Failed to write pending requests file")?;
|
||||
let parent = path.parent().context("Pending requests parent missing")?;
|
||||
let temporary = parent.join(format!(".pending-{}.tmp", uuid::Uuid::new_v4()));
|
||||
let result = async {
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.write(true)
|
||||
.create_new(true)
|
||||
.mode(0o600)
|
||||
.open(&temporary)
|
||||
.await?;
|
||||
file.write_all(content.as_bytes()).await?;
|
||||
file.sync_all().await?;
|
||||
drop(file);
|
||||
fs::rename(&temporary, &path).await?;
|
||||
fs::File::open(parent).await?.sync_all().await?;
|
||||
Ok::<_, anyhow::Error>(())
|
||||
}
|
||||
.await;
|
||||
if result.is_err() {
|
||||
let _ = fs::remove_file(&temporary).await;
|
||||
}
|
||||
result.context("Failed to atomically save pending requests")?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -102,7 +124,10 @@ pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> R
|
||||
fn expire_stale(requests: &mut Vec<PendingPeerRequest>) {
|
||||
let cutoff = chrono::Utc::now() - chrono::Duration::days(PENDING_EXPIRY_DAYS);
|
||||
for r in requests.iter_mut() {
|
||||
if !matches!(r.state, PendingState::Pending | PendingState::Sent) {
|
||||
if !matches!(
|
||||
r.state,
|
||||
PendingState::Pending | PendingState::Sent | PendingState::Approved
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
if let Ok(ts) = chrono::DateTime::parse_from_rfc3339(&r.received_at) {
|
||||
@@ -131,6 +156,7 @@ pub async fn insert_inbound(
|
||||
from_name: Option<String>,
|
||||
message: Option<String>,
|
||||
) -> Result<Option<PendingPeerRequest>> {
|
||||
let _guard = PENDING_STORE_LOCK.lock().await;
|
||||
let mut requests = load_pending(data_dir).await?;
|
||||
expire_stale(&mut requests);
|
||||
|
||||
@@ -189,6 +215,7 @@ pub async fn insert_outbound(
|
||||
to_name: Option<String>,
|
||||
message: Option<String>,
|
||||
) -> Result<PendingPeerRequest> {
|
||||
let _guard = PENDING_STORE_LOCK.lock().await;
|
||||
let mut requests = load_pending(data_dir).await?;
|
||||
expire_stale(&mut requests);
|
||||
requests.retain(|r| {
|
||||
@@ -218,6 +245,7 @@ pub async fn find_by_id(data_dir: &Path, id: &str) -> Result<Option<PendingPeerR
|
||||
}
|
||||
|
||||
pub async fn set_state(data_dir: &Path, id: &str, state: PendingState) -> Result<()> {
|
||||
let _guard = PENDING_STORE_LOCK.lock().await;
|
||||
let mut requests = load_pending(data_dir).await?;
|
||||
if let Some(r) = requests.iter_mut().find(|r| r.id == id) {
|
||||
r.state = state;
|
||||
@@ -228,10 +256,32 @@ pub async fn set_state(data_dir: &Path, id: &str, state: PendingState) -> Result
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Resolve a pending decision once; concurrent approval/rejection cannot
|
||||
/// overwrite each other after a slow network request.
|
||||
pub async fn decide(data_dir: &Path, id: &str, decision: PendingState) -> Result<()> {
|
||||
anyhow::ensure!(
|
||||
matches!(decision, PendingState::Approved | PendingState::Rejected),
|
||||
"Invalid pending decision"
|
||||
);
|
||||
let _guard = PENDING_STORE_LOCK.lock().await;
|
||||
let mut requests = load_pending(data_dir).await?;
|
||||
let row = requests
|
||||
.iter_mut()
|
||||
.find(|row| row.id == id)
|
||||
.context("Pending request not found")?;
|
||||
anyhow::ensure!(
|
||||
!row.outbound && row.state == PendingState::Pending,
|
||||
"Request has already been decided"
|
||||
);
|
||||
row.state = decision;
|
||||
save_pending(data_dir, &requests).await
|
||||
}
|
||||
|
||||
/// Remove a pending request entirely. Used when the sender cancels an
|
||||
/// outbound request they initiated and we want it gone (not just marked
|
||||
/// Rejected/Cancelled — those states fill up the UI audit trail).
|
||||
pub async fn delete(data_dir: &Path, id: &str) -> Result<()> {
|
||||
let _guard = PENDING_STORE_LOCK.lock().await;
|
||||
let mut requests = load_pending(data_dir).await?;
|
||||
let before = requests.len();
|
||||
requests.retain(|r| r.id != id);
|
||||
@@ -372,3 +422,140 @@ mod tests {
|
||||
assert_eq!(reloaded.state, PendingState::Approved);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod persistence_regressions {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn concurrent_requests_are_not_lost() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let mut tasks = Vec::new();
|
||||
for i in 0..24 {
|
||||
let path = dir.path().to_path_buf();
|
||||
tasks.push(tokio::spawn(async move {
|
||||
insert_inbound(
|
||||
&path,
|
||||
format!("key-{i}"),
|
||||
format!("npub-{i}"),
|
||||
format!("did:key:{i}"),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
}));
|
||||
}
|
||||
for task in tasks {
|
||||
assert!(task.await.unwrap().is_some());
|
||||
}
|
||||
let rows = load_pending(dir.path()).await.unwrap();
|
||||
assert_eq!(rows.len(), 24);
|
||||
assert_eq!(
|
||||
rows.iter()
|
||||
.map(|r| &r.id)
|
||||
.collect::<std::collections::HashSet<_>>()
|
||||
.len(),
|
||||
24
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn malformed_store_is_preserved_instead_of_replaced_with_one_request() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let path = dir.path().join(PENDING_FILE);
|
||||
fs::create_dir_all(path.parent().unwrap()).await.unwrap();
|
||||
let damaged = b"{incomplete existing requests";
|
||||
fs::write(&path, damaged).await.unwrap();
|
||||
assert!(insert_inbound(
|
||||
dir.path(),
|
||||
"key".into(),
|
||||
"npub".into(),
|
||||
"did:key:test".into(),
|
||||
None,
|
||||
None
|
||||
)
|
||||
.await
|
||||
.is_err());
|
||||
assert_eq!(fs::read(path).await.unwrap(), damaged);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod decision_regressions {
|
||||
use super::*;
|
||||
#[tokio::test]
|
||||
async fn only_one_concurrent_operator_decision_wins() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let row = insert_inbound(
|
||||
dir.path(),
|
||||
"key".into(),
|
||||
"npub".into(),
|
||||
"did:key:peer".into(),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
let (approve, reject) = tokio::join!(
|
||||
decide(dir.path(), &row.id, PendingState::Approved),
|
||||
decide(dir.path(), &row.id, PendingState::Rejected)
|
||||
);
|
||||
assert_eq!(
|
||||
usize::from(approve.is_ok()) + usize::from(reject.is_ok()),
|
||||
1
|
||||
);
|
||||
let saved = find_by_id(dir.path(), &row.id).await.unwrap().unwrap();
|
||||
assert_eq!(
|
||||
saved.state,
|
||||
if approve.is_ok() {
|
||||
PendingState::Approved
|
||||
} else {
|
||||
PendingState::Rejected
|
||||
}
|
||||
);
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn an_expired_approval_does_not_block_a_new_request_forever() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let mut row = insert_inbound(
|
||||
dir.path(),
|
||||
"key".into(),
|
||||
"npub".into(),
|
||||
"did:key:peer".into(),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
row.state = PendingState::Approved;
|
||||
row.received_at = (chrono::Utc::now() - chrono::Duration::days(31)).to_rfc3339();
|
||||
save_pending(dir.path(), &[row]).await.unwrap();
|
||||
let renewed = insert_inbound(
|
||||
dir.path(),
|
||||
"key".into(),
|
||||
"npub".into(),
|
||||
"did:key:peer".into(),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(renewed.is_some());
|
||||
let rows = load_pending(dir.path()).await.unwrap();
|
||||
assert_eq!(
|
||||
rows.iter()
|
||||
.filter(|r| r.state == PendingState::Expired)
|
||||
.count(),
|
||||
1
|
||||
);
|
||||
assert_eq!(
|
||||
rows.iter()
|
||||
.filter(|r| r.state == PendingState::Pending)
|
||||
.count(),
|
||||
1
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -230,12 +230,12 @@ async fn merge_transitive_peers(
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub fn build_local_state(
|
||||
apps: Vec<AppStatus>,
|
||||
cpu: f64,
|
||||
mem_used: u64,
|
||||
mem_total: u64,
|
||||
disk_used: u64,
|
||||
disk_total: u64,
|
||||
uptime: u64,
|
||||
cpu: Option<f64>,
|
||||
mem_used: Option<u64>,
|
||||
mem_total: Option<u64>,
|
||||
disk_used: Option<u64>,
|
||||
disk_total: Option<u64>,
|
||||
uptime: Option<u64>,
|
||||
tor_active: bool,
|
||||
server_name: Option<String>,
|
||||
nostr_npub: Option<String>,
|
||||
@@ -261,12 +261,12 @@ pub fn build_local_state(
|
||||
timestamp: chrono::Utc::now().to_rfc3339(),
|
||||
node_name: server_name,
|
||||
apps,
|
||||
cpu_usage_percent: Some(cpu),
|
||||
mem_used_bytes: Some(mem_used),
|
||||
mem_total_bytes: Some(mem_total),
|
||||
disk_used_bytes: Some(disk_used),
|
||||
disk_total_bytes: Some(disk_total),
|
||||
uptime_secs: Some(uptime),
|
||||
cpu_usage_percent: cpu,
|
||||
mem_used_bytes: mem_used,
|
||||
mem_total_bytes: mem_total,
|
||||
disk_used_bytes: disk_used,
|
||||
disk_total_bytes: disk_total,
|
||||
uptime_secs: uptime,
|
||||
tor_active: Some(tor_active),
|
||||
nostr_npub,
|
||||
own_fips_npub,
|
||||
@@ -355,12 +355,12 @@ mod tests {
|
||||
status: "running".to_string(),
|
||||
version: Some("0.18".to_string()),
|
||||
}],
|
||||
25.5,
|
||||
2_000_000_000,
|
||||
8_000_000_000,
|
||||
100_000_000_000,
|
||||
500_000_000_000,
|
||||
3600,
|
||||
Some(25.5),
|
||||
Some(2_000_000_000),
|
||||
Some(8_000_000_000),
|
||||
Some(100_000_000_000),
|
||||
Some(500_000_000_000),
|
||||
Some(3600),
|
||||
true,
|
||||
Some("Test Node".to_string()),
|
||||
None,
|
||||
@@ -430,12 +430,12 @@ mod tests {
|
||||
];
|
||||
let state = build_local_state(
|
||||
vec![],
|
||||
0.0,
|
||||
0,
|
||||
0,
|
||||
0,
|
||||
0,
|
||||
0,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
None,
|
||||
true,
|
||||
None,
|
||||
None,
|
||||
|
||||
@@ -13,11 +13,11 @@ pub async fn collect_snapshot() -> Result<MetricSnapshot> {
|
||||
read_loadavg(),
|
||||
);
|
||||
|
||||
let cpu = cpu.unwrap_or(0.0);
|
||||
let (mem_used, mem_total) = mem.unwrap_or((0, 0));
|
||||
let (disk_used, disk_total) = disk.unwrap_or((0, 0));
|
||||
let (net_rx, net_tx) = net.unwrap_or((0, 0));
|
||||
let (l1, l5, l15) = load.unwrap_or((0.0, 0.0, 0.0));
|
||||
let cpu = cpu?;
|
||||
let (mem_used, mem_total) = mem?;
|
||||
let (disk_used, disk_total) = disk?;
|
||||
let (net_rx, net_tx) = net?;
|
||||
let (l1, l5, l15) = load?;
|
||||
|
||||
let system = SystemMetrics {
|
||||
cpu_percent: cpu,
|
||||
@@ -120,10 +120,9 @@ async fn read_disk_usage() -> Result<(u64, u64)> {
|
||||
} else {
|
||||
"/"
|
||||
};
|
||||
let output = tokio::process::Command::new("df")
|
||||
.args(["--block-size=1", "--output=used,size", target])
|
||||
.output()
|
||||
.await
|
||||
let mut command = tokio::process::Command::new("df");
|
||||
command.args(["--block-size=1", "--output=used,size", target]);
|
||||
let output = bounded_output(command, std::time::Duration::from_secs(3)).await
|
||||
.context("Failed to run df")?;
|
||||
|
||||
if !output.status.success() {
|
||||
@@ -215,12 +214,23 @@ async fn read_network_totals() -> Result<(u64, u64)> {
|
||||
Ok((rx_total, tx_total))
|
||||
}
|
||||
|
||||
/// A wedged runtime or filesystem must not freeze every monitoring snapshot.
|
||||
/// Dropping a timed-out child kills it, so repeated polls cannot leak processes.
|
||||
async fn bounded_output(
|
||||
mut command: tokio::process::Command,
|
||||
timeout: std::time::Duration,
|
||||
) -> Result<std::process::Output> {
|
||||
command.kill_on_drop(true);
|
||||
tokio::time::timeout(timeout, command.output())
|
||||
.await.context("Metrics subprocess timed out")?
|
||||
.context("Metrics subprocess failed")
|
||||
}
|
||||
|
||||
/// Get per-container resource stats via `podman stats --no-stream --format json`.
|
||||
async fn read_container_stats() -> Result<Vec<ContainerMetrics>> {
|
||||
let output = tokio::process::Command::new("podman")
|
||||
.args(["stats", "--no-stream", "--format", "json"])
|
||||
.output()
|
||||
.await
|
||||
let mut command = tokio::process::Command::new("podman");
|
||||
command.args(["stats", "--no-stream", "--format", "json"]);
|
||||
let output = bounded_output(command, std::time::Duration::from_secs(8)).await
|
||||
.context("Failed to run podman stats")?;
|
||||
|
||||
if !output.status.success() {
|
||||
@@ -391,3 +401,35 @@ mod tests {
|
||||
assert_eq!(parse_bytes_field(&obj, "mem"), Some(268435456));
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod subprocess_deadline_tests {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_stalled_metrics_command_is_bounded_and_killed() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let pid_file = dir.path().join("pid");
|
||||
let mut command = tokio::process::Command::new("sh");
|
||||
command.arg("-c").arg("echo $$ > \"$1\"; exec sleep 30").arg("metrics-test").arg(&pid_file);
|
||||
let start = std::time::Instant::now();
|
||||
let error = bounded_output(command, std::time::Duration::from_millis(500)).await.unwrap_err();
|
||||
assert!(error.to_string().contains("timed out"));
|
||||
assert!(start.elapsed() < std::time::Duration::from_secs(3));
|
||||
let pid = tokio::fs::read_to_string(pid_file).await.unwrap();
|
||||
for _ in 0..40 {
|
||||
if !std::path::Path::new(&format!("/proc/{}", pid.trim())).exists() { return; }
|
||||
tokio::time::sleep(std::time::Duration::from_millis(25)).await;
|
||||
}
|
||||
panic!("Timed-out metrics subprocess was not reaped");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn successful_metrics_output_is_preserved() {
|
||||
let mut command = tokio::process::Command::new("printf");
|
||||
command.arg("metrics-ok");
|
||||
let output = bounded_output(command, std::time::Duration::from_secs(1)).await.unwrap();
|
||||
assert!(output.status.success());
|
||||
assert_eq!(output.stdout, b"metrics-ok");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,21 +14,21 @@ use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
/// Spawn the background metrics collector (runs every 300 seconds / 5 minutes).
|
||||
/// Spawn the background metrics collector at the store's one-minute resolution.
|
||||
/// Evaluates alert rules on each snapshot and dispatches notifications.
|
||||
/// Note: health_monitor.rs handles container state polling at 120s intervals.
|
||||
/// This collector handles system-level metrics (CPU, disk, network) and only
|
||||
/// calls podman stats every 5 minutes to avoid duplicate subprocess overhead.
|
||||
/// Runtime commands have deadlines; unavailable container stats cannot hold
|
||||
/// system readings indefinitely. Missed ticks are skipped, never replayed.
|
||||
pub fn spawn_metrics_collector(
|
||||
store: Arc<MetricsStore>,
|
||||
state: Option<Arc<crate::state::StateManager>>,
|
||||
data_dir: Option<PathBuf>,
|
||||
) {
|
||||
tokio::spawn(async move {
|
||||
// Wait 60s for system to stabilize after boot
|
||||
tokio::time::sleep(std::time::Duration::from_secs(60)).await;
|
||||
// Start promptly without competing with the very first boot tasks.
|
||||
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
|
||||
|
||||
let mut interval = tokio::time::interval(std::time::Duration::from_secs(300));
|
||||
let mut interval = tokio::time::interval(std::time::Duration::from_secs(60));
|
||||
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
|
||||
|
||||
loop {
|
||||
|
||||
@@ -291,15 +291,15 @@ impl Server {
|
||||
);
|
||||
|
||||
// Background handshake poll: fetch inbound nostr peer requests every
|
||||
// 5 minutes instead of only when a user presses the Federation Poll
|
||||
// 30 seconds instead of only when a user presses the Federation Poll
|
||||
// button (requests used to sit on relays unseen — 2026-07-22). The
|
||||
// handler's own discoverability gate makes this a no-op until the
|
||||
// user opts in.
|
||||
{
|
||||
let rpc = api_handler.rpc_handler().clone();
|
||||
tokio::spawn(async move {
|
||||
let mut tick = tokio::time::interval(std::time::Duration::from_secs(300));
|
||||
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
|
||||
let mut tick = tokio::time::interval(std::time::Duration::from_secs(30));
|
||||
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
|
||||
loop {
|
||||
tick.tick().await;
|
||||
rpc.background_handshake_poll().await;
|
||||
|
||||
@@ -9,3 +9,5 @@ pub mod session_policy;
|
||||
pub mod transport;
|
||||
|
||||
pub mod bitcoin_storage;
|
||||
|
||||
pub mod model_provider;
|
||||
|
||||
@@ -0,0 +1,178 @@
|
||||
//! Owner-selected chat provider. API keys remain in the node's private secret
|
||||
//! ledger and are never returned by settings or included in chat context.
|
||||
use anyhow::{Context, Result};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::path::Path;
|
||||
use tokio::{fs, io::AsyncWriteExt};
|
||||
|
||||
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum Provider {
|
||||
#[default]
|
||||
Auto,
|
||||
Claude,
|
||||
Openai,
|
||||
Local,
|
||||
Routstr,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub struct ModelProvider {
|
||||
#[serde(default)]
|
||||
pub provider: Provider,
|
||||
#[serde(default)]
|
||||
pub openai_model: String,
|
||||
}
|
||||
|
||||
impl ModelProvider {
|
||||
pub fn validate(&self) -> Result<()> {
|
||||
anyhow::ensure!(
|
||||
!self.openai_model.starts_with("sk-")
|
||||
&& self.openai_model.len() <= 128
|
||||
&& self
|
||||
.openai_model
|
||||
.bytes()
|
||||
.all(|b| b.is_ascii_alphanumeric() || b"-_.:".contains(&b)),
|
||||
"Invalid OpenAI model name"
|
||||
);
|
||||
anyhow::ensure!(
|
||||
self.provider != Provider::Openai || !self.openai_model.is_empty(),
|
||||
"Choose an OpenAI model before connecting"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
pub async fn load(data_dir: &Path) -> Result<Self> {
|
||||
match fs::read(data_dir.join("settings/model-provider.json")).await {
|
||||
Ok(bytes) => {
|
||||
let settings: Self = serde_json::from_slice(&bytes)
|
||||
.context("Invalid AI provider settings; preserved for recovery")?;
|
||||
settings.validate()?;
|
||||
Ok(settings)
|
||||
}
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(Self::default()),
|
||||
Err(error) => Err(error.into()),
|
||||
}
|
||||
}
|
||||
pub async fn save(&self, data_dir: &Path) -> Result<()> {
|
||||
self.validate()?;
|
||||
write_private(
|
||||
&data_dir.join("settings/model-provider.json"),
|
||||
&serde_json::to_vec(self)?,
|
||||
)
|
||||
.await
|
||||
}
|
||||
}
|
||||
|
||||
pub fn key_name(provider: &str) -> Result<&'static str> {
|
||||
match provider {
|
||||
"claude" => Ok("claude-api-key"),
|
||||
"openai" => Ok("openai-api-key"),
|
||||
_ => anyhow::bail!("Unsupported AI provider"),
|
||||
}
|
||||
}
|
||||
pub async fn has_key(data_dir: &Path, provider: &str) -> bool {
|
||||
let Ok(name) = key_name(provider) else {
|
||||
return false;
|
||||
};
|
||||
fs::read_to_string(data_dir.join("secrets").join(name))
|
||||
.await
|
||||
.is_ok_and(|key| !key.trim().is_empty())
|
||||
}
|
||||
pub async fn save_key(data_dir: &Path, provider: &str, value: &str) -> Result<()> {
|
||||
let path = data_dir.join("secrets").join(key_name(provider)?);
|
||||
let value = value.trim();
|
||||
anyhow::ensure!(
|
||||
value.len() <= 4096 && value.bytes().all(|b| b.is_ascii_graphic()),
|
||||
"Invalid API key format"
|
||||
);
|
||||
if value.is_empty() {
|
||||
match fs::remove_file(path).await {
|
||||
Ok(()) => Ok(()),
|
||||
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
|
||||
Err(e) => Err(e.into()),
|
||||
}
|
||||
} else {
|
||||
write_private(&path, value.as_bytes()).await
|
||||
}
|
||||
}
|
||||
async fn write_private(path: &Path, bytes: &[u8]) -> Result<()> {
|
||||
let parent = path.parent().context("Missing settings directory")?;
|
||||
fs::create_dir_all(parent).await?;
|
||||
let temporary = parent.join(format!(".provider-{}.tmp", uuid::Uuid::new_v4()));
|
||||
let result = async {
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.create_new(true)
|
||||
.write(true)
|
||||
.mode(0o600)
|
||||
.open(&temporary)
|
||||
.await?;
|
||||
file.write_all(bytes).await?;
|
||||
file.sync_all().await?;
|
||||
drop(file);
|
||||
fs::rename(&temporary, path).await?;
|
||||
fs::File::open(parent).await?.sync_all().await?;
|
||||
Ok::<_, anyhow::Error>(())
|
||||
}
|
||||
.await;
|
||||
if result.is_err() {
|
||||
let _ = fs::remove_file(temporary).await;
|
||||
}
|
||||
result
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
#[tokio::test]
|
||||
async fn private_keys_replace_atomically_and_never_enter_public_settings() {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
assert!(!has_key(dir.path(), "openai").await);
|
||||
save_key(dir.path(), "openai", "test-key-one")
|
||||
.await
|
||||
.unwrap();
|
||||
save_key(dir.path(), "openai", "test-key-two")
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(has_key(dir.path(), "openai").await);
|
||||
let key_path = dir.path().join("secrets/openai-api-key");
|
||||
assert_eq!(
|
||||
fs::metadata(&key_path).await.unwrap().permissions().mode() & 0o777,
|
||||
0o600
|
||||
);
|
||||
assert_eq!(fs::read_to_string(&key_path).await.unwrap(), "test-key-two");
|
||||
let settings = ModelProvider {
|
||||
provider: Provider::Openai,
|
||||
openai_model: "test-model".into(),
|
||||
};
|
||||
settings.save(dir.path()).await.unwrap();
|
||||
let body = serde_json::to_string(&ModelProvider::load(dir.path()).await.unwrap()).unwrap();
|
||||
assert!(!body.contains("test-key"));
|
||||
assert!(save_key(dir.path(), "../openai", "key").await.is_err());
|
||||
assert!(save_key(dir.path(), "openai", "key\nInjected: bad")
|
||||
.await
|
||||
.is_err());
|
||||
assert_eq!(fs::read_to_string(&key_path).await.unwrap(), "test-key-two");
|
||||
save_key(dir.path(), "openai", "").await.unwrap();
|
||||
assert!(!has_key(dir.path(), "openai").await);
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn invalid_settings_preserve_existing_configuration() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
ModelProvider::default().save(dir.path()).await.unwrap();
|
||||
let invalid = ModelProvider {
|
||||
provider: Provider::Openai,
|
||||
openai_model: String::new(),
|
||||
};
|
||||
assert!(invalid.save(dir.path()).await.is_err());
|
||||
assert_eq!(
|
||||
ModelProvider::load(dir.path()).await.unwrap().provider,
|
||||
Provider::Auto
|
||||
);
|
||||
let path = dir.path().join("settings/model-provider.json");
|
||||
fs::write(&path, b"broken").await.unwrap();
|
||||
assert!(ModelProvider::load(dir.path()).await.is_err());
|
||||
assert_eq!(fs::read(&path).await.unwrap(), b"broken");
|
||||
}
|
||||
}
|
||||
@@ -75,11 +75,15 @@ pub fn open(data: &[u8], key: &[u8; 32]) -> Result<Vec<u8>> {
|
||||
.map_err(|_| anyhow::anyhow!("decryption failed — key mismatch or corruption"))
|
||||
}
|
||||
|
||||
/// Heuristic: does this look like legacy plaintext JSON (starts with `{`/`[`)?
|
||||
/// Encrypted blobs start with a random nonce byte, so a `{`/`[` first byte is a
|
||||
/// reliable migration signal.
|
||||
/// Recognize a complete legacy JSON object/array, including leading whitespace.
|
||||
/// A random nonce can start with `{` or `[`; checking only that byte misclassifies
|
||||
/// valid ciphertext and can trigger an empty-store migration. Validate the entire
|
||||
/// document without allocating a second copy of the store's object tree.
|
||||
pub fn is_plaintext_json(raw: &[u8]) -> bool {
|
||||
matches!(raw.first(), Some(b'{') | Some(b'['))
|
||||
matches!(
|
||||
raw.iter().copied().find(|byte| !byte.is_ascii_whitespace()),
|
||||
Some(b'{') | Some(b'[')
|
||||
) && serde_json::from_slice::<serde::de::IgnoredAny>(raw).is_ok()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -110,9 +114,44 @@ mod tests {
|
||||
fn detects_plaintext_vs_ciphertext() {
|
||||
assert!(is_plaintext_json(b"{\"a\":1}"));
|
||||
assert!(is_plaintext_json(b"[]"));
|
||||
assert!(is_plaintext_json(b" \r\n\t{\"messages\": []}\n"));
|
||||
for invalid in [
|
||||
b"{".as_slice(),
|
||||
b"[",
|
||||
b"{}trailing",
|
||||
b"[\xff]",
|
||||
b"null",
|
||||
b"\"text\"",
|
||||
] {
|
||||
assert!(!is_plaintext_json(invalid));
|
||||
}
|
||||
assert!(!is_plaintext_json(&seal(b"x", &[3u8; 32]).unwrap()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn json_prefix_nonce_does_not_trigger_plaintext_migration() {
|
||||
use chacha20poly1305::aead::{Aead, KeyInit};
|
||||
let key = [3u8; 32];
|
||||
let plaintext = br#"{"messages":[{"message":"preserve me"}]}"#;
|
||||
let cipher = chacha20poly1305::ChaCha20Poly1305::new_from_slice(&key).unwrap();
|
||||
// Deterministically reproduce both collisions instead of relying on
|
||||
// OsRng to happen to pick one during a test run.
|
||||
for prefix in [b'{', b'['] {
|
||||
let mut nonce = [0xffu8; 12];
|
||||
nonce[0] = prefix;
|
||||
let encrypted = cipher
|
||||
.encrypt(
|
||||
chacha20poly1305::aead::generic_array::GenericArray::from_slice(&nonce),
|
||||
plaintext.as_slice(),
|
||||
)
|
||||
.unwrap();
|
||||
let mut envelope = nonce.to_vec();
|
||||
envelope.extend(encrypted);
|
||||
assert!(!is_plaintext_json(&envelope));
|
||||
assert_eq!(open(&envelope, &key).unwrap(), plaintext);
|
||||
}
|
||||
}
|
||||
|
||||
/// KEY-05 regression: a blob written by the pre-migration `seal` must still
|
||||
/// open after the migration.
|
||||
///
|
||||
|
||||
Reference in New Issue
Block a user