registerJobs method
Registers delivery and backfill on the durable job layer, with the codecs a queue shared between processes stores them under. A refused batch throws, so the queue retries it with its backoff.
Implementation
void registerJobs(DVQueues queues) {
const DVJobPayloadCodecs()
..register<DVCaptureDeliveryJob>(
const DVJobPayloadCodec<DVCaptureDeliveryJob>(
name: DVCaptureDeliveryJob.codecName,
encode: DVCaptureDeliveryJob.encode,
decode: DVCaptureDeliveryJob.decode,
),
)
..register<DVCaptureBackfillJob>(
const DVJobPayloadCodec<DVCaptureBackfillJob>(
name: DVCaptureBackfillJob.codecName,
encode: DVCaptureBackfillJob.encode,
decode: DVCaptureBackfillJob.decode,
),
);
queues
..register<DVCaptureDeliveryJob>((DVCaptureDeliveryJob job) async {
await _consumer(job.consumer).deliverAll();
})
..register<DVCaptureBackfillJob>((DVCaptureBackfillJob job) async {
final DVRecordTable? table = _tables[job.model];
if (table == null) {
throw StateError('${job.model} is not a captured table here.');
}
final DVCaptureBackfillProgress progress =
await _consumer(job.consumer).backfill(
table,
chunkSize: job.chunkSize,
maxChunks: 1,
tenantColumn: job.tenantColumn,
);
// One chunk per run is the rate limit: the queue's own pacing sits
// between chunks, and the work never holds a worker for the whole
// table.
if (!progress.done) {
await queues.dispatch<DVCaptureBackfillJob>(job, queue: job.queue);
}
});
}