Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 22 additions & 5 deletions .github/workflows/phpunit.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,30 +7,47 @@ jobs:
runs-on: ubuntu-latest

strategy:
fail-fast: false
matrix:
php: [ '8.1', '8.2', '8.3', '8.4', '8.5' ]
symfony: [ '6.4.*', '7.4.*' ]
doctrine: [ 'latest' ]
exclude:
- php: 8.1
symfony: '7.4.*'
include:
# JobOutputWriter and the test support branch between DBAL 3 and 4 - keep the old majors covered
- php: '8.2'
symfony: '6.4.*'
doctrine: 'dbal3'
- php: '8.4'
symfony: '7.4.*'
doctrine: 'dbal3'

name: PHP ${{ matrix.php }} / Symfony ${{ matrix.symfony }}
name: PHP ${{ matrix.php }} / Symfony ${{ matrix.symfony }} / Doctrine ${{ matrix.doctrine }}

steps:
- uses: actions/checkout@v3
- uses: actions/checkout@v4

- name: Setup PHP
uses: shivammathur/setup-php@v2
with:
php-version: ${{ matrix.php }}
extensions: mbstring, intl, pdo_mysql
extensions: mbstring, intl, pdo_mysql, pdo_sqlite
ini-values: post_max_size=256M, upload_max_filesize=256M
# SYMFONY_REQUIRE is applied by symfony/flex only
tools: flex

- name: Set Symfony Version
run: echo "SYMFONY_REQUIRE=${{ matrix.symfony }}" >> $GITHUB_ENV

- name: Pin DBAL 3 / ORM 2
if: matrix.doctrine == 'dbal3'
run: composer require --no-update --no-interaction "doctrine/dbal:^3.9" "doctrine/orm:^2.20" "doctrine/doctrine-bundle:^2.13"

- name: Install Dependencies
run: composer install --prefer-dist --no-progress
# composer.lock is not committed - resolve for the PHP version and the pinned Symfony version of this leg
run: composer update --prefer-dist --no-progress --no-interaction

- name: Run PHPUnit Tests
run: ./vendor/bin/phpunit
run: ./vendor/bin/phpunit
54 changes: 54 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
# Changelog

## Unreleased (2.2.0)

### Fixed

- The job message handler no longer holds database transactions while the command runs. Output, status and result are
written through DBAL in autocommit mode, so a connection dropped in the middle of a write can no longer leave a
half-open transaction on the shared connection. Previously the worker crashed afterwards in the messenger retry with
`There is already an active transaction` and the message was redelivered (and the command run again) every
`redeliver_timeout`.
- The wait loop is throttled (100 ms sleep, database writes and the cancellation check at most once per
`poll_interval_ms`) instead of a busy loop which reloaded the whole job, including its output, and rewrote the whole
output on every iteration.
- A short database outage no longer kills the running command; the output is buffered and the write retried.
- A failing command is stored as a `failed` job and no longer makes the handler throw. Storing the result is retried;
when it still fails, an `UnrecoverableMessageHandlingException` is thrown so the command is not run again by a retry.
- A message of a deleted job is rejected with `UnrecoverableMessageHandlingException` instead of a `TypeError`.
- Output params printed in several output chunks are all kept (previously only those of the last chunk), and a
`OUTPUT PARAMS:` line split between two chunks is found.
- A job cancelled while it was finishing keeps the `cancelled` status.
- Output containing multi-byte characters can no longer fail to be stored on a strict MySQL (`Incorrect string value`)
when a read or the output cap splits a character; output a working database rejects is replaced by a note and can
no longer kill a healthy command or block storing the job result.
- A database error before the command starts no longer loses the message; a job deleted while running stops its
command.
- Symfony Process no longer keeps the whole command output in a temporary file.

### Added

- Configuration `job_queue.processing`: `poll_interval_ms` (1000), `output_max_bytes` (4 MB, 0 = unlimited),
`db_failure_tolerance` (30), `rerun_on_redelivery` (false).
- Integration tests with SQLite and real subprocesses; CI legs with DBAL 3 / ORM 2.

### Changed

