use super::*;
use crate::peer::{self, Assembly, Chunk, Identity, LineSaid, Peer, Waiting};
use anyhow::Context;
use futures_util::{SinkExt, StreamExt};
use std::collections::HashMap;
use std::sync::atomic::Ordering::Relaxed;
use std::time::Duration;
use tokio::time::Instant;
use tokio_tungstenite::tungstenite::{client::IntoClientRequest, Message as WsMessage};
type Socket =
tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>;
const CONNECT_TIMEOUT: Duration = Duration::from_secs(20);
const SILENCE: Duration = Duration::from_secs(60);
const RETRY_EVERY: Duration = Duration::from_secs(600);
const FALLBACK_AFTER: u32 = 6;
const LOOK_EVERY: Duration = Duration::from_secs(60);
const LINE_QUEUE: usize = 8;
const QUIET_BACKOFF: Duration = Duration::from_secs(60);
enum Left {
Done,
Closed(String),
Budget(i64),
}
pub(crate) fn spawn_peer_link(app: Arc<App>) {
spawn_prints(app.clone());
tokio::spawn(supervise(app));
}
async fn supervise(app: Arc<App>) {
let mut lines: HashMap<i64, tokio::task::JoinHandle<()>> = HashMap::new();
let mut mailbox: Option<tokio::task::JoinHandle<()>> = None;
let mut wake = app.peers.wake.subscribe();
let mut tick = tokio::time::interval(RETRY_EVERY);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
tick.tick().await;
loop {
lines.retain(|_, h| !h.is_finished());
if mailbox.as_ref().is_some_and(|h| h.is_finished()) {
mailbox = None;
}
let friends = live_friends(&app);
if !friends.is_empty() {
if let Some(me) = friend_identity(&app).await {
for p in &friends {
lines.entry(p.id).or_insert_with(|| {
tokio::spawn(line_task(app.clone(), me.clone(), p.clone()))
});
}
if mailbox.is_none() && friends.iter().any(|p| p.v < peer::LINE_V) {
mailbox = Some(tokio::spawn(mailbox_task(app.clone(), me.clone())));
}
}
}
tokio::select! {
_ = tokio::time::sleep(LOOK_EVERY) => {},
_ = wake.changed() => {
if !friends.is_empty() {
let app_b = app.clone();
let _ = tokio::task::spawn_blocking(move || sweep_offers(&app_b)).await;
flush_all(&app);
}
},
_ = tick.tick() => {
if !friends.is_empty() {
let app_b = app.clone();
let _ = tokio::task::spawn_blocking(move || {
let _ = app_b.store.prune_peer_taken();
}).await;
flush_all(&app);
}
},
}
}
}
const PRINTS_FIRST: Duration = Duration::from_secs(20);
const PRINTS_EVERY: Duration = Duration::from_secs(600);
const PRINTS_ROUND: usize = 40;
fn spawn_prints(app: Arc<App>) {
tokio::spawn(async move {
tokio::time::sleep(PRINTS_FIRST).await;
loop {
if has_friends(&app) {
let app_b = app.clone();
let _ = tokio::task::spawn_blocking(move || print_folders(&app_b)).await;
}
tokio::time::sleep(PRINTS_EVERY).await;
}
});
}
pub(crate) fn print_folders(app: &App) -> usize {
let before = crate::store::now() - PRINT_KEPT;
let Ok(roots) = app.store.roots_to_print(before) else {
return 0;
};
let mut n = 0;
for root in roots.into_iter().take(PRINTS_ROUND) {
let dir = std::path::Path::new(&root);
let print = if dir.is_dir() {
crate::git::print(dir)
} else {
None
};
n += print.is_some() as usize;
let _ = app.store.set_print(&root, print.as_ref());
}
n
}
async fn friend_identity(app: &Arc<App>) -> Option<Identity> {
if !has_friends(app) {
return None;
}
let app_b = app.clone();
tokio::task::spawn_blocking(move || identity_if_any(&app_b))
.await
.ok()
.flatten()
}
fn live_friends(app: &App) -> Vec<Peer> {
app.store
.peers()
.map(|ps| ps.into_iter().filter(|p| p.removed_at == 0).collect())
.unwrap_or_default()
}
fn has_friends(app: &App) -> bool {
!live_friends(app).is_empty()
}
fn friend(app: &App, id: i64) -> Option<Peer> {
app.store
.peer(id)
.ok()
.flatten()
.filter(|p| p.removed_at == 0)
}
fn needs_mailbox(app: &App) -> bool {
live_friends(app).iter().any(|p| p.v < peer::LINE_V)
}
fn wait_before(app: &App, attempt: u32) -> Duration {
let b = peer::backoff(attempt);
let quiet = app.peers.quiet_until.load(Relaxed);
if quiet > crate::store::now() {
b.max(QUIET_BACKOFF)
} else {
b
}
}
pub(crate) fn attempt_after(up_for: Duration, attempt: u32) -> u32 {
if up_for >= SILENCE {
0
} else {
attempt
}
}
async fn sleep_until_secs(until: i64) {
let now = crate::store::now();
if until > now {
tokio::time::sleep(Duration::from_secs((until - now) as u64)).await;
}
}
async fn connect(me: &Identity, url: String, path: &str) -> anyhow::Result<Socket> {
let mut req = url.into_client_request().context("the relay's address")?;
req.headers_mut().insert(
"x-snyvi-auth",
me.relay_auth("GET", path)
.parse()
.context("the signature as a header")?,
);
let (socket, _) = tokio::time::timeout(CONNECT_TIMEOUT, tokio_tungstenite::connect_async(req))
.await
.map_err(|_| anyhow::anyhow!("no answer to the link in {} s", CONNECT_TIMEOUT.as_secs()))?
.context("opening the link")?;
Ok(socket)
}
fn budget_close(
frame: &Option<tokio_tungstenite::tungstenite::protocol::CloseFrame>,
) -> Option<i64> {
let f = frame.as_ref()?;
(u16::from(f.code) == peer::CLOSE_BUDGET).then(|| peer::until_of(&f.reason))
}
fn heard_quiet(app: &App, t: &str) {
if let Ok(LineSaid::Quiet { until }) = serde_json::from_str::<LineSaid>(t) {
app.peers.quiet_until.store(until / 1000, Relaxed);
}
}
async fn line_task(app: Arc<App>, me: Identity, p: Peer) {
let mut attempt: u32 = 0;
let mut wake = app.peers.wake.subscribe();
let path = format!("/line/{}/{}", me.address(), p.sign_key);
loop {
if friend(&app, p.id).is_none() {
return;
}
match connect(&me, peer::line_ws(&me.address(), &p.sign_key), &path).await {
Ok(socket) => {
if attempt > 0 {
eprintln!("snyvi: the line to {} is back", p.name);
}
let opened = Instant::now();
match on_line(&app, &me, &p, socket, &mut wake).await {
Ok(Left::Done) => return,
Ok(Left::Closed(why)) => {
eprintln!("snyvi: the line to {} closed: {why}", p.name)
}
Ok(Left::Budget(until)) => {
eprintln!("snyvi: the relay has had this snyvi's share of sockets for today; the line to {} waits", p.name);
sleep_until_secs(until).await;
attempt = 0;
continue;
}
Err(e) => eprintln!("snyvi: the line to {}: {e:#}", p.name),
}
attempt = attempt_after(opened.elapsed(), attempt);
}
Err(e) => {
if attempt == 0 {
eprintln!("snyvi: the line to {}: {e:#}", p.name);
}
}
}
attempt += 1;
tokio::select! {
_ = tokio::time::sleep(wait_before(&app, attempt)) => {},
_ = wake.changed() => {},
}
}
}
async fn on_line(
app: &Arc<App>,
me: &Identity,
p: &Peer,
mut socket: Socket,
wake: &mut tokio::sync::watch::Receiver<u64>,
) -> anyhow::Result<Left> {
let (tx, mut rx) = tokio::sync::mpsc::channel::<Outbound>(LINE_QUEUE);
app.peers
.lines
.lock()
.unwrap()
.insert(p.id, LineHandle { tx: tx.clone() });
let app_b = app.clone();
let pid = p.id;
let _ = tokio::task::spawn_blocking(move || app_b.store.peer_due_now(pid)).await;
flush_all(app);
let mut ping = tokio::time::interval(peer::PING_EVERY);
ping.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
ping.tick().await;
let mut heard = Instant::now();
let mut assembly = Assembly::default();
let left = loop {
let silence = tokio::time::sleep_until(heard + SILENCE);
tokio::select! {
msg = socket.next() => {
let Some(msg) = msg else {
break Ok(Left::Closed("the relay hung up".into()));
};
let msg = msg.context("reading the line")?;
heard = Instant::now();
match msg {
WsMessage::Binary(b) => {
let Some(c) = Chunk::parse(&b) else { continue };
if let Some((id, bytes)) = assembly.take(&c) {
let w = Waiting { id, sender: p.sign_key.clone(), size: bytes.len() as u64 };
took(app, me, &mut socket, w, bytes, true).await?;
}
}
WsMessage::Text(t) => {
if t.as_str() == "pong" {
continue;
}
let Ok(said) = serde_json::from_str::<LineSaid>(t.as_str()) else { continue };
let app_b = app.clone();
let p_b = p.clone();
tokio::task::spawn_blocking(move || line_said(&app_b, &p_b, said)).await?;
}
WsMessage::Close(f) => {
break Ok(match budget_close(&f) {
Some(until) => Left::Budget(until),
None => Left::Closed("closed by the relay".into()),
});
}
_ => {}
}
}
out = rx.recv() => {
let Some(out) = out else { break Ok(Left::Closed("the line's queue closed".into())) };
for c in out.chunks {
socket.send(WsMessage::Binary(c.into())).await.with_context(|| format!("sending {} down the line", out.id))?;
}
}
_ = ping.tick() => {
socket.send(WsMessage::Text("ping".into())).await.context("the ping")?;
}
_ = silence => {
let _ = socket.close(None).await;
break Ok(Left::Closed(format!("nothing heard in {} s", SILENCE.as_secs())));
}
_ = wake.changed() => {
if friend(app, p.id).is_none() {
let _ = socket.close(None).await;
break Ok(Left::Done);
}
}
}
};
let mut lines = app.peers.lines.lock().unwrap();
if lines.get(&p.id).is_some_and(|h| h.tx.same_channel(&tx)) {
lines.remove(&p.id);
}
left
}
fn line_said(app: &Arc<App>, p: &Peer, said: LineSaid) {
let now = crate::store::now();
match said {
LineSaid::Held(id) => {
let _ = app.store.peer_sent(&id);
let _ = app.store.touch_peer(p.id, false);
peers_moved(app);
}
LineSaid::Arrived(id) => {
let _ = app.store.peer_sent(&id);
let _ = app.store.touch_peer(p.id, false);
let _ = app.store.peer_receipt(p.id, &id, "arrived");
peers_moved(app);
}
LineSaid::Resend(id) => {
let _ = app.store.peer_waiting(&id, "sending again", Some(now));
flush_all(app);
}
LineSaid::Full(id) => {
let _ = app.store.peer_waiting(&id, "their line is full", None);
}
LineSaid::Failed { id, why } => {
let _ = app.store.peer_failed(&id, &why);
}
LineSaid::Friend { on: true } => {
if p.v < peer::LINE_V {
let _ = app.store.peer_set_v(p.id, peer::LINE_V);
}
let _ = app.store.peer_due_now(p.id);
flush_all(app);
}
LineSaid::Friend { on: false } => {}
LineSaid::Quiet { until } => app.peers.quiet_until.store(until / 1000, Relaxed),
}
}
async fn mailbox_task(app: Arc<App>, me: Identity) {
let mut attempt: u32 = 0;
let mut wake = app.peers.wake.subscribe();
let path = format!("/inbox/{}", me.address());
loop {
if !needs_mailbox(&app) {
return;
}
match connect(&me, peer::relay_ws(&me.address()), &path).await {
Ok(socket) => {
if attempt > 0 {
eprintln!("snyvi: the relay is back");
}
let opened = Instant::now();
match linked(&app, &me, socket, &mut wake).await {
Ok(Left::Done) => return,
Ok(Left::Closed(why)) => eprintln!("snyvi: the relay link closed: {why}"),
Ok(Left::Budget(until)) => {
eprintln!("snyvi: the relay has had this snyvi's share of sockets for today; the mailbox link waits");
sleep_until_secs(until).await;
attempt = 0;
continue;
}
Err(e) => eprintln!("snyvi: the relay link: {e:#}"),
}
attempt = attempt_after(opened.elapsed(), attempt);
}
Err(e) => {
if attempt == 0 {
eprintln!("snyvi: the relay: {e:#}");
}
}
}
attempt += 1;
if attempt >= FALLBACK_AFTER {
flush_all(&app);
let app_b = app.clone();
let me_b = me.clone();
let _ = tokio::task::spawn_blocking(move || {
if let Err(e) = bring_in(&app_b, &me_b) {
eprintln!("snyvi: the relay, over HTTP: {e:#}");
}
})
.await;
}
tokio::select! {
_ = tokio::time::sleep(wait_before(&app, attempt)) => {},
_ = wake.changed() => {},
}
}
}
async fn linked(
app: &Arc<App>,
me: &Identity,
mut socket: Socket,
wake: &mut tokio::sync::watch::Receiver<u64>,
) -> anyhow::Result<Left> {
let app_b = app.clone();
let me_b = me.clone();
tokio::task::spawn_blocking(move || {
sweep_offers(&app_b);
if !app_b.peers.held_read.swap(true, Relaxed) {
read_held(&app_b, &me_b);
}
})
.await?;
flush_all(app);
let mut ping = tokio::time::interval(peer::PING_EVERY);
ping.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
ping.tick().await;
let mut heard = Instant::now();
let mut pending: Option<Waiting> = None;
loop {
let silence = tokio::time::sleep_until(heard + SILENCE);
tokio::select! {
msg = socket.next() => {
let Some(msg) = msg else {
return Ok(Left::Closed("the relay hung up".into()));
};
let msg = msg.context("reading the link")?;
heard = Instant::now();
match msg {
WsMessage::Text(t) => {
if t.as_str() == "pong" {
continue;
}
let Ok(peer::Pushed::Frame(w)) = serde_json::from_str::<peer::Pushed>(t.as_str()) else {
heard_quiet(app, t.as_str());
continue;
};
if w.size as usize <= peer::INLINE_MAX {
pending = Some(w);
continue;
}
pending = None;
let me_b = me.clone();
let id = w.id.clone();
match tokio::task::spawn_blocking(move || peer::fetch(&me_b, &id)).await? {
Ok(bytes) => took(app, me, &mut socket, w, bytes, false).await?,
Err(e) => eprintln!("snyvi: a frame from the relay could not be fetched: {e:#}"),
}
}
WsMessage::Binary(b) => {
if let Some(w) = pending.take() {
took(app, me, &mut socket, w, b.to_vec(), false).await?;
}
}
WsMessage::Close(f) => {
return Ok(match budget_close(&f) {
Some(until) => Left::Budget(until),
None => Left::Closed("closed by the relay".into()),
});
}
_ => {}
}
}
_ = ping.tick() => {
socket.send(WsMessage::Text("ping".into())).await.context("the ping")?;
}
_ = silence => {
let _ = socket.close(None).await;
return Ok(Left::Closed(format!("nothing heard in {} s", SILENCE.as_secs())));
}
_ = wake.changed() => {
if !needs_mailbox(app) {
let _ = socket.close(None).await;
return Ok(Left::Done);
}
}
}
}
}
async fn took(
app: &Arc<App>,
me: &Identity,
socket: &mut Socket,
w: Waiting,
bytes: Vec<u8>,
via_line: bool,
) -> anyhow::Result<()> {
let app_b = app.clone();
let me_b = me.clone();
let id = w.id.clone();
match tokio::task::spawn_blocking(move || take_in(&app_b, &me_b, &w, bytes, via_line)).await? {
Ok(_) => socket
.send(WsMessage::Text(peer::ack_message(&id).into()))
.await
.context("the ack")?,
Err(e) => eprintln!("snyvi: a frame from the relay was left there: {e:#}"),
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_short_link_keeps_counting_and_a_long_one_starts_again() {
assert_eq!(attempt_after(Duration::from_secs(5), 4), 4);
assert_eq!(attempt_after(Duration::from_secs(120), 4), 0);
}
}