diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md
index f35edc0..895a6c4 100644
--- a/CONTRIBUTING.md
+++ b/CONTRIBUTING.md
@@ -96,13 +96,17 @@ Please include:
Before submitting your PR, make sure you've:
- [ ] Written clear and concise commit messages
-- [ ] Followed existing code style and naming conventions
+- [ ] Followed the [code style](docs/code-style.md), in particular: code is self-explaining and comments are a last resort
- [ ] Added or updated relevant documentation (if applicable)
- [ ] Added or updated unit tests (if applicable)
- [ ] Verified that all existing tests pass (`npm test`)
- [ ] Run linting and formatting (`npm run lint`, `npm run prettier`)
-- [ ] Updated the documentation site if needed
+- [ ] Updated the documentation site if needed, and for any public API change also `website/ai-usage.md`, the one-page contract agents read (take signatures from the code, not from memory)
+- [ ] Added a consumer-perspective test in `package-tests/consumer-app` for any new `global` surface
- [ ] Checked [API evolution rules](docs/api-evolution.md) if you touched a `global` type or added an extension point
+- [ ] Added every new field to its page layout in `force-app/main/default/layouts` and, for `AsyncResult__c`, to `AsyncResultAccess` (a field missing from either is invisible to admins)
+- [ ] Added any new `extras/` class to `extras/README.md` and to the tables in `website/introduction/packaged-install.md`
+- [ ] Added any new website page to the sidebar in `website/.vitepress/config.mts` (`llms.txt` is generated from it)
## 📝 Types of Contributions
diff --git a/README.md b/README.md
index 4b576da..c598f1b 100644
--- a/README.md
+++ b/README.md
@@ -43,13 +43,35 @@ Visit https://async.beyondthecloud.dev/ to view the full documentation.
- **Custom Metadata Configuration**: Configure the QueueableJob settings using the `QueueableJobSettings__mdt` custom metadata type to enable or disable jobs, and to control the creation of Async Result records.
- **Custom Object for Async Results**: The `AsyncResult__c` custom object is created for each processed queueable job, allowing you to track the chained job status and details.
-## Deploy to Salesforce
+## Installation
-
+Two ways in. Pick one.
+
+### Unlocked package
+
+
+
+
+
+Versioned, uninstallable, upgrades by installing the next version. Every class carries the `btcdev.` prefix and a few features need a class copied from [`extras/`](./extras). Guide: [Installing as a Package](https://async.beyondthecloud.dev/introduction/packaged-install).
+
+### Source deploy
+
+
+Or with the CLI:
+
+```bash
+git clone https://github.com/beyond-the-cloud-dev/async-lib.git
+cd async-lib
+sf project deploy start --source-dir force-app --target-org your-org
+```
+
+No namespace, nothing extra to set up, upgrades by redeploying the next tag. Guide: [Deploying the Source](https://async.beyondthecloud.dev/introduction/source-deploy).
+
## Contributors
diff --git a/docs/code-style.md b/docs/code-style.md
new file mode 100644
index 0000000..812b5eb
--- /dev/null
+++ b/docs/code-style.md
@@ -0,0 +1,93 @@
+# Code Style
+
+## Comments
+
+**Code must be self-explaining. A comment is a last resort, not a default.**
+
+Classes, methods, fields and variables carry the meaning. If a comment feels
+necessary, that is almost always a naming or structure problem, so fix the code
+instead:
+
+| Instead of a comment saying | Do this |
+| -------------------------------------------- | ------------------------------------------------------------ |
+| what a block does | extract it into a method whose name says it |
+| why a `catch` swallows | name the handler method, `reportWithoutAffectingTheJob(...)` |
+| that a static resets each transaction | name the field, `loggerCacheForThisTransaction` |
+| that a class must be `global` to be resolved | name the method, `newInstanceOfGlobalClass(...)` |
+| what a flag means | name the variable, `retryWillRestore` |
+
+Delete on sight: comments restating the code, section banners, commented-out
+code, narration ("first we...", "now handle..."), and ApexDoc that only echoes
+the signature.
+
+### The bar for keeping one
+
+A comment earns its place only when a competent Apex developer would be
+**surprised or misled** without it, and no name or structure can carry it. In
+practice that means a platform quirk or a deliberate choice that looks wrong:
+
+```apex
+// A failed cast is the only way to read the runtime type with its namespace.
+String.valueOf((DateTime) job);
+```
+
+```apex
+// List.sort() does not define the order of equal elements, so equal priority needs an
+// explicit tiebreak to keep jobs running in the order they were chained.
+```
+
+Rules of thumb that stay in prose belong in `website/explanations/`, not in the
+source. If the reason is about **how consumers use the library**, document it
+there and link it from the error message. If it is about **how the framework may
+evolve**, it belongs in `docs/api-evolution.md`.
+
+### The two allowed exceptions
+
+**PMD suppression justification.** Every `@SuppressWarnings` carries a header
+block saying why the rule is a false positive here. Without it a suppression is
+indistinguishable from hiding a defect.
+
+```apex
+/**
+ * PMD False Positives:
+ * - ExcessivePublicCount: one fluent method per job option
+ **/
+@SuppressWarnings('PMD.ExcessivePublicCount')
+```
+
+**A member that must never be deleted.** Where the reason for keeping
+dead-looking code is invisible, say so, because the next maintainer will
+otherwise remove it.
+
+```apex
+/**
+ * Superseded by Async.Retryable.resetBeforeRetry(Integer). Nothing calls this any more.
+ * It cannot be deleted: dropping a global member makes the package install fail in every
+ * subscriber org that referenced it. See docs/api-evolution.md.
+ **/
+```
+
+## Why ApexDoc is not enforced
+
+`pmd/ruleset.xml` deliberately excludes `category/apex/documentation.xml`.
+Requiring `@description` and `@param` on every member produces exactly the
+restatement this policy exists to remove. Editor plugins ship that rule on by
+default, so expect warnings; ignore them.
+
+## Design
+
+Ordinary clean-code expectations apply, and they matter more here than in an org
+codebase because this is a library whose public surface is
+[frozen once shipped](/docs/api-evolution.md):
+
+- **KISS.** The smallest thing that solves the actual problem. No configuration
+ nobody asked for.
+- **DRY, within reason.** Duplication in tests is often clearer than a shared
+ helper. Duplication in framework logic is a bug waiting to diverge.
+- **SOLID.** Most relevant here is interface segregation: many small capability
+ interfaces beat one fat one, because Apex has no default methods, so a fat
+ interface can never gain a member.
+- **Composition over inheritance.** A consumer has one inheritance slot. Do not
+ spend it. Prefer a marker interface the consumer can add to any class.
+- **Guard clauses over nesting.** Early return, and let the shape of the method
+ show the flow.
diff --git a/extras/README.md b/extras/README.md
index ec03e2b..7c4a654 100644
--- a/extras/README.md
+++ b/extras/README.md
@@ -5,6 +5,16 @@ namespace, which is exactly why they cannot ship inside the package.
Copy what you need. Rename anything to suit your project.
+| File | Needed for |
+| ---- | ---------- |
+| `BaseQueueableJob`, `BaseChunkJob` | `deepClone()`, `restoreStateOnRetry()`, `restoreStateOnNextChunk()` |
+| `AsyncJobSerializer` | `Async.requeue()` |
+
+All of them exist for one reason: JSON cannot cross a namespace boundary, so the conversion has to
+run in your namespace. See
+[Installing as a Package](https://async.beyondthecloud.dev/introduction/packaged-install) for the
+full checklist.
+
## `BaseQueueableJob`
Only needed when Async Lib is installed as a **namespaced package**. If you deployed the source
@@ -69,3 +79,22 @@ public class OddJob extends BaseQueueableJob {
See [Deep Clone in Packages](https://async.beyondthecloud.dev/explanations/deep-clone-in-packages)
for the full explanation and for the error messages that point back here.
+
+## `AsyncJobSerializer`
+
+Only needed when Async Lib is installed as a **namespaced package** and you use `Async.requeue()`.
+
+Requeue stores a snapshot of a job on `AsyncResult__c` and rebuilds it later. Both the store and
+the rebuild are JSON conversions, and JSON cannot cross a namespace boundary in either direction,
+so both have to run in your code. This class is that code.
+
+Register it once, on the `All` record of `QueueableJobSetting__mdt`:
+
+```
+JobSerializerClass__c = AsyncJobSerializer
+```
+
+It must stay `global`. Async Lib resolves it by name from inside its own namespace, and
+`Type.forName` reaches nothing else across the boundary. The methods stay `public`.
+
+See [Requeue](https://async.beyondthecloud.dev/explanations/requeue).
diff --git a/extras/classes/AsyncJobSerializer.cls b/extras/classes/AsyncJobSerializer.cls
new file mode 100644
index 0000000..cfd80f4
--- /dev/null
+++ b/extras/classes/AsyncJobSerializer.cls
@@ -0,0 +1,20 @@
+/**
+ * Copy this when Async Lib is installed as a namespaced package and you use Async.requeue().
+ * Register it once: QueueableJobSetting__mdt.JobSerializerClass__c = 'AsyncJobSerializer' on the
+ * All record.
+ *
+ * JSON cannot cross a namespace boundary in either direction, so Async Lib can neither store nor
+ * rebuild your job from inside its own namespace. Both halves run here instead, in yours.
+ *
+ * `global` is required. Async Lib resolves this class by name, and Type.forName reaches nothing
+ * else across the boundary. The methods stay public.
+ **/
+global class AsyncJobSerializer implements btcdev.Async.JobSerializer {
+ public String serialize(btcdev.QueueableJob job) {
+ return JSON.serialize(job);
+ }
+
+ public btcdev.QueueableJob deserialize(String className, String payload) {
+ return (btcdev.QueueableJob) JSON.deserialize(payload, Type.forName(className));
+ }
+}
diff --git a/extras/classes/AsyncJobSerializer.cls-meta.xml b/extras/classes/AsyncJobSerializer.cls-meta.xml
new file mode 100644
index 0000000..cad713d
--- /dev/null
+++ b/extras/classes/AsyncJobSerializer.cls-meta.xml
@@ -0,0 +1,5 @@
+
+
+ 66.0
+ Active
+
diff --git a/force-app/main/default/classes/Async.cls b/force-app/main/default/classes/Async.cls
index 75165ce..b3d0f26 100644
--- a/force-app/main/default/classes/Async.cls
+++ b/force-app/main/default/classes/Async.cls
@@ -58,6 +58,14 @@ public inherited sharing class Async {
QueueableManager.get().skipJob(customJobId);
}
+ public static RequeueSummary requeue(Id resultId) {
+ return AsyncRequeue.run(new Set{ resultId });
+ }
+
+ public static RequeueSummary requeue(Set resultIds) {
+ return AsyncRequeue.run(resultIds);
+ }
+
// Inside a QueueableJob subclass the inherited `backoff` field shadows the Backoff type, so a
// bare `Backoff.exponential(1)` will not compile there. Reaching the factories through Async
// is collision-free in both packaged and source deployments.
@@ -67,6 +75,27 @@ public inherited sharing class Async {
}
}
+ public interface OnJobEnqueued {
+ void onJobEnqueued(JobContext ctx);
+ }
+
+ public interface OnJobSucceeded {
+ void onJobSucceeded(JobContext ctx);
+ }
+
+ public interface OnJobFailed {
+ void onJobFailed(FailureContext ctx);
+ }
+
+ public interface OnRetryEnqueued {
+ void onRetryEnqueued(FailureContext ctx);
+ }
+
+ public interface JobSerializer {
+ String serialize(QueueableJob job);
+ QueueableJob deserialize(String className, String payload);
+ }
+
public interface Retryable {
void resetBeforeRetry(Integer attempt);
}
@@ -171,6 +200,13 @@ public inherited sharing class Async {
}
}
+ @JsonAccess(serializable='always' deserializable='always')
+ public class RequeueSummary {
+ public List requeued = new List();
+ public Map skipReasonByResultId = new Map();
+ public Result enqueueResult;
+ }
+
public enum AsyncType {
QUEUEABLE,
BATCHABLE,
@@ -189,6 +225,27 @@ public inherited sharing class Async {
EXHAUSTED
}
+ @JsonAccess(serializable='always' deserializable='always')
+ public class JobContext {
+ public String customJobId;
+ public String className;
+ public Id salesforceJobId;
+ public String chainId;
+ public Integer priority;
+ public Integer retryAttempt;
+ public Map info;
+
+ public JobContext(QueueableJob job) {
+ this.customJobId = job.customJobId;
+ this.className = job.className;
+ this.salesforceJobId = job.salesforceJobId;
+ this.chainId = job.chainId;
+ this.priority = job.priority;
+ this.retryAttempt = job.retryAttempt;
+ this.info = job.info ?? new Map();
+ }
+ }
+
@JsonAccess(serializable='always' deserializable='always')
public class FailureContext {
public RetryOutcome retryOutcome;
@@ -198,6 +255,8 @@ public inherited sharing class Async {
public Integer retryAttempt;
public Integer maxRetries;
public String retryHistory;
+ public Integer nextAttemptDelayMinutes;
+ public Map info;
public FailureContext(QueueableJob job) {
this.retryOutcome = retryOutcomeFor(job);
@@ -207,6 +266,7 @@ public inherited sharing class Async {
this.retryAttempt = job.retryAttempt;
this.maxRetries = job.maxRetries;
this.retryHistory = job.retryHistory;
+ this.info = job.info ?? new Map();
}
private Async.RetryOutcome retryOutcomeFor(QueueableJob job) {
diff --git a/force-app/main/default/classes/AsyncTest.cls b/force-app/main/default/classes/AsyncTest.cls
index 79dd9e9..aeca4ae 100644
--- a/force-app/main/default/classes/AsyncTest.cls
+++ b/force-app/main/default/classes/AsyncTest.cls
@@ -17,11 +17,13 @@ private class AsyncTest implements Database.Batchable {
private static Integer chunkCalloutPages = 0;
private static Integer chunkPageFailures = 0;
private static List chunkPagesRun = new List();
+ private static List loggedEvents = new List();
private static final String DUPLICATE_SIGNATURE_ERROR_MESSAGE = 'Attempt to enqueue job with duplicate queueable signature';
private static final String TEST_SIGNATURE_NAME = 'SignatureName';
private static final String TEST_SCHEDULABLE_JOB_NAME = 'SchedulableTestJob';
private static final String PRIMITIVE_VALUE_INITIAL = 'INITIAL_VALUE';
private static final String CHUNK_PROCESSED = 'CHUNK_PROCESSED';
+ private static final String LOGGER_CONSTRUCTOR_FAILURE = 'the logger needs a setting that is not there';
@IsTest
private static void shouldEnqueue60QueueablesSuccessfully() {
@@ -383,7 +385,7 @@ private class AsyncTest implements Database.Batchable {
job3.uniqueName = 'job3';
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
IsDisabled__c = true
@@ -412,7 +414,7 @@ private class AsyncTest implements Database.Batchable {
QueueableJobTest8 job8 = new QueueableJobTest8();
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
getClassNameWithNamespaceDotPrefix(
'AsyncTest.QueueableJobTest1'
) => new QueueableJobSetting__mdt(
@@ -463,7 +465,7 @@ private class AsyncTest implements Database.Batchable {
QueueableJobTest2 job2 = new QueueableJobTest2();
QueueableChain chain1 = new QueueableChain();
- chain1.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
getClassNameWithNamespaceDotPrefix(
'AsyncTest.QueueableJobTest1'
) => new QueueableJobSetting__mdt(
@@ -472,7 +474,7 @@ private class AsyncTest implements Database.Batchable {
)
};
QueueableChain chain2 = new QueueableChain();
- chain2.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
getClassNameWithNamespaceDotPrefix(
'AsyncTest.QueueableJobTest1'
) => new QueueableJobSetting__mdt(
@@ -2229,22 +2231,299 @@ private class AsyncTest implements Database.Batchable {
}
@IsTest
- private static void shouldGateAJobThatGetsItsRetryFromCustomMetadata() {
+ private static void shouldRouteToALoggerRegisteredThroughAsyncMock() {
+ AsyncMock.jobSettings(
+ new List{
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ LoggerClass__c = loggerName('RecordingLogger')
+ )
+ }
+ );
+ QueueableChain chain = new QueueableChain();
+ QueueableManager.get().setChain(chain);
+
+ chain.addJob(new SuccessfulQueueableTest());
+
+ Assert.isTrue(
+ loggedEvents.contains('enqueued:' + loggerName('SuccessfulQueueableTest') + ':null'),
+ 'A consumer must be able to register a logger without reaching into the framework: ' +
+ loggedEvents
+ );
+ }
+
+ @IsTest
+ private static void shouldApplyRetryDefaultsRegisteredThroughAsyncMock() {
+ AsyncMock.jobSettings(
+ new List{
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ MaxRetries__c = 3,
+ BackoffStrategy__c = 'FIXED',
+ BackoffBaseMinutes__c = 5
+ )
+ }
+ );
+ QueueableChain chain = new QueueableChain();
+ StateCarryingRetryJob job = new StateCarryingRetryJob();
+
+ chain.addJob(job);
+
+ Assert.areEqual(3, job.maxRetries, 'Retry settings must reach the chain through the mock.');
+ Assert.areEqual(5, job.backoff.delayMinutes(1), 'Backoff settings must reach it too.');
+ }
+
+ @IsTest
+ private static void shouldForgetMockedJobSettingsOnReset() {
+ AsyncMock.jobSettings(
+ new List{
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ MaxRetries__c = 3
+ )
+ }
+ );
+
+ AsyncMock.reset();
+
+ QueueableChain chain = new QueueableChain();
+ StateCarryingRetryJob job = new StateCarryingRetryJob();
+ chain.addJob(job);
+
+ Assert.areEqual(
+ 0,
+ job.maxRetries,
+ 'reset() must clear mocked settings with everything else.'
+ );
+ }
+
+ @IsTest
+ private static void shouldSendEveryLifecycleEventToTheRegisteredLogger() {
+ QueueableChain chain = chainWithLogger(loggerName('RecordingLogger'));
+ QueueableManager.get().setChain(chain);
+ StateCarryingRetryJob job = (StateCarryingRetryJob) Async.queueable(
+ new StateCarryingRetryJob()
+ )
+ .info('team', 'A')
+ .chain()
+ .job;
+ job.maxRetries = 1;
+ job.retryDecision = true;
+ job.continueOnJobExecuteFail = true;
+
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(committed());
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(committed());
+ List events = new List(loggedEvents);
+ Test.stopTest();
+
+ Assert.isTrue(
+ events.contains('enqueued:' + loggerName('StateCarryingRetryJob') + ':A'),
+ 'The enqueue event must fire and carry info(), but was: ' + events
+ );
+ Assert.isTrue(
+ events.contains('retry:0:delay=null'),
+ 'A queued retry must fire onRetryEnqueued, but was: ' + events
+ );
+ Assert.isTrue(
+ events.contains('failed:' + loggerName('StateCarryingRetryJob') + ':1'),
+ 'The terminal failure must fire onJobFailed, but was: ' + events
+ );
+ }
+
+ @IsTest
+ private static void shouldLetALoggerSubscribeToOneEventOnly() {
+ QueueableChain chain = chainWithLogger(loggerName('FailureOnlyLogger'));
+ QueueableManager.get().setChain(chain);
+ SuccessfulQueueableTest job = new SuccessfulQueueableTest();
+ chain.addJob(job);
+
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(committed());
+ List events = new List(loggedEvents);
+ Test.stopTest();
+
+ Assert.isTrue(events.isEmpty(), 'A logger must only receive what it implements: ' + events);
+ }
+
+ @IsTest
+ private static void shouldFireTheJobsOwnListenerAndTheGlobalLoggerTogether() {
+ QueueableChain chain = chainWithLogger(loggerName('RecordingLogger'));
+ QueueableManager.get().setChain(chain);
+ chain.addJob(new SelfLoggingJob());
+
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(committed());
+ List events = new List(loggedEvents);
+ Test.stopTest();
+
+ Assert.isTrue(
+ events.contains('self:' + loggerName('SelfLoggingJob')),
+ 'The job own listener must fire, but was: ' + events
+ );
+ Assert.isTrue(
+ events.contains('succeeded:' + loggerName('SelfLoggingJob')),
+ 'The registered logger must fire too, but was: ' + events
+ );
+ }
+
+ @IsTest
+ private static void shouldNotFailTheJobWhenTheLoggerThrows() {
+ QueueableChain chain = chainWithLogger(loggerName('ThrowingLogger'));
+ QueueableManager.get().setChain(chain);
+ SuccessfulQueueableTest job = new SuccessfulQueueableTest();
+ chain.addJob(job);
+
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(committed());
+ Test.stopTest();
+
+ Assert.isFalse(job.hasFailed, 'A logger that throws must not fail the job it observed.');
+ }
+
+ @IsTest
+ private static void shouldRunNormallyWhenNoLoggerIsConfigured() {
+ QueueableChain chain = new QueueableChain();
+ QueueableManager.get().setChain(chain);
+ SuccessfulQueueableTest job = new SuccessfulQueueableTest();
+ chain.addJob(job);
+
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(committed());
+ Test.stopTest();
+
+ Assert.isFalse(job.hasFailed, 'No logger configured must be a silent no-op.');
+ Assert.isTrue(loggedEvents.isEmpty(), 'Nothing should have been logged.');
+ }
+
+ @IsTest
+ private static void shouldWarnRatherThanHaltWhenTheLoggerClassCannotBeResolved() {
+ QueueableChain chain = chainWithLogger('NoSuchLoggerClass');
+ QueueableManager.get().setChain(chain);
+ SuccessfulQueueableTest job = new SuccessfulQueueableTest();
+
+ chain.addJob(job);
+
+ Assert.isTrue(
+ job.retryHistory.contains('NoSuchLoggerClass'),
+ 'The unresolvable class must be named, but was: ' + job.retryHistory
+ );
+ Assert.isTrue(
+ job.retryHistory.contains('global'),
+ 'The warning must say the class has to be global, but was: ' + job.retryHistory
+ );
+ }
+
+ @IsTest
+ private static void shouldPreferAJobSpecificLoggerOverTheOrgWideDefault() {
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
- MaxRetries__c = 2
+ LoggerClass__c = loggerName('RecordingLogger')
+ ),
+ loggerName('SuccessfulQueueableTest') => new QueueableJobSetting__mdt(
+ QueueableJobName__c = loggerName('SuccessfulQueueableTest'),
+ LoggerClass__c = loggerName('FailureOnlyLogger')
)
};
+ QueueableManager.get().setChain(chain);
+
+ chain.addJob(new SuccessfulQueueableTest());
+
+ Assert.isTrue(
+ loggedEvents.isEmpty(),
+ 'The job specific logger only implements OnJobFailed, so enqueue logs nothing: ' +
+ loggedEvents
+ );
+ }
+
+ @IsTest
+ private static void shouldNotApplyCustomMetadataRetryToAJobThatDeclaresNoReset() {
+ QueueableChain chain = chainWithSettings(
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ MaxRetries__c = 2
+ )
+ );
+ UngatedRetryJob job = new UngatedRetryJob();
+
+ chain.addJob(job);
+
+ Assert.areEqual(
+ 0,
+ job.maxRetries,
+ 'An admin turning retry on org-wide must not enable it for a job that never declared a reset.'
+ );
+ Assert.isTrue(
+ job.retryHistory.contains('Async.Retryable'),
+ 'The job must carry the reason retry was skipped, but was: ' + job.retryHistory
+ );
+ }
+
+ @IsTest
+ private static void shouldNotHaltAJobWhenCustomMetadataNamesAnUnknownBackoffStrategy() {
+ QueueableChain chain = chainWithSettings(
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ MaxRetries__c = 2,
+ BackoffStrategy__c = 'expnential'
+ )
+ );
+ StateCarryingRetryJob job = new StateCarryingRetryJob();
+
+ chain.addJob(job);
+
+ Assert.areEqual(2, job.maxRetries, 'Retry still applies; only the bad backoff is dropped.');
+ Assert.isNull(job.backoff, 'An unknown strategy must leave backoff unset, not throw.');
+ Assert.isTrue(
+ job.retryHistory.contains('expnential'),
+ 'The typo must be named in the job history, but was: ' + job.retryHistory
+ );
+ Assert.isTrue(
+ job.retryHistory.contains('EXPONENTIAL_JITTER'),
+ 'The valid strategies must still be listed, but was: ' + job.retryHistory
+ );
+ }
+ @IsTest
+ private static void shouldClampCustomMetadataRetriesAboveTheCapInsteadOfThrowing() {
+ QueueableChain chain = chainWithSettings(
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ MaxRetries__c = QueueableManager.MAX_RETRY_CAP + 40
+ )
+ );
+ StateCarryingRetryJob job = new StateCarryingRetryJob();
+
+ chain.addJob(job);
+
+ Assert.areEqual(
+ QueueableManager.MAX_RETRY_CAP,
+ job.maxRetries,
+ 'A number above the cap must clamp, not stop the job from enqueueing.'
+ );
+ Assert.isTrue(
+ job.retryHistory.contains(String.valueOf(QueueableManager.MAX_RETRY_CAP)),
+ 'The clamp must be recorded, but was: ' + job.retryHistory
+ );
+ }
+
+ @IsTest
+ private static void shouldStillThrowWhenRetryIsAskedForInApexWithoutAReset() {
try {
- chain.addJob(new UngatedRetryJob());
- Assert.fail('Retry configured by CMDT must be gated the same as retry().');
+ Async.queueable(new UngatedRetryJob()).retry(2).enqueue();
+ Assert.fail('Retry written in Apex is the developer own code and must still throw.');
} catch (IllegalArgumentException ex) {
Assert.isTrue(
ex.getMessage().contains('Async.Retryable'),
- 'CMDT-configured retry must reach the same gate, but was: ' + ex.getMessage()
+ 'Code-triggered misconfiguration keeps the loud gate, but was: ' + ex.getMessage()
);
}
}
@@ -2391,7 +2670,7 @@ private class AsyncTest implements Database.Batchable {
QueueableJobTest1 job = new QueueableJobTest1();
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
MaxRetries__c = 3,
@@ -2414,7 +2693,7 @@ private class AsyncTest implements Database.Batchable {
job.maxRetries = 5;
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
MaxRetries__c = 3
@@ -2425,29 +2704,6 @@ private class AsyncTest implements Database.Batchable {
Assert.areEqual(5, job.maxRetries, 'Explicit fluent config must win over CMDT.');
}
- @IsTest
- private static void shouldRejectCmdtMaxRetriesAboveCap() {
- QueueableJobTest1 job = new QueueableJobTest1();
-
- QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
- QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
- DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
- MaxRetries__c = QueueableManager.MAX_RETRY_CAP + 1
- )
- };
-
- try {
- chain.addJob(job);
- Assert.fail('Should reject a CMDT retry count above the framework cap.');
- } catch (Exception ex) {
- Assert.areEqual(
- QueueableManager.ERROR_MESSAGE_MAX_RETRIES_EXCEEDS_CAP,
- ex.getMessage()
- );
- }
- }
-
@IsTest
private static void shouldReEnqueueRetryJobOnFailure() {
FailureQueueableTest job = new FailureQueueableTest();
@@ -2505,13 +2761,6 @@ private class AsyncTest implements Database.Batchable {
);
}
- private static Async.Dependency dependencyOn(String customJobId, Async.Outcome outcome) {
- Async.Dependency dependency = new Async.Dependency();
- dependency.resolvedTargetCustomJobId = customJobId;
- dependency.requiredOutcome = outcome;
- return dependency;
- }
-
@IsTest
private static void shouldRunDependentJobWhenRequiredOutcomeMatches() {
SuccessfulQueueableTest jobA = new SuccessfulQueueableTest();
@@ -2808,21 +3057,12 @@ private class AsyncTest implements Database.Batchable {
);
}
- private static Map resultsEnabledForAll() {
- return new Map{
- QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
- DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
- CreateResult__c = true
- )
- };
- }
-
@IsTest
private static void shouldCreateCompletedResultWithIdentityFields() {
SuccessfulQueueableTest job = new SuccessfulQueueableTest();
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
chain.addJob(job);
QueueableManager.get().setChain(chain);
@@ -2852,7 +3092,7 @@ private class AsyncTest implements Database.Batchable {
SuccessfulQueueableTest jobB = new SuccessfulQueueableTest();
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
chain.addJob(jobA);
jobB.dependencies = new List{
dependencyOn(jobA.customJobId, Async.Outcome.SUCCESS)
@@ -2911,7 +3151,7 @@ private class AsyncTest implements Database.Batchable {
IsDisabled__c = true
)
);
- chain.queueableJobSettingByJobName = settings;
+ QueueableChain.jobSettingByName = settings;
chain.addJob(job);
Test.startTest();
@@ -2927,17 +3167,6 @@ private class AsyncTest implements Database.Batchable {
Assert.isNotNull(result.SkipReason__c);
}
- private static MarkerJob markerJob(String tag, Boolean shouldFail) {
- MarkerJob job = new MarkerJob();
- job.tag = tag;
- job.shouldFail = shouldFail;
- return job;
- }
-
- private static Integer accountCount(String name) {
- return [SELECT COUNT() FROM Account WHERE Name = :name];
- }
-
@IsTest
private static void shouldGateChainOnDependencyOutcomesEndToEnd() {
Test.startTest();
@@ -2954,7 +3183,7 @@ private class AsyncTest implements Database.Batchable {
@IsTest
private static void shouldRecordChainResultsEndToEnd() {
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
QueueableManager.get().setChain(chain);
Test.startTest();
@@ -3031,7 +3260,7 @@ private class AsyncTest implements Database.Batchable {
@IsTest
private static void shouldRecordRetryHistoryOnExhaustionEndToEnd() {
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
QueueableManager.get().setChain(chain);
Test.startTest();
@@ -3560,30 +3789,6 @@ private class AsyncTest implements Database.Batchable {
);
}
- @IsTest
- private static void shouldFailFastOnAnUnknownBackoffStrategyInMetadata() {
- QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
- QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
- DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
- QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
- MaxRetries__c = 2,
- BackoffStrategy__c = 'EXPONENTAIL'
- )
- };
-
- try {
- chain.addJob(new SuccessfulQueueableTest());
- Assert.fail('A typo in BackoffStrategy__c must not silently disable backoff.');
- } catch (IllegalArgumentException ex) {
- Assert.isTrue(ex.getMessage().contains('EXPONENTAIL'), 'The bad value is named.');
- Assert.isTrue(
- ex.getMessage().contains('EXPONENTIAL_JITTER'),
- 'The valid strategies are listed.'
- );
- }
- }
-
@IsTest
private static void shouldRecordWhyAFailureWasNotRetriedWhenItsTypeIsNotListed() {
SuccessfulQueueableTest job = new SuccessfulQueueableTest();
@@ -3602,7 +3807,7 @@ private class AsyncTest implements Database.Batchable {
@IsTest
private static void shouldKeepEveryAttemptInRetryHistoryAfterDiscardingChainChanges() {
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
QueueableManager.get().setChain(chain);
Test.startTest();
@@ -3692,7 +3897,7 @@ private class AsyncTest implements Database.Batchable {
@IsTest
private static void shouldKeepChainRunningWhenOnFinalFailureThrows() {
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
QueueableManager.get().setChain(chain);
Test.startTest();
@@ -3728,49 +3933,18 @@ private class AsyncTest implements Database.Batchable {
Assert.areEqual(2, accountCount('HOOK'), 'Each failed page is its own final failure.');
}
- private static Account hookRecord() {
- return [SELECT Description FROM Account WHERE Name = 'HOOK' LIMIT 1];
- }
+ @IsTest
+ private static void shouldOrderFinalizerBeforeRegularJob() {
+ QueueableJobTest1 finalizerJob = new QueueableJobTest1();
+ finalizerJob.parentCustomJobId = 'parent';
+ QueueableJobTest1 regularJob = new QueueableJobTest1();
- private static String getClassNameWithNamespaceDotPrefix(String className) {
- return getNamespaceDotPrefix() + className;
- }
-
- private static String getNamespaceDotPrefix() {
- String className = AsyncTest.class.getName();
- return className.contains('.') ? className.substringBefore('.') + '.' : '';
- }
-
- public Iterable start(Database.BatchableContext bc) {
- // This is just a placeholder to start the batch.
- return new List{ new Account() };
- }
-
- public void execute(Database.BatchableContext ctx, List scope) {
- for (Account acc : scope) {
- acc.Description = 'Processed by: ' + ctx.getJobId();
- }
- if (!scope.isEmpty() && scope[0].Id != null) {
- update scope;
- }
- }
-
- public void finish(Database.BatchableContext bc) {
- insert new Account(Name = 'Batch Complete', Description = 'Job: ' + bc.getJobId());
- }
-
- @IsTest
- private static void shouldOrderFinalizerBeforeRegularJob() {
- QueueableJobTest1 finalizerJob = new QueueableJobTest1();
- finalizerJob.parentCustomJobId = 'parent';
- QueueableJobTest1 regularJob = new QueueableJobTest1();
-
- Assert.areEqual(-1, finalizerJob.compareTo(regularJob), 'Finalizer sorts first.');
- Assert.areEqual(
- 1,
- regularJob.compareTo(finalizerJob),
- 'Regular job sorts after finalizer.'
- );
+ Assert.areEqual(-1, finalizerJob.compareTo(regularJob), 'Finalizer sorts first.');
+ Assert.areEqual(
+ 1,
+ regularJob.compareTo(finalizerJob),
+ 'Regular job sorts after finalizer.'
+ );
}
@IsTest
@@ -3862,7 +4036,7 @@ private class AsyncTest implements Database.Batchable {
QueueableJobTest1 job = new QueueableJobTest1();
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = new Map{
+ QueueableChain.jobSettingByName = new Map{
QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
MaxRetries__c = 2,
@@ -4093,1653 +4267,2474 @@ private class AsyncTest implements Database.Batchable {
}
}
- private static AsyncResult__c insertAsyncResultWithAge(String status, Integer ageDays) {
- AsyncResult__c result = new AsyncResult__c(Status__c = status);
- insert result;
- Test.setCreatedDate(result.Id, System.now().addDays(-ageDays));
- return result;
- }
+ @IsTest
+ private static void shouldProcessFirstChunkThroughEnqueue() {
+ List accounts = createAccounts(5);
- private static Integer followUpsLeftIn(QueueableChain chain) {
- Integer followUps = 0;
- for (QueueableJob job : chain.jobs) {
- if (job instanceof MarkerJob) {
- followUps++;
- }
- }
- return followUps;
- }
+ Test.startTest();
+ Async.Result result = Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(5)
+ .enqueue();
+ Test.stopTest();
- private static QueueableChain chainRunning(QueueableJob job) {
- QueueableChain chain = new QueueableChain();
- chain.addJob(job);
- QueueableManager.get().setChain(chain);
- return chain;
+ Assert.areNotEqual(null, result.salesforceJobId);
+ Assert.areEqual(
+ 5,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'The single chunk should process every record in the source.'
+ );
}
- private static QueueableChain chainWithEnqueuedJob(QueueableJob job) {
- QueueableChain chain = chainRunning(job);
- job.chain = chain;
- return chain;
- }
+ @IsTest
+ private static void shouldProcessAllInMemoryChunksAcrossPages() {
+ List accounts = createAccounts(6);
- private static AsyncMock.MockFinalizerContext rolledBack() {
- return new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.UNHANDLED_EXCEPTION);
- }
+ Test.startTest();
+ Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts)).chunkSize(3).enqueue();
+ Test.stopTest();
- private static AsyncMock.MockFinalizerContext committed() {
- return new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.SUCCESS);
+ Assert.areEqual(
+ 6,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'Self-chaining must process every page of the run.'
+ );
}
- private class DeepCloneFailJob extends QueueableJob {
- public override void work() {
- }
- public override QueueableJob cloneForDeepCopy() {
- throw new JSONException('forced deep clone failure');
- }
- }
+ @IsTest
+ private static void shouldProcessAllCursorChunksAcrossPages() {
+ createAccounts(6);
- public class SelfReferencingJob extends QueueableJob {
- public SelfReferencingJob self;
- public override void work() {
- }
- }
+ Test.startTest();
+ Async.chunk(new MarkingChunkJob(), ChunkSource.query('SELECT Id FROM Account'))
+ .chunkSize(3)
+ .enqueue();
+ Test.stopTest();
- public class AbstractFieldHoldingJob extends QueueableJob {
- public Comparable held;
- public override void work() {
- }
+ Assert.areEqual(
+ 6,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'A cursor must survive re-enqueue and fetch every page.'
+ );
}
- public class AsyncLibTypeHoldingJob extends QueueableJob {
- public Async.Result heldResult;
- public Backoff heldBackoff;
- public Async.Dependency heldDependency;
- public override void work() {
- }
- }
+ @IsTest
+ private static void shouldCarryJobMemberStateAcrossChunks() {
+ List accounts = createAccounts(6);
- public class DeepCloneRetryJob extends QueueableJob implements Async.Retryable {
- public List attempts = new List();
- public override void work() {
- attempts.add('attempt' + retryAttempt);
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
+ Test.startTest();
+ Async.chunk(new LabelingChunkJob('RUN_LABEL'), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .enqueue();
+ Test.stopTest();
- public void resetBeforeRetry(Integer attempt) {
- }
+ Assert.areEqual(
+ 6,
+ [SELECT COUNT() FROM Account WHERE Site = 'RUN_LABEL'],
+ 'Job member state must survive cloning and serialization across every page.'
+ );
}
- public class DeepCloneFailRetryJob extends QueueableJob implements Async.Retryable {
- public void resetBeforeRetry(Integer attempt) {
- }
+ @IsTest
+ private static void shouldContinueToNextChunkAfterAFailedChunk() {
+ List accounts = accountsWithOneFailingRecord();
- public override void work() {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- public override QueueableJob cloneForDeepCopy() {
- throw new JSONException('forced deep clone failure');
- }
- }
+ Test.startTest();
+ Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts)).chunkSize(2).enqueue();
+ Test.stopTest();
- private class SuccessfulQueueableTest extends QueueableJob {
- public override void work() {
- insert new Account(Name = Async.getQueueableJobContext()?.currentJob?.uniqueName);
- }
+ Assert.areEqual(
+ 4,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'A failed chunk must not stop the remaining chunks by default.'
+ );
}
- private class FailureQueueableTest extends QueueableJob.AllowsCallouts implements Async.Retryable {
- public void resetBeforeRetry(Integer attempt) {
- }
+ @IsTest
+ private static void shouldStopAfterFailedChunkWhenConfigured() {
+ List accounts = new List{
+ new Account(Name = 'FAIL first'),
+ new Account(Name = 'ok second'),
+ new Account(Name = 'ok third'),
+ new Account(Name = 'ok fourth')
+ };
+ insert accounts;
- public override void work() {
- insert new Account(Name = Async.getQueueableJobContext()?.currentJob?.uniqueName);
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ Test.startTest();
+ Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .stopRemainingChunksOnFailure()
+ .enqueue();
+ Test.stopTest();
- private class SelfConfiguringRetryJob extends QueueableJob implements Async.Retryable {
- public void resetBeforeRetry(Integer attempt) {
- }
+ Assert.areEqual(
+ 0,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'stopRemainingChunksOnFailure must halt the run after a failed chunk.'
+ );
+ }
- private SelfConfiguringRetryJob() {
- this.maxRetries = 2;
- this.backoff = Async.Backoff.fixed(4);
- this.continueOnJobExecuteFail = true;
- }
+ @IsTest
+ private static void shouldApplyAllChunkBuilderOptions() {
+ List accounts = createAccounts(2);
- public override void work() {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ Test.startTest();
+ Async.Result result = Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .priority(3)
+ .delay(1)
+ .delayBetweenChunks(1)
+ .retry(2)
+ .backoff(Backoff.fixed(1))
+ .retryOn(CustomException.class)
+ .mockId('chunk-mock')
+ .keepChunkPages()
+ .stopRemainingChunksOnFailure()
+ .enqueue();
+ Test.stopTest();
- private class ChainStoppingJob extends QueueableJob {
- public override void work() {
- Async.queueable(new StopChainFinalizer()).attachFinalizer();
- }
+ Assert.areNotEqual(
+ null,
+ result.salesforceJobId,
+ 'Every builder option should still enqueue.'
+ );
}
- private class MarkerJob extends QueueableJob implements Async.Retryable {
- public String tag;
- public Boolean shouldFail = false;
- public override void work() {
- insert new Account(Name = tag);
- if (shouldFail) {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ @IsTest
+ private static void shouldProcessFirstChunkInBulk() {
+ List accounts = createAccounts(200);
- public void resetBeforeRetry(Integer attempt) {
- }
- }
+ Test.startTest();
+ Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts)).chunkSize(200).enqueue();
+ Test.stopTest();
- private class StopChainFinalizer extends QueueableJob.Finalizer {
- public override void work() {
- Async.stopChain();
- }
+ Assert.areEqual(
+ 200,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'A 200-record chunk should process in one bulk-safe page.'
+ );
}
- private class FinalizerAttachingRetryJob extends QueueableJob implements Async.Retryable {
- public void resetBeforeRetry(Integer attempt) {
- }
-
- public override void work() {
- insert new Account(Name = 'RETRY-PARENT');
- Async.queueable(new MarkerFinalizer()).attachFinalizer();
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ @IsTest
+ private static void shouldFetchCorrectOffsetForLaterChunk() {
+ List accounts = createAccounts(10);
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.getRun().configure(ChunkSource.of(accounts), 5, false, null, false);
+ job.getRun().position = 5;
- private class MarkerFinalizer extends QueueableJob.Finalizer {
- public override void work() {
- insert new Account(Name = 'RETRY-FINALIZER');
- }
- }
+ job.work();
- private abstract class HookRecordingJob extends QueueableJob implements Async.Retryable {
- public void resetBeforeRetry(Integer attempt) {
+ Set laterHalf = new Set();
+ for (Integer i = 5; i < 10; i++) {
+ laterHalf.add(accounts[i].Id);
}
-
- public override void onFinalFailure(Async.FailureContext failureCtx) {
- insert new Account(
- Name = 'HOOK',
- Description = failureCtx.retryOutcome.name() +
- '|' +
- (String.isBlank(failureCtx.failure?.stackTrace) ? 'NO-TRACE' : 'HAS-TRACE') +
- '|' +
- failureCtx.retryAttempt +
- '/' +
- failureCtx.maxRetries
+ List processed = [SELECT Id FROM Account WHERE Site = :CHUNK_PROCESSED];
+ Assert.areEqual(
+ 5,
+ processed.size(),
+ 'Only the second page of records should be processed.'
+ );
+ for (Account processedAccount : processed) {
+ Assert.isTrue(
+ laterHalf.contains(processedAccount.Id),
+ 'Offset fetch must return records from position 5 onward.'
);
}
}
- private class FailureHookJob extends HookRecordingJob {
- public override void work() {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ @IsTest
+ private static void shouldRejectChunkJobRunningOutsideAChunkRun() {
+ try {
+ new MarkingChunkJob().work();
+ Assert.fail('A ChunkJob without a configured run must not execute.');
+ } catch (IllegalArgumentException ex) {
+ Assert.areEqual(ChunkRun.ERROR_MESSAGE_RUN_NOT_STARTED, ex.getMessage());
}
}
- private class SucceedingHookJob extends HookRecordingJob {
- public override void work() {
- insert new Account(Name = 'SUCCEEDED');
- }
- }
+ @IsTest
+ private static void shouldCreateResultRowForChunk() {
+ List accounts = createAccounts(3);
+ QueueableChain chain = new QueueableChain();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
- private class ExplodingHookJob extends QueueableJob {
- public override void work() {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
+ Test.startTest();
+ QueueableManager.get().setChain(chain);
+ Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts)).chunkSize(3).enqueue();
+ Test.stopTest();
- public override void onFinalFailure(Async.FailureContext failureCtx) {
- throw new CustomException('hook exploded');
- }
+ List results = [
+ SELECT Id, Status__c, ChainId__c, ClassName__c
+ FROM AsyncResult__c
+ ];
+ Assert.areEqual(1, results.size(), 'A completed chunk page should record one AsyncResult.');
+ Assert.areEqual(QueueableManager.STATUS_COMPLETED, results[0].Status__c);
+ Assert.areNotEqual(null, results[0].ChainId__c);
}
- private class FailingChunkHookJob extends ChunkJob implements Async.ChunkResettable {
- public void resetBeforeNextChunk(Integer pageNumber) {
- }
+ @IsTest
+ private static void shouldReturnNextChunkWhenRecordsRemain() {
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
- public override void work(List chunk) {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
+ ChunkJob nextPage = job.nextPageOrNull();
- public override void onFinalFailure(Async.FailureContext failureCtx) {
- insert new Account(Name = 'HOOK', Description = failureCtx.retryOutcome.name());
- }
+ Assert.areNotEqual(null, nextPage, 'A next page is expected while records remain.');
+ Assert.areEqual(5, nextPage.getRun().position, 'The next page must advance by chunkSize.');
}
- private class JobChainingRetryJob extends QueueableJob implements Async.Retryable {
- public void resetBeforeRetry(Integer attempt) {
- }
+ @IsTest
+ private static void shouldNotCarryInitialDelayToLaterChunks() {
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
+ job.delay = 5;
- public override void work() {
- insert new Account(Name = 'CHAIN-PARENT');
- Async.queueable(markerJob('CHAIN-FOLLOWUP', false)).chain();
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ ChunkJob nextPage = job.nextPageOrNull();
- private class ChainingFailureJob extends QueueableJob {
- public override void work() {
- Async.queueable(markerJob('ROLLED-BACK-FOLLOWUP', false)).chain();
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
+ Assert.areEqual(null, nextPage.delay, 'An initial delay must not repeat on every page.');
}
- private class FinalizerAndJobChainingFailureJob extends QueueableJob {
- public override void work() {
- Async.queueable(new MarkerFinalizer()).attachFinalizer();
- Async.queueable(markerJob('ROLLED-BACK-FOLLOWUP', false)).chain();
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ @IsTest
+ private static void shouldApplyDelayBetweenChunks() {
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, 3, false);
- private class ChainStoppingFailureJob extends QueueableJob implements Async.Retryable {
- public void resetBeforeRetry(Integer attempt) {
- }
+ ChunkJob nextPage = job.nextPageOrNull();
- public override void work() {
- Async.stopChain();
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
+ Assert.areEqual(3, nextPage.delay, 'delayBetweenChunks must throttle each later page.');
}
- private class JobSkippingFailureJob extends QueueableJob {
- public String targetCustomJobId;
- public override void work() {
- Async.skipJob(targetCustomJobId);
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ @IsTest
+ private static void shouldRejectDeepCloneChunkJob() {
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.deepClone = true;
- private class NonRetryableJob extends QueueableJob {
- public override void work() {
- }
- public override Boolean isRetryable(Exception ex) {
- return false;
+ try {
+ Async.chunk(job, ChunkSource.of(createAccounts(2)));
+ Assert.fail('A deepClone ChunkJob should be rejected at build time.');
+ } catch (Exception ex) {
+ Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_DEEP_CLONE_UNSUPPORTED, ex.getMessage());
}
}
- private class MessageVetoJob extends QueueableJob {
- public override void work() {
- }
- public override Boolean isRetryable(Exception ex) {
- return !ex.getMessage().containsIgnoreCase('permanent');
- }
- }
+ @IsTest
+ private static void shouldPruneSettledPagesByDefault() {
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.getRun().configure(ChunkSource.of(createAccounts(4)), 2, false, null, false);
- private class ThrowingClassifierJob extends QueueableJob {
- public override void work() {
- }
- public override Boolean isRetryable(Exception ex) {
- throw new CustomException('classifier blew up');
- }
+ Assert.isFalse(job.getRun().keepPages, 'Settled pages are pruned by default.');
}
- private class ResettableRetryJob extends QueueableJob implements Async.Retryable {
- public Boolean wasReset = false;
- public Integer resetAttempt;
- public override void work() {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- public void resetBeforeRetry(Integer attempt) {
- this.wasReset = true;
- this.resetAttempt = attempt;
- }
- }
+ @IsTest
+ private static void shouldKeepSettledPagesWhenRequested() {
+ List accounts = createAccounts(6);
- private class UngatedRetryJob extends QueueableJob {
- public override void work() {
- }
+ Test.startTest();
+ Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .keepChunkPages()
+ .enqueue();
+ Test.stopTest();
+
+ Assert.areEqual(
+ 6,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'A kept-pages run still processes every page.'
+ );
}
- public class RetryingChunkJob extends ChunkJob implements Async.Retryable, Async.ChunkResettable {
- public List touched = new List();
+ @IsTest
+ private static void shouldReturnNoNextChunkAtEndOfSource() {
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
+ job.getRun().position = 5;
- public override void work(List chunk) {
- String firstOfPage = (String) chunk[0].get('Name');
- AsyncTest.chunkPagesRun.add(firstOfPage + ':' + touched.size());
- touched.add('page');
- if (firstOfPage == 'P3' && AsyncTest.chunkPageFailures == 0) {
- AsyncTest.chunkPageFailures++;
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- }
+ Assert.areEqual(null, job.nextPageOrNull(), 'No page should follow the final chunk.');
+ }
- public void resetBeforeRetry(Integer attempt) {
- }
+ @IsTest
+ private static void shouldHaltRunAfterFailedPageWhenConfigured() {
+ MarkingChunkJob stopping = new MarkingChunkJob();
+ stopping.getRun().configure(ChunkSource.of(createAccounts(10)), 5, true, null, false);
- public void resetBeforeNextChunk(Integer pageNumber) {
- }
+ Assert.isTrue(
+ stopping.getRun().isHaltedBy(true),
+ 'A failed page must halt the run when stopRemainingChunksOnFailure is set.'
+ );
+ Assert.isFalse(
+ stopping.getRun().isHaltedBy(false),
+ 'A page that succeeded must not halt the run.'
+ );
}
- private class GateTrippingParentJob extends QueueableJob {
- public override void work() {
- insert new Account(Name = 'GATE-PARENT-RAN');
- Async.queueable(new UngatedRetryJob()).retry(2).chain();
- }
- }
+ @IsTest
+ private static void shouldContinueRemainingChunksOnFailureByDefault() {
+ MarkingChunkJob continuing = new MarkingChunkJob();
+ continuing.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
- public class UngatedChunkJob extends ChunkJob {
- public override void work(List chunk) {
- }
+ Assert.isFalse(
+ continuing.getRun().isHaltedBy(true),
+ 'A failed chunk must not halt the run by default.'
+ );
+ Assert.areNotEqual(
+ null,
+ continuing.nextPageOrNull(),
+ 'The run continues to the next page after a failed page.'
+ );
}
- private class OkCalloutMock implements HttpCalloutMock {
- public HttpResponse respond(HttpRequest request) {
- HttpResponse response = new HttpResponse();
- response.setStatusCode(200);
- return response;
- }
- }
+ @IsTest
+ private static void shouldCarryFailedPageFlagToLaterPages() {
+ MarkingChunkJob job = new MarkingChunkJob();
+ job.getRun().configure(ChunkSource.of(createAccounts(6)), 2, false, null, false);
+ job.hasFailed = true;
- public class MarkerCalloutChunkJob extends ChunkJob implements Database.AllowsCallouts, Async.ChunkResettable {
- public override void work(List chunk) {
- HttpRequest request = new HttpRequest();
- request.setEndpoint('https://example.com');
- request.setMethod('GET');
- AsyncTest.chunkCalloutStatus = new Http().send(request).getStatusCode();
- AsyncTest.chunkCalloutPages++;
- }
+ ChunkJob nextPage = job.nextPageOrNull();
- public void resetBeforeNextChunk(Integer pageNumber) {
- }
+ Assert.isTrue(
+ nextPage.getRun().hasFailedPage,
+ 'The run remembers a failed page so the run outcome stays FAILURE.'
+ );
+ Assert.isFalse(nextPage.hasFailed, 'The next page starts clean.');
}
- public class StateCarryingRetryJob extends QueueableJob implements Async.Retryable {
- public List seen = new List();
-
- public override void work() {
- insert new Account(Name = 'RETRY-SIZE-' + seen.size());
- seen.add('attempt');
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
+ @IsTest
+ private static void shouldTreatEmptySourceAsNoOp() {
+ Async.Result result = Async.chunk(
+ new MarkingChunkJob(),
+ ChunkSource.of(new List())
+ )
+ .enqueue();
- public void resetBeforeRetry(Integer attempt) {
- }
+ Assert.areEqual(null, result.salesforceJobId, 'An empty source should enqueue nothing.');
+ Assert.isTrue(result.queueableChainState.jobs.isEmpty());
}
- public class StateCarryingChunkJob extends ChunkJob implements Async.ChunkResettable {
- public List seen = new List();
- public List pagesReset = new List();
-
- public override void work(List chunk) {
- insert new Account(Name = 'CHUNK-SIZE-' + seen.size());
- seen.add('page');
- }
+ @IsTest
+ private static void shouldStillRunChainedJobsWhenSourceIsEmpty() {
+ Test.startTest();
+ Async.queueable(new ProcessedCountMarkerJob('EMPTY_SOURCE'))
+ .chunk(new MarkingChunkJob(), ChunkSource.of(new List()))
+ .enqueue();
+ Test.stopTest();
- public void resetBeforeNextChunk(Integer pageNumber) {
- pagesReset.add(pageNumber);
- }
+ Assert.areEqual(
+ 1,
+ [SELECT COUNT() FROM Account WHERE Name = 'EMPTY_SOURCE'],
+ 'An empty chunk source must not swallow the jobs already chained.'
+ );
}
- private class LegacyResetForRetryJob extends QueueableJob implements Async.Retryable {
- public Boolean legacyHookRan = false;
- public override void work() {
- throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
- }
- public override void resetForRetry() {
- this.legacyHookRan = true;
- }
- public void resetBeforeRetry(Integer attempt) {
+ @IsTest
+ private static void shouldRejectInvalidChunkSize() {
+ ChunkBuilder builder = Async.chunk(
+ new MarkingChunkJob(),
+ ChunkSource.of(createAccounts(1))
+ );
+ for (Integer invalid : new List{ 0, -1, null }) {
+ try {
+ builder.chunkSize(invalid);
+ Assert.fail('chunkSize ' + invalid + ' should be rejected.');
+ } catch (Exception ex) {
+ Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_INVALID_CHUNK_SIZE, ex.getMessage());
+ }
}
}
- private class QueueableTestFinalizer extends QueueableJob.Finalizer {
- public override void work() {
- FinalizerContext finalizerCtx = Async.getQueueableJobContext()?.finalizerCtx;
- insert new Account(
- Name = Async.getQueueableJobContext()?.currentJob?.uniqueName,
- Description = finalizerCtx?.getResult() == ParentJobResult.SUCCESS
- ? 'Success'
- : finalizerCtx?.getException()?.getMessage()
+ @IsTest
+ private static void shouldRejectChunkSizeAboveTheCursorFetchCap() {
+ createAccounts(1);
+ ChunkBuilder builder = Async.chunk(
+ new MarkingChunkJob(),
+ ChunkSource.query('SELECT Id FROM Account')
+ );
+
+ try {
+ builder.chunkSize(2001);
+ Assert.fail('A cursor source cannot page more than 2000 records at a time.');
+ } catch (Exception ex) {
+ Assert.areEqual(
+ 'chunkSize must be between 1 and 2000 for this source',
+ ex.getMessage()
);
}
+ Assert.areNotEqual(
+ null,
+ builder.chunkSize(2000),
+ 'A chunk size at the cursor fetch cap is allowed.'
+ );
}
- private class FinalizerErrorQueueableTest extends QueueableJob {
- public override void work() {
- Async.queueable(new SuccessfulQueueableTest()).attachFinalizer();
- }
- }
+ @IsTest
+ private static void shouldNotCapChunkSizeForAnInMemorySource() {
+ ChunkBuilder builder = Async.chunk(
+ new MarkingChunkJob(),
+ ChunkSource.of(createAccounts(1))
+ );
- private class ChainedQueueableJob extends QueueableJob {
- private Integer chainDepthLimit = 1;
- private Integer currentChainDepth = 0;
+ Assert.areEqual(
+ null,
+ ChunkSource.of(new List()).maxChunkSize(),
+ 'An in-memory source has no cursor fetch ceiling.'
+ );
+ Assert.areNotEqual(
+ null,
+ builder.chunkSize(5000),
+ 'An in-memory source is bounded by the job body, not by a page ceiling.'
+ );
+ }
- public ChainedQueueableJob(Integer chainDepthLimit) {
- this.chainDepthLimit = chainDepthLimit;
- }
+ @IsTest
+ private static void shouldExposeCursorFetchCapOnCursorSource() {
+ createAccounts(1);
- public override void work() {
- currentChainDepth++;
- insert new Account(Name = Async.getQueueableJobContext()?.currentJob?.uniqueName);
- if (currentChainDepth < chainDepthLimit) {
- Async.queueable(this).enqueue();
- }
- }
+ Assert.areEqual(2000, ChunkSource.query('SELECT Id FROM Account').maxChunkSize());
}
- public class QueueableJobTest1 extends QueueableJob implements Async.Retryable {
- public QueueableJobTest2 complexMember = new QueueableJobTest2();
- public override void work() {
+ @IsTest
+ private static void shouldRejectNullJobAndSource() {
+ try {
+ Async.chunk(null, ChunkSource.of(createAccounts(1)));
+ Assert.fail('A null job should be rejected.');
+ } catch (Exception ex) {
+ Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_NULL_JOB, ex.getMessage());
}
-
- public void resetBeforeRetry(Integer attempt) {
+ try {
+ Async.chunk(new MarkingChunkJob(), null);
+ Assert.fail('A null source should be rejected.');
+ } catch (Exception ex) {
+ Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_NULL_SOURCE, ex.getMessage());
}
}
- private class QueueableJobTest2 extends QueueableJob {
- public String primitiveMember = PRIMITIVE_VALUE_INITIAL;
- public override void work() {
+ @IsTest
+ private static void shouldRejectRetryAboveCapOnChunk() {
+ ChunkBuilder builder = Async.chunk(
+ new MarkingChunkJob(),
+ ChunkSource.of(createAccounts(1))
+ );
+ try {
+ builder.retry(QueueableManager.MAX_RETRY_CAP + 1);
+ Assert.fail('Retry above the cap should be rejected.');
+ } catch (Exception ex) {
+ Assert.areEqual(
+ QueueableManager.ERROR_MESSAGE_MAX_RETRIES_EXCEEDS_CAP,
+ ex.getMessage()
+ );
}
}
- private class QueueableJobTest3 extends QueueableJob {
- public override void work() {
- }
- }
+ @IsTest
+ private static void shouldAccumulateJobStateAcrossChunks() {
+ List accounts = createAccounts(6);
- private class QueueableJobTest4 extends QueueableJob {
- public override void work() {
- }
- }
+ Test.startTest();
+ Async.chunk(new TotallingChunkJob(), ChunkSource.of(accounts)).chunkSize(2).enqueue();
+ Test.stopTest();
- private class QueueableJobTest5 extends QueueableJob {
- public override void work() {
+ Set runningTotals = new Set();
+ for (Account marker : [SELECT Site FROM Account WHERE Name = 'RUNNING_TOTAL']) {
+ runningTotals.add(marker.Site);
}
+ Assert.areEqual(
+ new Set{ '2', '4', '6' },
+ runningTotals,
+ 'A member mutated in work() keeps accumulating on every later page.'
+ );
}
- private class QueueableJobTest6 extends QueueableJob {
- public override void work() {
- }
- }
-
- private class QueueableJobTest7 extends QueueableJob {
- public override void work() {
- }
- }
-
- private class QueueableJobTest8 extends QueueableJob {
- public override void work() {
- }
- }
-
- private class SchedulableTest implements Schedulable {
- public void execute(SchedulableContext ctx) {
- }
- }
-
- private class CustomException extends Exception {
- }
-
- public class ParentJobWithFinalizer extends QueueableJob {
- private String mockId;
-
- public ParentJobWithFinalizer(String mockId) {
- this.mockId = mockId;
- }
-
- public override void work() {
- Async.queueable(new ErrorHandlerFinalizer()).mockId(mockId).attachFinalizer();
- }
- }
-
- public class ErrorHandlerFinalizer extends QueueableJob.Finalizer {
- public override void work() {
- FinalizerContext ctx = this.finalizerCtx;
- if (ctx?.getResult() == ParentJobResult.UNHANDLED_EXCEPTION) {
- insert new Account(
- Name = 'Error Log',
- Description = ctx.getException()?.getMessage()
- );
- }
- }
- }
-
- public class AccountCreatorJob extends QueueableJob {
- private String accountName;
-
- public AccountCreatorJob(String accountName) {
- this.accountName = accountName;
- }
-
- public override void work() {
- Id jobId = this.queueableCtx?.getJobId();
- insert new Account(Name = accountName, Description = 'Job: ' + jobId);
+ @IsTest
+ private static void shouldPageRecordsFromAConsumerDefinedSource() {
+ Test.startTest();
+ Async.chunk(new CountingChunkJob(), new SyntheticChunkSource(2500))
+ .chunkSize(1000)
+ .enqueue();
+ Test.stopTest();
+
+ Set runningTotals = new Set();
+ for (Account marker : [SELECT Site FROM Account WHERE Name = 'SYNTHETIC_TOTAL']) {
+ runningTotals.add(marker.Site);
}
+ Assert.areEqual(
+ new Set{ '1000', '2000', '2500' },
+ runningTotals,
+ 'A custom ChunkSource pages 2500 fabricated records without any DML behind it.'
+ );
}
@IsTest
- private static void shouldProcessFirstChunkThroughEnqueue() {
- List accounts = createAccounts(5);
+ private static void shouldInjectQueueableMockIntoAChunkPage() {
+ List accounts = createAccounts(2);
+ Id mockJobId = fakeAsyncApexJobId();
+ AsyncMock.whenQueueable('chunk-page')
+ .thenReturn(new AsyncMock.MockQueueableContext().setJobId(mockJobId));
Test.startTest();
- Async.Result result = Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(5)
+ Async.chunk(new ContextReadingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .mockId('chunk-page')
.enqueue();
Test.stopTest();
- Assert.areNotEqual(null, result.salesforceJobId);
Assert.areEqual(
- 5,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'The single chunk should process every record in the source.'
+ String.valueOf(mockJobId),
+ [SELECT Site FROM Account WHERE Name = 'MOCKED_CONTEXT' LIMIT 1].Site,
+ 'A chunk page reads the mocked QueueableContext like any other job.'
);
}
@IsTest
- private static void shouldProcessAllInMemoryChunksAcrossPages() {
- List accounts = createAccounts(6);
+ private static void shouldFailAJobFromAQueueableMock() {
+ QueueableChain chain = new QueueableChain();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
+ AsyncMock.whenQueueable('mocked-failure')
+ .thenThrow(new CustomException(CUSTOM_ERROR_MESSAGE));
Test.startTest();
- Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts)).chunkSize(3).enqueue();
+ QueueableManager.get().setChain(chain);
+ Async.queueable(new AccountCreatorJob('never runs'))
+ .mockId('mocked-failure')
+ .continueOnJobExecuteFail()
+ .enqueue();
Test.stopTest();
+ AsyncResult__c result = [
+ SELECT Status__c, ExceptionType__c, ExceptionMessage__c
+ FROM AsyncResult__c
+ LIMIT 1
+ ];
Assert.areEqual(
- 6,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'Self-chaining must process every page of the run.'
+ QueueableManager.STATUS_FAILED,
+ result.Status__c,
+ 'A mocked failure settles through the same path a real one does.'
+ );
+ Assert.areEqual(CUSTOM_ERROR_MESSAGE, result.ExceptionMessage__c);
+ Assert.areEqual(
+ 0,
+ [SELECT COUNT() FROM Account],
+ 'The job body never runs when its execution is mocked to throw.'
);
}
@IsTest
- private static void shouldProcessAllCursorChunksAcrossPages() {
- createAccounts(6);
+ private static void shouldFailOnlyTheMockedChunkPage() {
+ List accounts = createAccounts(6);
+ AsyncMock.whenQueueable('recalc-run')
+ .thenReturn(new AsyncMock.MockQueueableContext())
+ .thenThrow(new CustomException(CUSTOM_ERROR_MESSAGE))
+ .thenReturn(new AsyncMock.MockQueueableContext());
Test.startTest();
- Async.chunk(new MarkingChunkJob(), ChunkSource.query('SELECT Id FROM Account'))
- .chunkSize(3)
+ Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .mockId('recalc-run')
.enqueue();
Test.stopTest();
Assert.areEqual(
- 6,
+ 4,
[SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'A cursor must survive re-enqueue and fetch every page.'
+ 'The mock queue picks which page fails; the rest of the run still processes.'
);
}
@IsTest
- private static void shouldCarryJobMemberStateAcrossChunks() {
+ private static void shouldHaltRunFromAMockedPageFailure() {
List accounts = createAccounts(6);
+ AsyncMock.whenQueueable('recalc-run')
+ .thenReturn(new AsyncMock.MockQueueableContext())
+ .thenThrow(new CustomException(CUSTOM_ERROR_MESSAGE));
Test.startTest();
- Async.chunk(new LabelingChunkJob('RUN_LABEL'), ChunkSource.of(accounts))
+ Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
.chunkSize(2)
+ .mockId('recalc-run')
+ .stopRemainingChunksOnFailure()
.enqueue();
Test.stopTest();
Assert.areEqual(
- 6,
- [SELECT COUNT() FROM Account WHERE Site = 'RUN_LABEL'],
- 'Job member state must survive cloning and serialization across every page.'
+ 2,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'A mocked page failure halts the run like a real one.'
);
}
@IsTest
- private static void shouldContinueToNextChunkAfterAFailedChunk() {
- List accounts = accountsWithOneFailingRecord();
+ private static void shouldMockFinalizerOutcomeForAChunkPage() {
+ List accounts = createAccounts(2);
+ AsyncMock.whenFinalizer('chunk-error-handler').thenThrow(new DmlException('Page blew up'));
Test.startTest();
- Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts)).chunkSize(2).enqueue();
+ Async.chunk(new FinalizerAttachingChunkJob('chunk-error-handler'), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .enqueue();
Test.stopTest();
Assert.areEqual(
- 4,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'A failed chunk must not stop the remaining chunks by default.'
+ 'Page blew up',
+ [SELECT Description FROM Account WHERE Name = 'Error Log' LIMIT 1].Description,
+ 'A page finalizer can be told its page failed without the page actually failing.'
);
}
@IsTest
- private static void shouldStopAfterFailedChunkWhenConfigured() {
- List accounts = new List{
- new Account(Name = 'FAIL first'),
- new Account(Name = 'ok second'),
- new Account(Name = 'ok third'),
- new Account(Name = 'ok fourth')
- };
- insert accounts;
+ private static void shouldFailAndRetryAPageWhenTheSourceThrows() {
+ QueueableChain chain = new QueueableChain();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
Test.startTest();
- Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts))
+ QueueableManager.get().setChain(chain);
+ Async.chunk(new MarkingChunkJob(), new ExplodingChunkSource())
.chunkSize(2)
- .stopRemainingChunksOnFailure()
+ .retry(1)
.enqueue();
Test.stopTest();
+ AsyncResult__c result = [
+ SELECT Status__c, ExceptionMessage__c, RetryAttempts__c
+ FROM AsyncResult__c
+ LIMIT 1
+ ];
Assert.areEqual(
- 0,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'stopRemainingChunksOnFailure must halt the run after a failed chunk.'
+ QueueableManager.STATUS_FAILED,
+ result.Status__c,
+ 'A source that throws fails the page instead of silently ending the run.'
);
+ Assert.areEqual(CUSTOM_ERROR_MESSAGE, result.ExceptionMessage__c);
+ Assert.areEqual(1, result.RetryAttempts__c, 'The page retried before it settled.');
}
@IsTest
- private static void shouldApplyAllChunkBuilderOptions() {
+ private static void shouldChainRunWithoutEnqueueingIt() {
List accounts = createAccounts(2);
- Test.startTest();
Async.Result result = Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
.chunkSize(2)
- .priority(3)
- .delay(1)
- .delayBetweenChunks(1)
- .retry(2)
- .backoff(Backoff.fixed(1))
- .retryOn(CustomException.class)
- .mockId('chunk-mock')
- .keepChunkPages()
- .stopRemainingChunksOnFailure()
- .enqueue();
- Test.stopTest();
+ .chain();
Assert.areNotEqual(
null,
- result.salesforceJobId,
- 'Every builder option should still enqueue.'
+ result.customJobId,
+ 'chain() returns the run handle so later jobs can depend on it.'
+ );
+ Assert.areEqual(
+ 0,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'chain() adds the run to the chain without enqueuing anything.'
);
}
@IsTest
- private static void shouldProcessFirstChunkInBulk() {
- List accounts = createAccounts(200);
+ private static void shouldRunASecondChunkRunAfterTheFirstOneFinishes() {
+ List accounts = createAccounts(4);
Test.startTest();
- Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts)).chunkSize(200).enqueue();
+ Async.chunk(new LabelingChunkJob('FIRST'), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .chunk(new LabelingChunkJob('SECOND'), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .enqueue();
Test.stopTest();
Assert.areEqual(
- 200,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'A 200-record chunk should process in one bulk-safe page.'
+ 4,
+ [SELECT COUNT() FROM Account WHERE Site = 'SECOND'],
+ 'The second run overwrites the first, so it ran after every page of run one.'
);
}
@IsTest
- private static void shouldFetchCorrectOffsetForLaterChunk() {
- List accounts = createAccounts(10);
- MarkingChunkJob job = new MarkingChunkJob();
- job.getRun().configure(ChunkSource.of(accounts), 5, false, null, false);
- job.getRun().position = 5;
+ private static void shouldGateASecondChunkRunOnTheFirstRunOutcome() {
+ List accounts = accountsWithOneFailingRecord();
- job.work();
+ Test.startTest();
+ Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .chunk(new LabelingChunkJob('SECOND'), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .dependsOn(Async.afterPrevious().succeeded())
+ .enqueue();
+ Test.stopTest();
- Set laterHalf = new Set();
- for (Integer i = 5; i < 10; i++) {
- laterHalf.add(accounts[i].Id);
- }
- List processed = [SELECT Id FROM Account WHERE Site = :CHUNK_PROCESSED];
Assert.areEqual(
- 5,
- processed.size(),
- 'Only the second page of records should be processed.'
+ 0,
+ [SELECT COUNT() FROM Account WHERE Site = 'SECOND'],
+ 'dependsOn(previous) resolves against the first run outcome, which failed.'
);
- for (Account processedAccount : processed) {
- Assert.isTrue(
- laterHalf.contains(processedAccount.Id),
- 'Offset fetch must return records from position 5 onward.'
- );
- }
}
@IsTest
- private static void shouldRejectChunkJobRunningOutsideAChunkRun() {
- try {
- new MarkingChunkJob().work();
- Assert.fail('A ChunkJob without a configured run must not execute.');
- } catch (IllegalArgumentException ex) {
- Assert.areEqual(ChunkRun.ERROR_MESSAGE_RUN_NOT_STARTED, ex.getMessage());
- }
- }
-
- @IsTest
- private static void shouldCreateResultRowForChunk() {
- List accounts = createAccounts(3);
- QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
-
- Test.startTest();
- QueueableManager.get().setChain(chain);
- Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts)).chunkSize(3).enqueue();
- Test.stopTest();
-
- List results = [
- SELECT Id, Status__c, ChainId__c, ClassName__c
- FROM AsyncResult__c
- ];
- Assert.areEqual(1, results.size(), 'A completed chunk page should record one AsyncResult.');
- Assert.areEqual(QueueableManager.STATUS_COMPLETED, results[0].Status__c);
- Assert.areNotEqual(null, results[0].ChainId__c);
- }
-
- @IsTest
- private static void shouldReturnNextChunkWhenRecordsRemain() {
- MarkingChunkJob job = new MarkingChunkJob();
- job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
-
- ChunkJob nextPage = job.nextPageOrNull();
-
- Assert.areNotEqual(null, nextPage, 'A next page is expected while records remain.');
- Assert.areEqual(5, nextPage.getRun().position, 'The next page must advance by chunkSize.');
- }
-
- @IsTest
- private static void shouldNotCarryInitialDelayToLaterChunks() {
- MarkingChunkJob job = new MarkingChunkJob();
- job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
- job.delay = 5;
+ private static void shouldOpenACursorInTheRequestedAccessLevel() {
+ createAccounts(2);
- ChunkJob nextPage = job.nextPageOrNull();
+ ChunkSource userMode = ChunkSource.query('SELECT Id FROM Account', AccessLevel.USER_MODE);
+ ChunkSource systemMode = ChunkSource.query(
+ 'SELECT Id FROM Account WHERE Name LIKE :prefix',
+ new Map{ 'prefix' => 'Chunk%' },
+ AccessLevel.SYSTEM_MODE
+ );
- Assert.areEqual(null, nextPage.delay, 'An initial delay must not repeat on every page.');
+ Assert.areEqual(2, userMode.getNumRecords(), 'The admin running the test sees both.');
+ Assert.areEqual(2, systemMode.getNumRecords(), 'Binds and access level combine.');
}
@IsTest
- private static void shouldApplyDelayBetweenChunks() {
- MarkingChunkJob job = new MarkingChunkJob();
- job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, 3, false);
-
- ChunkJob nextPage = job.nextPageOrNull();
+ private static void shouldChainNothingWhenSourceIsEmpty() {
+ Async.Result result = Async.chunk(
+ new MarkingChunkJob(),
+ ChunkSource.of(new List())
+ )
+ .chain();
- Assert.areEqual(3, nextPage.delay, 'delayBetweenChunks must throttle each later page.');
+ Assert.areEqual(null, result.customJobId, 'An empty source chains no run.');
+ Assert.isTrue(result.queueableChainState.jobs.isEmpty());
}
@IsTest
- private static void shouldRejectDeepCloneChunkJob() {
- MarkingChunkJob job = new MarkingChunkJob();
- job.deepClone = true;
-
- try {
- Async.chunk(job, ChunkSource.of(createAccounts(2)));
- Assert.fail('A deepClone ChunkJob should be rejected at build time.');
- } catch (Exception ex) {
- Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_DEEP_CLONE_UNSUPPORTED, ex.getMessage());
- }
- }
+ private static void shouldRunChunkRunWhenItsDependencySucceeded() {
+ List accounts = createAccounts(4);
- @IsTest
- private static void shouldPruneSettledPagesByDefault() {
- MarkingChunkJob job = new MarkingChunkJob();
- job.getRun().configure(ChunkSource.of(createAccounts(4)), 2, false, null, false);
+ Test.startTest();
+ Async.queueable(new ProcessedCountMarkerJob('GATE'))
+ .chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .dependsOn(Async.afterPrevious().succeeded())
+ .enqueue();
+ Test.stopTest();
- Assert.isFalse(job.getRun().keepPages, 'Settled pages are pruned by default.');
+ Assert.areEqual(
+ 4,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'The run starts once the job it depends on succeeded.'
+ );
}
@IsTest
- private static void shouldKeepSettledPagesWhenRequested() {
- List accounts = createAccounts(6);
+ private static void shouldSkipChunkRunWhenItsDependencyFailed() {
+ List accounts = createAccounts(4);
Test.startTest();
- Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ Async.queueable(new FailureQueueableTest())
+ .continueOnJobExecuteFail()
+ .chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
.chunkSize(2)
- .keepChunkPages()
+ .dependsOn(Async.afterPrevious().succeeded())
.enqueue();
Test.stopTest();
Assert.areEqual(
- 6,
+ 0,
[SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'A kept-pages run still processes every page.'
+ 'The whole run is skipped when the job it depends on failed.'
);
}
@IsTest
- private static void shouldReturnNoNextChunkAtEndOfSource() {
- MarkingChunkJob job = new MarkingChunkJob();
- job.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
- job.getRun().position = 5;
-
- Assert.areEqual(null, job.nextPageOrNull(), 'No page should follow the final chunk.');
+ private static void shouldRejectInvalidDependencyOnChunk() {
+ ChunkBuilder builder = Async.chunk(
+ new MarkingChunkJob(),
+ ChunkSource.of(createAccounts(1))
+ );
+ try {
+ builder.dependsOn(null);
+ Assert.fail('A dependency without an outcome should be rejected.');
+ } catch (Exception ex) {
+ Assert.areEqual(QueueableManager.ERROR_MESSAGE_INVALID_DEPENDENCY, ex.getMessage());
+ }
+ try {
+ builder.dependsOn(Async.afterPrevious().succeeded());
+ Assert.fail('afterPrevious() without a previously chained job should be rejected.');
+ } catch (Exception ex) {
+ Assert.areEqual(
+ QueueableManager.ERROR_MESSAGE_DEPENDS_ON_PREVIOUS_WITHOUT_JOB,
+ ex.getMessage()
+ );
+ }
}
@IsTest
- private static void shouldHaltRunAfterFailedPageWhenConfigured() {
- MarkingChunkJob stopping = new MarkingChunkJob();
- stopping.getRun().configure(ChunkSource.of(createAccounts(10)), 5, true, null, false);
+ private static void shouldYieldIdOnlyShellsWhenSourceIsBuiltFromIds() {
+ List accounts = createAccounts(5);
+ ChunkSource source = ChunkSource.ofIds(new Map(accounts).keySet());
- Assert.isTrue(
- stopping.getRun().isHaltedBy(true),
- 'A failed page must halt the run when stopRemainingChunksOnFailure is set.'
- );
- Assert.isFalse(
- stopping.getRun().isHaltedBy(false),
- 'A page that succeeded must not halt the run.'
+ List page = source.fetch(0, 5);
+
+ Assert.areEqual(5, page.size());
+ Assert.areEqual(
+ new Map(accounts).keySet(),
+ new Map(page).keySet(),
+ 'A job reads ids straight off the page with new Map(chunk).keySet().'
);
}
@IsTest
- private static void shouldContinueRemainingChunksOnFailureByDefault() {
- MarkingChunkJob continuing = new MarkingChunkJob();
- continuing.getRun().configure(ChunkSource.of(createAccounts(10)), 5, false, null, false);
+ private static void shouldExposeInMemorySourceRecordsAndSize() {
+ List accounts = createAccounts(5);
+ ChunkSource source = ChunkSource.of(accounts);
- Assert.isFalse(
- continuing.getRun().isHaltedBy(true),
- 'A failed chunk must not halt the run by default.'
- );
- Assert.areNotEqual(
- null,
- continuing.nextPageOrNull(),
- 'The run continues to the next page after a failed page.'
+ Assert.areEqual(5, source.getNumRecords());
+ Assert.areEqual(2, source.fetch(0, 2).size());
+ Assert.areEqual(
+ 1,
+ source.fetch(4, 5).size(),
+ 'A short final page must clamp to the remainder.'
);
}
@IsTest
- private static void shouldCarryFailedPageFlagToLaterPages() {
- MarkingChunkJob job = new MarkingChunkJob();
- job.getRun().configure(ChunkSource.of(createAccounts(6)), 2, false, null, false);
- job.hasFailed = true;
+ private static void shouldBuildIdSourceRecordsWithIds() {
+ List accounts = createAccounts(3);
+ Set ids = new Map(accounts).keySet();
+ ChunkSource source = ChunkSource.ofIds(ids);
- ChunkJob nextPage = job.nextPageOrNull();
+ Assert.areEqual(3, source.getNumRecords());
+ Set fetched = new Set();
+ for (SObject record : source.fetch(0, 3)) {
+ fetched.add(record.Id);
+ }
+ Assert.areEqual(ids, fetched);
+ }
- Assert.isTrue(
- nextPage.getRun().hasFailedPage,
- 'The run remembers a failed page so the run outcome stays FAILURE.'
- );
- Assert.isFalse(nextPage.hasFailed, 'The next page starts clean.');
+ @IsTest
+ private static void shouldTreatNullRecordsAsEmptySource() {
+ Assert.areEqual(0, ChunkSource.ofIds(null).getNumRecords());
+ Assert.areEqual(0, ChunkSource.of(null).getNumRecords());
}
@IsTest
- private static void shouldTreatEmptySourceAsNoOp() {
- Async.Result result = Async.chunk(
- new MarkingChunkJob(),
- ChunkSource.of(new List())
- )
- .enqueue();
+ private static void shouldDelegateToCursorSource() {
+ createAccounts(5);
+ ChunkSource source = ChunkSource.cursor(Database.getCursor('SELECT Id FROM Account'));
- Assert.areEqual(null, result.salesforceJobId, 'An empty source should enqueue nothing.');
- Assert.isTrue(result.queueableChainState.jobs.isEmpty());
+ Assert.areEqual(5, source.getNumRecords());
+ Assert.areEqual(2, source.fetch(0, 2).size());
}
@IsTest
- private static void shouldStillRunChainedJobsWhenSourceIsEmpty() {
- Test.startTest();
- Async.queueable(new ProcessedCountMarkerJob('EMPTY_SOURCE'))
- .chunk(new MarkingChunkJob(), ChunkSource.of(new List()))
- .enqueue();
- Test.stopTest();
+ private static void shouldOpenCursorFromQuery() {
+ createAccounts(4);
+ ChunkSource source = ChunkSource.query('SELECT Id FROM Account');
- Assert.areEqual(
- 1,
- [SELECT COUNT() FROM Account WHERE Name = 'EMPTY_SOURCE'],
- 'An empty chunk source must not swallow the jobs already chained.'
- );
+ Assert.areEqual(4, source.getNumRecords());
}
@IsTest
- private static void shouldRejectInvalidChunkSize() {
- ChunkBuilder builder = Async.chunk(
- new MarkingChunkJob(),
- ChunkSource.of(createAccounts(1))
+ private static void shouldOpenCursorFromQueryWithBinds() {
+ createAccounts(4);
+ ChunkSource source = ChunkSource.query(
+ 'SELECT Id FROM Account WHERE Name = :accountName',
+ new Map{ 'accountName' => 'Chunk 1' }
);
- for (Integer invalid : new List{ 0, -1, null }) {
- try {
- builder.chunkSize(invalid);
- Assert.fail('chunkSize ' + invalid + ' should be rejected.');
- } catch (Exception ex) {
- Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_INVALID_CHUNK_SIZE, ex.getMessage());
- }
- }
+
+ Assert.areEqual(1, source.getNumRecords(), 'Bind variables must reach the cursor.');
}
@IsTest
- private static void shouldRejectChunkSizeAboveTheCursorFetchCap() {
- createAccounts(1);
- ChunkBuilder builder = Async.chunk(
- new MarkingChunkJob(),
- ChunkSource.query('SELECT Id FROM Account')
- );
-
+ private static void shouldRejectNullCursorAndBlankQuery() {
try {
- builder.chunkSize(2001);
- Assert.fail('A cursor source cannot page more than 2000 records at a time.');
- } catch (Exception ex) {
- Assert.areEqual(
- 'chunkSize must be between 1 and 2000 for this source',
- ex.getMessage()
- );
+ ChunkSource.cursor(null);
+ Assert.fail('A null cursor should be rejected.');
+ } catch (Exception ex) {
+ Assert.areEqual(ChunkSource.ERROR_MESSAGE_NULL_CURSOR, ex.getMessage());
}
- Assert.areNotEqual(
- null,
- builder.chunkSize(2000),
- 'A chunk size at the cursor fetch cap is allowed.'
- );
- }
-
- @IsTest
- private static void shouldNotCapChunkSizeForAnInMemorySource() {
- ChunkBuilder builder = Async.chunk(
- new MarkingChunkJob(),
- ChunkSource.of(createAccounts(1))
- );
-
- Assert.areEqual(
- null,
- ChunkSource.of(new List()).maxChunkSize(),
- 'An in-memory source has no cursor fetch ceiling.'
- );
- Assert.areNotEqual(
- null,
- builder.chunkSize(5000),
- 'An in-memory source is bounded by the job body, not by a page ceiling.'
- );
- }
-
- @IsTest
- private static void shouldExposeCursorFetchCapOnCursorSource() {
- createAccounts(1);
-
- Assert.areEqual(2000, ChunkSource.query('SELECT Id FROM Account').maxChunkSize());
- }
-
- @IsTest
- private static void shouldRejectNullJobAndSource() {
try {
- Async.chunk(null, ChunkSource.of(createAccounts(1)));
- Assert.fail('A null job should be rejected.');
+ ChunkSource.query(' ');
+ Assert.fail('A blank query should be rejected.');
} catch (Exception ex) {
- Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_NULL_JOB, ex.getMessage());
+ Assert.areEqual(ChunkSource.ERROR_MESSAGE_BLANK_QUERY, ex.getMessage());
}
try {
- Async.chunk(new MarkingChunkJob(), null);
- Assert.fail('A null source should be rejected.');
+ ChunkSource.query(' ', new Map());
+ Assert.fail('A blank query with binds should be rejected.');
} catch (Exception ex) {
- Assert.areEqual(ChunkBuilder.ERROR_MESSAGE_NULL_SOURCE, ex.getMessage());
+ Assert.areEqual(ChunkSource.ERROR_MESSAGE_BLANK_QUERY, ex.getMessage());
}
}
@IsTest
- private static void shouldRejectRetryAboveCapOnChunk() {
- ChunkBuilder builder = Async.chunk(
- new MarkingChunkJob(),
- ChunkSource.of(createAccounts(1))
- );
- try {
- builder.retry(QueueableManager.MAX_RETRY_CAP + 1);
- Assert.fail('Retry above the cap should be rejected.');
- } catch (Exception ex) {
- Assert.areEqual(
- QueueableManager.ERROR_MESSAGE_MAX_RETRIES_EXCEEDS_CAP,
- ex.getMessage()
- );
- }
+ private static void shouldOrderJobsOfEqualPriorityByChainSequence() {
+ QueueableJobTest1 first = new QueueableJobTest1();
+ first.uniqueName = 'first';
+ first.chainSequence = 1;
+ QueueableJobTest2 second = new QueueableJobTest2();
+ second.uniqueName = 'second';
+ second.chainSequence = 2;
+
+ List jobs = new List{ second, first };
+ jobs.sort();
+
+ Assert.areEqual('first', jobs[0].uniqueName, 'Equal priority keeps the chained order.');
+ Assert.areEqual('second', jobs[1].uniqueName);
}
@IsTest
- private static void shouldAccumulateJobStateAcrossChunks() {
+ private static void shouldRunEveryChunkBeforeTheNextChainedJob() {
List accounts = createAccounts(6);
Test.startTest();
- Async.chunk(new TotallingChunkJob(), ChunkSource.of(accounts)).chunkSize(2).enqueue();
+ Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
+ .chain(new ProcessedCountMarkerJob('AFTER_RUN'))
+ .enqueue();
Test.stopTest();
- Set runningTotals = new Set();
- for (Account marker : [SELECT Site FROM Account WHERE Name = 'RUNNING_TOTAL']) {
- runningTotals.add(marker.Site);
- }
Assert.areEqual(
- new Set{ '2', '4', '6' },
- runningTotals,
- 'A member mutated in work() keeps accumulating on every later page.'
+ '6',
+ processedCountSeenBy('AFTER_RUN'),
+ 'A job chained after the run must wait for every page.'
);
}
@IsTest
- private static void shouldPageRecordsFromAConsumerDefinedSource() {
+ private static void shouldRunChunkRunAfterAnEarlierChainedJob() {
+ List accounts = createAccounts(4);
+
Test.startTest();
- Async.chunk(new CountingChunkJob(), new SyntheticChunkSource(2500))
- .chunkSize(1000)
+ Async.queueable(new ProcessedCountMarkerJob('BEFORE_RUN'))
+ .chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ .chunkSize(2)
.enqueue();
Test.stopTest();
- Set runningTotals = new Set();
- for (Account marker : [SELECT Site FROM Account WHERE Name = 'SYNTHETIC_TOTAL']) {
- runningTotals.add(marker.Site);
- }
Assert.areEqual(
- new Set{ '1000', '2000', '2500' },
- runningTotals,
- 'A custom ChunkSource pages 2500 fabricated records without any DML behind it.'
+ '0',
+ processedCountSeenBy('BEFORE_RUN'),
+ 'A job chained before the run must run first.'
);
+ Assert.areEqual(4, [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED]);
}
@IsTest
- private static void shouldInjectQueueableMockIntoAChunkPage() {
- List accounts = createAccounts(2);
- Id mockJobId = fakeAsyncApexJobId();
- AsyncMock.whenQueueable('chunk-page')
- .thenReturn(new AsyncMock.MockQueueableContext().setJobId(mockJobId));
+ private static void shouldLetHigherPriorityJobPreemptTheNextChunk() {
+ List accounts = createAccounts(6);
+ PreemptingChunkJob job = new PreemptingChunkJob();
+ job.preemptPriority = 1;
Test.startTest();
- Async.chunk(new ContextReadingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .mockId('chunk-page')
- .enqueue();
+ Async.chunk(job, ChunkSource.of(accounts)).chunkSize(2).priority(5).enqueue();
Test.stopTest();
Assert.areEqual(
- String.valueOf(mockJobId),
- [SELECT Site FROM Account WHERE Name = 'MOCKED_CONTEXT' LIMIT 1].Site,
- 'A chunk page reads the mocked QueueableContext like any other job.'
+ '2',
+ processedCountSeenBy('PREEMPT'),
+ 'A higher priority job added mid-run runs before the next page.'
+ );
+ Assert.areEqual(
+ 6,
+ [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
+ 'The run resumes after the preempting job.'
);
}
@IsTest
- private static void shouldFailAJobFromAQueueableMock() {
- QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
- AsyncMock.whenQueueable('mocked-failure')
- .thenThrow(new CustomException(CUSTOM_ERROR_MESSAGE));
+ private static void shouldRunLowerPriorityJobAfterTheWholeRun() {
+ List accounts = createAccounts(6);
+ PreemptingChunkJob job = new PreemptingChunkJob();
+ job.preemptPriority = 9;
Test.startTest();
- QueueableManager.get().setChain(chain);
- Async.queueable(new AccountCreatorJob('never runs'))
- .mockId('mocked-failure')
- .continueOnJobExecuteFail()
- .enqueue();
+ Async.chunk(job, ChunkSource.of(accounts)).chunkSize(2).priority(5).enqueue();
Test.stopTest();
- AsyncResult__c result = [
- SELECT Status__c, ExceptionType__c, ExceptionMessage__c
- FROM AsyncResult__c
- LIMIT 1
- ];
- Assert.areEqual(
- QueueableManager.STATUS_FAILED,
- result.Status__c,
- 'A mocked failure settles through the same path a real one does.'
- );
- Assert.areEqual(CUSTOM_ERROR_MESSAGE, result.ExceptionMessage__c);
Assert.areEqual(
- 0,
- [SELECT COUNT() FROM Account],
- 'The job body never runs when its execution is mocked to throw.'
+ '6',
+ processedCountSeenBy('PREEMPT'),
+ 'A lower priority job added mid-run waits for the whole run.'
);
}
@IsTest
- private static void shouldFailOnlyTheMockedChunkPage() {
+ private static void shouldRunDependentJobWhenEveryChunkSucceeded() {
List accounts = createAccounts(6);
- AsyncMock.whenQueueable('recalc-run')
- .thenReturn(new AsyncMock.MockQueueableContext())
- .thenThrow(new CustomException(CUSTOM_ERROR_MESSAGE))
- .thenReturn(new AsyncMock.MockQueueableContext());
Test.startTest();
Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
.chunkSize(2)
- .mockId('recalc-run')
+ .chain(new ProcessedCountMarkerJob('DEPENDENT'))
+ .dependsOn(Async.afterPrevious().succeeded())
.enqueue();
Test.stopTest();
Assert.areEqual(
- 4,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'The mock queue picks which page fails; the rest of the run still processes.'
+ '6',
+ processedCountSeenBy('DEPENDENT'),
+ 'The run outcome is SUCCESS when every page passed.'
);
}
@IsTest
- private static void shouldHaltRunFromAMockedPageFailure() {
- List accounts = createAccounts(6);
- AsyncMock.whenQueueable('recalc-run')
- .thenReturn(new AsyncMock.MockQueueableContext())
- .thenThrow(new CustomException(CUSTOM_ERROR_MESSAGE));
+ private static void shouldSkipDependentJobWhenAnyChunkFailed() {
+ List accounts = accountsWithOneFailingRecord();
Test.startTest();
- Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
+ Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts))
.chunkSize(2)
- .mockId('recalc-run')
- .stopRemainingChunksOnFailure()
+ .chain(new ProcessedCountMarkerJob('DEPENDENT'))
+ .dependsOn(Async.afterPrevious().succeeded())
.enqueue();
Test.stopTest();
Assert.areEqual(
- 2,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'A mocked page failure halts the run like a real one.'
+ 0,
+ [SELECT COUNT() FROM Account WHERE Name = 'DEPENDENT'],
+ 'One failed page makes the whole run outcome FAILURE.'
);
}
@IsTest
- private static void shouldMockFinalizerOutcomeForAChunkPage() {
- List accounts = createAccounts(2);
- AsyncMock.whenFinalizer('chunk-error-handler').thenThrow(new DmlException('Page blew up'));
+ private static void shouldRunDependentJobOnFinishedRunEvenAfterAFailedChunk() {
+ List accounts = accountsWithOneFailingRecord();
Test.startTest();
- Async.chunk(new FinalizerAttachingChunkJob('chunk-error-handler'), ChunkSource.of(accounts))
+ Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts))
.chunkSize(2)
+ .chain(new ProcessedCountMarkerJob('DEPENDENT'))
+ .dependsOn(Async.afterPrevious().finished())
.enqueue();
Test.stopTest();
Assert.areEqual(
- 'Page blew up',
- [SELECT Description FROM Account WHERE Name = 'Error Log' LIMIT 1].Description,
- 'A page finalizer can be told its page failed without the page actually failing.'
+ 1,
+ [SELECT COUNT() FROM Account WHERE Name = 'DEPENDENT'],
+ 'finished() runs the dependent job whatever the run outcome was.'
);
}
@IsTest
- private static void shouldFailAndRetryAPageWhenTheSourceThrows() {
+ private static void shouldRecordSummaryResultWhenRunHaltsOnFailure() {
+ List accounts = accountsWithOneFailingRecord();
QueueableChain chain = new QueueableChain();
- chain.queueableJobSettingByJobName = resultsEnabledForAll();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
Test.startTest();
QueueableManager.get().setChain(chain);
- Async.chunk(new MarkingChunkJob(), new ExplodingChunkSource())
+ Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts))
.chunkSize(2)
- .retry(1)
+ .stopRemainingChunksOnFailure()
.enqueue();
Test.stopTest();
- AsyncResult__c result = [
- SELECT Status__c, ExceptionMessage__c, RetryAttempts__c
+ List stopped = [
+ SELECT SkipReason__c
FROM AsyncResult__c
- LIMIT 1
+ WHERE Status__c = :QueueableManager.STATUS_SKIPPED_CHUNK_STOPPED
];
- Assert.areEqual(
- QueueableManager.STATUS_FAILED,
- result.Status__c,
- 'A source that throws fails the page instead of silently ending the run.'
+ Assert.areEqual(1, stopped.size(), 'A halted run records one summary result.');
+ Assert.isTrue(
+ stopped[0].SkipReason__c.contains('page 3 of 3'),
+ 'The summary states where the run stopped: ' + stopped[0].SkipReason__c
+ );
+ Assert.isTrue(
+ stopped[0].SkipReason__c.contains('2 record(s) were not processed'),
+ 'The summary states how much work was left: ' + stopped[0].SkipReason__c
);
- Assert.areEqual(CUSTOM_ERROR_MESSAGE, result.ExceptionMessage__c);
- Assert.areEqual(1, result.RetryAttempts__c, 'The page retried before it settled.');
}
@IsTest
- private static void shouldChainRunWithoutEnqueueingIt() {
- List accounts = createAccounts(2);
+ private static void shouldRecordSummaryResultWhenChainStopsDuringRun() {
+ List accounts = createAccounts(6);
+ QueueableChain chain = new QueueableChain();
+ QueueableChain.jobSettingByName = resultsEnabledForAll();
- Async.Result result = Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .chain();
+ Test.startTest();
+ QueueableManager.get().setChain(chain);
+ Async.chunk(new ChainStoppingChunkJob(), ChunkSource.of(accounts)).chunkSize(2).enqueue();
+ Test.stopTest();
- Assert.areNotEqual(
- null,
- result.customJobId,
- 'chain() returns the run handle so later jobs can depend on it.'
- );
Assert.areEqual(
- 0,
+ 2,
[SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'chain() adds the run to the chain without enqueuing anything.'
+ 'Async.stopChain() ends the run wherever it is.'
+ );
+ List stopped = [
+ SELECT SkipReason__c
+ FROM AsyncResult__c
+ WHERE
+ Status__c = :QueueableManager.STATUS_SKIPPED_CHAIN_STOPPED
+ AND SkipReason__c LIKE '%not processed%'
+ ];
+ Assert.areEqual(1, stopped.size(), 'A stopped chain records what the run left behind.');
+ Assert.isTrue(
+ stopped[0].SkipReason__c.contains('page 2 of 3'),
+ 'The summary points at the first page that never ran: ' + stopped[0].SkipReason__c
);
}
@IsTest
- private static void shouldRunASecondChunkRunAfterTheFirstOneFinishes() {
- List accounts = createAccounts(4);
+ private static void shouldStoreAPayloadWhenCustomMetadataTurnsItOn() {
+ RequeueableJob job = new RequeueableJob('STORED');
- Test.startTest();
- Async.chunk(new LabelingChunkJob('FIRST'), ChunkSource.of(accounts))
- .chunkSize(2)
- .chunk(new LabelingChunkJob('SECOND'), ChunkSource.of(accounts))
- .chunkSize(2)
- .enqueue();
- Test.stopTest();
+ chainStoringPayloads().addJob(job);
- Assert.areEqual(
- 4,
- [SELECT COUNT() FROM Account WHERE Site = 'SECOND'],
- 'The second run overwrites the first, so it ran after every page of run one.'
+ Assert.areEqual(AsyncRequeue.STATUS_STORED, job.requeueStatus);
+ Assert.isTrue(
+ job.requeuePayload.contains('STORED'),
+ 'The payload must carry the state the job was enqueued with: ' + job.requeuePayload
);
+ Assert.areEqual(job.requeuePayload.length(), job.requeuePayloadSize);
}
@IsTest
- private static void shouldGateASecondChunkRunOnTheFirstRunOutcome() {
- List accounts = accountsWithOneFailingRecord();
+ private static void shouldNotTouchTheJobWhenPayloadStorageIsOff() {
+ RequeueableJob job = new RequeueableJob('NOT-STORED');
- Test.startTest();
- Async.chunk(new FailMarkedChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .chunk(new LabelingChunkJob('SECOND'), ChunkSource.of(accounts))
- .chunkSize(2)
- .dependsOn(Async.afterPrevious().succeeded())
- .enqueue();
- Test.stopTest();
+ new QueueableChain().addJob(job);
- Assert.areEqual(
- 0,
- [SELECT COUNT() FROM Account WHERE Site = 'SECOND'],
- 'dependsOn(previous) resolves against the first run outcome, which failed.'
- );
+ Assert.isNull(job.requeueStatus, 'Storing payloads has to be opted into.');
+ Assert.isNull(job.requeuePayload);
}
@IsTest
- private static void shouldOpenACursorInTheRequestedAccessLevel() {
- createAccounts(2);
-
- ChunkSource userMode = ChunkSource.query('SELECT Id FROM Account', AccessLevel.USER_MODE);
- ChunkSource systemMode = ChunkSource.query(
- 'SELECT Id FROM Account WHERE Name LIKE :prefix',
- new Map{ 'prefix' => 'Chunk%' },
- AccessLevel.SYSTEM_MODE
+ private static void shouldLetAJobRecordOptOutOfOrgWidePayloadStorage() {
+ RequeueableJob job = new RequeueableJob('SENSITIVE');
+ QueueableChain chain = new QueueableChain();
+ QueueableChain.jobSettingByName = payloadStorageOn();
+ QueueableChain.jobSettingByName.put(
+ job.className,
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = job.className,
+ StoreJobPayload__c = 'No'
+ )
);
- Assert.areEqual(2, userMode.getNumRecords(), 'The admin running the test sees both.');
- Assert.areEqual(2, systemMode.getNumRecords(), 'Binds and access level combine.');
+ chain.addJob(job);
+
+ Assert.isNull(
+ job.requeueStatus,
+ 'A job carrying sensitive data must be able to opt out of an org-wide Yes.'
+ );
}
@IsTest
- private static void shouldChainNothingWhenSourceIsEmpty() {
- Async.Result result = Async.chunk(
- new MarkingChunkJob(),
- ChunkSource.of(new List())
+ private static void shouldNotStoreAPayloadForChunkJobs() {
+ QueueableChain.jobSettingByName = payloadStorageOn();
+
+ Async.Result chunked = Async.chunk(
+ new StateCarryingChunkJob(),
+ ChunkSource.of(new List{ new Account(Name = 'PAGE') })
)
.chain();
- Assert.areEqual(null, result.customJobId, 'An empty source chains no run.');
- Assert.isTrue(result.queueableChainState.jobs.isEmpty());
+ Assert.areEqual(
+ AsyncRequeue.STATUS_NOT_SERIALIZABLE,
+ chunked.job.requeueStatus,
+ 'A chunk run holds a source that cannot be stored, and a half-read cursor is not ' +
+ 'a meaningful thing to replay.'
+ );
}
@IsTest
- private static void shouldRunChunkRunWhenItsDependencySucceeded() {
- List accounts = createAccounts(4);
+ private static void shouldRecordTooLargeRatherThanBreakTheInsert() {
+ BulkyRequeueableJob job = new BulkyRequeueableJob();
+ for (Integer i = 0; i < 70; i++) {
+ job.payloadParts.add('x'.repeat(2000));
+ }
+
+ chainStoringPayloads().addJob(job);
+
+ Assert.areEqual(AsyncRequeue.STATUS_TOO_LARGE, job.requeueStatus);
+ Assert.isNull(
+ job.requeuePayload,
+ 'An oversized payload would fail the insert and lose the failure record with it.'
+ );
+ Assert.isTrue(
+ job.requeuePayloadSize > JobPayload.MAX_CHARS,
+ 'The size is recorded even when the payload is not, so it can be reported on.'
+ );
+ }
+
+ @IsTest
+ private static void shouldWriteAFailureRowEvenWhenResultCreationIsOff() {
+ FailingRequeueableJob job = new FailingRequeueableJob();
+ job.continueOnJobExecuteFail = true;
+ QueueableChain chain = chainStoringPayloads();
+ chain.addJob(job);
+ QueueableManager.get().setChain(chain);
Test.startTest();
- Async.queueable(new ProcessedCountMarkerJob('GATE'))
- .chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .dependsOn(Async.afterPrevious().succeeded())
- .enqueue();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(
+ new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.UNHANDLED_EXCEPTION)
+ );
Test.stopTest();
+ AsyncResult__c result = [
+ SELECT Status__c, RequeueStatus__c, JobPayload__c
+ FROM AsyncResult__c
+ WHERE CustomJobId__c = :job.customJobId
+ ];
+ Assert.areEqual(QueueableManager.STATUS_FAILED, result.Status__c);
Assert.areEqual(
- 4,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'The run starts once the job it depends on succeeded.'
+ AsyncRequeue.STATUS_STORED,
+ result.RequeueStatus__c,
+ 'A payload with no row to sit on could never be replayed.'
);
+ Assert.isNotNull(result.JobPayload__c);
}
@IsTest
- private static void shouldSkipChunkRunWhenItsDependencyFailed() {
- List accounts = createAccounts(4);
+ private static void shouldNotWriteASuccessRowWhenOnlyPayloadStorageIsOn() {
+ RequeueableJob job = new RequeueableJob('SUCCEEDS');
+ QueueableChain chain = chainStoringPayloads();
+ chain.addJob(job);
+ QueueableManager.get().setChain(chain);
Test.startTest();
- Async.queueable(new FailureQueueableTest())
- .continueOnJobExecuteFail()
- .chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .dependsOn(Async.afterPrevious().succeeded())
- .enqueue();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(
+ new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.SUCCESS)
+ );
Test.stopTest();
Assert.areEqual(
0,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'The whole run is skipped when the job it depends on failed.'
+ [SELECT COUNT() FROM AsyncResult__c],
+ 'Jobs that succeeded still obey CreateResult__c, so the added volume stays bounded ' +
+ 'by the failure rate.'
);
}
@IsTest
- private static void shouldRejectInvalidDependencyOnChunk() {
- ChunkBuilder builder = Async.chunk(
- new MarkingChunkJob(),
- ChunkSource.of(createAccounts(1))
+ private static void shouldWriteASkippedRowSoAChainWideReplayStaysPossible() {
+ FailureQueueableTest blocker = new FailureQueueableTest();
+ blocker.continueOnJobExecuteFail = true;
+ RequeueableJob blocked = new RequeueableJob('NEVER-RAN');
+ QueueableChain chain = chainStoringPayloads();
+ chain.addJob(blocker);
+ blocked.dependencies = new List{
+ dependencyOn(blocker.customJobId, Async.Outcome.SUCCESS)
+ };
+ chain.addJob(blocked);
+ QueueableManager.get().setChain(chain);
+
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(
+ new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.UNHANDLED_EXCEPTION)
+ );
+ Test.stopTest();
+
+ AsyncResult__c skipped = [
+ SELECT Status__c, RequeueStatus__c
+ FROM AsyncResult__c
+ WHERE CustomJobId__c = :blocked.customJobId
+ ];
+ Assert.areEqual(QueueableManager.STATUS_SKIPPED_DEPENDENCY, skipped.Status__c);
+ Assert.areEqual(
+ AsyncRequeue.STATUS_STORED,
+ skipped.RequeueStatus__c,
+ 'A job blocked by a dependency never failed, and it is exactly the one a chain-wide ' +
+ 'replay has to rebuild.'
);
- try {
- builder.dependsOn(null);
- Assert.fail('A dependency without an outcome should be rejected.');
- } catch (Exception ex) {
- Assert.areEqual(QueueableManager.ERROR_MESSAGE_INVALID_DEPENDENCY, ex.getMessage());
- }
- try {
- builder.dependsOn(Async.afterPrevious().succeeded());
- Assert.fail('afterPrevious() without a previously chained job should be rejected.');
- } catch (Exception ex) {
- Assert.areEqual(
- QueueableManager.ERROR_MESSAGE_DEPENDS_ON_PREVIOUS_WITHOUT_JOB,
- ex.getMessage()
- );
- }
}
@IsTest
- private static void shouldYieldIdOnlyShellsWhenSourceIsBuiltFromIds() {
- List accounts = createAccounts(5);
- ChunkSource source = ChunkSource.ofIds(new Map(accounts).keySet());
+ private static void shouldDropThePayloadFromTheJobOnceItIsOnTheRecord() {
+ FailingRequeueableJob job = new FailingRequeueableJob();
+ job.continueOnJobExecuteFail = true;
+ QueueableChain chain = chainStoringPayloads();
+ chain.addJob(job);
+ QueueableManager.get().setChain(chain);
- List page = source.fetch(0, 5);
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(
+ new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.UNHANDLED_EXCEPTION)
+ );
+ Test.stopTest();
- Assert.areEqual(5, page.size());
- Assert.areEqual(
- new Map(accounts).keySet(),
- new Map(page).keySet(),
- 'A job reads ids straight off the page with new Map(chunk).keySet().'
+ Assert.isNull(
+ job.requeuePayload,
+ 'The payload rides the serialized job into every later context, so keeping it would ' +
+ 'carry the job size twice for the rest of the chain.'
);
}
@IsTest
- private static void shouldExposeInMemorySourceRecordsAndSize() {
- List accounts = createAccounts(5);
- ChunkSource source = ChunkSource.of(accounts);
+ private static void shouldRequeueAStoredJobAndLinkItBackToTheOriginal() {
+ RequeueableJob job = new RequeueableJob('REPLAYED');
+ chainStoringPayloads().addJob(job);
+ QueueableChain.jobSettingByName.get(QueueableManager.QUEUEABLE_JOB_SETTING_ALL)
+ .CreateResult__c = true;
+ AsyncResult__c original = storedResultFor(job);
- Assert.areEqual(5, source.getNumRecords());
- Assert.areEqual(2, source.fetch(0, 2).size());
+ Test.startTest();
+ Async.RequeueSummary summary = Async.requeue(original.Id);
+ Test.stopTest();
+
+ Assert.areEqual(new List{ original.Id }, summary.requeued);
+ Assert.isNotNull(
+ summary.enqueueResult?.salesforceJobId,
+ 'A requeue is an enqueue, so the caller gets the same handle on the chain.'
+ );
Assert.areEqual(
1,
- source.fetch(4, 5).size(),
- 'A short final page must clamp to the remainder.'
+ [SELECT COUNT() FROM Account WHERE Name = 'REPLAYED'],
+ 'The rebuilt job has to run with the state it was enqueued with.'
+ );
+ Assert.areEqual(
+ AsyncRequeue.STATUS_REQUEUED,
+ [SELECT RequeueStatus__c FROM AsyncResult__c WHERE Id = :original.Id].RequeueStatus__c,
+ 'Marking the source is what stops a scheduled replay picking it up again.'
+ );
+ Assert.areEqual(
+ original.Id,
+ [
+ SELECT RequeuedFrom__c
+ FROM AsyncResult__c
+ WHERE Id != :original.Id
+ LIMIT 1
+ ]
+ .RequeuedFrom__c,
+ 'A replay runs in a new chain, so this lookup is the only link back.'
);
}
@IsTest
- private static void shouldBuildIdSourceRecordsWithIds() {
- List accounts = createAccounts(3);
- Set ids = new Map(accounts).keySet();
- ChunkSource source = ChunkSource.ofIds(ids);
-
- Assert.areEqual(3, source.getNumRecords());
- Set fetched = new Set();
- for (SObject record : source.fetch(0, 3)) {
- fetched.add(record.Id);
+ private static void shouldRequeueTwoHundredResultsAsOneChainInOneCall() {
+ RequeueableJob job = new RequeueableJob('BULK');
+ chainStoringPayloads().addJob(job);
+ List stored = new List();
+ for (Integer i = 0; i < 200; i++) {
+ stored.add(
+ new AsyncResult__c(
+ ClassName__c = job.className,
+ Status__c = QueueableManager.STATUS_FAILED,
+ JobPayload__c = job.requeuePayload,
+ PayloadSize__c = job.requeuePayloadSize,
+ RequeueStatus__c = AsyncRequeue.STATUS_STORED
+ )
+ );
}
- Assert.areEqual(ids, fetched);
- }
-
- @IsTest
- private static void shouldTreatNullRecordsAsEmptySource() {
- Assert.areEqual(0, ChunkSource.ofIds(null).getNumRecords());
- Assert.areEqual(0, ChunkSource.of(null).getNumRecords());
- }
+ insert stored;
+ Map storedById = new Map(stored);
- @IsTest
- private static void shouldDelegateToCursorSource() {
- createAccounts(5);
- ChunkSource source = ChunkSource.cursor(Database.getCursor('SELECT Id FROM Account'));
+ Async.RequeueSummary summary = Async.requeue(storedById.keySet());
- Assert.areEqual(5, source.getNumRecords());
- Assert.areEqual(2, source.fetch(0, 2).size());
+ Assert.areEqual(200, summary.requeued.size());
+ Assert.isTrue(summary.skipReasonByResultId.isEmpty(), '' + summary.skipReasonByResultId);
+ Assert.areEqual(
+ 200,
+ summary.enqueueResult.queueableChainState.jobs.size(),
+ 'Every replay has to land in the one chain, not in 200 separate enqueues.'
+ );
+ Assert.areEqual(
+ 200,
+ [
+ SELECT COUNT()
+ FROM AsyncResult__c
+ WHERE
+ Id IN :storedById.keySet()
+ AND RequeueStatus__c = :AsyncRequeue.STATUS_REQUEUED
+ ],
+ 'Every source row has to be marked in the same call.'
+ );
}
@IsTest
- private static void shouldOpenCursorFromQuery() {
- createAccounts(4);
- ChunkSource source = ChunkSource.query('SELECT Id FROM Account');
+ private static void shouldNotStoreAPayloadOnTheRowOfAJobThatSucceeded() {
+ RequeueableJob job = new RequeueableJob('SUCCEEDED');
+ QueueableChain chain = chainStoringPayloads();
+ QueueableChain.jobSettingByName.get(QueueableManager.QUEUEABLE_JOB_SETTING_ALL)
+ .CreateResult__c = true;
+ chain.addJob(job);
+ QueueableManager.get().setChain(chain);
- Assert.areEqual(4, source.getNumRecords());
+ Test.startTest();
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(
+ new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.SUCCESS)
+ );
+ Test.stopTest();
+
+ AsyncResult__c result = [
+ SELECT Status__c, RequeueStatus__c, JobPayload__c
+ FROM AsyncResult__c
+ WHERE CustomJobId__c = :job.customJobId
+ ];
+ Assert.areEqual(QueueableManager.STATUS_COMPLETED, result.Status__c);
+ Assert.isNull(
+ result.RequeueStatus__c,
+ 'There is nothing to replay about a job that worked, and a Stored status here would ' +
+ 'let a bulk replay re-run it.'
+ );
+ Assert.isNull(result.JobPayload__c);
}
@IsTest
- private static void shouldOpenCursorFromQueryWithBinds() {
- createAccounts(4);
- ChunkSource source = ChunkSource.query(
- 'SELECT Id FROM Account WHERE Name = :accountName',
- new Map{ 'accountName' => 'Chunk 1' }
+ private static void shouldKeepThePayloadThroughARetryThatRestoresEnqueuedState() {
+ QueueableChain.jobSettingByName = payloadStorageOn();
+
+ Test.startTest();
+ Async.queueable(new FailingRequeueableJob())
+ .retry(1)
+ .restoreStateOnRetry()
+ .continueOnJobExecuteFail()
+ .enqueue();
+ Test.stopTest();
+
+ AsyncResult__c result = [
+ SELECT RetryAttempts__c, RequeueStatus__c, JobPayload__c
+ FROM AsyncResult__c
+ LIMIT 1
+ ];
+ Assert.areEqual(1, result.RetryAttempts__c, 'The retry has to have run.');
+ Assert.areEqual(
+ AsyncRequeue.STATUS_STORED,
+ result.RequeueStatus__c,
+ 'A job that exhausted its retries after an outage is the one requeue exists for.'
);
+ Assert.isNotNull(result.JobPayload__c);
+ }
- Assert.areEqual(1, source.getNumRecords(), 'Bind variables must reach the cursor.');
+ @IsTest
+ private static void shouldKeepThePayloadThroughADeepClonedRetry() {
+ QueueableChain.jobSettingByName = payloadStorageOn();
+
+ Test.startTest();
+ Async.queueable(new FailingRequeueableJob())
+ .deepClone()
+ .retry(1)
+ .continueOnJobExecuteFail()
+ .enqueue();
+ Test.stopTest();
+
+ AsyncResult__c result = [
+ SELECT RetryAttempts__c, RequeueStatus__c, JobPayload__c
+ FROM AsyncResult__c
+ LIMIT 1
+ ];
+ Assert.areEqual(1, result.RetryAttempts__c, 'The retry has to have run.');
+ Assert.areEqual(AsyncRequeue.STATUS_STORED, result.RequeueStatus__c);
+ Assert.isNotNull(
+ result.JobPayload__c,
+ 'The deep copy strips the payload out of the JSON it makes; the copy still has to get it back.'
+ );
}
@IsTest
- private static void shouldRejectNullCursorAndBlankQuery() {
- try {
- ChunkSource.cursor(null);
- Assert.fail('A null cursor should be rejected.');
- } catch (Exception ex) {
- Assert.areEqual(ChunkSource.ERROR_MESSAGE_NULL_CURSOR, ex.getMessage());
- }
- try {
- ChunkSource.query(' ');
- Assert.fail('A blank query should be rejected.');
- } catch (Exception ex) {
- Assert.areEqual(ChunkSource.ERROR_MESSAGE_BLANK_QUERY, ex.getMessage());
- }
- try {
- ChunkSource.query(' ', new Map());
- Assert.fail('A blank query with binds should be rejected.');
- } catch (Exception ex) {
- Assert.areEqual(ChunkSource.ERROR_MESSAGE_BLANK_QUERY, ex.getMessage());
- }
+ private static void shouldDropThePayloadFromASupersededAttempt() {
+ FailureQueueableTest job = new FailureQueueableTest();
+ job.maxRetries = 1;
+ job.continueOnJobExecuteFail = true;
+ QueueableChain chain = chainStoringPayloads();
+ chain.addJob(job);
+ QueueableManager.get().setChain(chain);
+ String payload = job.requeuePayload;
+
+ chain.executeCurrentJob(new AsyncMock.MockQueueableContext());
+ chain.enqueueNextJobIfAnyFromFinalizer(
+ new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.UNHANDLED_EXCEPTION)
+ );
+
+ Assert.isNull(
+ job.requeuePayload,
+ 'The failed attempt stays in the chain; it must not carry the payload too.'
+ );
+ Assert.areEqual(
+ payload,
+ chain.jobs[0].requeuePayload,
+ 'The retry attempt is the one that will produce the row, so it carries the payload.'
+ );
}
@IsTest
- private static void shouldOrderJobsOfEqualPriorityByChainSequence() {
- QueueableJobTest1 first = new QueueableJobTest1();
- first.uniqueName = 'first';
- first.chainSequence = 1;
- QueueableJobTest2 second = new QueueableJobTest2();
- second.uniqueName = 'second';
- second.chainSequence = 2;
+ private static void shouldDegradeWhenTheRegisteredSerializerReturnsNull() {
+ QueueableChain chain = chainStoringPayloads();
+ QueueableChain.jobSettingByName.get(QueueableManager.QUEUEABLE_JOB_SETTING_ALL)
+ .JobSerializerClass__c = loggerName('NullReturningJobSerializer');
+ RequeueableJob job = new RequeueableJob('NULL-SERIALIZER');
- List jobs = new List{ second, first };
- jobs.sort();
+ chain.addJob(job);
- Assert.areEqual('first', jobs[0].uniqueName, 'Equal priority keeps the chained order.');
- Assert.areEqual('second', jobs[1].uniqueName);
+ Assert.areEqual(AsyncRequeue.STATUS_NOT_SERIALIZABLE, job.requeueStatus);
+ Assert.isTrue(
+ job.retryHistory.contains(QueueableManager.CAUSE_SERIALIZER_RETURNED_NULL),
+ 'A registered class is configuration, and configuration mistakes degrade with a ' +
+ 'reason rather than stopping every enqueue: ' +
+ job.retryHistory
+ );
}
@IsTest
- private static void shouldRunEveryChunkBeforeTheNextChainedJob() {
- List accounts = createAccounts(6);
+ private static void shouldCaptureTheJobBeforeAnyChainBookkeepingLandsOnIt() {
+ RequeueableJob job = new RequeueableJob('PRISTINE');
+ QueueableChain.jobSettingByName = payloadStorageOn();
+ QueueableChain.jobSettingByName.put(
+ job.className,
+ new QueueableJobSetting__mdt(QueueableJobName__c = job.className, MaxRetries__c = 3)
+ );
- Test.startTest();
- Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .chain(new ProcessedCountMarkerJob('AFTER_RUN'))
- .enqueue();
- Test.stopTest();
+ Async.queueable(job).chain();
+ RequeueableJob rebuilt = (RequeueableJob) JSON.deserialize(
+ QueueableManager.get().getChain().getJobs()[0].requeuePayload,
+ RequeueableJob.class
+ );
+ Assert.isNull(rebuilt.customJobId, 'The replay gets a fresh id.');
+ Assert.isNull(rebuilt.chainSequence, 'The replay takes its place in its own chain.');
Assert.areEqual(
- '6',
- processedCountSeenBy('AFTER_RUN'),
- 'A job chained after the run must wait for every page.'
+ 0,
+ rebuilt.maxRetries,
+ 'Custom Metadata defaults are applied again at replay, from the metadata of that day.'
);
}
@IsTest
- private static void shouldRunChunkRunAfterAnEarlierChainedJob() {
- List accounts = createAccounts(4);
+ private static void shouldSkipAJobWhoseClassNoLongerDeclaresHowItsStateResets() {
+ RequeueableJob job = new RequeueableJob('CHANGED-CLASS');
+ chainStoringPayloads().addJob(job);
+ job.requeuePayload = job.requeuePayload.replace('"maxRetries":0', '"maxRetries":2');
+ AsyncResult__c stored = storedResultFor(job);
+
+ Async.RequeueSummary summary = Async.requeue(new Set{ stored.Id });
+
+ Assert.isTrue(summary.requeued.isEmpty(), 'The gate has to hold on replay too.');
+ Assert.isTrue(
+ summary.skipReasonByResultId.get(stored.Id).contains('Async.Retryable'),
+ 'One bad class is one skipped row with the gate message, not an abandoned batch: ' +
+ summary.skipReasonByResultId.get(stored.Id)
+ );
+ }
+
+ @IsTest
+ private static void shouldStillWriteTheReplayRowWhenTheSourceWasDeletedMeanwhile() {
+ FailingRequeueableJob job = new FailingRequeueableJob();
+ job.continueOnJobExecuteFail = true;
+ chainStoringPayloads().addJob(job);
+ AsyncResult__c original = storedResultFor(job);
Test.startTest();
- Async.queueable(new ProcessedCountMarkerJob('BEFORE_RUN'))
- .chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .enqueue();
+ Async.requeue(new Set{ original.Id });
+ delete original;
Test.stopTest();
+ AsyncResult__c replay = [
+ SELECT Status__c, RequeuedFrom__c
+ FROM AsyncResult__c
+ LIMIT 1
+ ];
+ Assert.areEqual(QueueableManager.STATUS_FAILED, replay.Status__c);
+ Assert.isNull(
+ replay.RequeuedFrom__c,
+ 'A lookup to a row the cleanup batch removed would fail the insert and lose the row.'
+ );
+ }
+
+ @IsTest
+ private static void shouldDoNothingWhenAskedToRequeueNothing() {
+ Async.RequeueSummary summary = Async.requeue(new Set());
+
+ Assert.isTrue(summary.requeued.isEmpty());
+ Assert.isTrue(summary.skipReasonByResultId.isEmpty());
+ Assert.isNull(summary.enqueueResult, 'Nothing was replayed, so there is no chain.');
+ }
+
+ @IsTest
+ private static void shouldMergeInfoMapsOnBothBuilders() {
+ Async.Result queued = Async.queueable(new SuccessfulQueueableTest())
+ .info('team', 'platform')
+ .info(new Map{ 'package' => 'billing', 'team' => 'billing' })
+ .chain();
+ Async.Result chunked = Async.chunk(
+ new StateCarryingChunkJob(),
+ ChunkSource.of(new List{ new Account(Name = 'PAGE') })
+ )
+ .info('team', 'platform')
+ .info(new Map{ 'package' => 'billing' })
+ .chain();
+
Assert.areEqual(
- '0',
- processedCountSeenBy('BEFORE_RUN'),
- 'A job chained before the run must run first.'
+ new Map{ 'team' => 'billing', 'package' => 'billing' },
+ queued.job.info,
+ 'A later info() call adds keys and overwrites the ones it repeats.'
+ );
+ Assert.areEqual(
+ new Map{ 'team' => 'platform', 'package' => 'billing' },
+ chunked.job.info
);
- Assert.areEqual(4, [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED]);
}
@IsTest
- private static void shouldLetHigherPriorityJobPreemptTheNextChunk() {
- List accounts = createAccounts(6);
- PreemptingChunkJob job = new PreemptingChunkJob();
- job.preemptPriority = 1;
+ private static void shouldWarnOnceWhenTheLoggerConstructorThrows() {
+ QueueableChain chain = chainWithLogger(loggerName('UnconstructableLogger'));
+ QueueableManager.get().setChain(chain);
+ SuccessfulQueueableTest first = new SuccessfulQueueableTest();
+ SuccessfulQueueableTest second = new SuccessfulQueueableTest();
+
+ chain.addJob(first);
+ chain.addJob(second);
+
+ Assert.isTrue(
+ first.retryHistory.contains('could not be constructed'),
+ 'The warning has to say the class was found but not built: ' + first.retryHistory
+ );
+ Assert.isTrue(
+ first.retryHistory.contains(LOGGER_CONSTRUCTOR_FAILURE),
+ 'The constructor message has to reach the job: ' + first.retryHistory
+ );
+ Assert.isNull(
+ second.retryHistory,
+ 'The same broken class is reported once per transaction, not once per job.'
+ );
+ }
+
+ @IsTest
+ private static void shouldExplainEveryResultItRefusedToRequeue() {
+ AsyncResult__c noPayload = new AsyncResult__c(Status__c = QueueableManager.STATUS_FAILED);
+ AsyncResult__c alreadyRequeued = new AsyncResult__c(
+ Status__c = QueueableManager.STATUS_FAILED,
+ RequeueStatus__c = AsyncRequeue.STATUS_REQUEUED
+ );
+ AsyncResult__c tooLarge = new AsyncResult__c(
+ Status__c = QueueableManager.STATUS_FAILED,
+ RequeueStatus__c = AsyncRequeue.STATUS_TOO_LARGE
+ );
+ insert new List{ noPayload, alreadyRequeued, tooLarge };
+ AsyncResult__c deleted = new AsyncResult__c(Status__c = QueueableManager.STATUS_FAILED);
+ insert deleted;
+ Id deletedId = deleted.Id;
+ delete deleted;
Test.startTest();
- Async.chunk(job, ChunkSource.of(accounts)).chunkSize(2).priority(5).enqueue();
+ Async.RequeueSummary summary = Async.requeue(
+ new Set{ noPayload.Id, alreadyRequeued.Id, tooLarge.Id, deletedId }
+ );
Test.stopTest();
+ Assert.isTrue(summary.requeued.isEmpty(), 'None of these can be replayed.');
Assert.areEqual(
- '2',
- processedCountSeenBy('PREEMPT'),
- 'A higher priority job added mid-run runs before the next page.'
+ QueueableManager.REQUEUE_SKIPPED_NO_PAYLOAD,
+ summary.skipReasonByResultId.get(noPayload.Id)
);
Assert.areEqual(
- 6,
- [SELECT COUNT() FROM Account WHERE Site = :CHUNK_PROCESSED],
- 'The run resumes after the preempting job.'
+ QueueableManager.REQUEUE_SKIPPED_ALREADY_REQUEUED,
+ summary.skipReasonByResultId.get(alreadyRequeued.Id)
+ );
+ Assert.isTrue(
+ summary.skipReasonByResultId.get(tooLarge.Id).contains(AsyncRequeue.STATUS_TOO_LARGE),
+ 'The reason has to name the status, but was: ' +
+ summary.skipReasonByResultId.get(tooLarge.Id)
+ );
+ Assert.areEqual(
+ QueueableManager.REQUEUE_SKIPPED_NO_RECORD,
+ summary.skipReasonByResultId.get(deletedId)
);
}
@IsTest
- private static void shouldRunLowerPriorityJobAfterTheWholeRun() {
- List accounts = createAccounts(6);
- PreemptingChunkJob job = new PreemptingChunkJob();
- job.preemptPriority = 9;
+ private static void shouldRefuseToRequeueMoreCharactersThanItCanHold() {
+ List heavy = new List();
+ for (Integer i = 0; i < 20; i++) {
+ heavy.add(
+ new AsyncResult__c(
+ Status__c = QueueableManager.STATUS_FAILED,
+ RequeueStatus__c = AsyncRequeue.STATUS_STORED,
+ PayloadSize__c = JobPayload.MAX_CHARS
+ )
+ );
+ }
+ insert heavy;
+ Map heavyById = new Map(heavy);
+
+ try {
+ Async.requeue(heavyById.keySet());
+ Assert.fail('Deserializing this many payloads at once would blow the heap.');
+ } catch (Async.IllegalArgumentException expected) {
+ Assert.isTrue(
+ expected.getMessage().contains(String.valueOf(AsyncRequeue.MAX_CHARS_PER_REQUEUE)),
+ 'The caller has to be told the limit, but was: ' + expected.getMessage()
+ );
+ }
+ }
+
+ @IsTest
+ private static void shouldRouteBothHalvesThroughARegisteredSerializer() {
+ QueueableChain chain = chainStoringPayloads();
+ QueueableChain.jobSettingByName.get(QueueableManager.QUEUEABLE_JOB_SETTING_ALL)
+ .JobSerializerClass__c = loggerName('RecordingJobSerializer');
+ RequeueableJob job = new RequeueableJob('VIA-SERIALIZER');
+ chain.addJob(job);
+ AsyncResult__c original = storedResultFor(job);
Test.startTest();
- Async.chunk(job, ChunkSource.of(accounts)).chunkSize(2).priority(5).enqueue();
+ Async.requeue(new Set{ original.Id });
+ Test.stopTest();
+
+ Assert.isTrue(
+ loggedEvents.contains('serialize:' + job.className),
+ 'Capture has to run in the consumer namespace too: ' + loggedEvents
+ );
+ Assert.isTrue(
+ loggedEvents.contains('deserialize:' + job.className),
+ 'Rebuild has to run in the consumer namespace: ' + loggedEvents
+ );
+ }
+
+ @IsTest
+ private static void shouldRejectARegisteredClassThatIsNotASerializer() {
+ QueueableChain chain = chainStoringPayloads();
+ QueueableChain.jobSettingByName.get(QueueableManager.QUEUEABLE_JOB_SETTING_ALL)
+ .JobSerializerClass__c = loggerName('RecordingLogger');
+ RequeueableJob job = new RequeueableJob('NOT-A-SERIALIZER');
+
+ chain.addJob(job);
+
+ Assert.areEqual(AsyncRequeue.STATUS_NOT_SERIALIZABLE, job.requeueStatus);
+ Assert.isTrue(
+ job.retryHistory.contains('JobSerializerClass__c'),
+ 'The misconfigured field has to be named, but was: ' + job.retryHistory
+ );
+ }
+
+ @IsTest
+ private static void shouldReportWhenTheStoredClassNoLongerExists() {
+ AsyncResult__c topLevel = orphanResultFor('NoSuchJobClass');
+ AsyncResult__c nested = orphanResultFor('NoSuch.JobClass');
+ insert new List{ topLevel, nested };
+
+ Test.startTest();
+ Async.RequeueSummary summary = Async.requeue(new Set{ topLevel.Id, nested.Id });
Test.stopTest();
+ Assert.isTrue(
+ summary.skipReasonByResultId.get(topLevel.Id).contains('NoSuchJobClass'),
+ 'A renamed or deleted class must be named, but was: ' +
+ summary.skipReasonByResultId.get(topLevel.Id)
+ );
+ Assert.isTrue(
+ summary.skipReasonByResultId.get(nested.Id).contains('NoSuch.JobClass'),
+ 'An inner class name goes through the two-argument lookup too: ' +
+ summary.skipReasonByResultId.get(nested.Id)
+ );
Assert.areEqual(
- '6',
- processedCountSeenBy('PREEMPT'),
- 'A lower priority job added mid-run waits for the whole run.'
+ 2,
+ [
+ SELECT COUNT()
+ FROM AsyncResult__c
+ WHERE RequeueStatus__c = :AsyncRequeue.STATUS_STORED
+ ],
+ 'A row that could not be replayed must stay replayable once the class is back.'
+ );
+ }
+
+ private static QueueableChain chainWithLogger(String loggerClassName) {
+ return chainWithSettings(
+ new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ LoggerClass__c = loggerClassName
+ )
);
}
- @IsTest
- private static void shouldRunDependentJobWhenEveryChunkSucceeded() {
- List accounts = createAccounts(6);
-
- Test.startTest();
- Async.chunk(new MarkingChunkJob(), ChunkSource.of(accounts))
- .chunkSize(2)
- .chain(new ProcessedCountMarkerJob('DEPENDENT'))
- .dependsOn(Async.afterPrevious().succeeded())
- .enqueue();
- Test.stopTest();
+ private static String loggerName(String simpleName) {
+ return getClassNameWithNamespaceDotPrefix('AsyncTest.' + simpleName);
+ }
+
+ private static Async.Dependency dependencyOn(String customJobId, Async.Outcome outcome) {
+ Async.Dependency dependency = new Async.Dependency();
+ dependency.resolvedTargetCustomJobId = customJobId;
+ dependency.requiredOutcome = outcome;
+ return dependency;
+ }
+
+ private static Map resultsEnabledForAll() {
+ return new Map{
+ QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
+ DeveloperName = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ CreateResult__c = true
+ )
+ };
+ }
+
+ private static MarkerJob markerJob(String tag, Boolean shouldFail) {
+ MarkerJob job = new MarkerJob();
+ job.tag = tag;
+ job.shouldFail = shouldFail;
+ return job;
+ }
+
+ private static Integer accountCount(String name) {
+ return [SELECT COUNT() FROM Account WHERE Name = :name];
+ }
+
+ private static Account hookRecord() {
+ return [SELECT Description FROM Account WHERE Name = 'HOOK' LIMIT 1];
+ }
+
+ private static String getClassNameWithNamespaceDotPrefix(String className) {
+ return getNamespaceDotPrefix() + className;
+ }
+
+ private static String getNamespaceDotPrefix() {
+ String className = AsyncTest.class.getName();
+ return className.contains('.') ? className.substringBefore('.') + '.' : '';
+ }
+
+ public Iterable start(Database.BatchableContext bc) {
+ // This is just a placeholder to start the batch.
+ return new List{ new Account() };
+ }
+
+ public void execute(Database.BatchableContext ctx, List scope) {
+ for (Account acc : scope) {
+ acc.Description = 'Processed by: ' + ctx.getJobId();
+ }
+ if (!scope.isEmpty() && scope[0].Id != null) {
+ update scope;
+ }
+ }
+
+ public void finish(Database.BatchableContext bc) {
+ insert new Account(Name = 'Batch Complete', Description = 'Job: ' + bc.getJobId());
+ }
+
+ private static AsyncResult__c insertAsyncResultWithAge(String status, Integer ageDays) {
+ AsyncResult__c result = new AsyncResult__c(Status__c = status);
+ insert result;
+ Test.setCreatedDate(result.Id, System.now().addDays(-ageDays));
+ return result;
+ }
+
+ private static Integer followUpsLeftIn(QueueableChain chain) {
+ Integer followUps = 0;
+ for (QueueableJob job : chain.jobs) {
+ if (job instanceof MarkerJob) {
+ followUps++;
+ }
+ }
+ return followUps;
+ }
+
+ private static QueueableChain chainWithSettings(QueueableJobSetting__mdt setting) {
+ QueueableChain chain = new QueueableChain();
+ QueueableChain.jobSettingByName = new Map{
+ setting.QueueableJobName__c => setting
+ };
+ return chain;
+ }
+
+ private static QueueableChain chainRunning(QueueableJob job) {
+ QueueableChain chain = new QueueableChain();
+ chain.addJob(job);
+ QueueableManager.get().setChain(chain);
+ return chain;
+ }
+
+ private static QueueableChain chainWithEnqueuedJob(QueueableJob job) {
+ QueueableChain chain = chainRunning(job);
+ job.chain = chain;
+ return chain;
+ }
+
+ private static AsyncMock.MockFinalizerContext rolledBack() {
+ return new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.UNHANDLED_EXCEPTION);
+ }
+
+ private static AsyncMock.MockFinalizerContext committed() {
+ return new AsyncMock.MockFinalizerContext().setResult(ParentJobResult.SUCCESS);
+ }
+
+ private static String processedCountSeenBy(String tag) {
+ return [SELECT Site FROM Account WHERE Name = :tag LIMIT 1].Site;
+ }
+
+ private static Id fakeAsyncApexJobId() {
+ return AsyncApexJob.SObjectType.getDescribe().getKeyPrefix() + '0'.repeat(11) + '1AAA';
+ }
+
+ private static List createAccounts(Integer count) {
+ List accounts = new List();
+ for (Integer i = 0; i < count; i++) {
+ accounts.add(new Account(Name = 'Chunk ' + i));
+ }
+ insert accounts;
+ return accounts;
+ }
+
+ private static List accountsWithOneFailingRecord() {
+ List accounts = new List{
+ new Account(Name = 'ok first'),
+ new Account(Name = 'ok second'),
+ new Account(Name = 'FAIL third'),
+ new Account(Name = 'ok fourth'),
+ new Account(Name = 'ok fifth'),
+ new Account(Name = 'ok sixth')
+ };
+ insert accounts;
+ return accounts;
+ }
+
+ private static void markProcessed(List chunk, String marker) {
+ List toUpdate = new List();
+ for (SObject record : chunk) {
+ toUpdate.add(new Account(Id = record.Id, Site = marker));
+ }
+ update toUpdate;
+ }
+
+ private static Map payloadStorageOn() {
+ return new Map{
+ QueueableManager.QUEUEABLE_JOB_SETTING_ALL => new QueueableJobSetting__mdt(
+ QueueableJobName__c = QueueableManager.QUEUEABLE_JOB_SETTING_ALL,
+ StoreJobPayload__c = 'Yes'
+ )
+ };
+ }
+
+ private static QueueableChain chainStoringPayloads() {
+ QueueableChain chain = new QueueableChain();
+ QueueableChain.jobSettingByName = payloadStorageOn();
+ return chain;
+ }
+
+ private static AsyncResult__c orphanResultFor(String missingClassName) {
+ return new AsyncResult__c(
+ ClassName__c = missingClassName,
+ Status__c = QueueableManager.STATUS_FAILED,
+ JobPayload__c = '{}',
+ PayloadSize__c = 2,
+ RequeueStatus__c = AsyncRequeue.STATUS_STORED
+ );
+ }
+
+ private static AsyncResult__c storedResultFor(QueueableJob capturedJob) {
+ AsyncResult__c stored = new AsyncResult__c(
+ ClassName__c = capturedJob.className,
+ CustomJobId__c = capturedJob.customJobId,
+ Status__c = QueueableManager.STATUS_FAILED,
+ JobPayload__c = capturedJob.requeuePayload,
+ PayloadSize__c = capturedJob.requeuePayloadSize,
+ RequeueStatus__c = AsyncRequeue.STATUS_STORED
+ );
+ insert stored;
+ return stored;
+ }
+
+ private class DeepCloneFailJob extends QueueableJob {
+ public override void work() {
+ }
+ public override QueueableJob cloneForDeepCopy() {
+ throw new JSONException('forced deep clone failure');
+ }
+ }
+
+ public class SelfReferencingJob extends QueueableJob {
+ public SelfReferencingJob self;
+ public override void work() {
+ }
+ }
+
+ public class AbstractFieldHoldingJob extends QueueableJob {
+ public Comparable held;
+ public override void work() {
+ }
+ }
+
+ public class AsyncLibTypeHoldingJob extends QueueableJob {
+ public Async.Result heldResult;
+ public Backoff heldBackoff;
+ public Async.Dependency heldDependency;
+ public override void work() {
+ }
+ }
+
+ public class DeepCloneRetryJob extends QueueableJob implements Async.Retryable {
+ public List attempts = new List();
+ public override void work() {
+ attempts.add('attempt' + retryAttempt);
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+
+ public void resetBeforeRetry(Integer attempt) {
+ }
+ }
+
+ public class DeepCloneFailRetryJob extends QueueableJob implements Async.Retryable {
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ public override void work() {
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ public override QueueableJob cloneForDeepCopy() {
+ throw new JSONException('forced deep clone failure');
+ }
+ }
+
+ public class RecordingLogger implements Async.OnJobEnqueued, Async.OnJobSucceeded, Async.OnJobFailed, Async.OnRetryEnqueued {
+ public void onJobEnqueued(Async.JobContext ctx) {
+ AsyncTest.loggedEvents.add('enqueued:' + ctx.className + ':' + ctx.info.get('team'));
+ }
+ public void onJobSucceeded(Async.JobContext ctx) {
+ AsyncTest.loggedEvents.add('succeeded:' + ctx.className);
+ }
+ public void onJobFailed(Async.FailureContext ctx) {
+ AsyncTest.loggedEvents.add('failed:' + ctx.className + ':' + ctx.retryAttempt);
+ }
+ public void onRetryEnqueued(Async.FailureContext ctx) {
+ AsyncTest.loggedEvents.add(
+ 'retry:' + ctx.retryAttempt + ':delay=' + ctx.nextAttemptDelayMinutes
+ );
+ }
+ }
+
+ public class UnconstructableLogger implements Async.OnJobFailed {
+ public UnconstructableLogger() {
+ throw new CustomException(LOGGER_CONSTRUCTOR_FAILURE);
+ }
+
+ public void onJobFailed(Async.FailureContext ctx) {
+ }
+ }
+
+ public class FailureOnlyLogger implements Async.OnJobFailed {
+ public void onJobFailed(Async.FailureContext ctx) {
+ AsyncTest.loggedEvents.add('failed-only:' + ctx.className);
+ }
+ }
+
+ public class ThrowingLogger implements Async.OnJobSucceeded {
+ public void onJobSucceeded(Async.JobContext ctx) {
+ throw new CustomException('logger exploded');
+ }
+ }
+
+ private class SelfLoggingJob extends QueueableJob implements Async.OnJobSucceeded {
+ public override void work() {
+ }
+ public void onJobSucceeded(Async.JobContext ctx) {
+ AsyncTest.loggedEvents.add('self:' + ctx.className);
+ }
+ }
+
+ private class SuccessfulQueueableTest extends QueueableJob {
+ public override void work() {
+ insert new Account(Name = Async.getQueueableJobContext()?.currentJob?.uniqueName);
+ }
+ }
+
+ private class FailureQueueableTest extends QueueableJob.AllowsCallouts implements Async.Retryable {
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ public override void work() {
+ insert new Account(Name = Async.getQueueableJobContext()?.currentJob?.uniqueName);
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class SelfConfiguringRetryJob extends QueueableJob implements Async.Retryable {
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ private SelfConfiguringRetryJob() {
+ this.maxRetries = 2;
+ this.backoff = Async.Backoff.fixed(4);
+ this.continueOnJobExecuteFail = true;
+ }
+
+ public override void work() {
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class ChainStoppingJob extends QueueableJob {
+ public override void work() {
+ Async.queueable(new StopChainFinalizer()).attachFinalizer();
+ }
+ }
+
+ private class MarkerJob extends QueueableJob implements Async.Retryable {
+ public String tag;
+ public Boolean shouldFail = false;
+ public override void work() {
+ insert new Account(Name = tag);
+ if (shouldFail) {
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ public void resetBeforeRetry(Integer attempt) {
+ }
+ }
+
+ private class StopChainFinalizer extends QueueableJob.Finalizer {
+ public override void work() {
+ Async.stopChain();
+ }
+ }
+
+ private class FinalizerAttachingRetryJob extends QueueableJob implements Async.Retryable {
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ public override void work() {
+ insert new Account(Name = 'RETRY-PARENT');
+ Async.queueable(new MarkerFinalizer()).attachFinalizer();
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class MarkerFinalizer extends QueueableJob.Finalizer {
+ public override void work() {
+ insert new Account(Name = 'RETRY-FINALIZER');
+ }
+ }
+
+ private abstract class HookRecordingJob extends QueueableJob implements Async.Retryable {
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ public override void onFinalFailure(Async.FailureContext failureCtx) {
+ insert new Account(
+ Name = 'HOOK',
+ Description = failureCtx.retryOutcome.name() +
+ '|' +
+ (String.isBlank(failureCtx.failure?.stackTrace) ? 'NO-TRACE' : 'HAS-TRACE') +
+ '|' +
+ failureCtx.retryAttempt +
+ '/' +
+ failureCtx.maxRetries
+ );
+ }
+ }
+
+ private class FailureHookJob extends HookRecordingJob {
+ public override void work() {
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class SucceedingHookJob extends HookRecordingJob {
+ public override void work() {
+ insert new Account(Name = 'SUCCEEDED');
+ }
+ }
+
+ private class ExplodingHookJob extends QueueableJob {
+ public override void work() {
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+
+ public override void onFinalFailure(Async.FailureContext failureCtx) {
+ throw new CustomException('hook exploded');
+ }
+ }
+
+ private class FailingChunkHookJob extends ChunkJob implements Async.ChunkResettable {
+ public void resetBeforeNextChunk(Integer pageNumber) {
+ }
+
+ public override void work(List chunk) {
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+
+ public override void onFinalFailure(Async.FailureContext failureCtx) {
+ insert new Account(Name = 'HOOK', Description = failureCtx.retryOutcome.name());
+ }
+ }
+
+ private class JobChainingRetryJob extends QueueableJob implements Async.Retryable {
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ public override void work() {
+ insert new Account(Name = 'CHAIN-PARENT');
+ Async.queueable(markerJob('CHAIN-FOLLOWUP', false)).chain();
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class ChainingFailureJob extends QueueableJob {
+ public override void work() {
+ Async.queueable(markerJob('ROLLED-BACK-FOLLOWUP', false)).chain();
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class FinalizerAndJobChainingFailureJob extends QueueableJob {
+ public override void work() {
+ Async.queueable(new MarkerFinalizer()).attachFinalizer();
+ Async.queueable(markerJob('ROLLED-BACK-FOLLOWUP', false)).chain();
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class ChainStoppingFailureJob extends QueueableJob implements Async.Retryable {
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ public override void work() {
+ Async.stopChain();
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class JobSkippingFailureJob extends QueueableJob {
+ public String targetCustomJobId;
+ public override void work() {
+ Async.skipJob(targetCustomJobId);
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ private class NonRetryableJob extends QueueableJob {
+ public override void work() {
+ }
+ public override Boolean isRetryable(Exception ex) {
+ return false;
+ }
+ }
+
+ private class MessageVetoJob extends QueueableJob {
+ public override void work() {
+ }
+ public override Boolean isRetryable(Exception ex) {
+ return !ex.getMessage().containsIgnoreCase('permanent');
+ }
+ }
+
+ private class ThrowingClassifierJob extends QueueableJob {
+ public override void work() {
+ }
+ public override Boolean isRetryable(Exception ex) {
+ throw new CustomException('classifier blew up');
+ }
+ }
+
+ private class ResettableRetryJob extends QueueableJob implements Async.Retryable {
+ public Boolean wasReset = false;
+ public Integer resetAttempt;
+ public override void work() {
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ public void resetBeforeRetry(Integer attempt) {
+ this.wasReset = true;
+ this.resetAttempt = attempt;
+ }
+ }
+
+ private class UngatedRetryJob extends QueueableJob {
+ public override void work() {
+ }
+ }
+
+ public class RetryingChunkJob extends ChunkJob implements Async.Retryable, Async.ChunkResettable {
+ public List touched = new List();
+
+ public override void work(List chunk) {
+ String firstOfPage = (String) chunk[0].get('Name');
+ AsyncTest.chunkPagesRun.add(firstOfPage + ':' + touched.size());
+ touched.add('page');
+ if (firstOfPage == 'P3' && AsyncTest.chunkPageFailures == 0) {
+ AsyncTest.chunkPageFailures++;
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+ }
+
+ public void resetBeforeRetry(Integer attempt) {
+ }
+
+ public void resetBeforeNextChunk(Integer pageNumber) {
+ }
+ }
+
+ private class GateTrippingParentJob extends QueueableJob {
+ public override void work() {
+ insert new Account(Name = 'GATE-PARENT-RAN');
+ Async.queueable(new UngatedRetryJob()).retry(2).chain();
+ }
+ }
+
+ public class UngatedChunkJob extends ChunkJob {
+ public override void work(List chunk) {
+ }
+ }
+
+ private class OkCalloutMock implements HttpCalloutMock {
+ public HttpResponse respond(HttpRequest request) {
+ HttpResponse response = new HttpResponse();
+ response.setStatusCode(200);
+ return response;
+ }
+ }
+
+ public class MarkerCalloutChunkJob extends ChunkJob implements Database.AllowsCallouts, Async.ChunkResettable {
+ public override void work(List chunk) {
+ HttpRequest request = new HttpRequest();
+ request.setEndpoint('https://example.com');
+ request.setMethod('GET');
+ AsyncTest.chunkCalloutStatus = new Http().send(request).getStatusCode();
+ AsyncTest.chunkCalloutPages++;
+ }
+
+ public void resetBeforeNextChunk(Integer pageNumber) {
+ }
+ }
+
+ public class StateCarryingRetryJob extends QueueableJob implements Async.Retryable {
+ public List seen = new List();
+
+ public override void work() {
+ insert new Account(Name = 'RETRY-SIZE-' + seen.size());
+ seen.add('attempt');
+ throw new CustomException(AsyncTest.CUSTOM_ERROR_MESSAGE);
+ }
+
+ public void resetBeforeRetry(Integer attempt) {
+ }
+ }
+
+ public class StateCarryingChunkJob extends ChunkJob implements Async.ChunkResettable {
+ public List