/** * Persistent Job Queue — Resumable task orchestration with durable storage * * Provides a job queue that persists jobs to disk (JSON files) with support * for status tracking, retries, priority ordering, and crash recovery. * Designed for long-running swarm orchestration tasks that must survive * process restarts. * * Backends: * - FileJobStore (built-in) — JSON files in a directory * - IJobStore interface — implement for SQLite, Postgres, etc. * * Features: * - Priority-based FIFO queue * - Configurable retries with exponential backoff * - Job timeout enforcement * - Crash recovery: stale "running" jobs are re-queued on startup * - Pluggable storage backends * * Usage: * const queue = new JobQueue({ store: new FileJobStore('./data/jobs') }); * await queue.start(); * await queue.enqueue({ type: 'delegateTask', payload: { agentId: 'analyzer', ... } }); * * @module JobQueue * @version 1.0.0 */ import { EventEmitter } from 'events'; /** Job status */ export type JobStatus = 'pending' | 'running' | 'completed' | 'failed' | 'cancelled'; /** Priority level (lower = higher priority) */ export type JobPriority = 0 | 1 | 2 | 3 | 4 | 5; /** A persistent job record */ export interface JobRecord { /** Unique job ID */ id: string; /** Job type (e.g. 'delegateTask', 'batchProcess') */ type: string; /** Job payload — any serializable data */ payload: Record; /** Current status */ status: JobStatus; /** Priority (0 = highest, 5 = lowest; default: 2) */ priority: JobPriority; /** Number of attempts so far */ attempts: number; /** Maximum attempts (default: 3) */ maxAttempts: number; /** When the job was created */ createdAt: number; /** When the job was last updated */ updatedAt: number; /** When the job started running */ startedAt?: number; /** When the job completed or failed */ completedAt?: number; /** Result data (on success) */ result?: unknown; /** Error message (on failure) */ error?: string; /** Job timeout in ms (default: 300000 = 5 min) */ timeoutMs: number; /** Metadata for tracking */ metadata?: Record; } /** Options for creating a new job */ export interface JobCreateOptions { /** Job type */ type: string; /** Job payload */ payload: Record; /** Priority (default: 2) */ priority?: JobPriority; /** Max attempts (default: 3) */ maxAttempts?: number; /** Timeout in ms (default: 300000) */ timeoutMs?: number; /** Additional metadata */ metadata?: Record; } /** * Job handler function — processes a single job. * Return a result on success, or throw to fail. */ export type JobHandler = (job: JobRecord) => Promise; /** Job queue configuration */ export interface JobQueueConfig { /** Storage backend */ store: IJobStore; /** Polling interval in ms (default: 1000) */ pollIntervalMs?: number; /** Maximum concurrent jobs (default: 5) */ concurrency?: number; /** Base retry delay in ms (default: 1000) */ retryBaseDelayMs?: number; /** Maximum retry delay in ms (default: 60000) */ retryMaxDelayMs?: number; /** Stale job threshold in ms — re-queue "running" jobs older than this (default: 600000 = 10 min) */ staleThresholdMs?: number; } /** Stats snapshot */ export interface JobQueueStats { pending: number; running: number; completed: number; failed: number; cancelled: number; total: number; } /** * Pluggable storage backend for the job queue. * Implement this for SQLite, Postgres, Redis, etc. */ export interface IJobStore { /** Initialize the store (create tables/dirs) */ init(): Promise; /** Save or update a job */ save(job: JobRecord): Promise; /** Get a job by ID */ get(id: string): Promise; /** Delete a job by ID */ delete(id: string): Promise; /** List jobs by status, ordered by priority ASC then createdAt ASC */ listByStatus(status: JobStatus, limit?: number): Promise; /** Count jobs by status */ countByStatus(status: JobStatus): Promise; /** Find stale running jobs (startedAt < threshold) */ findStale(thresholdMs: number): Promise; } /** * File-system job store — persists jobs as individual JSON files. * Simple and dependency-free. Suitable for single-node deployments. */ export declare class FileJobStore implements IJobStore { private readonly dir; constructor(dir: string); init(): Promise; save(job: JobRecord): Promise; get(id: string): Promise; delete(id: string): Promise; listByStatus(status: JobStatus, limit?: number): Promise; countByStatus(status: JobStatus): Promise; findStale(thresholdMs: number): Promise; } /** * Persistent job queue with priority ordering, retries, and crash recovery. * * Register handlers for job types, then start the queue. Jobs are * persisted to the configured store and survive process restarts. */ export declare class JobQueue extends EventEmitter { private handlers; private activeJobs; private pollTimer; private running; private readonly store; private readonly pollIntervalMs; private readonly concurrency; private readonly retryBaseDelayMs; private readonly retryMaxDelayMs; private readonly staleThresholdMs; constructor(config: JobQueueConfig); /** Register a handler for a job type */ handle(type: string, handler: JobHandler): void; /** Start the queue: init store, recover stale jobs, begin polling */ start(): Promise; /** Stop the queue (active jobs continue but no new ones are dequeued) */ stop(): Promise; /** Enqueue a new job */ enqueue(options: JobCreateOptions): Promise; /** Cancel a pending or running job */ cancel(jobId: string): Promise; /** Get a job by ID */ getJob(jobId: string): Promise; /** Get queue stats */ stats(): Promise; /** Whether the queue is running */ get isRunning(): boolean; private poll; private processJob; private recoverStaleJobs; } //# sourceMappingURL=job-queue.d.ts.map