Skip to content

Stop searching for work when stolen timers filled the local queue - #4688

Open
ichraf7 wants to merge 1 commit into
typelevel:series/3.xfrom
ichraf7:fix/searching-worker-nonempty-queue
Open

ichraf7 wants to merge 1 commit into
typelevel:series/3.xfrom
ichraf7:fix/searching-worker-nonempty-queue

Conversation

@ichraf7

@ichraf7 ichraf7 commented Sep 20, 2026

Copy link
Copy Markdown

Add a check that prevent searcher thread from stealling from other thread or from external queue when it queue is no longer empty. The current behavior cause a forver spin of searcher thread when expiring timer fill all the local queue

issue discussed here #issues/4674

@reardonj

Copy link
Copy Markdown
Contributor

@armanbilge , @djspiewak , looks like we have the answer to this question!

We should probably steal both timers and fibers; why only steal one?
-- #4247 (comment)

@reardonj reardonj left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This makes sense to me.

// First try to steal some expired timers.
val stoleTimers = pool.stealTimers(now, rnd)

// Stolen timer callbacks resume fibers, and resumed fibers are

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do stolen timer callbacks always resume fibers? Wondering if we could just check stoleTimers instead of queue.nonEmpty()

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't know if we can have a race condition if we have a cancel that happens in firing moment making it don't resume fiber ( just hypothesis)
I think stoleTimers is a sufficient guard and it is the shape the commit originally had,
both are correct just matter of clarity and taste , I can change it if you want

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The nonEmpty call is likely going to be a bit more expensive since it ends up reading and comparing several AtomicIntegers, instead of the boolean we just computed, so it would be nice to avoid it if possible.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I updated

@ichraf7
ichraf7 force-pushed the fix/searching-worker-nonempty-queue branch 2 times, most recently from 860befe to d899276 Compare September 28, 2026 12:34

@djspiewak djspiewak left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a really good catch, but as written I think it runs the risk of creating permanently asymmetric workloads by always prioritizing timer theft, particularly if an application has a lot of timer expiries that schedule new timers. It may be better to alternate this prioritization, either swapping back and forth (steal fibers first, then timers; steal timers first, then fibers) or by efficiently randomly choosing which one to do first.

@ichraf7

ichraf7 commented Sep 30, 2026

Copy link
Copy Markdown
Author

This is a really good catch, but as written I think it runs the risk of creating permanently asymmetric workloads by always prioritizing timer theft, particularly if an application has a lot of timer expiries that schedule new timers. It may be better to alternate this prioritization, either swapping back and forth (steal fibers first, then timers; steal timers first, then fibers) or by efficiently randomly choosing which one to do first.

Valid point, for having alternating the prio per worker it add state to the worker and more branches, also for app with periodic load having a lot of timers in case of alternation it makes the worker search for work first where there is nothing to steal from other worker queue, the good part is that simple to test, deterministic and costs almost nothing.
For adding randomness , testing will be harder as it not deterministic, I don't know if you're ok with
maybe what we can do keep doing both but limit the number of stolen expiring timer fibers as it done in stealInto, so we have a room for both expiring timers and work stolen from other thread, and anyway fix the hang in enqueueBatch

@ichraf7
ichraf7 force-pushed the fix/searching-worker-nonempty-queue branch 2 times, most recently from 7508c51 to 75c91ee Compare October 1, 2026 21:47
@ichraf7

ichraf7 commented Oct 1, 2026 •

Copy link
Copy Markdown
Author

I checked the git history and my understanding is that enqueueBatch has condition that never changed which is that it's only called when the queue has spare capacity , otherwise it will spin forever.
After some changes this precondition was broken in stealFromOtherWorkerThread , it assumes implicitly that the queue is still empty when it don't stole from other threads, or in our case it stole fibers of the expired timer and made the queue full.
So I think it is adequate to keep the same contract for enqueueBatch , and just before stealing from external queue check that there is enough capacity.
We stay fair and avoid asymmetry by resolving the bug that broke method contract, without adding any complexity of randomness or alternate between both of them. So ignore what I said in my previous suggestion.

`stealFromOtherWorkerThread` falls back to polling the external queue
and enqueueing a batch on the local queue of the searching worker
thread, assuming that this local queue is empty. However, a searching
worker thread steals expired timers before stealing fibers, and the
timer callbacks may resume enough fibers onto its local queue to leave
no room for a batch. `enqueueBatch` then spins forever and block the worker

Add `LocalQueue.hasCapacityForBatch`, and use it in `stealFromOtherWorkerThread`
to skip the external queue when a batch does not fit.
Timers and fibers are still stolen together, to avoid asymmetry

Fixes typelevel#4674
@ichraf7
ichraf7 force-pushed the fix/searching-worker-nonempty-queue branch from 75c91ee to 848d3d9 Compare October 1, 2026 22:02

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants