diff --git a/fluxer_gateway/src/push/push.erl b/fluxer_gateway/src/push/push.erl index f9967d44a..f5f9f46ee 100644 --- a/fluxer_gateway/src/push/push.erl +++ b/fluxer_gateway/src/push/push.erl @@ -34,6 +34,9 @@ -define(CNT_DISPATCH_DROPPED, push_loss_dispatch_dropped). -define(CNT_DISPATCH_DROPPED_USERS, push_loss_dispatch_dropped_users). -define(CNT_CLEAR_DROPPED, push_loss_clear_dropped). +-define(CNT_CLEAR_RETRIED, push_clear_retried). +-define(CLEAR_RETRY_ATTEMPTS, 3). +-define(CLEAR_RETRY_BASE_MS, 250). -define(CNT_QUEUE_FULL, push_loss_queue_full). -define(CNT_INVALID_JOB, push_loss_invalid_job). -define(CNT_ENQUEUE_TIMEOUT, push_loss_enqueue_timeout). @@ -135,6 +138,8 @@ handle_info(evict_caches, State) -> }), schedule_eviction(), {noreply, State}; +handle_info({retry_clear_notifications, UserId, ChannelId, MessageId, Attempt}, State) -> + clear_via_dispatcher(UserId, ChannelId, MessageId, Attempt, State); handle_info(_Info, State) -> {noreply, State}. @@ -855,6 +860,12 @@ clear_via_service(UserId, ChannelId, MessageId, ConfigVersion, State) -> -spec clear_via_dispatcher(integer(), integer(), integer(), state()) -> {noreply, state()}. clear_via_dispatcher(UserId, ChannelId, MessageId, State) -> + clear_via_dispatcher(UserId, ChannelId, MessageId, 0, State). + +-spec clear_via_dispatcher( + integer(), integer(), integer(), non_neg_integer(), state() +) -> {noreply, state()}. +clear_via_dispatcher(UserId, ChannelId, MessageId, Attempt, State) -> BadgeCountsTtl = maps:get(badge_counts_ttl_seconds, State), case push_dispatcher:enqueue_clear_notifications( @@ -864,11 +875,24 @@ clear_via_dispatcher(UserId, ChannelId, MessageId, State) -> ok -> ok; dropped -> - count_clear_dropped(), - log_clear_drop(loss_logging_enabled(), UserId, ChannelId, MessageId) + retry_or_drop_clear(UserId, ChannelId, MessageId, Attempt) end, {noreply, State}. +-spec retry_or_drop_clear(integer(), integer(), integer(), non_neg_integer()) -> ok. +retry_or_drop_clear(UserId, ChannelId, MessageId, Attempt) when + Attempt < ?CLEAR_RETRY_ATTEMPTS +-> + bump_counter(?CNT_CLEAR_RETRIED), + Delay = ?CLEAR_RETRY_BASE_MS bsl Attempt, + _ = erlang:send_after( + Delay, self(), {retry_clear_notifications, UserId, ChannelId, MessageId, Attempt + 1} + ), + ok; +retry_or_drop_clear(UserId, ChannelId, MessageId, _Attempt) -> + count_clear_dropped(), + log_clear_drop(loss_logging_enabled(), UserId, ChannelId, MessageId). + -spec log_clear_drop(boolean(), integer(), integer(), integer()) -> ok. log_clear_drop(true, UserId, ChannelId, MessageId) -> logger:warning("Push: dispatcher saturated, dropping clear notification job", #{ diff --git a/fluxer_gateway/test/push_tests.erl b/fluxer_gateway/test/push_tests.erl index 7d66a3ff5..d7bd49eb0 100644 --- a/fluxer_gateway/test/push_tests.erl +++ b/fluxer_gateway/test/push_tests.erl @@ -9,6 +9,44 @@ clear_channel_notifications_disabled_by_default_test() -> erase_persistent_term(push_clear_notifications_enabled), ?assertEqual(ok, push:clear_channel_notifications(1, 2, 3)). +a_saturated_dispatcher_retries_the_clear_before_dropping_it_test() -> + ok = meck:new(push_dispatcher, [passthrough, no_link]), + try + ok = meck:expect( + push_dispatcher, enqueue_clear_notifications, fun(_U, _C, _M, _T) -> dropped end + ), + State = #{badge_counts_ttl_seconds => 0}, + ?assertEqual( + {noreply, State}, + push:handle_info({retry_clear_notifications, 1, 2, 3, 0}, State) + ), + receive + {retry_clear_notifications, 1, 2, 3, 1} -> ok + after 2000 -> erlang:error(no_retry_scheduled) + end + after + meck:unload(push_dispatcher) + end. + +a_clear_is_dropped_only_after_the_retry_budget_is_spent_test() -> + ok = meck:new(push_dispatcher, [passthrough, no_link]), + try + ok = meck:expect( + push_dispatcher, enqueue_clear_notifications, fun(_U, _C, _M, _T) -> dropped end + ), + State = #{badge_counts_ttl_seconds => 0}, + ?assertEqual( + {noreply, State}, + push:handle_info({retry_clear_notifications, 1, 2, 3, 3}, State) + ), + receive + {retry_clear_notifications, _, _, _, _} -> erlang:error(retried_past_budget) + after 700 -> ok + end + after + meck:unload(push_dispatcher) + end. + push_owner_key_prefers_first_recipient_test() -> ?assertEqual( 42, diff --git a/fluxer_push/src/payload.rs b/fluxer_push/src/payload.rs index ea7d79fc2..b74b72f57 100644 --- a/fluxer_push/src/payload.rs +++ b/fluxer_push/src/payload.rs @@ -14,7 +14,7 @@ const FALLBACK_TITLE: &str = "Fluxer"; const APNS_CATEGORY: &str = "FLUXER_MESSAGE"; const APNS_SOUND: &str = "default"; const APNS_ALERT_EXPIRATION_SECONDS: i64 = 86_400; -const APNS_BACKGROUND_EXPIRATION_SECONDS: i64 = 3_600; +const APNS_BACKGROUND_EXPIRATION_SECONDS: i64 = 86_400; const APNS_COLLAPSE_ID_MAX_BYTES: usize = 64; const SHRUNK_BODY_MAX_BYTES: usize = 40; const MINIMAL_TITLE_MAX_BYTES: usize = 120; @@ -168,7 +168,7 @@ fn fcm_clear_message(device_token: &str, envelope: &Value) -> Value { "data": data, "android": { "priority": "NORMAL", - "ttl": "3600s", + "ttl": "86400s", "collapse_key": format!("clear:{tag}"), }, "fcm_options": {"analytics_label": CLEAR_TYPE}, diff --git a/fluxer_push/src/providers/web_push.rs b/fluxer_push/src/providers/web_push.rs index f2c503e4e..46b50de16 100644 --- a/fluxer_push/src/providers/web_push.rs +++ b/fluxer_push/src/providers/web_push.rs @@ -27,7 +27,7 @@ const MAX_TRANSIENT_RETRIES: u32 = 2; const BASE_RETRY_DELAY_MS: u64 = 200; const MAX_RETRY_DELAY_MS: u64 = 2_000; const ALERT_TTL_SECONDS: &str = "86400"; -const CLEAR_TTL_SECONDS: &str = "3600"; +const CLEAR_TTL_SECONDS: &str = "86400"; const RING_TTL_SECONDS: &str = "0"; const TTL_HEADER: &str = "TTL"; const URGENCY_HEADER: &str = "Urgency"; diff --git a/fluxer_push/src/relay/envelope.rs b/fluxer_push/src/relay/envelope.rs index 3ae808cca..12e950f3c 100644 --- a/fluxer_push/src/relay/envelope.rs +++ b/fluxer_push/src/relay/envelope.rs @@ -6,7 +6,7 @@ use serde_json::{Value, json}; pub const FORMAT_VERSION: u64 = 1; pub const ALERT_TTL_CAP_SECONDS: i64 = 86_400; -pub const BACKGROUND_TTL_CAP_SECONDS: i64 = 3_600; +pub const BACKGROUND_TTL_CAP_SECONDS: i64 = 86_400; pub const APNS_BODY_MAX_BYTES: usize = 4_096; pub const FCM_DATA_MAX_BYTES: usize = 4_096; diff --git a/fluxer_push/src/relay/mod.rs b/fluxer_push/src/relay/mod.rs index ae878057a..2192d551e 100644 --- a/fluxer_push/src/relay/mod.rs +++ b/fluxer_push/src/relay/mod.rs @@ -539,6 +539,29 @@ mod tests { assert_eq!(digest(""), "-"); } + #[test] + fn a_clear_survives_a_device_that_is_offline_for_a_day() { + assert_eq!( + Urgency::Background.ttl_cap_seconds(), + Urgency::Alert.ttl_cap_seconds(), + "a clear must outlive the alert it removes" + ); + } + + #[test] + fn a_clear_stays_silent_on_the_apns_leg() { + let body = envelope::apns_body("payload", Urgency::Background).expect("body fits"); + let parsed: serde_json::Value = serde_json::from_slice(&body).expect("body is json"); + assert_eq!(parsed["aps"]["content-available"], 1); + assert!(parsed["aps"].get("alert").is_none()); + let headers = envelope::apns_headers(Urgency::Background, 0, 86_400); + let push_type = headers + .iter() + .find(|(name, _)| name == "apns-push-type") + .map(|(_, value)| value.as_str()); + assert_eq!(push_type, Some("background")); + } + #[test] fn payload_too_large_answers_413() { assert_eq!(