- **Behaviour change:** a message redelivered for a job which is already `running` (its previous worker died) marks the
job `failed` instead of running the command again. Set `rerun_on_redelivery: true` for the previous behaviour.
Messages for `completed`, `failed` and `cancelled` jobs are ignored.
- **Behaviour change:** the stored job output is capped at `output_max_bytes` (4 MB by default, previously unlimited);
a truncation marker is appended once.
- **Behaviour change:** output params are written once, when the job finishes (previously during the run), and the
values of all chunks are kept (previously only those of the last chunk).
- **Behaviour change:** the stored output is valid UTF-8 - bytes which are not UTF-8 are replaced by `?`.
- **Behaviour change:** the `Job` entity is detached from the entity manager while the handler runs and is not
refreshed afterwards. When the handler runs inside a unit of work (sync transport, direct invocation), a `Job`
object loaded before is stale (still `planned`) and detached - re-load it with `find()` before using it, e.g. as
`$parentJob` of `createCommandJob()`.
- The job is switched to `running` by an atomic claim before the command starts (previously right after it started).
- A job is cancelled by a conditional update (only while it is `running`), so a result stored by the handler in the
meantime is no longer overwritten by the cancellation.
- `symfony/dependency-injection` 8 is declared as a conflict (the bundle loads `config/services.xml`).
- Recommended messenger setup for the job transport: `retry_strategy: { max_retries: 0 }`, `--keepalive` (or a higher
`redeliver_timeout`) for jobs longer than an hour. See README.
81 changes: 78 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,44 @@ framework:
TomAtom\JobQueueBundle\Message\JobMessage: job_message # or async
```

Recommendations for the transport of the job messages:

- **Disable retries** - `retry_strategy: { max_retries: 0 }`. A retry of a job message means running the command
again. The handler itself never throws because a command failed (a failed command is stored as a `failed` job), the
only exception it throws once the command has started is an `UnrecoverableMessageHandlingException` when the job
result cannot be stored. A database error before the command starts (loading or claiming the job) is retried by the
handler a few times and then thrown as a `RecoverableMessageHandlingException`, which Messenger retries even with
`max_retries: 0` - nothing has run yet, so the job is not lost. A message of a job which does not exist (also after a
short wait for the creating transaction) is rejected without retry.
- **Long jobs and several consumers** - the stock Doctrine transport redelivers a message which is not acknowledged
within `redeliver_timeout` (default 3600 s). For jobs running longer than that, run the consumer with `--keepalive`
(Symfony >= 7.2, `messenger:consume job_message --keepalive`) or raise `redeliver_timeout`. A redelivered message of
a job which is already `running` does not run the command again (see `rerun_on_redelivery` below).
**Without `--keepalive`, a job running longer than `redeliver_timeout` is redelivered to another consumer while the
first one still runs it.** The handler cannot tell this from a dead worker: the job is marked `failed` with
"previous worker died" (the first worker later overwrites the status with its result, the message stays in the
output), and with `rerun_on_redelivery: true` the command runs twice at the same time.
- **Own DBAL connection (optional)** - the handler keeps no transaction open while the command runs, so sharing the
default connection is fine. If you want to be completely sure the handler's connection state never meets the
transport's send/reject, use a separate connection to the same database, e.g. `dsn: "doctrine://jobs"` with a
`doctrine.dbal.connections.jobs` entry.
- Do not wrap the job handler in the `doctrine_transaction` middleware - the job state must be visible (and the
cancellation readable) while the command runs. (A transaction opened by the caller is not rolled back by the
handler's connection resets, but with a lost connection it is gone anyway.)

```yaml
framework:
messenger:
transports:
job_message:
dsn: "%env(MESSENGER_TRANSPORT_DSN)%"
retry_strategy:
max_retries: 0
options:
queue_name: job_message
# redeliver_timeout: 3600 # raise for jobs longer than an hour if you do not use --keepalive
```

<hr>

#### config/packages/security.yaml:
Expand Down Expand Up @@ -140,8 +178,42 @@ job_queue:
job_recurring_table_name: "your_job_recurring_table_name" # Default = job_recurring_queue
scheduling:
heartbeat_interval: "1 hour" # Default = 1 minute
processing:
poll_interval_ms: 1000 # Default = 1000 - how often the output is written and the cancellation checked
output_max_bytes: 4194304 # Default = 4 MB - cap of the stored output, 0 = unlimited (not recommended)
db_failure_tolerance: 30 # Default = 30 - consecutive failed database polls (~ seconds) before the job is given up
rerun_on_redelivery: false # Default = false - run the command again when a message of a RUNNING job is redelivered
```

How a job is processed:

