use rust_socketio::asynchronous::Client;
use serde_json::{json, Value};
use std::time::{Duration, Instant};
const SYNC_MIN_INTERVAL: Duration = Duration::from_secs(5);
const CONFIRM_TTL: Duration = Duration::from_secs(30);
pub struct Session {
pub room: Option<String>,
pub name: String,
pub uid: String,
pub channels: Vec<String>,
confirmed: Option<(String, Instant)>,
last_attempt: Option<Instant>,
rng: u64, loss: f64,
}
#[cfg(feature = "test-hooks")]
fn emit_loss_params() -> (f64, u64) {
let loss = std::env::var("FILAMENT_TEST_EMIT_LOSS")
.ok()
.and_then(|v| v.parse::<f64>().ok())
.filter(|l| (0.0..1.0).contains(l))
.unwrap_or(0.0);
let seed = std::env::var("FILAMENT_TEST_EMIT_SEED")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(0xF11A_C30D);
(loss, seed)
}
#[cfg(not(feature = "test-hooks"))]
#[inline]
fn emit_loss_params() -> (f64, u64) {
(0.0, 0xF11A_C30D)
}
impl Session {
pub fn new(name: &str, uid: &str) -> Self {
let (loss, seed) = emit_loss_params();
Session {
room: None,
name: name.to_string(),
uid: uid.to_string(),
channels: Vec::new(),
confirmed: None,
last_attempt: None,
rng: seed | 1,
loss,
}
}
pub fn sync_payload(&self) -> Value {
json!({
"v": 1,
"room": self.room,
"name": self.name,
"uid": self.uid,
"channels": self.channels,
})
}
pub fn heartbeat_payload(&self) -> Value {
json!({ "v": 1, "room": Value::Null, "name": self.name, "uid": self.uid })
}
fn desired_digest(&self) -> String {
let mut chans = self.channels.clone();
chans.sort();
format!("{}|{}", self.room.as_deref().unwrap_or(""), chans.join(","))
}
pub fn on_synced(&mut self, v: &Value) -> Option<Vec<Value>> {
if v["ok"].as_bool() != Some(true) {
return None; }
self.confirmed = Some((self.desired_digest(), Instant::now()));
v["peers"].as_array().cloned()
}
pub fn touch(&mut self) {
self.last_attempt = None;
}
pub fn invalidate(&mut self) {
self.confirmed = None;
self.last_attempt = None;
}
fn roll(&mut self) -> f64 {
let mut x = self.rng;
x ^= x >> 12;
x ^= x << 25;
x ^= x >> 27;
self.rng = x;
(x.wrapping_mul(0x2545F4914F6CDD1D) >> 11) as f64 / (1u64 << 53) as f64
}
pub async fn emit(&mut self, sio: &Client, event: &str, payload: Value) {
if self.loss > 0.0 && self.roll() < self.loss {
return; }
let _ = sio.emit(event, payload).await;
}
pub async fn tick(&mut self, sio: &Client) {
let Some(room) = self.room.clone() else { return };
let now = Instant::now();
let due = match (&self.confirmed, &self.last_attempt) {
(None, Some(at)) => now.duration_since(*at) >= SYNC_MIN_INTERVAL,
(None, None) => true,
(Some((digest, at)), _) => {
let stale = now.duration_since(*at) >= CONFIRM_TTL;
let diverged = *digest != self.desired_digest();
(stale || diverged)
&& self
.last_attempt
.map(|a| now.duration_since(a) >= SYNC_MIN_INTERVAL)
.unwrap_or(true)
}
};
if !due {
return;
}
self.last_attempt = Some(now);
let _ = room; self.emit(sio, "sync", self.sync_payload()).await;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn digest_orders_channels() {
let mut s = Session::new("n", "u");
s.room = Some("r".into());
s.channels = vec!["b".into(), "a".into()];
let d1 = s.desired_digest();
s.channels = vec!["a".into(), "b".into()];
assert_eq!(d1, s.desired_digest());
}
#[test]
fn loss_shim_deterministic() {
let mut a = Session::new("n", "u");
let mut b = Session::new("n", "u");
let ra: Vec<u64> = (0..8).map(|_| (a.roll() * 1e9) as u64).collect();
let rb: Vec<u64> = (0..8).map(|_| (b.roll() * 1e9) as u64).collect();
assert_eq!(ra, rb, "same seed must replay the same drop pattern");
}
}