Skip to content

refactor(scheduler): run ticks on one daemon thread instead of chained Timers - #1339

Open
davidberenstein1957 wants to merge 2 commits into
masterfrom
scaling/04-scheduler-single-thread
Open

davidberenstein1957 wants to merge 2 commits into
masterfrom
scaling/04-scheduler-single-thread

Conversation

@davidberenstein1957

@davidberenstein1957 davidberenstein1957 commented Aug 12, 2026 •

Copy link
Copy Markdown
Collaborator

Description

PeriodicScheduler chained a fresh threading.Timer per tick and armed the successor before running the payload. This replaces it with a single daemon thread waiting on an Event against an absolute time.monotonic() deadline, with a catch-up guard, in-loop exception logging, and a stop() that sets the event and joins the worker. A wedged (unresponsive) scheduler thread now refuses to let start() double-start a second loop rather than silently creating two concurrent loops.

Related Issue

Part of #1338

Motivation and Context

With chained Timers, a measurement slower than the tick interval re-entered the payload, letting two concurrent invocations mutate _total_energy / _last_measured_time at once and corrupt reported numbers. This is routine for _scheduler_monitor_power (1s, hard-coded) when powermetrics or RAPL is slow. There was also a stop() race where measurements and API pushes could continue after the tracker had written its final row.

Behaviour changes worth flagging: tracker.stop() (and start_task(), which calls it internally) can now block for up to min(measure_power_secs, 5.0) seconds instead of returning immediately, though the normal case is ~0.1 ms — this is a real behaviour change, not a refactor-only tweak, since a wedged in-flight tick can hold stop() for that whole window; a slow callback now yields missing measurements instead of overlapping ones, so the existing "Background scheduler didn't run for a long period" warning will fire where master previously produced corrupted concurrent measurements; and from_run was removed as a parameter since nothing outside the removed _run passed it.

Known limit: the final flush on stop() can still overlap a wedged in-flight tick, since the 5 s join timeout doesn't guarantee the tick has released its locks before stop() proceeds.

How Has This Been Tested?

Measured against master: overlapping entries under a payload 2.4x the interval went from 17 to 0; extra calls after stop() in a gated race repro went from 8 (and never stopping) to 0; drift per tick at 1s interval went from +4.3ms (0.4%) to +0.1ms. New tests/test_scheduler.py, 7 tests, all passing; 4 of them fail against master's scheduler, including test_ticks_run_on_a_single_thread, test_slow_function_is_never_re_entered, and test_start_is_idempotent_and_restartable. test_start_after_a_timed_out_stop_never_runs_two_loops guards the wedged-thread case and fails if the one-line stop() fix is reverted. Full suite: 633 passed, 21 skipped (excluding tests/test_viz_data.py, which fails to import on master too since dash is not installed). uv run pre-commit run --all-files passes.

Screenshots (if appropriate):

N/A

Types of changes

  • Bug fix (non-breaking change which fixes an issue)
  • New feature (non-breaking change which adds functionality)
  • Breaking change (fix or feature that would cause existing functionality to change)

Refactor with a behaviour change on stop() latency (see Motivation and Context) plus tests.

AI Usage Disclosure

  • 🟥 AI-vibecoded
  • 🟠 AI-generated
  • ⭐ AI-assisted
  • ♻️ No AI used

Checklist:

  • My code follows the code style of this project.
  • My change requires a change to the documentation.
  • I have updated the documentation accordingly.
  • I have read the docs/how-to/contributing.md document.
  • I have added tests to cover my changes.
  • All new and existing tests passed.

This should merge before #1340, which builds on this scheduler.

Not addressed here: the final flush on stop() can still overlap a wedged in-flight tick (see Motivation and Context).

@codecov

codecov Bot commented Aug 12, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 91.73%. Comparing base (3ec31a0) to head (61bc10d).
⚠️ Report is 41 commits behind head on master.

