Merge pull request #296: Fan the unread room notice out from a job

Reviewed and merged by GPT on behalf of DHH.
This commit is contained in:
GPT on behalf of DHH
2026-10-07 10:33:57 +02:00
9 changed files with 127 additions and 22 deletions
+1 -1
View File
@@ -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
@@ -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 } })
}
}
@@ -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)
}
@@ -0,0 +1,5 @@
class Message::BroadcastUnreadRoomJob < ApplicationJob
def perform(message)
message.broadcast_unread_room
end
end
+17 -9
View File
@@ -1,21 +1,29 @@
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. 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, at: created_at.to_fs(:epoch))
room.memberships.pluck(:user_id).each do |user_id|
ActionCable.server.broadcast UnreadRoomsChannel.stream_name_for(user_id), payload, coder: nil
end
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
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
+10
View File
@@ -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
+9 -3
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
@@ -24,12 +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
direct.messages.create!(body: "Private", creator: users(:bender), client_message_id: "member").broadcast_create
perform_enqueued_jobs only: Message::BroadcastUnreadRoomJob do
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
+32 -2
View File
@@ -70,20 +70,50 @@ 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
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 "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"
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
+37 -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
@@ -24,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