// Extract outcomes + tool_use rows from an action payload and write
// them to the indexed tables. Called from writer.ts at ingest time so
// /api/thomas and /api/tools never have to scan payload BLOBs.
//
// The detection logic is the same as the legacy query-time path
// (detectorFor(agent)), just moved upstream. Cross-action verification
// (matching tool_use to its tool_result) is also done here: when a new
// model_call lands whose messages contain tool_result blocks, we
// UPDATE outcomes.verified=1 and tool_uses.is_error for the matching
// tool_use ids.

import type Database from 'better-sqlite3'
import { logger } from '../../core/logger.ts'
import type { AgentPacket, ModelCallPayload } from '../../core/packet.ts'
import { SYNTHETIC_TOOLS } from '../orchestrator/vision-loop.ts'
import { detectorFor } from '../outcomes/detectors/index.ts'
import { detectMcpFrame } from '../outcomes/detectors/mcp.ts'
import { deriveToolTarget } from './tool-targets.ts'

const DEDUP_WINDOW_MS = 5 * 60 * 1000

// Prepared statements are cached per Database instance to avoid the
// parse cost on every ingest call.
type Stmts = {
  outcomeDedupCheck: Database.Statement
  outcomeInsert: Database.Statement
  toolUseInsert: Database.Statement
  outcomeVerifyById: Database.Statement
  toolUseErrorById: Database.Statement
  toolTargetInsert: Database.Statement
}
const cache = new WeakMap<Database.Database, Stmts>()

function prepare(raw: Database.Database): Stmts {
  let s = cache.get(raw)
  if (s) return s
  s = {
    outcomeDedupCheck: raw.prepare(
      'SELECT 1 FROM outcomes WHERE fingerprint = ? AND ts > ? LIMIT 1',
    ),
    outcomeInsert: raw.prepare(
      'INSERT INTO outcomes (action_id, run_id, thread_id, agent, ts, kind, value_usd_micro, fingerprint, tool_use_id, sub_agent, instance_id, verified) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0)',
    ),
    toolUseInsert: raw.prepare(
      'INSERT INTO tool_uses (action_id, agent, ts, name, tool_use_id, cost_micro_share, is_error) VALUES (?, ?, ?, ?, ?, ?, 0)',
    ),
    outcomeVerifyById: raw.prepare('UPDATE outcomes SET verified = 1 WHERE tool_use_id = ?'),
    toolUseErrorById: raw.prepare('UPDATE tool_uses SET is_error = 1 WHERE tool_use_id = ?'),
    toolTargetInsert: raw.prepare(
      'INSERT INTO tool_targets (action_id, run_id, thread_id, agent, ts, kind, op, target, tool_name, tool_use_id) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)',
    ),
  }
  cache.set(raw, s)
  return s
}

/**
 * Hook called after a packet is written to the actions table.
 * Best-effort: errors are logged but never propagate (we don't want
 * ingest to fail just because detection threw).
 */
export function extractAfterIngest(raw: Database.Database, packet: AgentPacket): void {
  try {
    runExtraction(raw, packet)
  } catch (err) {
    logger.warn(`ingest extract failed (${packet.id}): ${(err as Error).message}`)
  }
}

