From fdd65eea7ca7da2ed5111311b7c7aebf7b7f947f Mon Sep 17 00:00:00 2001 From: Patrick Ecord Date: Sat, 20 Jun 2026 20:38:12 -0500 Subject: [PATCH] =?UTF-8?q?core:=20in-core=20RTC=20transport=20=E2=80=94?= =?UTF-8?q?=20direct=20P2P=20clipboard=20over=20WebRTC=20data=20channel?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implement RtcTransport (webrtc-rs) behind the Transport trait, so every native platform that links tethercore gets direct P2P alongside SSE from one impl. Signaling rides the existing server bus and matches the web client byte-for- byte (presence chirp, lower-id-offers tie-break, {kind, sdp|candidate} signal envelopes, "tether" data channel carrying raw clipboard text) — so Rust peers interoperate with the browser's native RTCPeerConnection. Dedup now keys on (from, text) with a 3s TTL so the SSE and data-channel copies of the same clipboard collapse to one; the faster RTC copy wins. Proven by examples/rtc_p2p.rs: two engines bring up a data channel and a clipboard is delivered over it (received source == peer from-id, confirming the RTC path beat the SSE relay). Co-Authored-By: Claude Opus 4.8 --- core/examples/rtc_p2p.rs | 61 ++++++ core/src/lib.rs | 39 ++-- core/src/rtc.rs | 446 +++++++++++++++++++++++++++++++++++++++ 3 files changed, 532 insertions(+), 14 deletions(-) create mode 100644 core/examples/rtc_p2p.rs create mode 100644 core/src/rtc.rs 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; + } + } + }); + } +}