import 'server-only'

import { createHash } from 'node:crypto'
import { createAdminClient } from '@/lib/supabase/admin'
import { extractPermit, type PermitExtraction } from '@/lib/adapters/flask'
import { resolveDocumentUrl, type DocumentLocation } from '@/lib/data/document-urls'
import { costFieldsFromExtraction } from '@/lib/data/permit-cost'
import { sendPermitDimensionAlerts } from '@/lib/data/permit-dimension-alert'
import { reconcilePermitProfiles } from '@/lib/data/permit-profiles'
import { activeDimensionWarnings, workspaceDimensionAlertsEnabled } from '@/lib/domain/dimension-alerts'
import { computePermitWarnings } from '@/lib/domain/warnings'
import { statusAfterPermitProcessed } from '@/lib/domain/status'
import { syncTripToHha } from '@/lib/integrations/hha-trip-sync'
import { isSynchronPermitUrl } from '@/lib/integrations/synchron-contract'
import { logOperationEvent } from '@/lib/observability/operation-events'
import type { Permit, Trip } from '@/types/db'

export interface PermitJob {
  permit_id: string
  revision: number
  lease_token: string
  attempts: number
  notify_contacts: boolean
}

export function permitDownloadAllowed(value: string): boolean {
  if (isSynchronPermitUrl(value)) return true
  try {
    const url = new URL(value)
    const hosts = ['heavy-haul-agent.s3.us-east-1.amazonaws.com']
    if (process.env.NEXT_PUBLIC_SUPABASE_URL) hosts.push(new URL(process.env.NEXT_PUBLIC_SUPABASE_URL).hostname)
    return url.protocol === 'https:' && !url.username && !url.password && !url.port && hosts.includes(url.hostname)
  } catch { return false }
}

async function downloadFile(document: DocumentLocation): Promise<Blob> {
  const url = await resolveDocumentUrl(document)
  if (!url || !permitDownloadAllowed(url)) throw new Error('permit_file_unavailable')
  const response = await fetch(url, { redirect: 'error', signal: AbortSignal.timeout(30_000) })
  if (!response.ok || !response.body) throw new Error('permit_download_failed')
  const reader = response.body.getReader()
  const chunks: Uint8Array<ArrayBuffer>[] = []
  let size = 0
  try {
    while (true) {
      const { done, value } = await reader.read()
      if (done) break
      size += value.byteLength
      if (size > 25 * 1024 * 1024) { await reader.cancel(); throw new Error('permit_too_large') }
      chunks.push(new Uint8Array(value))
    }
  } finally { reader.releaseLock() }
  return new Blob(chunks)
}

export function permitExtractionFields(permit: Permit, p: PermitExtraction) {
  const manualCost = ['manually_entered', 'corrected'].includes(permit.permit_cost_status || '')
  return {
    state_code: p.state_code || permit.state_code,
    permit_number: p.permit_number ?? permit.permit_number,
    effective_date: p.effective_date ?? permit.effective_date,
    expiration_date: p.expiration_date ?? permit.expiration_date,
    permit_length_in: p.length_in ?? null, permit_width_in: p.width_in ?? null,
    permit_height_in: p.height_in ?? null, permit_weight_lbs: p.weight_lbs ?? null,
    ...(!manualCost ? costFieldsFromExtraction(p) : {}),
    route_token: p.route_token ?? null,
    route_processing_status: p.route_status || 'failed',
    route_links: p.route_links || [], route_gpx_url: p.route_gpx_url || null,
    route_hummer_url: p.route_hummer_url || null,
    extraction: p, extraction_status: 'processed' as const,
  }
}

