openrtc 1.0.8

OpenRTC: a Rust-first P2P runtime for device discovery, signaling, and iroh/QUIC networking.
Documentation
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};

pub(crate) type TokenProvider = Arc<Mutex<Box<dyn Fn() -> Option<String> + Send + Sync>>>;
pub(crate) type TrustTokenProvider =
    Arc<Mutex<Box<dyn Fn() -> Option<crate::client::NativeTrustToken> + Send + Sync>>>;

const RENEW_EARLY_MS: u64 = 10_000;
const TRUST_TOKEN_EXPIRY_SKEW_MS: u64 = 5_000;
const APP_CHECK_HEADER: &str = "X-Firebase-AppCheck";
static NEXT_LEASE_ID: AtomicU64 = AtomicU64::new(1);

#[derive(Clone)]
pub(crate) struct LogicalUsageLeaseClient {
    http: reqwest::Client,
    project_id: String,
    token_provider: TokenProvider,
    trust_token_provider: TrustTokenProvider,
    instance_id: String,
    valid_until_by_key: Arc<Mutex<HashMap<String, u64>>>,
    disabled_for_unit_tests: bool,
}

#[derive(Serialize)]
struct CallableRequest<'a> {
    data: LeaseRequest<'a>,
}

#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct LeaseRequest<'a> {
    operation_id: &'a str,
    lease_id: &'a str,
    route: &'a str,
}

#[derive(Deserialize)]
struct CallableResponse {
    #[serde(alias = "data")]
    result: LeaseResponse,
}

#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct LeaseResponse {
    valid_until_ms: u64,
}

impl LogicalUsageLeaseClient {
    pub(crate) fn new(project_id: impl Into<String>, token_provider: TokenProvider) -> Self {
        Self::new_with_trust_token_provider(
            project_id,
            token_provider,
            Arc::new(Mutex::new(Box::new(|| None))),
        )
    }

    pub(crate) fn new_with_trust_token_provider(
        project_id: impl Into<String>,
        token_provider: TokenProvider,
        trust_token_provider: TrustTokenProvider,
    ) -> Self {
        let sequence = NEXT_LEASE_ID.fetch_add(1, Ordering::Relaxed);
        let now = crate::firebase::now_millis_u64();
        Self {
            http: reqwest::Client::new(),
            project_id: project_id.into(),
            token_provider,
            trust_token_provider,
            instance_id: format!("{now:x}-{sequence:x}"),
            valid_until_by_key: Arc::new(Mutex::new(HashMap::new())),
            disabled_for_unit_tests: cfg!(test),
        }
    }

    pub(crate) async fn ensure(
        &self,
        operation_id: &str,
        resource_key: &str,
        route: &str,
    ) -> Result<u64> {
        if self.disabled_for_unit_tests {
            return Ok(crate::firebase::now_millis_u64());
        }
        let cache_key = format!("{operation_id}\0{resource_key}\0{route}");
        let now = crate::firebase::now_millis_u64();
        if self
            .valid_until_by_key
            .lock()
            .ok()
            .and_then(|cache| cache.get(&cache_key).copied())
            .is_some_and(|valid_until| valid_until > now.saturating_add(RENEW_EARLY_MS))
        {
            return Ok(now);
        }

        let token = (self.token_provider.lock().unwrap())()
            .map(|value| value.trim().to_string())
            .filter(|value| !value.is_empty())
            .context("Logical usage lease requires an auth token")?;
        let lease_id = format!(
            "{}-{}",
            self.instance_id,
            sanitize_resource_key(resource_key),
        );
        let callable_url = self.callable_url();
        let trust_evidence = (self.trust_token_provider.lock().unwrap())();
        let trust_token = usable_trust_token(trust_evidence, now, is_emulator_url(&callable_url))?;
        let response = self
            .lease_request(
                callable_url,
                token,
                trust_token,
                CallableRequest {
                    data: LeaseRequest {
                        operation_id,
                        lease_id: &lease_id,
                        route,
                    },
                },
            )
            .send()
            .await
            .context("Logical usage lease request failed")?;
        let status = response.status();
        if !status.is_success() {
            let body = response.text().await.unwrap_or_default();
            anyhow::bail!(
                "Logical usage lease denied with status {}: {}",
                status,
                bounded_error(&body),
            );
        }
        let result = response
            .json::<CallableResponse>()
            .await
            .context("Logical usage lease returned an invalid response")?;
        self.valid_until_by_key
            .lock()
            .unwrap()
            .insert(cache_key, result.result.valid_until_ms);
        Ok(result.result.valid_until_ms)
    }

    fn callable_url(&self) -> String {
        #[cfg(target_arch = "wasm32")]
        {
            if let Some(host) = wasm_global_string("__OPENRTC_FUNCTIONS_EMULATOR_HOST__") {
                return format!(
                    "http://{host}/{}/us-central1/renewNativeUsageLease",
                    self.project_id
                );
            }
        }
        if let Ok(host) = std::env::var("OPENRTC_FUNCTIONS_EMULATOR_HOST") {
            let host = host.trim();
            if !host.is_empty() {
                return format!(
                    "http://{host}/{}/us-central1/renewNativeUsageLease",
                    self.project_id
                );
            }
        }
        format!(
            "https://us-central1-{}.cloudfunctions.net/renewNativeUsageLease",
            self.project_id
        )
    }

