Fan the unread room notice out from a job

Since the unread rooms stream was scoped per user, posting a message
publishes the room id once to each member of the room. Each publish is a
round trip to Redis, made one after another inside the request, so the time
to post a message grows with the size of the room and in a big room is far
more than the rest of the request.

The fanout now runs in Message::BroadcastUnreadRoomJob, so the poster no
longer waits for it. Who gets the notice and what it says are unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NuTtwJKhb77Dv7EQqv2C3C
This commit is contained in:
Marcello Costagliola
2026-10-05 01:17:59 +02:00
parent 254dd1d46f
commit 3212a683ee
5 changed files with 47 additions and 15 deletions
@@ -0,0 +1,5 @@
class Message::BroadcastUnreadRoomJob < ApplicationJob
def perform(message)
message.broadcast_unread_room
end
end
+15 -9
View File
@@ -1,21 +1,27 @@
module Message::Broadcasts
def broadcast_create
broadcast_append_to room, :messages, target: [ room, :messages ]
broadcast_unread_room
broadcast_unread_room_later
end
def broadcast_remove
broadcast_remove_to room, :messages
end
private
# Fanned out to the room's members rather than published on one global stream, so
# that the timing of activity in a room only reaches people who are in it.
def broadcast_unread_room
payload = ActiveSupport::JSON.encode(roomId: room.id)
# Fanned out to the room's members rather than published on one global stream, so
# that the timing of activity in a room only reaches people who are in it.
def broadcast_unread_room
payload = ActiveSupport::JSON.encode(roomId: room.id)
room.memberships.pluck(:user_id).each do |user_id|
ActionCable.server.broadcast UnreadRoomsChannel.stream_name_for(user_id), payload, coder: nil
end
room.memberships.pluck(:user_id).each do |user_id|
ActionCable.server.broadcast UnreadRoomsChannel.stream_name_for(user_id), payload, coder: nil
end
end
private
# The fanout is one publish per member, which in a big room takes far longer than the
# rest of posting a message, so it runs in a job rather than while the poster waits.
def broadcast_unread_room_later
Message::BroadcastUnreadRoomJob.perform_later(self)
end
end
+6 -2
View File
@@ -16,7 +16,9 @@ class UnreadRoomsChannelTest < ActionCable::Channel::TestCase
assert_not direct.users.include?(users(:jz)), "jz must be an outsider for this test to mean anything"
broadcasts = capture_unread_broadcasts_for(users(:jz)) do
direct.messages.create!(body: "Private", creator: users(:kevin), client_message_id: "outsider").broadcast_create
perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do
direct.messages.create!(body: "Private", creator: users(:kevin), client_message_id: "outsider").broadcast_create
end
end
assert_empty broadcasts
@@ -26,7 +28,9 @@ class UnreadRoomsChannelTest < ActionCable::Channel::TestCase
direct = rooms(:bender_and_kevin)
broadcasts = capture_unread_broadcasts_for(users(:kevin)) do
direct.messages.create!(body: "Private", creator: users(:bender), client_message_id: "member").broadcast_create
perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do
direct.messages.create!(body: "Private", creator: users(:bender), client_message_id: "member").broadcast_create
end
end
assert_equal [ direct.id ], broadcasts.collect { |broadcast| broadcast["roomId"] }
+16 -2
View File
@@ -60,18 +60,32 @@ class MessagesControllerTest < ActionDispatch::IntegrationTest
test "creating a message broadcasts unread room to each member" do
@room.users.each do |member|
assert_broadcasts UnreadRoomsChannel.stream_name_for(member.id), 1 do
post room_messages_url(@room, format: :turbo_stream), params: { message: { body: "New one #{member.id}", client_message_id: member.id } }
perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do
post room_messages_url(@room, format: :turbo_stream), params: { message: { body: "New one #{member.id}", client_message_id: member.id } }
end
end
end
end
test "creating a message leaves the unread fanout to a job" do
member = @room.users.excluding(users(:david)).first
assert_no_broadcasts UnreadRoomsChannel.stream_name_for(member.id) do
post room_messages_url(@room, format: :turbo_stream), params: { message: { body: "New one", client_message_id: 999 } }
end
assert_enqueued_with job: Message::BroadcastUnreadRoomJob, args: [ Message.last ]
end
test "creating a message doesn't broadcast unread room to non-members" do
outsiders = User.where.not(id: @room.users.map(&:id))
assert outsiders.any?, "need someone outside the room for this test to mean anything"
outsiders.each do |outsider|
assert_no_broadcasts UnreadRoomsChannel.stream_name_for(outsider.id) do
post room_messages_url(@room, format: :turbo_stream), params: { message: { body: "New one", client_message_id: 999 } }
perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do
post room_messages_url(@room, format: :turbo_stream), params: { message: { body: "New one", client_message_id: 999 } }
end
end
end
end
+5 -2
View File
@@ -15,8 +15,11 @@ class UnreadRoomsTest < ApplicationSystemTestCase
using_session("Kevin") do
sign_in "kevin@37signals.com"
join_room designers_room
send_message("Hello!!")
send_message("Talking to myself?")
perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do
send_message("Hello!!")
send_message("Talking to myself?")
end
end
assert_room_unread designers_room