import { readFile } from "node:fs/promises";
import type { DatabaseAdapter } from "../db.js";
import { hmacHex } from "../privacy.js";
import { FileImportError, type FileImportWarning } from "./file-import-contract.js";

/**
 * Cursor keeps no token usage in its local storage (bubble `tokenCount` fields are always zero),
 * but the Cursor dashboard exports per-request usage events as CSV ("Usage" -> export). An export
 * dropped in `<analytics dir>/imports/cursor/*.csv` is imported as telemetry rows:
 * - input_tokens = "Input (w/o Cache Write)", cache_write_tokens = "Input (w/ Cache Write)",
 *   cache_read_tokens = "Cache Read", output_tokens = "Output Tokens";
 * - each event is attached to the local Cursor session whose message is closest in time (events
 *   are written 1-3 s after the local bubble); events without a local conversation go to one
 *   technical `task` session per export file, so their tokens are still counted;
 * - rows are keyed by the event values, so overlapping exports never count an event twice.
 * Only technical columns are read: the user column is used solely to refuse multi-user team
 * exports and is never stored; cost, repository, PR and agent identifiers are ignored.
 */

interface ImportContext {
  db: DatabaseAdapter;
  sourceFileId: string;
  machineId: string;
  hmacSalt: string;
}

interface ImportCounters {
  sessions_upserted: number;
  messages_upserted: number;
  runtime_events_upserted: number;
}

const REQUIRED_COLUMNS = ["Date", "User", "Model", "Input (w/ Cache Write)", "Input (w/o Cache Write)", "Cache Read", "Output Tokens"] as const;
const MAX_EXPORT_BYTES = 64 * 1024 * 1024;
// A usage event is recorded after the request completes: look back for the local message that
// started it, and allow a small forward skew between the local clock and the server timestamp.
const MATCH_LOOKBACK_MS = 3 * 60 * 60 * 1000;
const MATCH_LOOKAHEAD_MS = 60 * 1000;
// Usage rows share the session with bubble rows (seq 1..n): offset their seq far above them.
const USAGE_SEQ_OFFSET = 10_000_000_000_000;

/** Minimal RFC 4180 parser: quoted fields, escaped quotes, commas and newlines inside quotes. */
export function parseCsv(raw: string): string[][] {
  const rows: string[][] = [];
  let row: string[] = [];
  let field = "";
  let quoted = false;
  const text = raw.charCodeAt(0) === 0xfeff ? raw.slice(1) : raw;
  for (let index = 0; index < text.length; index += 1) {
    const char = text[index];
    if (quoted) {
      if (char === '"') {
        if (text[index + 1] === '"') {
          field += '"';
          index += 1;
        } else {
          quoted = false;
        }
      } else {
        field += char;
      }
      continue;
    }
    if (char === '"') quoted = true;
    else if (char === ",") {
      row.push(field);
      field = "";
    } else if (char === "\n" || char === "\r") {
      if (char === "\r" && text[index + 1] === "\n") index += 1;
      row.push(field);
      field = "";
      if (row.some((value) => value !== "")) rows.push(row);
      row = [];
    } else field += char;
  }
  row.push(field);
  if (row.some((value) => value !== "")) rows.push(row);
  return rows;
}

function tokenValue(value: string | undefined): number | null {
  if (value === undefined) return null;
  const trimmed = value.trim();
  if (trimmed === "" || trimmed === "-") return 0;
  if (!/^\d+$/u.test(trimmed)) return null;
  return Number(trimmed);
}

function cleanModel(value: string | undefined): string | null {
  const model = (value ?? "").trim();
  if (!model || model.length > 160 || /[\s\\{}]/u.test(model)) return null;
  return model;
}

function warning(code: string, sourceFileId: string, severity: FileImportWarning["severity"], count?: number): FileImportWarning {
  return { code, message_code: code, severity, source_file_id: sourceFileId, ...(count !== undefined ? { details: { count } } : {}) };
}

export function isCursorUsageExportPath(fullPath: string): boolean {
  return /\/imports\/cursor\/[^/]+\.csv$/iu.test(fullPath.replace(/\\/gu, "/"));
}

