import * as _alepha_core2 from "alepha";
import * as _alepha_core3 from "alepha";
import * as _alepha_core1 from "alepha";
import * as _alepha_core0 from "alepha";
import { Alepha, Descriptor, KIND, Service, Static, TSchema } from "alepha";
import { DateTimeProvider } from "alepha/datetime";

//#region src/providers/QueueProvider.d.ts
/**
 * Minimalist Queue interface.
 *
 * Will be probably enhanced in the future to support more advanced features. But for now, it's enough!
 */
declare abstract class QueueProvider {
  /**
   * Push a message to the queue.
   *
   * @param queue Name of the queue to push the message to.
   * @param message String message to be pushed to the queue. Buffer messages are not supported for now.
   */
  abstract push(queue: string, message: string): Promise<void>;
  /**
   * Pop a message from the queue.
   *
   * @param queue Name of the queue to pop the message from.
   *
   * @returns The message popped or `undefined` if the queue is empty.
   */
  abstract pop(queue: string): Promise<string | undefined>;
}
//# sourceMappingURL=QueueProvider.d.ts.map
//#endregion
//#region src/providers/MemoryQueueProvider.d.ts
declare class MemoryQueueProvider implements QueueProvider {
  protected readonly log: _alepha_core2.Logger;
  protected queueList: Record<string, string[]>;
  push(queue: string, ...messages: string[]): Promise<void>;
  pop(queue: string): Promise<string | undefined>;
}
//# sourceMappingURL=MemoryQueueProvider.d.ts.map
//#endregion
//#region src/providers/WorkerProvider.d.ts
declare const envSchema: _alepha_core3.TObject<{
  /**
   * The interval in milliseconds to wait before checking for new messages.
   */
  QUEUE_WORKER_INTERVAL: _alepha_core3.TNumber;
  /**
   * The maximum interval in milliseconds to wait before checking for new messages.
   */
  QUEUE_WORKER_MAX_INTERVAL: _alepha_core3.TNumber;
  /**
   * The number of workers to run concurrently. Defaults to 1.
   * Useful only if you are doing a lot of I/O.
   */
  QUEUE_WORKER_CONCURRENCY: _alepha_core3.TNumber;
}>;
declare module "alepha" {
  interface Env extends Partial<Static<typeof envSchema>> {}
}
declare class WorkerProvider {
  protected readonly log: _alepha_core3.Logger;
  protected readonly env: {
    QUEUE_WORKER_INTERVAL: number;
    QUEUE_WORKER_MAX_INTERVAL: number;
    QUEUE_WORKER_CONCURRENCY: number;
  };
  protected readonly alepha: Alepha;
  protected readonly queueProvider: QueueProvider;
  protected readonly dateTimeProvider: DateTimeProvider;
  protected workerPromises: Array<Promise<void>>;
  protected isWorkersRunning: boolean;
  protected abortController: AbortController;
  protected workerIntervals: Record<number, number>;
  protected consumers: Array<Consumer>;
  protected readonly start: _alepha_core3.HookDescriptor<"start">;
  /**
   * Start the workers.
   * This method will create an endless loop that will check for new messages!
   */
  protected startWorkers(): void;
  protected readonly stop: _alepha_core3.HookDescriptor<"stop">;
  /**
   * Wait for the next message, where `n` is the worker number.
   *
   * This method will wait for a certain amount of time, increasing the wait time again if no message is found.
   */
  protected waitForNextMessage(n: number): Promise<void>;
  /**
   * Get the next message.
   */
  protected getNextMessage(): Promise<undefined | NextMessage>;
  /**
   * Process a message from a queue.
   */
  protected processMessage(response: {
    message: any;
    consumer: Consumer;
  }): Promise<void>;
  /**
   * Stop the workers.
   *
   * This method will stop the workers and wait for them to finish processing.
   */
  protected stopWorkers(): Promise<void>;
  /**
   * Force the workers to get back to work. zug zug!
   */
  wakeUp(): void;
}
interface Consumer<T extends TSchema = TSchema> {
  queue: QueueDescriptor<T>;
  handler: (message: QueueMessage<T>) => Promise<void>;
}
interface NextMessage {
  consumer: Consumer;
  message: string;
}
//#endregion
//#region src/descriptors/$queue.d.ts
/**
 * Create a new queue.
 */
declare const $queue: {
  <T extends TSchema>(options: QueueDescriptorOptions<T>): QueueDescriptor<T>;
  [KIND]: typeof QueueDescriptor;
};
interface QueueDescriptorOptions<T extends TSchema> {
  name?: string;
  description?: string;
  provider?: "memory" | Service<QueueProvider>;
  schema: T;
  handler?: (message: QueueMessage<T>) => Promise<void>;
}
declare class QueueDescriptor<T extends TSchema> extends Descriptor<QueueDescriptorOptions<T>> {
  protected readonly log: _alepha_core1.Logger;
  protected readonly workerProvider: WorkerProvider;
  readonly provider: QueueProvider | MemoryQueueProvider;
  push(...payloads: Array<Static<T>>): Promise<void>;
  get name(): string;
  protected $provider(): QueueProvider | MemoryQueueProvider;
}
interface QueueMessageSchema {
  payload: TSchema;
}
interface QueueMessage<T extends TSchema> {
  payload: Static<T>;
}
//# sourceMappingURL=$queue.d.ts.map
//#endregion
//#region src/descriptors/$consumer.d.ts
/**
 * Consumer descriptor.
 */
declare const $consumer: {
  <T extends TSchema>(options: ConsumerDescriptorOptions<T>): ConsumerDescriptor<T>;
  [KIND]: typeof ConsumerDescriptor;
};
interface ConsumerDescriptorOptions<T extends TSchema> {
  queue: QueueDescriptor<T>;
  handler: (message: {
    payload: Static<T["payload"]>;
  }) => Promise<void>;
}
declare class ConsumerDescriptor<T extends TSchema> extends Descriptor<ConsumerDescriptorOptions<T>> {}
//# sourceMappingURL=$consumer.d.ts.map
//#endregion
//#region src/index.d.ts
/**
 * Provides asynchronous message queuing and processing capabilities through declarative queue descriptors.
 *
 * The queue module enables reliable background job processing and message passing using the `$queue` descriptor
 * on class properties. It supports schema validation, automatic retries, and multiple queue backends for
 * building scalable, decoupled applications with robust error handling.
 *
 * @see {@link $queue}
 * @see {@link $consumer}
 * @module alepha.queue
 */
declare const AlephaQueue: _alepha_core0.Service<_alepha_core0.Module>;
//# sourceMappingURL=index.d.ts.map

//#endregion
export { $consumer, $queue, AlephaQueue, ConsumerDescriptor, ConsumerDescriptorOptions, MemoryQueueProvider, QueueDescriptor, QueueDescriptorOptions, QueueMessage, QueueMessageSchema, QueueProvider };
//# sourceMappingURL=index.d.ts.map