216 lines
7 KiB
Rust
216 lines
7 KiB
Rust
|
|
//! 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};
|
||
|
|
|
||
|
|
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::qr;
|
||
|
|
use crate::session::{self, Keys};
|
||
|
|
use crate::store::{Message, Store};
|
||
|
|
|
||
|
|
pub struct Node {
|
||
|
|
home: PathBuf,
|
||
|
|
store: Mutex<Store>,
|
||
|
|
client: Client,
|
||
|
|
_svc: Arc<RunningOnionService>,
|
||
|
|
onion: 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 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("onionwire")?)
|
||
|
|
.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,
|
||
|
|
_svc: svc,
|
||
|
|
onion,
|
||
|
|
keys,
|
||
|
|
});
|
||
|
|
let serve = Arc::clone(&node);
|
||
|
|
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(&serve);
|
||
|
|
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}");
|
||
|
|
}
|
||
|
|
});
|
||
|
|
}
|
||
|
|
});
|
||
|
|
Ok(node)
|
||
|
|
}
|
||
|
|
|
||
|
|
pub fn onion(&self) -> &str {
|
||
|
|
&self.onion
|
||
|
|
}
|
||
|
|
|
||
|
|
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())?;
|
||
|
|
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 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 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 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)?;
|
||
|
|
store
|
||
|
|
.lock()
|
||
|
|
.map_err(|e| e.to_string())?
|
||
|
|
.append_message(&sess.peer_identity, "in", &pt)
|
||
|
|
.map_err(|e| e.to_string())?;
|
||
|
|
Ok(())
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
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
|
||
|
|
}
|