/** One entry point for uploads, callbacks, and the import queue. No billing. */
export async function processWorkspacePermit(params: {
  permitId?: string
  origin: string
  file?: Blob
}): Promise<{ ok: boolean; permitId?: string; status: string; error?: string }> {
  const admin = createAdminClient()
  if (params.permitId) {
    // Native inserts already have a job via the trigger. Older permits may
    // not: recover them without unexpectedly emailing historical contacts.
    const queued = await admin.from('permit_processing_jobs').upsert({
      permit_id: params.permitId, notify_contacts: false,
    }, { onConflict: 'permit_id', ignoreDuplicates: true })
    if (queued.error) return { ok: false, status: 'pending', error: 'permit_queue_unavailable' }
  }
  const claim = await admin.rpc('claim_workspace_permit', { p_permit_id: params.permitId || null })
  if (claim.error) return { ok: false, status: 'pending', error: 'permit_queue_unavailable' }
  const job = (claim.data as PermitJob[] | null)?.[0]
  if (!job) return { ok: true, status: 'queued_or_already_processed', permitId: params.permitId }
  let logTripId: string | null = null
  await logOperationEvent({ area: 'permit', event: 'processing', outcome: 'started', permitId: job.permit_id, detail: { revision: job.revision, attempt: job.attempts, notify_contacts: job.notify_contacts } })
  const finish = (status: 'done' | 'retry', error?: string) => admin.from('permit_processing_jobs').update({
    status, lease_until: null, lease_token: null, last_error: error || null,
    available_at: new Date(Date.now() + Math.min(3600, 60 * 2 ** Math.min(job.attempts, 6)) * 1000).toISOString(),
    updated_at: new Date().toISOString(),
  }).eq('permit_id', job.permit_id).eq('lease_token', job.lease_token).eq('revision', job.revision)
  try {
    const { data: permitRow, error: permitError } = await admin.from('permits').select('*').eq('id', job.permit_id).single()
    if (permitError || !permitRow?.document_id) throw new Error('permit_not_found')
    const permit = permitRow as Permit
    logTripId = permit.trip_id
    const { data: document, error: documentError } = await admin.from('documents').select('*')
      .eq('id', permit.document_id).eq('trip_id', permit.trip_id).single()
    if (documentError || !document) throw new Error('permit_document_not_found')
    const file = params.file || await downloadFile(document)
    if (!file.size || file.size > 25 * 1024 * 1024) throw new Error('invalid_permit_size')
    const bytes = Buffer.from(await file.arrayBuffer())
    const pdf = bytes.subarray(0, 5).toString() === '%PDF-'
    const png = bytes[0] === 0x89 && bytes[1] === 0x50 && bytes[2] === 0x4e && bytes[3] === 0x47
    const jpeg = bytes[0] === 0xff && bytes[1] === 0xd8 && bytes[2] === 0xff
    if (!pdf && !png && !jpeg) throw new Error('unsupported_permit_file')
    const hash = createHash('sha256').update(bytes).digest('hex')
    let p: PermitExtraction
    if (permit.extraction?._workspace_file_sha256 === hash && permit.extraction_status === 'processed'
        && permit.route_processing_status !== 'failed') {
      p = { ...permit.extraction } as PermitExtraction
    } else {
      const sync = await syncTripToHha(permit.trip_id)
      if (!sync.ok) throw new Error('trip_sync_unavailable')
      const extracted = await extractPermit(file, document.file_name || 'permit.pdf', {
        tripId: permit.trip_id, documentId: document.id, permitId: permit.id,
        sourceUrl: document.external_url || undefined,
      })
      if (!extracted.ok || !extracted.data) throw new Error('permit_extraction_unavailable')
      p = extracted.data
    }
    // Do not silently assign a provider's state item to a different state.
    if (permit.source_system === 'synchron' && permit.state_code && p.state_code
        && permit.state_code !== p.state_code) throw new Error('permit_state_mismatch')
    const { data: tripRow, error: tripError } = await admin.from('trips').select('*').eq('id', permit.trip_id).single()
    if (tripError || !tripRow) throw new Error('trip_not_found')
    const trip = tripRow as Trip
    const { data: siblings, error: siblingError } = await admin.from('permits')
      .select('id,state_code,effective_date,expiration_date').eq('trip_id', trip.id)
    if (siblingError) throw new Error('permit_siblings_unavailable')
    const fields = permitExtractionFields(permit, p)
    fields.extraction = { ...p, _workspace_file_sha256: hash } as PermitExtraction
    const saved = { ...permit, ...fields } as Permit
    const computedWarnings = computePermitWarnings(trip, saved, new Date(), siblings || [])
    const warnings = activeDimensionWarnings(computedWarnings)
    const publish = await admin.rpc('publish_workspace_permit', {
      p_permit_id: permit.id, p_lease: job.lease_token, p_revision: job.revision,
      p_result: fields, p_warnings: warnings,
    })
    if (publish.error) throw new Error('permit_result_save_failed')
    if (!publish.data) {
      await logOperationEvent({ area: 'permit', event: 'processing', outcome: 'skipped', tripId: trip.id, permitId: permit.id, documentId: permit.document_id, detail: { reason: 'superseded', revision: job.revision } })
      return { ok: true, permitId: permit.id, status: 'superseded' }
    }
    const profileObservation = p.raw?.profile_observation
    if (profileObservation && typeof profileObservation === 'object' && !Array.isArray(profileObservation)) {
      try {
        await reconcilePermitProfiles({
          source_system: 'workspace_permit', source_key: permit.id,
          trip_id: trip.id, permit_id: permit.id, document_id: document.id,
          synchron_order_token: trip.synchron_order_token || null,
          // The current Synchron callback supplies an item ID, not a permit token.
          synchron_permit_token: null,
          file_sha256: hash, observation: profileObservation as Record<string, unknown>,
        })
      } catch (error) {
        console.error('permit_profile_reconciliation', { permit_id: permit.id, error: error instanceof Error ? error.message : 'unknown' })
      }
    }
    const status = statusAfterPermitProcessed(trip.status)
    if (status !== trip.status) await admin.from('trips').update({ status }).eq('id', trip.id).eq('status', trip.status)
    const sync = await syncTripToHha(trip.id)
    if (!sync.ok) throw new Error('trip_sync_unavailable')
    if (job.notify_contacts && workspaceDimensionAlertsEnabled()) {
      const alert = await sendPermitDimensionAlerts({ trip, permit: saved, warnings, origin: params.origin })
      if (alert.failed) throw new Error('dimension_email_unavailable')
    } else {
      await logOperationEvent({ area: 'permit', event: 'dimension_alert_checked', outcome: 'skipped', tripId: trip.id, permitId: permit.id, detail: { reason: workspaceDimensionAlertsEnabled() ? 'notifications_disabled' : 'dimension_alert_paused', mismatch_count: computedWarnings.filter((warning) => warning.kind === 'dimension_mismatch').length } })
    }
    if (p.route_status === 'failed') throw new Error('route_processing_unavailable')
    // Token generation is complete. The existing GPT API has no map-refresh
    // endpoint; do not repeat paid extraction just to look for map links.
    const result = await finish('done')
    if (result.error) throw new Error('permit_job_save_failed')
    await logOperationEvent({ area: 'permit', event: 'processing', outcome: 'success', tripId: trip.id, permitId: permit.id, documentId: permit.document_id, detail: { revision: job.revision, attempt: job.attempts, route_status: p.route_status || 'needs_preparation', route_link_count: p.route_links?.length || 0, warning_count: warnings.length, notify_contacts: job.notify_contacts } })
    return { ok: true, permitId: permit.id, status: p.route_status || 'needs_preparation' }
  } catch (error) {
    // Only fixed internal error codes reach logs/responses, never URLs, tokens,
    // provider bodies or secrets.
    const code = error instanceof Error && /^[a-z_]+$/.test(error.message)
      ? error.message : 'permit_processing_unavailable'
    await finish('retry', code)
    console.error('workspace_permit_processing', { permit_id: job.permit_id, revision: job.revision, error: code })
    await logOperationEvent({ area: 'permit', event: 'processing', outcome: 'failed', tripId: logTripId, permitId: job.permit_id, detail: { revision: job.revision, attempt: job.attempts, error_code: code } })
    return { ok: false, permitId: job.permit_id, status: 'pending', error: code }
  }
}
