import { basename } from "node:path";
import sqlite3 from "sqlite3";
import type { DatabaseAdapter } from "../db.js";
import { hmacHex } from "../privacy.js";
import { FileImportError, type FileImportWarning } from "./file-import-contract.js";

/**
 * Antigravity keeps one SQLite store per conversation (`conversations/<conversation-id>.db`) whose
 * `gen_metadata` rows are protobuf generation records. The layout is reverse-engineered, not an
 * official schema, so only one narrowly validated message is read:
 *   1.4.2 uncached input, 1.4.3 output, 1.4.5 cached input, 1.4.9 thinking, 1.4.10 visible output,
 *   1.9.4 {1: epoch seconds, 2: nanos} generation time, 1.21 model display name.
 * A record is accepted only when output == thinking + visible output (whenever both are present)
 * and a valid timestamp exists; anything else is skipped and counted, never guessed. Prompt or
 * response content is never decoded: only varints and the short model name are read.
 * Rows attach to the transcript session of the same conversation id (see antigravity.ts); if the
 * transcript is missing, a technical `task` session holds them until the transcript upserts it.
 */

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

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

const USAGE_SEQ_OFFSET = 10_000_000_000_000;
const CONVERSATION_DB_PATTERN = /\/conversations\/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})\.db$/iu;

export function isAntigravityConversationDbPath(fullPath: string): boolean {
  return CONVERSATION_DB_PATTERN.test(fullPath.replace(/\\/gu, "/"));
}

type ProtoValue = number | Buffer;

function readVarint(buffer: Buffer, offset: number): [number, number] {
  let value = 0;
  let multiplier = 1;
  for (let index = offset; index < buffer.length && index < offset + 10; index += 1) {
    const byte = buffer[index];
    value += (byte & 0x7f) * multiplier;
    if ((byte & 0x80) === 0) return [value, index + 1];
    multiplier *= 128;
  }
  throw new RangeError("varint");
}

/** Decodes one protobuf message level: field number -> values (varints and length-delimited). */
function decodeMessage(buffer: Buffer): Map<number, ProtoValue[]> {
  const fields = new Map<number, ProtoValue[]>();
  let offset = 0;
  while (offset < buffer.length) {
    const [key, afterKey] = readVarint(buffer, offset);
    offset = afterKey;
    const field = Math.floor(key / 8);
    const wireType = key % 8;
    let value: ProtoValue;
    if (wireType === 0) {
      [value, offset] = readVarint(buffer, offset);
    } else if (wireType === 2) {
      const [length, afterLength] = readVarint(buffer, offset);
      if (afterLength + length > buffer.length) throw new RangeError("length");
      value = buffer.subarray(afterLength, afterLength + length);
      offset = afterLength + length;
    } else if (wireType === 1 || wireType === 5) {
      offset += wireType === 1 ? 8 : 4;
      continue;
    } else {
      throw new RangeError("wire type");
    }
    const list = fields.get(field) ?? [];
    list.push(value);
    fields.set(field, list);
  }
  return fields;
}

const firstNumber = (fields: Map<number, ProtoValue[]> | null, field: number): number | undefined => {
  const value = fields?.get(field)?.[0];
  return typeof value === "number" && Number.isSafeInteger(value) ? value : undefined;
};
const firstMessage = (fields: Map<number, ProtoValue[]> | null, field: number): Map<number, ProtoValue[]> | null => {
  const value = fields?.get(field)?.[0];
  return Buffer.isBuffer(value) ? decodeMessage(value) : null;
};

export interface AntigravityUsageRecord { at: string; model: string | null; input: number; output: number; cacheRead: number; reasoning: number }