export async function importCursorUsageExport(
  context: ImportContext,
  absolutePath: string
): Promise<{ counters: ImportCounters; warnings: FileImportWarning[] }> {
  const counters: ImportCounters = { sessions_upserted: 0, messages_upserted: 0, runtime_events_upserted: 0 };
  const warnings: FileImportWarning[] = [];
  let raw: string;
  try {
    const buffer = await readFile(absolutePath);
    if (buffer.byteLength > MAX_EXPORT_BYTES) throw new FileImportError("SOURCE_SCHEMA_UNSUPPORTED", "SOURCE_SCHEMA_UNSUPPORTED", [warning("CURSOR_USAGE_EXPORT_TOO_LARGE", context.sourceFileId, "error")]);
    raw = buffer.toString("utf8");
  } catch (error) {
    if (error instanceof FileImportError) throw error;
    throw new FileImportError("SOURCE_LOCKED_OR_UNREADABLE");
  }

  const [header, ...records] = parseCsv(raw);
  const columns = new Map((header ?? []).map((name, index) => [name.trim(), index]));
  if (!REQUIRED_COLUMNS.every((name) => columns.has(name))) {
    throw new FileImportError("SOURCE_SCHEMA_UNSUPPORTED", "SOURCE_SCHEMA_UNSUPPORTED", [warning("CURSOR_USAGE_EXPORT_COLUMNS_MISSING", context.sourceFileId, "error")]);
  }
  const cell = (record: string[], name: (typeof REQUIRED_COLUMNS)[number]): string | undefined => record[columns.get(name) as number];

  // A team export mixes users: the local machine cannot tell which rows are its own.
  const users = new Set(records.map((record) => (cell(record, "User") ?? "").trim().toLowerCase()).filter(Boolean));
  if (users.size > 1) {
    throw new FileImportError("SOURCE_SCHEMA_UNSUPPORTED", "SOURCE_SCHEMA_UNSUPPORTED", [warning("CURSOR_USAGE_EXPORT_MULTIPLE_USERS", context.sourceFileId, "error", users.size)]);
  }

  interface UsageEvent { at: string; ms: number; model: string | null; input: number; cacheWrite: number; cacheRead: number; output: number }
  const events: UsageEvent[] = [];
  let invalidRows = 0;
  for (const record of records) {
    const ms = Date.parse(cell(record, "Date") ?? "");
    const input = tokenValue(cell(record, "Input (w/o Cache Write)"));
    const cacheWrite = tokenValue(cell(record, "Input (w/ Cache Write)"));
    const cacheRead = tokenValue(cell(record, "Cache Read"));
    const output = tokenValue(cell(record, "Output Tokens"));
    if (!Number.isFinite(ms) || input === null || cacheWrite === null || cacheRead === null || output === null) {
      invalidRows += 1;
      continue;
    }
    if (input + cacheWrite + cacheRead + output === 0) continue;
    events.push({ at: new Date(ms).toISOString(), ms, model: cleanModel(cell(record, "Model")), input, cacheWrite, cacheRead, output });
  }
  if (invalidRows > 0) warnings.push(warning("CURSOR_USAGE_EXPORT_ROW_INVALID", context.sourceFileId, "warning", invalidRows));
  if (records.length > 0 && events.length === 0 && invalidRows === records.length) {
    throw new FileImportError("SOURCE_SCHEMA_UNSUPPORTED", "SOURCE_SCHEMA_UNSUPPORTED", warnings);
  }

  const unmatchedSessionId = hmacHex(`${context.machineId}:cursor:usage-export:${context.sourceFileId}`, context.hmacSalt);
  let matched = 0;
  let unmatched = 0;
  let duplicates = 0;
  await context.db.exec("BEGIN IMMEDIATE TRANSACTION");
  try {
    // Re-importing an export replaces its own projection only.
    await context.db.run("DELETE FROM message_metrics WHERE source_file_id = ?", [context.sourceFileId]);
    await context.db.run("DELETE FROM sessions WHERE id = ? AND source = 'cursor' AND source_file_id = ?", [unmatchedSessionId, context.sourceFileId]);
    const unmatchedEvents: UsageEvent[] = [];
    const insertUsage = async (sessionId: string, event: UsageEvent): Promise<void> => {
      const id = hmacHex(`${context.machineId}:cursor-usage:${event.at}|${event.model ?? ""}|${event.input}|${event.cacheWrite}|${event.cacheRead}|${event.output}`, context.hmacSalt);
      const result = await context.db.run(
        `INSERT OR IGNORE INTO message_metrics(
           id, session_id, source_file_id, seq, role, model, input_tokens, output_tokens, cache_read_tokens, cache_write_tokens,
           reasoning_tokens, token_available, partial_token_data, created_at, timestamp_source, metadata_json
         ) VALUES (?, ?, ?, ?, 'assistant', ?, ?, ?, ?, ?, NULL, 1, 0, ?, 'native', ?)`,
        [id, sessionId, context.sourceFileId, USAGE_SEQ_OFFSET + event.ms, event.model, event.input, event.output, event.cacheRead, event.cacheWrite, event.at, JSON.stringify({ metric_kind: "telemetry", origin: "cursor_usage_export" })]
      );
      if (result.changes > 0) counters.messages_upserted += 1;
      else duplicates += 1;
    };
    for (const event of events) {
      const candidates = await context.db.all<{ session_id: string; created_at: string }>(
        `SELECT mm.session_id, mm.created_at FROM message_metrics mm JOIN sessions s ON s.id = mm.session_id
          WHERE s.source = 'cursor' AND s.session_kind <> 'task' AND mm.created_at >= ? AND mm.created_at <= ?
            AND mm.metadata_json IS NOT '{"metric_kind":"telemetry","origin":"cursor_usage_export"}'`,
        [new Date(event.ms - MATCH_LOOKBACK_MS).toISOString(), new Date(event.ms + MATCH_LOOKAHEAD_MS).toISOString()]
      );
      let best: { session_id: string; distance: number } | null = null;
      for (const candidate of candidates) {
        const distance = Math.abs(event.ms - Date.parse(candidate.created_at));
        if (Number.isFinite(distance) && (!best || distance < best.distance)) best = { session_id: candidate.session_id, distance };
      }
      if (best) {
        matched += 1;
        await insertUsage(best.session_id, event);
      } else {
        unmatchedEvents.push(event);
      }
    }
    if (unmatchedEvents.length > 0) {
      unmatched = unmatchedEvents.length;
      const times = unmatchedEvents.map((event) => event.at).sort();
      await context.db.run(
        `INSERT INTO sessions(
           id, source, source_session_hash, source_file_id, machine_id, project_id, mode, model_primary, token_available,
           message_count, user_message_count, assistant_message_count, created_at, updated_at, created_at_source, client_surface, session_kind, metadata_json
         ) VALUES (?, 'cursor', ?, ?, ?, NULL, 'chat', NULL, 1, 0, 0, 0, ?, ?, 'native', 'desktop', 'task', ?)`,
        [unmatchedSessionId, hmacHex(`usage-export:${context.sourceFileId}`, context.hmacSalt), context.sourceFileId, context.machineId, times[0], times[times.length - 1], JSON.stringify({ classification_reason: "cursor_usage_export_unmatched" })]
      );
      counters.sessions_upserted += 1;
      for (const event of unmatchedEvents) await insertUsage(unmatchedSessionId, event);
      // Events already imported by an overlapping export leave this container empty.
      const kept = await context.db.get<{ n: number }>("SELECT COUNT(*) AS n FROM message_metrics WHERE session_id = ?", [unmatchedSessionId]);
      if ((kept?.n ?? 0) === 0) {
        await context.db.run("DELETE FROM sessions WHERE id = ?", [unmatchedSessionId]);
        counters.sessions_upserted -= 1;
      }
    }
    await context.db.exec("COMMIT");
  } catch (error) {
    try {
      await context.db.exec("ROLLBACK");
    } catch {
      // Preserve the original SQL/import error.
    }
    throw error;
  }
  warnings.push(warning("CURSOR_USAGE_EXPORT_EVENTS_MATCHED", context.sourceFileId, "info", matched));
  if (unmatched > 0) warnings.push(warning("CURSOR_USAGE_EXPORT_EVENTS_UNMATCHED", context.sourceFileId, "info", unmatched));
  if (duplicates > 0) warnings.push(warning("CURSOR_USAGE_EXPORT_DUPLICATE_EVENTS_IGNORED", context.sourceFileId, "info", duplicates));
  return { counters, warnings };
}
