All files / core/src queue.ts

0% Statements 0/65
100% Branches 1/1
100% Functions 1/1
0% Lines 0/65

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89                                                                                                                                                                                 
import type { Job, JobData, JobRepositoryDriver } from './jobs/jobs.types';
import type { TaskDefinition } from './tasks/tasks.types';
import { customAlphabet, nanoid } from 'nanoid';
import { serializeError } from './errors/errors.models';
import { processJob } from './jobs/jobs.usecases';
import { createTaskDefinitionRegistry } from './tasks/task-definition.registry';
 
export function createQueue({
  driver,
  processingTimeoutMs = 10000,
  generateJobId = customAlphabet('1234567890abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ', 32),
}: {
  driver: JobRepositoryDriver;
  processingTimeoutMs?: number;
  generateJobId?: () => string;
}) {
  const taskDefinitionRegistry = createTaskDefinitionRegistry();
  let status: 'running' | 'stopped' | 'stopping' = 'stopped';
  let stopPromiseResolver: (() => void) | null = null;
 
  return {
    registerTask: (taskDefinition: TaskDefinition) => {
      taskDefinitionRegistry.add(taskDefinition);
    },
 
    async scheduleJob({
      taskName,
      data,
      now = new Date(),
      scheduleAt = now,
      maxRetries,
    }: {
      taskName: string;
      data: JobData;
      scheduleAt?: Date;
      now?: Date;
      maxRetries?: number;
    }) {
      const job: Job = {
        id: generateJobId(),
        taskName,
        data,
        scheduleAt,
        status: 'pending',
        maxRetries,
      };
 
      await driver.saveJob({ job, now });
    },
 
    startWorker({ workerId: _ }: { workerId: string }) {
      status = 'running';
 
      (async () => {
      // eslint-disable-next-line no-unmodified-loop-condition
        while (status === 'running') {
          const { job } = await driver.fetchNextJob({ processingTimeoutMs });
 
          if (!job) {
            await new Promise(resolve => setTimeout(resolve, 1000));
            continue;
          }
 
          const jobId = job.id;
 
          try {
            const result = await processJob({ job, taskDefinitionRegistry });
            await driver.markJobAsCompleted({ jobId, result: result ?? undefined });
          } catch (error) {
            await driver.markJobAsFailed({ jobId, error: serializeError({ error }) });
          }
        }
 
        status = 'stopped';
        stopPromiseResolver?.();
        stopPromiseResolver = null;
      })();
    },
 
    stopWorker(): Promise<void> {
      const { promise, resolve } = Promise.withResolvers<void>();
      stopPromiseResolver = resolve;
      status = 'stopping';
 
      return promise;
    },
  };
}