diff --git a/app/channels/presence_channel.rb b/app/channels/presence_channel.rb index 3ad3bc9..249f6ff 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, at: Time.current.to_fs(:epoch) } + ActionCable.server.broadcast "user_#{current_user.id}_reads", { room_id: membership.room_id, at: Time.current.to_f } end end diff --git a/app/models/message/broadcasts.rb b/app/models/message/broadcasts.rb index 7118f82..e3f8c00 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, at: created_at.to_fs(:epoch)) + payload = ActiveSupport::JSON.encode(roomId: room.id, at: created_at.to_f) 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/lib/web_push/connections.rb b/lib/web_push/connections.rb index df783fe..f151701 100644 --- a/lib/web_push/connections.rb +++ b/lib/web_push/connections.rb @@ -26,8 +26,9 @@ class WebPush::Connections end end - def initialize(keep_alive_timeout: 30) + def initialize(keep_alive_timeout: 30, max_idle: 150) @keep_alive_timeout = keep_alive_timeout + @max_idle = max_idle @idle = Hash.new { |idle, address| idle[address] = [] } @mutex = Mutex.new @pid = Process.pid @@ -85,7 +86,7 @@ class WebPush::Connections def checkin(address, http) @mutex.synchronize do close_expired - if @shut_down + if @shut_down || @idle.values.sum(&:size) >= @max_idle close(http) else @idle[address].push [ http, now ] diff --git a/test/channels/presence_channel_test.rb b/test/channels/presence_channel_test.rb index 881f255..719d26b 100644 --- a/test/channels/presence_channel_test.rb +++ b/test/channels/presence_channel_test.rb @@ -44,7 +44,7 @@ class PresenceChannelTest < ActionCable::Channel::TestCase 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 + assert_broadcast_on "user_#{users(:david).id}_reads", { room_id: membership.room_id, at: Time.current.to_f } do subscribe room_id: membership.room_id end end diff --git a/test/channels/unread_rooms_channel_test.rb b/test/channels/unread_rooms_channel_test.rb index 079234b..e8be2f7 100644 --- a/test/channels/unread_rooms_channel_test.rb +++ b/test/channels/unread_rooms_channel_test.rb @@ -35,7 +35,7 @@ class UnreadRoomsChannelTest < ActionCable::Channel::TestCase end end - assert_equal [ { "roomId" => direct.id, "at" => message.created_at.to_fs(:epoch) } ], broadcasts + assert_equal [ { "roomId" => direct.id, "at" => message.created_at.to_f } ], broadcasts end private diff --git a/test/lib/web_push/connections_test.rb b/test/lib/web_push/connections_test.rb index d465702..b93ced3 100644 --- a/test/lib/web_push/connections_test.rb +++ b/test/lib/web_push/connections_test.rb @@ -135,6 +135,24 @@ class WebPush::ConnectionsTest < ActiveSupport::TestCase end end + test "the idle pool closes connections beyond its bound" do + connections = WebPush::Connections.new(max_idle: 1) + + with_push_service do |server| + other = Server.new + connections.request(pinned_connection(server), push_request) + connections.request(pinned_connection(other), push_request) + assert other.hung_up? + assert_not server.hung_up?(within: 0.1) + + connections.request(pinned_connection(server), push_request) + assert_equal 1, server.connections + ensure + connections.shutdown + other&.stop + end + end + test "idle connections are closed after the keep-alive timeout and on shutdown" do connections = WebPush::Connections.new(keep_alive_timeout: 0) diff --git a/test/system/unread_rooms_test.rb b/test/system/unread_rooms_test.rb index f2164c4..d4b170a 100644 --- a/test/system/unread_rooms_test.rb +++ b/test/system/unread_rooms_test.rb @@ -55,8 +55,23 @@ class UnreadRoomsTest < ApplicationSystemTestCase end end + test "a message after a read in the same second marks the room unread" do + room = rooms(:designers) + user = users(:jz) + join_room rooms(:hq) + read_at = Time.current.change(usec: 100_000) + + broadcast_unread_notice room, to: user, at: read_at - 1 + assert_room_unread room + ActionCable.server.broadcast "user_#{user.id}_reads", { room_id: room.id, at: read_at.to_f } + assert_room_read room + + broadcast_unread_notice room, to: user, at: read_at + 0.2 + assert_room_unread room + 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) } + ActionCable.server.broadcast UnreadRoomsChannel.stream_name_for(to.id), { roomId: room.id, at: at.to_f } end end