import { executeRows } from '../execute-rows.js';
import { and, eq, isNull, sql } from 'drizzle-orm';
import type { DrizzleClient } from '@/server/db/client.js';
import { carrierWebhookEvents } from '@/server/db/schema.js';

export async function claimCarrierWebhookEvent(
  db: DrizzleClient,
  { carrier, eventId }: { carrier: string; eventId: string },
): Promise<{ claimed: boolean; duplicate: boolean }> {
  const result = await db.execute(sql`
    INSERT INTO carrier_webhook_events (carrier, event_id, claimed_at)
    VALUES (${carrier}, ${eventId}, now())
    ON CONFLICT (carrier, event_id) DO UPDATE SET claimed_at = now()
    WHERE carrier_webhook_events.processed_at IS NULL
      AND (carrier_webhook_events.claimed_at IS NULL OR carrier_webhook_events.claimed_at < now() - INTERVAL '5 minutes')
    RETURNING event_id
  `);

  if (executeRows<{ event_id: string }>(result).length > 0) {
    return { claimed: true, duplicate: false };
  }

  const existing = await db
    .select({ processedAt: carrierWebhookEvents.processedAt })
    .from(carrierWebhookEvents)
    .where(
      and(eq(carrierWebhookEvents.carrier, carrier), eq(carrierWebhookEvents.eventId, eventId)),
    )
    .limit(1);

  if (existing[0]?.processedAt != null) {
    return { claimed: false, duplicate: true };
  }

  return { claimed: false, duplicate: false };
}

export async function markCarrierWebhookProcessed(
  db: DrizzleClient,
  { carrier, eventId }: { carrier: string; eventId: string },
): Promise<void> {
  await db
    .update(carrierWebhookEvents)
    .set({ processedAt: new Date() })
    .where(
      and(eq(carrierWebhookEvents.carrier, carrier), eq(carrierWebhookEvents.eventId, eventId)),
    );
}

/** Clear in-flight claim so carrier redeliveries can re-enter after handler failure. */
export async function releaseCarrierWebhookEventClaim(
  db: DrizzleClient,
  { carrier, eventId }: { carrier: string; eventId: string },
): Promise<void> {
  await db
    .update(carrierWebhookEvents)
    .set({ claimedAt: null })
    .where(
      and(
        eq(carrierWebhookEvents.carrier, carrier),
        eq(carrierWebhookEvents.eventId, eventId),
        isNull(carrierWebhookEvents.processedAt),
      ),
    );
}
