@@ -149,4 +149,111 @@ def reveal(secret:)
149149 broadcast_executor &.stop
150150 worker &.stop
151151 end
152+
153+ test "recovers the oldest broadcast across pending and stale work" do
154+ pending_reference = PublicCounterActor . ref ( "pending" ) . async . increment
155+ stale_reference = PublicCounterActor . ref ( "stale" ) . async . increment
156+ worker = SolidObjects ::Worker . new
157+ worker . run_until_idle
158+ stale_process_registry = SolidObjects ::ProcessRegistry . new
159+ stale_process = stale_process_registry . register ( kind : "broadcast" )
160+ now = SolidObjects . database_adapter . database_now
161+ pending_broadcast = SolidObjects ::Broadcast . find_by! ( message_id : pending_reference . id )
162+ pending_broadcast . update! ( available_at : now - 1 . minute )
163+ stale_broadcast = SolidObjects ::Broadcast . find_by! ( message_id : stale_reference . id )
164+ stale_broadcast . update! (
165+ status : "processing" ,
166+ available_at : now - 2 . minutes ,
167+ claimed_by : stale_process . id ,
168+ claimed_at : now - SolidObjects . configuration . process_alive_threshold - 1 . second
169+ )
170+ delivered = Queue . new
171+ SolidObjects . configuration . broadcast_adapter = -> ( broadcast ) { delivered << broadcast . id }
172+ broadcast_executor = SolidObjects ::BroadcastExecutor . new
173+
174+ assert broadcast_executor . run_once
175+
176+ assert_equal stale_broadcast . id , delivered . pop
177+ assert_equal "delivered" , stale_broadcast . reload . status
178+ assert_equal "pending" , pending_broadcast . reload . status
179+ ensure
180+ broadcast_executor &.stop
181+ stale_process_registry &.stop
182+ worker &.stop
183+ end
184+
185+ test "concurrent executors claim different broadcasts" do
186+ pending_reference = PublicCounterActor . ref ( "pending" ) . async . increment
187+ stale_reference = PublicCounterActor . ref ( "stale" ) . async . increment
188+ worker = SolidObjects ::Worker . new
189+ worker . run_until_idle
190+ stale_process_registry = SolidObjects ::ProcessRegistry . new
191+ stale_process = stale_process_registry . register ( kind : "broadcast" )
192+ now = SolidObjects . database_adapter . database_now
193+ pending_broadcast = SolidObjects ::Broadcast . find_by! ( message_id : pending_reference . id )
194+ pending_broadcast . update! ( available_at : now - 1 . minute )
195+ stale_broadcast = SolidObjects ::Broadcast . find_by! ( message_id : stale_reference . id )
196+ stale_broadcast . update! (
197+ status : "processing" ,
198+ available_at : now - 2 . minutes ,
199+ claimed_by : stale_process . id ,
200+ claimed_at : now - SolidObjects . configuration . process_alive_threshold - 1 . second
201+ )
202+ claims = Queue . new
203+ release = Queue . new
204+ SolidObjects . configuration . broadcast_adapter = lambda do |broadcast |
205+ claims << broadcast . id
206+ release . pop
207+ end
208+ executor_a = SolidObjects ::BroadcastExecutor . new
209+ executor_b = SolidObjects ::BroadcastExecutor . new
210+
211+ thread_a = Thread . new { executor_a . run_once }
212+ assert_equal stale_broadcast . id , Timeout . timeout ( 5 ) { claims . pop }
213+ thread_b = Thread . new { executor_b . run_once }
214+ assert_equal pending_broadcast . id , Timeout . timeout ( 5 ) { claims . pop }
215+ 2 . times { release << true }
216+
217+ assert thread_a . value
218+ assert thread_b . value
219+ assert_equal %w[ delivered delivered ] , SolidObjects ::Broadcast . order ( :id ) . pluck ( :status )
220+ ensure
221+ 2 . times { release << true } if release
222+ thread_a &.join ( 2 )
223+ thread_b &.join ( 2 )
224+ executor_a &.stop
225+ executor_b &.stop
226+ stale_process_registry &.stop
227+ worker &.stop
228+ end
229+
230+ test "polls pending and stale broadcasts separately" do
231+ PublicCounterActor . ref ( "one" ) . async . increment
232+ worker = SolidObjects ::Worker . new
233+ worker . run_until_idle
234+ SolidObjects . configuration . broadcast_adapter = -> ( broadcast ) { broadcast }
235+ broadcast_executor = SolidObjects ::BroadcastExecutor . new
236+ polling_queries = [ ]
237+ subscription = ActiveSupport ::Notifications . subscribe ( "sql.active_record" ) do |event |
238+ query = event . payload . fetch ( :sql ) . squish
239+ if query . match? ( /SELECT .* FROM ["`]solid_objects_broadcasts["`]/ ) &&
240+ query . match? ( /ORDER BY .*available_at.*id.*LIMIT/i )
241+ polling_queries << query
242+ end
243+ end
244+
245+ assert broadcast_executor . run_once
246+
247+ assert_equal 2 , polling_queries . length
248+ assert polling_queries . one? { |query | !query . include? ( "claimed_at" ) }
249+ assert polling_queries . one? { |query | query . include? ( "claimed_at" ) }
250+ polling_queries . each { |query | refute_match ( /\s OR\s /i , query ) }
251+ if database_family != :sqlite
252+ polling_queries . each { |query | assert_match ( /FOR UPDATE SKIP LOCKED\z /i , query ) }
253+ end
254+ ensure
255+ ActiveSupport ::Notifications . unsubscribe ( subscription ) if subscription
256+ broadcast_executor &.stop
257+ worker &.stop
258+ end
152259end
0 commit comments