import type { StandardSchemaV1 } from '@standard-schema/spec'

/** Type-discriminated job envelope — dispatch keys on `type`. */
export type JobEnvelope<T = unknown> = {
  type: string
  payload: T
}

/**
 * Structurally matches `@cloudflare/workers-types` `Message` — declared locally so
 * core needs no CF types at the `.` entry.
 */
export type QueueMessage<T = unknown> = {
  readonly id: string
  readonly timestamp: Date
  readonly body: T
  readonly attempts: number
  ack(): void
  retry(options?: { delaySeconds?: number }): void
}

/**
 * Structurally matches `@cloudflare/workers-types` `MessageBatch` — declared locally.
 */
export type QueueBatch<T = unknown> = {
  readonly queue: string
  readonly messages: readonly QueueMessage<T>[]
  retryAll(options?: { delaySeconds?: number }): void
  ackAll(): void
}

/** Marker for domain-terminal failures — ack after writing FAILED state, never retry. */
export class TerminalJobError extends Error {
  override readonly name = 'TerminalJobError'
}

/** Injected idempotency strategy — typical impls use KV, table, unique-index, processedAt, status-claim. */
export interface IdempotencyStore {
  seen(key: string): Promise<boolean>
  mark(key: string, ttl?: number): Promise<void>
}

export type DispatchOpts<E> = {
  /** Backstop path: missing-secret / unknown-type → throw. Live path: log + ack. */
  strict?: boolean
  idempotency?: { key: string; store: IdempotencyStore }
  /** Fired before ack on `TerminalJobError` — write FAILED state here. */
  onTerminalFailure?: (
    err: TerminalJobError,
    env: E,
    msg: JobEnvelope,
  ) => void | Promise<void>
}

type HandlerEntry<E> = {
  schema: StandardSchemaV1
  handle: (env: E, payload: unknown, msg: JobEnvelope) => Promise<void>
}

export type JobRegistry<E> = {
  register<T>(
    type: string,
    schema: StandardSchemaV1<unknown, T>,
    handle: (env: E, payload: T, msg: JobEnvelope<T>) => Promise<void>,
  ): void
  dispatch(env: E, msg: JobEnvelope, opts?: DispatchOpts<E>): Promise<void>
}

export function createJobRegistry<E>(): JobRegistry<E> {
  const handlers = new Map<string, HandlerEntry<E>>()

  return {
    register(type, schema, handle) {
      handlers.set(type, {
        schema,
        handle: handle as HandlerEntry<E>['handle'],
      })
    },

    async dispatch(env, msg, opts) {
      const entry = handlers.get(msg.type)
      if (!entry) {
        const detail = `unknown job type: ${msg.type}`
        if (opts?.strict) throw new Error(detail)
        console.warn(detail)
        return
      }

      if (opts?.idempotency) {
        const { key, store } = opts.idempotency
        if (await store.seen(key)) return
      }

      const validated = await entry.schema['~standard'].validate(msg.payload)
      if ('issues' in validated) {
        console.warn(`malformed job payload for type ${msg.type}`, validated.issues)
        return
      }

      try {
        await entry.handle(env, validated.value, msg)
        if (opts?.idempotency) {
          await opts.idempotency.store.mark(opts.idempotency.key)
        }
      } catch (err) {
        if (err instanceof TerminalJobError) {
          await opts?.onTerminalFailure?.(err, env, msg)
          return
        }
        throw err
      }
    },
  }
}
