diff --git a/core/examples/rtc_p2p.rs b/core/examples/rtc_p2p.rs new file mode 100644 index 0000000..5a99ce1 --- /dev/null +++ b/core/examples/rtc_p2p.rs @@ -0,0 +1,61 @@ +//! Proves two engines establish a WebRTC data channel and clipboard flows P2P. +//! cargo run --example rtc_p2p -- http://192.168.11.233:8765 +//! +//! Path discriminator: an RTC-delivered message has source == the peer's +//! from-id; the SSE copy has source == the peer's source-label. Since RTC is +//! direct (no server round-trip) it should win the dedup race, so a received +//! source equal to the from-id proves delivery came over the data channel. + +use std::sync::mpsc; +use std::time::Duration; + +use tethercore::{Engine, Message, MessageHandler}; + +struct Collector(mpsc::Sender); +impl MessageHandler for Collector { + fn on_message(&self, msg: Message) { + let _ = self.0.send(msg); + } + fn on_status(&self, _connected: bool) {} +} + +fn main() { + let server = std::env::args() + .nth(1) + .unwrap_or_else(|| "http://192.168.11.233:8765".into()); + let room = "rtc-p2p-spike".to_string(); + + let (tx, rx) = mpsc::channel(); + let a = Engine::new(server.clone(), room.clone(), "peer-a".into(), "label-a".into()); + a.start(Box::new(Collector(tx))); + + // B's from-id and source-label are deliberately different so we can tell + // which transport delivered. + let b_from = "peer-b"; + let b = Engine::new(server, room, b_from.into(), "label-b-sse".into()); + b.start(Box::new(Collector(mpsc::channel().0))); // B needs to be live to negotiate + + // Give presence + ICE + DTLS time to bring the data channel up. + eprintln!("waiting for RTC negotiation…"); + std::thread::sleep(Duration::from_secs(6)); + + let payload = "hello-over-rtc"; + b.send(payload.to_string()); + + match rx.recv_timeout(Duration::from_secs(5)) { + Ok(m) if m.text == payload => { + let via = if m.source == b_from { "RTC data channel ✅" } else { "SSE relay (RTC lost the race)" }; + println!("received {:?}\n from={} source={} → delivered via {}", m.text, m.from, m.source, via); + if m.source != b_from { + println!(" note: message arrived, but over SSE — check the 'data channel open' log above to confirm RTC linked"); + } + } + Ok(m) => println!("⚠️ unexpected: {m:?}"), + Err(_) => { + println!("❌ nothing received in time"); + std::process::exit(1); + } + } + a.stop(); + b.stop(); +} diff --git a/core/src/lib.rs b/core/src/lib.rs index 107bf6c..c72e107 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -16,12 +16,14 @@ use std::collections::VecDeque; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; -use std::time::Duration; +use std::time::{Duration, Instant}; use futures_util::StreamExt; use serde::{Deserialize, Serialize}; use tokio::runtime::Handle; +mod rtc; + uniffi::setup_scaffolding!(); /// 1:1 with the Go server `Message`, minus the RTC `signal` blob. @@ -99,10 +101,13 @@ impl Engine { .build() .expect("tokio runtime"); let cfg = Config { server, room, from, source }; - let transports: Vec> = vec![Arc::new(SseTransport { - cfg: cfg.clone(), - http: reqwest::Client::new(), - })]; + let transports: Vec> = vec![ + Arc::new(SseTransport { + cfg: cfg.clone(), + http: reqwest::Client::new(), + }), + Arc::new(rtc::RtcTransport::new(cfg.clone())), + ]; Arc::new(Self { cfg, rt: Arc::new(rt), @@ -163,23 +168,29 @@ impl Engine { } } -/// Short-window de-dup keyed on the fields the server preserves. Once RTC lands, -/// a clipboard echoed over both SSE and the data channel collapses to one. +/// Short-window de-dup so a clipboard echoed over both SSE and the RTC data +/// channel surfaces once. Keyed on (from, text) — NOT ts: the SSE copy carries +/// the server's stamp while the data-channel copy has none, so ts must be +/// ignored. The 3s window is the cost: identical text re-sent within it is +/// suppressed, which for clipboard is an acceptable trade. #[derive(Default)] struct Dedup { - recent: VecDeque<(String, i64, String)>, + recent: VecDeque<(String, String, Instant)>, } impl Dedup { fn seen(&mut self, m: &Message) -> bool { - let key = (m.from.clone(), m.ts, m.text.clone()); - if self.recent.contains(&key) { + let now = Instant::now(); + self.recent + .retain(|(_, _, t)| now.duration_since(*t) < Duration::from_secs(3)); + if self + .recent + .iter() + .any(|(f, x, _)| f == &m.from && x == &m.text) + { return true; } - self.recent.push_back(key); - if self.recent.len() > 256 { - self.recent.pop_front(); - } + self.recent.push_back((m.from.clone(), m.text.clone(), now)); false } } diff --git a/core/src/rtc.rs b/core/src/rtc.rs new file mode 100644 index 0000000..1048029 --- /dev/null +++ b/core/src/rtc.rs @@ -0,0 +1,446 @@ +//! RTC transport — direct P2P clipboard over a WebRTC data channel. +//! +//! Signaling rides the existing server bus (`type:"signal"` / `"presence"`), +//! so this interoperates with the browser web client byte-for-byte: +//! - presence chirp: {type:"presence", role, from, room} +//! - deterministic initiator: the lower `from` id creates the offer +//! - signal: {type:"signal", to, from, signal:{kind, sdp|candidate}} +//! - data channel "tether" carries RAW clipboard text, both directions +//! +//! The data channel bypasses the relay entirely; the engine's dedup collapses +//! the SSE and RTC copies of the same clipboard into one. + +use std::collections::HashMap; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex as StdMutex}; +use std::time::Duration; + +use futures_util::StreamExt; +use serde_json::{json, Value}; +use tokio::runtime::Handle; +use tokio::sync::Mutex; + +use webrtc::api::media_engine::MediaEngine; +use webrtc::api::{APIBuilder, API}; +use webrtc::data_channel::data_channel_message::DataChannelMessage; +use webrtc::data_channel::RTCDataChannel; +use webrtc::ice_transport::ice_candidate::{RTCIceCandidate, RTCIceCandidateInit}; +use webrtc::ice_transport::ice_server::RTCIceServer; +use webrtc::peer_connection::configuration::RTCConfiguration; +use webrtc::peer_connection::peer_connection_state::RTCPeerConnectionState; +use webrtc::peer_connection::sdp::session_description::RTCSessionDescription; +use webrtc::peer_connection::RTCPeerConnection; + +use crate::{Config, Message, Sink, Status, Transport}; + +struct Peer { + pc: Arc, + dc: Mutex>>, + offerer: bool, +} + +pub(crate) struct RtcTransport { + cfg: Config, + http: reqwest::Client, + api: Arc, + peers: Arc>>>, + sink: StdMutex>, +} + +impl RtcTransport { + pub(crate) fn new(cfg: Config) -> Self { + // Data-channel-only: a bare media engine, no codecs/interceptors needed. + let api = APIBuilder::new() + .with_media_engine(MediaEngine::default()) + .build(); + RtcTransport { + cfg, + http: reqwest::Client::new(), + api: Arc::new(api), + peers: Arc::new(Mutex::new(HashMap::new())), + sink: StdMutex::new(None), + } + } + + fn sink(&self) -> Option { + self.sink.lock().unwrap().clone() + } + + fn rtc_config() -> RTCConfiguration { + RTCConfiguration { + ice_servers: vec![RTCIceServer { + urls: vec!["stun:stun.l.google.com:19302".to_owned()], + ..Default::default() + }], + ..Default::default() + } + } + + // ── bus helpers ──────────────────────────────────────────────────────── + + async fn post(&self, body: Value) { + let url = format!("{}/api/send", self.cfg.server.trim_end_matches('/')); + let _ = self + .http + .post(&url) + .header("X-Tether-Source", &self.cfg.source) + .json(&body) + .send() + .await; + } + + async fn post_presence(&self) { + self.post(json!({ + "type": "presence", "role": self.cfg.source, + "from": self.cfg.from, "room": self.cfg.room, + })) + .await; + } + + async fn post_signal(&self, to: &str, signal: Value) { + self.post(json!({ + "type": "signal", "to": to, "signal": signal, + "from": self.cfg.from, "room": self.cfg.room, + })) + .await; + } + + // ── peer setup ───────────────────────────────────────────────────────── + + /// Build a peer connection, wire ICE/state/data-channel callbacks, register + /// it in the map, and return it. + async fn new_peer(self: &Arc, remote: &str, offerer: bool) -> Arc { + let pc = Arc::new( + self.api + .new_peer_connection(Self::rtc_config()) + .await + .expect("peer connection"), + ); + + // Trickle ICE → post each candidate as it gathers. + { + let this = self.clone(); + let remote = remote.to_owned(); + pc.on_ice_candidate(Box::new(move |c: Option| { + let this = this.clone(); + let remote = remote.clone(); + Box::pin(async move { + if let Some(c) = c { + if let Ok(init) = c.to_json() { + this.post_signal( + &remote, + json!({"kind": "ice", "candidate": { + "candidate": init.candidate, + "sdpMid": init.sdp_mid, + "sdpMLineIndex": init.sdp_mline_index, + "usernameFragment": init.username_fragment, + }}), + ) + .await; + } + } + }) + })); + } + + // Drop the peer when the connection dies. + { + let peers = self.peers.clone(); + let remote = remote.to_owned(); + pc.on_peer_connection_state_change(Box::new(move |s: RTCPeerConnectionState| { + let peers = peers.clone(); + let remote = remote.clone(); + Box::pin(async move { + if matches!( + s, + RTCPeerConnectionState::Failed + | RTCPeerConnectionState::Closed + | RTCPeerConnectionState::Disconnected + ) { + peers.lock().await.remove(&remote); + } + }) + })); + } + + let peer = Arc::new(Peer { + pc: pc.clone(), + dc: Mutex::new(None), + offerer, + }); + + if offerer { + // We create the channel; the answerer receives it via on_data_channel. + let dc = pc + .create_data_channel("tether", None) + .await + .expect("data channel"); + self.wire_channel(remote, dc, &peer).await; + } else { + let this = self.clone(); + let remote = remote.to_owned(); + let peer_ref = peer.clone(); + pc.on_data_channel(Box::new(move |dc: Arc| { + let this = this.clone(); + let remote = remote.clone(); + let peer_ref = peer_ref.clone(); + Box::pin(async move { + this.wire_channel(&remote, dc, &peer_ref).await; + }) + })); + } + + self.peers + .lock() + .await + .insert(remote.to_owned(), peer.clone()); + peer + } + + /// Attach handlers to a data channel: inbound text → sink. + async fn wire_channel(self: &Arc, remote: &str, dc: Arc, peer: &Arc) { + { + let remote_log = remote.to_owned(); + dc.on_open(Box::new(move || { + let remote_log = remote_log.clone(); + Box::pin(async move { + eprintln!("tethercore[rtc]: data channel open ↔ {remote_log}"); + }) + })); + } + let this = self.clone(); + let remote_id = remote.to_owned(); + dc.on_message(Box::new(move |msg: DataChannelMessage| { + let this = this.clone(); + let remote_id = remote_id.clone(); + Box::pin(async move { + if let Ok(text) = String::from_utf8(msg.data.to_vec()) { + if let Some(sink) = this.sink() { + // from = the peer id, matching its SSE copy so dedup collapses them. + sink(Message { + kind: "clipboard".into(), + text, + from: remote_id.clone(), + to: String::new(), + role: String::new(), + source: remote_id, + room: this.cfg.room.clone(), + ts: 0, + }); + } + } + }) + })); + *peer.dc.lock().await = Some(dc); + } + + // ── signaling handlers ─────────────────────────────────────────────────── + + async fn initiate_offer(self: &Arc, remote: &str) { + if self.peers.lock().await.contains_key(remote) { + return; + } + let peer = self.new_peer(remote, true).await; + let offer = match peer.pc.create_offer(None).await { + Ok(o) => o, + Err(_) => return, + }; + if peer.pc.set_local_description(offer.clone()).await.is_err() { + return; + } + self.post_signal(remote, json!({"kind": "offer", "sdp": {"type": "offer", "sdp": offer.sdp}})) + .await; + } + + async fn handle_offer(self: &Arc, remote: &str, sdp: &str) { + // Glare tie-break: if we're already offering and we sort lower, ignore + // their offer — they should answer ours. + { + let map = self.peers.lock().await; + if let Some(p) = map.get(remote) { + if p.offerer && self.cfg.from.as_str() < remote { + return; + } + } + } + if let Some(p) = self.peers.lock().await.remove(remote) { + let _ = p.pc.close().await; + } + let peer = self.new_peer(remote, false).await; + let desc = match RTCSessionDescription::offer(sdp.to_owned()) { + Ok(d) => d, + Err(_) => return, + }; + if peer.pc.set_remote_description(desc).await.is_err() { + return; + } + let answer = match peer.pc.create_answer(None).await { + Ok(a) => a, + Err(_) => return, + }; + if peer.pc.set_local_description(answer.clone()).await.is_err() { + return; + } + self.post_signal(remote, json!({"kind": "answer", "sdp": {"type": "answer", "sdp": answer.sdp}})) + .await; + } + + async fn handle_answer(&self, remote: &str, sdp: &str) { + let peer = self.peers.lock().await.get(remote).cloned(); + if let Some(p) = peer { + if p.offerer { + if let Ok(desc) = RTCSessionDescription::answer(sdp.to_owned()) { + let _ = p.pc.set_remote_description(desc).await; + } + } + } + } + + async fn handle_ice(&self, remote: &str, candidate: &Value) { + let peer = self.peers.lock().await.get(remote).cloned(); + if let Some(p) = peer { + let init = RTCIceCandidateInit { + candidate: candidate["candidate"].as_str().unwrap_or("").to_owned(), + sdp_mid: candidate["sdpMid"].as_str().map(str::to_owned), + sdp_mline_index: candidate["sdpMLineIndex"].as_u64().map(|x| x as u16), + username_fragment: candidate["usernameFragment"].as_str().map(str::to_owned), + }; + let _ = p.pc.add_ice_candidate(init).await; + } + } + + /// One pass of the signaling SSE subscription; returns on stream end/error. + async fn signal_stream(self: &Arc, running: &AtomicBool) -> Result<(), String> { + let url = format!( + "{}/api/stream?room={}", + self.cfg.server.trim_end_matches('/'), + self.cfg.room + ); + let resp = self + .http + .get(&url) + .header("X-Tether-Client", format!("{}-rtc", self.cfg.from)) + .send() + .await + .map_err(|e| e.to_string())?; + if !resp.status().is_success() { + return Err(format!("HTTP {}", resp.status())); + } + let mut stream = resp.bytes_stream(); + let mut buf: Vec = Vec::new(); + let mut data: Vec = Vec::new(); + while let Some(chunk) = stream.next().await { + if !running.load(Ordering::SeqCst) { + return Ok(()); + } + buf.extend_from_slice(&chunk.map_err(|e| e.to_string())?); + while let Some(nl) = buf.iter().position(|&b| b == b'\n') { + let mut line: Vec = buf.drain(..=nl).collect(); + line.pop(); + if line.last() == Some(&b'\r') { + line.pop(); + } + if line.is_empty() { + if !data.is_empty() { + if let Ok(v) = serde_json::from_slice::(&data) { + self.dispatch(&v).await; + } + data.clear(); + } + } else if let Some(rest) = line.strip_prefix(b"data:") { + let rest = rest.strip_prefix(b" ").unwrap_or(rest); + data.extend_from_slice(rest); + } + } + } + Ok(()) + } + + async fn dispatch(self: &Arc, v: &Value) { + let from = v["from"].as_str().unwrap_or(""); + if from.is_empty() || from == self.cfg.from { + return; + } + match v["type"].as_str().unwrap_or("") { + "presence" => { + // Deterministic initiator: the smaller id offers. + if self.cfg.from.as_str() < from { + self.initiate_offer(from).await; + } + } + "signal" => { + let to = v["to"].as_str().unwrap_or(""); + if !to.is_empty() && to != self.cfg.from { + return; + } + let p = &v["signal"]; + match p["kind"].as_str() { + Some("offer") => { + if let Some(sdp) = p["sdp"]["sdp"].as_str() { + self.handle_offer(from, sdp).await; + } + } + Some("answer") => { + if let Some(sdp) = p["sdp"]["sdp"].as_str() { + self.handle_answer(from, sdp).await; + } + } + Some("ice") => self.handle_ice(from, &p["candidate"]).await, + _ => {} + } + } + _ => {} + } + } +} + +impl Transport for RtcTransport { + fn name(&self) -> &'static str { + "rtc" + } + + fn start(self: Arc, rt: Handle, running: Arc, sink: Sink, _status: Status) { + *self.sink.lock().unwrap() = Some(sink); + + // Presence chirp — announce ourselves so peers initiate. + { + let this = self.clone(); + let running = running.clone(); + rt.spawn(async move { + while running.load(Ordering::SeqCst) { + this.post_presence().await; + tokio::time::sleep(Duration::from_secs(8)).await; + } + }); + } + + // Signaling receive loop with reconnect/backoff. + rt.spawn(async move { + let min = Duration::from_millis(500); + let max = Duration::from_secs(10); + let mut backoff = min; + while running.load(Ordering::SeqCst) { + match self.signal_stream(&running).await { + Ok(()) => backoff = min, + Err(e) => eprintln!("tethercore[rtc]: {e}"), + } + if running.load(Ordering::SeqCst) { + tokio::time::sleep(backoff).await; + backoff = (backoff * 2).min(max); + } + } + }); + } + + fn publish(&self, rt: &Handle, msg: Message) { + let peers = self.peers.clone(); + rt.spawn(async move { + let map = peers.lock().await; + for peer in map.values() { + if let Some(dc) = peer.dc.lock().await.as_ref() { + let _ = dc.send_text(msg.text.clone()).await; + } + } + }); + } +}