import { mkdtempSync, rmSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { afterAll, describe, expect, it } from 'vitest'
import { newActionId, newRunId } from '../../core/ids.ts'
import type { AgentPacket } from '../../core/packet.ts'

// Point the store at a throwaway home BEFORE anything imports paths.ts.
const tmp = mkdtempSync(join(tmpdir(), 'thomas-threads-'))
process.env.THOMAS_HOME = tmp

const { correlatedThreadId, deriveThreadTitle } = await import('../decoders/correlate.ts')
const { appendPacket } = await import('./writer.ts')
const { listThreads, getThreadDetail, countThreads, getRunDetail } = await import('./queries.ts')
const { closeDb } = await import('./db.ts')
const { api } = await import('../api/index.ts')
const { computeThomas } = await import('./thomas-queries.ts')

afterAll(() => {
  closeDb()
  rmSync(tmp, { recursive: true, force: true })
})

let ts = 1_700_000_000_000

function turn(agent: string, messages: unknown[], usd: number): AgentPacket {
  ts += 1000
  const title = deriveThreadTitle(messages)
  return {
    id: newActionId(),
    runId: newRunId(),
    threadId: correlatedThreadId(agent, messages),
    ...(title ? { threadTitle: title } : {}),
    ts,
    durMs: 500,
    sourceAgent: agent,
    cost: { usd, tokensIn: 100, tokensOut: 50 },
    payload: {
      kind: 'model_call',
      protocol: 'anthropic',
      endpoint: 'https://api.anthropic.com/v1/messages',
      model: 'claude-opus-4-7',
      messages,
      stream: false,
      response: [],
      status: 200,
    },
  }
}

describe('thread correlation → task aggregation (end to end)', () => {
  it('collapses a multi-turn conversation into one task and keeps a separate one apart', async () => {
    // Conversation #1 — three requests, history grows as a prefix.
    const c1t1 = [{ role: 'user', content: 'build a todo app' }]
    const c1t2 = [
      { role: 'user', content: 'build a todo app' },
      { role: 'assistant', content: 'ok' },
      { role: 'user', content: 'add dark mode' },
    ]
    const c1t3 = [
      ...c1t2,
      { role: 'assistant', content: 'done' },
      { role: 'user', content: 'ship it' },
    ]
    await appendPacket(turn('claude-code', c1t1, 0.01))
    await appendPacket(turn('claude-code', c1t2, 0.02))
    await appendPacket(turn('claude-code', c1t3, 0.03))

    // Conversation #2 — a different goal.
    await appendPacket(turn('claude-code', [{ role: 'user', content: 'fix the flaky test' }], 0.05))

    const threads = await listThreads({ limit: 50 })
    expect(threads.length).toBe(2)
    expect(await countThreads()).toBe(2)

    const todo = threads.find((t) => t.title === 'build a todo app')
    expect(todo).toBeDefined()
    expect(todo?.runCount).toBe(3)
    expect(todo?.actionCount).toBe(3)
    // 0.01 + 0.02 + 0.03 — one task's total spend.
    expect(todo?.costUsd).toBeCloseTo(0.06, 6)
    expect(todo?.tokensIn).toBe(300)

    const flaky = threads.find((t) => t.title === 'fix the flaky test')
    expect(flaky?.runCount).toBe(1)
    expect(flaky?.costUsd).toBeCloseTo(0.05, 6)

    // Task detail lists each underlying run for drill-down.
    if (!todo) throw new Error('missing todo task')
    const detail = await getThreadDetail(todo.id)
    expect(detail).not.toBeNull()
    expect(detail?.thread.runCount).toBe(3)
    expect(detail?.runs.length).toBe(3)
    expect(detail?.thread.title).toBe('build a todo app')
    // Plain single-call runs are not flagged as routed.
    expect(detail?.runs.every((r) => r.routed === false)).toBe(true)
  })

  it('a later turn lands in the same task (deterministic, restart-safe)', async () => {
    const before = (await listThreads({ limit: 50 })).find((t) => t.title === 'build a todo app')
    expect(before?.runCount).toBe(3)

    // Same conversation resumes (same first message) — even if the daemon had
    // restarted, the id is recomputed identically.
    const resume = [
      { role: 'user', content: 'build a todo app' },
      { role: 'assistant', content: 'shipped' },
      { role: 'user', content: 'now add tests' },
    ]
    await appendPacket(turn('claude-code', resume, 0.04))

    const after = (await listThreads({ limit: 50 })).find((t) => t.title === 'build a todo app')
    expect(after?.runCount).toBe(4)
    expect(after?.costUsd).toBeCloseTo(0.1, 6)
  })

  it('serves tasks over the HTTP API (route registration + handlers + errored flag)', async () => {
    // Bare array (no limit/offset) — the CLI shape.
    const bare = await api.request('/threads')
    expect(bare.status).toBe(200)
    const list = (await bare.json()) as Array<{
      id: string
      title: string | null
      runCount: number
      errored: boolean
      thomas: number
    }>
    expect(list.length).toBe(2)
    for (const t of list) expect(typeof t.errored).toBe('boolean')

    // Paged envelope.
    const paged = await api.request('/threads?limit=1&offset=0')
    const env = (await paged.json()) as { rows: unknown[]; total: number }
    expect(env.total).toBe(2)
    expect(env.rows.length).toBe(1)

    // Detail with nested runs.
    const todo = list.find((t) => t.title === 'build a todo app')
    if (!todo) throw new Error('missing todo task in API response')
    const detailRes = await api.request(`/threads/${encodeURIComponent(todo.id)}`)
    expect(detailRes.status).toBe(200)
    const detail = (await detailRes.json()) as {
      thread: { runCount: number; errored: boolean }
      runs: Array<{ errored: boolean }>
    }
    expect(detail.thread.runCount).toBe(4)
    expect(detail.runs.length).toBe(4)
    // errored flag enriched on both the task and each nested run.
    expect(typeof detail.thread.errored).toBe('boolean')
    for (const r of detail.runs) expect(typeof r.errored).toBe('boolean')

    // Unknown id → 404.
    const missing = await api.request('/threads/th_does_not_exist')
    expect(missing.status).toBe(404)
  })

  it('round-trips raw req/res through gzip storage, gated by includeRaw', async () => {
    const reqBytes = new TextEncoder().encode(JSON.stringify({ model: 'm', messages: [] }))
    const resBytes = new TextEncoder().encode('data: {"hello":"world"}\n\n')
    const packet: AgentPacket = {
      id: newActionId(),
      runId: newRunId(),
      threadId: correlatedThreadId('claude-code', [{ role: 'user', content: 'raw test' }]),
      threadTitle: 'raw test',
      rawReq: reqBytes,
      rawRes: resBytes,
      ts: (ts += 1000),
      durMs: 100,
      sourceAgent: 'claude-code',
      cost: { usd: 0.01, tokensIn: 1, tokensOut: 1 },
      payload: {
        kind: 'model_call',
        protocol: 'anthropic',
        endpoint: 'https://api.anthropic.com/v1/messages',
        model: 'm',
        messages: [{ role: 'user', content: 'raw test' }],
        stream: true,
        response: [],
        status: 200,
      },
    }
    await appendPacket(packet)

    // Without includeRaw: no raw fields leak into the default detail.
    const plain = await getRunDetail(packet.runId)
    expect(plain?.actions[0]?.rawReq).toBeUndefined()

    // With includeRaw: gunzipped back to the exact original text.
    const withRaw = await getRunDetail(packet.runId, { includeRaw: true })
    expect(withRaw?.actions[0]?.rawReq).toBe(new TextDecoder().decode(reqBytes))
    expect(withRaw?.actions[0]?.rawRes).toBe(new TextDecoder().decode(resBytes))
  })

  it('derives file ops + shell execs from tool_use blocks at ingest', async () => {
    const packet: AgentPacket = {
      id: newActionId(),
      runId: newRunId(),
      threadId: correlatedThreadId('claude-code', [{ role: 'user', content: 'do work' }]),
      ts: (ts += 1000),
      durMs: 100,
      sourceAgent: 'claude-code',
      cost: { usd: 0.02, tokensIn: 10, tokensOut: 10 },
      payload: {
        kind: 'model_call',
        protocol: 'anthropic',
        endpoint: 'https://api.anthropic.com/v1/messages',
        model: 'claude-opus-4-7',
        messages: [{ role: 'user', content: 'do work' }],
        stream: false,
        status: 200,
        response: [
          { type: 'tool_use', id: 't1', name: 'Edit', input: { file_path: '/src/a.ts' } },
          { type: 'tool_use', id: 't2', name: 'Bash', input: { command: 'npm test' } },
          { type: 'tool_use', id: 't3', name: 'web_search', input: { query: 'x' } },
        ],
      },
    }
    await appendPacket(packet)

    const detail = await getRunDetail(packet.runId)
    const targets = detail?.toolTargets ?? []
    expect(targets.length).toBe(2) // web_search is not a file/shell op
    const fs = targets.find((t) => t.kind === 'fs')
    const sh = targets.find((t) => t.kind === 'shell')
    expect(fs).toMatchObject({ op: 'edit', target: '/src/a.ts', toolName: 'Edit' })
    expect(sh).toMatchObject({ op: 'exec', target: 'npm test', toolName: 'Bash' })
  })

  it('splits the same prompt across instances into separate tasks + per-instance breakdown', async () => {
    const msgs = [{ role: 'user', content: 'fleet identical prompt' }]
    const inst = (instanceId: string, usd: number): AgentPacket => ({
      id: newActionId(),
      runId: newRunId(),
      threadId: correlatedThreadId('openclaw', msgs, instanceId),
      threadTitle: deriveThreadTitle(msgs),
      instanceId,
      ts: (ts += 1000),
      durMs: 100,
      sourceAgent: 'openclaw',
      cost: { usd, tokensIn: 10, tokensOut: 5 },
      payload: {
        kind: 'model_call',
        protocol: 'openai-chat',
        endpoint: 'http://up',
        model: `thomas-openclaw-${instanceId}`,
        messages: msgs,
        stream: false,
        status: 200,
        response: [],
      },
    })
    await appendPacket(inst('main', 0.03))
    await appendPacket(inst('casual', 0.01))

    // Same prompt, two instances → two distinct tasks (instance-salted).
    const oc = (await listThreads({ agent: ['openclaw'], limit: 50 })).filter(
      (t) => t.title === 'fleet identical prompt',
    )
    expect(oc.length).toBe(2)
    expect(new Set(oc.map((t) => t.id)).size).toBe(2)

    // Per-instance cost breakdown surfaces both.
    const thomas = await computeThomas('lifetime')
    const agent = thomas.perAgent.find((a) => a.agent === 'openclaw')
    const byInst = new Map(agent?.instances.map((i) => [i.instanceId, i]))
    expect(byInst.get('main')?.costUsd).toBeCloseTo(0.03, 6)
    expect(byInst.get('casual')?.costUsd).toBeCloseTo(0.01, 6)
    expect(byInst.get('main')?.taskCount).toBe(1)

    // Home's spend ranking is per-task (thread), ordered by cost desc.
    expect(thomas.topTasks.length).toBeGreaterThan(0)
    for (const t of thomas.topTasks) {
      expect(typeof t.threadId).toBe('string')
      expect(t.runCount).toBeGreaterThanOrEqual(1)
    }
    const costs = thomas.topTasks.map((t) => t.costUsd)
    expect([...costs].sort((a, b) => b - a)).toEqual(costs) // already sorted desc
  })

  it('flags a run with fan-out children as routed', async () => {
    const runId = newRunId()
    const parentId = newActionId()
    const act = (id: `ac_${string}`, parent?: `ac_${string}`): AgentPacket => ({
      id,
      runId,
      threadId: 'th_routed',
      ...(parent ? { parentActionId: parent } : {}),
      ts: (ts += 1000),
      durMs: 50,
      sourceAgent: 'claude-code',
      cost: { usd: 0.01, tokensIn: 1, tokensOut: 1 },
      payload: {
        kind: 'model_call',
        protocol: 'anthropic',
        endpoint: 'x',
        model: 'm',
        messages: [],
        stream: false,
        status: 200,
        response: [],
      },
    })
    await appendPacket(act(parentId)) // parent (no parentActionId)
    await appendPacket(act(newActionId(), parentId)) // child → run is "routed"

    const d = await getRunDetail(runId)
    expect(d?.run.routed).toBe(true)
  })
})
