use anyhow::{Context, Result};
use serde::Serialize;
use serde_json::{Value, json};
#[derive(Debug, Clone, Serialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum SyncDelivery {
Delivered {
event_id: String,
relay_url: String,
slot_id: String,
},
Duplicate {
event_id: String,
relay_url: String,
slot_id: String,
},
DeliveredNostr {
event_id: String,
relay_url: String,
npub: String,
},
PeerUnknown { event_id: String },
SlotStale {
event_id: String,
relay_url: String,
slot_id: String,
detail: String,
},
TransportError {
event_id: String,
relay_url: String,
slot_id: String,
detail: String,
},
}
impl SyncDelivery {
pub fn status_str(&self) -> &'static str {
match self {
SyncDelivery::Delivered { .. } => "delivered",
SyncDelivery::Duplicate { .. } => "duplicate",
SyncDelivery::DeliveredNostr { .. } => "delivered_nostr",
SyncDelivery::PeerUnknown { .. } => "peer_unknown",
SyncDelivery::SlotStale { .. } => "slot_stale",
SyncDelivery::TransportError { .. } => "transport_error",
}
}
pub fn reached_relay(&self) -> bool {
matches!(
self,
SyncDelivery::Delivered { .. }
| SyncDelivery::Duplicate { .. }
| SyncDelivery::DeliveredNostr { .. }
)
}
pub fn event_id(&self) -> &str {
match self {
SyncDelivery::Delivered { event_id, .. }
| SyncDelivery::Duplicate { event_id, .. }
| SyncDelivery::DeliveredNostr { event_id, .. }
| SyncDelivery::PeerUnknown { event_id }
| SyncDelivery::SlotStale { event_id, .. }
| SyncDelivery::TransportError { event_id, .. } => event_id,
}
}
}
pub fn attempt_deliver(peer_handle: &str, signed_event: &Value) -> Result<SyncDelivery> {
let event_id = signed_event
.get("event_id")
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
let state = crate::config::read_relay_state().context("reading relay state")?;
let endpoints = crate::endpoints::peer_endpoints_in_priority_order(&state, peer_handle);
let mut last_failure: Option<SyncDelivery> = None;
for ep in endpoints {
if ep.relay_url.is_empty() || ep.slot_id.is_empty() || ep.slot_token.is_empty() {
continue;
}
let client = crate::relay_client::RelayClient::new(&ep.relay_url);
match client.post_event(&ep.slot_id, &ep.slot_token, signed_event) {
Ok(resp) => {
let now = time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap_or_default();
if let Err(e) = crate::config::append_pushed_log(peer_handle, &event_id, &now) {
eprintln!(
"wire send: pushed-log append for {peer_handle}/{event_id} failed (non-fatal): {e:#}"
);
}
return Ok(if resp.status == "duplicate" {
SyncDelivery::Duplicate {
event_id,
relay_url: ep.relay_url,
slot_id: ep.slot_id,
}
} else {
SyncDelivery::Delivered {
event_id,
relay_url: ep.relay_url,
slot_id: ep.slot_id,
}
});
}
Err(e) => {
let detail = crate::relay_client::format_transport_error(&e);
last_failure = Some(if crate::cli::error_smells_like_slot_4xx(&detail) {
SyncDelivery::SlotStale {
event_id: event_id.clone(),
relay_url: ep.relay_url,
slot_id: ep.slot_id,
detail,
}
} else {
SyncDelivery::TransportError {
event_id: event_id.clone(),
relay_url: ep.relay_url,
slot_id: ep.slot_id,
detail,
}
});
}
}
}
if let Some((peer_npub, nostr_relay)) =
crate::endpoints::peer_nostr_transport(&state, peer_handle)
&& let Ok(nsk) = crate::config::read_nostr_key()
{
match deliver_over_nostr(&peer_npub, &nostr_relay, signed_event, &nsk) {
Ok(true) => {
return Ok(SyncDelivery::DeliveredNostr {
event_id,
relay_url: nostr_relay,
npub: peer_npub,
});
}
Ok(false) => {
last_failure = Some(SyncDelivery::TransportError {
event_id: event_id.clone(),
relay_url: nostr_relay,
slot_id: String::new(),
detail: "nostr relay rejected the event (OK=false)".to_string(),
});
}
Err(e) => {
last_failure = Some(SyncDelivery::TransportError {
event_id: event_id.clone(),
relay_url: nostr_relay,
slot_id: String::new(),
detail: format!("nostr publish failed: {e:#}"),
});
}
}
}
Ok(last_failure.unwrap_or(SyncDelivery::PeerUnknown { event_id }))
}
fn deliver_over_nostr(
peer_npub_hex: &str,
relay_url: &str,
signed_event: &Value,
nsk: &[u8; 32],
) -> Result<bool> {
let ev = build_addressed_nostr(signed_event, nsk, peer_npub_hex)?;
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.context("build nostr runtime")?;
rt.block_on(async {
let mut ws = crate::nostr_ws::NostrWs::connect(relay_url)
.await
.with_context(|| format!("connect {relay_url}"))?;
ws.publish(&ev).await.context("publish over nostr")
})
}
fn build_addressed_nostr(
signed_event: &Value,
nsk: &[u8; 32],
peer_npub_hex: &str,
) -> Result<crate::nostr_event::NostrEvent> {
crate::nostr_event::wire_to_nostr_addressed(signed_event, nsk, peer_npub_hex)
.map_err(|e| anyhow::anyhow!("encode wire event as nostr: {e}"))
}
fn peer_unknown_reason(
peer: &str,
trusted: bool,
has_endpoint: bool,
has_usable_slot: bool,
) -> String {
if !trusted {
format!(
"peer '{peer}' is not pinned — run `wire dial {peer}@<relay>` to pair, or pass --queue (CLI) / queue:true (MCP) to buffer for the daemon to attempt later"
)
} else if !has_endpoint {
format!(
"peer '{peer}' IS pinned but has no relay endpoint recorded — re-register with a FULL `wire dial {peer}@<relay>` (the bare nickname reports `already_pinned` WITHOUT re-registering the slot)"
)
} else if !has_usable_slot {
format!(
"peer '{peer}' IS pinned but its relay slot has no token yet — their pair_drop_ack hasn't landed (common right after a daemon/MCP restart). Re-run the FULL `wire dial {peer}@<relay>` (NOT the bare nickname) to re-register, then resend"
)
} else {
format!(
"peer '{peer}' could not be reached on any recorded endpoint — check `wire status`, then re-dial `{peer}@<relay>`"
)
}
}
pub(crate) fn unsendable_reason(peer: &str) -> Option<String> {
let trust = crate::config::read_trust().unwrap_or_default();
let state = crate::config::read_relay_state().unwrap_or_default();
let trusted = trust.get("agents").and_then(|a| a.get(peer)).is_some();
let eps = crate::endpoints::peer_endpoints_in_priority_order(&state, peer);
let has_endpoint = !eps.is_empty();
let has_usable_slot = eps
.iter()
.any(|e| !e.relay_url.is_empty() && !e.slot_id.is_empty() && !e.slot_token.is_empty());
let nostr_reachable = crate::endpoints::peer_nostr_transport(&state, peer).is_some()
&& crate::config::read_nostr_key().is_ok();
if has_usable_slot || nostr_reachable {
None
} else {
Some(peer_unknown_reason(
peer,
trusted,
has_endpoint,
has_usable_slot,
))
}
}
pub fn delivery_json(d: &SyncDelivery, peer: &str) -> Value {
let base = json!({
"status": d.status_str(),
"peer": peer,
"event_id": d.event_id(),
});
let mut obj = base.as_object().cloned().unwrap_or_default();
match d {
SyncDelivery::Delivered {
relay_url, slot_id, ..
}
| SyncDelivery::Duplicate {
relay_url, slot_id, ..
} => {
obj.insert("relay_url".into(), json!(relay_url));
obj.insert("slot_id".into(), json!(slot_id));
}
SyncDelivery::DeliveredNostr {
relay_url, npub, ..
} => {
obj.insert("relay_url".into(), json!(relay_url));
obj.insert("transport".into(), json!("nostr"));
obj.insert("npub".into(), json!(npub));
}
SyncDelivery::SlotStale {
relay_url,
slot_id,
detail,
..
}
| SyncDelivery::TransportError {
relay_url,
slot_id,
detail,
..
} => {
obj.insert("relay_url".into(), json!(relay_url));
obj.insert("slot_id".into(), json!(slot_id));
obj.insert("reason".into(), json!(detail));
}
SyncDelivery::PeerUnknown { .. } => {
let reason = unsendable_reason(peer)
.unwrap_or_else(|| peer_unknown_reason(peer, true, true, true));
obj.insert("reason".into(), json!(reason));
}
}
Value::Object(obj)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn peer_unknown_reason_classifies_the_three_states() {
let r = peer_unknown_reason("p", false, false, false);
assert!(r.contains("is not pinned"), "{r}");
let r = peer_unknown_reason("p", true, false, false);
assert!(r.contains("IS pinned"), "{r}");
assert!(r.contains("no relay endpoint"), "{r}");
assert!(r.contains("@<relay>"), "{r}");
assert!(r.contains("bare nickname"), "{r}");
let r = peer_unknown_reason("p", true, true, false);
assert!(r.contains("no token yet"), "{r}");
assert!(r.contains("daemon/MCP restart"), "{r}");
assert!(r.contains("NOT the bare nickname"), "{r}");
let r = peer_unknown_reason("p", true, true, true);
assert!(r.contains("could not be reached"), "{r}");
}
#[test]
fn unsendable_reason_reads_live_state() {
use crate::endpoints::{Endpoint, EndpointScope, pin_peer_endpoints};
crate::config::test_support::with_temp_home(|| {
let r = unsendable_reason("ghost").expect("unknown peer is unsendable");
assert!(r.contains("is not pinned"), "{r}");
let mut st = crate::config::read_relay_state().unwrap();
pin_peer_endpoints(
&mut st,
"live",
&[Endpoint {
relay_url: "https://wireup.net".into(),
slot_id: "slot-live".into(),
slot_token: "tok-abc".into(),
scope: EndpointScope::Federation,
}],
)
.unwrap();
crate::config::write_relay_state(&st).unwrap();
assert!(
unsendable_reason("live").is_none(),
"peer with a non-empty slot_token must read as sendable"
);
let mut st2 = crate::config::read_relay_state().unwrap();
pin_peer_endpoints(
&mut st2,
"pending",
&[Endpoint {
relay_url: "https://wireup.net".into(),
slot_id: "slot-pending".into(),
slot_token: String::new(),
scope: EndpointScope::Federation,
}],
)
.unwrap();
crate::config::write_relay_state(&st2).unwrap();
crate::config::update_trust(|t| {
t.get_mut("agents")
.and_then(Value::as_object_mut)
.unwrap()
.insert(
"pending".into(),
json!({"did": "did:wire:pending-0000", "tier": "VERIFIED"}),
);
Ok(())
})
.unwrap();
let r = unsendable_reason("pending").expect("empty-token peer is unsendable");
assert!(r.contains("no token yet"), "{r}");
let mut st3 = crate::config::read_relay_state().unwrap();
pin_peer_endpoints(
&mut st3,
"nostronly",
&[Endpoint {
relay_url: "https://wireup.net".into(),
slot_id: "slot-n".into(),
slot_token: String::new(),
scope: EndpointScope::Federation,
}],
)
.unwrap();
st3["peers"]["nostronly"]["nostr_transport"] =
json!({"npub": "npub1xxx", "relay": "wss://relay.example"});
crate::config::write_relay_state(&st3).unwrap();
crate::config::write_nostr_key(&[3u8; 32]).unwrap();
assert!(
unsendable_reason("nostronly").is_none(),
"a Nostr-reachable peer must read as sendable despite an empty HTTP token"
);
let mut st4 = crate::config::read_relay_state().unwrap();
pin_peer_endpoints(
&mut st4,
"malformed",
&[Endpoint {
relay_url: String::new(),
slot_id: String::new(),
slot_token: "tok-orphan".into(),
scope: EndpointScope::Federation,
}],
)
.unwrap();
crate::config::write_relay_state(&st4).unwrap();
assert!(
unsendable_reason("malformed").is_some(),
"a token on an endpoint with empty relay_url/slot_id is not usable"
);
});
}
#[test]
fn status_str_matches_variant() {
let d = SyncDelivery::Delivered {
event_id: "x".into(),
relay_url: "https://r".into(),
slot_id: "s".into(),
};
assert_eq!(d.status_str(), "delivered");
assert!(d.reached_relay());
let d = SyncDelivery::Duplicate {
event_id: "x".into(),
relay_url: "https://r".into(),
slot_id: "s".into(),
};
assert_eq!(d.status_str(), "duplicate");
assert!(
d.reached_relay(),
"duplicate counts as relay-reached: peer can pull it"
);
let d = SyncDelivery::PeerUnknown {
event_id: "x".into(),
};
assert_eq!(d.status_str(), "peer_unknown");
assert!(!d.reached_relay());
let d = SyncDelivery::SlotStale {
event_id: "x".into(),
relay_url: "https://r".into(),
slot_id: "s".into(),
detail: "410".into(),
};
assert_eq!(d.status_str(), "slot_stale");
assert!(!d.reached_relay());
let d = SyncDelivery::TransportError {
event_id: "x".into(),
relay_url: "https://r".into(),
slot_id: "s".into(),
detail: "tls".into(),
};
assert_eq!(d.status_str(), "transport_error");
assert!(!d.reached_relay());
}
#[test]
fn delivered_nostr_counts_as_reached_and_renders_transport() {
let d = SyncDelivery::DeliveredNostr {
event_id: "ev1".into(),
relay_url: "wss://relay.example".into(),
npub: "ab".repeat(32),
};
assert_eq!(d.status_str(), "delivered_nostr");
assert!(
d.reached_relay(),
"nostr delivery means the peer can pull it"
);
assert_eq!(d.event_id(), "ev1");
let v = delivery_json(&d, "alice");
assert_eq!(v["status"], "delivered_nostr");
assert_eq!(v["peer"], "alice");
assert_eq!(v["event_id"], "ev1");
assert_eq!(v["relay_url"], "wss://relay.example");
assert_eq!(v["transport"], "nostr");
assert_eq!(v["npub"], "ab".repeat(32));
assert!(v.get("slot_id").is_none(), "nostr send has no slot_id");
assert!(v.get("reason").is_none(), "success has no reason");
}
#[test]
fn build_addressed_nostr_is_verifiable_and_addressed() {
use crate::nostr_key::generate_transport_key;
use crate::signing::{generate_keypair, sign_message_v31};
let (sk, pk) = generate_keypair();
let wire = sign_message_v31(
&json!({
"v": "3.1",
"timestamp": "2026-06-14T12:00:00Z",
"from": "did:wire:slate-lotus-88232017",
"to": "did:wire:raven-kettle-1234",
"kind": 1,
"body": {"content": "routed over nostr"},
}),
&sk,
&pk,
"slate-lotus",
)
.unwrap();
let (nsk, _x) = generate_transport_key();
let (_psk, peer_x) = generate_transport_key();
let peer_hex = hex::encode(peer_x);
let ev = build_addressed_nostr(&wire, &nsk, &peer_hex).unwrap();
assert!(
ev.tags
.iter()
.any(|t| t.first().map(String::as_str) == Some("p") && t.get(1) == Some(&peer_hex))
);
assert_eq!(crate::nostr_event::verify_and_decode(&ev).unwrap(), wire);
}
#[test]
fn delivery_json_includes_reason_only_for_failures() {
let ok = SyncDelivery::Delivered {
event_id: "abc".into(),
relay_url: "https://r".into(),
slot_id: "s".into(),
};
let v = delivery_json(&ok, "alice");
assert_eq!(v["status"], "delivered");
assert_eq!(v["event_id"], "abc");
assert_eq!(v["peer"], "alice");
assert_eq!(v["relay_url"], "https://r");
assert!(v.get("reason").is_none(), "happy path has no reason field");
let bad = SyncDelivery::TransportError {
event_id: "abc".into(),
relay_url: "https://r".into(),
slot_id: "s".into(),
detail: "TLS error: UnknownIssuer".into(),
};
let v = delivery_json(&bad, "alice");
assert_eq!(v["status"], "transport_error");
assert_eq!(v["reason"], "TLS error: UnknownIssuer");
let unknown = SyncDelivery::PeerUnknown {
event_id: "abc".into(),
};
let v = delivery_json(&unknown, "alice");
assert_eq!(v["status"], "peer_unknown");
assert!(
v["reason"]
.as_str()
.unwrap_or("")
.contains("wire dial alice")
);
}
}