import { z } from "zod"; /** * Supported queue backend providers. * * - `qstash` — Upstash QStash (serverless, HTTP-based, zero infra) * - `bullmq` — BullMQ over self-hosted Redis (full-featured, self-managed) * - `sqs` — AWS SQS (managed, pull-based, ideal for AWS-native deployments) * - `memory` — In-memory queue for local dev & testing (NOT for production) */ export type QueueProviderType = "qstash" | "bullmq" | "sqs" | "memory"; export declare const JobOptionsSchema: z.ZodObject<{ idempotencyKey: z.ZodOptional; delaySec: z.ZodOptional; maxRetries: z.ZodOptional; priority: z.ZodOptional; cron: z.ZodOptional; tenantId: z.ZodOptional; metadata: z.ZodOptional>; }, z.core.$strip>; export type JobOptions = z.infer; export declare const JobPayloadSchema: z.ZodObject<{ id: z.ZodString; queue: z.ZodString; type: z.ZodString; data: z.ZodRecord; options: z.ZodOptional; delaySec: z.ZodOptional; maxRetries: z.ZodOptional; priority: z.ZodOptional; cron: z.ZodOptional; tenantId: z.ZodOptional; metadata: z.ZodOptional>; }, z.core.$strip>>; createdAt: z.ZodString; }, z.core.$strip>; export type JobPayload = z.infer; export interface JobResult { /** Provider-specific job/message ID */ jobId: string; /** Acknowledged by the provider? */ accepted: boolean; /** Provider name that handled the enqueue */ provider: QueueProviderType; } export type JobStatus = "pending" | "active" | "completed" | "failed" | "dead-lettered" | "delayed" | "waiting" | "canceled"; export type JobLifecycleAction = "enqueued" | "started" | "attempt-failed" | "completed" | "failed" | "dead-lettered" | "canceled" | "retried"; export interface JobLifecycleEvent { id: string; jobId: string; queue: string; action: JobLifecycleAction; at: string; attempt?: number; reason?: string; actor?: string; metadata?: Record; } export interface JobLifecycleActionResult { jobId: string; queue: string; accepted: boolean; action: "cancel" | "retry"; reason?: string; } export interface JobStatusInfo { id: string; queue: string; status: JobStatus; attempts: number; maxRetries: number; failedReason?: string; canceledReason?: string; createdAt: string; processedAt?: string; completedAt?: string; deadLetteredAt?: string; canceledAt?: string; } export interface DeadLetterJob { id: string; queue: string; type: string; originalJob: JobPayload; attempts: number; maxRetries: number; failedReason: string; provider: QueueProviderType; failedAt: string; } export type QueueDeadLetterRecord = Record; export interface QueueDeadLetterFetchContext { token: string; endpoint?: string; queue?: string; } export type QueueDeadLetterFetcher = (context: QueueDeadLetterFetchContext) => Promise; /** * A job handler receives the payload and returns void on success * or throws to trigger a retry. */ export type JobHandler = Record> = (job: JobPayload & { data: T; }) => Promise; /** * Every queue backend must implement this interface. * The factory function (`createQueue`) returns a `QueueProvider`. */ export interface QueueProvider { readonly name: QueueProviderType; /** * Enqueue a job for async processing. */ enqueue(job: JobPayload): Promise; /** * Enqueue multiple jobs in a single round-trip (where supported). * Falls back to sequential enqueue if the provider has no native batch API. */ enqueueBatch(jobs: JobPayload[]): Promise; /** * Register a handler for a given job type. * - BullMQ: starts a Worker that polls Redis. * - QStash: returns an HTTP handler to mount in your API route. * - Memory: processes inline. */ registerHandler>(queue: string, type: string, handler: JobHandler): void; /** * Query the status of a previously enqueued job (best-effort). * Not all providers support this; returns `undefined` if unsupported. */ getJobStatus?(jobId: string, queue: string): Promise; /** * Inspect jobs that exhausted retries and require operator attention. * Providers without a durable DLQ may omit this method. */ getDeadLetteredJobs?(queue?: string): Promise; /** * Return provider-neutral lifecycle logs for operator tooling. * Providers without durable inspection may omit this method. */ getJobLogs?(jobId: string, queue: string): Promise; /** * Cancel a pending/delayed job when supported. */ cancelJob?(jobId: string, queue: string, reason?: string): Promise; /** * Replay a failed or dead-lettered job when supported. */ retryJob?(jobId: string, queue: string, reason?: string): Promise; /** * Graceful shutdown — drain in-flight jobs, close connections. */ close(): Promise; } export interface QStashProviderConfig { provider: "qstash"; /** QStash REST token (defaults to `process.env.QSTASH_TOKEN`) */ token?: string; /** * Base URL of your API that QStash will POST callbacks to. * e.g. "https://api.nebutra.com" or "https://my-tunnel.ngrok.io" */ callbackBaseUrl: string; /** Optional signing keys for webhook verification */ currentSigningKey?: string; nextSigningKey?: string; /** Optional provider-side DLQ endpoint for injected fetchers. */ dlqEndpoint?: string; /** Injectable provider-side DLQ fetcher. No Upstash SDK DLQ APIs are assumed. */ dlqFetcher?: QueueDeadLetterFetcher; } export interface BullMQProviderConfig { provider: "bullmq"; /** Redis connection URL (defaults to `process.env.REDIS_URL`) */ redisUrl?: string; /** Default concurrency per worker (default: 5) */ concurrency?: number; /** Key prefix for Redis keys (default: "nebutra:queue") */ prefix?: string; } export interface MemoryProviderConfig { provider: "memory"; } export interface SQSProviderConfig { provider: "sqs"; /** AWS region (defaults to `process.env.AWS_REGION`) */ region?: string; /** SQS queue URL (defaults to `process.env.AWS_SQS_QUEUE_URL`) */ queueUrl?: string; /** AWS access key (defaults to `process.env.AWS_ACCESS_KEY_ID`) */ accessKeyId?: string; /** AWS secret key (defaults to `process.env.AWS_SECRET_ACCESS_KEY`) */ secretAccessKey?: string; /** Long-poll wait time in seconds when receiving (1-20, default 20) */ waitTimeSeconds?: number; /** Max messages per receive call (1-10, default 10) */ maxMessages?: number; /** Visibility timeout in seconds — how long a received message is hidden (default 30) */ visibilityTimeoutSeconds?: number; } export type QueueConfig = QStashProviderConfig | BullMQProviderConfig | SQSProviderConfig | MemoryProviderConfig; //# sourceMappingURL=types.d.ts.map