    fn lease_request(
        &self,
        callable_url: String,
        auth_token: String,
        trust_token: Option<crate::client::NativeTrustToken>,
        body: CallableRequest<'_>,
    ) -> reqwest::RequestBuilder {
        let mut request = self.http.post(callable_url).bearer_auth(auth_token);
        if let Some(trust_token) = trust_token {
            request = request.header(APP_CHECK_HEADER, trust_token.token);
        }
        request.json(&body)
    }
}

fn is_emulator_url(url: &str) -> bool {
    url.starts_with("http://127.0.0.1:")
        || url.starts_with("http://localhost:")
        || url.starts_with("http://[::1]:")
}

fn usable_trust_token(
    evidence: Option<crate::client::NativeTrustToken>,
    now: u64,
    emulator: bool,
) -> Result<Option<crate::client::NativeTrustToken>> {
    let trust_required = evidence.is_some();
    let usable = evidence.filter(|value| {
        !value.token.trim().is_empty()
            && value.expires_at_ms > now.saturating_add(TRUST_TOKEN_EXPIRY_SKEW_MS)
    });
    if trust_required && usable.is_none() && !emulator {
        anyhow::bail!("Logical usage lease requires a valid native trust token");
    }
    Ok(usable)
}

fn sanitize_resource_key(value: &str) -> String {
    let sanitized = value
        .chars()
        .map(|character| {
            if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
                character
            } else {
                '-'
            }
        })
        .take(72)
        .collect::<String>();
    if sanitized.is_empty() {
        "default".to_string()
    } else {
        sanitized
    }
}

fn bounded_error(value: &str) -> String {
    value.chars().take(256).collect()
}

#[cfg(target_arch = "wasm32")]
fn wasm_global_string(key: &str) -> Option<String> {
    let global = js_sys::global();
    js_sys::Reflect::get(&global, &wasm_bindgen::JsValue::from_str(key))
        .ok()
        .and_then(|value| value.as_string())
        .map(|value| value.trim().to_string())
        .filter(|value| !value.is_empty())
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn resource_keys_are_bounded_and_callable_safe() {
        assert_eq!(sanitize_resource_key("user/a:b"), "user-a-b");
        assert_eq!(sanitize_resource_key(""), "default");
        assert_eq!(sanitize_resource_key(&"x".repeat(100)).len(), 72);
    }

    #[test]
    fn legacy_live_usage_leases_remain_unconfigured_without_native_trust() {
        let auth: TokenProvider = Arc::new(Mutex::new(Box::new(|| Some("auth-token".to_string()))));
        let client = LogicalUsageLeaseClient::new("project-id", auth);
        let url = client.callable_url();
        assert!(!is_emulator_url(&url));
        let trust = (client.trust_token_provider.lock().unwrap())();
        assert!(trust.is_none());
    }

    #[test]
    fn native_trust_tokens_are_expiry_fenced() {
        let now = crate::firebase::now_millis_u64();
        let token = crate::client::NativeTrustToken {
            token: "opaque-token".to_string(),
            expires_at_ms: now + TRUST_TOKEN_EXPIRY_SKEW_MS + 1,
        };
        assert!(token.expires_at_ms > now + TRUST_TOKEN_EXPIRY_SKEW_MS);
        assert_eq!(APP_CHECK_HEADER, "X-Firebase-AppCheck");
    }

    #[test]
    fn configured_native_trust_fails_closed_after_clear_or_expiry() {
        let configured_but_cleared = crate::client::NativeTrustToken {
            token: String::new(),
            expires_at_ms: 0,
        };
        let now = crate::firebase::now_millis_u64();
        assert!(usable_trust_token(Some(configured_but_cleared), now, false).is_err());
        assert!(usable_trust_token(None, now, false).unwrap().is_none());
        assert!(usable_trust_token(
            Some(crate::client::NativeTrustToken {
                token: String::new(),
                expires_at_ms: 0,
            }),
            now,
            true,
        )
        .unwrap()
        .is_none());
    }

    #[test]
    fn native_trust_token_is_forwarded_only_in_a_header() {
        let auth: TokenProvider = Arc::new(Mutex::new(Box::new(|| None)));
        let client = LogicalUsageLeaseClient::new("project-id", auth);
        let trust_token = crate::client::NativeTrustToken {
            token: "opaque-trust-token".to_string(),
            expires_at_ms: u64::MAX,
        };
        let request = client
            .lease_request(
                "https://example.test/renew".to_string(),
                "firebase-auth-token".to_string(),
                Some(trust_token),
                CallableRequest {
                    data: LeaseRequest {
                        operation_id: "presence.lease",
                        lease_id: "lease-1",
                        route: "native-rtdb",
                    },
                },
            )
            .build()
            .expect("request should build");
        assert_eq!(
            request
                .headers()
                .get(APP_CHECK_HEADER)
                .and_then(|value| value.to_str().ok()),
            Some("opaque-trust-token"),
        );
        assert!(!request.url().as_str().contains("opaque-trust-token"));
        let body = request
            .body()
            .and_then(reqwest::Body::as_bytes)
            .expect("JSON body should be buffered");
        assert!(!String::from_utf8_lossy(body).contains("opaque-trust-token"));
    }
}