Additional details and impacted files
@@            Coverage Diff             @@
##           master    #1339      +/-   ##
==========================================
+ Coverage   91.43%   91.73%   +0.29%     
==========================================
  Files          49       49              
  Lines        5057     5176     +119     
==========================================
+ Hits         4624     4748     +124     
+ Misses        433      428       -5     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@davidberenstein1957

Copy link
Copy Markdown
Collaborator Author

Note for whoever reviews this: #1324 fixes the same scheduler defect by a different route. It is a narrow fix to the existing chained-Timer design, while this PR replaces the design with a single daemon thread. Worth picking one before reviewing either, otherwise they conflict on merge.

…d Timers

Every tick spawned a fresh threading.Timer and the next one was armed
before the callback ran, so a callback slower than the interval overlapped
with itself and the thread count grew with the run.

Run the loop on a single daemon thread waiting on an Event, with an
absolute deadline so the cadence does not drift and an overrun skips ahead
instead of firing catch-up ticks. Each run owns its Event, so a thread left
behind by a timed-out stop() keeps its own set event and exits after its
callback returns, while start() takes effect immediately on a new thread.

Note: stop() now blocks up to min(interval, 5.0) waiting for the in-flight
call. start_task() calls _scheduler.stop() (emissions_tracker.py:755), so
start_task() becomes potentially multi-second where it used to be instant.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@davidberenstein1957
davidberenstein1957 force-pushed the scaling/04-scheduler-single-thread branch from 611879b to 1e90504 Compare August 19, 2026 14:19
@benoit-cty

Copy link
Copy Markdown
Contributor

🤖 This review comment was written and posted by Claude Opus 5.5 (AI assistant), at the request of @benoit-cty. Findings were checked by reading the code and running tests locally (merged with current master where relevant), but please double-check before acting on them.

Verdict: ✅ Approve with nits

Replacing the chained Timers with one daemon thread on a monotonic deadline is a real improvement:

  • no drift
  • overrunning ticks are skipped rather than overlapping
  • exceptions are logged inside the loop
  • the final measure and flush always still run

The 7 new scheduler tests pass, and the PR merges cleanly with master.

Issues:

  1. The "wedged-thread handling" in the description isn't implemented (medium).
    • The description says a timed-out stop() keeps its thread reference, logs a warning, and start() then refuses to run a second loop.
    • The code (scheduler.py:69) always sets _thread = None, logs nothing, and uses one Event per run.
    • So after a stop() that timed out (e.g. via start_task(), emissions_tracker.py:754), start() launches a new loop while the old callback is still running. That gives two concurrent _measure_power_and_energy calls, which is exactly the overlap this PR wants to remove.
    • test_start_after_a_timed_out_stop_leaves_only_the_new_loop explicitly accepts that overlap.
    • Fix: if join() times out, keep the reference and log a warning. In start(), if the previous thread is still alive, wait for it (or refuse) before starting a new loop. Update the test.
  2. The final flush can still overlap an in-flight tick, silently (medium).
    • stop() waits at most min(interval, 5 s), but a tick can make network calls: Electricity Maps (30 s timeout), and the API (up to about 26 s with feat(api): reuse connections, add timeouts, and stamp emissions with measurement time #1340).
    • If one is in flight, tracker.stop() (emissions_tracker.py:909-930) runs the final measurement and _persist_data concurrently with it.
    • At minimum, log a warning when the join times out. Better, make the final measurement wait for the in-flight tick, e.g. with a lock around the tick body.

Nits:

  • Pre-existing: start_task never restarts _scheduler.
  • The tests rely on real sleeps and may flake on shared CI runners. Consider an injectable clock/event.

…unning

stop() keeps the thread reference and logs a warning when the join times
out; start() waits for it and refuses to start a second loop if it is
still busy.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@davidberenstein1957

Copy link
Copy Markdown
Collaborator Author

Made the changes in 61bc10d: timed-out stop() now warns and blocks a second loop.

Not done: making the final flush wait on the in-flight tick (needs a lock shared between scheduler and tracker, bigger change); the pre-existing start_task restart and the real-sleep tests are left as is.

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants