onionwire/src/node.rs

396 lines
13 KiB
Rust
Raw Normal View History

//! Two-node onion session: host HS, Noise IK, persist plaintext locally.
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use futures::StreamExt;
use futures::io::{AsyncRead, AsyncWrite};
use tor_cell::relaycell::msg::Connected;
use tor_hsservice::{RunningOnionService, handle_rend_requests};
use crate::frame;
use crate::hs::{self, Client, HS_PORT};
use crate::loc;
use crate::qr;
use crate::session::{self, Keys};
use crate::store::{Friend, Message, Store};
pub struct RotateResult {
pub notified: usize,
pub friends: usize,
}
struct HsHandle {
_svc: Arc<RunningOnionService>,
rend: tokio::task::JoinHandle<()>,
}
impl Drop for HsHandle {
fn drop(&mut self) {
self.rend.abort();
}
}
pub struct Node {
home: PathBuf,
store: Mutex<Store>,
client: Client,
hs: Mutex<Option<HsHandle>>,
onion: Mutex<String>,
keys: Keys,
}
impl Node {
pub async fn start(home: PathBuf) -> Result<Arc<Self>, String> {
let store = Store::open_at(&home).map_err(|e| e.to_string())?;
let me = store.self_identity().map_err(|e| e.to_string())?;
let keys = Keys::from_self(&me).map_err(|e| e.to_string())?;
let nickname = store.hs_nickname().map_err(|e| e.to_string())?;
let state = home.join("arti");
let cache = home.join("cache");
let client = hs::bootstrapped(&state, &cache).await?;
let launched = client
.launch_onion_service(hs::hs_config(&nickname)?)
.map_err(|e| format!("launch_onion_service: {e}"))?
.ok_or_else(|| "onion service disabled in config — fail closed".to_string())?;
let (svc, rend) = launched;
let onion = hs::onion_string(&svc)?;
hs::wait_until_published(&svc, &onion).await?;
store.set_onion(&onion).map_err(|e| e.to_string())?;
let node = Arc::new(Self {
home,
store: Mutex::new(store),
client,
hs: Mutex::new(None),
onion: Mutex::new(onion),
keys,
});
let rend = spawn_rend(Arc::clone(&node), rend);
*node.hs.lock().map_err(|e| e.to_string())? = Some(HsHandle { _svc: svc, rend });
Ok(node)
}
pub fn onion(&self) -> String {
self.onion.lock().map(|g| g.clone()).unwrap_or_default()
}
pub fn identity_pk(&self) -> [u8; 32] {
self.keys.identity_pk
}
pub fn arti_dir(&self) -> PathBuf {
self.home.join("arti")
}
pub fn cache_dir(&self) -> PathBuf {
self.home.join("cache")
}
pub fn qr_payload(&self) -> Result<String, String> {
qr::encode(&self.keys.identity_sk, &self.onion(), &self.keys.prekey_pk)
.map_err(|e| e.to_string())
}
pub fn add_friend_from_qr(&self, raw: &str) -> Result<(), String> {
let p = qr::decode(raw).map_err(|e| e.to_string())?;
self.add_friend_payload(&p)
}
pub fn add_friend_payload(&self, p: &qr::QrPayload) -> Result<(), String> {
let store = self.store.lock().map_err(|e| e.to_string())?;
store
.upsert_friend(&p.pubkey, &p.onion, None)
.map_err(|e| e.to_string())?;
store
.set_friend_prekey(&p.pubkey, &p.signed_prekey)
.map_err(|e| e.to_string())?;
Ok(())
}
pub fn get_friend(&self, pubkey: &[u8]) -> Result<Option<Friend>, String> {
self.store
.lock()
.map_err(|e| e.to_string())?
.get_friend(pubkey)
.map_err(|e| e.to_string())
}
pub fn list_friends(&self) -> Result<Vec<Friend>, String> {
self.store
.lock()
.map_err(|e| e.to_string())?
.list_friends()
.map_err(|e| e.to_string())
}
pub fn friend(&self, pubkey: &[u8]) -> Result<Friend, String> {
self.store
.lock()
.map_err(|e| e.to_string())?
.get_friend(pubkey)
.map_err(|e| e.to_string())?
.ok_or_else(|| "friend not found".into())
}
pub fn list_messages(&self, friend_pk: &[u8]) -> Result<Vec<Message>, String> {
self.store
.lock()
.map_err(|e| e.to_string())?
.list_messages(friend_pk)
.map_err(|e| e.to_string())
}
pub fn wipe_messages(&self) -> Result<(), String> {
self.store
.lock()
.map_err(|e| e.to_string())?
.wipe_messages()
.map_err(|e| e.to_string())
}
pub async fn connect_onion(&self, onion: &str) -> Result<(), String> {
tokio::time::timeout(
Duration::from_secs(20),
self.client.connect((onion, HS_PORT)),
)
.await
.map_err(|_| format!("connect {onion}:{HS_PORT} timed out"))?
.map_err(|e| format!("connect {onion}:{HS_PORT}: {e}"))?;
Ok(())
}
pub async fn rotate(self: &Arc<Self>) -> Result<RotateResult, String> {
let old_nick = self
.store
.lock()
.map_err(|e| e.to_string())?
.hs_nickname()
.map_err(|e| e.to_string())?;
let new_nick = next_hs_nickname(&old_nick);
let launched = self
.client
.launch_onion_service(hs::hs_config(&new_nick)?)
.map_err(|e| format!("launch_onion_service: {e}"))?
.ok_or_else(|| "onion service disabled in config — fail closed".to_string())?;
let (svc, rend) = launched;
let onion = hs::onion_string(&svc)?;
// ponytail: hard-cut old HS before waiting; keeping both stalled ow1 at Bootstrapping.
*self.hs.lock().map_err(|e| e.to_string())? = None;
hs::wait_until_published(&svc, &onion).await?;
let rend = spawn_rend(Arc::clone(self), rend);
*self.hs.lock().map_err(|e| e.to_string())? = Some(HsHandle { _svc: svc, rend });
{
let store = self.store.lock().map_err(|e| e.to_string())?;
store.set_onion(&onion).map_err(|e| e.to_string())?;
store
.set_hs_nickname(&new_nick)
.map_err(|e| e.to_string())?;
}
*self.onion.lock().map_err(|e| e.to_string())? = onion.clone();
let ts = unix_now();
let loc = loc::sign(&self.keys.identity_sk, &onion, ts).map_err(|e| e.to_string())?;
let loc_pt = loc::encode(&loc);
let friends = self
.store
.lock()
.map_err(|e| e.to_string())?
.list_friends()
.map_err(|e| e.to_string())?;
let mut notified = 0;
for f in &friends {
if self
.push_loc(&f.pubkey, &f.onion, &f.prekey, &loc_pt)
.await
.is_ok()
{
notified += 1;
}
}
Ok(RotateResult {
notified,
friends: friends.len(),
})
}
pub async fn send(&self, friend_pk: &[u8], plaintext: &[u8]) -> Result<(), String> {
let (onion, prekey) = {
let store = self.store.lock().map_err(|e| e.to_string())?;
let f = store
.get_friend(friend_pk)
.map_err(|e| e.to_string())?
.ok_or_else(|| "unknown friend".to_string())?;
if f.prekey.len() != 32 {
return Err("friend missing prekey".into());
}
(f.onion, f.prekey)
};
let deadline = Instant::now() + Duration::from_secs(180);
let mut last = None::<String>;
while Instant::now() < deadline {
match self.try_send(&onion, friend_pk, &prekey, plaintext).await {
Ok(()) => {
self.store
.lock()
.map_err(|e| e.to_string())?
.append_message(friend_pk, "out", plaintext)
.map_err(|e| e.to_string())?;
return Ok(());
}
Err(e) if e.contains("fingerprint mismatch") => return Err(e),
Err(e) => last = Some(e),
}
tokio::time::sleep(Duration::from_secs(3)).await;
}
Err(last.unwrap_or_else(|| "send timed out".into()))
}
async fn push_loc(
&self,
friend_pk: &[u8],
onion: &str,
prekey: &[u8],
loc_pt: &[u8],
) -> Result<(), String> {
if prekey.len() != 32 {
return Err("friend missing prekey".into());
}
let deadline = Instant::now() + Duration::from_secs(180);
let mut last = None::<String>;
while Instant::now() < deadline {
match self.try_send(onion, friend_pk, prekey, loc_pt).await {
Ok(()) => return Ok(()),
Err(e) if e.contains("fingerprint mismatch") => return Err(e),
Err(e) => last = Some(e),
}
tokio::time::sleep(Duration::from_secs(3)).await;
}
Err(last.unwrap_or_else(|| "loc push timed out".into()))
}
async fn try_send(
&self,
onion: &str,
pinned_id: &[u8],
remote_prekey: &[u8],
plaintext: &[u8],
) -> Result<(), String> {
let mut stream = self
.client
.connect((onion, HS_PORT))
.await
.map_err(|e| format!("connect {onion}:{HS_PORT}: {e}"))?;
let mut sess =
session::handshake_initiator(&mut stream, &self.keys, pinned_id, remote_prekey)
.await
.map_err(session_err)?;
let ct = sess.encrypt(plaintext).map_err(session_err)?;
frame::write_frame(&mut stream, &ct)
.await
.map_err(|e| e.to_string())?;
Ok(())
}
async fn handle_incoming<S>(&self, stream: &mut S) -> Result<(), String>
where
S: AsyncRead + AsyncWrite + Unpin,
{
let store = &self.store;
let mut sess = session::handshake_responder(stream, &self.keys, |spk| {
let Ok(g) = store.lock() else {
return None;
};
g.get_friend_by_prekey(spk).ok().flatten().map(|f| f.pubkey)
})
.await
.map_err(session_err)?;
let ct = frame::read_frame(stream).await.map_err(|e| e.to_string())?;
let pt = sess.decrypt(&ct).map_err(session_err)?;
if let Some(loc) = loc::decode(&pt) {
let applied = store.lock().map_err(|e| e.to_string())?.apply_loc(
&sess.peer_identity,
&loc.onion,
loc.ts,
&loc.sig,
);
match applied {
Ok(true) => {}
Ok(false) => eprintln!("loc dropped (bad sig, stale ts, or unknown friend)"),
Err(e) => return Err(e.to_string()),
}
return Ok(());
}
store
.lock()
.map_err(|e| e.to_string())?
.append_message(&sess.peer_identity, "in", &pt)
.map_err(|e| e.to_string())?;
Ok(())
}
}
fn spawn_rend(
node: Arc<Node>,
rend: impl futures::Stream<Item = tor_hsservice::RendRequest> + Send + 'static,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut requests = std::pin::pin!(handle_rend_requests(rend));
while let Some(req) = requests.next().await {
let serve = Arc::clone(&node);
tokio::spawn(async move {
let Ok(mut stream) = req.accept(Connected::new_empty()).await else {
return;
};
if let Err(e) = serve.handle_incoming(&mut stream).await {
eprintln!("incoming: {e}");
}
});
}
})
}
fn next_hs_nickname(cur: &str) -> String {
let n = cur
.strip_prefix("ow")
.and_then(|s| s.parse::<u32>().ok())
.unwrap_or(0);
format!("ow{}", n.saturating_add(1))
}
fn unix_now() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs() as i64)
.unwrap_or(0)
}
fn session_err(e: session::Error) -> String {
if e.is_fingerprint_mismatch() {
crate::tui::fingerprint_mismatch_banner().to_string()
} else {
e.to_string()
}
}
pub fn dir_contains_bytes(dir: &Path, needle: &[u8]) -> bool {
let mut stack = vec![dir.to_path_buf()];
while let Some(p) = stack.pop() {
let Ok(rd) = std::fs::read_dir(&p) else {
continue;
};
for ent in rd.flatten() {
let path = ent.path();
if path.is_dir() {
stack.push(path);
} else if let Ok(bytes) = std::fs::read(&path)
&& bytes.windows(needle.len()).any(|w| w == needle)
{
return true;
}
}
}
false
}