registerJobs method

void registerJobs(
  1. DVQueues queues
)

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);
      }
    });
}