- The job is claimed atomically (`planned` -> `running`), so a duplicate message never runs the command twice.
- While the command runs, its output is appended to the job and the cancellation is checked at most once per
`poll_interval_ms`. All these writes run in autocommit mode through DBAL - no transaction is ever held open by the
handler, and the `Job` entity is not managed by the entity manager during the run.
- A short database outage does not stop the command: the output stays buffered in memory and the write is retried.
After `db_failure_tolerance` consecutive failed polls the command is stopped and the job is marked `failed`.
- Output above `output_max_bytes` is dropped and a `[... output truncated by JobQueueBundle at N bytes ...]` marker is
appended once. **Keep `output_max_bytes` well below MySQL `max_allowed_packet`** (64 MB by default on MySQL 8) -
the `output` column is a LONGTEXT and a larger value cannot be written or read in one packet. With
`output_max_bytes: 0` nothing guards the column size, not even the final status message.
- The stored output is always valid UTF-8 (a strict MySQL rejects anything else in a `utf8mb4` column): a character
split between two reads is completed first, the cap never cuts a character, and bytes which are not UTF-8 (binary
output, another encoding) are replaced by `?`.
- Output which the database rejects although the connection works (e.g. a value or packet size error) is replaced by
a `[JobQueueBundle: N bytes of output could not be stored: ...]` note - it never fails a healthy command or blocks
storing the job result.
- A job deleted while its command runs stops the command (like a cancellation); nothing is recorded.
- When a message is redelivered for a job which is already `running` (the previous worker died), the job is marked
`failed` with an explanation instead of running the command again. Set `rerun_on_redelivery: true` to run it again.
Messages of `completed`, `failed` or `cancelled` jobs are ignored, a message of a deleted job is rejected without
retry.

**Upgrading from 2.1 with a job stuck in `running`** (redelivered over and over because its output grew close to
`max_allowed_packet`): truncate its output before deploying, e.g.
`UPDATE job_queue SET output = LEFT(output, 1048576) WHERE id = <id>`. Its next redelivery marks it `failed`; the
status message is only appended while the output is below the cap, and later output chunks of such a job are dropped.

#### Update your database so the job tables are created

```shell
Expand Down Expand Up @@ -356,7 +428,9 @@ translations/messages.{locale}.yaml:

## Testing

The bundle has ready tests for job creations in the tests/ folder.
The bundle has tests for job creation and integration tests of the job processing in the tests/ folder.
The integration tests run the real message handler with real subprocesses against a temporary SQLite database, so
they need the `pdo_sqlite` PHP extension.
Running tests in your app can be done like this:

```bash
Expand All @@ -368,11 +442,12 @@ The tests are also run on every push / pull request on GitHub.
## Dependencies

* "php": ">=8.1",
* "doctrine/doctrine-bundle": "^2",
* "doctrine/doctrine-bundle": "^2|^3",
* "doctrine/orm": "^2|^3",
* "dragonmantank/cron-expression": "^3",
* "knplabs/knp-paginator-bundle": "^6",
* "spiriitlabs/form-filter-bundle": "^11",
* "psr/log": "^1|^2|^3",
* "spiriitlabs/form-filter-bundle": "^10|^11|^12",
* "symfony/form": "^6.4 || ^7.4",
* "symfony/framework-bundle": "^6.4 || ^7.4",
* "symfony/lock": "^6.4 || ^7.4",
Expand Down
5 changes: 5 additions & 0 deletions composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
"doctrine/orm": "^2|^3",
"dragonmantank/cron-expression": "^3",
"knplabs/knp-paginator-bundle": "^6",
"psr/log": "^1|^2|^3",
"spiriitlabs/form-filter-bundle": "^10|^11|^12",
"symfony/form": "^6.4 || ^7.4",
"symfony/framework-bundle": "^6.4 || ^7.4",
Expand Down Expand Up @@ -57,7 +58,11 @@
},
"public-dir": "public"
},
"conflict": {
"symfony/dependency-injection": ">=8.0"
},
"require-dev": {
"ext-pdo_sqlite": "*",
"phpunit/phpunit": "^10.5"
}
}
6 changes: 5 additions & 1 deletion config/services.xml
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,14 @@
<defaults autowire="true" autoconfigure="true">
<bind key="$jobTableName">%job_queue.database.job_table_name%</bind>
<bind key="$jobRecurringTableName">%job_queue.database.job_recurring_table_name%</bind>
<bind key="$pollIntervalMs">%job_queue.processing.poll_interval_ms%</bind>
<bind key="$outputMaxBytes">%job_queue.processing.output_max_bytes%</bind>
<bind key="$dbFailureTolerance">%job_queue.processing.db_failure_tolerance%</bind>
<bind key="$rerunOnRedelivery">%job_queue.processing.rerun_on_redelivery%</bind>
</defaults>

