Repository navigation
feat: Poll every group of an installation in one process - #44
Draft
stephanschuler wants to merge 21 commits into
Draft
stephanschuler wants to merge 21 commits into
stephanschuler wants to merge 21 commits into
Conversation
Groups were plain booleans below "groupNames". They now are objects below "groups", holding how many jobs run in parallel, how often to poll, how many workers to keep ready and after how long a job counts as abandoned. Every value falls back to a constant of the new Group class, so a group spells out only what deviates. A truthy scalar keeps working and means "enabled, all defaults", so existing registrations need no change. Whether a group is enabled is decided in one place, Group::activeNames(), which the scheduler now asks instead of filtering the settings itself - an array carrying "enabled: false" would have passed its array_filter. The pool is built on first access: freeing stale jobs and counting them for metrics need the configuration but must not start worker processes.
next() takes more than one group now and claims the oldest due job of all of them at once. Every group gets its own UNION ALL branch with its own LIMIT 1, and only those results are sorted afterwards. A "groupname IN (...)" would have been the obvious form but gives up the ordering of idx_for_update: duedate is sorted within one group only, so sorting across groups falls back to a filesort over every due row. PostgreSQL needs one CTE per branch because FOR UPDATE is rejected in a query carrying a UNION. Since the claimed group is unknown until the row has been read, the following select and release filter on the claim value alone - a UUID - and the group name is read from the row. Queries are built in methods instead of constants, abstract where they differ per dialect. The constants held an empty string in the base class for the sole purpose of existing, and a claim spanning a variable number of groups cannot be a constant anyway. Stale jobs are freed with the timeout of their own group, which is why resetStaleJobs() handles one group per statement: a single UPDATE cannot apply a different threshold per row. The deprecated $minutes argument and the global staleJobTimeout setting are gone with it - the crontabs passed --minutes everywhere, so the configured value never applied.
The poll command used to handle a single group, so every group needed its own supervisor program, each of them with numprocs copies polling the same rows. It now resolves all active groups - or those given by --group-names - and drives them from one event loop, with one pool per group so a busy group cannot use up another one's capacity. How many jobs run in parallel and how often to poll therefore belong to the group, not to the command line: --parallel, --prefork-size and --polling-interval-in-seconds are gone. One poll scheduler serves all groups with the smallest configured interval, which satisfies every group since the interval is an upper bound of the waiting time, not a beat. The program shipped with this package no longer names a group and thus covers every group of an installation, so a package adding one needs no program of its own.
The status service carried the same construction as the scheduler did: five constants holding an empty string in the base class, overwritten per dialect. It now builds its queries in methods, abstract where they differ. fetchOne() states that it returns an int, in the base class and in the MySQL override alike - leaving the override untyped would have been a fatal error against the new signature. The service had no tests at all, although it decides what the metrics report. Three of them now cover the buckets, the group boundary and the fact that the border between "running" and "stale" follows the group's own staleJobTimeout.
Static analysis had a standing set of findings in these files. The two factories returned whatever the object manager handed them, and the connection's two read methods declared neither their return nor the shape of their parameter arrays. Both factories also tested "instanceof PostgreSQL94Platform" next to "instanceof PostgreSqlPlatform". The 94 class descends from that base over three steps, so the second test could never add a match - the detection was redundant, not broken. A match expression states the mapping once, and an assert pins the type the object manager returns. Findings that carry behaviour rather than types are deliberately left: the interface hierarchy around ScheduledJobInterface, the serialisation path in ScheduledJob and the retry shape in SchedulingCoordinator each need a decision of their own.
groupname is a CHAR(36) column. MySQL strips the trailing spaces on read, PostgreSQL returns them. Since the claimed job now takes its group from the row instead of the argument, the poll command looked up a padded name in its map of groups, found nothing and crashed on every job.
The helpers of ResetStaleJobsTest and JobStatusServiceTest wrote the activity timestamp with DATE_SUB(NOW(), INTERVAL … SECOND), which PostgreSQL rejects as a syntax error. The platform builds the expression now, so both suites run on either database; the seconds are an int and go into the SQL directly, because an interval literal cannot carry a bound parameter.
A Group carried its configuration and, built lazily, its pool. Because a second instance of the same name would have carried a second pool, groups had to be handed out through a static registry of weak references, with pool-owning groups pinned against collection. Groups are now immutable configuration, created once by the GroupRepository singleton from the settings, which is also the single place that decides whether a group is enabled. The pools belong to the poll command, the only code that runs jobs, and live in a local map there. The registry, the weak references and the pinning are gone, and readers of the configuration can no longer trigger prefork workers by accident.
Since one process polls every group of an installation, its periodic restart affected all of them: once stopPollingAfter had passed, it stopped claiming work and waited for the longest running job of any group, so a short job of another group could sit due for as long as that one took. The process now keeps polling after the timeout and exits at the first moment no job is running. The default grows from ten to thirty minutes, which makes those restarts rarer as well.
Both commands document --group-names as comma separated, but Flow's CLI appends every occurrence of an array option as one string and never splits it: "--group-names=a,b" arrived as the single, unknown group "a,b". The commands now split the values themselves, so a comma separated list and a repeated option both work.
claimed is a CHAR(36) column, and PostgreSQL compares its value including the padding. "LIKE 'failed(%)'" requires the value to end in a parenthesis, so it never matched there: the failed job count was always zero. A prefix pattern matches with and without the padding. MySQL ignores trailing spaces and keeps its queries.
next() always read back the claimed row, even when the claim UPDATE had matched nothing, which is the common case for an idle poller. The affected row count of the UPDATE already tells both apart, so the SELECT now only runs when a row was claimed.
Contributor
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
A critical worker-pool shutdown defect and multiple configuration and polling regressions remain unresolved.
Get a fresh assessment by requesting another Copilot review.
Review effort: Balanced
Findings: 1
Open (4)
What changed in this PR
Adds configurable job groups so one polling process can handle multiple groups with per-group concurrency, preforking, polling, and stale-job settings.
Changes:
- Introduces group configuration and active-group resolution.
- Adds multi-group scheduler and status queries.
- Updates polling, Supervisor configuration, and functional tests.
Required fixes:
- Critical (1 vote): Shut down preforked worker pools when polling stops.
- Moderate (1 vote): Eagerly create pools with nonzero
preforkSize. - Moderate (1 vote): Poll immediately on startup.
- Moderate (3 votes): Respect an empty active-group list instead of enabling
default. - Moderate (2 votes): Validate positive child-process polling intervals before claiming jobs.
- Moderate (1 vote): Preserve or explicitly migrate the previous
staleJobTimeoutsetting. - Nit (2 votes): Correct documentation claiming an empty configuration array enables a group.
| File | Description |
|---|---|
Tests/Functional/JobStatusServiceTest.php |
Tests group-aware status counts. |
Tests/Functional/GroupTest.php |
Tests group configuration and activation. |
Tests/Functional/Groups/ResetStaleJobsTest.php |
Tests per-group stale-job resets. |
Tests/Functional/Groups/MultipleGroupsTest.php |
Tests multi-group job claiming. |
Tests/Functional/Command/GroupNamesOptionTest.php |
Tests CLI group parsing. |
Configuration/Testing/Settings.Groups.yaml |
Adds test group configurations. |
Configuration/Settings.Timeout.yaml |
Removes the global stale timeout. |
Configuration/Settings.Supervisor.yaml |
Configures one multi-group poller. |
Configuration/Settings.Groups.yaml |
Documents per-group settings. |
Classes/Service/PostgreSQLJobStatusService.php |
Adds PostgreSQL group-aware status queries. |
Classes/Service/MySQLJobStatusService.php |
Adds MySQL group-aware status queries. |
Classes/Service/JobStatusServiceFactory.php |
Selects database-specific status services. |
Classes/Service/JobStatusService.php |
Uses per-group stale timeouts. |
Classes/Service/Connection.php |
Adds database helper typing. |
Classes/Domain/SchedulerFactory.php |
Selects database-specific schedulers. |
Classes/Domain/Scheduler.php |
Extends scheduler APIs for groups. |
Classes/Domain/PostgreSQLScheduler.php |
Implements PostgreSQL multi-group claims. |
Classes/Domain/MySQLScheduler.php |
Implements MySQL multi-group claims. |
Classes/Domain/GroupRepository.php |
Resolves active groups. |
Classes/Domain/Group.php |
Models per-group configuration. |
Classes/Domain/AbstractScheduler.php |
Implements group-aware scheduling. |
Classes/Command/SchedulerCommandController.php |
Polls multiple groups using separate pools. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
A prefork worker waiting for a job keeps a periodic timer on the event loop to watch its exit status. Stopping the poll and ping timers therefore never emptied the loop: as soon as a group with preforkSize > 0 had run a job, the command did not return after stopPollingAfter. The drain now terminates every pool's workers as well, which lets the loop run out.
The scheduler kept its own list of active group names and fell back to "default" when the GroupRepository returned none. With every group disabled, the poll command reported that nothing was active and exited, while jobs could still be scheduled into "default" and then waited for a poller that never came. The scheduler now asks the GroupRepository directly, so disabling every group is respected everywhere.
The settings claimed "my-group: true" to be the same as an empty configuration array. The GroupRepository treats [] as falsy, though, and GroupTest expects exactly that, so a group configured that way was silently disabled.
A pollingInterval of zero made the poll loop spin without pause, and since one process polls with the smallest interval of its groups, a single group configured that way affected all of them. A childProcessPollInterval of zero or less was rejected by the Pool, but only once a job had been claimed: the exception ended the poll process on every claim in that group, and each job waited for resetStaleJobs to free it again. Both intervals are now raised to 0.01 seconds, the way parallel and preforkSize are raised to their lower bounds.
A group's pool was built on its first job, so preforkSize had no effect until then: the first job of every group, after each restart of the poll process, still waited for a worker to boot. The pools now exist from the start, one per polled group, and the lookups no longer need to cope with a missing pool.
Member
Author
|
On the items from the review overview:
|
next() claims a job as running = 2 and releases it to running = 1 right after reading it back. A poll process dying in between, or a final failure of the read or the release, left the row at running = 2 for good: the claim only picks running = 0, and resetStaleJobs only freed running = 1. Worse, rescheduling the same identifier keeps claimed and duedate of such a row, so every later run of that job vanished as well. next() passes that state within milliseconds, so a row still in it after staleJobTimeout is orphaned. resetStaleJobs now frees it like a stale running job.
The status service counted every row at running = 2 as pending, so a job that never got past its claim looked like one about to start and never showed up as stale. Once its activity is older than staleJobTimeout, it now counts as stale, matching what resetStaleJobs frees.
An exception thrown in a loop callback ended $loop->run(), so a database outage outlasting the retries of Connection killed the poll process. Its children went on as orphans without anyone writing their activity, and once resetStaleJobs freed them, the restarted poller ran them a second time in parallel. Database errors in loop callbacks are now logged instead, other errors still end the process. A release that fails is repeated until it succeeds, since a lost release lets the job run again. And once the last successful activity of a running job is older than its group's staleJobTimeout, the job's process gets terminated: from then on resetStaleJobs may already have freed it for another run.
Crontabs used to pass --minutes=30 to resetstalejobs, which overrode the configured 60 seconds. Without that option the 60 seconds took effect: a job orphaned by a dying poll process was freed and started again while it still ran, and a database outage of about a minute now terminates every running job. 30 minutes restores the threshold the crontabs had in effect; a group of short jobs can still ask for less.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.



