From 3212a683ee3d65b297c84293ce2200a4ee246895 Mon Sep 17 00:00:00 2001 From: Marcello Costagliola Date: Mon, 5 Oct 2026 01:17:59 +0200 Subject: [PATCH 1/3] 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 Claude-Session: https://claude.ai/code/session_01NuTtwJKhb77Dv7EQqv2C3C --- app/jobs/message/broadcast_unread_room_job.rb | 5 ++++ app/models/message/broadcasts.rb | 24 ++++++++++++------- test/channels/unread_rooms_channel_test.rb | 8 +++++-- test/controllers/messages_controller_test.rb | 18 ++++++++++++-- test/system/unread_rooms_test.rb | 7 ++++-- 5 files changed, 47 insertions(+), 15 deletions(-) create mode 100644 app/jobs/message/broadcast_unread_room_job.rb diff --git a/app/jobs/message/broadcast_unread_room_job.rb b/app/jobs/message/broadcast_unread_room_job.rb new file mode 100644 index 0000000..4b0bc2c --- /dev/null +++ b/app/jobs/message/broadcast_unread_room_job.rb @@ -0,0 +1,5 @@ +class Message::BroadcastUnreadRoomJob < ApplicationJob + def perform(message) + message.broadcast_unread_room + end +end diff --git a/app/models/message/broadcasts.rb b/app/models/message/broadcasts.rb index f82e4e3..b9660c9 100644 --- a/app/models/message/broadcasts.rb +++ b/app/models/message/broadcasts.rb @@ -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 diff --git a/test/channels/unread_rooms_channel_test.rb b/test/channels/unread_rooms_channel_test.rb index 4444eb1..b3f36ed 100644 --- a/test/channels/unread_rooms_channel_test.rb +++ b/test/channels/unread_rooms_channel_test.rb @@ -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"] } diff --git a/test/controllers/messages_controller_test.rb b/test/controllers/messages_controller_test.rb index 8088694..63266cb 100644 --- a/test/controllers/messages_controller_test.rb +++ b/test/controllers/messages_controller_test.rb @@ -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 diff --git a/test/system/unread_rooms_test.rb b/test/system/unread_rooms_test.rb index 235e912..ea68ba0 100644 --- a/test/system/unread_rooms_test.rb +++ b/test/system/unread_rooms_test.rb @@ -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 From 506c2771fe0900632d39652cdf1422afc435a78a Mon Sep 17 00:00:00 2001 From: Marcello Costagliola Date: Mon, 5 Oct 2026 02:08:01 +0200 Subject: [PATCH 2/3] Skip the unread notice for members who have caught up Now that the unread fanout runs in a job, it can run after a member has already opened the room and moved on to another one. Their sidebar would then mark the room unread for a message they have seen. The job now notifies only members who still have the room unread or are in it now. Room#receive marks members who aren't in the room unread, and opening the room clears it, so a member who caught up in the meantime is skipped. When the job runs right away this is the same set of people as before, except members who have hidden the room, whose sidebar doesn't list it. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01JVFo3Lt9T8M5NR7KxVvsZ2 --- app/models/message/broadcasts.rb | 6 ++++-- test/controllers/messages_controller_test.rb | 16 ++++++++++++++++ 2 files changed, 20 insertions(+), 2 deletions(-) diff --git a/app/models/message/broadcasts.rb b/app/models/message/broadcasts.rb index b9660c9..6ae7569 100644 --- a/app/models/message/broadcasts.rb +++ b/app/models/message/broadcasts.rb @@ -9,11 +9,13 @@ module Message::Broadcasts end # 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. + # that the timing of activity in a room only reaches people who are in it. Of those, + # only the ones who still have it unread or are in it now: by the time this runs, + # someone may have opened the room and moved on, and would see it marked unread again. def broadcast_unread_room payload = ActiveSupport::JSON.encode(roomId: room.id) - room.memberships.pluck(:user_id).each do |user_id| + room.memberships.unread.or(room.memberships.connected).pluck(:user_id).each do |user_id| ActionCable.server.broadcast UnreadRoomsChannel.stream_name_for(user_id), payload, coder: nil end end diff --git a/test/controllers/messages_controller_test.rb b/test/controllers/messages_controller_test.rb index 63266cb..83cedda 100644 --- a/test/controllers/messages_controller_test.rb +++ b/test/controllers/messages_controller_test.rb @@ -58,6 +58,8 @@ class MessagesControllerTest < ActionDispatch::IntegrationTest end test "creating a message broadcasts unread room to each member" do + memberships(:david_watercooler).present # the poster is in the room + @room.users.each do |member| assert_broadcasts UnreadRoomsChannel.stream_name_for(member.id), 1 do perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do @@ -77,6 +79,20 @@ class MessagesControllerTest < ActionDispatch::IntegrationTest assert_enqueued_with job: Message::BroadcastUnreadRoomJob, args: [ Message.last ] end + test "the unread fanout skips a member who has opened the room since" do + membership = memberships(:jason_watercooler) + + post room_messages_url(@room, format: :turbo_stream), params: { message: { body: "New one", client_message_id: 999 } } + assert membership.reload.unread? + + membership.present + membership.reload.disconnected + + assert_no_broadcasts UnreadRoomsChannel.stream_name_for(membership.user_id) do + perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob + end + 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" From a740581f208775dd46238969cd246fb9e257b777 Mon Sep 17 00:00:00 2001 From: Marcello Costagliola Date: Mon, 5 Oct 2026 16:11:37 +0200 Subject: [PATCH 3/3] Keep an unread notice older than the latest read from marking the room The job picks the room's members once and then publishes to them one at a time, as the request did before it. A member who opens the room during that loop gets the read event in their other tabs and then the older notice, which marked the room unread again. Both events now carry the server time, and the sidebar ignores a notice dated before the room's latest read. The notice still moves a direct room to the top. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_0142qgjggdJ2KDdGk7RF9Xm9 --- app/channels/presence_channel.rb | 2 +- .../controllers/read_rooms_controller.js | 4 +-- .../controllers/rooms_list_controller.js | 17 ++++++++-- app/models/message/broadcasts.rb | 2 +- test/channels/presence_channel_test.rb | 10 ++++++ test/channels/unread_rooms_channel_test.rb | 6 ++-- test/system/unread_rooms_test.rb | 32 +++++++++++++++++++ 7 files changed, 64 insertions(+), 9 deletions(-) diff --git a/app/channels/presence_channel.rb b/app/channels/presence_channel.rb index 9c49c13..3ad3bc9 100644 --- a/app/channels/presence_channel.rb +++ b/app/channels/presence_channel.rb @@ -22,6 +22,6 @@ class PresenceChannel < RoomChannel end def broadcast_read_room - ActionCable.server.broadcast "user_#{current_user.id}_reads", { room_id: membership.room_id } + ActionCable.server.broadcast "user_#{current_user.id}_reads", { room_id: membership.room_id, at: Time.current.to_fs(:epoch) } end end diff --git a/app/javascript/controllers/read_rooms_controller.js b/app/javascript/controllers/read_rooms_controller.js index 422808d..ca8f6e2 100644 --- a/app/javascript/controllers/read_rooms_controller.js +++ b/app/javascript/controllers/read_rooms_controller.js @@ -16,7 +16,7 @@ export default class extends Controller { }) } - #read = ({ room_id }) => { - this.dispatch("read", { detail: { roomId: room_id } }) + #read = ({ room_id, at }) => { + this.dispatch("read", { detail: { roomId: room_id, at } }) } } diff --git a/app/javascript/controllers/rooms_list_controller.js b/app/javascript/controllers/rooms_list_controller.js index afd9981..4868904 100644 --- a/app/javascript/controllers/rooms_list_controller.js +++ b/app/javascript/controllers/rooms_list_controller.js @@ -7,6 +7,7 @@ export default class extends Controller { static classes = [ "unread" ] #disconnected = true + #readAt = new Map() async connect() { this.channel ??= await cable.subscribeTo({ channel: "UnreadRoomsChannel" }, { @@ -27,9 +28,13 @@ export default class extends Controller { this.read({ detail: { roomId: Current.room.id } }) } - read({ detail: { roomId } }) { + read({ detail: { roomId, at } }) { const room = this.#findRoomTarget(roomId) + if (at) { + this.#readAt.set(Number(roomId), Math.max(Number(at), this.#readAt.get(Number(roomId)) ?? 0)) + } + if (room) { room.classList.remove(this.unreadClass) this.dispatch("read", { detail: { targetId: roomId } }) @@ -47,11 +52,11 @@ export default class extends Controller { this.#disconnected = true } - #unread({ roomId }) { + #unread({ roomId, at }) { const unreadRoom = this.#findRoomTarget(roomId) if (unreadRoom) { - if (Current.room.id != roomId) { + if (Current.room.id != roomId && !this.#readSince(roomId, at)) { unreadRoom.classList.add(this.unreadClass) } @@ -59,6 +64,12 @@ export default class extends Controller { } } + // Notices fan out one member at a time, so one can arrive after the member has already + // read the room in another tab. It still reorders the room, but doesn't mark it unread. + #readSince(roomId, at) { + return Number(at) <= this.#readAt.get(Number(roomId)) + } + #findRoomTarget(roomId) { return this.roomTargets.find(roomTarget => roomTarget.dataset.roomId == roomId) } diff --git a/app/models/message/broadcasts.rb b/app/models/message/broadcasts.rb index 6ae7569..7118f82 100644 --- a/app/models/message/broadcasts.rb +++ b/app/models/message/broadcasts.rb @@ -13,7 +13,7 @@ module Message::Broadcasts # only the ones who still have it unread or are in it now: by the time this runs, # someone may have opened the room and moved on, and would see it marked unread again. def broadcast_unread_room - payload = ActiveSupport::JSON.encode(roomId: room.id) + payload = ActiveSupport::JSON.encode(roomId: room.id, at: created_at.to_fs(:epoch)) room.memberships.unread.or(room.memberships.connected).pluck(:user_id).each do |user_id| ActionCable.server.broadcast UnreadRoomsChannel.stream_name_for(user_id), payload, coder: nil diff --git a/test/channels/presence_channel_test.rb b/test/channels/presence_channel_test.rb index 99227a6..881f255 100644 --- a/test/channels/presence_channel_test.rb +++ b/test/channels/presence_channel_test.rb @@ -40,6 +40,16 @@ class PresenceChannelTest < ActionCable::Channel::TestCase end end + test "subscribing tells the user's other tabs when they read the room" do + membership = users(:david).memberships.first + + freeze_time do + assert_broadcast_on "user_#{users(:david).id}_reads", { room_id: membership.room_id, at: Time.current.to_fs(:epoch) } do + subscribe room_id: membership.room_id + end + end + end + test "unsubscribing marks the membership as disconnected" do membership = users(:david).memberships.first subscribe room_id: membership.room_id diff --git a/test/channels/unread_rooms_channel_test.rb b/test/channels/unread_rooms_channel_test.rb index b3f36ed..48ca520 100644 --- a/test/channels/unread_rooms_channel_test.rb +++ b/test/channels/unread_rooms_channel_test.rb @@ -26,14 +26,16 @@ class UnreadRoomsChannelTest < ActionCable::Channel::TestCase test "a member is told about activity in their own room" do direct = rooms(:bender_and_kevin) + message = nil broadcasts = capture_unread_broadcasts_for(users(:kevin)) do perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do - direct.messages.create!(body: "Private", creator: users(:bender), client_message_id: "member").broadcast_create + message = direct.messages.create!(body: "Private", creator: users(:bender), client_message_id: "member") + message.broadcast_create end end - assert_equal [ direct.id ], broadcasts.collect { |broadcast| broadcast["roomId"] } + assert_equal [ { "roomId" => direct.id, "at" => message.created_at.to_fs(:epoch) } ], broadcasts end private diff --git a/test/system/unread_rooms_test.rb b/test/system/unread_rooms_test.rb index ea68ba0..f2164c4 100644 --- a/test/system/unread_rooms_test.rb +++ b/test/system/unread_rooms_test.rb @@ -27,4 +27,36 @@ class UnreadRoomsTest < ApplicationSystemTestCase join_room designers_room assert_room_read designers_room end + + test "a notice about a message from before the room was read in another tab doesn't mark it unread again" do + using_session("Kevin in HQ") do + sign_in "kevin@37signals.com" + join_room rooms(:hq) + broadcast_unread_notice rooms(:designers), to: users(:kevin), at: Time.current + assert_room_unread rooms(:designers) + end + + sent_before_reading = Time.current + using_session("Kevin in Designers") do + sign_in "kevin@37signals.com" + join_room rooms(:designers) + end + + using_session("Kevin in HQ") do + assert_room_read rooms(:designers) + + broadcast_unread_notice rooms(:designers), to: users(:kevin), at: sent_before_reading + broadcast_unread_notice rooms(:bender_and_kevin), to: users(:kevin), at: Time.current + assert_selector "#" + dom_id(rooms(:bender_and_kevin), :list) + ".unread" + assert_room_read rooms(:designers) + + broadcast_unread_notice rooms(:designers), to: users(:kevin), at: Time.current + assert_room_unread rooms(:designers) + end + end + + private + def broadcast_unread_notice(room, to:, at:) + ActionCable.server.broadcast UnreadRoomsChannel.stream_name_for(to.id), { roomId: room.id, at: at.to_fs(:epoch) } + end end