import { sql, type SQLWrapper } from 'drizzle-orm'
import type { Schema } from '@platform-modules/db'
import { asBigint } from './internal/bigint.js'
import { applyDelta, lockCounter, type CounterKey } from './internal/counter.js'
import {
  isMeteringError,
  MeteringStorageError,
  MeteringValidationError,
} from './errors.js'
import type { MeteringDatabase } from './types.js'

export type SweepDeps<S extends Schema = Record<string, never>> = {
  db: MeteringDatabase<S>
}

type SqlRow = Record<string, unknown>
type SqlExecutor = { execute(query: SQLWrapper): Promise<unknown> }

function rows(result: unknown): SqlRow[] {
  return (Array.isArray(result) ? result : (result as { rows?: SqlRow[] }).rows) ?? []
}

function firstRow(result: unknown): SqlRow | undefined {
  return rows(result)[0]
}


function keyFromRow(row: SqlRow): CounterKey {
  return {
    tenantId: String(row.tenant_id),
    account: String(row.account),
    meter: String(row.meter),
    periodId: String(row.period_id),
  }
}

function validateInput(before: Date, limit: number): void {
  if (!(before instanceof Date) || Number.isNaN(before.getTime())) {
    throw new MeteringValidationError('before', 'must be a valid Date')
  }
  if (!Number.isSafeInteger(limit) || limit < 0) {
    throw new MeteringValidationError('limit', 'must be a non-negative safe integer')
  }
}

async function expireOne(tx: SqlExecutor, before: Date): Promise<boolean | null> {
  const beforeIso = before.toISOString()
  const candidate = firstRow(await tx.execute(sql`
    SELECT r.id, r.tenant_id, r.account, r.meter, r.period_id
    FROM metering_reservation AS r
    INNER JOIN metering_counter AS c
      ON c.tenant_id = r.tenant_id
      AND c.account = r.account
      AND c.meter = r.meter
      AND c.period_id = r.period_id
    WHERE r.status = 'reserved'
      AND r.expires_at < ${beforeIso}::timestamptz
    ORDER BY r.expires_at, r.id
    LIMIT 1
    FOR UPDATE OF c SKIP LOCKED
  `))
  if (!candidate) return null

  const key = keyFromRow(candidate)
  await lockCounter(tx, key)

  const reservation = firstRow(await tx.execute(sql`
    SELECT reserved_units
    FROM metering_reservation
    WHERE id = ${String(candidate.id)}
      AND status = 'reserved'
    FOR UPDATE
  `))
  if (!reservation) return false

  const reservedUnits = asBigint(reservation.reserved_units, 'reservation', 'reserved_units')
  const updated = firstRow(await tx.execute(sql`
    UPDATE metering_reservation
    SET status = 'expired'
    WHERE id = ${String(candidate.id)}
      AND status = 'reserved'
    RETURNING id
  `))
  if (!updated) return false

  await applyDelta(tx, key, { reserved: -reservedUnits })
  return true
}

export async function expireStale<S extends Schema = Record<string, never>>(
  deps: SweepDeps<S>,
  before: Date,
  limit: number,
): Promise<number> {
  validateInput(before, limit)
  if (limit === 0) return 0

  try {
    let expired = 0
    while (expired < limit) {
      const result = await deps.db.transaction((tx) => expireOne(tx, before))
      if (result === null) break
      if (result) expired += 1
    }
    return expired
  } catch (error) {
    if (isMeteringError(error)) throw error
    throw new MeteringStorageError('expire stale reservations', error)
  }
}