Upgrading to 6.0
One process now polls every group of an installation. Each group carries its own settings, so the configuration moves from command line options and global settings into the group itself.
Configuration
Netlogix.JobQueue.Scheduled.groupNamesNetlogix.JobQueue.Scheduled.groupsNetlogix.JobQueue.Scheduled.staleJobTimeoutstaleJobTimeoutper groupA leftover global
staleJobTimeoutis ignored without notice: every group falls back to 60 seconds. Move the value into each group that needs it.Every group setting is optional and falls back to the constants of
Netlogix\JobQueue\Scheduled\Domain\Group; seeConfiguration/Settings.Groups.yamlfor the full list.trueenables a group with all defaults, a falsy value orenabled: falsedisables it. An empty array counts as falsy.Commands
scheduler:pollforincomingjobs--groupNameis replaced by--groupNames, comma separated. Without it, every active group is polled.--parallel,--preforkSizeand--pollingIntervalInSecondsare gone. Configureparallel,preforkSizeandpollingIntervalper group instead.--stopPollingAfternow defaults to 1800 instead of 600 seconds. Once it has passed, the process keeps polling and exits at the first moment no job is running.scheduler:resetstalejobs--groupNameis replaced by--groupNames, comma separated; all active groups without it.--minutesis gone. Each group uses its ownstaleJobTimeout.--minutesdefaulted to 10 and overruled the configuredstaleJobTimeout, so the command never used the setting. A running job updates its activity every second, so only jobs of dead workers are affected.Supervisor
The shipped program
jobqueue-scheduled-jobsnow starts a single poller for all groups. Remove programs you added per group, or restrict them with--groupNamesif groups must stay in separate processes.Job status
scheduler:resetstalejobsalso frees jobs stuck between claim and release (running = 2) once their activity is older thanstaleJobTimeout. Before, such a job stayed forever and swallowed every later scheduling of its identifier.JobStatusServicecounts these jobs as stale instead of pending once they passstaleJobTimeout. Pending counts may drop, stale counts may rise.PHP API
Scheduler::next()accepts further group names:next(string $groupName, string ...$furtherGroupNames).Scheduler::resetStaleJobs()lost its$minutesparameter.Scheduler::getStaleJobTimeoutSeconds()is gone; useGroupRepository::get($name)->getStaleJobTimeout().AbstractScheduler::injectSettings()is replaced byinjectGroupRepository().🤖 Generated with Claude Code
https://claude.ai/code/session_01Q234jtPC8WNaFaF2o8psYm