// Claude Code session correlator (opt-in). The wire can't tell which window
// (instance) or parallel sub-agent made a call — they share path, key, and
// headers. But Claude Code's local transcripts can: each assistant record in
// `~/.claude/projects/<proj>/<sessionId>.jsonl` carries the upstream
// `message.id`, and Thomas already stores the upstream response bytes
// (raw_res) which contain that same id. Joining on `message.id` deterministic-
// ally attributes each captured call to its sessionId (= instance) and, for
// sidechain turns, its spawning sub-agent.
//
// Strictly local + opt-in (config.correlateSessions): reads files already on
// the user's disk, sends nothing, and only ever back-fills NULL columns. When
// off, absent, or cross-machine (BYOA), instance_id stays null and everything
// degrades to the per-type view. Privacy contract intact (PRIVACY.md).

import { readFileSync, readdirSync, statSync } from 'node:fs'
import { homedir } from 'node:os'
import { join } from 'node:path'
import { gunzipSync } from 'node:zlib'
import type Database from 'better-sqlite3'
import { readConfig } from '../../core/config.ts'
import { logger } from '../../core/logger.ts'
import { getRawDb } from '../store/db.ts'

export const CLAUDE_PROJECTS_DIR = join(homedir(), '.claude', 'projects')

// The upstream message id, both in the JSONL (`message.id`) and in the
// captured raw_res (`message_start` event / non-stream body). First match in
// a response is the message's own id.
const MSG_ID_RE = /"id"\s*:\s*"(msg_[A-Za-z0-9_-]+)"/

export type SessionHit = { sessionId: string; subAgentId: string | null }

/** Scan Claude Code transcripts modified since `sinceMs` and build a
 *  `message.id → { sessionId, subAgentId }` map. sessionId is the file's
 *  session (the instance); subAgentId is the sidechain's originating
 *  sub-agent (sourceToolAssistantUUID) when the turn is sub-agent work. */
export function buildSessionIndex(projectsDir: string, sinceMs: number): Map<string, SessionHit> {
  const index = new Map<string, SessionHit>()
  let files: string[]
  try {
    files = listJsonl(projectsDir)
  } catch {
    return index
  }
  for (const file of files) {
    let mtime: number
    try {
      mtime = statSync(file).mtimeMs
    } catch {
      continue
    }
    if (mtime < sinceMs) continue
    let text: string
    try {
      text = readFileSync(file, 'utf8')
    } catch {
      continue
    }
    for (const line of text.split('\n')) {
      if (!line.includes('"message"')) continue
      let rec: Record<string, unknown>
      try {
        rec = JSON.parse(line)
      } catch {
        continue
      }
      if (rec.type !== 'assistant') continue
      const msg = rec.message as { id?: unknown } | undefined
      const id = typeof msg?.id === 'string' ? msg.id : undefined
      const sessionId = typeof rec.sessionId === 'string' ? rec.sessionId : undefined
      if (!id || !sessionId) continue
      const sub =
        rec.isSidechain === true
          ? typeof rec.sourceToolAssistantUUID === 'string'
            ? rec.sourceToolAssistantUUID
            : typeof rec.parentUuid === 'string'
              ? rec.parentUuid
              : 'sidechain'
          : null
      index.set(id, { sessionId, subAgentId: sub })
    }
  }
  return index
}

function listJsonl(dir: string): string[] {
  const out: string[] = []
  for (const entry of readdirSync(dir, { withFileTypes: true })) {
    const p = join(dir, entry.name)
    if (entry.isDirectory()) out.push(...safeList(p))
    else if (entry.isFile() && entry.name.endsWith('.jsonl')) out.push(p)
  }
  return out
}
function safeList(dir: string): string[] {
  try {
    return readdirSync(dir, { withFileTypes: true })
      .filter((e) => e.isFile() && e.name.endsWith('.jsonl'))
      .map((e) => join(dir, e.name))
  } catch {
    return []
  }
}

/** Extract the upstream message id from a (gunzipped) raw_res blob. */
export function msgIdFromRawRes(blob: unknown): string | null {
  if (!Buffer.isBuffer(blob) || blob.length === 0) return null
  try {
    return MSG_ID_RE.exec(gunzipSync(blob).toString('utf8'))?.[1] ?? null
  } catch {
    return null
  }
}

/** Back-fill instance_id / sub_agent_id on still-uncorrelated claude-code
 *  actions (and their outcomes) by joining stored raw_res to the session
 *  index on message.id. Returns the number of actions correlated. Best-
 *  effort: only NULL columns are written, so re-runs are idempotent. */
export function correlateActions(
  raw: Database.Database,
  index: Map<string, SessionHit>,
  sinceMs: number,
): number {
  if (index.size === 0) return 0
  const rows = raw
    .prepare(
      `SELECT a.id AS id, ap.raw_res AS raw_res
       FROM actions a
       JOIN action_payloads ap ON ap.action_id = a.id
       WHERE a.source_agent = 'claude-code'
         AND a.instance_id IS NULL
         AND a.kind = 'model_call'
         AND ap.raw_res IS NOT NULL
         AND a.ts >= ?`,
    )
    .all(sinceMs) as Array<{ id: string; raw_res: unknown }>

  const updAction = raw.prepare(
    'UPDATE actions SET instance_id = ?, sub_agent_id = ? WHERE id = ? AND instance_id IS NULL',
  )
  const updOutcome = raw.prepare(
    'UPDATE outcomes SET instance_id = ?, sub_agent_id = ? WHERE action_id = ? AND instance_id IS NULL',
  )

  let correlated = 0
  const apply = raw.transaction((batch: Array<{ id: string; raw_res: unknown }>) => {
    for (const r of batch) {
      const msgId = msgIdFromRawRes(r.raw_res)
      if (!msgId) continue
      const hit = index.get(msgId)
      if (!hit) continue
      updAction.run(hit.sessionId, hit.subAgentId, r.id)
      updOutcome.run(hit.sessionId, hit.subAgentId, r.id)
      correlated++
    }
  })
  apply(rows)
  return correlated
}

const LOOKBACK_MS = 2 * 24 * 3_600_000 // only scan/correlate the last ~2 days
let timer: ReturnType<typeof setInterval> | null = null

/** Opt-in background loop: every few minutes, join recent Claude Code
 *  transcripts to uncorrelated captures. Disabled unless
 *  config.correlateSessions. */
export function startSessionCorrelatorLoop(): void {
  if (timer) return
  const tick = async (): Promise<void> => {
    try {
      if (!(await readConfig()).correlateSessions) return
      const since = Date.now() - LOOKBACK_MS
      const index = buildSessionIndex(CLAUDE_PROJECTS_DIR, since)
      if (index.size === 0) return
      const raw = await getRawDb()
      const n = correlateActions(raw, index, since)
      if (n > 0) logger.info(`fleet: correlated ${n} claude-code call(s) to their session`)
    } catch (err) {
      logger.debug(`fleet: session correlate failed — ${(err as Error).message}`)
    }
  }
  const firstDelay = 90_000 + Math.floor(Math.random() * 60_000)
  setTimeout(() => void tick(), firstDelay).unref()
  timer = setInterval(() => void tick(), 5 * 60_000)
  timer.unref()
  logger.info('fleet: session correlator enabled (opt-in)')
}

export function stopSessionCorrelatorLoop(): void {
  if (timer) {
    clearInterval(timer)
    timer = null
  }
}
