feat(push): deliver our own relay endpoints in process (#2917)

This commit is contained in:
Hampus
2026-09-24 01:07:09 +02:00
committed by GitHub
parent b6e504f68c
commit b85e975fb5
6 changed files with 220 additions and 4 deletions
Generated
+1
View File
@@ -1894,6 +1894,7 @@ dependencies = [
"futures",
"hmac 0.13.0",
"p256",
"percent-encoding",
"rand 0.10.2",
"reqwest",
"ring",
+1
View File
@@ -23,6 +23,7 @@ base64 = "0.23.1"
clap = { version = "4.6.7", features = ["derive"] }
futures = "0.3.34"
hmac = "0.13.0"
percent-encoding = "2.3.2"
p256 = { version = "0.13.2", default-features = false, features = ["ecdh", "ecdsa", "pkcs8", "pem", "std"] }
rand = "0.10.2"
reqwest = { version = "0.13.5", default-features = false, features = ["http2", "json", "rustls"] }
+11
View File
@@ -183,6 +183,7 @@ pub struct DeliveryConfig {
pub vapid: VapidConfig,
pub apns: Option<ApnsConfig>,
pub fcm: Option<FcmConfig>,
pub own_relay_hosts: Vec<String>,
}
#[derive(Clone, Copy, Debug)]
@@ -262,10 +263,20 @@ impl DeliveryConfig {
vapid: vapid_config(&env)?,
apns: apns_config(&env)?,
fcm: fcm_config(&env)?,
own_relay_hosts: own_relay_hosts(&env),
})
}
}
fn own_relay_hosts(env: &Env) -> Vec<String> {
env.get("FLUXER_PUSH_SERVICE_OWN_RELAY_HOSTS")
.unwrap_or_default()
.split(',')
.map(|host| host.trim().to_ascii_lowercase())
.filter(|host| !host.is_empty())
.collect()
}
impl RelayConfig {
pub fn load_from_iter<I, K, V>(vars: I) -> anyhow::Result<Self>
where
+11
View File
@@ -354,6 +354,7 @@ pub struct Metrics {
sends: [[AtomicU64; SEND_RESULT_COUNT]; PROVIDER_COUNT],
token_deletions: [AtomicU64; PROVIDER_COUNT],
payload_shrinks: [AtomicU64; PAYLOAD_SHRINK_COUNT],
own_relay_shortcuts: AtomicU64,
auth_tokens_minted: [AtomicU64; AUTH_PROVIDER_COUNT],
rpc_requests: [[AtomicU64; RPC_OUTCOME_COUNT]; RPC_METHOD_COUNT],
rollout_updates: [AtomicU64; ROLLOUT_OUTCOME_COUNT],
@@ -384,6 +385,7 @@ impl Metrics {
sends: [const { [const { AtomicU64::new(0) }; SEND_RESULT_COUNT] }; PROVIDER_COUNT],
token_deletions: [const { AtomicU64::new(0) }; PROVIDER_COUNT],
payload_shrinks: [const { AtomicU64::new(0) }; PAYLOAD_SHRINK_COUNT],
own_relay_shortcuts: AtomicU64::new(0),
auth_tokens_minted: [const { AtomicU64::new(0) }; AUTH_PROVIDER_COUNT],
rpc_requests: [const { [const { AtomicU64::new(0) }; RPC_OUTCOME_COUNT] };
RPC_METHOD_COUNT],
@@ -438,6 +440,10 @@ impl Metrics {
self.token_deletions[provider as usize].fetch_add(1, ORDERING);
}
pub fn record_own_relay_shortcut(&self) {
self.own_relay_shortcuts.fetch_add(1, ORDERING);
}
pub fn record_payload_shrink(&self, step: PayloadShrink) {
self.payload_shrinks[step as usize].fetch_add(1, ORDERING);
}
@@ -539,6 +545,11 @@ impl Metrics {
Provider::ALL.map(Provider::label),
&self.token_deletions,
)?;
render_counter(
out,
"fluxer_push_own_relay_shortcuts_total",
&self.own_relay_shortcuts,
)?;
render_labelled_counter(
out,
"fluxer_push_payload_shrinks_total",
+14 -4
View File
@@ -2,6 +2,7 @@
pub mod apns;
pub mod fcm;
pub mod own_relay;
pub mod web_push;
use crate::metrics::{DeliveryRoute, Provider, SendResult, elapsed_ms};
@@ -90,10 +91,19 @@ pub async fn send(state: &AppState, sub: &Subscription, envelope: &Value) -> Sen
return SendOutcome::permanent("unsupported_platform");
};
let started_ms = now_ms();
let outcome = match route {
Route::WebPush => web_push::send(state, sub, envelope).await,
Route::LegacyApns => apns::send(state, sub, envelope).await,
Route::LegacyFcm => fcm::send(state, sub, envelope).await,
let direct = own_relay::parse(&sub.endpoint, &state.cfg.own_relay_hosts);
let outcome = match (route, direct) {
(Route::WebPush, Some(hop)) => {
let hopped = hop.as_subscription(sub);
state.metrics.record_own_relay_shortcut();
match hop.leg {
own_relay::Leg::Fcm => fcm::send(state, &hopped, envelope).await,
_ => apns::send(state, &hopped, envelope).await,
}
}
(Route::WebPush, None) => web_push::send(state, sub, envelope).await,
(Route::LegacyApns, _) => apns::send(state, sub, envelope).await,
(Route::LegacyFcm, _) => fcm::send(state, sub, envelope).await,
};
state.metrics.record_send(
provider_of(platform),
+182
View File
@@ -0,0 +1,182 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use crate::config::ProviderEnvironment;
use crate::subscription::Subscription;
use url::Url;
const APNS_SEGMENT: &str = "apns";
const APNS_VOIP_SEGMENT: &str = "apns-voip";
const FCM_SEGMENT: &str = "fcm";
const RELAY_PREFIX: [&str; 2] = ["relay", "v1"];
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Leg {
Apns,
ApnsVoip,
Fcm,
}
#[derive(Debug, Eq, PartialEq)]
pub struct Hop {
pub leg: Leg,
pub app_id: String,
pub environment: Option<ProviderEnvironment>,
pub device_token: String,
}
impl Hop {
pub fn as_subscription(&self, sub: &Subscription) -> Subscription {
Subscription {
subscription_id: sub.subscription_id.clone(),
endpoint: self.device_token.clone(),
p256dh_key: None,
auth_key: None,
platform: Some(
match self.leg {
Leg::Apns => "ios_apns",
Leg::ApnsVoip => "ios_apns_voip",
Leg::Fcm => "android_fcm",
}
.to_owned(),
),
app_id: Some(self.app_id.clone()),
provider_environment: self.environment.map(|environment| {
match environment {
ProviderEnvironment::Production => "production",
ProviderEnvironment::Development => "development",
}
.to_owned()
}),
}
}
}
pub fn parse(endpoint: &str, hosts: &[String]) -> Option<Hop> {
if hosts.is_empty() {
return None;
}
let url = Url::parse(endpoint).ok()?;
if url.scheme() != "https" {
return None;
}
let host = url.host_str()?.to_ascii_lowercase();
if !hosts.iter().any(|allowed| allowed == &host) {
return None;
}
let mut segments = url.path_segments()?;
for expected in RELAY_PREFIX {
if segments.next()? != expected {
return None;
}
}
let leg = match segments.next()? {
APNS_SEGMENT => Leg::Apns,
APNS_VOIP_SEGMENT => Leg::ApnsVoip,
FCM_SEGMENT => Leg::Fcm,
_ => return None,
};
let app_id = decode(segments.next()?)?;
let (environment, device_token) = match leg {
Leg::Fcm => (None, decode(segments.next()?)?),
_ => (
ProviderEnvironment::from_label(&decode(segments.next()?)?),
decode(segments.next()?)?,
),
};
if segments.next().is_some() || app_id.is_empty() || device_token.is_empty() {
return None;
}
if !matches!(leg, Leg::Fcm) && environment.is_none() {
return None;
}
Some(Hop {
leg,
app_id,
environment,
device_token,
})
}
fn decode(segment: &str) -> Option<String> {
percent_encoding::percent_decode_str(segment)
.decode_utf8()
.ok()
.map(|value| value.into_owned())
.filter(|value| !value.contains('/'))
}
#[cfg(test)]
mod tests {
use super::*;
const TOKEN: &str = "3dbc5a5ef1a1c1666afc26f466e1b3ebaaf4c66d92dddeb0fd1b69c49641d4cd";
fn ours() -> Vec<String> {
vec!["push.fluxer.com".to_owned()]
}
#[test]
fn an_apns_endpoint_on_our_own_relay_is_taken_in_process() {
let hop = parse(
&format!("https://push.fluxer.com/relay/v1/apns/canary/production/{TOKEN}"),
&ours(),
)
.expect("our own relay endpoint parses");
assert_eq!(hop.leg, Leg::Apns);
assert_eq!(hop.app_id, "canary");
assert_eq!(hop.environment, Some(ProviderEnvironment::Production));
assert_eq!(hop.device_token, TOKEN);
}
#[test]
fn an_fcm_endpoint_keeps_its_percent_encoded_token() {
let hop = parse(
"https://push.fluxer.com/relay/v1/fcm/canary/dYC_x9gXTjyyrG8_Aw3nUM%3AAPA91bExample",
&ours(),
)
.expect("an fcm endpoint parses");
assert_eq!(hop.leg, Leg::Fcm);
assert_eq!(hop.environment, None);
assert_eq!(hop.device_token, "dYC_x9gXTjyyrG8_Aw3nUM:APA91bExample");
}
#[test]
fn a_relay_we_do_not_operate_is_left_on_the_network_path() {
let endpoint = format!("https://push.example.org/relay/v1/apns/canary/production/{TOKEN}");
assert!(parse(&endpoint, &ours()).is_none());
}
#[test]
fn a_self_hosted_deployment_configures_no_hosts_and_never_shortcuts() {
let endpoint = format!("https://push.fluxer.com/relay/v1/apns/canary/production/{TOKEN}");
assert!(parse(&endpoint, &[]).is_none());
}
#[test]
fn a_third_party_web_push_endpoint_is_not_mistaken_for_a_relay_hop() {
assert!(
parse(
"https://updates.push.services.mozilla.com/wpush/v2/gAAAAA",
&ours()
)
.is_none()
);
assert!(parse("https://ntfy.sh/upZzH87cT9jJCc?up=1", &ours()).is_none());
}
#[test]
fn a_malformed_relay_path_is_refused() {
for endpoint in [
"https://push.fluxer.com/relay/v1/apns/canary/production",
"https://push.fluxer.com/relay/v1/apns/canary/production/tok/extra",
"https://push.fluxer.com/relay/v1/sms/canary/production/tok",
"https://push.fluxer.com/relay/v2/apns/canary/production/tok",
"http://push.fluxer.com/relay/v1/apns/canary/production/tok",
] {
assert!(
parse(endpoint, &ours()).is_none(),
"{endpoint} must not parse"
);
}
}
}