use super::client::FirestoreClient;
use super::device_registry::{
is_device_record_expired, offline_device_expires_at_ms, FirebaseDeviceRegistryContract,
SharedDeviceRecord, DEVICE_STALE_HEARTBEAT_MS,
};
use super::models::{bool_val, int_val, map_val, str_val, ts_val, DeviceFields};
use super::schema::AppNamespace;
use crate::logging::{env_flag_enabled, log_error, log_info};
use crate::signaling::{Device, SignalingBackend};
use anyhow::Result;
use async_trait::async_trait;
use iroh_tickets::endpoint::EndpointTicket;
use std::str::FromStr;
use std::sync::Arc;
fn signaling_verbose() -> bool {
env_flag_enabled("OPENRTC_SIGNALING_VERBOSE")
}
#[derive(Clone)]
pub struct FirebaseSignalingBackend {
client: FirestoreClient,
namespace: AppNamespace,
app_backgrounded_provider: Arc<dyn Fn() -> bool + Send + Sync>,
#[cfg(not(target_arch = "wasm32"))]
rtdb: super::rtdb_presence::RtdbPresenceClient,
}
impl FirebaseSignalingBackend {
pub fn new(client: FirestoreClient, app_tag: String) -> Self {
Self::new_with_background_provider(client, app_tag, Arc::new(|| false))
}
pub fn new_with_background_provider(
client: FirestoreClient,
app_tag: String,
app_backgrounded_provider: Arc<dyn Fn() -> bool + Send + Sync>,
) -> Self {
#[cfg(not(target_arch = "wasm32"))]
let rtdb = super::rtdb_presence::RtdbPresenceClient::new(
client.project_id(),
app_tag.clone(),
client.token_provider(),
);
Self {
client,
namespace: AppNamespace::new(app_tag),
app_backgrounded_provider,
#[cfg(not(target_arch = "wasm32"))]
rtdb,
}
}
fn parse_subscription_scope(scope: &str) -> (Option<String>, String) {
if let Some((user_id, device_id)) = scope.split_once("::") {
if !user_id.is_empty() && !device_id.is_empty() {
return (Some(user_id.to_string()), device_id.to_string());
}
}
(None, scope.to_string())
}
fn parse_session_scope(scope: &str) -> (Option<String>, String) {
if let Some((user_id, session_id)) = scope.split_once("::") {
if !user_id.is_empty() && !session_id.is_empty() {
return (Some(user_id.to_string()), session_id.to_string());
}
}
(None, scope.to_string())
}
fn extract_device_id(metadata: Option<&str>, fallback_node_id: &str) -> String {
let parsed = metadata
.and_then(|raw| serde_json::from_str::<serde_json::Value>(raw).ok())
.and_then(|value| value.as_object().cloned());
let candidate = parsed
.as_ref()
.and_then(|obj| {
obj.get("deviceId")
.or_else(|| obj.get("device_id"))
.and_then(|value| value.as_str())
})
.map(str::trim)
.filter(|value| !value.is_empty());
candidate
.map(|value| value.to_string())
.unwrap_or_else(|| fallback_node_id.to_string())
}
fn capabilities_value(
capabilities: &crate::signaling::DeviceCapabilities,
) -> super::models::FirestoreValue {
let mut fields = std::collections::HashMap::new();
fields.insert("canHost".to_string(), bool_val(capabilities.can_host));
fields.insert("canSync".to_string(), bool_val(capabilities.can_sync));
fields.insert("readOnly".to_string(), bool_val(capabilities.read_only));
map_val(fields)
}
async fn resolve_device_doc_id(&self, user_id: &str, device_id: &str) -> Result<String> {
let collection = self.namespace.devices_collection(user_id);
let docs = self
.client
.list_documents::<DeviceFields>(&collection)
.await?;
for doc in docs {
let id = doc.name.split('/').last().unwrap_or("").to_string();
let fields = doc.fields;
let doc_device_id = fields.device_id.and_then(|v| v.string_value);
let doc_node_id = fields.node_id.and_then(|v| v.string_value);
if id == device_id
|| doc_device_id.as_deref() == Some(device_id)
|| doc_node_id.as_deref() == Some(device_id)
{
return Ok(id);
}
}
Ok(device_id.to_string())
}
async fn resolve_existing_device_doc(
&self,
user_id: &str,
device_id: &str,
) -> Result<Option<(String, DeviceFields)>> {
let collection = self.namespace.devices_collection(user_id);
let docs = self
.client
.list_documents::<DeviceFields>(&collection)
.await?;
for doc in docs {
let id = doc.name.split('/').last().unwrap_or("").to_string();
let fields = doc.fields;
let doc_device_id = fields
.device_id
.as_ref()
.and_then(|v| v.string_value.as_ref())
.map(|value| value.as_str());
let doc_node_id = fields
.node_id
.as_ref()
.and_then(|v| v.string_value.as_ref())
.map(|value| value.as_str());
if id == device_id || doc_device_id == Some(device_id) || doc_node_id == Some(device_id)
{
return Ok(Some((id, fields)));
}
}
Ok(None)
}
fn canonical_session_collection(&self) -> String {
self.namespace.sessions_collection()
}
fn parse_i64_field(value: Option<&super::models::FirestoreValue>) -> Option<i64> {
value
.and_then(|v| v.integer_value.as_ref())
.and_then(|raw| raw.parse::<i64>().ok())
}
fn default_capabilities() -> super::models::FirestoreValue {
let mut fields = std::collections::HashMap::new();
fields.insert("canHost".to_string(), bool_val(false));
fields.insert("canSync".to_string(), bool_val(false));
fields.insert("readOnly".to_string(), bool_val(false));
map_val(fields)
}
fn parse_capabilities(
value: Option<&super::models::FirestoreValue>,
) -> Option<crate::signaling::DeviceCapabilities> {
let map_fields = value
.and_then(|v| v.map_value.as_ref())
.map(|m| &m.fields)?;
let can_host = map_fields
.get("canHost")
.and_then(|v| v.boolean_value)
.unwrap_or(false);
let can_sync = map_fields
.get("canSync")
.and_then(|v| v.boolean_value)
.unwrap_or(false);
let read_only = map_fields
.get("readOnly")
.and_then(|v| v.boolean_value)
.unwrap_or(false);
Some(crate::signaling::DeviceCapabilities {
can_host,
can_sync,
read_only,
})
}
fn is_online(fields: &super::models::DeviceFields) -> bool {
fields
.online
.as_ref()
.and_then(|v| v.boolean_value)
.unwrap_or(false)
}
async fn mark_device_doc_offline(
client: &FirestoreClient,
collection: &str,
doc_id: &str,
now_ms: i64,
) -> Result<()> {
let now_ts = super::now_timestamp_rfc3339();
let fields = DeviceFields {
app_tag: None,
user_id: None,
device_id: None,
device_name: None,
node_id: None,
platform_type: None,
session_id: None,
tag: None,
kind: None,
capabilities: None,
metadata: None,
online: Some(bool_val(false)),
ticket: None,
last_seen_at: None,
created_at: None,
updated_at: Some(ts_val(&now_ts)),
last_seen: None,
expires_at: Some(int_val(offline_device_expires_at_ms(now_ms))),
excluded_peers: None,
};
client
.update_document(
collection,
doc_id,
&fields,
Some(vec!["online", "updatedAt", "expiresAt"]),
)
.await
}
fn shared_device_record_from_fields(
fields: &super::models::DeviceFields,
) -> SharedDeviceRecord {
SharedDeviceRecord {
canonical_device_id: fields
.device_id
.as_ref()
.and_then(|v| v.string_value.clone()),
app_tag: fields.app_tag.as_ref().and_then(|v| v.string_value.clone()),
user_id: fields.user_id.as_ref().and_then(|v| v.string_value.clone()),
device_name: fields
.device_name
.as_ref()
.and_then(|v| v.string_value.clone())
.unwrap_or_default(),
platform_type: fields
.platform_type
.as_ref()
.and_then(|v| v.string_value.clone()),
capabilities: Self::parse_capabilities(fields.capabilities.as_ref()),
session_id: fields
.session_id
.as_ref()
.and_then(|v| v.string_value.clone()),
node_id: fields.node_id.as_ref().and_then(|v| v.string_value.clone()),
tag: fields.tag.as_ref().and_then(|v| v.string_value.clone()),
kind: fields.kind.as_ref().and_then(|v| v.string_value.clone()),
metadata: fields
.metadata
.as_ref()
.and_then(|v| v.string_value.clone()),
online: fields
.online
.as_ref()
.and_then(|v| v.boolean_value)
.unwrap_or(false),
ticket: fields.ticket.as_ref().and_then(|v| v.string_value.clone()),
last_seen_at: fields
.last_seen_at
.as_ref()
.and_then(|v| v.timestamp_value.clone())
.map(serde_json::Value::String),
expires_at: Self::parse_i64_field(fields.expires_at.as_ref())
.map(|value| serde_json::Value::Number(value.into())),
created_at: fields
.created_at
.as_ref()
.and_then(|v| v.timestamp_value.clone())
.map(serde_json::Value::String),
updated_at: fields
.updated_at
.as_ref()
.and_then(|v| v.timestamp_value.clone())
.map(serde_json::Value::String),
excluded_peers: fields
.excluded_peers
.as_ref()
.and_then(|v| v.array_value.as_ref())
.map(|arr| {
arr.values
.iter()
.filter_map(|v| v.string_value.clone())
.collect()
})
.unwrap_or_default(),
}
}
}
impl FirebaseDeviceRegistryContract for FirebaseSignalingBackend {
fn registry_app_tag(&self) -> &str {
self.namespace.app_tag()
}
fn registry_retention(&self) -> crate::firebase::device_registry::DiscoveryRetention {
if self.namespace.is_space() {
crate::firebase::device_registry::DiscoveryRetention::Ephemeral
} else {
crate::firebase::device_registry::DiscoveryRetention::Persistent
}
}
}
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
impl SignalingBackend for FirebaseSignalingBackend {
async fn update_presence(
&self,
user_id: &str,
local_node_id: &str,
ticket_str: &str,
is_online: bool,
name: &str,
ttl_ms: u64,
metadata: Option<&str>,
) -> Result<()> {
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING] update_presence start user_id={} node_id={} name={} online={} ticket_len={} metadata_len={} ttl_ms={}",
user_id,
local_node_id,
name,
is_online,
ticket_str.len(),
metadata.map(|v| v.len()).unwrap_or(0),
ttl_ms
));
}
let (iroh_ticket_part, _) = crate::session_token::split_compound_ticket(ticket_str.trim());
EndpointTicket::from_str(iroh_ticket_part)
.map_err(|e| anyhow::anyhow!("invalid endpoint ticket for presence: {}", e))?;
let ttl_i64 = ttl_ms.min(i64::MAX as u64) as i64;
let now_ms = super::now_millis_u64() as i64;
let now_ts = super::now_timestamp_rfc3339();
let canonical_device_id = Self::extract_device_id(metadata, local_node_id);
let collection = self.namespace.devices_collection(user_id);
let mut target_doc_id = canonical_device_id.clone();
let mut should_set_created_at = true;
let mut should_set_device_name = true;
let mut existing_created_at: Option<String> = None;
let mut existing_device_name: Option<String> = None;
if let Ok(existing_docs) = self
.client
.list_documents::<DeviceFields>(&collection)
.await
{
for doc in existing_docs {
let id = doc.name.split('/').last().unwrap_or("").to_string();
let fields = doc.fields;
let doc_device_id = fields
.device_id
.as_ref()
.and_then(|v| v.string_value.as_ref())
.map(|v| v.as_str());
let doc_node_id = fields
.node_id
.as_ref()
.and_then(|v| v.string_value.as_ref())
.map(|v| v.as_str());
if id == canonical_device_id
|| doc_device_id == Some(canonical_device_id.as_str())
|| doc_node_id == Some(canonical_device_id.as_str())
{
target_doc_id = id;
should_set_created_at = false;
existing_created_at = fields
.created_at
.as_ref()
.and_then(|value| value.timestamp_value.clone());
existing_device_name = fields
.device_name
.as_ref()
.and_then(|value| value.string_value.clone())
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty());
should_set_device_name = fields
.device_name
.and_then(|v| v.string_value)
.map(|value| value.trim().is_empty())
.unwrap_or(true);
break;
}
}
}
let ticket = Some(str_val(ticket_str));
#[cfg(target_arch = "wasm32")]
let platform_type = "web";
#[cfg(not(target_arch = "wasm32"))]
let platform_type = crate::native_device::default_platform_type();
let fields = DeviceFields {
app_tag: Some(str_val(self.namespace.app_tag())),
user_id: Some(str_val(user_id)),
device_id: Some(str_val(&canonical_device_id)),
device_name: if should_set_device_name {
Some(str_val(name))
} else if self.namespace.is_space() {
existing_device_name.as_deref().map(str_val)
} else {
None
},
node_id: Some(str_val(local_node_id)),
platform_type: Some(str_val(platform_type)),
session_id: Some(str_val(&canonical_device_id)),
tag: Some(str_val(self.namespace.app_tag())),
kind: Some(str_val("device")),
capabilities: Some(Self::default_capabilities()),
metadata: metadata.map(|m| str_val(m)),
online: Some(bool_val(is_online)),
ticket,
last_seen_at: Some(ts_val(&now_ts)),
created_at: Some(ts_val(existing_created_at.as_deref().unwrap_or(&now_ts))),
updated_at: Some(ts_val(&now_ts)),
last_seen: None,
expires_at: Some(int_val(now_ms.saturating_add(ttl_i64))),
excluded_peers: Some(super::models::str_array_val(&[])),
};
let mut paths = vec![
"appTag",
"userId",
"deviceId",
"nodeId",
"platformType",
"sessionId",
"tag",
"kind",
"capabilities",
"online",
"lastSeenAt",
"updatedAt",
"expiresAt",
"ticket",
];
if should_set_device_name {
paths.push("deviceName");
}
if should_set_created_at {
paths.push("createdAt");
}
if metadata.is_some() {
paths.push("metadata");
}
let update_mask = if self.namespace.is_space() {
None
} else {
Some(paths)
};
self.client
.update_document(&collection, &target_doc_id, &fields, update_mask)
.await?;
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING] update_presence success user_id={} node_id={}",
user_id, local_node_id
));
}
Ok(())
}
async fn set_offline(&self, user_id: &str, local_node_id: &str) -> Result<()> {
let collection = self.namespace.devices_collection(user_id);
let Some((target_doc_id, existing)) = self
.resolve_existing_device_doc(user_id, local_node_id)
.await?
else {
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING] set_offline skipped user_id={} node_id={} reason=device-doc-missing",
user_id, local_node_id
));
}
return Ok(());
};
let now_ts = super::now_timestamp_rfc3339();
let canonical_device_id = existing
.device_id
.as_ref()
.and_then(|v| v.string_value.as_ref())
.map(ToOwned::to_owned)
.unwrap_or_else(|| target_doc_id.clone());
let node_id = existing
.node_id
.as_ref()
.and_then(|v| v.string_value.as_ref())
.map(ToOwned::to_owned)
.unwrap_or_else(|| local_node_id.to_string());
let fields = DeviceFields {
app_tag: Some(str_val(self.namespace.app_tag())),
user_id: Some(str_val(user_id)),
device_id: Some(str_val(&canonical_device_id)),
device_name: None,
node_id: Some(str_val(&node_id)),
platform_type: None,
session_id: Some(str_val(&canonical_device_id)),
tag: Some(str_val(self.namespace.app_tag())),
kind: Some(str_val("device")),
capabilities: None,
metadata: None,
online: Some(bool_val(false)),
ticket: None,
last_seen_at: Some(ts_val(&now_ts)),
created_at: None,
updated_at: Some(ts_val(&now_ts)),
last_seen: None,
expires_at: Some(int_val(offline_device_expires_at_ms(
super::now_millis_u64() as i64,
))),
excluded_peers: None,
};
self.client
.update_document(
&collection,
&target_doc_id,
&fields,
Some(vec![
"appTag",
"userId",
"deviceId",
"nodeId",
"sessionId",
"tag",
"kind",
"online",
"lastSeenAt",
"updatedAt",
"expiresAt",
]),
)
.await?;
Ok(())
}
async fn update_live_presence(
&self,
user_id: &str,
local_node_id: &str,
ticket_str: &str,
name: &str,
metadata: Option<&str>,
) -> Result<()> {
#[cfg(target_arch = "wasm32")]
{
let _ = (user_id, local_node_id, ticket_str, name, metadata);
return Ok(());
}
#[cfg(not(target_arch = "wasm32"))]
{
if (self.app_backgrounded_provider)() {
return Ok(());
}
let (iroh_ticket, _) = crate::session_token::split_compound_ticket(ticket_str.trim());
EndpointTicket::from_str(iroh_ticket).map_err(|error| {
anyhow::anyhow!("invalid endpoint ticket for RTDB live presence: {error}")
})?;
let device_id = Self::extract_device_id(metadata, local_node_id);
self.rtdb
.publish(
user_id,
local_node_id,
&device_id,
ticket_str,
name,
metadata,
)
.await?;
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING][rest-native][rtdb] publish user_id={} device_id={} node_id={} target={}",
user_id,
device_id,
local_node_id,
self.rtdb.target_label(),
));
}
Ok(())
}
}
async fn set_live_presence_offline(&self, user_id: &str, local_node_id: &str) -> Result<()> {
#[cfg(target_arch = "wasm32")]
{
let _ = (user_id, local_node_id);
return Ok(());
}
#[cfg(not(target_arch = "wasm32"))]
self.rtdb.set_offline(user_id, local_node_id).await
}
async fn update_device(
&self,
user_id: &str,
device_id: &str,
device_name: Option<&str>,
capabilities: Option<crate::signaling::DeviceCapabilities>,
metadata: Option<&str>,
) -> Result<()> {
let trimmed_name = device_name.map(str::trim);
if trimmed_name.is_some_and(|value| value.is_empty()) {
anyhow::bail!("device_name cannot be empty");
}
if trimmed_name.is_none() && capabilities.is_none() && metadata.is_none() {
return Ok(());
}
let collection = self.namespace.devices_collection(user_id);
let target_doc_id = self.resolve_device_doc_id(user_id, device_id).await?;
let now_ts = super::now_timestamp_rfc3339();
let fields = DeviceFields {
app_tag: None,
user_id: None,
device_id: None,
device_name: trimmed_name.map(str_val),
node_id: None,
platform_type: None,
session_id: None,
tag: None,
kind: None,
capabilities: capabilities.as_ref().map(Self::capabilities_value),
metadata: metadata.map(str_val),
online: None,
ticket: None,
last_seen_at: None,
created_at: None,
updated_at: Some(ts_val(&now_ts)),
last_seen: None,
expires_at: None,
excluded_peers: None,
};
let mut update_mask = vec!["updatedAt"];
if trimmed_name.is_some() {
update_mask.push("deviceName");
}
if capabilities.is_some() {
update_mask.push("capabilities");
}
if metadata.is_some() {
update_mask.push("metadata");
}
self.client
.update_document(&collection, &target_doc_id, &fields, Some(update_mask))
.await?;
Ok(())
}
async fn delete_device(&self, user_id: &str, device_id: &str) -> Result<()> {
let target_doc_id = self.resolve_device_doc_id(user_id, device_id).await?;
let doc_path = format!(
"{}/{}",
self.namespace.devices_collection(user_id),
target_doc_id
);
self.client.delete_document(&doc_path).await?;
Ok(())
}
async fn set_excluded_peers(
&self,
user_id: &str,
local_node_id: &str,
excluded_peers: &[String],
) -> Result<()> {
#[cfg(not(target_arch = "wasm32"))]
{
return self
.rtdb
.set_excluded_peers(user_id, local_node_id, excluded_peers)
.await;
}
#[cfg(target_arch = "wasm32")]
{
let _ = (user_id, local_node_id, excluded_peers);
anyhow::bail!(
"browser auto-connect exclusions are owned by RTDB presence through WasmPresenceManager"
)
}
}
async fn search_devices(
&self,
user_id: &str,
_exclude_node_id: Option<&str>,
) -> Result<Vec<Device>> {
self.list_devices(user_id, _exclude_node_id)
.await
.map(|devices| devices.into_iter().filter(|device| device.online).collect())
}
async fn list_devices(
&self,
user_id: &str,
_exclude_node_id: Option<&str>,
) -> Result<Vec<Device>> {
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING] list_devices start user_id={}",
user_id
));
}
let collection = self.namespace.devices_collection(user_id);
let docs = self
.client
.list_documents::<DeviceFields>(&collection)
.await?;
let mut devices = Vec::new();
let now = super::now_millis_u64() as i64;
for doc in docs {
let id = doc.name.split('/').last().unwrap_or("").to_string();
let fields = doc.fields;
let shared = Self::shared_device_record_from_fields(&fields);
if shared.effective_app_tag() != Some(self.registry_app_tag()) {
continue;
}
if is_device_record_expired(&shared, now, DEVICE_STALE_HEARTBEAT_MS)
&& (shared.online || Self::is_online(&fields))
{
if let Err(error) =
Self::mark_device_doc_offline(&self.client, &collection, &id, now).await
{
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING] failed to mark stale device offline user_id={} device_id={} error={}",
user_id, id, error
));
}
}
}
if let Some(device) = self.project_discovered_device(id, shared, _exclude_node_id, now)
{
devices.push(device);
}
}
#[cfg(not(target_arch = "wasm32"))]
match self.rtdb.read(user_id).await {
Ok(presence) => {
let stats = self.rtdb.overlay_devices(
&mut devices,
&presence,
user_id,
_exclude_node_id,
now,
);
if signaling_verbose() {
super::rtdb_presence::log_overlay(
"rest-native",
user_id,
&self.rtdb.target_label(),
presence.len(),
stats,
);
}
}
Err(error) => {
log_error(&format!(
"[OPENRTC][SIGNALING][rest-native][rtdb] presence overlay unavailable user_id={} target={} error={}",
user_id,
self.rtdb.target_label(),
error,
));
}
}
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING] list_devices success user_id={} count={}",
user_id,
devices.len()
));
}
Ok(devices)
}
async fn send_message(
&self,
sender_id: &str,
target_id: &str,
payload: &str,
state: Option<&str>,
reply_payload: Option<&str>,
) -> Result<String> {
let now = super::now_millis_u64() as i64;
const MESSAGE_TTL_MS: i64 = 24 * 60 * 60 * 1000;
let expires_at = now + MESSAGE_TTL_MS;
let auth_uid = self.client.token_provider().lock().unwrap()()
.as_deref()
.and_then(super::auth_claims::firebase_uid_from_id_token);
let (sender_user_id, target_user_id) = if !self.namespace.is_space() {
let uid = auth_uid.ok_or_else(|| {
anyhow::anyhow!(
"send_message requires an authenticated Firebase user for app-scoped signaling"
)
})?;
(Some(str_val(&uid)), Some(str_val(&uid)))
} else if let Some(uid) = auth_uid {
(Some(str_val(&uid)), Some(str_val(&uid)))
} else {
(None, None)
};
let fields = super::models::SignalingEnvelopeFields {
app_tag: Some(str_val(self.namespace.app_tag())),
sender_id: Some(str_val(sender_id)),
target_id: Some(str_val(target_id)),
payload: Some(str_val(payload)),
state: state.map(|s| str_val(s)),
reply_payload: reply_payload.map(|r| str_val(r)),
timestamp: Some(int_val(now)),
sender_user_id,
target_user_id,
expires_at: Some(int_val(expires_at)),
};
let collection = self.namespace.messages_collection();
let doc_name = self
.client
.create_document(&collection, None, &fields)
.await?;
Ok(doc_name)
}
async fn subscribe_devices(
&self,
user_id: &str,
) -> Result<futures::stream::BoxStream<'static, Result<Vec<crate::signaling::DeviceEvent>>>>
{
let client = self.client.clone();
let collection = self.namespace.devices_collection(user_id);
let expected_tag = self.namespace.app_tag().to_string();
let project_id = self.client.project_id().to_string();
let token_provider = self.client.token_provider();
let (parent_path, collection_id) = {
let parts: Vec<&str> = collection.rsplitn(2, '/').collect();
if parts.len() == 2 {
(parts[1].to_string(), parts[0].to_string())
} else {
(String::new(), collection.clone())
}
};
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
Self::spawn_device_listener(
tx,
client,
collection,
expected_tag,
self.registry_retention(),
project_id,
token_provider,
self.app_backgrounded_provider.clone(),
parent_path,
collection_id,
);
use futures::StreamExt;
use tokio_stream::wrappers::UnboundedReceiverStream;
Ok(UnboundedReceiverStream::new(rx).boxed())
}
async fn create_session(&self, session: crate::signaling::SignalingSession) -> Result<()> {
self.create_session_impl(session).await
}
async fn update_session(&self, session_id: &str, update_data: serde_json::Value) -> Result<()> {
self.update_session_impl(session_id, update_data).await
}
async fn subscribe_sessions(
&self,
local_device_id: &str,
) -> Result<futures::stream::BoxStream<'static, Result<Vec<crate::signaling::SessionEvent>>>>
{
self.subscribe_sessions_impl(local_device_id).await
}
}
impl FirebaseSignalingBackend {
fn spawn_device_listener(
tx: tokio::sync::mpsc::UnboundedSender<Result<Vec<crate::signaling::DeviceEvent>>>,
client: FirestoreClient,
collection: String,
expected_tag: String,
retention: crate::firebase::device_registry::DiscoveryRetention,
project_id: String,
token_provider: std::sync::Arc<
std::sync::Mutex<Box<dyn Fn() -> Option<String> + Send + Sync>>,
>,
app_backgrounded_provider: Arc<dyn Fn() -> bool + Send + Sync>,
parent_path: String,
collection_id: String,
) {
let spawn_body = async move {
let mut listener = super::listener::FirestoreListener::new_with_label(
project_id.clone(),
token_provider.clone(),
format!("signaling-devices:{}", collection),
);
listener.add_target(1, &parent_path, &collection_id);
Self::run_device_listener_loop(
&tx,
&mut listener,
&expected_tag,
retention,
&client,
&collection,
app_backgrounded_provider.as_ref(),
)
.await;
};
#[cfg(target_arch = "wasm32")]
wasm_bindgen_futures::spawn_local(spawn_body);
#[cfg(not(target_arch = "wasm32"))]
tokio::spawn(spawn_body);
}
async fn run_device_listener_loop(
tx: &tokio::sync::mpsc::UnboundedSender<Result<Vec<crate::signaling::DeviceEvent>>>,
listener: &mut super::listener::FirestoreListener,
expected_tag: &str,
retention: crate::firebase::device_registry::DiscoveryRetention,
client: &FirestoreClient,
collection: &str,
app_backgrounded: &(dyn Fn() -> bool + Send + Sync),
) {
loop {
if tx.is_closed() {
break;
}
match listener.poll().await {
super::listener::ListenOutcome::Events(listen_events) => {
let mut device_events = Vec::new();
for event in listen_events {
match event {
super::webchannel::ListenEvent::DocumentChange { document, .. } => {
if let Some(device_event) = Self::parse_device_change(
&document,
expected_tag,
retention,
client,
collection,
)
.await
{
device_events.push(device_event);
}
}
super::webchannel::ListenEvent::DocumentDelete {
document_path,
..
}
| super::webchannel::ListenEvent::DocumentRemove {
document_path,
..
} => {
let device_id = document_path
.split('/')
.last()
.unwrap_or(&document_path)
.to_string();
device_events
.push(crate::signaling::DeviceEvent::Removed { device_id });
}
super::webchannel::ListenEvent::TargetChange { .. }
| super::webchannel::ListenEvent::Filter { .. } => {
}
}
}
if !device_events.is_empty() {
if tx.send(Ok(device_events)).is_err() {
break;
}
}
}
super::listener::ListenOutcome::AuthError => {
let delay_ms =
Self::reconnect_delay_ms(listener, app_backgrounded(), Some("auth"));
log_info(&format!(
"[OPENRTC][SIGNALING][{}] auth error, reconnecting after {}ms...",
listener.debug_label(),
delay_ms
));
Self::sleep_ms(delay_ms).await;
}
super::listener::ListenOutcome::FatalError(msg) => {
log_error(&format!(
"[OPENRTC][SIGNALING][{}] fatal error: {}",
listener.debug_label(),
msg
));
let _ = tx.send(Err(anyhow::anyhow!(msg)));
break;
}
super::listener::ListenOutcome::TransientError(msg) => {
let delay_ms =
Self::reconnect_delay_ms(listener, app_backgrounded(), Some(&msg));
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING][{}] transient error: {}, backing off {}ms",
listener.debug_label(),
msg,
delay_ms
));
}
Self::sleep_ms(delay_ms).await;
}
}
}
}
async fn parse_device_change(
document: &serde_json::Value,
expected_tag: &str,
retention: crate::firebase::device_registry::DiscoveryRetention,
client: &FirestoreClient,
collection: &str,
) -> Option<crate::signaling::DeviceEvent> {
let doc_name = document.get("name")?.as_str()?;
let id = doc_name.split('/').last()?.to_string();
let fields: super::models::DeviceFields =
serde_json::from_value(document.get("fields")?.clone()).ok()?;
let shared = Self::shared_device_record_from_fields(&fields);
if shared.effective_app_tag() != Some(expected_tag) {
return None;
}
let now = super::now_millis_u64() as i64;
if is_device_record_expired(&shared, now, DEVICE_STALE_HEARTBEAT_MS)
&& (shared.online || FirebaseSignalingBackend::is_online(&fields))
{
let _ = FirebaseSignalingBackend::mark_device_doc_offline(client, collection, &id, now)
.await;
}
if !crate::firebase::device_registry::is_device_record_discoverable(&shared, now, retention)
{
return Some(crate::signaling::DeviceEvent::Removed {
device_id: shared.public_device_id(&id),
});
}
let device = shared.into_public_device(id);
Some(crate::signaling::DeviceEvent::Modified { device })
}
async fn sleep_ms(ms: u64) {
#[cfg(target_arch = "wasm32")]
{
let _ = gloo_timers::future::sleep(std::time::Duration::from_millis(ms)).await;
}
#[cfg(not(target_arch = "wasm32"))]
{
tokio::time::sleep(std::time::Duration::from_millis(ms)).await;
}
}
async fn create_session_impl(&self, session: crate::signaling::SignalingSession) -> Result<()> {
let primary_collection = self.canonical_session_collection();
let fields = super::models::SignalingSessionFields {
app_tag: Some(super::models::str_val(self.namespace.app_tag())),
connection_id: Some(super::models::str_val(&session.connection_id)),
initiator_device_id: Some(super::models::str_val(&session.initiator_device_id)),
target_device_id: Some(super::models::str_val(&session.target_device_id)),
state: Some(super::models::str_val(&session.state)),
offer: session
.offer
.as_ref()
.map(|s| super::models::str_val(&serde_json::to_string(s).unwrap_or_default())),
answer: session
.answer
.as_ref()
.map(|s| super::models::str_val(&serde_json::to_string(s).unwrap_or_default())),
created_at: session.created_at.map(|t| super::models::int_val(t)),
};
self.client
.update_document(&primary_collection, &session.connection_id, &fields, None)
.await?;
Ok(())
}
async fn update_session_impl(
&self,
session_id: &str,
update_data: serde_json::Value,
) -> Result<()> {
let (_, resolved_session_id) = Self::parse_session_scope(session_id);
let primary_collection = self.canonical_session_collection();
let mut fields = std::collections::HashMap::new();
let mut update_mask = Vec::new();
if let Some(obj) = update_data.as_object() {
for (k, v) in obj {
let mask_key = k.clone();
update_mask.push(mask_key);
if let Some(s) = v.as_str() {
fields.insert(k.clone(), super::models::str_val(s));
} else if let Some(i) = v.as_i64() {
fields.insert(k.clone(), super::models::int_val(i));
} else if let Some(b) = v.as_bool() {
fields.insert(k.clone(), super::models::bool_val(b));
}
}
}
let paths: Vec<&str> = update_mask.iter().map(|s| s.as_str()).collect();
self.client
.update_document(
&primary_collection,
&resolved_session_id,
&fields,
Some(paths),
)
.await?;
Ok(())
}
async fn subscribe_sessions_impl(
&self,
local_device_id: &str,
) -> Result<futures::stream::BoxStream<'static, Result<Vec<crate::signaling::SessionEvent>>>>
{
let (_, target_id) = Self::parse_subscription_scope(local_device_id);
let collection = self.canonical_session_collection();
let project_id = self.client.project_id().to_string();
let token_provider = self.client.token_provider();
let (parent_path, collection_id) = {
let parts: Vec<&str> = collection.rsplitn(2, '/').collect();
if parts.len() == 2 {
(parts[1].to_string(), parts[0].to_string())
} else {
(String::new(), collection.clone())
}
};
let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
Self::spawn_session_listener(
tx,
collection,
target_id,
project_id,
token_provider,
self.app_backgrounded_provider.clone(),
parent_path,
collection_id,
);
use futures::StreamExt;
use tokio_stream::wrappers::UnboundedReceiverStream;
Ok(UnboundedReceiverStream::new(rx).boxed())
}
fn spawn_session_listener(
tx: tokio::sync::mpsc::UnboundedSender<Result<Vec<crate::signaling::SessionEvent>>>,
collection: String,
target_id: String,
project_id: String,
token_provider: std::sync::Arc<
std::sync::Mutex<Box<dyn Fn() -> Option<String> + Send + Sync>>,
>,
app_backgrounded_provider: Arc<dyn Fn() -> bool + Send + Sync>,
parent_path: String,
collection_id: String,
) {
let spawn_body = async move {
let mut listener = super::listener::FirestoreListener::new_with_label(
project_id.clone(),
token_provider.clone(),
format!("signaling-sessions:{}:{}", target_id, collection),
);
listener.add_target(2, &parent_path, &collection_id);
Self::run_session_listener_loop(
&tx,
&mut listener,
&target_id,
&collection,
app_backgrounded_provider.as_ref(),
)
.await;
};
#[cfg(target_arch = "wasm32")]
wasm_bindgen_futures::spawn_local(spawn_body);
#[cfg(not(target_arch = "wasm32"))]
tokio::spawn(spawn_body);
}
async fn run_session_listener_loop(
tx: &tokio::sync::mpsc::UnboundedSender<Result<Vec<crate::signaling::SessionEvent>>>,
listener: &mut super::listener::FirestoreListener,
target_id: &str,
collection: &str,
app_backgrounded: &(dyn Fn() -> bool + Send + Sync),
) {
loop {
if tx.is_closed() {
break;
}
match listener.poll().await {
super::listener::ListenOutcome::Events(listen_events) => {
let mut session_events = Vec::new();
for event in listen_events {
match event {
super::webchannel::ListenEvent::DocumentChange { document, .. } => {
if let Some(session_event) =
Self::parse_session_change(&document, target_id, collection)
{
session_events.push(session_event);
}
}
super::webchannel::ListenEvent::DocumentDelete {
document_path,
..
}
| super::webchannel::ListenEvent::DocumentRemove {
document_path,
..
} => {
let session_id = document_path
.split('/')
.last()
.unwrap_or(&document_path)
.to_string();
session_events
.push(crate::signaling::SessionEvent::Removed { session_id });
}
super::webchannel::ListenEvent::TargetChange { .. }
| super::webchannel::ListenEvent::Filter { .. } => {}
}
}
if !session_events.is_empty() {
if tx.send(Ok(session_events)).is_err() {
break;
}
}
}
super::listener::ListenOutcome::AuthError => {
let delay_ms =
Self::reconnect_delay_ms(listener, app_backgrounded(), Some("auth"));
log_info(&format!(
"[OPENRTC][SIGNALING][{}] auth error, reconnecting after {}ms...",
listener.debug_label(),
delay_ms
));
Self::sleep_ms(delay_ms).await;
}
super::listener::ListenOutcome::FatalError(msg) => {
log_error(&format!(
"[OPENRTC][SIGNALING][{}] fatal error: {}",
listener.debug_label(),
msg
));
let _ = tx.send(Err(anyhow::anyhow!(msg)));
break;
}
super::listener::ListenOutcome::TransientError(msg) => {
let delay_ms =
Self::reconnect_delay_ms(listener, app_backgrounded(), Some(&msg));
if signaling_verbose() {
log_info(&format!(
"[OPENRTC][SIGNALING][{}] transient error: {}, backing off {}ms",
listener.debug_label(),
msg,
delay_ms
));
}
Self::sleep_ms(delay_ms).await;
}
}
}
}
fn reconnect_delay_ms(
listener: &super::listener::FirestoreListener,
app_backgrounded: bool,
error_message: Option<&str>,
) -> u64 {
const BACKGROUND_RECONNECT_DELAY_MS: u64 = 30_000;
if app_backgrounded && Self::should_suppress_background_reconnect(error_message) {
return BACKGROUND_RECONNECT_DELAY_MS;
}
listener.backoff_duration().as_millis() as u64
}
fn should_suppress_background_reconnect(error_message: Option<&str>) -> bool {
let normalized = error_message
.map(|message| message.trim().to_ascii_lowercase())
.unwrap_or_default();
normalized.is_empty()
|| normalized.contains("parse")
|| normalized.contains("backward-channel")
|| normalized.contains("webchannel")
|| normalized.contains("auth")
}
fn parse_session_change(
document: &serde_json::Value,
target_id: &str,
collection: &str,
) -> Option<crate::signaling::SessionEvent> {
let doc_name = document.get("name")?.as_str()?;
let id = doc_name.split('/').last()?.to_string();
let fields: super::models::SignalingSessionFields =
serde_json::from_value(document.get("fields")?.clone()).ok()?;
let initiator = fields
.initiator_device_id
.as_ref()
.and_then(|v| v.string_value.as_deref())
.unwrap_or_default()
.to_string();
let target_dst = fields
.target_device_id
.as_ref()
.and_then(|v| v.string_value.as_deref())
.unwrap_or_default()
.to_string();
if initiator != target_id && target_dst != target_id {
return None;
}
let session = crate::signaling::SignalingSession {
connection_id: id,
initiator: "".to_string(),
target: "".to_string(),
initiator_device_id: initiator,
target_device_id: target_dst,
connection_type: None,
offer: fields
.offer
.and_then(|v| v.string_value)
.and_then(|s| serde_json::from_str(&s).ok()),
offer_e2ee: None,
answer: fields
.answer
.and_then(|v| v.string_value)
.and_then(|s| serde_json::from_str(&s).ok()),
answer_e2ee: None,
ice_candidates: vec![],
initiator_node_id: None,
target_node_id: None,
initiator_endpoint_addr: None,
target_endpoint_addr: None,
intent: None,
app_tag: fields
.app_tag
.and_then(|v| v.string_value)
.or_else(|| Some(collection.split('/').nth(1).unwrap_or_default().to_string()))
.filter(|value| !value.is_empty()),
created_at: fields
.created_at
.and_then(|v| v.integer_value)
.unwrap_or_default()
.parse()
.ok(),
expires_at: None,
state: fields
.state
.and_then(|v| v.string_value)
.unwrap_or_default(),
};
Some(crate::signaling::SessionEvent::Modified { session })
}
}
#[cfg(test)]
mod tests {
use super::FirebaseSignalingBackend;
#[test]
fn parses_scoped_subscription_key() {
let (user_id, device_id) =
FirebaseSignalingBackend::parse_subscription_scope("user-1::device-1");
assert_eq!(user_id.as_deref(), Some("user-1"));
assert_eq!(device_id, "device-1");
}
#[test]
fn keeps_legacy_subscription_key() {
let (user_id, device_id) =
FirebaseSignalingBackend::parse_subscription_scope("device-only");
assert!(user_id.is_none());
assert_eq!(device_id, "device-only");
}
#[test]
fn parses_scoped_session_key() {
let (user_id, session_id) =
FirebaseSignalingBackend::parse_session_scope("user-1::session-1");
assert_eq!(user_id.as_deref(), Some("user-1"));
assert_eq!(session_id, "session-1");
}
#[test]
fn keeps_legacy_session_key() {
let (user_id, session_id) = FirebaseSignalingBackend::parse_session_scope("session-only");
assert!(user_id.is_none());
assert_eq!(session_id, "session-only");
}
}