/**
 * `ai_index_update` queue consumer — ai-assistant.
 *
 * Processes upsert/delete jobs enqueued by entity modules via `enqueueIndexUpsert`
 * / `enqueueIndexDelete` in `@zync/ai/rag`.
 *
 * - upsert: embed the entity text → upsert vector into namespace `tenant:{tenantId}`
 * - delete: remove vector by stable ID from the namespace
 *
 * Permanent failures (malformed payload) are acked to avoid poison loops.
 * Transient failures (Workers AI / Vectorize unavailable) trigger queue retry.
 */
import type { MessageBatch } from '@cloudflare/workers-types'
import type { Env } from '@zync/types'
import { embedText, vectorIdFor } from '@zync/ai/rag'
import type { AiIndexUpdateJob } from '@zync/ai/rag'

export async function handleAiIndexUpdate(
  batch: MessageBatch<AiIndexUpdateJob>,
  env: Env,
): Promise<void> {
  for (const msg of batch.messages) {
    const job = msg.body

    // Validate basic shape — ack malformed to avoid poison loop
    if (!job || typeof job.op !== 'string' || !job.tenantId || !job.entityId || !job.entityType) {
      console.warn('[ai-index-update] Malformed job, acking to discard:', job)
      msg.ack()
      continue
    }

    try {
      if (job.op === 'upsert') {
        // Validate upsert has text
        if (typeof job.text !== 'string' || job.text.length === 0) {
          console.warn('[ai-index-update] Upsert job missing text, discarding:', job)
          msg.ack()
          continue
        }

        const values = await embedText(env, job.text)
        const vectorId = vectorIdFor(job.entityType, job.entityId)

        await env.VECTORIZE.upsert([
          {
            id: vectorId,
            values,
            namespace: `tenant:${job.tenantId}`,
            metadata: {
              text: job.text,
              entity_type: job.entityType,
              entity_id: job.entityId,
              tenant_id: job.tenantId,
            },
          },
        ])

        msg.ack()
      } else if (job.op === 'delete') {
        const vectorId = vectorIdFor(job.entityType, job.entityId)
        await env.VECTORIZE.deleteByIds([vectorId])
        msg.ack()
      } else {
        // Unknown op — ack to discard
        console.warn('[ai-index-update] Unknown op, discarding:', (job as { op: string }).op)
        msg.ack()
      }
    } catch (err) {
      // Transient error — do NOT ack; let Queue retry with backoff
      console.error('[ai-index-update] Transient failure, will retry:', err)
      msg.retry()
    }
  }
}
