mirror of
https://github.com/basecamp/once-campfire.git
synced 2026-10-09 00:00:12 +09:00
d485db6038
Each push notification carries the subscriber's unread room count as its badge. The pool built it per subscription, loading the user and counting their unread memberships: two queries for every subscriber, all in the job before the deliveries reach the threads. With 1,000 subscribed members the job spent ~180 ms and 2,000 queries there; with 5,000, a second. The pool now counts the unread rooms of a whole batch with one grouped query and hands each subscription its badge; nothing else in the notification needs the user. The queries still run before the work is posted to the threads, which run outside the Rails executor. Push::Subscription#notification still counts by itself when no badge is given, as for the test notification. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0142qgjggdJ2KDdGk7RF9Xm9
61 lines
2.1 KiB
Ruby
61 lines
2.1 KiB
Ruby
# This is in lib so we can use it in a thread pool without the Rails executor
|
|
class WebPush::Pool
|
|
attr_reader :delivery_pool, :invalidation_pool, :connection, :invalid_subscription_handler
|
|
|
|
def initialize(invalid_subscription_handler:)
|
|
@delivery_pool = Concurrent::ThreadPoolExecutor.new(max_threads: 50, max_queue: 10000)
|
|
@invalidation_pool = Concurrent::FixedThreadPool.new(1)
|
|
@connection = Net::HTTP::Persistent.new(name: "web_push", pool_size: 150)
|
|
@invalid_subscription_handler = invalid_subscription_handler
|
|
end
|
|
|
|
def queue(payload, subscriptions)
|
|
subscriptions.find_in_batches do |batch|
|
|
unread_counts = Membership.unread.where(user_id: batch.map(&:user_id)).group(:user_id).count
|
|
|
|
batch.each do |subscription|
|
|
deliver_later(payload, subscription, badge: unread_counts.fetch(subscription.user_id, 0))
|
|
end
|
|
end
|
|
end
|
|
|
|
def shutdown
|
|
connection.shutdown
|
|
shutdown_pool(delivery_pool)
|
|
shutdown_pool(invalidation_pool)
|
|
end
|
|
|
|
private
|
|
def deliver_later(payload, subscription, badge:)
|
|
# Ensure any AR operations happen before we post to the thread pool
|
|
notification = subscription.notification(**payload, badge: badge)
|
|
subscription_id = subscription.id
|
|
|
|
delivery_pool.post do
|
|
deliver(notification, subscription_id)
|
|
rescue Exception => e
|
|
Rails.logger.error "Error in WebPush::Pool.deliver: #{e.class} #{e.message}"
|
|
end
|
|
rescue Concurrent::RejectedExecutionError
|
|
end
|
|
|
|
def deliver(notification, id)
|
|
notification.deliver(connection: connection)
|
|
rescue WebPush::ExpiredSubscription, OpenSSL::OpenSSLError => ex
|
|
invalidate_subscription_later(id) if invalid_subscription_handler
|
|
end
|
|
|
|
def invalidate_subscription_later(id)
|
|
invalidation_pool.post do
|
|
invalid_subscription_handler.call(id)
|
|
rescue Exception => e
|
|
Rails.logger.error "Error in WebPush::Pool.invalid_subscription_handler: #{e.class} #{e.message}"
|
|
end
|
|
end
|
|
|
|
def shutdown_pool(pool)
|
|
pool.shutdown
|
|
pool.kill unless pool.wait_for_termination(1)
|
|
end
|
|
end
|