use crate::firebase::schema::AppNamespace;
use crate::logging::{log_info, log_warn};
use crate::signaling::{Device, DeviceCapabilities};
use anyhow::Result;
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex};
const REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(4);
const REQUEST_ATTEMPTS: usize = 4;
const RETRY_BASE_MS: u64 = 125;
const RETRY_MAX_MS: u64 = 2_000;
const LEASE_MS_DEFAULT: i64 = 90_000;
const LEASE_MS_MIN: i64 = 15_000;
const LEASE_MS_MAX: i64 = 10 * 60 * 1000;
const SPACE_PREFIX: &str = "space::";
type TokenProvider = Arc<Mutex<Box<dyn Fn() -> Option<String> + Send + Sync>>>;
#[derive(Clone, Debug, serde::Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub(crate) struct RtdbPresenceState {
#[serde(skip)]
pub(crate) runtime_instance_id: Option<String>,
#[serde(default)]
pub(crate) online: bool,
#[serde(default)]
pub(crate) last_seen_at: Option<i64>,
#[serde(default)]
pub(crate) ticket_version: Option<i64>,
#[serde(default)]
pub(crate) lease_expires_at: Option<i64>,
#[serde(default)]
pub(crate) ticket: Option<String>,
#[serde(default)]
pub(crate) node_id: Option<String>,
#[serde(default)]
pub(crate) device_name: Option<String>,
#[serde(default)]
pub(crate) platform_type: Option<String>,
#[serde(default)]
pub(crate) metadata: Option<String>,
#[serde(default)]
pub(crate) excluded_peers: Vec<String>,
}
#[derive(Clone, Debug, serde::Deserialize, Default)]
struct RtdbDevicePresenceNode {
#[serde(default)]
instances: HashMap<String, RtdbPresenceState>,
}
#[derive(Clone, Debug, serde::Serialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
struct RtdbPresenceWriteState {
online: bool,
last_seen_at: i64,
ticket_version: i64,
lease_expires_at: i64,
#[serde(skip_serializing_if = "Option::is_none")]
ticket: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
node_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
device_name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
platform_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
metadata: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
excluded_peers: Vec<String>,
}
#[derive(Clone, Debug, serde::Serialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
struct RtdbExclusionWriteState {
excluded_peers: Vec<String>,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(crate) struct RtdbOverlayStats {
pub(crate) matched: usize,
pub(crate) promoted_online: usize,
pub(crate) demoted_offline: usize,
pub(crate) rtdb_only_added: usize,
}
#[derive(Clone)]
pub(crate) struct RtdbPresenceClient {
http: reqwest::Client,
project_id: String,
namespace: AppNamespace,
token_provider: TokenProvider,
device_ids_by_node: Arc<tokio::sync::RwLock<HashMap<String, String>>>,
excluded_peers_by_node: Arc<tokio::sync::RwLock<HashMap<String, Vec<String>>>>,
runtime_instance_id: String,
platform_type: String,
}
impl RtdbPresenceClient {
pub(crate) fn new(
project_id: impl Into<String>,
app_tag: impl Into<String>,
token_provider: TokenProvider,
) -> Self {
Self::new_with_platform_type(
project_id,
app_tag,
token_provider,
crate::native_device::default_platform_type(),
)
}
fn new_with_platform_type(
project_id: impl Into<String>,
app_tag: impl Into<String>,
token_provider: TokenProvider,
platform_type: impl Into<String>,
) -> Self {
Self {
http: reqwest::Client::new(),
project_id: project_id.into(),
namespace: AppNamespace::new(app_tag.into()),
token_provider,
device_ids_by_node: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
excluded_peers_by_node: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
runtime_instance_id: new_runtime_instance_id(),
platform_type: platform_type.into(),
}
}
pub(crate) fn target_label(&self) -> String {
let (base, namespace) = self.database_url();
namespace
.map(|value| format!("{base}?ns={value}"))
.unwrap_or(base)
}
pub(crate) async fn publish(
&self,
user_id: &str,
local_node_id: &str,
device_id: &str,
ticket: &str,
device_name: &str,
metadata: Option<&str>,
) -> Result<()> {
let now = now_ms()?;
let excluded_peers = self
.excluded_peers_by_node
.read()
.await
.get(local_node_id)
.cloned()
.unwrap_or_default();
let state = self.online_state(
now,
ticket,
local_node_id,
device_name,
metadata,
excluded_peers,
);
self.write(user_id, device_id, state).await?;
self.device_ids_by_node
.write()
.await
.insert(local_node_id.to_string(), device_id.to_string());
Ok(())
}
fn online_state(
&self,
now: i64,
ticket: &str,
local_node_id: &str,
device_name: &str,
metadata: Option<&str>,
excluded_peers: Vec<String>,
) -> RtdbPresenceWriteState {
RtdbPresenceWriteState {
online: true,
last_seen_at: now,
ticket_version: now,
lease_expires_at: now.saturating_add(lease_ms()),
ticket: Some(ticket.to_string()),
node_id: Some(local_node_id.to_string()),
device_name: Some(device_name.to_string()),
platform_type: Some(self.platform_type.clone()),
metadata: metadata.map(str::to_string),
excluded_peers,
}
}
pub(crate) async fn set_offline(&self, user_id: &str, local_node_id: &str) -> Result<()> {
let device_id = self
.device_ids_by_node
.read()
.await
.get(local_node_id)
.cloned()
.ok_or_else(|| {
anyhow::anyhow!(
"RTDB live presence must be published before updating peer exclusions"
)
})?;
self.set_device_offline(user_id, &device_id, Some(local_node_id))
.await
}
pub(crate) async fn set_device_offline(
&self,
user_id: &str,
device_id: &str,
local_node_id: Option<&str>,
) -> Result<()> {
let now = now_ms()?;
self.write(
user_id,
device_id,
RtdbPresenceWriteState {
online: false,
last_seen_at: now,
ticket_version: now,
lease_expires_at: now,
ticket: None,
node_id: local_node_id.map(str::to_string),
device_name: None,
platform_type: Some(self.platform_type.clone()),
metadata: None,
excluded_peers: Vec::new(),
},
)
.await
}
pub(crate) async fn set_excluded_peers(
&self,
user_id: &str,
local_node_id: &str,
excluded_peers: &[String],
) -> Result<()> {
let device_id = self
.device_ids_by_node
.read()
.await
.get(local_node_id)
.cloned()
.unwrap_or_else(|| local_node_id.to_string());
let mut normalized = excluded_peers
.iter()
.map(|value| value.trim().to_ascii_lowercase())
.filter(|value| !value.is_empty())
.collect::<Vec<_>>();
normalized.sort();
normalized.dedup();
self.patch(
user_id,
&device_id,
&RtdbExclusionWriteState {
excluded_peers: normalized.clone(),
},
)
.await?;
self.excluded_peers_by_node
.write()
.await
.insert(local_node_id.to_string(), normalized);
Ok(())
}
pub(crate) async fn read(&self, user_id: &str) -> Result<HashMap<String, RtdbPresenceState>> {
let token = self.token()?;
let path = self.scope_path(user_id);
let (base, emulator_namespace) = self.database_url();
let mut url =
reqwest::Url::parse(&format!("{}/{}.json", base.trim_end_matches('/'), path))?;
{
let mut query = url.query_pairs_mut();
if let Some(namespace) = emulator_namespace.as_deref() {
query.append_pair("ns", namespace);
}
query.append_pair("auth", &token);
}
let response = self
.send_with_retry("presence-read", || Ok(self.http.get(url.clone())))
.await?;
if response.status() == reqwest::StatusCode::NOT_FOUND {
return Ok(HashMap::new());
}
if !response.status().is_success() {
anyhow::bail!("RTDB presence read failed status={}", response.status());
}
let raw: serde_json::Value = response
.json()
.await
.map_err(|_| anyhow::anyhow!("RTDB presence response was not valid JSON"))?;
let Some(object) = raw.as_object() else {
return Ok(HashMap::new());
};
Ok(object
.iter()
.filter_map(|(device_id, value)| {
let node: RtdbDevicePresenceNode = serde_json::from_value(value.clone()).ok()?;
aggregate_device_presence(node.instances, now_ms().ok()?)
.map(|state| (device_id.clone(), state))
})
.collect())
}
pub(crate) fn overlay_devices(
&self,
devices: &mut Vec<Device>,
presence: &HashMap<String, RtdbPresenceState>,
user_id: &str,
exclude_node_id: Option<&str>,
now_ms: i64,
) -> RtdbOverlayStats {
let mut stats = RtdbOverlayStats::default();
for device in devices.iter_mut() {
let Some(state) = presence.get(device.device_id.trim()) else {
continue;
};
stats.matched += 1;
if is_live(state, now_ms) {
if !device.online {
stats.promoted_online += 1;
}
device.online = true;
if let Some(value) = trimmed(&state.ticket) {
device.ticket = Some(value);
}
if let Some(value) = trimmed(&state.node_id) {
device.node_id = Some(value);
}
if let Some(value) = trimmed(&state.device_name) {
device.device_name = value;
}
if let Some(value) = trimmed(&state.platform_type) {
device.platform_type = Some(value);
}
if let Some(value) = trimmed(&state.metadata) {
device.metadata = Some(value);
}
if let Some(value) = trimmed(&state.runtime_instance_id) {
device.session_id = Some(value);
}
device.excluded_peers = state.excluded_peers.clone();
} else if device.online {
stats.demoted_offline += 1;
device.online = false;
}
}
let known: HashSet<String> = devices
.iter()
.map(|device| device.device_id.trim().to_string())
.filter(|value| !value.is_empty())
.collect();
for (device_id, state) in presence {
if known.contains(device_id.trim()) {
continue;
}
if let Some(device) =
self.device_from_presence(user_id, device_id, state, exclude_node_id, now_ms)
{
devices.push(device);
stats.rtdb_only_added += 1;
}
}
stats
}
async fn write(
&self,
user_id: &str,
device_id: &str,
state: RtdbPresenceWriteState,
) -> Result<()> {
let token = self.token()?;
let path = self.instance_path(user_id, device_id);
let (base, emulator_namespace) = self.database_url();
let url = reqwest::Url::parse(&format!("{}/{}.json", base.trim_end_matches('/'), path))?;
let response = self
.send_with_retry("presence-write", || {
let mut request_url = url.clone();
{
let mut query = request_url.query_pairs_mut();
if let Some(namespace) = emulator_namespace.as_deref() {
query.append_pair("ns", namespace);
}
query.append_pair("auth", &token);
query.append_pair("print", "silent");
}
Ok(self.http.put(request_url).json(&state))
})
.await?;
if !response.status().is_success() {
anyhow::bail!("RTDB presence write failed status={}", response.status());
}
Ok(())
}
async fn patch<T: serde::Serialize + ?Sized>(
&self,
user_id: &str,
device_id: &str,
state: &T,
) -> Result<()> {
let token = self.token()?;
let path = self.instance_path(user_id, device_id);
let (base, emulator_namespace) = self.database_url();
let url = reqwest::Url::parse(&format!("{}/{}.json", base.trim_end_matches('/'), path))?;
let response = self
.send_with_retry("presence-patch", || {
let mut request_url = url.clone();
{
let mut query = request_url.query_pairs_mut();
if let Some(namespace) = emulator_namespace.as_deref() {
query.append_pair("ns", namespace);
}
query.append_pair("auth", &token);
query.append_pair("print", "silent");
}
Ok(self.http.patch(request_url).json(state))
})
.await?;
if !response.status().is_success() {
anyhow::bail!("RTDB presence patch failed status={}", response.status());
}
Ok(())
}
async fn send_with_retry<F>(&self, operation: &str, mut request: F) -> Result<reqwest::Response>
where
F: FnMut() -> Result<reqwest::RequestBuilder>,
{
for attempt in 0..REQUEST_ATTEMPTS {
match tokio::time::timeout(REQUEST_TIMEOUT, request()?.send()).await {
Ok(Ok(response))
if retryable_status(response.status()) && attempt + 1 < REQUEST_ATTEMPTS =>
{
self.retry(operation, attempt, Some(response.status().to_string()))
.await;
}
Ok(Ok(response)) => return Ok(response),
Ok(Err(error)) if attempt + 1 < REQUEST_ATTEMPTS => {
self.retry(operation, attempt, Some(error.to_string()))
.await;
}
Ok(Err(error)) => anyhow::bail!(
"RTDB {operation} failed after {REQUEST_ATTEMPTS} attempts: {error}"
),
Err(_) if attempt + 1 < REQUEST_ATTEMPTS => {
self.retry(operation, attempt, Some("timeout".to_string()))
.await;
}
Err(_) => {
anyhow::bail!("RTDB {operation} timed out after {REQUEST_ATTEMPTS} attempts")
}
}
}
unreachable!("RTDB retry loop always returns or fails")
}
async fn retry(&self, operation: &str, attempt: usize, reason: Option<String>) {
let delay = retry_delay(attempt);
log_warn(&format!(
"[OPENRTC][RTDB] retry operation={} attempt={}/{} delay_ms={} reason={}",
operation,
attempt + 1,
REQUEST_ATTEMPTS,
delay.as_millis(),
reason.unwrap_or_else(|| "unknown".to_string()),
));
tokio::time::sleep(delay).await;
}
fn token(&self) -> Result<String> {
let token = (self
.token_provider
.lock()
.map_err(|_| anyhow::anyhow!("RTDB token provider lock poisoned"))?)(
)
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty());
token.ok_or_else(|| anyhow::anyhow!("RTDB live presence requires an auth token"))
}
fn scope_path(&self, user_id: &str) -> String {
if self.namespace.is_space() {
let namespace_id = self
.namespace
.app_tag()
.strip_prefix(SPACE_PREFIX)
.unwrap_or_else(|| self.namespace.app_tag());
format!("presence/spaces/{}/devices", encode(namespace_id))
} else {
format!(
"presence/{}/users/{}/devices",
encode(self.namespace.app_tag()),
encode(user_id)
)
}
}
fn instance_path(&self, user_id: &str, device_id: &str) -> String {
format!(
"{}/{}/instances/{}",
self.scope_path(user_id),
encode(device_id),
encode(&self.runtime_instance_id),
)
}
fn database_url(&self) -> (String, Option<String>) {
if let Some(host) = emulator_host() {
let base = if host.starts_with("http://") || host.starts_with("https://") {
host
} else {
format!("http://{host}")
};
let namespace = std::env::var("OPENRTC_RTDB_EMULATOR_NAMESPACE")
.ok()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
.unwrap_or_else(|| format!("{}-default-rtdb", self.project_id));
return (base, Some(namespace));
}
let configured = [
"OPENRTC_DATABASE_URL",
"VITE_OPENRTC_DATABASE_URL",
"FIREBASE_DATABASE_URL",
]
.into_iter()
.find_map(|key| {
std::env::var(key)
.ok()
.map(|value| value.trim().trim_end_matches('/').to_string())
.filter(|value| !value.is_empty())
});
(
configured.unwrap_or_else(|| {
format!("https://{}-default-rtdb.firebaseio.com", self.project_id)
}),
None,
)
}
fn device_from_presence(
&self,
user_id: &str,
device_id: &str,
state: &RtdbPresenceState,
exclude_node_id: Option<&str>,
now_ms: i64,
) -> Option<Device> {
if !is_live(state, now_ms) || device_id.trim().is_empty() {
return None;
}
let ticket = trimmed(&state.ticket);
let node_id = trimmed(&state.node_id);
if ticket.is_none() && node_id.is_none() {
return None;
}
if let Some(exclude) = exclude_node_id
.map(str::trim)
.filter(|value| !value.is_empty())
{
if node_id.as_deref() == Some(exclude) || device_id == exclude {
return None;
}
}
let timestamp = state.last_seen_at.unwrap_or(now_ms);
let timestamp_value = Some(serde_json::Value::Number(timestamp.into()));
Some(Device {
app_tag: Some(self.namespace.app_tag().to_string()),
device_id: device_id.to_string(),
user_id: Some(user_id.to_string()),
device_name: trimmed(&state.device_name).unwrap_or_else(|| "Device".to_string()),
platform_type: trimmed(&state.platform_type),
capabilities: Some(DeviceCapabilities {
can_host: false,
can_sync: false,
read_only: false,
}),
session_id: trimmed(&state.runtime_instance_id).or_else(|| Some(device_id.to_string())),
node_id,
tag: Some(self.namespace.app_tag().to_string()),
kind: Some("device".to_string()),
metadata: trimmed(&state.metadata),
online: true,
ticket,
last_seen_at: timestamp_value.clone(),
expires_at: None,
created_at: timestamp_value.clone(),
updated_at: timestamp_value,
excluded_peers: state.excluded_peers.clone(),
})
}
}
fn emulator_host() -> Option<String> {
[
"OPENRTC_RTDB_EMULATOR_HOST",
"FIREBASE_DATABASE_EMULATOR_HOST",
"VITE_OPENRTC_RTDB_EMULATOR_HOST",
]
.into_iter()
.find_map(|key| {
std::env::var(key)
.ok()
.map(|value| value.trim().trim_end_matches('/').to_string())
.filter(|value| !value.is_empty())
})
}
fn encode(segment: &str) -> String {
segment
.as_bytes()
.iter()
.map(|byte| {
let ch = *byte as char;
if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
ch.to_string()
} else {
format!("%{byte:02X}")
}
})
.collect()
}
fn is_live(state: &RtdbPresenceState, now_ms: i64) -> bool {
state.online
&& state
.lease_expires_at
.is_none_or(|expires_at| expires_at > now_ms)
}
fn aggregate_device_presence(
instances: HashMap<String, RtdbPresenceState>,
now_ms: i64,
) -> Option<RtdbPresenceState> {
let mut instances = instances.into_iter().collect::<Vec<_>>();
instances.sort_by(|left, right| right.0.cmp(&left.0));
instances
.iter()
.find(|(_, state)| is_live(state, now_ms))
.map(|(runtime_instance_id, state)| {
let mut state = state.clone();
state.runtime_instance_id = Some(runtime_instance_id.clone());
state
})
.or_else(|| {
instances
.into_iter()
.next()
.map(|(runtime_instance_id, mut state)| {
state.runtime_instance_id = Some(runtime_instance_id);
state
})
})
}
fn new_runtime_instance_id() -> String {
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis())
.unwrap_or_default();
let mut random = [0_u8; 12];
let _ = getrandom::getrandom(&mut random);
format!("{timestamp:013}-{}", hex::encode(random))
}
fn trimmed(value: &Option<String>) -> Option<String> {
value
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
fn now_ms() -> Result<i64> {
Ok(std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)?
.as_millis() as i64)
}
fn lease_ms() -> i64 {
std::env::var("OPENRTC_RTDB_PRESENCE_LEASE_MS")
.ok()
.and_then(|value| value.trim().parse::<i64>().ok())
.map(|value| value.clamp(LEASE_MS_MIN, LEASE_MS_MAX))
.unwrap_or(LEASE_MS_DEFAULT)
}
fn retryable_status(status: reqwest::StatusCode) -> bool {
status == reqwest::StatusCode::REQUEST_TIMEOUT
|| status == reqwest::StatusCode::TOO_MANY_REQUESTS
|| status.is_server_error()
}
fn retry_delay(attempt: usize) -> std::time::Duration {
let exponent = attempt.min(8) as u32;
let base = RETRY_BASE_MS
.saturating_mul(1_u64 << exponent)
.min(RETRY_MAX_MS);
let mut bytes = [0_u8; 2];
let jitter = if getrandom::getrandom(&mut bytes).is_ok() {
u64::from(u16::from_le_bytes(bytes)) % ((base / 2).max(1) + 1)
} else {
0
};
std::time::Duration::from_millis(base + jitter)
}
pub(crate) fn log_overlay(
backend: &str,
user_id: &str,
target: &str,
presence_count: usize,
stats: RtdbOverlayStats,
) {
log_info(&format!(
"[OPENRTC][SIGNALING][{backend}][rtdb] user_id={user_id} target={target} count={presence_count} matched={} promoted={} demoted={} rtdb_only_added={}",
stats.matched,
stats.promoted_online,
stats.demoted_offline,
stats.rtdb_only_added,
));
}
#[cfg(test)]
mod tests {
use super::*;
fn client(token: Option<&str>) -> RtdbPresenceClient {
client_with_platform(token, "desktop")
}
fn client_with_platform(token: Option<&str>, platform_type: &str) -> RtdbPresenceClient {
let token = token.map(str::to_string);
RtdbPresenceClient::new_with_platform_type(
"test-project",
"app_test",
Arc::new(Mutex::new(Box::new(move || token.clone()))),
platform_type,
)
}
fn device(id: &str, online: bool) -> Device {
Device {
app_tag: Some("app_test".to_string()),
device_id: id.to_string(),
user_id: Some("user-1".to_string()),
device_name: "Durable name".to_string(),
platform_type: Some("desktop".to_string()),
capabilities: None,
session_id: Some(id.to_string()),
node_id: None,
tag: Some("app_test".to_string()),
kind: Some("device".to_string()),
metadata: None,
online,
ticket: None,
last_seen_at: None,
expires_at: None,
created_at: None,
updated_at: None,
excluded_peers: Vec::new(),
}
}
#[test]
fn overlay_promotes_live_records_and_adds_rtdb_only_devices() {
let client = client(Some("token"));
let mut devices = vec![device("known", false)];
let presence = HashMap::from([
(
"known".to_string(),
RtdbPresenceState {
runtime_instance_id: Some("runtime-known".to_string()),
online: true,
last_seen_at: Some(10),
ticket_version: Some(10),
lease_expires_at: Some(200),
ticket: Some("known-ticket".to_string()),
node_id: Some("known-node".to_string()),
device_name: Some("Live name".to_string()),
platform_type: Some("desktop".to_string()),
metadata: None,
excluded_peers: vec!["local-device".to_string()],
},
),
(
"rtdb-only".to_string(),
RtdbPresenceState {
runtime_instance_id: Some("runtime-new".to_string()),
online: true,
last_seen_at: Some(20),
ticket_version: Some(20),
lease_expires_at: Some(200),
ticket: Some("new-ticket".to_string()),
node_id: Some("new-node".to_string()),
device_name: Some("New device".to_string()),
platform_type: Some("mobile".to_string()),
metadata: None,
excluded_peers: Vec::new(),
},
),
]);
let stats = client.overlay_devices(&mut devices, &presence, "user-1", None, 100);
assert_eq!(stats.matched, 1);
assert_eq!(stats.promoted_online, 1);
assert_eq!(stats.rtdb_only_added, 1);
assert_eq!(devices.len(), 2);
assert!(devices.iter().all(|item| item.online));
assert_eq!(devices[0].ticket.as_deref(), Some("known-ticket"));
assert_eq!(devices[0].session_id.as_deref(), Some("runtime-known"));
assert_eq!(devices[0].excluded_peers, vec!["local-device"]);
assert_eq!(devices[0].platform_type.as_deref(), Some("desktop"));
assert_eq!(devices[1].platform_type.as_deref(), Some("mobile"));
}
#[test]
fn overlay_demotes_an_expired_lease() {
let client = client(Some("token"));
let mut devices = vec![device("known", true)];
let presence = HashMap::from([(
"known".to_string(),
RtdbPresenceState {
runtime_instance_id: Some("runtime-expired".to_string()),
online: true,
last_seen_at: Some(10),
ticket_version: Some(10),
lease_expires_at: Some(99),
ticket: Some("ticket".to_string()),
node_id: Some("node".to_string()),
device_name: None,
platform_type: None,
metadata: None,
excluded_peers: Vec::new(),
},
)]);
let stats = client.overlay_devices(&mut devices, &presence, "user-1", None, 100);
assert_eq!(stats.demoted_offline, 1);
assert!(!devices[0].online);
}
#[test]
fn aggregate_keeps_device_live_when_one_runtime_instance_stops() {
let instances = HashMap::from([
(
"runtime-z".to_string(),
RtdbPresenceState {
runtime_instance_id: None,
online: false,
last_seen_at: Some(20),
ticket_version: Some(20),
lease_expires_at: Some(20),
ticket: None,
node_id: Some("stopped-node".to_string()),
device_name: None,
platform_type: Some("desktop".to_string()),
metadata: None,
excluded_peers: Vec::new(),
},
),
(
"runtime-a".to_string(),
RtdbPresenceState {
runtime_instance_id: None,
online: true,
last_seen_at: Some(10),
ticket_version: Some(10),
lease_expires_at: Some(200),
ticket: Some("surviving-ticket".to_string()),
node_id: Some("surviving-node".to_string()),
device_name: None,
platform_type: Some("mobile".to_string()),
metadata: None,
excluded_peers: Vec::new(),
},
),
]);
let aggregate = aggregate_device_presence(instances, 100).unwrap();
assert!(aggregate.online);
assert_eq!(aggregate.node_id.as_deref(), Some("surviving-node"));
assert_eq!(aggregate.runtime_instance_id.as_deref(), Some("runtime-a"));
assert_eq!(aggregate.platform_type.as_deref(), Some("mobile"));
}
#[test]
fn native_live_presence_publishes_the_explicit_target_platform() {
let mobile = client_with_platform(Some("token"), "mobile").online_state(
10,
"ticket",
"node",
"iOS Device",
None,
Vec::new(),
);
let desktop = client_with_platform(Some("token"), "desktop").online_state(
10,
"ticket",
"node",
"Desktop Device",
None,
Vec::new(),
);
assert_eq!(mobile.platform_type.as_deref(), Some("mobile"));
assert_eq!(desktop.platform_type.as_deref(), Some("desktop"));
}
#[tokio::test]
async fn read_fails_closed_without_an_auth_token() {
let error = client(None).read("user-1").await.unwrap_err();
assert!(error.to_string().contains("requires an auth token"));
}
}