diff --git a/CHANGELOG.md b/CHANGELOG.md index 1bd337be..7082dd4c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,9 @@ ### Unreleased +* `Noticed::EventJob` enqueues delivery jobs in batches with `ActiveJob.perform_all_later` (Rails 7.1+) and loads notifications with `find_in_batches`. Note that `perform_all_later` does not run ActiveJob enqueue callbacks. +* [Bugfix] iOS delivery opens a connection per delivery instead of a shared pool, so `error_handler` runs for the notification actually being delivered and notifiers with different APNs credentials no longer share a connection. The `pool_size` option is removed. +* [Bugfix] FCM only treats a 400 as an invalid token when the error refers to the registration token, so malformed payloads no longer delete valid tokens. The access token is fetched once per delivery. +* Request and response bodies are no longer logged by `post_request` since they can contain credentials and access tokens * [Bugfix] Support Rails 8.2's `GlobalID::Locator::RecordNotFound` in `Noticed::Coder` * [Bugfix] Notifier `Notification` classes no longer inherit from a top-level `Notification` model in the host app * [Bugfix] Ephemeral notifiers: bulk delivery methods read their config, `before_enqueue` and `required_params` are honored, `deliver` accepts job options, and `notification_methods` are inherited diff --git a/app/jobs/noticed/event_job.rb b/app/jobs/noticed/event_job.rb index 52228c4d..f8e7f68c 100644 --- a/app/jobs/noticed/event_job.rb +++ b/app/jobs/noticed/event_job.rb @@ -6,11 +6,26 @@ def perform(event) deliver_by.perform_later(event) if deliver_by.perform?(event) end - # Enqueue individual deliveries - event.notifications.each do |notification| - event.delivery_methods.each_value do |deliver_by| - deliver_by.perform_later(notification) if deliver_by.perform?(notification) - end + # Enqueue individual deliveries in batches so large recipient lists don't load into memory all at once + event.notifications.find_in_batches do |notifications| + enqueue_all notifications.flat_map { |notification| delivery_jobs_for(event, notification) } + end + end + + private + + def delivery_jobs_for(event, notification) + event.delivery_methods.values.filter_map do |deliver_by| + deliver_by.job(notification) if deliver_by.perform?(notification) + end + end + + # perform_all_later was added in Rails 7.1 + def enqueue_all(jobs) + if ActiveJob.respond_to?(:perform_all_later) + ActiveJob.perform_all_later(jobs) + else + jobs.each(&:enqueue) end end end diff --git a/app/models/noticed/deliverable/deliver_by.rb b/app/models/noticed/deliverable/deliver_by.rb index 6299f0fd..07bda36c 100644 --- a/app/models/noticed/deliverable/deliver_by.rb +++ b/app/models/noticed/deliverable/deliver_by.rb @@ -18,8 +18,13 @@ def validate! end end + # Builds the delivery job without enqueuing it, so jobs can be enqueued in bulk + def job(event_or_notification, options = {}) + constant.new(name, event_or_notification).set(computed_options(options, event_or_notification)) + end + def perform_later(event_or_notification, options = {}) - constant.set(computed_options(options, event_or_notification)).perform_later(name, event_or_notification) + job(event_or_notification, options).enqueue end # Ephemeral notifiers aren't persisted, so the notifier name, recipient(s), and params are passed to the job instead. diff --git a/docs/delivery_methods/fcm.md b/docs/delivery_methods/fcm.md index 19ae5853..c3d820f9 100644 --- a/docs/delivery_methods/fcm.md +++ b/docs/delivery_methods/fcm.md @@ -139,8 +139,9 @@ end ## Handling Failures Firebase Cloud Messaging Notifications may fail delivery if the user has removed the app from their device. -In this case, FCM will return a [400 or 404 HTTP Status Code](https://firebase.google.com/docs/reference/fcm/rest/v1/ErrorCode) -and the delivery method will call the `invalid_token` handler, if you configure one. +In this case, FCM will return a [404 UNREGISTERED or 400 INVALID_ARGUMENT error](https://firebase.google.com/docs/reference/fcm/rest/v1/ErrorCode) +and the delivery method will call the `invalid_token` handler, if you configure one. A 400 is only treated as an invalid token +when the error refers to the registration token; a malformed payload goes to the `error_handler` instead. ```ruby class CommentNotification diff --git a/docs/delivery_methods/ios.md b/docs/delivery_methods/ios.md index ce5a0c2a..1313ef12 100644 --- a/docs/delivery_methods/ios.md +++ b/docs/delivery_methods/ios.md @@ -62,10 +62,6 @@ end Your APN Team ID -* `pool_size: 5` - *Optional* - - The connection pool size for Apnotic - * `development` - *Optional* Set this to `true` to use the APNS sandbox environment for sending notifications. This is required when running the app to your device via Xcode. Running the app via TestFlight or the App Store should not use development. diff --git a/lib/noticed/api_client.rb b/lib/noticed/api_client.rb index 8ad6ef87..9843c7aa 100644 --- a/lib/noticed/api_client.rb +++ b/lib/noticed/api_client.rb @@ -33,10 +33,10 @@ def post_request(url, args = {}) request.body = body end + # Bodies aren't logged since they can contain credentials and access tokens logger.debug("POST #{url}") - logger.debug(request.body) response = http.request(request) - logger.debug("Response: #{response.code}: #{response.body.inspect}") + logger.debug("Response: #{response.code}") raise ResponseUnsuccessful.new(response, url, args) unless response.code.start_with?("20") diff --git a/lib/noticed/bulk_delivery_methods/bluesky.rb b/lib/noticed/bulk_delivery_methods/bluesky.rb index 9d0ff319..167c28b2 100644 --- a/lib/noticed/bulk_delivery_methods/bluesky.rb +++ b/lib/noticed/bulk_delivery_methods/bluesky.rb @@ -10,7 +10,6 @@ class Bluesky < BulkDeliveryMethod # end def deliver - Rails.logger.debug(evaluate_option(:json)) post_request( "https://#{host}/xrpc/com.atproto.repo.createRecord", headers: {"Authorization" => "Bearer #{token}"}, diff --git a/lib/noticed/bulk_delivery_methods/webhook.rb b/lib/noticed/bulk_delivery_methods/webhook.rb index 04da59e5..42445bd5 100644 --- a/lib/noticed/bulk_delivery_methods/webhook.rb +++ b/lib/noticed/bulk_delivery_methods/webhook.rb @@ -4,7 +4,6 @@ class Webhook < BulkDeliveryMethod required_options :url def deliver - Rails.logger.debug(evaluate_option(:json)) post_request( evaluate_option(:url), basic_auth: evaluate_option(:basic_auth), diff --git a/lib/noticed/delivery_methods/fcm.rb b/lib/noticed/delivery_methods/fcm.rb index c05c631e..a527f04b 100644 --- a/lib/noticed/delivery_methods/fcm.rb +++ b/lib/noticed/delivery_methods/fcm.rb @@ -34,8 +34,17 @@ def format_notification(device_token) end end + # FCM returns 404 UNREGISTERED for tokens that are no longer valid. 400 INVALID_ARGUMENT is also + # returned for malformed payloads, so only treat it as a bad token when the error is about the token. + # https://firebase.google.com/docs/reference/fcm/rest/v1/ErrorCode def bad_token?(response) - response.code == "404" || response.code == "400" + response.code == "404" || (response.code == "400" && token_error?(response)) + end + + def token_error?(response) + JSON.parse(response.body).dig("error", "message").to_s.match?(/registration token/i) + rescue JSON::ParserError + false end def credentials @@ -59,11 +68,14 @@ def load_json(path) end def access_token + @access_token ||= authorizer.fetch_access_token!["access_token"] + end + + def authorizer @authorizer ||= (evaluate_option(:authorizer) || Google::Auth::ServiceAccountCredentials).make_creds( json_key_io: StringIO.new(credentials.to_json), scope: "https://www.googleapis.com/auth/firebase.messaging" ) - @authorizer.fetch_access_token!["access_token"] end end end diff --git a/lib/noticed/delivery_methods/ios.rb b/lib/noticed/delivery_methods/ios.rb index ab3cd95f..6788fb2c 100644 --- a/lib/noticed/delivery_methods/ios.rb +++ b/lib/noticed/delivery_methods/ios.rb @@ -3,29 +3,29 @@ module Noticed module DeliveryMethods class Ios < DeliveryMethod - cattr_accessor :development_connection_pool, :production_connection_pool - required_options :bundle_identifier, :key_id, :team_id, :apns_key, :device_tokens + # A connection is opened per delivery and closed afterwards. Long-lived connections + # get reset by Apple when idle, which stalls the next push for over a minute. def deliver + connection = new_connection + evaluate_option(:device_tokens).each do |device_token| apn = Apnotic::Notification.new(device_token) format_notification(apn) - connection_pool = (!!evaluate_option(:development)) ? development_pool : production_pool - connection_pool.with do |connection| - response = connection.push(apn) - raise "Timeout sending iOS push notification" unless response - connection.close + response = connection.push(apn) + raise "Timeout sending iOS push notification" unless response - if bad_token?(response) && config[:invalid_token] - # Allow notification to cleanup invalid iOS device tokens - notification.instance_exec(device_token, &config[:invalid_token]) - elsif !response.ok? - raise "Request failed #{response.body}" - end + if bad_token?(response) && config[:invalid_token] + # Allow notification to cleanup invalid iOS device tokens + notification.instance_exec(device_token, &config[:invalid_token]) + elsif !response.ok? + raise "Request failed #{response.body}" end end + ensure + connection&.close end private @@ -51,32 +51,16 @@ def bad_token?(response) response.status == "410" || (response.status == "400" && response.body["reason"] == "BadDeviceToken") end - def development_pool - self.class.development_connection_pool ||= new_connection_pool(development: true) - end - - def production_pool - self.class.production_connection_pool ||= new_connection_pool(development: false) - end - - def new_connection_pool(development:) - handler = proc do |connection| - connection.on(:error) do |exception| - Rails.logger.info "Apnotic exception raised: #{exception}" - if config[:error_handler].respond_to?(:call) - notification.instance_exec(exception, &config[:error_handler]) - end - end - end - - if development - Apnotic::ConnectionPool.development(connection_pool_options, pool_options, &handler) - else - Apnotic::ConnectionPool.new(connection_pool_options, pool_options, &handler) + def new_connection + connection = evaluate_option(:development) ? Apnotic::Connection.development(connection_options) : Apnotic::Connection.new(connection_options) + connection.on(:error) do |exception| + Rails.logger.info "Apnotic exception raised: #{exception}" + notification.instance_exec(exception, &config[:error_handler]) if config[:error_handler] end + connection end - def connection_pool_options + def connection_options { auth_method: :token, cert_path: StringIO.new(evaluate_option(:apns_key)), @@ -84,10 +68,6 @@ def connection_pool_options team_id: evaluate_option(:team_id) } end - - def pool_options - {size: evaluate_option(:pool_size) || 5} - end end end end diff --git a/test/delivery_methods/fcm_test.rb b/test/delivery_methods/fcm_test.rb index 8f8bd460..9f1e2788 100644 --- a/test/delivery_methods/fcm_test.rb +++ b/test/delivery_methods/fcm_test.rb @@ -92,12 +92,83 @@ def fetch_access_token! invalid_token: ->(device_token) { cleanups += 1 } ) - stub_request(:post, "https://fcm.googleapis.com/v1/projects/p_1234/messages:send").to_return(status: 400, body: "", headers: {}) + body = { + error: { + code: 400, + message: "The registration token is not a valid FCM registration token", + status: "INVALID_ARGUMENT", + details: [ + {"@type": "type.googleapis.com/google.firebase.fcm.v1.FcmError", errorCode: "INVALID_ARGUMENT"}, + {"@type": "type.googleapis.com/google.rpc.BadRequest", fieldViolations: [{field: "message.token", description: "Invalid registration token"}]} + ] + } + } + stub_request(:post, "https://fcm.googleapis.com/v1/projects/p_1234/messages:send").to_return(status: 400, body: body.to_json, headers: {}) @delivery_method.deliver assert_equal 2, cleanups end + test "invalid payloads are not treated as invalid tokens" do + cleanups = 0 + errors = 0 + + set_config( + authorizer: FakeAuthorizer, + credentials: { + "type" => "service_account", + "project_id" => "p_1234", + "private_key_id" => "private_key" + }, + device_tokens: [:a, :b], + json: ->(device_token) { + { + message: { + token: device_token, + notification: {title: "Title", body: "Body"} + } + } + }, + invalid_token: ->(device_token) { cleanups += 1 }, + error_handler: ->(response) { errors += 1 } + ) + + body = { + error: { + code: 400, + message: "Invalid JSON payload received. Unknown name \"foo\" at 'message.notification': Cannot find field.", + status: "INVALID_ARGUMENT", + details: [{"@type": "type.googleapis.com/google.rpc.BadRequest", fieldViolations: [{field: "message.notification", description: "Invalid JSON payload received."}]}] + } + } + stub_request(:post, "https://fcm.googleapis.com/v1/projects/p_1234/messages:send").to_return(status: 400, body: body.to_json, headers: {}) + + @delivery_method.deliver + assert_equal 0, cleanups + assert_equal 2, errors + end + + test "fetches the access token once per delivery" do + fetches = 0 + authorizer = Class.new(FakeAuthorizer) do + define_method(:fetch_access_token!) do + fetches += 1 + super() + end + end + + set_config( + authorizer: authorizer, + credentials: {"type" => "service_account", "project_id" => "p_1234"}, + device_tokens: [:a, :b], + json: ->(device_token) { {message: {token: device_token}} } + ) + stub_request(:post, "https://fcm.googleapis.com/v1/projects/p_1234/messages:send") + + @delivery_method.deliver + assert_equal 1, fetches + end + test "notifies error handler if exists for other errors" do error_notifications = 0 diff --git a/test/delivery_methods/ios_test.rb b/test/delivery_methods/ios_test.rb index 2a36d8cd..314ddcb9 100644 --- a/test/delivery_methods/ios_test.rb +++ b/test/delivery_methods/ios_test.rb @@ -1,17 +1,18 @@ require "test_helper" class IosTest < ActiveSupport::TestCase - class FakeConnectionPool + class FakeConnection class_attribute :invalid_tokens, default: [] - attr_reader :deliveries + attr_reader :deliveries, :closed, :error_handler - def initialize(response) + def initialize(response = nil) @response = response @deliveries = [] + @closed = false end - def with - yield self + def on(event, &block) + @error_handler = block if event == :error end def push(apn) @@ -20,14 +21,16 @@ def push(apn) end def close + @closed = true end end class FakeResponse - attr_reader :status + attr_reader :status, :body def initialize(status, body = {}) @status = status + @body = body end def ok? @@ -36,7 +39,7 @@ def ok? end setup do - FakeConnectionPool.invalid_tokens = [] + FakeConnection.invalid_tokens = [] @delivery_method = Noticed::DeliveryMethods::Ios.new @delivery_method.instance_variable_set :@notification, noticed_notifications(:one) @@ -51,30 +54,68 @@ def ok? apn.custom_payload = {url: root_url(host: "example.org")} }, invalid_token: ->(device_token) { - FakeConnectionPool.invalid_tokens << device_token + FakeConnection.invalid_tokens << device_token } ) end test "notifies each device token" do - connection_pool = FakeConnectionPool.new(FakeResponse.new("200")) - @delivery_method.stub(:production_pool, connection_pool) do + connection = FakeConnection.new(FakeResponse.new("200")) + @delivery_method.stub(:new_connection, connection) do @delivery_method.deliver end - assert_equal 2, connection_pool.deliveries.count - assert_equal 0, FakeConnectionPool.invalid_tokens.count + assert_equal 2, connection.deliveries.count + assert_equal 0, FakeConnection.invalid_tokens.count end test "notifies of invalid tokens for cleanup" do - connection_pool = FakeConnectionPool.new(FakeResponse.new("410")) - @delivery_method.stub(:production_pool, connection_pool) do + connection = FakeConnection.new(FakeResponse.new("410")) + @delivery_method.stub(:new_connection, connection) do @delivery_method.deliver end - # Our fake connection pool doesn't understand these wouldn't be delivered in the real world - assert_equal 2, connection_pool.deliveries.count - assert_equal 2, FakeConnectionPool.invalid_tokens.count + # Our fake connection doesn't understand these wouldn't be delivered in the real world + assert_equal 2, connection.deliveries.count + assert_equal 2, FakeConnection.invalid_tokens.count + end + + test "closes the connection after delivery" do + connection = FakeConnection.new(FakeResponse.new("200")) + @delivery_method.stub(:new_connection, connection) do + @delivery_method.deliver + end + + assert connection.closed + end + + test "closes the connection when delivery fails" do + connection = FakeConnection.new(FakeResponse.new("500")) + @delivery_method.stub(:new_connection, connection) do + assert_raises(RuntimeError) { @delivery_method.deliver } + end + + assert connection.closed + end + + test "error handler is bound to the notification being delivered" do + handled = nil + set_config( + bundle_identifier: "bundle_id", + key_id: "key_id", + team_id: "team_id", + apns_key: "apns_key", + device_tokens: [], + error_handler: ->(exception) { handled = [self, exception] } + ) + + connection = FakeConnection.new + Apnotic::Connection.stub(:new, connection) do + @delivery_method.send(:new_connection) + end + connection.error_handler.call("boom") + + assert_equal [noticed_notifications(:one), "boom"], handled end private