Files
tether/core/src/rtc.rs
Patrick Ecord 479dda22d6 core: file transfer over LAN; share file frames between RTC and LAN
Rework the LAN transport from newline-delimited JSON to length-prefixed framing
([u32 len][u8 tag][payload]; tags H/J/M/C/E) so it can carry binary file chunks
alongside clipboard messages. send_file now goes over BOTH direct transports
(RTC + LAN) — so same-network file transfer works even when RTC isn't the path
(e.g. dead/unreachable server). Per-peer write locks instead of a whole-map lock
so a large file send doesn't stall clipboard.

Extract the M/C/E frame build/parse + IncomingFile into shared crate helpers
(file_frames / apply_file_frame), used by both RTC and LAN.

Since a file can now arrive over both transports, de-dup in the engine's file
sink (sha256 of name+bytes, 15s window) so on_file fires once.

Verified: clipboard still flows over the new LAN framing; a file transfers over
LAN with a dead server (RTC can't connect); RTC file transfer still works; no
duplicate delivery.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-21 12:39:58 -05:00

615 lines
23 KiB
Rust

//! 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, Instant};
use bytes::Bytes;
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, DirectSet, FileSink, Message, PresentSet, Sink, Status, Transport};
struct Peer {
pc: Arc<RTCPeerConnection>,
dc: Mutex<Option<Arc<RTCDataChannel>>>,
offerer: bool,
}
pub(crate) struct RtcTransport {
cfg: Config,
http: reqwest::Client,
api: Arc<API>,
peers: Arc<Mutex<HashMap<String, Arc<Peer>>>>,
sink: StdMutex<Option<Sink>>,
files: StdMutex<Option<FileSink>>,
incoming: Arc<Mutex<HashMap<String, crate::IncomingFile>>>, // in-flight inbound files
ice: StdMutex<Vec<RTCIceServer>>, // refreshed from /api/turn-cred at start
direct: DirectSet, // peers with an open data channel
present: PresentSet, // peers seen via presence chirps
}
fn default_ice() -> Vec<RTCIceServer> {
vec![RTCIceServer {
urls: vec!["stun:stun.l.google.com:19302".to_owned()],
..Default::default()
}]
}
impl RtcTransport {
pub(crate) fn new(cfg: Config, direct: DirectSet, present: PresentSet) -> 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),
files: StdMutex::new(None),
incoming: Arc::new(Mutex::new(HashMap::new())),
ice: StdMutex::new(default_ice()),
direct,
present,
}
}
fn sink(&self) -> Option<Sink> {
self.sink.lock().unwrap().clone()
}
fn file_sink(&self) -> Option<FileSink> {
self.files.lock().unwrap().clone()
}
fn rtc_config(&self) -> RTCConfiguration {
RTCConfiguration {
ice_servers: self.ice.lock().unwrap().clone(),
..Default::default()
}
}
/// Fetch server-issued ICE servers (STUN + short-lived TURN creds) from
/// `/api/turn-cred`. This is what makes cross-network P2P actually connect
/// when hole-punching fails (symmetric NAT / CGNAT). Falls back to STUN.
async fn refresh_ice(&self) {
let url = format!("{}/api/turn-cred", self.cfg.server.trim_end_matches('/'));
let resp = match self.http.get(&url).send().await {
Ok(r) if r.status().is_success() => r,
_ => return, // keep the STUN fallback
};
let Ok(v) = resp.json::<serde_json::Value>().await else {
return;
};
let Some(arr) = v["iceServers"].as_array() else {
return;
};
let mut servers = Vec::new();
for s in arr {
let urls: Vec<String> = s["urls"]
.as_array()
.map(|a| a.iter().filter_map(|u| u.as_str().map(String::from)).collect())
.unwrap_or_default();
if urls.is_empty() {
continue;
}
servers.push(RTCIceServer {
urls,
username: s["username"].as_str().unwrap_or("").to_string(),
credential: s["credential"].as_str().unwrap_or("").to_string(),
..Default::default()
});
}
if !servers.is_empty() {
*self.ice.lock().unwrap() = servers;
}
}
// ── 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, "source": self.cfg.source,
"text": self.cfg.name, // friendly name for the "who's here" list
"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<Self>, remote: &str, offerer: bool) -> Arc<Peer> {
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<RTCIceCandidate>| {
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 direct = self.direct.clone();
let remote = remote.to_owned();
pc.on_peer_connection_state_change(Box::new(move |s: RTCPeerConnectionState| {
let peers = peers.clone();
let direct = direct.clone();
let remote = remote.clone();
Box::pin(async move {
if matches!(
s,
RTCPeerConnectionState::Failed
| RTCPeerConnectionState::Closed
| RTCPeerConnectionState::Disconnected
) {
peers.lock().await.remove(&remote);
direct.lock().unwrap().remove(&remote); // relay needed again
}
})
}));
}
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<RTCDataChannel>| {
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
}
/// Reassemble inbound file frames (M meta · C chunks · E end) → file sink.
async fn handle_file_frame(self: &Arc<Self>, data: &[u8]) {
let Some((&tag, rest)) = data.split_first() else {
return;
};
let done = {
let mut map = self.incoming.lock().await;
crate::apply_file_frame(&mut map, tag, rest)
};
if let Some((name, mime, bytes)) = done {
eprintln!("tethercore[rtc]: received file {name:?} ({} bytes)", bytes.len());
if let Some(fs) = self.file_sink() {
fs(name, mime, bytes);
}
}
}
/// Attach handlers to a data channel: inbound text → sink, binary → files.
async fn wire_channel(self: &Arc<Self>, remote: &str, dc: Arc<RTCDataChannel>, peer: &Arc<Peer>) {
{
let remote_log = remote.to_owned();
let direct = self.direct.clone();
dc.on_open(Box::new(move || {
let remote_log = remote_log.clone();
let direct = direct.clone();
Box::pin(async move {
direct.lock().unwrap().insert(remote_log.clone()); // now P2P-reachable
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 msg.is_string {
// Clipboard text (string channel message).
if let Ok(text) = String::from_utf8(msg.data.to_vec()) {
if !text.is_empty() {
if let Some(sink) = this.sink() {
// from = 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,
});
}
}
}
} else {
// Binary file frame.
this.handle_file_frame(&msg.data).await;
}
})
}));
*peer.dc.lock().await = Some(dc);
}
// ── signaling handlers ───────────────────────────────────────────────────
async fn initiate_offer(self: &Arc<Self>, 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<Self>, 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<Self>, 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<u8> = Vec::new();
let mut data: Vec<u8> = 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<u8> = 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::<Value>(&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<Self>, 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" => {
// Track who's in the room (with their name) for send()'s coverage
// check and the engine's present() "who's here" list.
let name = v["text"].as_str().filter(|s| !s.is_empty()).unwrap_or(from);
let source = v["source"].as_str().unwrap_or("").to_string();
let room = v["room"].as_str().unwrap_or("").to_string();
self.present.lock().unwrap().insert(
from.to_string(),
(
crate::Peer {
id: from.to_string(),
name: name.to_string(),
source,
room,
account: String::new(),
},
Instant::now(),
),
);
// 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<Self>,
rt: Handle,
running: Arc<AtomicBool>,
sink: Sink,
_status: Status,
files: FileSink,
) {
*self.sink.lock().unwrap() = Some(sink);
*self.files.lock().unwrap() = Some(files);
// 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 {
// Pull TURN creds before any peer is built so they're in the ICE
// config (presence→offer takes a beat, so this lands in time).
self.refresh_ice().await;
let min = Duration::from_millis(500);
let max = Duration::from_secs(10);
let mut backoff = min;
// Log a signaling error only once per outage streak (LAN-only / dead
// server must not spam).
let mut logged = false;
while running.load(Ordering::SeqCst) {
match self.signal_stream(&running).await {
Ok(()) => {
backoff = min;
logged = false;
}
Err(e) => {
if !logged {
eprintln!("tethercore[rtc]: {e}");
logged = true;
}
}
}
if running.load(Ordering::SeqCst) {
tokio::time::sleep(backoff).await;
backoff = (backoff * 2).min(max);
}
}
});
}
fn publish(&self, rt: &Handle, msg: Message) {
// The "tether" data channel carries raw clipboard text only; receipts
// and other typed messages go over SSE/LAN.
if !(msg.kind.is_empty() || msg.kind == "clipboard") {
return;
}
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;
}
}
});
}
/// Stream a file over each open data channel as binary frames:
/// M<json meta> · C<16-byte id><chunk>… · E<16-byte id>
/// The channel is reliable+ordered, so chunks reassemble in arrival order.
fn send_file(&self, rt: &Handle, id: String, name: String, mime: String, data: Vec<u8>) {
let peers = self.peers.clone();
rt.spawn(async move {
let dcs: Vec<Arc<RTCDataChannel>> = {
let map = peers.lock().await;
let mut v = Vec::new();
for peer in map.values() {
if let Some(dc) = peer.dc.lock().await.as_ref() {
v.push(dc.clone());
}
}
v
};
if dcs.is_empty() {
return;
}
// Each frame is sent as one binary message: tag byte + payload.
let frames = crate::file_frames(&id, &name, &mime, &data);
for dc in dcs {
for (tag, payload) in &frames {
let mut f = Vec::with_capacity(1 + payload.len());
f.push(*tag);
f.extend_from_slice(payload);
if dc.send(&Bytes::from(f)).await.is_err() {
break;
}
}
}
});
}
}