import { $hook, $inject, Alepha, AlephaError } from "alepha";
import { $logger } from "alepha/logger";
import { QueueProvider } from "./QueueProvider.ts";

// ---------------------------------------------------------------------------------------------------------------------

/**
 * Cloudflare Queue interface matching the CF Workers Queue API.
 */
export interface CloudflareQueue {
  send(message: unknown): Promise<void>;
  sendBatch(messages: Array<{ body: unknown }>): Promise<void>;
}

// ---------------------------------------------------------------------------------------------------------------------

/**
 * Default queue binding name used in wrangler configuration.
 */
export const QUEUE_DEFAULT_BINDING = "JOBS_QUEUE";

/**
 * How many times Cloudflare redelivers a message before routing it to the
 * consumer's dead-letter queue. Matches Cloudflare's own default.
 */
export const QUEUE_DEFAULT_MAX_RETRIES = 3;

// ---------------------------------------------------------------------------------------------------------------------

/**
 * Cloudflare Queue provider.
 *
 * Uses a Queue binding for message dispatch. Messages are wrapped with the
 * logical queue name so the consumer can route them to the correct handler.
 *
 * **Required Cloudflare binding:**
 * - `JOBS_QUEUE` - A Queue binding in wrangler configuration
 *
 * @example
 * ```toml
 * # wrangler.toml - automatically generated by alepha build
 * [[queues.producers]]
 * binding = "JOBS_QUEUE"
 * queue = "my-app-queue"
 * ```
 */
export class CloudflareQueueProvider extends QueueProvider {
  protected readonly alepha = $inject(Alepha);
  protected readonly log = $logger();

  protected queue?: CloudflareQueue;

  protected readonly onStart = $hook({
    on: "start",
    handler: async () => {
      const cloudflareEnv = this.alepha.store.get("cloudflare.env") as
        | Record<string, unknown>
        | undefined;
      if (!cloudflareEnv) {
        throw new AlephaError(
          "Cloudflare Workers environment not found in Alepha store under 'cloudflare.env'.",
        );
      }

      const binding = cloudflareEnv[QUEUE_DEFAULT_BINDING] as
        | CloudflareQueue
        | undefined;
      if (!binding) {
        throw new AlephaError(
          `Queue binding '${QUEUE_DEFAULT_BINDING}' not found in Cloudflare Workers environment.`,
        );
      }

      this.queue = binding;
      this.log.info("Cloudflare Queue ready");
    },
  });

  public async push(queue: string, message: string): Promise<void> {
    await this.getQueue().send({ queue, message });
  }

  /**
   * Not used on Cloudflare — queue consumption is push-based via the `queue` handler.
   */
  public async pop(_queue: string): Promise<string | undefined> {
    return undefined;
  }

  protected getQueue(): CloudflareQueue {
    if (!this.queue) {
      throw new AlephaError(
        "Queue binding not initialized. Call start() first.",
      );
    }
    return this.queue;
  }
}
