Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
25 changes: 20 additions & 5 deletions app/jobs/noticed/event_job.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 6 additions & 1 deletion app/models/noticed/deliverable/deliver_by.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
5 changes: 3 additions & 2 deletions docs/delivery_methods/fcm.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 0 additions & 4 deletions docs/delivery_methods/ios.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
4 changes: 2 additions & 2 deletions lib/noticed/api_client.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down
1 change: 0 additions & 1 deletion lib/noticed/bulk_delivery_methods/bluesky.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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}"},
Expand Down
1 change: 0 additions & 1 deletion lib/noticed/bulk_delivery_methods/webhook.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
16 changes: 14 additions & 2 deletions lib/noticed/delivery_methods/fcm.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
60 changes: 20 additions & 40 deletions lib/noticed/delivery_methods/ios.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -51,43 +51,23 @@ 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)),
key_id: evaluate_option(:key_id),
team_id: evaluate_option(:team_id)
}
end

def pool_options
{size: evaluate_option(:pool_size) || 5}
end
end
end
end
Expand Down
73 changes: 72 additions & 1 deletion test/delivery_methods/fcm_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Loading
Loading