<prototype namespace="TomAtom\JobQueueBundle\" resource="../src/*"
exclude="../src/{DependencyInjection,Entity,Tests,Kernel.php}"/>
exclude="../src/{DependencyInjection,Entity,Output,Tests,Kernel.php}"/>

<service id="tomatom.job_queue.override_mapping_listener"
class="TomAtom\JobQueueBundle\EventListener\OverrideMappingListener"
Expand Down
20 changes: 10 additions & 10 deletions src/Controller/JobController.php
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
use TomAtom\JobQueueBundle\Form\JobFilterType;
use TomAtom\JobQueueBundle\Security\JobQueuePermissions;
use TomAtom\JobQueueBundle\Service\CommandJobFactory;
use TomAtom\JobQueueBundle\Service\JobOutputWriter;

#[Route(path: '/job')]
class JobController extends AbstractController
Expand Down Expand Up @@ -215,7 +216,7 @@ public function deleteRecurring(?JobRecurring $job): Response

#[IsGranted(JobQueuePermissions::ROLE_JOB_CANCEL)]
#[Route(path: '/cancel/{id<\d+>}', name: 'job_queue_cancel')]
public function cancel(?Job $job, Request $request): Response
public function cancel(?Job $job, Request $request, JobOutputWriter $writer): Response
{
if (empty($job)) {
$this->addFlash('warning', $this->translator->trans('job.detail.error.not_found'));
Expand All @@ -225,17 +226,16 @@ public function cancel(?Job $job, Request $request): Response
]);
}

if ($job->isCancellable()) {
// Try to cancel job if job is running
try {
$job->setCancelledAt(new DateTimeImmutable());
$this->entityManager->flush();
// Cancel the job only if it is (still) running - a conditional update, so a result the handler stored after
// the job was loaded here is never overwritten by the cancellation
try {
if ($job->isCancellable() && $writer->cancel($job->getId(), new DateTimeImmutable())) {
$this->addFlash('success', $this->translator->trans('job.cancellation.success'));
} catch (Exception $e) {
$this->addFlash('danger', $this->translator->trans('job.cancellation.error') . $e->getMessage());
} else {
$this->addFlash('warning', $this->translator->trans('job.cancellation.not_cancellable'));
}
} else {
$this->addFlash('warning', $this->translator->trans('job.cancellation.not_cancellable'));
} catch (Exception $e) {
$this->addFlash('danger', $this->translator->trans('job.cancellation.error') . $e->getMessage());
}

return $this->redirectToRoute('job_queue_detail', [
Expand Down
16 changes: 16 additions & 0 deletions src/JobQueueBundle.php
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,18 @@ public function configure(DefinitionConfigurator $definition): void
->scalarNode('heartbeat_interval')->defaultValue('1 minute')->end()
->end()
->end()
->arrayNode('processing')->addDefaultsIfNotSet()
->children()
->integerNode('poll_interval_ms')->defaultValue(1000)->min(0)
->info('How often the job output is written and the cancellation checked while the command runs.')->end()
->integerNode('output_max_bytes')->defaultValue(4194304)->min(0)
->info('Cap of the stored job output in bytes (keep it well below MySQL max_allowed_packet), 0 = unlimited.')->end()
->integerNode('db_failure_tolerance')->defaultValue(30)->min(1)
->info('Consecutive failed database polls (about one per poll interval) after which a running job is given up.')->end()
->booleanNode('rerun_on_redelivery')->defaultFalse()
->info('Run the command again when the message of a job which is already RUNNING is redelivered (e.g. after a worker crash). Default: mark the job as failed.')->end()
->end()
->end()
->end();
}

Expand Down Expand Up @@ -71,5 +83,9 @@ public function loadExtension(array $config, ContainerConfigurator $container, C
$builder->setParameter('job_queue.database.job_table_name', $config['database']['job_table_name']);
$builder->setParameter('job_queue.database.job_recurring_table_name', $config['database']['job_recurring_table_name']);
$builder->setParameter('job_queue.scheduling.heartbeat_interval', $config['scheduling']['heartbeat_interval']);
$builder->setParameter('job_queue.processing.poll_interval_ms', $config['processing']['poll_interval_ms']);
$builder->setParameter('job_queue.processing.output_max_bytes', $config['processing']['output_max_bytes']);
$builder->setParameter('job_queue.processing.db_failure_tolerance', $config['processing']['db_failure_tolerance']);
$builder->setParameter('job_queue.processing.rerun_on_redelivery', $config['processing']['rerun_on_redelivery']);
}
}
Loading
Loading