function runExtraction(raw: Database.Database, packet: AgentPacket): void {
  const stmts = prepare(raw)
  const payload = packet.payload

  // ── 1. tool_result verification pass ────────────────────────────
  // Whichever side carries tool_results (model_call user-content blocks
  // for Anthropic, role=tool messages for OpenAI, mcp_call responses for
  // MCP), iterate and mark matching outcomes / tool_uses.
  collectToolResults(payload).forEach((r) => {
    if (!r.isError) stmts.outcomeVerifyById.run(r.id)
    if (r.isError) stmts.toolUseErrorById.run(r.id)
  })

  // ── 2. detect outcomes for this action ──────────────────────────
  let events: ReturnType<ReturnType<typeof detectorFor>> = []
  if (packet.payload.kind === 'mcp_call') {
    events = detectMcpFrame({
      id: packet.id,
      runId: packet.runId,
      threadId: packet.threadId,
      kind: packet.payload.kind,
      sourceAgent: packet.sourceAgent,
      ts: packet.ts,
      payload: packet.payload,
    })
  } else {
    const detector = detectorFor(packet.sourceAgent)
    events = detector({
      id: packet.id,
      runId: packet.runId,
      threadId: packet.threadId,
      kind: packet.payload.kind,
      sourceAgent: packet.sourceAgent,
      ts: packet.ts,
      payload: packet.payload,
    })
  }

  for (const ev of events) {
    // 5-minute dedup window against the same fingerprint.
    const dup = stmts.outcomeDedupCheck.get(ev.fingerprint, ev.detectedAt - DEDUP_WINDOW_MS)
    if (dup) continue
    stmts.outcomeInsert.run(
      packet.id,
      packet.runId,
      packet.threadId,
      packet.sourceAgent,
      ev.detectedAt,
      ev.kind,
      Math.round(ev.valueUsd * 1_000_000),
      ev.fingerprint,
      ev.toolUseId ?? null,
      ev.subAgent ?? null,
      packet.instanceId ?? null,
    )
    // If the matching tool_result already arrived (rare but possible —
    // the agent could batch them), backfill verified now.
    if (ev.toolUseId) {
      const existing = raw
        .prepare('SELECT 1 FROM tool_uses WHERE tool_use_id = ? AND is_error = 0 LIMIT 1')
        .get(ev.toolUseId)
      if (existing) stmts.outcomeVerifyById.run(ev.toolUseId)
    }
  }

  // ── 3. record tool_use invocations for /api/tools ───────────────
  if (packet.payload.kind === 'model_call') {
    const mcp = packet.payload as ModelCallPayload
    const response = mcp.response
    if (Array.isArray(response)) {
      const toolUses: Array<{ name: string; id?: string }> = []
      for (const block of response) {
        if (!block || typeof block !== 'object') continue
        const b = block as { type?: unknown; name?: unknown; id?: unknown; input?: unknown }
        if (b.type !== 'tool_use') continue
        if (typeof b.name !== 'string') continue
        // Tools Thomas injects (e.g. view_image) are not the agent's own.
        if (SYNTHETIC_TOOLS.has(b.name)) continue
        const id = typeof b.id === 'string' ? b.id : undefined
        toolUses.push({ name: b.name, id })

        // Derive a file op / shell exec from the call's input, if any.
        const target = deriveToolTarget(b.name, b.input)
        if (target) {
          stmts.toolTargetInsert.run(
            packet.id,
            packet.runId,
            packet.threadId,
            packet.sourceAgent,
            packet.ts,
            target.kind,
            target.op,
            target.target,
            b.name,
            id ?? null,
          )
        }
      }
      const costMicro =
        typeof packet.cost?.usd === 'number' ? Math.round(packet.cost.usd * 1_000_000) : 0
      const sharePerTool = toolUses.length > 0 ? Math.round(costMicro / toolUses.length) : 0
      for (const tu of toolUses) {
        stmts.toolUseInsert.run(
          packet.id,
          packet.sourceAgent,
          packet.ts,
          tu.name,
          tu.id ?? null,
          sharePerTool,
        )
      }
    }
  }
}

type ToolResultInfo = { id: string; isError: boolean }

// biome-ignore lint/suspicious/noExplicitAny: payload shapes vary by agent
function collectToolResults(payload: any): ToolResultInfo[] {
  const out: ToolResultInfo[] = []
  if (!payload) return out

  if (payload.kind === 'model_call') {
    // Anthropic-style tool_result blocks live in user messages of the
    // next model_call's request.messages.
    const messages = payload.messages
    if (Array.isArray(messages)) {
      for (const m of messages) {
        if (!m) continue
        // OpenAI-style: role='tool' messages.
        if (m.role === 'tool' && typeof m.tool_call_id === 'string') {
          out.push({ id: m.tool_call_id, isError: false })
          continue
        }
        if (m.role !== 'user' || !Array.isArray(m.content)) continue
        for (const block of m.content) {
          if (!block || block.type !== 'tool_result') continue
          const id = block.tool_use_id ?? block.tool_call_id
          if (typeof id !== 'string') continue
          out.push({ id, isError: block.is_error === true })
        }
      }
    }
  } else if (payload.kind === 'mcp_call' && payload.direction === 'response') {
    const id =
      typeof payload.jsonrpcId === 'number' || typeof payload.jsonrpcId === 'string'
        ? `mcp:${typeof payload.server === 'string' ? payload.server : ''}:${payload.jsonrpcId}`
        : null
    if (id) out.push({ id, isError: !!payload.error })
  }
  return out
}