export function parseAntigravityGenerationMetadata(data: Buffer): AntigravityUsageRecord | null {
  try {
    const generation = firstMessage(decodeMessage(data), 1);
    const usage = firstMessage(generation, 4);
    const output = firstNumber(usage, 3);
    if (output === undefined) return null;
    const thinking = firstNumber(usage, 9);
    const visible = firstNumber(usage, 10);
    if (thinking !== undefined && visible !== undefined && thinking + visible !== output) return null;
    const seconds = firstNumber(firstMessage(firstMessage(generation, 9), 4), 1);
    if (seconds === undefined || seconds < 1_000_000_000 || seconds > 10_000_000_000) return null;
    const modelRaw = generation?.get(21)?.[0];
    const modelText = Buffer.isBuffer(modelRaw) ? modelRaw.toString("utf8").trim() : "";
    const model = /^[\w .()/+-]{1,80}$/u.test(modelText) ? modelText : null;
    return { at: new Date(seconds * 1000).toISOString(), model, input: firstNumber(usage, 2) ?? 0, output, cacheRead: firstNumber(usage, 5) ?? 0, reasoning: thinking ?? 0 };
  } catch {
    return null;
  }
}

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 async function importAntigravityConversationDb(
  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[] = [];
  const conversationId = basename(absolutePath).replace(/\.db$/iu, "");

  let rows: Array<{ idx: number; data: Buffer }>;
  try {
    rows = await new Promise((resolvePromise, reject) => {
      const handle = new sqlite3.Database(absolutePath, sqlite3.OPEN_READONLY, (openError) => {
        if (openError) {
          reject(new FileImportError("SOURCE_LOCKED_OR_UNREADABLE"));
          return;
        }
        handle.all("SELECT idx, data FROM gen_metadata ORDER BY idx", (queryError: Error | null, result: Array<{ idx: number; data: Buffer }>) => {
          handle.close(() => (queryError ? reject(new FileImportError("SOURCE_SCHEMA_UNSUPPORTED")) : resolvePromise(result ?? [])));
        });
      });
    });
  } catch (error) {
    throw error instanceof FileImportError ? error : new FileImportError("SOURCE_LOCKED_OR_UNREADABLE");
  }

  const records: Array<AntigravityUsageRecord & { idx: number }> = [];
  let unrecognized = 0;
  for (const row of rows) {
    const record = Buffer.isBuffer(row.data) && Number.isSafeInteger(row.idx) ? parseAntigravityGenerationMetadata(row.data) : null;
    if (!record) {
      unrecognized += 1;
      continue;
    }
    if (record.input + record.output + record.cacheRead > 0) records.push({ ...record, idx: row.idx });
  }
  if (unrecognized > 0) warnings.push(warning("ANTIGRAVITY_GENERATION_RECORD_UNRECOGNIZED", context.sourceFileId, "info", unrecognized));
  if (rows.length > 0 && records.length === 0 && unrecognized === rows.length) {
    throw new FileImportError("SOURCE_SCHEMA_UNSUPPORTED", "SOURCE_SCHEMA_UNSUPPORTED", warnings);
  }

  const sessionId = hmacHex(`${context.machineId}:antigravity:${conversationId}`, context.hmacSalt);
  await context.db.exec("BEGIN IMMEDIATE TRANSACTION");
  try {
    await context.db.run("DELETE FROM message_metrics WHERE session_id = ? AND source_file_id = ?", [sessionId, context.sourceFileId]);
    if (records.length > 0) {
      const times = records.map((record) => record.at).sort();
      // The transcript adapter upserts this row into the main conversation session.
      const inserted = await context.db.run(
        `INSERT OR IGNORE 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 (?, 'antigravity', ?, ?, ?, NULL, 'chat', NULL, 1, 0, 0, 0, ?, ?, 'native', 'code', 'task', ?)`,
        [sessionId, hmacHex(conversationId, context.hmacSalt), context.sourceFileId, context.machineId, times[0], times[times.length - 1], JSON.stringify({ classification_reason: "antigravity_usage_without_transcript" })]
      );
      if (inserted.changes > 0) counters.sessions_upserted += 1;
      for (const record of records) {
        await context.db.run(
          `INSERT OR REPLACE 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', ?, ?, ?, ?, 0, ?, 1, 0, ?, 'native', ?)`,
          [
            hmacHex(`${context.machineId}:antigravity-usage:${conversationId}:${record.idx}`, context.hmacSalt),
            sessionId, context.sourceFileId, USAGE_SEQ_OFFSET + record.idx, record.model,
            record.input, record.output, record.cacheRead, record.reasoning, record.at,
            JSON.stringify({ metric_kind: "telemetry", origin: "antigravity_generation_metadata" })
          ]
        );
        counters.messages_upserted += 1;
      }
      await context.db.run("UPDATE sessions SET token_available = 1 WHERE id = ?", [sessionId]);
    }
    await context.db.exec("COMMIT");
  } catch (error) {
    try {
      await context.db.exec("ROLLBACK");
    } catch {
      // Preserve the original SQL/import error.
    }
    throw error;
  }
  return { counters, warnings };
}
