mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
refactor(app-proxy): drop invite metadata and database access (#2464)
This commit is contained in:
Generated
-3
@@ -1990,13 +1990,10 @@ dependencies = [
|
|||||||
"anyhow",
|
"anyhow",
|
||||||
"axum",
|
"axum",
|
||||||
"base64",
|
"base64",
|
||||||
"fluxer-svc",
|
|
||||||
"fluxer_common",
|
"fluxer_common",
|
||||||
"hex",
|
"hex",
|
||||||
"moka",
|
|
||||||
"rand 0.10.1",
|
"rand 0.10.1",
|
||||||
"reqwest",
|
"reqwest",
|
||||||
"scylla",
|
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
"tokio",
|
"tokio",
|
||||||
|
|||||||
@@ -140,7 +140,7 @@ FLUXER_DISCOVERY_ENABLED=true
|
|||||||
# budget roughly shared_buffers + (server max_connections x 12 MB) +
|
# 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
|
# (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
|
# 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_SERVER_MAX_CONNECTIONS=150
|
||||||
#FLUXER_POSTGRES_SHARED_BUFFERS=512MB
|
#FLUXER_POSTGRES_SHARED_BUFFERS=512MB
|
||||||
#FLUXER_POSTGRES_EFFECTIVE_CACHE_SIZE=2GB
|
#FLUXER_POSTGRES_EFFECTIVE_CACHE_SIZE=2GB
|
||||||
|
|||||||
@@ -460,7 +460,6 @@ services:
|
|||||||
limits:
|
limits:
|
||||||
memory: ${FLUXER_APP_PROXY_MEMORY_LIMIT:-256mb}
|
memory: ${FLUXER_APP_PROXY_MEMORY_LIMIT:-256mb}
|
||||||
environment:
|
environment:
|
||||||
<<: *fluxer-postgres-env
|
|
||||||
FLUXER_APP_PROXY_HOST: 0.0.0.0
|
FLUXER_APP_PROXY_HOST: 0.0.0.0
|
||||||
FLUXER_APP_PROXY_PORT: "8080"
|
FLUXER_APP_PROXY_PORT: "8080"
|
||||||
DISCOVERY_UPSTREAM_URL: http://caddy:8088/api/.well-known/fluxer
|
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_WORKER_SRC: ${FLUXER_CSP_EXTRA_WORKER_SRC:-}
|
||||||
FLUXER_CSP_EXTRA_MANIFEST_SRC: ${FLUXER_CSP_EXTRA_MANIFEST_SRC:-}
|
FLUXER_CSP_EXTRA_MANIFEST_SRC: ${FLUXER_CSP_EXTRA_MANIFEST_SRC:-}
|
||||||
FLUXER_CSP_REPORT_URI: ${FLUXER_CSP_REPORT_URI:-}
|
FLUXER_CSP_REPORT_URI: ${FLUXER_CSP_REPORT_URI:-}
|
||||||
FLUXER_POSTGRES_MAX_CONNECTIONS: "5"
|
|
||||||
depends_on:
|
depends_on:
|
||||||
api: {condition: service_healthy}
|
api: {condition: service_healthy}
|
||||||
caddy: {condition: service_healthy}
|
caddy: {condition: service_healthy}
|
||||||
postgres: {condition: service_healthy}
|
|
||||||
|
|
||||||
snowflakes:
|
snowflakes:
|
||||||
<<: *fluxer-service
|
<<: *fluxer-service
|
||||||
|
|||||||
@@ -10,12 +10,9 @@ anyhow = "1.0.104"
|
|||||||
axum = { version = "0.8.9", features = ["macros"] }
|
axum = { version = "0.8.9", features = ["macros"] }
|
||||||
base64 = "0.22"
|
base64 = "0.22"
|
||||||
fluxer_common = { path = "../fluxer_common" }
|
fluxer_common = { path = "../fluxer_common" }
|
||||||
fluxer-svc = { path = "../fluxer_svc" }
|
|
||||||
hex = "0.4"
|
hex = "0.4"
|
||||||
moka = { version = "0.12.15", features = ["future"] }
|
|
||||||
rand = "0.10"
|
rand = "0.10"
|
||||||
reqwest = { version = "0.13.4", default-features = false, features = ["json", "rustls", "stream", "gzip", "brotli", "deflate"] }
|
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 = { version = "1.0.228", features = ["derive"] }
|
||||||
serde_json = "1.0.150"
|
serde_json = "1.0.150"
|
||||||
tokio = { version = "1.52.3", features = ["macros", "net", "rt-multi-thread", "signal", "time", "fs"] }
|
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]
|
[features]
|
||||||
default = ["time-freeze"]
|
default = ["time-freeze"]
|
||||||
scylla = ["dep:scylla", "fluxer-svc/scylla"]
|
|
||||||
time-freeze = []
|
time-freeze = []
|
||||||
|
|||||||
@@ -118,14 +118,12 @@ WORKDIR /usr/src/app
|
|||||||
COPY Cargo.lock ./
|
COPY Cargo.lock ./
|
||||||
COPY fluxer_common/ fluxer_common/
|
COPY fluxer_common/ fluxer_common/
|
||||||
COPY fluxer_app_proxy/ fluxer_app_proxy/
|
COPY fluxer_app_proxy/ fluxer_app_proxy/
|
||||||
COPY fluxer_svc/ fluxer_svc/
|
|
||||||
|
|
||||||
RUN printf '%s\n' \
|
RUN printf '%s\n' \
|
||||||
'[workspace]' \
|
'[workspace]' \
|
||||||
'members = [' \
|
'members = [' \
|
||||||
' "fluxer_app_proxy",' \
|
' "fluxer_app_proxy",' \
|
||||||
' "fluxer_common",' \
|
' "fluxer_common",' \
|
||||||
' "fluxer_svc",' \
|
|
||||||
']' \
|
']' \
|
||||||
'resolver = "2"' \
|
'resolver = "2"' \
|
||||||
'' \
|
'' \
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
use fluxer_common::config::{self as cfg, GeoipS3Config, GeoipSourceConfig};
|
use fluxer_common::config::{self as cfg, GeoipS3Config, GeoipSourceConfig};
|
||||||
use fluxer_svc::config::{DatabaseBackend, normalize_host, parse_hosts};
|
|
||||||
use reqwest::Url;
|
use reqwest::Url;
|
||||||
use std::env;
|
use std::env;
|
||||||
use std::fmt;
|
use std::fmt;
|
||||||
@@ -368,25 +367,6 @@ pub struct AppProxyConfig {
|
|||||||
pub geoip_s3_config: Option<GeoipS3Config>,
|
pub geoip_s3_config: Option<GeoipS3Config>,
|
||||||
pub trust_client_ip_header: bool,
|
pub trust_client_ip_header: bool,
|
||||||
pub client_ip_header_name: String,
|
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<String>,
|
|
||||||
pub scylla_keyspace: String,
|
|
||||||
pub scylla_username: Option<String>,
|
|
||||||
pub scylla_password: Option<String>,
|
|
||||||
pub postgres_url: Option<String>,
|
|
||||||
pub postgres_host: String,
|
|
||||||
pub postgres_port: u16,
|
|
||||||
pub postgres_database: String,
|
|
||||||
pub postgres_username: String,
|
|
||||||
pub postgres_password: Option<String>,
|
|
||||||
pub postgres_ssl: bool,
|
|
||||||
pub postgres_ssl_ca: Option<String>,
|
|
||||||
pub postgres_max_connections: usize,
|
|
||||||
pub postgres_kv_table: String,
|
|
||||||
pub postgres_prepared_statements: bool,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
#[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 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::<Vec<_>>()
|
|
||||||
})
|
|
||||||
.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(
|
let s3_public_endpoint = parse_optional_http_endpoint(
|
||||||
"FLUXER_S3_PUBLIC_ENDPOINT",
|
"FLUXER_S3_PUBLIC_ENDPOINT",
|
||||||
cfg::non_empty_env("FLUXER_S3_PUBLIC_ENDPOINT"),
|
cfg::non_empty_env("FLUXER_S3_PUBLIC_ENDPOINT"),
|
||||||
@@ -581,48 +533,10 @@ impl AppProxyConfig {
|
|||||||
)
|
)
|
||||||
.trim()
|
.trim()
|
||||||
.to_ascii_lowercase(),
|
.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 {
|
fn resolve_discovery_upstream_url_from_env() -> String {
|
||||||
resolve_discovery_upstream_url(|name| env::var(name).ok())
|
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())
|
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<String> {
|
fn resolve_bootstrap_api_public_endpoint_from_env() -> Option<String> {
|
||||||
resolve_bootstrap_api_public_endpoint(|name| env::var(name).ok())
|
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<F>(mut read_var: F) -> bool
|
|
||||||
where
|
|
||||||
F: FnMut(&str) -> Option<String>,
|
|
||||||
{
|
|
||||||
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<F>(mut read_var: F) -> String
|
fn resolve_discovery_upstream_url<F>(mut read_var: F) -> String
|
||||||
where
|
where
|
||||||
F: FnMut(&str) -> Option<String>,
|
F: FnMut(&str) -> Option<String>,
|
||||||
@@ -749,11 +634,6 @@ mod tests {
|
|||||||
resolve_time_freeze_enabled(|name| env.get(name).map(|value| value.to_string()))
|
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<String> {
|
fn resolve_bootstrap_endpoint_from_pairs(pairs: &[(&str, &str)]) -> Option<String> {
|
||||||
let env: HashMap<&str, &str> = pairs.iter().copied().collect();
|
let env: HashMap<&str, &str> = pairs.iter().copied().collect();
|
||||||
resolve_bootstrap_api_public_endpoint(|name| env.get(name).map(|value| value.to_string()))
|
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]
|
#[test]
|
||||||
fn time_freeze_enabled_by_default_for_hosted_runtime() {
|
fn time_freeze_enabled_by_default_for_hosted_runtime() {
|
||||||
assert!(resolve_time_freeze_from_pairs(&[]));
|
assert!(resolve_time_freeze_from_pairs(&[]));
|
||||||
|
|||||||
@@ -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<String>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[derive(Clone, Debug, Default)]
|
|
||||||
pub struct InviteMetaEndpoints {
|
|
||||||
pub media_endpoint: Option<String>,
|
|
||||||
pub static_cdn_endpoint: Option<String>,
|
|
||||||
}
|
|
||||||
|
|
||||||
pub struct InviteMetaResolver {
|
|
||||||
storage: InviteMetaStorage,
|
|
||||||
cache: Cache<String, Option<InvitePageMeta>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
enum InviteMetaStorage {
|
|
||||||
Postgres(PostgresInviteMetaStorage),
|
|
||||||
#[cfg(feature = "scylla")]
|
|
||||||
Scylla(Box<ScyllaInviteMetaStorage>),
|
|
||||||
}
|
|
||||||
|
|
||||||
struct PostgresInviteMetaStorage {
|
|
||||||
kv: postgres::KvClient,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(feature = "scylla")]
|
|
||||||
struct ScyllaInviteMetaStorage {
|
|
||||||
db: Arc<Session>,
|
|
||||||
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<i64>,
|
|
||||||
channel_id: Option<i64>,
|
|
||||||
inviter_id: Option<i64>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg_attr(feature = "scylla", derive(DeserializeRow))]
|
|
||||||
#[derive(Debug, Deserialize)]
|
|
||||||
struct GuildDbRow {
|
|
||||||
guild_id: i64,
|
|
||||||
name: String,
|
|
||||||
icon_hash: Option<String>,
|
|
||||||
member_count: Option<i32>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg_attr(feature = "scylla", derive(DeserializeRow))]
|
|
||||||
#[derive(Debug, Deserialize)]
|
|
||||||
struct ChannelDbRow {
|
|
||||||
channel_id: i64,
|
|
||||||
r#type: i32,
|
|
||||||
name: Option<String>,
|
|
||||||
icon_hash: Option<String>,
|
|
||||||
recipient_ids: Option<HashSet<i64>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg_attr(feature = "scylla", derive(DeserializeRow))]
|
|
||||||
#[derive(Debug, Deserialize)]
|
|
||||||
struct UserDbRow {
|
|
||||||
user_id: i64,
|
|
||||||
username: String,
|
|
||||||
discriminator: i32,
|
|
||||||
global_name: Option<String>,
|
|
||||||
avatar_hash: Option<String>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl InviteMetaResolver {
|
|
||||||
pub async fn connect(config: &AppProxyConfig) -> anyhow::Result<Self> {
|
|
||||||
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<Self> {
|
|
||||||
let cache = invite_meta_cache(config);
|
|
||||||
|
|
||||||
Ok(Self {
|
|
||||||
storage: InviteMetaStorage::Postgres(PostgresInviteMetaStorage { kv }),
|
|
||||||
cache,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(feature = "scylla")]
|
|
||||||
async fn new_scylla(db: Arc<Session>, config: &AppProxyConfig) -> anyhow::Result<Self> {
|
|
||||||
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<Option<InvitePageMeta>> {
|
|
||||||
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<Option<InvitePageMeta>> {
|
|
||||||
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<Option<InvitePageMeta>> {
|
|
||||||
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<Option<InvitePageMeta>> {
|
|
||||||
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<Option<InviteDbRow>> {
|
|
||||||
self.storage.fetch_invite(code).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_guild(&self, guild_id: i64) -> anyhow::Result<Option<GuildDbRow>> {
|
|
||||||
self.storage.fetch_guild(guild_id).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_channel(&self, channel_id: i64) -> anyhow::Result<Option<ChannelDbRow>> {
|
|
||||||
self.storage.fetch_channel(channel_id).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_user(&self, user_id: i64) -> anyhow::Result<Option<UserDbRow>> {
|
|
||||||
self.storage.fetch_user(user_id).await
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
pub fn start_background_connect(
|
|
||||||
slot: &Arc<OnceLock<InviteMetaResolver>>,
|
|
||||||
config: Arc<AppProxyConfig>,
|
|
||||||
) -> 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<String, Option<InvitePageMeta>> {
|
|
||||||
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<Option<InviteDbRow>> {
|
|
||||||
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<Option<GuildDbRow>> {
|
|
||||||
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<Option<ChannelDbRow>> {
|
|
||||||
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<Option<UserDbRow>> {
|
|
||||||
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<Option<InviteDbRow>> {
|
|
||||||
self.fetch_row("invites", &[KeyPart::String(code)]).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_guild(&self, guild_id: i64) -> anyhow::Result<Option<GuildDbRow>> {
|
|
||||||
self.fetch_row("guilds", &[KeyPart::BigInt(guild_id)]).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_channel(&self, channel_id: i64) -> anyhow::Result<Option<ChannelDbRow>> {
|
|
||||||
self.fetch_row(
|
|
||||||
"channels",
|
|
||||||
&[KeyPart::BigInt(channel_id), KeyPart::Bool(false)],
|
|
||||||
)
|
|
||||||
.await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_user(&self, user_id: i64) -> anyhow::Result<Option<UserDbRow>> {
|
|
||||||
self.fetch_row("users", &[KeyPart::BigInt(user_id)]).await
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_row<T>(&self, table_name: &str, key: &[KeyPart<'_>]) -> anyhow::Result<Option<T>>
|
|
||||||
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<Option<InviteDbRow>> {
|
|
||||||
let result = self.db.execute_unpaged(&self.stmt_invite, (code,)).await?;
|
|
||||||
let rows = result.into_rows_result()?;
|
|
||||||
Ok(rows.maybe_first_row::<InviteDbRow>()?)
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_guild(&self, guild_id: i64) -> anyhow::Result<Option<GuildDbRow>> {
|
|
||||||
let result = self
|
|
||||||
.db
|
|
||||||
.execute_unpaged(&self.stmt_guild, (guild_id,))
|
|
||||||
.await?;
|
|
||||||
let rows = result.into_rows_result()?;
|
|
||||||
Ok(rows.maybe_first_row::<GuildDbRow>()?)
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_channel(&self, channel_id: i64) -> anyhow::Result<Option<ChannelDbRow>> {
|
|
||||||
let result = self
|
|
||||||
.db
|
|
||||||
.execute_unpaged(&self.stmt_channel, (channel_id,))
|
|
||||||
.await?;
|
|
||||||
let rows = result.into_rows_result()?;
|
|
||||||
Ok(rows.maybe_first_row::<ChannelDbRow>()?)
|
|
||||||
}
|
|
||||||
|
|
||||||
async fn fetch_user(&self, user_id: i64) -> anyhow::Result<Option<UserDbRow>> {
|
|
||||||
let result = self.db.execute_unpaged(&self.stmt_user, (user_id,)).await?;
|
|
||||||
let rows = result.into_rows_result()?;
|
|
||||||
Ok(rows.maybe_first_row::<UserDbRow>()?)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
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<String> {
|
|
||||||
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<String> {
|
|
||||||
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<String> {
|
|
||||||
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<String> {
|
|
||||||
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#"<meta name="description" content="{description}">
|
|
||||||
<meta property="og:title" content="{title}">
|
|
||||||
<meta property="og:description" content="{description}">
|
|
||||||
<meta property="og:type" content="website">
|
|
||||||
<meta name="twitter:card" content="summary">
|
|
||||||
<meta name="twitter:title" content="{title}">
|
|
||||||
<meta name="twitter:description" content="{description}">"#
|
|
||||||
);
|
|
||||||
if let Some(image_url) = &meta.image_url {
|
|
||||||
let image = escape_html_attr(image_url);
|
|
||||||
tags.push_str(&format!(
|
|
||||||
r#"
|
|
||||||
<meta property="og:image" content="{image}">
|
|
||||||
<meta name="twitter:image" content="{image}">"#
|
|
||||||
));
|
|
||||||
}
|
|
||||||
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("<title>") else {
|
|
||||||
return html.to_owned();
|
|
||||||
};
|
|
||||||
let content_start = start + "<title>".len();
|
|
||||||
let Some(relative_end) = lower[content_start..].find("</title>") 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("<meta") {
|
|
||||||
let start = cursor + relative_start;
|
|
||||||
let Some(relative_end) = lower[start..].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("</title>") {
|
|
||||||
let insert_at = title_end + "</title>".len();
|
|
||||||
return insert_at_pos(html, insert_at, tags);
|
|
||||||
}
|
|
||||||
if let Some(head_start) = lower
|
|
||||||
.find("<head>")
|
|
||||||
.map(|pos| pos + "<head>".len())
|
|
||||||
.or_else(|| {
|
|
||||||
lower
|
|
||||||
.find("<head ")
|
|
||||||
.and_then(|pos| lower[pos..].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#"<html><head><title>Fluxer</title><meta name="description" content="Default"><meta property="og:title" content="Old"></head><body></body></html>"#;
|
|
||||||
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("<title>Join A & B</title>"));
|
|
||||||
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""#)
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -7,7 +7,6 @@ pub mod discovery_cache;
|
|||||||
#[cfg(feature = "time-freeze")]
|
#[cfg(feature = "time-freeze")]
|
||||||
pub mod frozen_snapshots;
|
pub mod frozen_snapshots;
|
||||||
pub mod geoip;
|
pub mod geoip;
|
||||||
pub mod invite_meta;
|
|
||||||
pub mod routes;
|
pub mod routes;
|
||||||
pub mod state;
|
pub mod state;
|
||||||
pub mod time_freeze;
|
pub mod time_freeze;
|
||||||
|
|||||||
@@ -5,13 +5,13 @@ use fluxer_app_proxy::{
|
|||||||
config::AppProxyConfig,
|
config::AppProxyConfig,
|
||||||
csp::CompiledCspPolicy,
|
csp::CompiledCspPolicy,
|
||||||
discovery_cache::DiscoveryCache,
|
discovery_cache::DiscoveryCache,
|
||||||
geoip, invite_meta,
|
geoip,
|
||||||
routes::build_router,
|
routes::build_router,
|
||||||
state::{
|
state::{
|
||||||
AppProxyBudgets, AppState, MAX_SPA_INDEX_BYTES, build_http_client, read_bounded_text_file,
|
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 tokio::{net::TcpListener, runtime::Builder};
|
||||||
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
||||||
|
|
||||||
@@ -56,11 +56,6 @@ fn main() -> anyhow::Result<()> {
|
|||||||
config.discovery_refresh_interval_ms,
|
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_html = if config.index_upstream_url.is_none() {
|
||||||
let index_path = std::path::Path::new(&config.static_dir).join("index.html");
|
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 {
|
match read_bounded_text_file(&index_path, MAX_SPA_INDEX_BYTES).await {
|
||||||
@@ -80,7 +75,6 @@ fn main() -> anyhow::Result<()> {
|
|||||||
http_client,
|
http_client,
|
||||||
discovery_cache,
|
discovery_cache,
|
||||||
geoip,
|
geoip,
|
||||||
invite_meta,
|
|
||||||
index_html,
|
index_html,
|
||||||
budgets: AppProxyBudgets::default(),
|
budgets: AppProxyBudgets::default(),
|
||||||
};
|
};
|
||||||
@@ -97,9 +91,6 @@ fn main() -> anyhow::Result<()> {
|
|||||||
.context("app proxy server exited unexpectedly")?;
|
.context("app proxy server exited unexpectedly")?;
|
||||||
|
|
||||||
cancel.abort();
|
cancel.abort();
|
||||||
if let Some(handle) = invite_meta_connect {
|
|
||||||
handle.abort();
|
|
||||||
}
|
|
||||||
Ok(())
|
Ok(())
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -436,8 +436,8 @@ mod tests {
|
|||||||
use axum::http::header::HeaderName;
|
use axum::http::header::HeaderName;
|
||||||
use fluxer_common::config::GeoipSourceConfig;
|
use fluxer_common::config::GeoipSourceConfig;
|
||||||
use fluxer_common::geoip::{GeoipConfig, GeoipResolver};
|
use fluxer_common::geoip::{GeoipConfig, GeoipResolver};
|
||||||
|
use std::sync::Arc;
|
||||||
use std::sync::atomic::{AtomicU64, Ordering};
|
use std::sync::atomic::{AtomicU64, Ordering};
|
||||||
use std::sync::{Arc, OnceLock};
|
|
||||||
use tower::ServiceExt;
|
use tower::ServiceExt;
|
||||||
|
|
||||||
async fn spawn_upstream(status: StatusCode, cache_control: &'static str) -> String {
|
async fn spawn_upstream(status: StatusCode, cache_control: &'static str) -> String {
|
||||||
@@ -491,7 +491,6 @@ mod tests {
|
|||||||
trust_client_ip_header: false,
|
trust_client_ip_header: false,
|
||||||
client_ip_header_name: "x-forwarded-for".to_owned(),
|
client_ip_header_name: "x-forwarded-for".to_owned(),
|
||||||
})),
|
})),
|
||||||
invite_meta: Arc::new(OnceLock::new()),
|
|
||||||
index_html: None,
|
index_html: None,
|
||||||
budgets: crate::state::AppProxyBudgets::default(),
|
budgets: crate::state::AppProxyBudgets::default(),
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,35 +12,14 @@ pub async fn health() -> &'static str {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub async fn ready(State(state): State<AppState>) -> Response {
|
pub async fn ready(State(state): State<AppState>) -> Response {
|
||||||
readiness_report(
|
readiness_report(state.discovery_cache.has_snapshot().await).into_response()
|
||||||
state.discovery_cache.has_snapshot().await,
|
|
||||||
state.config.invite_meta_enabled,
|
|
||||||
state.invite_meta.get().is_some(),
|
|
||||||
)
|
|
||||||
.into_response()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn readiness_report(
|
fn readiness_report(discovery_cached: bool) -> (StatusCode, &'static str) {
|
||||||
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(", "))
|
|
||||||
};
|
|
||||||
if discovery_cached {
|
if discovery_cached {
|
||||||
(StatusCode::OK, format!("OK{suffix}"))
|
(StatusCode::OK, "OK")
|
||||||
} else {
|
} else {
|
||||||
(
|
(StatusCode::SERVICE_UNAVAILABLE, "NOT READY: discovery")
|
||||||
StatusCode::SERVICE_UNAVAILABLE,
|
|
||||||
format!("NOT READY: discovery{suffix}"),
|
|
||||||
)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -55,7 +34,7 @@ mod tests {
|
|||||||
use axum::http::{HeaderValue, Request as HttpRequest, header};
|
use axum::http::{HeaderValue, Request as HttpRequest, header};
|
||||||
use fluxer_common::config::GeoipSourceConfig;
|
use fluxer_common::config::GeoipSourceConfig;
|
||||||
use fluxer_common::geoip::{GeoipConfig, GeoipResolver};
|
use fluxer_common::geoip::{GeoipConfig, GeoipResolver};
|
||||||
use std::sync::{Arc, OnceLock};
|
use std::sync::Arc;
|
||||||
use tower::ServiceExt;
|
use tower::ServiceExt;
|
||||||
|
|
||||||
const DISCOVERY_BODY: &str = r#"{"api_code_version":"proxy-test"}"#;
|
const DISCOVERY_BODY: &str = r#"{"api_code_version":"proxy-test"}"#;
|
||||||
@@ -77,9 +56,8 @@ mod tests {
|
|||||||
format!("http://{addr}/")
|
format!("http://{addr}/")
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn probe_state(invite_meta_enabled: bool) -> AppState {
|
async fn probe_state() -> AppState {
|
||||||
let mut config = AppProxyConfig::from_env();
|
let mut config = AppProxyConfig::from_env();
|
||||||
config.invite_meta_enabled = invite_meta_enabled;
|
|
||||||
config.discovery_upstream_url = spawn_discovery_origin().await;
|
config.discovery_upstream_url = spawn_discovery_origin().await;
|
||||||
let csp = Arc::new(
|
let csp = Arc::new(
|
||||||
crate::csp::CompiledCspPolicy::from_config(&config)
|
crate::csp::CompiledCspPolicy::from_config(&config)
|
||||||
@@ -98,7 +76,6 @@ mod tests {
|
|||||||
trust_client_ip_header: false,
|
trust_client_ip_header: false,
|
||||||
client_ip_header_name: "x-forwarded-for".to_owned(),
|
client_ip_header_name: "x-forwarded-for".to_owned(),
|
||||||
})),
|
})),
|
||||||
invite_meta: Arc::new(OnceLock::new()),
|
|
||||||
index_html: None,
|
index_html: None,
|
||||||
budgets: crate::state::AppProxyBudgets::default(),
|
budgets: crate::state::AppProxyBudgets::default(),
|
||||||
}
|
}
|
||||||
@@ -131,7 +108,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn liveness_stays_constant_while_the_proxy_cannot_serve() {
|
async fn liveness_stays_constant_while_the_proxy_cannot_serve() {
|
||||||
let state = probe_state(false).await;
|
let state = probe_state().await;
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
probe(state, "/_health").await,
|
probe(state, "/_health").await,
|
||||||
(StatusCode::OK, "OK".to_owned())
|
(StatusCode::OK, "OK".to_owned())
|
||||||
@@ -140,7 +117,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn readiness_fails_while_the_discovery_cache_is_empty() {
|
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;
|
let (status, body) = probe(state, "/_ready").await;
|
||||||
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
|
assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE);
|
||||||
assert!(body.contains("discovery"), "{body}");
|
assert!(body.contains("discovery"), "{body}");
|
||||||
@@ -148,7 +125,7 @@ mod tests {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn readiness_passes_once_a_discovery_snapshot_is_cached() {
|
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;
|
warm_discovery(&state).await;
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
probe(state, "/_ready").await,
|
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]
|
#[test]
|
||||||
fn invite_metadata_is_reported_as_degraded_instead_of_gating_readiness() {
|
fn readiness_tracks_only_the_discovery_snapshot() {
|
||||||
assert_eq!(readiness_report(true, false, false).0, StatusCode::OK);
|
assert_eq!(readiness_report(true), (StatusCode::OK, "OK"));
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
readiness_report(true, true, true),
|
readiness_report(false),
|
||||||
(StatusCode::OK, "OK".to_owned())
|
(StatusCode::SERVICE_UNAVAILABLE, "NOT READY: discovery")
|
||||||
);
|
|
||||||
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()
|
|
||||||
)
|
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,9 +5,6 @@ use crate::config::HttpEndpoint;
|
|||||||
use crate::csp::{RuntimeCspSources, generate_nonce};
|
use crate::csp::{RuntimeCspSources, generate_nonce};
|
||||||
use crate::discovery_cache::{DiscoveryResponse, discovery_endpoint};
|
use crate::discovery_cache::{DiscoveryResponse, discovery_endpoint};
|
||||||
use crate::geoip::build_geoip_response;
|
use crate::geoip::build_geoip_response;
|
||||||
use crate::invite_meta::{
|
|
||||||
InviteMetaEndpoints, InvitePageMeta, inject_invite_meta, invite_code_from_path,
|
|
||||||
};
|
|
||||||
use crate::state::{
|
use crate::state::{
|
||||||
AppProxyBudgets, AppState, MAX_RENDERED_SPA_INDEX_BYTES, MAX_SPA_INDEX_BYTES,
|
AppProxyBudgets, AppState, MAX_RENDERED_SPA_INDEX_BYTES, MAX_SPA_INDEX_BYTES,
|
||||||
read_bounded_text_file,
|
read_bounded_text_file,
|
||||||
@@ -60,7 +57,7 @@ pub async fn spa_catch_all(
|
|||||||
.await;
|
.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";
|
const CRAWL_CONTROL_CACHE_CONTROL: &str = "public, max-age=300, must-revalidate";
|
||||||
@@ -144,7 +141,7 @@ async fn serve_static_file(
|
|||||||
response
|
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 time_freeze = load_time_freeze_config_for_request(&state.config, headers);
|
||||||
let debug_header = time_freeze_debug_header(&time_freeze);
|
let debug_header = time_freeze_debug_header(&time_freeze);
|
||||||
let should_bust_dev_assets = state.config.index_upstream_url.is_some();
|
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 nonce = generate_nonce();
|
||||||
let runtime_csp_sources = build_runtime_csp_sources(state, &discovery);
|
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
|
let static_cdn_endpoint = runtime_csp_sources
|
||||||
.static_cdn_endpoint
|
.static_cdn_endpoint
|
||||||
.as_ref()
|
.as_ref()
|
||||||
@@ -188,7 +184,6 @@ async fn serve_spa_index(state: &AppState, headers: &HeaderMap, request_path: &s
|
|||||||
&script_tag,
|
&script_tag,
|
||||||
static_cdn_endpoint,
|
static_cdn_endpoint,
|
||||||
media_endpoint,
|
media_endpoint,
|
||||||
invite_meta.as_ref(),
|
|
||||||
dev_buster.as_deref(),
|
dev_buster.as_deref(),
|
||||||
) {
|
) {
|
||||||
Ok(html) => html,
|
Ok(html) => html,
|
||||||
@@ -233,7 +228,6 @@ fn render_spa_document(
|
|||||||
script_tag: &str,
|
script_tag: &str,
|
||||||
static_cdn_endpoint: &str,
|
static_cdn_endpoint: &str,
|
||||||
media_endpoint: &str,
|
media_endpoint: &str,
|
||||||
invite_meta: Option<&InvitePageMeta>,
|
|
||||||
dev_asset_cache_buster: Option<&str>,
|
dev_asset_cache_buster: Option<&str>,
|
||||||
) -> Result<String, SpaDocumentSizeLimitError> {
|
) -> Result<String, SpaDocumentSizeLimitError> {
|
||||||
let mut document = bounded_document(inject_bootstrap(
|
let mut document = bounded_document(inject_bootstrap(
|
||||||
@@ -243,9 +237,6 @@ fn render_spa_document(
|
|||||||
static_cdn_endpoint,
|
static_cdn_endpoint,
|
||||||
media_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 {
|
if let Some(buster) = dev_asset_cache_buster {
|
||||||
document = bounded_document(append_dev_asset_cache_buster(&document, buster))?;
|
document = bounded_document(append_dev_asset_cache_buster(&document, buster))?;
|
||||||
}
|
}
|
||||||
@@ -259,33 +250,6 @@ async fn refresh_discovery_for_spa(state: &AppState) -> Option<DiscoveryResponse
|
|||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn resolve_invite_meta(
|
|
||||||
state: &AppState,
|
|
||||||
request_path: &str,
|
|
||||||
runtime_csp_sources: &RuntimeCspSources,
|
|
||||||
) -> Option<InvitePageMeta> {
|
|
||||||
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 {
|
fn build_runtime_csp_sources(state: &AppState, discovery: &DiscoveryResponse) -> RuntimeCspSources {
|
||||||
RuntimeCspSources {
|
RuntimeCspSources {
|
||||||
static_cdn_endpoint: discovery_endpoint(discovery, "static_cdn")
|
static_cdn_endpoint: discovery_endpoint(discovery, "static_cdn")
|
||||||
@@ -625,7 +589,7 @@ mod tests {
|
|||||||
use axum::body::Body;
|
use axum::body::Body;
|
||||||
use fluxer_common::config::GeoipSourceConfig;
|
use fluxer_common::config::GeoipSourceConfig;
|
||||||
use fluxer_common::geoip::{GeoipConfig, GeoipResolver};
|
use fluxer_common::geoip::{GeoipConfig, GeoipResolver};
|
||||||
use std::sync::{Arc, OnceLock};
|
use std::sync::Arc;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn dev_asset_cache_buster_rewrites_script_and_link_assets() {
|
fn dev_asset_cache_buster_rewrites_script_and_link_assets() {
|
||||||
@@ -677,14 +641,6 @@ mod tests {
|
|||||||
|
|
||||||
const SHELL_WITH_A_NONCE_HOLE: &str = r#"<!doctype html><html><head><title>Fluxer</title><script nonce="{{CSP_NONCE_PLACEHOLDER}}"></script><script src="/assets/app.js"></script></head><body></body></html>"#;
|
const SHELL_WITH_A_NONCE_HOLE: &str = r#"<!doctype html><html><head><title>Fluxer</title><script nonce="{{CSP_NONCE_PLACEHOLDER}}"></script><script src="/assets/app.js"></script></head><body></body></html>"#;
|
||||||
|
|
||||||
fn sample_invite_meta() -> InvitePageMeta {
|
|
||||||
InvitePageMeta {
|
|
||||||
title: "Join Sample Space".to_owned(),
|
|
||||||
description: "A sample invite".to_owned(),
|
|
||||||
image_url: None,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn the_rendered_document_always_carries_the_bootstrap_and_a_real_nonce() {
|
fn the_rendered_document_always_carries_the_bootstrap_and_a_real_nonce() {
|
||||||
let rendered = render_spa_document(
|
let rendered = render_spa_document(
|
||||||
@@ -694,7 +650,6 @@ mod tests {
|
|||||||
"https://static.example.test",
|
"https://static.example.test",
|
||||||
"",
|
"",
|
||||||
None,
|
None,
|
||||||
None,
|
|
||||||
)
|
)
|
||||||
.expect("test SPA document must render within its size limit");
|
.expect("test SPA document must render within its size limit");
|
||||||
|
|
||||||
@@ -703,36 +658,6 @@ mod tests {
|
|||||||
assert!(rendered.contains("<script>booted</script>"));
|
assert!(rendered.contains("<script>booted</script>"));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[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",
|
|
||||||
"<script>booted</script>",
|
|
||||||
"",
|
|
||||||
"",
|
|
||||||
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",
|
|
||||||
"<script>booted</script>",
|
|
||||||
"",
|
|
||||||
"",
|
|
||||||
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]
|
#[test]
|
||||||
fn the_dev_cache_buster_reaches_the_rendered_document_only_when_supplied() {
|
fn the_dev_cache_buster_reaches_the_rendered_document_only_when_supplied() {
|
||||||
let busted = render_spa_document(
|
let busted = render_spa_document(
|
||||||
@@ -741,7 +666,6 @@ mod tests {
|
|||||||
"<script>booted</script>",
|
"<script>booted</script>",
|
||||||
"",
|
"",
|
||||||
"",
|
"",
|
||||||
None,
|
|
||||||
Some("9911"),
|
Some("9911"),
|
||||||
)
|
)
|
||||||
.expect("test SPA document must render within its size limit");
|
.expect("test SPA document must render within its size limit");
|
||||||
@@ -752,7 +676,6 @@ mod tests {
|
|||||||
"",
|
"",
|
||||||
"",
|
"",
|
||||||
None,
|
None,
|
||||||
None,
|
|
||||||
)
|
)
|
||||||
.expect("test SPA document must render within its size limit");
|
.expect("test SPA document must render within its size limit");
|
||||||
|
|
||||||
@@ -778,7 +701,6 @@ mod tests {
|
|||||||
"https://fluxerstatic.com",
|
"https://fluxerstatic.com",
|
||||||
"",
|
"",
|
||||||
None,
|
None,
|
||||||
None,
|
|
||||||
)
|
)
|
||||||
.expect("test SPA document must render within its size limit");
|
.expect("test SPA document must render within its size limit");
|
||||||
|
|
||||||
@@ -810,7 +732,6 @@ mod tests {
|
|||||||
"https://cdn.example.test/",
|
"https://cdn.example.test/",
|
||||||
"https://media.example.test",
|
"https://media.example.test",
|
||||||
None,
|
None,
|
||||||
None,
|
|
||||||
)
|
)
|
||||||
.expect("test SPA document must render within its size limit");
|
.expect("test SPA document must render within its size limit");
|
||||||
|
|
||||||
@@ -844,7 +765,6 @@ mod tests {
|
|||||||
"https://cdn.example.test",
|
"https://cdn.example.test",
|
||||||
"https://media.example.test/",
|
"https://media.example.test/",
|
||||||
None,
|
None,
|
||||||
None,
|
|
||||||
)
|
)
|
||||||
.expect("test SPA document must render within its size limit");
|
.expect("test SPA document must render within its size limit");
|
||||||
assert!(
|
assert!(
|
||||||
@@ -861,7 +781,6 @@ mod tests {
|
|||||||
"https://cdn.example.test",
|
"https://cdn.example.test",
|
||||||
"https://cdn.example.test",
|
"https://cdn.example.test",
|
||||||
None,
|
None,
|
||||||
None,
|
|
||||||
)
|
)
|
||||||
.expect("test SPA document must render within its size limit");
|
.expect("test SPA document must render within its size limit");
|
||||||
assert!(
|
assert!(
|
||||||
@@ -967,7 +886,6 @@ mod tests {
|
|||||||
trust_client_ip_header: false,
|
trust_client_ip_header: false,
|
||||||
client_ip_header_name: "x-forwarded-for".to_owned(),
|
client_ip_header_name: "x-forwarded-for".to_owned(),
|
||||||
})),
|
})),
|
||||||
invite_meta: Arc::new(OnceLock::new()),
|
|
||||||
index_html: cached_shell.map(Arc::from),
|
index_html: cached_shell.map(Arc::from),
|
||||||
budgets: crate::state::AppProxyBudgets::default(),
|
budgets: crate::state::AppProxyBudgets::default(),
|
||||||
}
|
}
|
||||||
@@ -1002,7 +920,7 @@ mod tests {
|
|||||||
let state =
|
let state =
|
||||||
spa_state_serving(ReleaseChannel::Canary, Some(SHELL_WITH_ENDPOINT_HOLES)).await;
|
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);
|
assert_eq!(response.status(), StatusCode::OK);
|
||||||
let granted_nonce = nonce_granted_by(&response);
|
let granted_nonce = nonce_granted_by(&response);
|
||||||
let served = read_document(response).await;
|
let served = read_document(response).await;
|
||||||
@@ -1062,7 +980,7 @@ mod tests {
|
|||||||
|
|
||||||
let state = spa_state_serving(ReleaseChannel::Stable, None).await;
|
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);
|
assert_eq!(response.status(), StatusCode::OK);
|
||||||
let served = read_document(response).await;
|
let served = read_document(response).await;
|
||||||
|
|
||||||
@@ -1099,7 +1017,7 @@ mod tests {
|
|||||||
|
|
||||||
let state = spa_state_serving(ReleaseChannel::Stable, None).await;
|
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);
|
assert_eq!(response.status(), StatusCode::OK);
|
||||||
let granted_nonce = nonce_granted_by(&response);
|
let granted_nonce = nonce_granted_by(&response);
|
||||||
let served = read_document(response).await;
|
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"
|
"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!(
|
assert_ne!(
|
||||||
nonce_granted_by(&second),
|
nonce_granted_by(&second),
|
||||||
granted_nonce,
|
granted_nonce,
|
||||||
@@ -1134,7 +1052,7 @@ mod tests {
|
|||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn every_branch_announces_which_snapshot_decision_it_took() {
|
async fn every_branch_announces_which_snapshot_decision_it_took() {
|
||||||
let frozen_state = spa_state_serving(ReleaseChannel::Stable, None).await;
|
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!(
|
assert_eq!(
|
||||||
frozen
|
frozen
|
||||||
.headers()
|
.headers()
|
||||||
@@ -1151,7 +1069,7 @@ mod tests {
|
|||||||
|
|
||||||
let live_state =
|
let live_state =
|
||||||
spa_state_serving(ReleaseChannel::Canary, Some(SHELL_WITH_ENDPOINT_HOLES)).await;
|
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!(
|
assert_eq!(
|
||||||
live.headers()
|
live.headers()
|
||||||
.get("x-time-freeze")
|
.get("x-time-freeze")
|
||||||
@@ -1168,7 +1086,7 @@ mod tests {
|
|||||||
async fn the_frozen_shell_is_never_served_with_the_asset_lifetime() {
|
async fn the_frozen_shell_is_never_served_with_the_asset_lifetime() {
|
||||||
let state = spa_state_serving(ReleaseChannel::Stable, None).await;
|
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!(
|
assert_eq!(
|
||||||
response
|
response
|
||||||
@@ -1247,7 +1165,7 @@ mod tests {
|
|||||||
let state =
|
let state =
|
||||||
spa_state_serving(ReleaseChannel::Canary, Some(SHELL_WITH_ENDPOINT_HOLES)).await;
|
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
|
let cache_control = response
|
||||||
.headers()
|
.headers()
|
||||||
@@ -1270,7 +1188,7 @@ mod tests {
|
|||||||
let index_upstream_url = spawn_local_origin(SHELL_WITH_ENDPOINT_HOLES, "text/html").await;
|
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 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.status(), StatusCode::OK);
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
response
|
response
|
||||||
@@ -1304,7 +1222,7 @@ mod tests {
|
|||||||
let state =
|
let state =
|
||||||
spa_state_without_discovered_endpoints(Some("https://fallbackcdn.example.test")).await;
|
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);
|
assert_eq!(response.status(), StatusCode::OK);
|
||||||
let served = read_document(response).await;
|
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() {
|
async fn an_endpoint_neither_discovered_nor_configured_warms_no_socket_at_all() {
|
||||||
let state = spa_state_without_discovered_endpoints(None).await;
|
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);
|
assert_eq!(response.status(), StatusCode::OK);
|
||||||
let served = read_document(response).await;
|
let served = read_document(response).await;
|
||||||
|
|
||||||
|
|||||||
@@ -3,9 +3,8 @@
|
|||||||
use crate::config::AppProxyConfig;
|
use crate::config::AppProxyConfig;
|
||||||
use crate::csp::CompiledCspPolicy;
|
use crate::csp::CompiledCspPolicy;
|
||||||
use crate::discovery_cache::DiscoveryCache;
|
use crate::discovery_cache::DiscoveryCache;
|
||||||
use crate::invite_meta::InviteMetaResolver;
|
|
||||||
use fluxer_common::geoip::GeoipResolver;
|
use fluxer_common::geoip::GeoipResolver;
|
||||||
use std::sync::{Arc, OnceLock};
|
use std::sync::Arc;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
use tokio::io::AsyncReadExt;
|
use tokio::io::AsyncReadExt;
|
||||||
use tokio::sync::Semaphore;
|
use tokio::sync::Semaphore;
|
||||||
@@ -44,7 +43,6 @@ pub struct AppState {
|
|||||||
pub http_client: reqwest::Client,
|
pub http_client: reqwest::Client,
|
||||||
pub discovery_cache: Arc<DiscoveryCache>,
|
pub discovery_cache: Arc<DiscoveryCache>,
|
||||||
pub geoip: Arc<GeoipResolver>,
|
pub geoip: Arc<GeoipResolver>,
|
||||||
pub invite_meta: Arc<OnceLock<InviteMetaResolver>>,
|
|
||||||
pub index_html: Option<Arc<str>>,
|
pub index_html: Option<Arc<str>>,
|
||||||
pub budgets: AppProxyBudgets,
|
pub budgets: AppProxyBudgets,
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user