diff --git a/Cargo.lock b/Cargo.lock index 3ed3d6e7a..d56e8d70b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1990,13 +1990,10 @@ dependencies = [ "anyhow", "axum", "base64", - "fluxer-svc", "fluxer_common", "hex", - "moka", "rand 0.10.1", "reqwest", - "scylla", "serde", "serde_json", "tokio", diff --git a/deploy/self-hosting/.env.example b/deploy/self-hosting/.env.example index 403332646..2ee7a2a40 100644 --- a/deploy/self-hosting/.env.example +++ b/deploy/self-hosting/.env.example @@ -140,7 +140,7 @@ FLUXER_DISCOVERY_ENABLED=true # budget roughly shared_buffers + (server max_connections x 12 MB) + # (3 x autovacuum_work_mem) + 300 MB for page cache and WAL. Note this is the # server setting, distinct from the per-service FLUXER_POSTGRES_MAX_CONNECTIONS -# pool sizes used by the api, worker, app-proxy and shards. +# pool sizes used by the api, worker and shards. #FLUXER_POSTGRES_SERVER_MAX_CONNECTIONS=150 #FLUXER_POSTGRES_SHARED_BUFFERS=512MB #FLUXER_POSTGRES_EFFECTIVE_CACHE_SIZE=2GB diff --git a/deploy/self-hosting/docker-compose.yml b/deploy/self-hosting/docker-compose.yml index eef279bad..229f97a83 100644 --- a/deploy/self-hosting/docker-compose.yml +++ b/deploy/self-hosting/docker-compose.yml @@ -460,7 +460,6 @@ services: limits: memory: ${FLUXER_APP_PROXY_MEMORY_LIMIT:-256mb} environment: - <<: *fluxer-postgres-env FLUXER_APP_PROXY_HOST: 0.0.0.0 FLUXER_APP_PROXY_PORT: "8080" DISCOVERY_UPSTREAM_URL: http://caddy:8088/api/.well-known/fluxer @@ -477,11 +476,9 @@ services: FLUXER_CSP_EXTRA_WORKER_SRC: ${FLUXER_CSP_EXTRA_WORKER_SRC:-} FLUXER_CSP_EXTRA_MANIFEST_SRC: ${FLUXER_CSP_EXTRA_MANIFEST_SRC:-} FLUXER_CSP_REPORT_URI: ${FLUXER_CSP_REPORT_URI:-} - FLUXER_POSTGRES_MAX_CONNECTIONS: "5" depends_on: api: {condition: service_healthy} caddy: {condition: service_healthy} - postgres: {condition: service_healthy} snowflakes: <<: *fluxer-service diff --git a/fluxer_app_proxy/Cargo.toml b/fluxer_app_proxy/Cargo.toml index 752ebc599..76a9b6213 100644 --- a/fluxer_app_proxy/Cargo.toml +++ b/fluxer_app_proxy/Cargo.toml @@ -10,12 +10,9 @@ anyhow = "1.0.104" axum = { version = "0.8.9", features = ["macros"] } base64 = "0.22" fluxer_common = { path = "../fluxer_common" } -fluxer-svc = { path = "../fluxer_svc" } hex = "0.4" -moka = { version = "0.12.15", features = ["future"] } rand = "0.10" reqwest = { version = "0.13.4", default-features = false, features = ["json", "rustls", "stream", "gzip", "brotli", "deflate"] } -scylla = { version = "1.6.0", features = ["chrono-04"], optional = true } serde = { version = "1.0.228", features = ["derive"] } serde_json = "1.0.150" tokio = { version = "1.52.3", features = ["macros", "net", "rt-multi-thread", "signal", "time", "fs"] } @@ -27,5 +24,4 @@ tracing-subscriber = { version = "0.3.23", features = ["env-filter"] } [features] default = ["time-freeze"] -scylla = ["dep:scylla", "fluxer-svc/scylla"] time-freeze = [] diff --git a/fluxer_app_proxy/Dockerfile b/fluxer_app_proxy/Dockerfile index 783a485ec..ad48bba90 100644 --- a/fluxer_app_proxy/Dockerfile +++ b/fluxer_app_proxy/Dockerfile @@ -118,14 +118,12 @@ WORKDIR /usr/src/app COPY Cargo.lock ./ COPY fluxer_common/ fluxer_common/ COPY fluxer_app_proxy/ fluxer_app_proxy/ -COPY fluxer_svc/ fluxer_svc/ RUN printf '%s\n' \ '[workspace]' \ 'members = [' \ ' "fluxer_app_proxy",' \ ' "fluxer_common",' \ - ' "fluxer_svc",' \ ']' \ 'resolver = "2"' \ '' \ diff --git a/fluxer_app_proxy/src/config.rs b/fluxer_app_proxy/src/config.rs index f69c67675..be52f8e9a 100644 --- a/fluxer_app_proxy/src/config.rs +++ b/fluxer_app_proxy/src/config.rs @@ -1,7 +1,6 @@ // SPDX-License-Identifier: AGPL-3.0-or-later use fluxer_common::config::{self as cfg, GeoipS3Config, GeoipSourceConfig}; -use fluxer_svc::config::{DatabaseBackend, normalize_host, parse_hosts}; use reqwest::Url; use std::env; use std::fmt; @@ -368,25 +367,6 @@ pub struct AppProxyConfig { pub geoip_s3_config: Option, pub trust_client_ip_header: bool, pub client_ip_header_name: String, - pub invite_meta_enabled: bool, - pub invite_meta_cache_max_entries: u64, - pub invite_meta_cache_ttl_ms: u64, - pub database_backend: DatabaseBackend, - pub scylla_hosts: Vec, - pub scylla_keyspace: String, - pub scylla_username: Option, - pub scylla_password: Option, - pub postgres_url: Option, - pub postgres_host: String, - pub postgres_port: u16, - pub postgres_database: String, - pub postgres_username: String, - pub postgres_password: Option, - pub postgres_ssl: bool, - pub postgres_ssl_ca: Option, - pub postgres_max_connections: usize, - pub postgres_kv_table: String, - pub postgres_prepared_statements: bool, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -488,34 +468,6 @@ impl AppProxyConfig { ); let geoip_s3_config = cfg::read_geoip_s3_config_from_env(&geoip_source); - let cassandra_port = parse_env_or_warn( - "FLUXER_CASSANDRA_PORT", - &cfg::read_env("FLUXER_CASSANDRA_PORT", "9042"), - 9042u16, - ); - let scylla_hosts = cfg::non_empty_env("FLUXER_CASSANDRA_HOSTS") - .map(|hosts| { - parse_hosts(&hosts) - .into_iter() - .map(|host| normalize_host(&host, cassandra_port)) - .collect::>() - }) - .filter(|hosts| !hosts.is_empty()) - .unwrap_or_else(|| vec![normalize_host("127.0.0.1", cassandra_port)]); - let database_backend = - parse_database_backend(&cfg::read_env("FLUXER_DATABASE_BACKEND", "postgres")); - let postgres_port = parse_env_or_warn( - "FLUXER_POSTGRES_PORT", - &cfg::read_env("FLUXER_POSTGRES_PORT", "5432"), - 5432u16, - ); - let postgres_max_connections = parse_env_or_warn( - "FLUXER_POSTGRES_MAX_CONNECTIONS", - &cfg::read_env("FLUXER_POSTGRES_MAX_CONNECTIONS", "20"), - 20usize, - ) - .max(1); - let s3_public_endpoint = parse_optional_http_endpoint( "FLUXER_S3_PUBLIC_ENDPOINT", cfg::non_empty_env("FLUXER_S3_PUBLIC_ENDPOINT"), @@ -581,48 +533,10 @@ impl AppProxyConfig { ) .trim() .to_ascii_lowercase(), - invite_meta_enabled: cfg::read_bool_env( - &["FLUXER_APP_PROXY_INVITE_META_ENABLED"], - true, - ), - invite_meta_cache_max_entries: parse_env_or_warn( - "FLUXER_APP_PROXY_INVITE_META_CACHE_MAX_ENTRIES", - &cfg::read_env("FLUXER_APP_PROXY_INVITE_META_CACHE_MAX_ENTRIES", "10000"), - 10_000u64, - ), - invite_meta_cache_ttl_ms: parse_env_or_warn( - "FLUXER_APP_PROXY_INVITE_META_CACHE_TTL_MS", - &cfg::read_env("FLUXER_APP_PROXY_INVITE_META_CACHE_TTL_MS", "30000"), - 30_000u64, - ), - database_backend, - scylla_hosts, - scylla_keyspace: cfg::read_env("FLUXER_CASSANDRA_KEYSPACE", "fluxer"), - scylla_username: cfg::non_empty_env("FLUXER_CASSANDRA_USERNAME"), - scylla_password: cfg::non_empty_env("FLUXER_CASSANDRA_PASSWORD"), - postgres_url: cfg::non_empty_env("FLUXER_POSTGRES_URL"), - postgres_host: cfg::read_env("FLUXER_POSTGRES_HOST", "127.0.0.1"), - postgres_port, - postgres_database: cfg::read_env("FLUXER_POSTGRES_DATABASE", "fluxer"), - postgres_username: cfg::read_env("FLUXER_POSTGRES_USERNAME", "fluxer"), - postgres_password: cfg::non_empty_env("FLUXER_POSTGRES_PASSWORD") - .or_else(|| Some("fluxer".to_owned())), - postgres_ssl: cfg::read_bool_env(&["FLUXER_POSTGRES_SSL"], false), - postgres_ssl_ca: cfg::non_empty_env("FLUXER_POSTGRES_SSL_CA"), - postgres_max_connections, - postgres_kv_table: cfg::read_env("FLUXER_POSTGRES_KV_TABLE", "fluxer_kv"), - postgres_prepared_statements: resolve_postgres_prepared_statements_from_env(), } } } -fn parse_database_backend(value: &str) -> DatabaseBackend { - match value.trim().to_ascii_lowercase().as_str() { - "cassandra" | "scylla" | "scylladb" => DatabaseBackend::Cassandra, - _ => DatabaseBackend::Postgres, - } -} - fn resolve_discovery_upstream_url_from_env() -> String { resolve_discovery_upstream_url(|name| env::var(name).ok()) } @@ -631,10 +545,6 @@ fn resolve_time_freeze_enabled_from_env() -> bool { resolve_time_freeze_enabled(|name| env::var(name).ok()) } -fn resolve_postgres_prepared_statements_from_env() -> bool { - resolve_postgres_prepared_statements(|name| env::var(name).ok()) -} - fn resolve_bootstrap_api_public_endpoint_from_env() -> Option { resolve_bootstrap_api_public_endpoint(|name| env::var(name).ok()) } @@ -674,31 +584,6 @@ fn parse_boolish(value: &str) -> bool { ) } -fn resolve_postgres_prepared_statements(mut read_var: F) -> bool -where - F: FnMut(&str) -> Option, -{ - let Some(value) = read_var("FLUXER_POSTGRES_PREPARED_STATEMENTS") - .map(|value| value.trim().to_ascii_lowercase()) - .filter(|value| !value.is_empty()) - else { - return true; - }; - - match value.as_str() { - "1" | "true" | "yes" | "y" | "on" => true, - "0" | "false" | "no" | "n" | "off" => false, - other => { - tracing::warn!( - env = "FLUXER_POSTGRES_PREPARED_STATEMENTS", - value = other, - "invalid value; falling back to default" - ); - true - } - } -} - fn resolve_discovery_upstream_url(mut read_var: F) -> String where F: FnMut(&str) -> Option, @@ -749,11 +634,6 @@ mod tests { resolve_time_freeze_enabled(|name| env.get(name).map(|value| value.to_string())) } - fn resolve_prepared_statements_from_pairs(pairs: &[(&str, &str)]) -> bool { - let env: HashMap<&str, &str> = pairs.iter().copied().collect(); - resolve_postgres_prepared_statements(|name| env.get(name).map(|value| value.to_string())) - } - fn resolve_bootstrap_endpoint_from_pairs(pairs: &[(&str, &str)]) -> Option { let env: HashMap<&str, &str> = pairs.iter().copied().collect(); resolve_bootstrap_api_public_endpoint(|name| env.get(name).map(|value| value.to_string())) @@ -905,43 +785,6 @@ mod tests { ); } - #[test] - fn a_set_but_empty_prepared_statements_value_keeps_the_shared_default() { - assert!(resolve_prepared_statements_from_pairs(&[])); - assert!( - resolve_prepared_statements_from_pairs(&[("FLUXER_POSTGRES_PREPARED_STATEMENTS", "")]), - "an empty value disabled named statements here while every other service kept them" - ); - assert!(resolve_prepared_statements_from_pairs(&[( - "FLUXER_POSTGRES_PREPARED_STATEMENTS", - " ", - )])); - } - - #[test] - fn an_explicit_prepared_statements_value_is_honoured() { - assert!(!resolve_prepared_statements_from_pairs(&[( - "FLUXER_POSTGRES_PREPARED_STATEMENTS", - "false", - )])); - assert!(!resolve_prepared_statements_from_pairs(&[( - "FLUXER_POSTGRES_PREPARED_STATEMENTS", - "OFF", - )])); - assert!(resolve_prepared_statements_from_pairs(&[( - "FLUXER_POSTGRES_PREPARED_STATEMENTS", - "yes", - )])); - } - - #[test] - fn a_non_boolean_prepared_statements_value_keeps_the_shared_default() { - assert!(resolve_prepared_statements_from_pairs(&[( - "FLUXER_POSTGRES_PREPARED_STATEMENTS", - "maybe", - )])); - } - #[test] fn time_freeze_enabled_by_default_for_hosted_runtime() { assert!(resolve_time_freeze_from_pairs(&[])); diff --git a/fluxer_app_proxy/src/invite_meta.rs b/fluxer_app_proxy/src/invite_meta.rs deleted file mode 100644 index 74e5e36a3..000000000 --- a/fluxer_app_proxy/src/invite_meta.rs +++ /dev/null @@ -1,922 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -use crate::config::AppProxyConfig; -#[cfg(feature = "scylla")] -use anyhow::Context; -use fluxer_svc::config::DatabaseBackend; -use fluxer_svc::{postgres, postgres::KeyPart}; -use moka::future::Cache; -#[cfg(feature = "scylla")] -use scylla::DeserializeRow; -#[cfg(feature = "scylla")] -use scylla::client::session::Session; -#[cfg(feature = "scylla")] -use scylla::statement::prepared::PreparedStatement; -use serde::Deserialize; -use std::collections::HashSet; -use std::sync::{Arc, OnceLock}; -use std::time::Duration; -use tokio::task::JoinHandle; - -const INVITE_TYPE_GUILD: i32 = 0; -const INVITE_TYPE_GROUP_DM: i32 = 1; -const CHANNEL_TYPE_GROUP_DM: i32 = 3; -const MEDIA_SIZE_DEFAULT: i32 = 160; -const DEFAULT_AVATAR_COUNT: i64 = 6; -const CONNECT_RETRY_BASE: Duration = Duration::from_secs(5); -const CONNECT_RETRY_MAX: Duration = Duration::from_secs(60); - -#[derive(Clone, Debug, Eq, PartialEq)] -pub struct InvitePageMeta { - pub title: String, - pub description: String, - pub image_url: Option, -} - -#[derive(Clone, Debug, Default)] -pub struct InviteMetaEndpoints { - pub media_endpoint: Option, - pub static_cdn_endpoint: Option, -} - -pub struct InviteMetaResolver { - storage: InviteMetaStorage, - cache: Cache>, -} - -enum InviteMetaStorage { - Postgres(PostgresInviteMetaStorage), - #[cfg(feature = "scylla")] - Scylla(Box), -} - -struct PostgresInviteMetaStorage { - kv: postgres::KvClient, -} - -#[cfg(feature = "scylla")] -struct ScyllaInviteMetaStorage { - db: Arc, - stmt_invite: PreparedStatement, - stmt_guild: PreparedStatement, - stmt_channel: PreparedStatement, - stmt_user: PreparedStatement, -} - -#[cfg_attr(feature = "scylla", derive(DeserializeRow))] -#[derive(Debug, Deserialize)] -struct InviteDbRow { - r#type: i32, - guild_id: Option, - channel_id: Option, - inviter_id: Option, -} - -#[cfg_attr(feature = "scylla", derive(DeserializeRow))] -#[derive(Debug, Deserialize)] -struct GuildDbRow { - guild_id: i64, - name: String, - icon_hash: Option, - member_count: Option, -} - -#[cfg_attr(feature = "scylla", derive(DeserializeRow))] -#[derive(Debug, Deserialize)] -struct ChannelDbRow { - channel_id: i64, - r#type: i32, - name: Option, - icon_hash: Option, - recipient_ids: Option>, -} - -#[cfg_attr(feature = "scylla", derive(DeserializeRow))] -#[derive(Debug, Deserialize)] -struct UserDbRow { - user_id: i64, - username: String, - discriminator: i32, - global_name: Option, - avatar_hash: Option, -} - -impl InviteMetaResolver { - pub async fn connect(config: &AppProxyConfig) -> anyhow::Result { - match config.database_backend { - DatabaseBackend::Postgres => { - let postgres_config = fluxer_svc::postgres::PostgresConfig { - url: config.postgres_url.clone(), - host: config.postgres_host.clone(), - port: config.postgres_port, - database: config.postgres_database.clone(), - username: config.postgres_username.clone(), - password: config.postgres_password.clone(), - ssl: config.postgres_ssl, - ssl_ca: config.postgres_ssl_ca.clone(), - max_connections: config.postgres_max_connections, - kv_table: config.postgres_kv_table.clone(), - prepared_statements: config.postgres_prepared_statements, - }; - let pool = fluxer_svc::postgres::connect(&postgres_config).await?; - let kv = postgres::KvClient::new(pool, &postgres_config)?; - Self::new_postgres(kv, config) - } - DatabaseBackend::Cassandra => { - #[cfg(feature = "scylla")] - { - let db = fluxer_svc::scylla::connect(&fluxer_svc::scylla::ScyllaConfig { - hosts: config.scylla_hosts.clone(), - keyspace: config.scylla_keyspace.clone(), - username: config.scylla_username.clone(), - password: config.scylla_password.clone(), - }) - .await?; - - Self::new_scylla(db, config).await - } - #[cfg(not(feature = "scylla"))] - { - anyhow::bail!("FLUXER_DATABASE_BACKEND=cassandra requires the scylla feature"); - } - } - } - } - - fn new_postgres(kv: postgres::KvClient, config: &AppProxyConfig) -> anyhow::Result { - let cache = invite_meta_cache(config); - - Ok(Self { - storage: InviteMetaStorage::Postgres(PostgresInviteMetaStorage { kv }), - cache, - }) - } - - #[cfg(feature = "scylla")] - async fn new_scylla(db: Arc, config: &AppProxyConfig) -> anyhow::Result { - let stmt_invite = db - .prepare( - "SELECT type, guild_id, channel_id, inviter_id \ - FROM invites WHERE code = ? LIMIT 1", - ) - .await - .context("failed to prepare invite lookup")?; - let stmt_guild = db - .prepare( - "SELECT guild_id, name, icon_hash, member_count \ - FROM guilds WHERE guild_id = ? LIMIT 1", - ) - .await - .context("failed to prepare guild lookup")?; - let stmt_channel = db - .prepare( - "SELECT channel_id, type, name, icon_hash, recipient_ids \ - FROM channels WHERE channel_id = ? AND soft_deleted = false LIMIT 1", - ) - .await - .context("failed to prepare channel lookup")?; - let stmt_user = db - .prepare( - "SELECT user_id, username, discriminator, global_name, avatar_hash \ - FROM users WHERE user_id = ? LIMIT 1", - ) - .await - .context("failed to prepare user lookup")?; - - let cache = invite_meta_cache(config); - - Ok(Self { - storage: InviteMetaStorage::Scylla(Box::new(ScyllaInviteMetaStorage { - db, - stmt_invite, - stmt_guild, - stmt_channel, - stmt_user, - })), - cache, - }) - } - - pub async fn resolve( - &self, - code: &str, - endpoints: &InviteMetaEndpoints, - ) -> anyhow::Result> { - let cache_key = format!( - "{}\n{}\n{}", - code, - endpoints.media_endpoint.as_deref().unwrap_or_default(), - endpoints.static_cdn_endpoint.as_deref().unwrap_or_default() - ); - if let Some(cached) = self.cache.get(&cache_key).await { - return Ok(cached); - } - - let meta = self.resolve_uncached(code, endpoints).await?; - self.cache.insert(cache_key, meta.clone()).await; - Ok(meta) - } - - async fn resolve_uncached( - &self, - code: &str, - endpoints: &InviteMetaEndpoints, - ) -> anyhow::Result> { - let Some(invite) = self.fetch_invite(code).await? else { - return Ok(None); - }; - - match invite.r#type { - INVITE_TYPE_GUILD => self.build_guild_meta(&invite, endpoints).await, - INVITE_TYPE_GROUP_DM => self.build_group_dm_meta(&invite, endpoints).await, - _ => Ok(None), - } - } - - async fn build_guild_meta( - &self, - invite: &InviteDbRow, - endpoints: &InviteMetaEndpoints, - ) -> anyhow::Result> { - let Some(guild_id) = invite.guild_id else { - return Ok(None); - }; - let Some(guild) = self.fetch_guild(guild_id).await? else { - return Ok(None); - }; - let channel = match invite.channel_id { - Some(channel_id) => self.fetch_channel(channel_id).await?, - None => None, - }; - let inviter = match invite.inviter_id { - Some(inviter_id) => self.fetch_user(inviter_id).await?, - None => None, - }; - - Ok(Some(build_guild_meta( - &guild, - channel.as_ref(), - inviter.as_ref(), - endpoints, - ))) - } - - async fn build_group_dm_meta( - &self, - invite: &InviteDbRow, - endpoints: &InviteMetaEndpoints, - ) -> anyhow::Result> { - let Some(channel_id) = invite.channel_id else { - return Ok(None); - }; - let Some(channel) = self.fetch_channel(channel_id).await? else { - return Ok(None); - }; - if channel.r#type != CHANNEL_TYPE_GROUP_DM { - return Ok(None); - } - let inviter = match invite.inviter_id { - Some(inviter_id) => self.fetch_user(inviter_id).await?, - None => None, - }; - - Ok(Some(build_group_dm_meta( - &channel, - inviter.as_ref(), - endpoints, - ))) - } - - async fn fetch_invite(&self, code: &str) -> anyhow::Result> { - self.storage.fetch_invite(code).await - } - - async fn fetch_guild(&self, guild_id: i64) -> anyhow::Result> { - self.storage.fetch_guild(guild_id).await - } - - async fn fetch_channel(&self, channel_id: i64) -> anyhow::Result> { - self.storage.fetch_channel(channel_id).await - } - - async fn fetch_user(&self, user_id: i64) -> anyhow::Result> { - self.storage.fetch_user(user_id).await - } -} - -pub fn start_background_connect( - slot: &Arc>, - config: Arc, -) -> JoinHandle<()> { - let slot = Arc::clone(slot); - tokio::spawn(async move { - let mut failures: u32 = 0; - loop { - match InviteMetaResolver::connect(&config).await { - Ok(resolver) => { - let _ = slot.set(resolver); - tracing::info!("invite metadata resolver connected"); - return; - } - Err(err) => { - failures = failures.saturating_add(1); - let backoff = connect_backoff(failures); - tracing::warn!( - %err, - failures, - backoff_ms = backoff.as_millis() as u64, - "invite metadata resolver connect failed; retrying" - ); - tokio::time::sleep(backoff).await; - } - } - } - }) -} - -fn connect_backoff(failures: u32) -> Duration { - CONNECT_RETRY_BASE - .saturating_mul(1u32 << failures.min(5)) - .min(CONNECT_RETRY_MAX) -} - -fn invite_meta_cache(config: &AppProxyConfig) -> Cache> { - Cache::builder() - .max_capacity(config.invite_meta_cache_max_entries) - .time_to_live(Duration::from_millis(config.invite_meta_cache_ttl_ms)) - .build() -} - -impl InviteMetaStorage { - async fn fetch_invite(&self, code: &str) -> anyhow::Result> { - match self { - InviteMetaStorage::Postgres(storage) => storage.fetch_invite(code).await, - #[cfg(feature = "scylla")] - InviteMetaStorage::Scylla(storage) => storage.fetch_invite(code).await, - } - } - - async fn fetch_guild(&self, guild_id: i64) -> anyhow::Result> { - match self { - InviteMetaStorage::Postgres(storage) => storage.fetch_guild(guild_id).await, - #[cfg(feature = "scylla")] - InviteMetaStorage::Scylla(storage) => storage.fetch_guild(guild_id).await, - } - } - - async fn fetch_channel(&self, channel_id: i64) -> anyhow::Result> { - match self { - InviteMetaStorage::Postgres(storage) => storage.fetch_channel(channel_id).await, - #[cfg(feature = "scylla")] - InviteMetaStorage::Scylla(storage) => storage.fetch_channel(channel_id).await, - } - } - - async fn fetch_user(&self, user_id: i64) -> anyhow::Result> { - match self { - InviteMetaStorage::Postgres(storage) => storage.fetch_user(user_id).await, - #[cfg(feature = "scylla")] - InviteMetaStorage::Scylla(storage) => storage.fetch_user(user_id).await, - } - } -} - -impl PostgresInviteMetaStorage { - async fn fetch_invite(&self, code: &str) -> anyhow::Result> { - self.fetch_row("invites", &[KeyPart::String(code)]).await - } - - async fn fetch_guild(&self, guild_id: i64) -> anyhow::Result> { - self.fetch_row("guilds", &[KeyPart::BigInt(guild_id)]).await - } - - async fn fetch_channel(&self, channel_id: i64) -> anyhow::Result> { - self.fetch_row( - "channels", - &[KeyPart::BigInt(channel_id), KeyPart::Bool(false)], - ) - .await - } - - async fn fetch_user(&self, user_id: i64) -> anyhow::Result> { - self.fetch_row("users", &[KeyPart::BigInt(user_id)]).await - } - - async fn fetch_row(&self, table_name: &str, key: &[KeyPart<'_>]) -> anyhow::Result> - where - T: serde::de::DeserializeOwned, - { - let key = postgres::kv_key(key)?; - let Some(row) = self.kv.get_row(table_name, &key).await? else { - return Ok(None); - }; - let row = postgres::decode_row_dates_as_millis(row)?; - Ok(Some(serde_json::from_value(row)?)) - } -} - -#[cfg(feature = "scylla")] -impl ScyllaInviteMetaStorage { - async fn fetch_invite(&self, code: &str) -> anyhow::Result> { - let result = self.db.execute_unpaged(&self.stmt_invite, (code,)).await?; - let rows = result.into_rows_result()?; - Ok(rows.maybe_first_row::()?) - } - - async fn fetch_guild(&self, guild_id: i64) -> anyhow::Result> { - let result = self - .db - .execute_unpaged(&self.stmt_guild, (guild_id,)) - .await?; - let rows = result.into_rows_result()?; - Ok(rows.maybe_first_row::()?) - } - - async fn fetch_channel(&self, channel_id: i64) -> anyhow::Result> { - let result = self - .db - .execute_unpaged(&self.stmt_channel, (channel_id,)) - .await?; - let rows = result.into_rows_result()?; - Ok(rows.maybe_first_row::()?) - } - - async fn fetch_user(&self, user_id: i64) -> anyhow::Result> { - let result = self.db.execute_unpaged(&self.stmt_user, (user_id,)).await?; - let rows = result.into_rows_result()?; - Ok(rows.maybe_first_row::()?) - } -} - -pub fn invite_code_from_path(path: &str) -> Option<&str> { - let mut segments = path.trim_start_matches('/').split('/'); - if segments.next()? != "invite" { - return None; - } - let code = segments.next()?; - if !is_valid_invite_code_segment(code) { - return None; - } - match segments.next() { - None => Some(code), - Some("login") if segments.next().is_none() => Some(code), - _ => None, - } -} - -pub fn inject_invite_meta(html: &str, meta: &InvitePageMeta) -> String { - let mut html = replace_title(html, &meta.title); - html = remove_meta_tags( - &html, - &[ - MetaSelector::Name("description"), - MetaSelector::Property("og:title"), - MetaSelector::Property("og:description"), - MetaSelector::Property("og:image"), - MetaSelector::Property("og:type"), - MetaSelector::Name("twitter:card"), - MetaSelector::Name("twitter:title"), - MetaSelector::Name("twitter:description"), - MetaSelector::Name("twitter:image"), - ], - ); - - let tags = build_meta_tags(meta); - insert_head_tags(&html, &tags) -} - -fn build_guild_meta( - guild: &GuildDbRow, - channel: Option<&ChannelDbRow>, - inviter: Option<&UserDbRow>, - endpoints: &InviteMetaEndpoints, -) -> InvitePageMeta { - let title = format!("Join {} on Fluxer", guild.name); - let mut description = format!("You've been invited to join {}.", guild.name); - if let Some(channel_name) = - channel.and_then(|channel| clean_optional_text(channel.name.as_deref())) - { - description.push_str(&format!(" Channel: {channel_name}.")); - } - if let Some(member_count) = guild.member_count { - description.push_str(&format!(" {}", format_member_count(member_count as i64))); - } - - let image_url = guild - .icon_hash - .as_deref() - .and_then(|hash| { - media_image_url( - endpoints.media_endpoint.as_deref(), - "icons", - guild.guild_id, - hash, - ) - }) - .or_else(|| inviter.and_then(|user| user_avatar_url(user, endpoints))); - - InvitePageMeta { - title, - description, - image_url, - } -} - -fn build_group_dm_meta( - channel: &ChannelDbRow, - inviter: Option<&UserDbRow>, - endpoints: &InviteMetaEndpoints, -) -> InvitePageMeta { - let channel_name = clean_optional_text(channel.name.as_deref()); - let inviter_name = inviter.map(display_user_name); - let title = if let Some(channel_name) = channel_name { - format!("Join {channel_name} on Fluxer") - } else if let Some(inviter_name) = inviter_name.as_deref() { - format!("Join {inviter_name}'s group DM on Fluxer") - } else { - "Join a group DM on Fluxer".to_owned() - }; - - let member_count = channel - .recipient_ids - .as_ref() - .map(|recipients| recipients.len() as i64); - let mut description = if let Some(inviter_name) = inviter_name { - format!("{inviter_name} invited you to a group DM.") - } else { - "You've been invited to a group DM.".to_owned() - }; - if let Some(member_count) = member_count { - description.push_str(&format!(" {}", format_member_count(member_count))); - } - - let image_url = channel - .icon_hash - .as_deref() - .and_then(|hash| { - media_image_url( - endpoints.media_endpoint.as_deref(), - "icons", - channel.channel_id, - hash, - ) - }) - .or_else(|| inviter.and_then(|user| user_avatar_url(user, endpoints))); - - InvitePageMeta { - title, - description, - image_url, - } -} - -fn user_avatar_url(user: &UserDbRow, endpoints: &InviteMetaEndpoints) -> Option { - user.avatar_hash - .as_deref() - .and_then(|hash| { - media_image_url( - endpoints.media_endpoint.as_deref(), - "avatars", - user.user_id, - hash, - ) - }) - .or_else(|| default_avatar_url(endpoints.static_cdn_endpoint.as_deref(), user.user_id)) -} - -fn media_image_url(endpoint: Option<&str>, path: &str, id: i64, hash: &str) -> Option { - let endpoint = clean_endpoint(endpoint?)?; - let (hash, animated) = parse_media_hash(hash); - if hash.is_empty() { - return None; - } - let mut url = format!("{endpoint}/{path}/{id}/{hash}.webp?size={MEDIA_SIZE_DEFAULT}"); - if animated { - url.push_str("&animated=false"); - } - Some(url) -} - -fn default_avatar_url(endpoint: Option<&str>, user_id: i64) -> Option { - let endpoint = clean_endpoint(endpoint?)?; - let index = user_id.rem_euclid(DEFAULT_AVATAR_COUNT); - Some(format!("{endpoint}/avatars/{index}.png")) -} - -fn parse_media_hash(value: &str) -> (&str, bool) { - value - .strip_prefix("a_") - .map(|hash| (hash, true)) - .unwrap_or((value, false)) -} - -fn clean_endpoint(value: &str) -> Option { - let value = value.trim().trim_end_matches('/'); - if value.is_empty() { - None - } else { - Some(value.to_owned()) - } -} - -fn display_user_name(user: &UserDbRow) -> String { - clean_optional_text(user.global_name.as_deref()) - .map(ToOwned::to_owned) - .unwrap_or_else(|| { - if user.discriminator > 0 { - format!("{}#{:04}", user.username, user.discriminator) - } else { - user.username.clone() - } - }) -} - -fn clean_optional_text(value: Option<&str>) -> Option<&str> { - value.map(str::trim).filter(|value| !value.is_empty()) -} - -fn format_member_count(value: i64) -> String { - let noun = if value == 1 { "member" } else { "members" }; - format!("{} {noun}", format_count(value)) -} - -fn format_count(value: i64) -> String { - let negative = value < 0; - let digits = value.abs().to_string(); - let mut out = String::with_capacity(digits.len() + digits.len() / 3 + usize::from(negative)); - if negative { - out.push('-'); - } - for (index, ch) in digits.chars().enumerate() { - if index > 0 && (digits.len() - index).is_multiple_of(3) { - out.push(','); - } - out.push(ch); - } - out -} - -fn is_valid_invite_code_segment(code: &str) -> bool { - !code.is_empty() - && code.len() <= 128 - && code - .bytes() - .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_')) -} - -fn build_meta_tags(meta: &InvitePageMeta) -> String { - let title = escape_html_attr(&meta.title); - let description = escape_html_attr(&meta.description); - let mut tags = format!( - r#" - - - - - -"# - ); - if let Some(image_url) = &meta.image_url { - let image = escape_html_attr(image_url); - tags.push_str(&format!( - r#" - -"# - )); - } - tags -} - -fn replace_title(html: &str, title: &str) -> String { - let escaped_title = escape_html_text(title); - let lower = html.to_ascii_lowercase(); - let Some(start) = lower.find("") else { - return html.to_owned(); - }; - let content_start = start + "<title>".len(); - let Some(relative_end) = lower[content_start..].find("") else { - return html.to_owned(); - }; - let content_end = content_start + relative_end; - - let mut result = String::with_capacity(html.len() + escaped_title.len()); - result.push_str(&html[..content_start]); - result.push_str(&escaped_title); - result.push_str(&html[content_end..]); - result -} - -#[derive(Clone, Copy)] -enum MetaSelector { - Name(&'static str), - Property(&'static str), -} - -fn remove_meta_tags(html: &str, selectors: &[MetaSelector]) -> String { - let lower = html.to_ascii_lowercase(); - let mut result = String::with_capacity(html.len()); - let mut cursor = 0; - - while let Some(relative_start) = lower[cursor..].find("') else { - break; - }; - let end = start + relative_end + 1; - let tag_lower = &lower[start..end]; - if selectors - .iter() - .any(|selector| meta_tag_matches(tag_lower, *selector)) - { - result.push_str(&html[cursor..start]); - } else { - result.push_str(&html[cursor..end]); - } - cursor = end; - } - - result.push_str(&html[cursor..]); - result -} - -fn meta_tag_matches(tag_lower: &str, selector: MetaSelector) -> bool { - match selector { - MetaSelector::Name(value) => attr_matches(tag_lower, "name", value), - MetaSelector::Property(value) => attr_matches(tag_lower, "property", value), - } -} - -fn attr_matches(tag_lower: &str, attr: &str, value: &str) -> bool { - tag_lower.contains(&format!(r#"{attr}="{value}""#)) - || tag_lower.contains(&format!("{attr}='{value}'")) -} - -fn insert_head_tags(html: &str, tags: &str) -> String { - let lower = html.to_ascii_lowercase(); - if let Some(title_end) = lower.find("") { - let insert_at = title_end + "".len(); - return insert_at_pos(html, insert_at, tags); - } - if let Some(head_start) = lower - .find("") - .map(|pos| pos + "".len()) - .or_else(|| { - lower - .find("').map(|end| pos + end + 1)) - }) - { - return insert_at_pos(html, head_start, tags); - } - html.to_owned() -} - -fn insert_at_pos(html: &str, insert_at: usize, tags: &str) -> String { - let mut result = String::with_capacity(html.len() + tags.len() + 2); - result.push_str(&html[..insert_at]); - result.push('\n'); - result.push_str(tags); - result.push_str(&html[insert_at..]); - result -} - -fn escape_html_text(value: &str) -> String { - value - .replace('&', "&") - .replace('<', "<") - .replace('>', ">") -} - -fn escape_html_attr(value: &str) -> String { - escape_html_text(value) - .replace('"', """) - .replace('\'', "'") -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn the_connect_backoff_grows_and_stops_at_the_ceiling() { - assert_eq!(connect_backoff(1), Duration::from_secs(10)); - assert_eq!(connect_backoff(2), Duration::from_secs(20)); - assert_eq!(connect_backoff(3), Duration::from_secs(40)); - assert_eq!(connect_backoff(4), CONNECT_RETRY_MAX); - assert_eq!(connect_backoff(64), CONNECT_RETRY_MAX); - } - - #[tokio::test] - async fn a_refused_database_never_fills_the_slot_and_never_gives_up() { - let mut config = AppProxyConfig::from_env(); - config.database_backend = DatabaseBackend::Postgres; - config.postgres_url = None; - config.postgres_host = "127.0.0.1".to_owned(); - config.postgres_port = 1; - config.postgres_ssl = false; - let slot = Arc::new(OnceLock::new()); - - let handle = start_background_connect(&slot, Arc::new(config)); - tokio::time::sleep(Duration::from_millis(200)).await; - - assert!(slot.get().is_none()); - assert!( - !handle.is_finished(), - "a failed connect ended the retry loop and disabled invite metadata for the process" - ); - handle.abort(); - } - - fn endpoints() -> InviteMetaEndpoints { - InviteMetaEndpoints { - media_endpoint: Some("https://media.example.test/media/".to_owned()), - static_cdn_endpoint: Some("https://static.example.test/".to_owned()), - } - } - - #[test] - fn extracts_invite_code_from_register_and_login_paths() { - assert_eq!(invite_code_from_path("/invite/abc123"), Some("abc123")); - assert_eq!( - invite_code_from_path("/invite/abc123/login"), - Some("abc123") - ); - assert_eq!(invite_code_from_path("/channels/@me"), None); - assert_eq!(invite_code_from_path("/invite/bad.code"), None); - assert_eq!(invite_code_from_path("/invite/abc123/settings"), None); - } - - #[test] - fn builds_guild_meta_with_icon_and_member_count() { - let guild = GuildDbRow { - guild_id: 42, - name: "Rust Friends".to_owned(), - icon_hash: Some("a_iconhash".to_owned()), - member_count: Some(12345), - }; - let channel = ChannelDbRow { - channel_id: 7, - r#type: 0, - name: Some("general".to_owned()), - icon_hash: None, - recipient_ids: None, - }; - - let meta = build_guild_meta(&guild, Some(&channel), None, &endpoints()); - - assert_eq!(meta.title, "Join Rust Friends on Fluxer"); - assert!(meta.description.contains("Channel: general.")); - assert!(meta.description.contains("12,345 members")); - assert_eq!( - meta.image_url.as_deref(), - Some("https://media.example.test/media/icons/42/iconhash.webp?size=160&animated=false") - ); - } - - #[test] - fn builds_group_dm_meta_with_inviter_avatar_fallback() { - let channel = ChannelDbRow { - channel_id: 55, - r#type: CHANNEL_TYPE_GROUP_DM, - name: None, - icon_hash: None, - recipient_ids: Some(HashSet::from([1, 2, 3])), - }; - let inviter = UserDbRow { - user_id: 99, - username: "ada".to_owned(), - discriminator: 7, - global_name: Some("Ada".to_owned()), - avatar_hash: Some("avatarhash".to_owned()), - }; - - let meta = build_group_dm_meta(&channel, Some(&inviter), &endpoints()); - - assert_eq!(meta.title, "Join Ada's group DM on Fluxer"); - assert!(meta.description.contains("3 members")); - assert_eq!( - meta.image_url.as_deref(), - Some("https://media.example.test/media/avatars/99/avatarhash.webp?size=160") - ); - } - - #[test] - fn injects_invite_meta_and_removes_default_description() { - let html = r#"Fluxer"#; - let meta = InvitePageMeta { - title: "Join A & B".to_owned(), - description: "A < B \"test\"".to_owned(), - image_url: Some("https://cdn.example.test/icon.png".to_owned()), - }; - - let result = inject_invite_meta(html, &meta); - - assert!(result.contains("Join A & B")); - assert!(!result.contains("content=\"Default\"")); - assert!(!result.contains("content=\"Old\"")); - assert!(result.contains(r#"name="description" content="A < B "test"""#)); - assert!( - result.contains(r#"property="og:image" content="https://cdn.example.test/icon.png""#) - ); - } -} diff --git a/fluxer_app_proxy/src/lib.rs b/fluxer_app_proxy/src/lib.rs index e128f9098..f20f11ac2 100644 --- a/fluxer_app_proxy/src/lib.rs +++ b/fluxer_app_proxy/src/lib.rs @@ -7,7 +7,6 @@ pub mod discovery_cache; #[cfg(feature = "time-freeze")] pub mod frozen_snapshots; pub mod geoip; -pub mod invite_meta; pub mod routes; pub mod state; pub mod time_freeze; diff --git a/fluxer_app_proxy/src/main.rs b/fluxer_app_proxy/src/main.rs index a743314d1..adac7c7db 100644 --- a/fluxer_app_proxy/src/main.rs +++ b/fluxer_app_proxy/src/main.rs @@ -5,13 +5,13 @@ use fluxer_app_proxy::{ config::AppProxyConfig, csp::CompiledCspPolicy, discovery_cache::DiscoveryCache, - geoip, invite_meta, + geoip, routes::build_router, state::{ AppProxyBudgets, AppState, MAX_SPA_INDEX_BYTES, build_http_client, read_bounded_text_file, }, }; -use std::sync::{Arc, OnceLock}; +use std::sync::Arc; use tokio::{net::TcpListener, runtime::Builder}; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt}; @@ -56,11 +56,6 @@ fn main() -> anyhow::Result<()> { config.discovery_refresh_interval_ms, ); - let invite_meta = Arc::new(OnceLock::new()); - let invite_meta_connect = config - .invite_meta_enabled - .then(|| invite_meta::start_background_connect(&invite_meta, Arc::clone(&config))); - let index_html = if config.index_upstream_url.is_none() { let index_path = std::path::Path::new(&config.static_dir).join("index.html"); match read_bounded_text_file(&index_path, MAX_SPA_INDEX_BYTES).await { @@ -80,7 +75,6 @@ fn main() -> anyhow::Result<()> { http_client, discovery_cache, geoip, - invite_meta, index_html, budgets: AppProxyBudgets::default(), }; @@ -97,9 +91,6 @@ fn main() -> anyhow::Result<()> { .context("app proxy server exited unexpectedly")?; cancel.abort(); - if let Some(handle) = invite_meta_connect { - handle.abort(); - } Ok(()) }) } diff --git a/fluxer_app_proxy/src/routes/assets_proxy.rs b/fluxer_app_proxy/src/routes/assets_proxy.rs index 664ed6c43..38c583d14 100644 --- a/fluxer_app_proxy/src/routes/assets_proxy.rs +++ b/fluxer_app_proxy/src/routes/assets_proxy.rs @@ -436,8 +436,8 @@ mod tests { use axum::http::header::HeaderName; use fluxer_common::config::GeoipSourceConfig; use fluxer_common::geoip::{GeoipConfig, GeoipResolver}; + use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; - use std::sync::{Arc, OnceLock}; use tower::ServiceExt; async fn spawn_upstream(status: StatusCode, cache_control: &'static str) -> String { @@ -491,7 +491,6 @@ mod tests { trust_client_ip_header: false, client_ip_header_name: "x-forwarded-for".to_owned(), })), - invite_meta: Arc::new(OnceLock::new()), index_html: None, budgets: crate::state::AppProxyBudgets::default(), } diff --git a/fluxer_app_proxy/src/routes/health.rs b/fluxer_app_proxy/src/routes/health.rs index 12667509a..34b580f6a 100644 --- a/fluxer_app_proxy/src/routes/health.rs +++ b/fluxer_app_proxy/src/routes/health.rs @@ -12,35 +12,14 @@ pub async fn health() -> &'static str { } pub async fn ready(State(state): State) -> Response { - readiness_report( - state.discovery_cache.has_snapshot().await, - state.config.invite_meta_enabled, - state.invite_meta.get().is_some(), - ) - .into_response() + readiness_report(state.discovery_cache.has_snapshot().await).into_response() } -fn readiness_report( - discovery_cached: bool, - invite_meta_configured: bool, - invite_meta_connected: bool, -) -> (StatusCode, String) { - let mut degraded = Vec::new(); - if invite_meta_configured && !invite_meta_connected { - degraded.push("invite_meta"); - } - let suffix = if degraded.is_empty() { - String::new() - } else { - format!(" (degraded: {})", degraded.join(", ")) - }; +fn readiness_report(discovery_cached: bool) -> (StatusCode, &'static str) { if discovery_cached { - (StatusCode::OK, format!("OK{suffix}")) + (StatusCode::OK, "OK") } else { - ( - StatusCode::SERVICE_UNAVAILABLE, - format!("NOT READY: discovery{suffix}"), - ) + (StatusCode::SERVICE_UNAVAILABLE, "NOT READY: discovery") } } @@ -55,7 +34,7 @@ mod tests { use axum::http::{HeaderValue, Request as HttpRequest, header}; use fluxer_common::config::GeoipSourceConfig; use fluxer_common::geoip::{GeoipConfig, GeoipResolver}; - use std::sync::{Arc, OnceLock}; + use std::sync::Arc; use tower::ServiceExt; const DISCOVERY_BODY: &str = r#"{"api_code_version":"proxy-test"}"#; @@ -77,9 +56,8 @@ mod tests { format!("http://{addr}/") } - async fn probe_state(invite_meta_enabled: bool) -> AppState { + async fn probe_state() -> AppState { let mut config = AppProxyConfig::from_env(); - config.invite_meta_enabled = invite_meta_enabled; config.discovery_upstream_url = spawn_discovery_origin().await; let csp = Arc::new( crate::csp::CompiledCspPolicy::from_config(&config) @@ -98,7 +76,6 @@ mod tests { trust_client_ip_header: false, client_ip_header_name: "x-forwarded-for".to_owned(), })), - invite_meta: Arc::new(OnceLock::new()), index_html: None, budgets: crate::state::AppProxyBudgets::default(), } @@ -131,7 +108,7 @@ mod tests { #[tokio::test] async fn liveness_stays_constant_while_the_proxy_cannot_serve() { - let state = probe_state(false).await; + let state = probe_state().await; assert_eq!( probe(state, "/_health").await, (StatusCode::OK, "OK".to_owned()) @@ -140,7 +117,7 @@ mod tests { #[tokio::test] async fn readiness_fails_while_the_discovery_cache_is_empty() { - let state = probe_state(false).await; + let state = probe_state().await; let (status, body) = probe(state, "/_ready").await; assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); assert!(body.contains("discovery"), "{body}"); @@ -148,7 +125,7 @@ mod tests { #[tokio::test] async fn readiness_passes_once_a_discovery_snapshot_is_cached() { - let state = probe_state(false).await; + let state = probe_state().await; warm_discovery(&state).await; assert_eq!( probe(state, "/_ready").await, @@ -156,48 +133,12 @@ mod tests { ); } - #[tokio::test] - async fn readiness_survives_a_configured_invite_resolver_that_never_connects() { - let state = probe_state(true).await; - warm_discovery(&state).await; - let (status, body) = probe(state, "/_ready").await; - assert_eq!(status, StatusCode::OK); - assert!(body.contains("invite_meta"), "{body}"); - } - - #[tokio::test] - async fn an_empty_discovery_cache_fails_readiness_even_while_invite_metadata_is_degraded() { - let state = probe_state(true).await; - let (status, body) = probe(state, "/_ready").await; - assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); - assert!(body.contains("discovery"), "{body}"); - assert!(body.contains("invite_meta"), "{body}"); - } - #[test] - fn invite_metadata_is_reported_as_degraded_instead_of_gating_readiness() { - assert_eq!(readiness_report(true, false, false).0, StatusCode::OK); + fn readiness_tracks_only_the_discovery_snapshot() { + assert_eq!(readiness_report(true), (StatusCode::OK, "OK")); assert_eq!( - readiness_report(true, true, true), - (StatusCode::OK, "OK".to_owned()) - ); - assert_eq!( - readiness_report(true, true, false), - (StatusCode::OK, "OK (degraded: invite_meta)".to_owned()) - ); - assert_eq!( - readiness_report(false, false, false), - ( - StatusCode::SERVICE_UNAVAILABLE, - "NOT READY: discovery".to_owned() - ) - ); - assert_eq!( - readiness_report(false, true, false), - ( - StatusCode::SERVICE_UNAVAILABLE, - "NOT READY: discovery (degraded: invite_meta)".to_owned() - ) + readiness_report(false), + (StatusCode::SERVICE_UNAVAILABLE, "NOT READY: discovery") ); } } diff --git a/fluxer_app_proxy/src/routes/spa_index.rs b/fluxer_app_proxy/src/routes/spa_index.rs index e0fe896e3..0339e7250 100644 --- a/fluxer_app_proxy/src/routes/spa_index.rs +++ b/fluxer_app_proxy/src/routes/spa_index.rs @@ -5,9 +5,6 @@ use crate::config::HttpEndpoint; use crate::csp::{RuntimeCspSources, generate_nonce}; use crate::discovery_cache::{DiscoveryResponse, discovery_endpoint}; use crate::geoip::build_geoip_response; -use crate::invite_meta::{ - InviteMetaEndpoints, InvitePageMeta, inject_invite_meta, invite_code_from_path, -}; use crate::state::{ AppProxyBudgets, AppState, MAX_RENDERED_SPA_INDEX_BYTES, MAX_SPA_INDEX_BYTES, read_bounded_text_file, @@ -60,7 +57,7 @@ pub async fn spa_catch_all( .await; } - serve_spa_index(&state, &headers, request_path).await + serve_spa_index(&state, &headers).await } const CRAWL_CONTROL_CACHE_CONTROL: &str = "public, max-age=300, must-revalidate"; @@ -144,7 +141,7 @@ async fn serve_static_file( response } -async fn serve_spa_index(state: &AppState, headers: &HeaderMap, request_path: &str) -> Response { +async fn serve_spa_index(state: &AppState, headers: &HeaderMap) -> Response { let time_freeze = load_time_freeze_config_for_request(&state.config, headers); let debug_header = time_freeze_debug_header(&time_freeze); let should_bust_dev_assets = state.config.index_upstream_url.is_some(); @@ -159,7 +156,6 @@ async fn serve_spa_index(state: &AppState, headers: &HeaderMap, request_path: &s let nonce = generate_nonce(); let runtime_csp_sources = build_runtime_csp_sources(state, &discovery); - let invite_meta = resolve_invite_meta(state, request_path, &runtime_csp_sources).await; let static_cdn_endpoint = runtime_csp_sources .static_cdn_endpoint .as_ref() @@ -188,7 +184,6 @@ async fn serve_spa_index(state: &AppState, headers: &HeaderMap, request_path: &s &script_tag, static_cdn_endpoint, media_endpoint, - invite_meta.as_ref(), dev_buster.as_deref(), ) { Ok(html) => html, @@ -233,7 +228,6 @@ fn render_spa_document( script_tag: &str, static_cdn_endpoint: &str, media_endpoint: &str, - invite_meta: Option<&InvitePageMeta>, dev_asset_cache_buster: Option<&str>, ) -> Result { let mut document = bounded_document(inject_bootstrap( @@ -243,9 +237,6 @@ fn render_spa_document( static_cdn_endpoint, media_endpoint, ))?; - if let Some(meta) = invite_meta { - document = bounded_document(inject_invite_meta(&document, meta))?; - } if let Some(buster) = dev_asset_cache_buster { document = bounded_document(append_dev_asset_cache_buster(&document, buster))?; } @@ -259,33 +250,6 @@ async fn refresh_discovery_for_spa(state: &AppState) -> Option Option { - let code = invite_code_from_path(request_path)?; - let resolver = state.invite_meta.get()?; - let endpoints = InviteMetaEndpoints { - media_endpoint: runtime_csp_sources - .media_endpoint - .as_ref() - .map(|endpoint| endpoint.as_str().to_owned()), - static_cdn_endpoint: runtime_csp_sources - .static_cdn_endpoint - .as_ref() - .map(|endpoint| endpoint.as_str().to_owned()), - }; - - match resolver.resolve(code, &endpoints).await { - Ok(meta) => meta, - Err(err) => { - tracing::warn!(%err, code, "failed to resolve invite metadata"); - None - } - } -} - fn build_runtime_csp_sources(state: &AppState, discovery: &DiscoveryResponse) -> RuntimeCspSources { RuntimeCspSources { static_cdn_endpoint: discovery_endpoint(discovery, "static_cdn") @@ -625,7 +589,7 @@ mod tests { use axum::body::Body; use fluxer_common::config::GeoipSourceConfig; use fluxer_common::geoip::{GeoipConfig, GeoipResolver}; - use std::sync::{Arc, OnceLock}; + use std::sync::Arc; #[test] fn dev_asset_cache_buster_rewrites_script_and_link_assets() { @@ -677,14 +641,6 @@ mod tests { const SHELL_WITH_A_NONCE_HOLE: &str = r#"Fluxer"#; - fn sample_invite_meta() -> InvitePageMeta { - InvitePageMeta { - title: "Join Sample Space".to_owned(), - description: "A sample invite".to_owned(), - image_url: None, - } - } - #[test] fn the_rendered_document_always_carries_the_bootstrap_and_a_real_nonce() { let rendered = render_spa_document( @@ -694,7 +650,6 @@ mod tests { "https://static.example.test", "", None, - None, ) .expect("test SPA document must render within its size limit"); @@ -703,36 +658,6 @@ mod tests { assert!(rendered.contains("")); } - #[test] - fn invite_metadata_reaches_the_rendered_document_only_when_resolved() { - let meta = sample_invite_meta(); - let with_meta = render_spa_document( - SHELL_WITH_A_NONCE_HOLE, - "reqnonce", - "", - "", - "", - Some(&meta), - None, - ) - .expect("test SPA document must render within its size limit"); - let without_meta = render_spa_document( - SHELL_WITH_A_NONCE_HOLE, - "reqnonce", - "", - "", - "", - None, - None, - ) - .expect("test SPA document must render within its size limit"); - - assert!(with_meta.contains("Join Sample Space")); - assert!(with_meta.contains("og:title")); - assert!(!without_meta.contains("Join Sample Space")); - assert!(!without_meta.contains("og:title")); - } - #[test] fn the_dev_cache_buster_reaches_the_rendered_document_only_when_supplied() { let busted = render_spa_document( @@ -741,7 +666,6 @@ mod tests { "", "", "", - None, Some("9911"), ) .expect("test SPA document must render within its size limit"); @@ -752,7 +676,6 @@ mod tests { "", "", None, - None, ) .expect("test SPA document must render within its size limit"); @@ -778,7 +701,6 @@ mod tests { "https://fluxerstatic.com", "", None, - None, ) .expect("test SPA document must render within its size limit"); @@ -810,7 +732,6 @@ mod tests { "https://cdn.example.test/", "https://media.example.test", None, - None, ) .expect("test SPA document must render within its size limit"); @@ -844,7 +765,6 @@ mod tests { "https://cdn.example.test", "https://media.example.test/", None, - None, ) .expect("test SPA document must render within its size limit"); assert!( @@ -861,7 +781,6 @@ mod tests { "https://cdn.example.test", "https://cdn.example.test", None, - None, ) .expect("test SPA document must render within its size limit"); assert!( @@ -967,7 +886,6 @@ mod tests { trust_client_ip_header: false, client_ip_header_name: "x-forwarded-for".to_owned(), })), - invite_meta: Arc::new(OnceLock::new()), index_html: cached_shell.map(Arc::from), budgets: crate::state::AppProxyBudgets::default(), } @@ -1002,7 +920,7 @@ mod tests { let state = spa_state_serving(ReleaseChannel::Canary, Some(SHELL_WITH_ENDPOINT_HOLES)).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; assert_eq!(response.status(), StatusCode::OK); let granted_nonce = nonce_granted_by(&response); let served = read_document(response).await; @@ -1062,7 +980,7 @@ mod tests { let state = spa_state_serving(ReleaseChannel::Stable, None).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; assert_eq!(response.status(), StatusCode::OK); let served = read_document(response).await; @@ -1099,7 +1017,7 @@ mod tests { let state = spa_state_serving(ReleaseChannel::Stable, None).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; assert_eq!(response.status(), StatusCode::OK); let granted_nonce = nonce_granted_by(&response); let served = read_document(response).await; @@ -1122,7 +1040,7 @@ mod tests { "the primary hosted path shipped without a bootstrap script the browser will run" ); - let second = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let second = serve_spa_index(&state, &HeaderMap::new()).await; assert_ne!( nonce_granted_by(&second), granted_nonce, @@ -1134,7 +1052,7 @@ mod tests { #[tokio::test] async fn every_branch_announces_which_snapshot_decision_it_took() { let frozen_state = spa_state_serving(ReleaseChannel::Stable, None).await; - let frozen = serve_spa_index(&frozen_state, &HeaderMap::new(), "/channels/@me").await; + let frozen = serve_spa_index(&frozen_state, &HeaderMap::new()).await; assert_eq!( frozen .headers() @@ -1151,7 +1069,7 @@ mod tests { let live_state = spa_state_serving(ReleaseChannel::Canary, Some(SHELL_WITH_ENDPOINT_HOLES)).await; - let live = serve_spa_index(&live_state, &HeaderMap::new(), "/channels/@me").await; + let live = serve_spa_index(&live_state, &HeaderMap::new()).await; assert_eq!( live.headers() .get("x-time-freeze") @@ -1168,7 +1086,7 @@ mod tests { async fn the_frozen_shell_is_never_served_with_the_asset_lifetime() { let state = spa_state_serving(ReleaseChannel::Stable, None).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; assert_eq!( response @@ -1247,7 +1165,7 @@ mod tests { let state = spa_state_serving(ReleaseChannel::Canary, Some(SHELL_WITH_ENDPOINT_HOLES)).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; let cache_control = response .headers() @@ -1270,7 +1188,7 @@ mod tests { let index_upstream_url = spawn_local_origin(SHELL_WITH_ENDPOINT_HOLES, "text/html").await; let state = spa_state_reading_its_shell_from(index_upstream_url).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; assert_eq!(response.status(), StatusCode::OK); assert_eq!( response @@ -1304,7 +1222,7 @@ mod tests { let state = spa_state_without_discovered_endpoints(Some("https://fallbackcdn.example.test")).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; assert_eq!(response.status(), StatusCode::OK); let served = read_document(response).await; @@ -1334,7 +1252,7 @@ mod tests { async fn an_endpoint_neither_discovered_nor_configured_warms_no_socket_at_all() { let state = spa_state_without_discovered_endpoints(None).await; - let response = serve_spa_index(&state, &HeaderMap::new(), "/channels/@me").await; + let response = serve_spa_index(&state, &HeaderMap::new()).await; assert_eq!(response.status(), StatusCode::OK); let served = read_document(response).await; diff --git a/fluxer_app_proxy/src/state.rs b/fluxer_app_proxy/src/state.rs index 21072b9ae..54a32bac6 100644 --- a/fluxer_app_proxy/src/state.rs +++ b/fluxer_app_proxy/src/state.rs @@ -3,9 +3,8 @@ use crate::config::AppProxyConfig; use crate::csp::CompiledCspPolicy; use crate::discovery_cache::DiscoveryCache; -use crate::invite_meta::InviteMetaResolver; use fluxer_common::geoip::GeoipResolver; -use std::sync::{Arc, OnceLock}; +use std::sync::Arc; use std::time::Duration; use tokio::io::AsyncReadExt; use tokio::sync::Semaphore; @@ -44,7 +43,6 @@ pub struct AppState { pub http_client: reqwest::Client, pub discovery_cache: Arc, pub geoip: Arc, - pub invite_meta: Arc>, pub index_html: Option>, pub budgets: AppProxyBudgets, }