core: LAN transport — serverless same-network sync over mDNS + TCP
Add a third transport behind the Transport trait: each engine advertises _tether._tcp.local. (TXT room/from/source) and browses for peers in the same room, then connects directly over plain TCP (lower `from` dials to avoid double-connects) exchanging newline-delimited JSON Messages. No server, no internet — the first truly serverless transport, complementing SSE (relay) and RTC (P2P-after-signaling). Engine dedup already collapses a clipboard arriving over LAN + SSE/RTC. Proven by examples/lan_p2p.rs: with a DEAD server URL (SSE/RTC can't deliver), two engines discover each other via mDNS and sync over TCP. Scope/caveats: room id is the only access control and the link is plaintext (rustls keyed on a room secret is the follow-up); iOS needs the multicast entitlement for mDNS, so the transport is inert there until that lands. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -16,7 +16,8 @@ same wire format.
|
||||
┌──────────────── tethercore (Rust) ────────────────┐
|
||||
│ Message codec · Transport trait · dedup │
|
||||
│ ├─ SSE (reqwest) reliable relay floor │
|
||||
│ └─ RTC (webrtc-rs) direct P2P data channel │
|
||||
│ ├─ RTC (webrtc-rs) direct P2P data channel │
|
||||
│ └─ LAN (mDNS + TCP) serverless same-network │
|
||||
└───────────────────────────────────────────────────┘
|
||||
iOS (Swift) ─ macOS (Swift) ─ agent (Rust, Linux/Win/Mac) ─ Android (Kotlin)
|
||||
└──── all via UniFFI bindings ────┘
|
||||
|
||||
64
core/Cargo.lock
generated
64
core/Cargo.lock
generated
@@ -712,6 +712,12 @@ dependencies = [
|
||||
"windows-sys 0.61.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "fastrand"
|
||||
version = "2.4.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9f1f227452a390804cdb637b74a86990f2a7d7ba4b7d5693aac9b4dd6defd8d6"
|
||||
|
||||
[[package]]
|
||||
name = "ff"
|
||||
version = "0.13.1"
|
||||
@@ -734,6 +740,17 @@ version = "0.1.9"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5baebc0774151f905a1a2cc41989300b1e6fbb29aff0ceffa1064fdd3088d582"
|
||||
|
||||
[[package]]
|
||||
name = "flume"
|
||||
version = "0.11.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095"
|
||||
dependencies = [
|
||||
"futures-core",
|
||||
"futures-sink",
|
||||
"spin",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "form_urlencoded"
|
||||
version = "1.2.2"
|
||||
@@ -1164,6 +1181,16 @@ dependencies = [
|
||||
"icu_properties",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "if-addrs"
|
||||
version = "0.15.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "c0a05c691e1fae256cf7013d99dad472dc52d5543322761f83ec8d47eab40d2b"
|
||||
dependencies = [
|
||||
"libc",
|
||||
"windows-sys 0.61.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "indexmap"
|
||||
version = "2.14.0"
|
||||
@@ -1283,6 +1310,21 @@ dependencies = [
|
||||
"digest",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "mdns-sd"
|
||||
version = "0.20.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "892f96f6d2ebe1ea641279f986ac52a2a6bac71e8f743bb258315cfe2bd7e88e"
|
||||
dependencies = [
|
||||
"fastrand",
|
||||
"flume",
|
||||
"if-addrs",
|
||||
"log",
|
||||
"mio",
|
||||
"socket-pktinfo",
|
||||
"socket2 0.6.4",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "memchr"
|
||||
version = "2.8.2"
|
||||
@@ -1327,6 +1369,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "02bd0af71c67b473010cbbc60715ee815645a4dc942899111f494b4b737d6fda"
|
||||
dependencies = [
|
||||
"libc",
|
||||
"log",
|
||||
"wasi",
|
||||
"windows-sys 0.61.2",
|
||||
]
|
||||
@@ -2212,6 +2255,17 @@ dependencies = [
|
||||
"serde",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "socket-pktinfo"
|
||||
version = "0.3.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "927136cc2ae6a1b0e66ac6b1210902b75c3f726db004a73bc18686dcd0dcd22f"
|
||||
dependencies = [
|
||||
"libc",
|
||||
"socket2 0.6.4",
|
||||
"windows-sys 0.60.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "socket2"
|
||||
version = "0.5.10"
|
||||
@@ -2232,6 +2286,15 @@ dependencies = [
|
||||
"windows-sys 0.61.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "spin"
|
||||
version = "0.9.8"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67"
|
||||
dependencies = [
|
||||
"lock_api",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "spki"
|
||||
version = "0.7.3"
|
||||
@@ -2330,6 +2393,7 @@ name = "tethercore"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"futures-util",
|
||||
"mdns-sd",
|
||||
"reqwest",
|
||||
"serde",
|
||||
"serde_json",
|
||||
|
||||
@@ -20,6 +20,7 @@ uniffi = { version = "0.28", features = ["cli"] }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "stream", "json"] }
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "time", "sync", "macros"] }
|
||||
tokio = { version = "1", features = ["rt-multi-thread", "time", "sync", "macros", "net", "io-util"] }
|
||||
futures-util = "0.3"
|
||||
webrtc = "0.17.1"
|
||||
mdns-sd = "0.20.0"
|
||||
|
||||
55
core/examples/lan_p2p.rs
Normal file
55
core/examples/lan_p2p.rs
Normal file
@@ -0,0 +1,55 @@
|
||||
//! Proves serverless LAN sync: two engines with a DEAD server URL (so SSE and
|
||||
//! RTC cannot deliver) discover each other over mDNS and exchange a clipboard
|
||||
//! over plain TCP.
|
||||
//! cargo run --example lan_p2p
|
||||
//!
|
||||
//! If the message arrives, it can only have come via the LAN transport.
|
||||
|
||||
use std::sync::mpsc;
|
||||
use std::time::Duration;
|
||||
|
||||
use tethercore::{Engine, Message, MessageHandler};
|
||||
|
||||
struct Collector(mpsc::Sender<Message>);
|
||||
impl MessageHandler for Collector {
|
||||
fn on_message(&self, msg: Message) {
|
||||
let _ = self.0.send(msg);
|
||||
}
|
||||
fn on_status(&self, _connected: bool) {}
|
||||
}
|
||||
|
||||
fn main() {
|
||||
// Unroutable server → SSE/RTC are dead; only LAN can carry the message.
|
||||
let dead = "http://127.0.0.1:1".to_string();
|
||||
let room = "lan-spike".to_string();
|
||||
|
||||
let (tx, rx) = mpsc::channel();
|
||||
let a = Engine::new(dead.clone(), room.clone(), "peer-a".into(), "lan-a".into());
|
||||
a.start(Box::new(Collector(tx)));
|
||||
|
||||
let b_from = "peer-b";
|
||||
let b = Engine::new(dead, room, b_from.into(), "lan-b".into());
|
||||
b.start(Box::new(Collector(mpsc::channel().0)));
|
||||
|
||||
eprintln!("waiting for mDNS discovery + TCP connect…");
|
||||
std::thread::sleep(Duration::from_secs(4));
|
||||
|
||||
let payload = "hello-over-lan";
|
||||
b.send(payload.to_string());
|
||||
|
||||
match rx.recv_timeout(Duration::from_secs(6)) {
|
||||
Ok(m) if m.text == payload => {
|
||||
println!(
|
||||
"✅ serverless LAN delivery: text={:?} from={} source={}",
|
||||
m.text, m.from, m.source
|
||||
);
|
||||
}
|
||||
Ok(m) => println!("⚠️ unexpected: {m:?}"),
|
||||
Err(_) => {
|
||||
println!("❌ nothing received — mDNS may be blocked, or peers didn't connect");
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
a.stop();
|
||||
b.stop();
|
||||
}
|
||||
221
core/src/lan.rs
Normal file
221
core/src/lan.rs
Normal file
@@ -0,0 +1,221 @@
|
||||
//! LAN transport — serverless same-network clipboard over mDNS + plain TCP.
|
||||
//!
|
||||
//! Each engine advertises `_tether._tcp.local.` with TXT {room, from, source}
|
||||
//! and browses for the same service. Peers in the *same room* connect directly
|
||||
//! over TCP (lower `from` dials, to avoid double-connects) and exchange
|
||||
//! newline-delimited JSON `Message`s. No server, no internet — just the LAN.
|
||||
//!
|
||||
//! Scope: room id is the only access control and the link is plaintext; fine
|
||||
//! for a trusted LAN, a rustls layer keyed on a room secret is the follow-up.
|
||||
//! iOS needs the multicast entitlement for mDNS — there this transport is inert.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::{Arc, Mutex as StdMutex};
|
||||
|
||||
use mdns_sd::{ServiceDaemon, ServiceEvent, ServiceInfo};
|
||||
use serde_json::json;
|
||||
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
||||
use tokio::net::tcp::OwnedWriteHalf;
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
use tokio::runtime::Handle;
|
||||
use tokio::sync::Mutex;
|
||||
|
||||
use crate::{Config, Message, Sink, Status, Transport};
|
||||
|
||||
const SERVICE: &str = "_tether._tcp.local.";
|
||||
|
||||
pub(crate) struct LanTransport {
|
||||
cfg: Config,
|
||||
peers: Arc<Mutex<HashMap<String, OwnedWriteHalf>>>, // peer `from` → its write half
|
||||
sink: StdMutex<Option<Sink>>,
|
||||
}
|
||||
|
||||
impl LanTransport {
|
||||
pub(crate) fn new(cfg: Config) -> Self {
|
||||
Self {
|
||||
cfg,
|
||||
peers: Arc::new(Mutex::new(HashMap::new())),
|
||||
sink: StdMutex::new(None),
|
||||
}
|
||||
}
|
||||
|
||||
fn sink(&self) -> Option<Sink> {
|
||||
self.sink.lock().unwrap().clone()
|
||||
}
|
||||
}
|
||||
|
||||
impl Transport for LanTransport {
|
||||
fn name(&self) -> &'static str {
|
||||
"lan"
|
||||
}
|
||||
|
||||
fn start(self: Arc<Self>, rt: Handle, running: Arc<AtomicBool>, sink: Sink, _status: Status) {
|
||||
*self.sink.lock().unwrap() = Some(sink);
|
||||
let this = self;
|
||||
rt.spawn(async move {
|
||||
if let Err(e) = this.run(running).await {
|
||||
eprintln!("tethercore[lan]: {e}");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
fn publish(&self, rt: &Handle, msg: Message) {
|
||||
let peers = self.peers.clone();
|
||||
rt.spawn(async move {
|
||||
let line = match serde_json::to_string(&msg) {
|
||||
Ok(s) => format!("{s}\n"),
|
||||
Err(_) => return,
|
||||
};
|
||||
let mut map = peers.lock().await;
|
||||
let mut dead = Vec::new();
|
||||
for (id, w) in map.iter_mut() {
|
||||
if w.write_all(line.as_bytes()).await.is_err() {
|
||||
dead.push(id.clone());
|
||||
}
|
||||
}
|
||||
for id in dead {
|
||||
map.remove(&id);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
impl LanTransport {
|
||||
async fn run(self: &Arc<Self>, running: Arc<AtomicBool>) -> Result<(), String> {
|
||||
let listener = TcpListener::bind("0.0.0.0:0")
|
||||
.await
|
||||
.map_err(|e| e.to_string())?;
|
||||
let port = listener.local_addr().map_err(|e| e.to_string())?.port();
|
||||
|
||||
// Advertise + browse over mDNS.
|
||||
let daemon = ServiceDaemon::new().map_err(|e| e.to_string())?;
|
||||
let instance = sanitize(&self.cfg.from);
|
||||
let props = [
|
||||
("room", self.cfg.room.as_str()),
|
||||
("from", self.cfg.from.as_str()),
|
||||
("source", self.cfg.source.as_str()),
|
||||
];
|
||||
let info = ServiceInfo::new(
|
||||
SERVICE,
|
||||
&instance,
|
||||
&format!("{instance}.local."),
|
||||
"",
|
||||
port,
|
||||
&props[..],
|
||||
)
|
||||
.map_err(|e| e.to_string())?
|
||||
.enable_addr_auto();
|
||||
daemon.register(info).map_err(|e| e.to_string())?;
|
||||
let browse = daemon.browse(SERVICE).map_err(|e| e.to_string())?;
|
||||
|
||||
// Accept inbound dials.
|
||||
{
|
||||
let this = self.clone();
|
||||
let running = running.clone();
|
||||
tokio::spawn(async move {
|
||||
while running.load(Ordering::SeqCst) {
|
||||
if let Ok((stream, _)) = listener.accept().await {
|
||||
let this = this.clone();
|
||||
tokio::spawn(async move { this.handle(stream).await });
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// Discover peers; the lower `from` dials.
|
||||
while running.load(Ordering::SeqCst) {
|
||||
match browse.recv_async().await {
|
||||
Ok(ServiceEvent::ServiceResolved(info)) => {
|
||||
let room = info.get_property_val_str("room").unwrap_or("");
|
||||
let from = info.get_property_val_str("from").unwrap_or("").to_string();
|
||||
if room != self.cfg.room || from.is_empty() || from == self.cfg.from {
|
||||
continue;
|
||||
}
|
||||
if self.cfg.from.as_str() >= from.as_str() {
|
||||
continue; // higher id waits to be dialed
|
||||
}
|
||||
if self.peers.lock().await.contains_key(&from) {
|
||||
continue;
|
||||
}
|
||||
// Prefer IPv4; fall back to any resolved address.
|
||||
let ip = info
|
||||
.get_addresses_v4()
|
||||
.into_iter()
|
||||
.next()
|
||||
.map(std::net::IpAddr::V4)
|
||||
.or_else(|| info.get_addresses().iter().next().map(|s| s.to_ip_addr()));
|
||||
if let Some(addr) = ip {
|
||||
let sockaddr = SocketAddr::new(addr, info.get_port());
|
||||
if let Ok(stream) = TcpStream::connect(sockaddr).await {
|
||||
let this = self.clone();
|
||||
tokio::spawn(async move { this.handle(stream).await });
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(_) => {}
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Handle one connection (dialed or accepted): hello handshake, then read
|
||||
/// `Message` lines into the sink. The write half is stored for `publish`.
|
||||
async fn handle(self: &Arc<Self>, stream: TcpStream) {
|
||||
let (rd, mut wr) = stream.into_split();
|
||||
|
||||
let hello = format!(
|
||||
"{}\n",
|
||||
json!({"t": "hello", "from": self.cfg.from, "room": self.cfg.room})
|
||||
);
|
||||
if wr.write_all(hello.as_bytes()).await.is_err() {
|
||||
return;
|
||||
}
|
||||
|
||||
let mut lines = BufReader::new(rd).lines();
|
||||
let peer_from = match lines.next_line().await {
|
||||
Ok(Some(l)) => match serde_json::from_str::<serde_json::Value>(&l) {
|
||||
Ok(v)
|
||||
if v["t"] == "hello"
|
||||
&& v["room"].as_str() == Some(self.cfg.room.as_str()) =>
|
||||
{
|
||||
v["from"].as_str().unwrap_or("").to_string()
|
||||
}
|
||||
_ => return,
|
||||
},
|
||||
_ => return,
|
||||
};
|
||||
if peer_from.is_empty() || peer_from == self.cfg.from {
|
||||
return;
|
||||
}
|
||||
{
|
||||
let mut map = self.peers.lock().await;
|
||||
if map.contains_key(&peer_from) {
|
||||
return; // already connected to this peer
|
||||
}
|
||||
map.insert(peer_from.clone(), wr);
|
||||
}
|
||||
eprintln!("tethercore[lan]: peer {peer_from} connected");
|
||||
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
if let Ok(m) = serde_json::from_str::<Message>(&line) {
|
||||
if m.from != self.cfg.from {
|
||||
if let Some(sink) = self.sink() {
|
||||
sink(m);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
self.peers.lock().await.remove(&peer_from);
|
||||
eprintln!("tethercore[lan]: peer {peer_from} disconnected");
|
||||
}
|
||||
}
|
||||
|
||||
fn sanitize(s: &str) -> String {
|
||||
s.chars()
|
||||
.map(|c| if c.is_ascii_alphanumeric() { c } else { '-' })
|
||||
.collect()
|
||||
}
|
||||
@@ -22,6 +22,7 @@ use futures_util::StreamExt;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::runtime::Handle;
|
||||
|
||||
mod lan;
|
||||
mod rtc;
|
||||
|
||||
uniffi::setup_scaffolding!();
|
||||
@@ -107,6 +108,7 @@ impl Engine {
|
||||
http: reqwest::Client::new(),
|
||||
}),
|
||||
Arc::new(rtc::RtcTransport::new(cfg.clone())),
|
||||
Arc::new(lan::LanTransport::new(cfg.clone())),
|
||||
];
|
||||
Arc::new(Self {
|
||||
cfg,
|
||||
|
||||
Reference in New Issue
Block a user