{"_id":"@alcyone-labs/workalot","name":"@alcyone-labs/workalot","dist-tags":{"latest":"1.0.0"},"versions":{"1.0.0":{"name":"@alcyone-labs/workalot","version":"1.0.0","description":"High-performance, multi-threaded job queue system for NodeJS/BunJS. Achieve high throughput jobs/sec with linear worker scaling.","main":"dist/src/index.js","types":"dist/src/index.d.ts","type":"module","keywords":["job-queue","multi-threaded","worker-threads","high-performance","task-management","scaling"],"author":{"name":"Nicolas Embleton","email":"nicolas.embleton@gmail.com"},"license":"MIT","homepage":"https://github.com/alcyone-labs/workalot#readme","repository":{"type":"git","url":"git+https://github.com/alcyone-labs/workalot.git"},"bugs":{"url":"https://github.com/alcyone-labs/workalot/issues"},"engines":{"node":">=18.0.0"},"publishConfig":{"access":"public"},"dependencies":{"ulidx":"^2.4.1"},"optionalDependencies":{"@electric-sql/pglite":"^0.3.4","better-sqlite3":"^12.1.1"},"devDependencies":{"@types/better-sqlite3":"^7.6.13","@types/cli-progress":"^3.11.6","@types/node":"^22.15.33","cli-progress":"^3.12.0","typescript":"^5.8.3","vite":"^6.3.5","vite-tsconfig-paths":"^5.1.4","vitest":"^3.2.4"},"scripts":{"build":"tsc && cp -r src/queue/migrations dist/src/queue/","dev":"tsc --watch","test":"vitest","test:run":"vitest run","test:coverage":"vitest run --coverage","benchmark":"pnpm run build && bun run benchmarks/run-benchmarks.ts","benchmark:easy":"pnpm run build && bun run benchmarks/run-benchmarks.ts --difficulty easy","benchmark:quick":"pnpm run build && bun run benchmarks/run-benchmarks.ts --configs 2-cores-10k-jobs,4-cores-10k-jobs","benchmark:quick:easy":"pnpm run build && bun run benchmarks/run-benchmarks.ts --configs 2-cores-10k-jobs,4-cores-10k-jobs --difficulty easy","benchmark:nodejs":"pnpm run build && node dist/benchmarks/run-benchmarks.js --difficulty easy","benchmark:nodejs:quick":"pnpm run build && node dist/benchmarks/run-benchmarks.js --configs 2-cores-10k-jobs,4-cores-10k-jobs --difficulty easy","benchmark:bun":"bun run benchmarks/run-benchmarks.ts --difficulty easy","benchmark:bun:quick":"bun run benchmarks/run-benchmarks.ts --configs 2-cores-10k-jobs,4-cores-10k-jobs --difficulty easy","benchmark:deno":"deno run --allow-all benchmarks/run-benchmarks.ts --difficulty easy","benchmark:deno:quick":"deno run --allow-all benchmarks/run-benchmarks.ts --configs 2-cores-10k-jobs,4-cores-10k-jobs --difficulty easy","benchmark:compare":"./benchmarks/compare-runtimes.sh easy 2-cores-10k-jobs,4-cores-10k-jobs","benchmark:compare:full":"./benchmarks/compare-runtimes.sh normal 2-cores-10k-jobs,4-cores-10k-jobs,6-cores-10k-jobs","benchmark:sqlite":"./benchmarks/run-sqlite-benchmarks.sh easy","benchmark:sqlite:normal":"./benchmarks/run-sqlite-benchmarks.sh normal","benchmark:backends":"./benchmarks/run-backend-comparison.sh easy 4 1k","benchmark:backends:10k":"./benchmarks/run-backend-comparison.sh easy 4 10k","example:quick-start":"pnpm run build && bun run examples/quick-start.ts","example:performance":"pnpm run build && bun run examples/performance-test.ts","example:errors":"pnpm run build && bun run examples/error-handling.ts","examples":"pnpm run example:quick-start","clean":"rm -rf dist"},"_id":"@alcyone-labs/workalot@1.0.0","_integrity":"sha512-P0qQ3D04ahzBIaXwzMPCBl0oKYlQg2dCU7evBN7ytl6lYNVvabi96AiQ6VwcGnzc/kJsVUugIyeAR9u1VNCqMw==","_resolved":"/private/var/folders/27/xlh5p6rd54vgnqk_t2kq551c0000gn/T/6f1c12aefa2c516e72cd36cd118c5cb3/alcyone-labs-workalot-1.0.0.tgz","_from":"file:alcyone-labs-workalot-1.0.0.tgz","_nodeVersion":"23.6.0","_npmVersion":"10.9.2","dist":{"integrity":"sha512-P0qQ3D04ahzBIaXwzMPCBl0oKYlQg2dCU7evBN7ytl6lYNVvabi96AiQ6VwcGnzc/kJsVUugIyeAR9u1VNCqMw==","shasum":"2ee0648fd9462c6a9678acd1e376916351f56448","tarball":"https://registry.npmjs.org/@alcyone-labs/workalot/-/workalot-1.0.0.tgz","fileCount":153,"unpackedSize":648189,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIF39pIvl1FUfn6nxA1vDe7hagxRFoRVCDCVCDLfsUq82AiEA9zefrV3C3ynisnrmafsOxOVa7S5k63b9ION2lvEf1Ms="}]},"_npmUser":{"name":"nembleton","email":"nicolas.embleton@gmail.com","actor":{"name":"nembleton","email":"nicolas.embleton@gmail.com","type":"user"}},"directories":{},"maintainers":[{"name":"nembleton","email":"nicolas.embleton@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/workalot_1.0.0_1751134647056_0.948413836364357"},"_hasShrinkwrap":false}},"time":{"created":"2025-06-28T18:17:26.949Z","1.0.0":"2025-06-28T18:17:27.288Z","modified":"2025-06-28T18:17:27.647Z"},"maintainers":[{"name":"nembleton","email":"nicolas.embleton@gmail.com"}],"description":"High-performance, multi-threaded job queue system for NodeJS/BunJS. Achieve high throughput jobs/sec with linear worker scaling.","homepage":"https://github.com/alcyone-labs/workalot#readme","keywords":["job-queue","multi-threaded","worker-threads","high-performance","task-management","scaling"],"repository":{"type":"git","url":"git+https://github.com/alcyone-labs/workalot.git"},"author":{"name":"Nicolas Embleton","email":"nicolas.embleton@gmail.com"},"bugs":{"url":"https://github.com/alcyone-labs/workalot/issues"},"license":"MIT","readme":"# @alcyone-labs/workalot\n\nA high-performance, multi-threaded job queue system for NodeJS/BunJS. Achieves high jobs/sec throughput with linear scaling across CPU cores. Includes comprehensive job recovery, fault tolerance, and monitoring.\n\n## Key Features\n\n- **High performance** - Thousands of jobs/sec with optimized job distribution and worker orchestration\n- **Linear scaling** - Scales across CPU cores (2→4→6→N cores) with significant performance improvements\n- **Job recovery system** - Automatic recovery of stalled/crashed jobs with configurable timeouts and retry limits\n- **Multiple backends** - In-memory, SQLite, PGLite, PostgreSQL with backend-specific optimizations\n- **Fault tolerant** - Worker crash detection, job recovery, graceful error handling, and automatic cleanup\n- **Developer-friendly** - Promise-based API, TypeScript support, dynamic job loading, job scheduling from within jobs\n- **Real-time monitoring** - Detailed statistics, worker utilization, performance metrics, and job recovery tracking\n- **Runtime flexibility** - Native BunJS and NodeJS support with TypeScript execution\n\n## Quick Start\n\n### Installation\n\n```bash\nnpm install @alcyone-labs/workalot\n# or\npnpm add @alcyone-labs/workalot\n# or\ndeno add npm:@alcyone-labs/workalot\n# or\nyarn add @alcyone-labs/workalot\n```\n\n#### Optional Database Dependencies\n\nFor specific queue backends, install the corresponding database driver:\n\n```bash\n# For PGLite backend (PostgreSQL-compatible embedded database)\nnpm install @electric-sql/pglite\n\n# For SQLite backend\n# - Bun runtime: Uses built-in bun:sqlite (no installation needed)\n# - Node.js runtime: Requires better-sqlite3\nnpm install better-sqlite3  # Only needed for Node.js\n\n# Memory backend requires no additional dependencies\n```\n\n### Basic Usage\n\n#### Option 1: PGLite In-Memory (Recommended)\n\n**PostgreSQL features with memory performance**\n\n```typescript\nimport { initializeTaskManager, scheduleAndWait, whenFree, shutdown } from '@alcyone-labs/workalot';\n\n// Initialize with PGLite in-memory backend\nawait initializeTaskManager({\n  backend: 'pglite',\n  databaseUrl: 'memory://',  // In-memory PostgreSQL-compatible database\n  maxThreads: 4\n});\n\n// Schedule a job and wait for completion\nconst result = await scheduleAndWait({\n  jobFile: 'jobs/ProcessDataJob.ts',\n  jobPayload: { data: [1, 2, 3, 4, 5] },\n  jobTimeout: 10000\n});\n\nconsole.log('Job completed:', result);\n\n// Get notified when queue is free\nwhenFree(() => {\n  console.log('All jobs completed!');\n});\n\n// Graceful shutdown\nawait shutdown();\n```\n\n**Why PGLite In-Memory?**\n- **Full PostgreSQL compatibility** - SQL queries, transactions, notifications\n- **Memory-level performance** - No disk I/O overhead\n- **Advanced features** - Real-time job monitoring, complex analytics\n- **No file dependencies** - Good for containers and serverless\n\n#### Option 2: Pure Memory (Maximum Speed)\n\n```typescript\n// Initialize with memory backend for maximum speed\nawait initializeTaskManager({\n  backend: 'memory',\n  persistenceFile: 'queue-state.json',  // Optional JSON persistence\n  maxThreads: 4\n});\n```\n\n#### Option 3: PGLite File-Based (Full Persistence)\n\n```typescript\n// Initialize with file-based PGLite for full persistence\nawait initializeTaskManager({\n  backend: 'pglite',\n  databaseUrl: './data/queue.db',  // File-based database\n  maxThreads: 4\n});\n```\n\n#### Performance Comparison Chart (a bit busy)\n\n![Performance Comparison](docs/Workalot/jobs-per-seconds-chart.png)\n\n## API Reference\n\n### Main Functions\n\n#### `initializeTaskManager(config?, projectRoot?)`\n\nInitialize the task management system.\n\n**Parameters:**\n- `config` (optional): Configuration object\n- `projectRoot` (optional): Project root directory\n\n**Configuration Options:**\n```typescript\ninterface QueueConfig {\n  maxThreads?: number;                          // Default: CPU cores - 2\n  maxInMemoryAge?: number;                      // Default: 24 hours (ms)\n  persistenceFile?: string;                     // Default: 'queue-state.json' (memory backend only)\n  healthCheckInterval?: number;                 // Default: 5000ms\n\n  // Backend configuration\n  backend?: 'memory' | 'sqlite' | 'pglite' | 'postgresql'; // Default: 'memory'\n  databaseUrl?: string;                         // For SQLite/PGLite/PostgreSQL backends\n\n  // Job recovery configuration\n  jobRecovery?: {\n    enabled?: boolean;                          // Default: true\n    checkInterval?: number;                     // Default: 60000ms (1 minute)\n    stalledTimeout?: number;                    // Default: 300000ms (5 minutes)\n    maxRecoveryAttempts?: number;               // Default: 3\n  };\n}\n```\n\n#### `scheduleAndWait(jobPayload)`\n\nSchedule a job and wait for it to complete. Returns a promise that resolves when **that specific job** completes.\n\n**Important:** `scheduleAndWait()` follows FIFO (first-in-first-out) queue order. If you have 1,000 jobs already queued, your `scheduleAndWait()` job becomes job #1,001 and will execute after all previous jobs complete. The function waits for completion, but does not provide immediate execution.\n\n**Parameters:**\n```typescript\ninterface JobPayload {\n  jobFile: string;                              // Path to job file\n  jobPayload: Record<string, any>;              // Data to pass to job\n  jobTimeout?: number;                          // Execution timeout (default: 5000ms)\n}\n```\n\n**Returns:** `Promise<JobResult>`\n\n```typescript\ninterface JobResult {\n  results: Record<string, any>; // Job output\n  executionTime: number;        // Time taken (ms)\n  queueTime: number;            // Time in queue (ms)\n}\n```\n\n**How it works:**\n1. Calls `schedule()` internally to add your job to the queue\n2. Waits for that specific job ID to complete (using event listeners)\n3. Returns the result when your job finishes, regardless of other jobs\n4. Properly handles timeouts, errors, and cleanup\n\n#### `schedule(jobPayload)`\n\nSchedule a job without waiting for completion (fire and forget).\n\n**Returns:** `Promise<string>` - Job ID\n\n#### `whenFree(callback)`\n\nRegister a callback to be called when the queue becomes free (no pending jobs).\n\n**Parameters:**\n- `callback`: Function to call when queue is free\n\n\n### Utility Functions\n\n- `getStatus()` - Get comprehensive system status including job recovery stats\n- `isIdle()` - Check if system is idle (no pending jobs)\n- `whenIdle(timeoutMs?)` - Wait for system to become idle\n- `getQueueStats()` - Get detailed queue statistics\n- `getWorkerStats()` - Get worker pool statistics\n- `getJobsByStatus(status)` - Get jobs by their status\n- `removeWhenFreeCallback(callback)` - Remove a whenFree callback\n- `getTaskManager()` - Get the underlying TaskManager instance\n- `isInitialized()` - Check if the system is initialized\n- `shutdown()` - Graceful shutdown with job recovery cleanup\n\n### Job Recovery Functions\n\n- `recoverStalledJobs()` - Manually trigger recovery of stalled jobs\n- `getJobRecoveryStats()` - Get detailed job recovery statistics\n\n## Creating Jobs\n\nJobs are TypeScript/JavaScript classes that implement the `IJob` interface by extending `BaseJob`. This section covers everything you need to know about creating, structuring, and using jobs.\n\n### Job File Requirements\n\n#### File Extensions and Runtime Support\n\nWorkalot supports multiple file formats with automatic runtime detection:\n\n```typescript\n// Supported file extensions\n'.ts'   // TypeScript - native support in BunJS/Deno, requires compilation for Node.js\n'.js'   // JavaScript - standard Node.js/Bun/Deno support\n'.mjs'  // ES Modules - explicit module format\n```\n\n**Runtime Behavior:**\n- **BunJS**: Executes TypeScript files directly without compilation\n- **Node.js**: Requires compilation to JavaScript first, then reference the compiled `.js` files\n- **Deno**: Executes TypeScript files directly without compilation, but requires specific flags\n\n#### File Path Resolution\n\nJob file paths in `scheduleAndWait()` are resolved relative to your project root:\n\n```typescript\n// Relative paths (recommended)\nawait scheduleAndWait({\n  jobFile: 'jobs/ProcessDataJob.ts',        // ./jobs/ProcessDataJob.ts\n  jobFile: 'src/workers/EmailJob.ts',       // ./src/workers/EmailJob.ts\n  jobFile: 'lib/tasks/ReportJob.ts',        // ./lib/tasks/ReportJob.ts\n});\n\n// Subdirectory paths work seamlessly\nawait scheduleAndWait({\n  jobFile: 'modules/analytics/MetricsJob.ts',\n  jobFile: 'services/notifications/SlackJob.ts'\n});\n```\n\n**Path Resolution Rules:**\n1. Paths are resolved relative to the project root (where you initialize TaskManager)\n2. Use forward slashes `/` on all platforms (Windows, macOS, Linux)\n3. File must be readable and accessible at runtime\n4. Extension is validated (`.ts`, `.js`, `.mjs` only)\n\n### Job Class Structure\n\n#### Basic Job Template\n\n```typescript\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\nexport class MyJob extends BaseJob {\n  constructor() {\n    super('MyJob'); // Optional: custom job name\n  }\n\n  async run(payload: Record<string, any>, context: JobExecutionContext): Promise<Record<string, any>> {\n    // 1. Validate input data\n    this.validatePayload(payload, ['requiredField']);\n\n    // 2. Process data\n    const result = await this.processData(payload.data);\n\n    // 3. Schedule follow-up jobs if needed\n    if (payload.needsFollowUp) {\n      const followUpJobId = await context.scheduleAndWait({\n        jobFile: 'jobs/FollowUpJob.ts',\n        jobPayload: { originalResult: result }\n      });\n      console.log('Scheduled follow-up job:', followUpJobId);\n    }\n\n    // 4. Return success response\n    return this.createSuccessResult({ result });\n  }\n\n  private async processData(data: any): Promise<any> {\n    // Your job logic here\n    return data;\n  }\n}\n```\n\n### Scheduling Jobs from Within Jobs\n\nJobs can schedule other jobs using the `context` parameter provided to the `run()` method. This enables powerful workflow patterns:\n\n```typescript\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\nexport class WorkflowJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext): Promise<Record<string, any>> {\n    this.validatePayload(payload, ['workflowType']);\n\n    // Schedule a job and wait for its completion\n    const processingRequestId = await context.scheduleAndWait({\n      jobFile: 'jobs/DataProcessorJob.ts',\n      jobPayload: { data: payload.inputData }\n    });\n\n    // Schedule background jobs (fire-and-forget)\n    const cleanupRequestId = context.schedule({\n      jobFile: 'jobs/CleanupJob.ts',\n      jobPayload: { workflowId: payload.id }\n    });\n\n    const notificationRequestId = context.schedule({\n      jobFile: 'jobs/NotificationJob.ts',\n      jobPayload: {\n        message: 'Workflow completed',\n        userId: payload.userId\n      }\n    });\n\n    return this.createSuccessResult({\n      workflowCompleted: true,\n      scheduledJobs: {\n        processing: processingRequestId,\n        cleanup: cleanupRequestId,\n        notification: notificationRequestId\n      }\n    });\n  }\n}\n```\n\n#### Scheduling API\n\n- **`context.scheduleAndWait(jobPayload)`**: Schedule a job for execution after current job completes. Returns a request ID.\n- **`context.schedule(jobPayload)`**: Schedule a job for fire-and-forget execution. Returns a request ID immediately.\n\nBoth functions accumulate scheduling requests during job execution and process them automatically after the job completes.\n\n**Note**: The `context` parameter is optional for backward compatibility. Existing jobs without the context parameter will continue to work unchanged.\n\n#### Required Implementation\n\nEvery job must:\n\n1. **Extend `BaseJob`**: Provides essential functionality and helper methods\n2. **Implement `run()` method**: Main job execution logic with optional context parameter\n3. **Export the class**: Must be exportable for dynamic loading\n\n```typescript\n// Correct - extends BaseJob and implements run()\nexport class ValidJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext): Promise<Record<string, any>> {\n    return this.createSuccessResult({ status: 'completed' });\n  }\n}\n\n// Incorrect - missing BaseJob extension\nexport class InvalidJob {\n  async run(payload: Record<string, any>) {\n    return { status: 'completed' };\n  }\n}\n```\n\n### Job Export Patterns\n\nWorkalot supports multiple export patterns for maximum flexibility:\n\n#### Default Export (Recommended)\n\n```typescript\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\nexport default class ProcessDataJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    return this.createSuccessResult({ processed: true });\n  }\n}\n```\n\n#### Named Export\n\n```typescript\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\nexport class ProcessDataJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    return this.createSuccessResult({ processed: true });\n  }\n}\n```\n\n**Auto-Discovery Rules:**\n1. **Default export**: Used if available and is a class\n2. **Named export**: Searches for class names matching:\n   - Exact filename: `ProcessDataJob.ts` → `ProcessDataJob`\n   - Capitalized filename: `processDataJob.ts` → `ProcessDataJob`\n   - With \"Job\" suffix: `processData.ts` → `ProcessDataJob`\n\n### Job Payload and Validation\n\n#### Input Validation\n\nUse `validatePayload()` to ensure required fields are present:\n\n```typescript\nexport class EmailJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    // Validate required fields\n    this.validatePayload(payload, ['to', 'subject', 'body']);\n\n    const { to, subject, body, attachments } = payload;\n\n    // Optional fields can be accessed safely\n    const priority = payload.priority || 'normal';\n\n    // Schedule follow-up notification if high priority\n    if (priority === 'high') {\n      const notificationRequestId = context.schedule({\n        jobFile: 'jobs/NotificationJob.ts',\n        jobPayload: {\n          type: 'email_sent',\n          emailId: 'email-123',\n          recipient: to\n        }\n      });\n    }\n\n    // Process email sending...\n    return this.createSuccessResult({\n      emailId: 'email-123',\n      sentAt: new Date().toISOString()\n    });\n  }\n}\n```\n\n#### TypeScript Type Safety\n\nFor better type safety, define payload interfaces:\n\n```typescript\ninterface EmailPayload {\n  to: string;\n  subject: string;\n  body: string;\n  attachments?: string[];\n  priority?: 'low' | 'normal' | 'high';\n}\n\nexport class EmailJob extends BaseJob {\n  async run(payload: EmailPayload, context: JobExecutionContext): Promise<Record<string, any>> {\n    this.validatePayload(payload, ['to', 'subject', 'body']);\n\n    // TypeScript now provides full type checking\n    const recipient = payload.to;\n    const isHighPriority = payload.priority === 'high';\n\n    // Schedule audit log for high priority emails\n    if (isHighPriority) {\n      context.schedule({\n        jobFile: 'jobs/AuditLogJob.ts',\n        jobPayload: {\n          action: 'high_priority_email_sent',\n          recipient,\n          timestamp: new Date().toISOString()\n        }\n      });\n    }\n\n    return this.createSuccessResult({\n      emailId: 'email-123',\n      recipient,\n      priority: payload.priority || 'normal'\n    });\n  }\n}\n```\n\n### Job ID Generation\n\n#### Default Behavior\n\nBy default, jobs generate monotonic, time-sortable ULID identifiers:\n\n```typescript\n// Default ID generation (ULID - monotonic, time-sortable)\nexport class DefaultJob extends BaseJob {\n  // Uses inherited getJobId() method\n  // ID = ULID (e.g., \"01ARZ3NDEKTSV4RRFFQ69G5FAV\")\n}\n```\n\n#### Custom ID Generation\n\nOverride `getJobId()` for custom ID logic:\n\n```typescript\nexport class CustomIdJob extends BaseJob {\n  getJobId(payload?: Record<string, any>): string | undefined {\n    if (!payload) return undefined;\n\n    // Custom ID based on business logic\n    if (payload.userId && payload.action) {\n      return `${payload.userId}-${payload.action}-${Date.now()}`;\n    }\n\n    // Fall back to default behavior\n    return super.getJobId(payload);\n  }\n\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    return this.createSuccessResult({ processed: true });\n  }\n}\n```\n\n**ID Generation Guidelines:**\n- Return `undefined` for auto-generated UUIDs\n- Return `string` for custom IDs\n- Ensure IDs are unique to prevent job collisions\n- Consider including timestamps for uniqueness\n\n### Helper Methods\n\nThe `BaseJob` class provides several helper methods for common operations:\n\n#### Success and Error Results\n\n```typescript\nexport class DataProcessorJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    try {\n      const result = await this.processData(payload.data);\n\n      // Schedule cleanup job after successful processing\n      context.schedule({\n        jobFile: 'jobs/CleanupJob.ts',\n        jobPayload: {\n          tempFiles: result.tempFiles,\n          processedAt: new Date().toISOString()\n        }\n      });\n\n      // Return success result\n      return this.createSuccessResult({\n        recordsProcessed: result.count,\n        outputFile: result.filename,\n        processingTime: result.duration\n      });\n\n    } catch (error) {\n      // Return error result (optional - you can also throw)\n      return this.createErrorResult(\n        'Data processing failed',\n        { originalError: error.message }\n      );\n    }\n  }\n}\n```\n\n#### Runtime Payload Validation\n\n```typescript\nexport class ValidationExampleJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    // Validate multiple required fields\n    this.validatePayload(payload, ['userId', 'action', 'timestamp']);\n\n    // Validate nested objects\n    if (payload.config) {\n      this.validatePayload(payload.config, ['apiKey', 'endpoint']);\n    }\n\n    // Custom validation logic\n    if (payload.timestamp < Date.now() - 86400000) {\n      throw new Error('Timestamp cannot be older than 24 hours');\n    }\n\n    return this.createSuccessResult({ validated: true });\n  }\n}\n```\n\n### Error Handling in Jobs\n\n#### Throwing Errors\n\nThe recommended approach is to throw errors for failures:\n\n```typescript\nexport class ErrorHandlingJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    this.validatePayload(payload, ['operation']);\n\n    if (payload.operation === 'invalid') {\n      // Throw error - will be caught and handled by the system\n      throw new Error('Invalid operation requested');\n    }\n\n    if (payload.timeout && payload.timeout < 1000) {\n      // Throw with specific error types\n      throw new Error('Timeout must be at least 1000ms');\n    }\n\n    return this.createSuccessResult({ status: 'completed' });\n  }\n}\n```\n\n#### Error Context and Debugging\n\n```typescript\nexport class DebuggableJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    try {\n      // Risky operation\n      const result = await this.performRiskyOperation(payload);\n      return this.createSuccessResult(result);\n\n    } catch (error) {\n      // Add context to errors for better debugging\n      const enhancedError = new Error(\n        `Job failed during ${payload.operation}: ${error.message}`\n      );\n      enhancedError.cause = error;\n      throw enhancedError;\n    }\n  }\n\n  private async performRiskyOperation(payload: any) {\n    // Simulate operation that might fail\n    if (Math.random() < 0.1) {\n      throw new Error('Random failure for testing');\n    }\n    return { success: true };\n  }\n}\n```\n\n### Complete Job Examples\n\n#### Simple Processing Job\n\n```typescript\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\nexport class SimpleProcessorJob extends BaseJob {\n  async run(payload: { items: string[] }, context: JobExecutionContext) {\n    this.validatePayload(payload, ['items']);\n\n    const processedItems = payload.items.map(item => ({\n      original: item,\n      processed: item.toUpperCase(),\n      timestamp: new Date().toISOString()\n    }));\n\n    // Schedule notification job if processing large batch\n    if (payload.items.length > 100) {\n      context.schedule({\n        jobFile: 'jobs/NotificationJob.ts',\n        jobPayload: {\n          message: `Processed ${payload.items.length} items`,\n          type: 'batch_complete'\n        }\n      });\n    }\n\n    return this.createSuccessResult({\n      totalItems: payload.items.length,\n      processedItems\n    });\n  }\n}\n```\n\n#### Complex Business Logic Job\n\n```typescript\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\ninterface ReportPayload {\n  userId: string;\n  reportType: 'daily' | 'weekly' | 'monthly';\n  dateRange: {\n    start: string;\n    end: string;\n  };\n  includeCharts?: boolean;\n}\n\nexport class ReportGeneratorJob extends BaseJob {\n  constructor() {\n    super('ReportGeneratorJob');\n  }\n\n  async run(payload: ReportPayload, context: JobExecutionContext): Promise<Record<string, any>> {\n    // Validate required fields\n    this.validatePayload(payload, ['userId', 'reportType', 'dateRange']);\n    this.validateDateRange(payload.dateRange);\n\n    console.log(`Generating ${payload.reportType} report for user ${payload.userId}`);\n\n    try {\n      // Simulate data collection\n      const data = await this.collectData(payload);\n\n      // Generate report\n      const report = await this.generateReport(data, payload);\n\n      // Optionally generate charts\n      let charts = null;\n      if (payload.includeCharts) {\n        charts = await this.generateCharts(data);\n      }\n\n      // Schedule follow-up jobs\n      const reportId = `report-${Date.now()}`;\n\n      // Schedule email notification\n      const emailRequestId = context.schedule({\n        jobFile: 'jobs/EmailNotificationJob.ts',\n        jobPayload: {\n          userId: payload.userId,\n          reportId,\n          reportType: payload.reportType,\n          downloadUrl: report.url\n        }\n      });\n\n      // Schedule cleanup job for temporary files\n      context.schedule({\n        jobFile: 'jobs/CleanupJob.ts',\n        jobPayload: {\n          reportId,\n          tempFiles: report.tempFiles,\n          scheduleAfter: Date.now() + 24 * 60 * 60 * 1000 // 24 hours\n        }\n      });\n\n      return this.createSuccessResult({\n        reportId,\n        reportType: payload.reportType,\n        userId: payload.userId,\n        generatedAt: new Date().toISOString(),\n        dataPoints: data.length,\n        reportSize: report.size,\n        chartsIncluded: !!charts,\n        downloadUrl: report.url\n      });\n\n    } catch (error) {\n      console.error(`Report generation failed for user ${payload.userId}:`, error);\n      throw new Error(`Report generation failed: ${error.message}`);\n    }\n  }\n\n  private validateDateRange(dateRange: { start: string; end: string }) {\n    const start = new Date(dateRange.start);\n    const end = new Date(dateRange.end);\n\n    if (isNaN(start.getTime()) || isNaN(end.getTime())) {\n      throw new Error('Invalid date format in dateRange');\n    }\n\n    if (start >= end) {\n      throw new Error('Start date must be before end date');\n    }\n  }\n\n  private async collectData(payload: ReportPayload) {\n    // Simulate data collection with processing time\n    await new Promise(resolve => setTimeout(resolve, 500));\n\n    // Return mock data based on report type\n    const baseCount = payload.reportType === 'daily' ? 24 :\n                     payload.reportType === 'weekly' ? 168 : 720;\n\n    return Array.from({ length: baseCount }, (_, i) => ({\n      timestamp: new Date(Date.now() - i * 3600000).toISOString(),\n      value: Math.floor(Math.random() * 100)\n    }));\n  }\n\n  private async generateReport(data: any[], payload: ReportPayload) {\n    // Simulate report generation\n    await new Promise(resolve => setTimeout(resolve, 300));\n\n    return {\n      size: `${Math.floor(data.length * 0.1)}KB`,\n      url: `https://reports.example.com/report-${Date.now()}.pdf`\n    };\n  }\n\n  private async generateCharts(data: any[]) {\n    // Simulate chart generation\n    await new Promise(resolve => setTimeout(resolve, 200));\n\n    return {\n      chartCount: 3,\n      types: ['line', 'bar', 'pie']\n    };\n  }\n\n  getJobId(payload?: ReportPayload): string | undefined {\n    if (!payload) return undefined;\n\n    // Create unique ID based on user and date range\n    const content = `${payload.userId}-${payload.reportType}-${payload.dateRange.start}-${payload.dateRange.end}`;\n    return require('crypto').createHash('sha1').update(content).digest('hex');\n  }\n}\n```\n\n### TypeScript vs JavaScript Jobs\n\n#### TypeScript Jobs (Recommended)\n\nTypeScript provides better development experience with type safety and IDE support:\n\n```typescript\n// jobs/TypedJob.ts\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\ninterface MyJobPayload {\n  userId: string;\n  action: 'create' | 'update' | 'delete';\n  data: Record<string, any>;\n}\n\nexport class TypedJob extends BaseJob {\n  async run(payload: MyJobPayload, context: JobExecutionContext): Promise<Record<string, any>> {\n    // Full type checking and IntelliSense support\n    this.validatePayload(payload, ['userId', 'action', 'data']);\n\n    // TypeScript catches errors at compile time\n    const isCreateAction = payload.action === 'create';\n\n    // Schedule audit log for all actions\n    context.schedule({\n      jobFile: 'jobs/AuditLogJob.ts',\n      jobPayload: {\n        userId: payload.userId,\n        action: payload.action,\n        timestamp: new Date().toISOString(),\n        data: payload.data\n      }\n    });\n\n    return this.createSuccessResult({\n      userId: payload.userId,\n      actionPerformed: payload.action,\n      success: true\n    });\n  }\n}\n```\n\n**Usage with TypeScript:**\n```typescript\n// BunJS - Direct TypeScript execution\nawait scheduleAndWait({\n  jobFile: 'jobs/TypedJob.ts',  // .ts extension works directly\n  jobPayload: { userId: '123', action: 'create', data: {} }\n});\n\n// Node.js - Must compile first and reference compiled output\n// 1. Compile: pnpm run build (or tsc, vite build, etc.)\n// 2. Reference the compiled .js file in dist/ (or your tsconfig outDir)\nawait scheduleAndWait({\n  jobFile: 'dist/jobs/TypedJob.js',  // Must point to compiled .js file\n  jobPayload: { userId: '123', action: 'create', data: {} }\n});\n```\n\n#### JavaScript Jobs\n\nPlain JavaScript works seamlessly but without compile-time type checking:\n\n```javascript\n// jobs/PlainJob.js\nconst { BaseJob } = require('@alcyone-labs/workalot');\n\nclass PlainJob extends BaseJob {\n  async run(payload, context) {\n    this.validatePayload(payload, ['userId', 'action']);\n\n    // No compile-time type checking\n    const result = await this.processAction(payload.action);\n\n    // Schedule follow-up job using context\n    context.schedule({\n      jobFile: 'jobs/LogJob.js',\n      jobPayload: {\n        userId: payload.userId,\n        action: payload.action,\n        result: result\n      }\n    });\n\n    return this.createSuccessResult({\n      userId: payload.userId,\n      result: result\n    });\n  }\n\n  async processAction(action) {\n    // Simulate processing\n    return { processed: true, action };\n  }\n}\n\nmodule.exports = { PlainJob };\n```\n\n**Usage with JavaScript:**\n```javascript\nawait taskManager.scheduleAndWait({\n  jobFile: 'jobs/PlainJob.js',  // .js extension\n  jobPayload: { userId: '123', action: 'process' }\n});\n```\n\n### Job File Organization\n\n#### Recommended Directory Structure\n\n```\nproject-root/\n├── jobs/                   # Job definitions\n│   ├── data/               # Data processing jobs\n│   │   ├── ProcessorJob.ts\n│   │   └── ValidatorJob.ts\n│   ├── notifications/      # Notification jobs\n│   │   ├── EmailJob.ts\n│   │   └── SlackJob.ts\n│   └── reports/           # Report generation jobs\n│       └── ReportJob.ts\n├── src/                   # Application code\n└── dist/                  # Compiled JavaScript (if needed)\n```\n\n#### Path Examples for Different Structures\n\n```typescript\n// Flat structure\nawait scheduleAndWait({ jobFile: 'jobs/EmailJob.ts' });\n\n// Nested structure\nawait scheduleAndWait({ jobFile: 'jobs/notifications/EmailJob.ts' });\nawait scheduleAndWait({ jobFile: 'jobs/data/ProcessorJob.ts' });\n\n// Mixed with source code\nawait scheduleAndWait({ jobFile: 'src/workers/BackgroundJob.ts' });\n\n// Compiled JavaScript (Node.js only - after compilation)\nawait scheduleAndWait({ jobFile: 'dist/jobs/EmailJob.js' });\n```\n\n### Common Patterns and Practices\n\n#### Job Naming Conventions\n\n```typescript\n// Good - Clear, descriptive names\nexport class UserRegistrationJob extends BaseJob { }\nexport class EmailNotificationJob extends BaseJob { }\nexport class DataBackupJob extends BaseJob { }\n\n// Avoid - Generic or unclear names\nexport class Job extends BaseJob { }\nexport class ProcessJob extends BaseJob { }\nexport class Handler extends BaseJob { }\n```\n\n#### Payload Design Patterns\n\n```typescript\n// Good - Structured, validated payloads\ninterface EmailJobPayload {\n  recipient: {\n    email: string;\n    name?: string;\n  };\n  message: {\n    subject: string;\n    body: string;\n    template?: string;\n  };\n  options?: {\n    priority: 'low' | 'normal' | 'high';\n    sendAt?: string;\n  };\n}\n\n// Avoid - Flat, unstructured payloads\ninterface BadPayload {\n  email: string;\n  subject: string;\n  body: string;\n  name: string;\n  priority: string;\n  sendAt: string;\n}\n```\n\n#### Error Handling Patterns\n\n```typescript\nexport class RobustJob extends BaseJob {\n  async run(payload: Record<string, any>, context: JobExecutionContext) {\n    try {\n      // Validate early\n      this.validatePayload(payload, ['operation']);\n\n      // Process with detailed error context\n      const result = await this.performOperation(payload);\n\n      return this.createSuccessResult(result);\n\n    } catch (error) {\n      // Log for debugging\n      console.error(`Job ${this.jobName} failed:`, {\n        payload,\n        error: error.message,\n        stack: error.stack\n      });\n\n      // Re-throw with context\n      throw new Error(`${this.jobName} failed: ${error.message}`);\n    }\n  }\n}\n```\n\n### Troubleshooting Job Issues\n\n#### Common Job Loading Errors\n\n```typescript\n// Error: \"Job file not found\"\nawait scheduleAndWait({\n  jobFile: 'wrong/path/Job.ts'  // File doesn't exist\n});\n\n// Solution: Verify file path relative to project root\nawait scheduleAndWait({\n  jobFile: 'jobs/MyJob.ts'  // Correct relative path\n});\n\n// Error: \"Unsupported job file extension\"\nawait scheduleAndWait({\n  jobFile: 'jobs/MyJob.py'  // Python not supported\n});\n\n// Solution: Use supported extensions\nawait scheduleAndWait({\n  jobFile: 'jobs/MyJob.ts'  // TypeScript\n});\n```\n\n#### Job Class Export Issues\n\n```typescript\n// Error: \"No valid job class found\"\n// Missing export\nclass MyJob extends BaseJob {\n  async run(payload) { return {}; }\n}\n\n// Solution: Export the class\nexport class MyJob extends BaseJob {\n  async run(payload, context) { return {}; }\n}\n\n// Alternative: Default export\nexport default class MyJob extends BaseJob {\n  async run(payload, context) { return {}; }\n}\n```\n\n#### Runtime Environment Issues\n\n```bash\n# BunJS - TypeScript works directly\nbun run app.ts\n# Job files: 'jobs/MyJob.ts'\n\n# Node.js - Requires compilation\nnpm run build\nnode dist/app.js\n# Job files: 'dist/jobs/MyJob.js'\n```\n\n## Queue Backends\n\n### PGLite In-Memory\n\n**The sweet spot: PostgreSQL features with memory performance**\n\n```typescript\nawait initializeTaskManager({\n  backend: 'pglite',\n  databaseUrl: 'memory://',  // In-memory PostgreSQL-compatible database\n  maxThreads: 4\n});\n```\n\n**Key Benefits:**\n- **Full SQL support** - Complex queries, joins, aggregations\n- **Real-time notifications** - PostgreSQL LISTEN/NOTIFY for job updates\n- **Memory performance** - No disk I/O, good for high-throughput\n- **Zero dependencies** - No external database required\n- **Container-friendly** - No persistent storage needed\n\n### Memory Backend (Maximum Speed)\n\nPure in-memory queue with optional JSON persistence:\n\n```typescript\nawait initializeTaskManager({\n  backend: 'memory',\n  persistenceFile: 'queue-state.json',  // Optional persistence\n  maxInMemoryAge: 24 * 60 * 60 * 1000   // 24 hours retention\n});\n```\n\n**Good for:** High throughput, temporary processing, simple use cases\n\n### SQLite Backend\n\nHigh-performance SQLite database with runtime-optimized drivers:\n\n```typescript\n// In-memory SQLite\nawait initializeTaskManager({\n  backend: 'sqlite',\n  databaseUrl: 'memory://',  // In-memory database\n  maxThreads: 4\n});\n\n// File-based SQLite with persistence\nawait initializeTaskManager({\n  backend: 'sqlite',\n  databaseUrl: './data/queue.db',  // File-based database\n  maxThreads: 4\n});\n```\n\n**Key Benefits:**\n- **Runtime-optimized** - Uses Bun's native SQLite when available, falls back to better-sqlite3 for Node.js\n- **High performance** - Good throughput with both in-memory and file-based storage\n- **SQL features** - Full SQLite support with indexes, transactions, and complex queries\n- **Cross-platform** - Works seamlessly with Bun, Node.js, and Deno\n- **Zero configuration** - Automatic driver selection and database setup\n\n**Good for:** Many use cases, development, high-performance applications\n\n### PGLite File-Based (Full Persistence)\n\nPostgreSQL-compatible embedded database with file persistence:\n\n```typescript\nawait initializeTaskManager({\n  backend: 'pglite',\n  databaseUrl: './data/queue.db',  // File-based database\n  maxThreads: 4\n});\n```\n\n**Good for:** Development, single-node deployments, full data persistence\n\n### Backend Performance Comparison\n\nRun comprehensive backend benchmarks:\n\n```bash\n# Test all backends with identical workloads\nbun run examples/performance-test.ts --backends\n\n# Or run both scaling and backend tests\nbun run examples/performance-test.ts --all\n```\n\n**Job Recovery Performance:**\n- **Recovery Detection**: < 1 minute (configurable check interval)\n- **Recovery Processing**: < 100ms per stalled job\n- **Zero Performance Impact**: Background monitoring with minimal overhead\n- **Automatic Cleanup**: Failed jobs after max retry attempts\n\n### Backend Selection Guide\n\n```typescript\n// For maximum throughput (temporary data)\n{ backend: 'memory' }\n\n// For high performance + SQL features (recommended)\n{ backend: 'sqlite', databaseUrl: 'memory://' }\n\n// For development with persistence\n{ backend: 'sqlite', databaseUrl: './dev-queue.db' }\n\n// For PostgreSQL compatibility\n{ backend: 'pglite', databaseUrl: 'memory://' }\n\n// For production with existing PostgreSQL\n{ backend: 'postgresql', databaseUrl: process.env.DATABASE_URL }\n```\n\n## BunJS Support\n\nThis library has native BunJS support and can run TypeScript files directly:\n\n```bash\n# Run with Bun\nbun run examples/sample-consumer/app.ts\n\n# Or with Node.js\npnpm run build\nnode dist/examples/sample-consumer/app.js\n```\n\n## Advanced Usage\n\n### Class-based API\n\nFor more control, use the class-based API:\n\n```typescript\nimport { TaskManager } from '@alcyone-labs/workalot';\n\nconst taskManager = new TaskManager({\n  backend: 'pglite',\n  databaseUrl: 'memory://',\n  maxThreads: 8\n});\n\nawait taskManager.initialize();\n\n// Use taskManager methods...\nawait taskManager.shutdown();\n```\n\n### Event Handling\n\n```typescript\nimport { TaskManager } from '@alcyone-labs/workalot';\n\nconst taskManager = new TaskManager();\n\ntaskManager.on('job-completed', (jobId, result) => {\n  console.log(`Job ${jobId} completed:`, result);\n});\n\ntaskManager.on('job-failed', (jobId, error) => {\n  console.log(`Job ${jobId} failed:`, error);\n});\n\ntaskManager.on('queue-empty', () => {\n  console.log('Queue is now empty');\n});\n\ntaskManager.on('all-workers-busy', () => {\n  console.log('All workers are currently busy');\n});\n\ntaskManager.on('workers-available', () => {\n  console.log('Workers are now available');\n});\n```\n\n### Statistics and Monitoring\n\n```typescript\n// Get comprehensive status including job recovery\nconst status = await getStatus();\nconsole.log('System status:', status);\n// Output includes: { queue: {...}, workers: {...}, jobRecovery: { isRunning: true, totalJobsWithAttempts: 0, ... } }\n\n// Get queue statistics\nconst queueStats = await getQueueStats();\nconsole.log('Queue stats:', queueStats);\n\n// Get detailed worker statistics with job distribution\nconst workerStats = await getWorkerStats();\nconsole.log('Worker stats:', workerStats);\n// Output: { total: 4, ready: 4, busy: 2, available: 2, distribution: [1250, 1245, 1255, 1250], totalJobsProcessed: 5000 }\n\n// Get job recovery statistics\nconst recoveryStats = await getJobRecoveryStats();\nconsole.log('Recovery stats:', recoveryStats);\n// Output: { isRunning: true, totalJobsWithAttempts: 2, totalRecoveredJobs: 5, totalFailedJobs: 1 }\n```\n\nThe monitoring system now includes:\n- **`distribution`**: Array showing jobs processed per worker\n- **`totalJobsProcessed`**: Total jobs completed across all workers\n- **`jobRecovery`**: Real-time job recovery statistics and status\n- **Real-time tracking** of worker utilization, load balancing, and fault tolerance\n\n### Batch Processing Configuration\n\n```typescript\nimport { setBatchConfig, getBatchConfig } from '@alcyone-labs/workalot';\n\n// Configure batch processing\nsetBatchConfig(20, true); // 20 jobs per batch, enabled\n\n// Get current configuration\nconst config = getBatchConfig();\nconsole.log(config); // { batchSize: 20, enabled: true }\n```\n\nBatch processing provides:\n- **Configurable batch sizes** (1-100 jobs per batch)\n- **Reduced communication overhead** between workers and orchestrator\n- **Maintained job isolation** and error handling\n- **Foundation for high throughput** processing\n\n### Architecture Strengths\n\n- **Optimized job distribution** - No O(n) bottlenecks for consistent performance\n- **Event-driven scheduling** - Immediate job processing without polling delays\n- **Backend-specific optimizations** - Each backend uses optimal data structures\n- **Real-time monitoring** - Precise worker utilization and job distribution tracking\n- **Race condition prevention** - Atomic job claiming and status management\n- **Intelligent throttling** - Prevents event loop overload while maximizing throughput\n- **Fault tolerance** - Automatic job recovery, worker crash detection, and stalled job handling\n- **Zero job loss** - Comprehensive recovery system ensures no jobs are lost due to failures\n\n### Performance Tips\n\n- **Use memory backend** for maximum throughput\n- **Scale worker count** with CPU cores for linear performance gains, keep at least 2 workers for the system\n- **Monitor worker utilization** - Target 90%+ for optimal efficiency\n- **Batch job submission** for reduced overhead on large workloads, tweak batch sizes for optimal results\n- **Configure job recovery** - Adjust check intervals and timeouts based on your job characteristics\n- **Monitor recovery stats** - Track stalled jobs and recovery performance for optimization\n\nSee `benchmarks/` directory for detailed performance testing.\n\n## Architecture\n\nWorkalot uses a multi-layered architecture designed for high performance, reliability, and fault tolerance:\n\n```\n┌─────────────────────────────────────────────────────────────┐\n│                     API Layer                               │\n│  ┌─────────────────┐  ┌─────────────────┐  ┌──────────────┐ │\n│  │ TaskManager     │  │ Singleton       │  │ Functions    │ │\n│  │ (Class)         │  │ Wrapper         │  │ (schedule)   │ │\n│  └─────────────────┘  └─────────────────┘  └──────────────┘ │\n└─────────────────────────────────────────────────────────────┘\n                                │\n┌─────────────────────────────────────────────────────────────┐\n│                  Worker System                              │\n│  ┌─────────────────┐  ┌─────────────────┐  ┌──────────────┐ │\n│  │ JobScheduler    │  │ WorkerManager   │  │ Worker       │ │\n│  │ (Coordination)  │  │ (Pool Mgmt)     │  │ Threads      │ │\n│  └─────────────────┘  └─────────────────┘  └──────────────┘ │\n└─────────────────────────────────────────────────────────────┘\n                                │\n┌─────────────────────────────────────────────────────────────┐\n│                Recovery System                              │\n│  ┌─────────────────┐  ┌─────────────────┐  ┌──────────────┐ │\n│  │ JobRecovery     │  │ Stalled Job     │  │ Crash        │ │\n│  │ Service         │  │ Detection       │  │ Detection    │ │\n│  └─────────────────┘  └─────────────────┘  └──────────────┘ │\n└─────────────────────────────────────────────────────────────┘\n                                │\n┌─────────────────────────────────────────────────────────────┐\n│                   Queue System                              │\n│  ┌─────────────────┐  ┌─────────────────┐  ┌──────────────┐ │\n│  │ QueueManager    │  │ IQueueBackend   │  │ JSON         │ │\n│  │ (In-memory)     │  │ (Interface)     │  │ Persistence  │ │\n│  └─────────────────┘  └─────────────────┘  └──────────────┘ │\n└─────────────────────────────────────────────────────────────┘\n                                │\n┌─────────────────────────────────────────────────────────────┐\n│                    Job System                               │\n│  ┌─────────────────┐  ┌─────────────────┐  ┌──────────────┐ │\n│  │ JobLoader       │  │ JobExecutor     │  │ JobRegistry  │ │\n│  │ (Dynamic Load)  │  │ (Execution)     │  │ (Discovery)  │ │\n│  └─────────────────┘  └─────────────────┘  └──────────────┘ │\n└─────────────────────────────────────────────────────────────┘\n```\n\n### Key Components\n\n- **API Layer**: Promise-based interface with singleton pattern for convenience\n- **Worker System**: Multi-threaded job processing with health monitoring and crash detection\n- **Recovery System**: Automatic job recovery, stalled job detection, and fault tolerance\n- **Queue System**: Pluggable backends (Memory, SQLite, PGLite, PostgreSQL) with recovery support\n- **Job System**: Dynamic job loading with TypeScript support and job scheduling from within jobs\n\n## Configuration\n\n### Environment Variables\n\n- `NODE_ENV=test` - Disables process signal handlers during testing\n- `VITEST=true` - Alternative test environment detection\n- `DATABASE_URL` - Automatically selects PostgreSQL backend when set\n\n### Backend Selection\n\nThe library automatically selects the appropriate backend:\n\n1. **Test Environment** (`NODE_ENV=test`): Uses in-memory PGLite\n2. **DATABASE_URL Set**: Uses PostgreSQL backend\n3. **Default**: Uses PGLite with file-based storage\n\n### Performance Tuning\n\n**Core Configuration**:\n- **maxThreads**: Set based on your CPU cores and workload (default: CPU cores - 2)\n  - **Recommendation**: Use all available cores for CPU-intensive jobs\n  - **Scaling**: Linear performance improvement up to 6+ cores verified\n- **maxInMemoryAge**: Adjust based on memory constraints (default: 24 hours)\n- **healthCheckInterval**: Lower for faster failure detection (default: 5000ms)\n- **jobTimeout**: Set appropriate timeouts for your jobs (default: 5000ms)\n\n**Backend Selection** (Performance Impact):\n- **`memory`**: **Fast** optional JSON persistence, good for high-throughput\n- **`pglite`**: **Moderate** (database-limited), PostgreSQL compatibility, good for development or single-VM execution\n- **`postgresql`**: **Full** features, good for persistence and complex queries\n\n**Optimization Tips**:\n- **Use memory backend** for maximum throughput\n- **Scale worker count** with CPU cores for linear performance gains\n- **Monitor worker utilization** - Target 90%+ efficiency\n- **Batch job submission** for reduced overhead on large workloads\n\n## Error Handling and Job Recovery\n\nThe library provides comprehensive error handling and automatic job recovery:\n\n### Error Handling\n\n```typescript\ntry {\n  const result = await scheduleAndWait({\n    jobFile: 'jobs/MyJob.ts',\n    jobPayload: { data: 'test' },\n    jobTimeout: 10000\n  });\n  console.log('Job completed:', result);\n} catch (error) {\n  if (error.message.includes('timed out')) {\n    console.log('Job timed out after 10 seconds');\n  } else if (error.message.includes('TaskManager is shutting down')) {\n    console.log('System is shutting down');\n  } else {\n    console.log('Job failed:', error.message);\n  }\n}\n```\n\n### Job Recovery System\n\nThe job recovery system automatically handles worker crashes and stalled jobs:\n\n```typescript\n// Configure job recovery (optional - enabled by default)\nawait initializeTaskManager({\n  backend: 'sqlite',\n  databaseUrl: 'memory://',\n  jobRecovery: {\n    enabled: true,                    // Enable/disable recovery\n    checkInterval: 60000,             // Check every minute\n    stalledTimeout: 300000,           // 5 minutes = stalled\n    maxRecoveryAttempts: 3            // Max retries before failing\n  }\n});\n\n// Monitor recovery events\ntaskManager.on('jobs-recovered', (data) => {\n  console.log(`Recovered ${data.count} stalled jobs`);\n});\n\n// Get recovery statistics\nconst recoveryStats = await getJobRecoveryStats();\nconsole.log('Recovery stats:', recoveryStats);\n```\n\n**What the recovery system handles:**\n- **Worker crashes** - Jobs are automatically recovered after timeout\n- **Worker stalls** - Background monitoring detects and recovers stuck jobs\n- **JobScheduler crashes** - Recovery service runs independently\n- **Database locks** - Each backend uses optimal locking strategy\n- **Infinite retries** - Max recovery attempts prevent endless loops\n\n**Recovery guarantees:**\n- **Zero job loss** - All jobs are either completed, failed with proper error handling, or automatically recovered\n- **Configurable timeouts** - Adjust based on your job characteristics\n- **Performance monitoring** - Track recovery performance and job failure patterns\n- **Graceful degradation** - System continues operating even during recovery operations\n\n## Understanding `scheduleAndWait()` Behavior\n\n### Queue Order and Execution\n\n`scheduleAndWait()` is **not** a priority queue - it follows strict FIFO (first-in-first-out) order:\n\n```typescript\n// If you have existing jobs in the queue:\nawait schedule({ jobFile: 'jobs/Job1.ts' }); // Position 1\nawait schedule({ jobFile: 'jobs/Job2.ts' }); // Position 2\nawait schedule({ jobFile: 'jobs/Job3.ts' }); // Position 3\n\n// Your scheduleAndWait() job goes to the END of the queue:\nconst result = await scheduleAndWait({\n  jobFile: 'jobs/UrgentJob.ts'  // Position 4 - waits for Jobs 1, 2, 3\n});\n```\n\n### What \"Now\" Means\n\n- **\"Now\"** = Wait for completion of your specific job\n- **NOT** = Execute immediately or skip the queue\n- **NOT** = Higher priority than other jobs\n\n### Practical Example\n\n```typescript\n// Scenario: 1000 jobs already queued\nconsole.log('Scheduling urgent job...');\nconst startTime = Date.now();\n\nconst result = await scheduleAndWait({\n  jobFile: 'jobs/UrgentJob.ts',\n  jobPayload: { urgent: true }\n});\n\nconst waitTime = Date.now() - startTime;\nconsole.log(`Job completed after ${waitTime}ms`);\n// Will be ~1000 jobs × average job time + your job execution time\n```\n\n### When to Use `scheduleAndWait()`\n\n**Good use cases:**\n- Synchronous workflow where next step depends on job result\n- Error handling where you need to know if the job succeeded\n\n**Consider alternatives:**\n- Fire-and-forget jobs → Use `schedule()`\n- Batch processing → Use `schedule()` + `whenFree()`\n\n## Good Practices\n\n1. **Job Design**\n   - Keep jobs focused and single-purpose\n   - Validate input data early\n   - Use appropriate timeouts\n   - Handle errors gracefully\n\n2. **Performance**\n   - Configure worker count based on workload\n   - Monitor queue statistics\n   - Use appropriate job timeouts\n   - Clean up old completed jobs\n\n3. **Reliability**\n   - Always call `shutdown()` for graceful cleanup\n   - Handle job failures appropriately\n   - Monitor worker health and job recovery statistics\n   - Use persistence for important queues\n   - Configure job recovery timeouts based on your job characteristics\n   - Monitor recovery events for system health insights\n\n4. **Testing**\n   - Test jobs in isolation\n   - Use in-memory backends for tests (`backend: 'sqlite', databaseUrl: 'memory://'`)\n   - Clean up resources in test teardown\n   - Use unique database URLs for parallel tests\n\n5. **BunJS Usage**\n   - Run TypeScript files directly with `bun run`\n   - Use absolute paths for job files when running with Bun\n   - Leverage Bun's native APIs when available\n\n## Benchmarking\n\nThe library includes a comprehensive benchmark suite for performance testing and validation:\n\n```bash\n# Quick performance test (recommended)\npnpm run benchmark:easy\n\n# Run all benchmarks with system info display\nbun run benchmarks/run-benchmarks.ts\n\n# Test scaling across different core counts\nbun run benchmarks/run-benchmarks.ts --difficulty easy --configs 2-cores-10k-jobs-memory,4-cores-10k-jobs-memory,6-cores-10k-jobs-memory\n\n# Run specific configurations\nbun run benchmarks/run-benchmarks.ts --configs 2-cores-10k-jobs,4-cores-10k-jobs\n```\n\n**Features:**\n- **Friendly CLI**: Clean progress bars, system information display\n- **Difficulty Scaling**: Adjustable CPU workload (easy, normal, hard, extreme)\n- **File Logging**: Detailed timestamped logs with structured JSON data\n- **Performance Metrics**: Throughput, latency, resource utilization, scaling efficiency\n- **Multiple Backends**: Test performance across memory, SQLite, PGLite, and PostgreSQL\n- **Visual Analysis**: Charts to help with analysis of linear scaling and bottlenecks (benchmarks/visualize-results.html)\n\n**Recent Benchmark Results** (M2 Max, 12 cores, easy difficulty):\n- **Memory Backend**: 4,366 jobs/sec (6 cores), nearly linear scaling\n- **Execution Speed**: Sub-second completion for 1,000 job batches in memory\n- **Job Recovery**: < 100ms recovery time per stalled job with zero performance impact\n- **Fault Tolerance**: 100% job recovery rate in crash simulation tests\n\nSee `benchmarks/README.md` for detailed documentation and `ROADMAP.md` for optimization targets.\n\n## Examples & Getting Started\n\n### Quick Start Example\n\n```typescript\nimport { TaskManager } from '@alcyone-labs/workalot';\n\n// Create a simple job\n// File: jobs/HelloJob.ts\nimport { BaseJob, JobExecutionContext } from '@alcyone-labs/workalot';\n\nexport class HelloJob extends BaseJob {\n  async run(payload: { name: string }, context: JobExecutionContext) {\n    this.validatePayload(payload, ['name']);\n\n    // Simulate some work\n    await new Promise(resolve => setTimeout(resolve, 100));\n\n    // Schedule a follow-up welcome email\n    if (payload.name !== 'Test') {\n      context.schedule({\n        jobFile: 'jobs/WelcomeEmailJob.ts',\n        jobPayload: {\n          recipientName: payload.name,\n          timestamp: new Date().toISOString()\n        }\n      });\n    }\n\n    return this.createSuccessResult({\n      message: `Hello, ${payload.name}!`,\n      timestamp: new Date().toISOString()\n    });\n  }\n}\n\n// Use the job\nconst taskManager = new TaskManager({\n  backend: 'memory',\n  maxThreads: 4\n});\n\nawait taskManager.initialize();\n\n// Schedule and wait for job completion\nconst result = await taskManager.scheduleAndWait({\n  jobFile: 'jobs/HelloJob.ts',\n  jobPayload: { name: 'World' }\n});\n\nconsole.log(result.results.message); // \"Hello, World!\"\n\nawait taskManager.shutdown();\n```\n\n### Complete Examples\n\nSee the `examples/` directory for comprehensive examples:\n\n- **`examples/quick-start.ts`** - Quick start guide with multiple examples\n- **`examples/pglite-inmemory.ts`** - PGLite in-memory features and performance comparison\n- **`examples/basic-usage.ts`** - Simple API usage and job creation\n- **`examples/sample-consumer/`** - Complete application with multiple job types\n- **`examples/performance-test.ts`** - Performance testing and backend benchmarking\n- **`examples/error-handling.ts`** - Comprehensive error handling patterns\n\n#### Try the PGLite In-Memory Example\n\n```bash\n# Run the dedicated PGLite in-memory example\nbun run examples/pglite-inmemory.ts\n\n# Compare performance across backends\nbun run examples/performance-test.ts --backends\n```\n\n### Common Use Cases\n\n#### High-Throughput Data Processing\n```typescript\n// Process large datasets efficiently\nfor (let i = 0; i < 10000; i++) {\n  taskManager.schedule({\n    jobFile: 'jobs/DataProcessorJob.ts',\n    jobPayload: { batchId: i, data: largeBatch[i] }\n  });\n}\n\n// Wait for all jobs to complete\nawait taskManager.whenIdle();\n```\n\n#### Background Task Processing\n```typescript\n// Fire-and-forget background tasks\nconst jobId = await taskManager.schedule({\n  jobFile: 'jobs/EmailSenderJob.ts',\n  jobPayload: {\n    to: 'user@example.com',\n    template: 'welcome',\n    data: { username: 'john' }\n  }\n});\n\nconsole.log(`Email job scheduled: ${jobId}`);\n```\n\n#### Real-time Monitoring\n```typescript\n// Monitor system performance and job recovery\ntaskManager.on('job-completed', (jobId, result) => {\n  console.log(`Job ${jobId} completed in ${result.executionTime}ms`);\n});\n\ntaskManager.on('jobs-recovered', (data) => {\n  console.log(`Recovered ${data.count} stalled jobs`);\n});\n\nsetInterval(async () => {\n  const stats = await taskManager.getStatus();\n  console.log(`Queue: ${stats.queue.pending} pending, ${stats.queue.processing} processing`);\n  console.log(`Recovery: ${stats.jobRecovery.totalJobsWithAttempts} jobs with recovery attempts`);\n}, 1000);\n```\n\n## Dependencies\n\n### Core Dependencies\n- **ulidx**: ULID generation for unique job IDs\n\n### Optional Dependencies\n- **@electric-sql/pglite**: PostgreSQL-compatible embedded database (for PGLite backend)\n- **better-sqlite3**: High-performance SQLite driver (for SQLite backend on Node.js)\n\n**Runtime-Specific SQLite Support:**\n- **Bun runtime**: Uses built-in `bun:sqlite` (no additional installation required)\n- **Node.js runtime**: Requires `better-sqlite3` package\n\nThe library automatically detects the runtime and uses the optimal SQLite driver. Install only the dependencies you need for your chosen backend and runtime.\n\n## Requirements\n\n- **Node.js**: >= 18.0.0\n- **BunJS**: >= 1.2.0 (optional, for native TypeScript support)\n- **Deno**: >= 2\n\n## License\n\nMIT\n\n## Contributing\n\n1. Fork the repository\n2. Create a feature branch\n3. Add tests for new functionality\n4. Ensure all tests pass with `pnpm test`\n5. Follow project conventions (no emojis, use PNPM)\n6. Submit a pull request\n\n## Roadmap\n\n### Current Status: V1 Ready\n\nWorkalot has achieved good performance with significant improvement and linear scaling. The library includes comprehensive features including automatic job recovery, fault tolerance, and zero job loss guarantees.\n\n### Upcoming Features\n\n**Database Backend Enhancement**\n- PGLite backend optimization for 1,000+ jobs/sec with database persistence\n- PostgreSQL production features with connection pooling and batch operations\n- Database-specific optimizations for high-throughput scenarios\n- Multi-tenant support with isolated job queues per namespace\n\n**Advanced Queue Management**\n- Dynamic queue sizing based on job complexity and worker performance\n- Auto-scaling worker queues during peak load periods\n- Enhanced monitoring with real-time queue depth and job latency tracking\n- Advanced error handling with automatic retry and dead letter queues\n\n**Performance Optimizations**\n- Object pooling for job context containers to reduce GC pressure\n- Binary serialization to replace JSON for faster worker communication\n- Memory-mapped job queues for low latency job distribution\n- Dynamic worker scaling based on system load and queue depth\n\n**Advanced Features**\n- Horizontal scaling with multi-node orchestrator and job distribution\n- Advanced analytics with performance dashboards and monitoring\n- Job prioritization with database-level priority queues\n- Comprehensive backup and recovery procedures for production deployments\n","readmeFilename":"README.md","_rev":"1-8c167edd3a73e665de8e026965e37715"}