/** * Workflow Scheduler * * Features: * - Priority queue for pending runs * - Cron-based scheduling * - Delayed execution * - Concurrency control */ import type { Workflow, WorkflowState, ScheduleOptions, RunStore } from '@cogitator-ai/types'; /** * Priority queue item */ export interface QueueItem { runId: string; workflowName: string; priority: number; scheduledFor: number; } /** * Scheduler configuration */ export interface SchedulerConfig { runStore: RunStore; maxConcurrency?: number; pollInterval?: number; onRunReady?: (runId: string) => void; } /** * Priority queue implementation using a heap */ export declare class PriorityQueue { private items; /** * Add item to queue (sorted by scheduledFor then priority) */ enqueue(item: QueueItem): void; /** * Remove and return highest priority item */ dequeue(): QueueItem | undefined; /** * Peek at highest priority item without removing */ peek(): QueueItem | undefined; /** * Remove item by runId */ remove(runId: string): boolean; /** * Get items ready to execute (scheduledFor <= now) */ getReady(now?: number): QueueItem[]; /** * Get queue size */ size(): number; /** * Clear the queue */ clear(): void; /** * Get all items (for inspection) */ getAll(): QueueItem[]; private bubbleUp; private bubbleDown; private compare; } /** * Cron-based scheduled job */ export interface CronJob { id: string; workflowName: string; workflow: Workflow; expression: string; timezone?: string; options?: Omit; nextRun: number; enabled: boolean; } /** * Job scheduler for workflow runs */ export declare class JobScheduler { private queue; private cronJobs; private runStore; private maxConcurrency; private pollInterval; private onRunReady?; private runningCount; private pollTimer?; private disposed; constructor(config: SchedulerConfig); /** * Start the scheduler */ start(): void; /** * Stop the scheduler */ stop(): void; /** * Schedule a workflow run */ scheduleRun(workflow: Workflow, options?: ScheduleOptions, existingRunId?: string): Promise; /** * Register a cron-based recurring job */ registerCronJob(workflow: Workflow, expression: string, options?: { id?: string; timezone?: string; jobOptions?: Omit; }): string; /** * Unregister a cron job */ unregisterCronJob(jobId: string): boolean; /** * Enable/disable a cron job */ setCronJobEnabled(jobId: string, enabled: boolean): boolean; /** * Cancel a scheduled run */ cancelRun(runId: string, reason?: string): Promise; /** * Mark run as started */ runStarted(_runId: string): void; /** * Mark run as completed */ runCompleted(_runId: string): void; /** * Get queue size */ getQueueSize(): number; /** * Get running count */ getRunningCount(): number; /** * Get all cron jobs */ getCronJobs(): CronJob[]; /** * Dispose the scheduler */ dispose(): void; private tick; } /** * Create a job scheduler */ export declare function createJobScheduler(config: SchedulerConfig): JobScheduler; //# sourceMappingURL=scheduler.d.ts.map