/**
 * index.ts — oc-prefix-cache: an opencode plugin that reuses opencode's
 * stable harness prefix (tool schemas + the stable head of the system
 * prompt) in a llama-server KV-cache slot across sessions, with optional
 * disk persistence of that slot (OC_PREFIX_CACHE_PERSIST=1).
 *
 * This is a port of pi-prefix-cache (an equivalent pi extension). The
 * llama-server side (warmup.ts / cache-store.ts / fingerprint.ts) is
 * unchanged in behaviour. What differs is the integration point:
 *
 *   pi         → `before_provider_request` event with the raw payload
 *   opencode   → the plugin `config` hook injects a `fetch` wrapper into
 *                the provider options, so we observe the exact bytes
 *                opencode puts on the wire
 *
 * Why the fetch wrapper rather than the `chat.params` /
 * `experimental.chat.system.transform` hooks:
 *   - `chat.params` never sees the message list or the tools;
 *   - `system.transform` gives the system string but the tools array is
 *     only reachable per-tool via `tool.definition` (zod schemas that
 *     would have to be re-serialized — a silent mismatch there is the
 *     exact failure that cost pi bug #1, "tools omitted from the warm
 *     request → lcp = 3");
 *   - the wrapper sees the final body after opencode's own middleware
 *     (`transformParams` message rewrite runs later still), and it can
 *     delay the real request while the slot is warmed — which is what
 *     makes the whole scheme work.
 *
 * The provider gate is inherent: the wrapper is only installed on
 * providers whose `options.baseURL` equals OC_PREFIX_CACHE_BASE_URL, so
 * we can only ever warm our own server. Sub-requests that are not the
 * agent turn (e.g. the title generator: no tools, ~2k chars) are
 * skipped by the "has tools + has system message" test.
 *
 * Flow per agent request:
 *   1. Parse the body; skip unless it has a tools array and a system
 *      message.
 *   2. Cut the system message at the first volatile marker (prefix.ts)
 *      → stable head.
 *   3. Fingerprint (model, gguf identity, server id, opencode version,
 *      plugin version, tools[], stable head).
 *   4. In-process HIT for that fingerprint? → verify the slot is still
 *      alive (a server restart zeroes the prompt counter) and forward.
 *   5. MISS → restore `oc-prefix-<fp16>.bin` (+ .ckpt sidecar) and
 *      verify in two stages (identical-prompt gate ≥ 90% from cache,
 *      then a divergent probe that must resume from a context
 *      checkpoint); invalid → erase + cold rebuild. No file → cold
 *      warm-up. With persistence on, save the fresh prefix.
 *   6. Forward the untouched request, then poll the slot counters for
 *      the observed reuse of the real request.
 *
 * Fail-open, always: any error is logged with the `[prefix-cache]`
 * prefix and the real request proceeds untouched. The body is never
 * modified.
 *
 * Configuration (env, read once at plugin load):
 *   OC_PREFIX_CACHE_BASE_URL     llama-server base URL, e.g.
 *                                http://127.0.0.1:8080/v1
 *                                (unset → plugin disables itself)
 *   OC_PREFIX_CACHE_SLOT         slot id for the shared prefix (0)
 *   OC_PREFIX_CACHE_PERSIST      "1" → enable disk save/restore
 *   OC_PREFIX_CACHE_SLOT_DIR     server's --slot-save-path; enables
 *                                stale-checkpoint cleanup
 *   OC_PREFIX_CACHE_CLEANUP      "0" disables cleanup (default on)
 *   OC_PREFIX_CACHE_TTL_DAYS     keep checkpoints used within N days (30)
 *   OC_PREFIX_CACHE_KEEP_RECENT  always keep N most recent (3)
 *   OC_PREFIX_CACHE_MAX_GB       max total checkpoint storage, 0 = off
 *   OC_PREFIX_CACHE_SERVER_ID    llama-server build identity (bump on
 *                                rebuild); forces one clean rebuild
 *   OC_PREFIX_CACHE_DIR          metadata index dir
 *                                (default ~/.cache/opencode/llama-prefix)
 *   OC_PREFIX_CACHE_LOG          "1" → also write the log lines to
 *                                /tmp/oc-prefix-cache.log (opencode's
 *                                TUI swallows plugin stdout)
 */

import { execFileSync } from "node:child_process";
import { appendFileSync, statSync } from "node:fs";
import { extractStablePrefix, findSystemMessage } from "./prefix";
import { computeFingerprint } from "./fingerprint";
import { eraseSlotFile, getSlotStats, probeSlot, restoreSlotFile, saveSlotFile, warmSlot } from "./warmup";
import { cleanupStale, recordUse } from "./cache-store";

const EXTENSION_VERSION = "1.0.0";
const LOG_FILE = "/tmp/oc-prefix-cache.log";

/** Integer env var with fallback to a default when unset/invalid. */
function numEnv(name: string, fallback: number): number {
  const n = Number.parseInt(process.env[name] ?? "", 10);
  return Number.isFinite(n) ? n : fallback;
}

const BASE_URL = (process.env.OC_PREFIX_CACHE_BASE_URL ?? "").replace(/\/+$/, "");
const SLOT = numEnv("OC_PREFIX_CACHE_SLOT", 0);
const PERSIST = (process.env.OC_PREFIX_CACHE_PERSIST ?? "") === "1";
const SLOT_DIR = process.env.OC_PREFIX_CACHE_SLOT_DIR ?? "";
const CLEANUP = (process.env.OC_PREFIX_CACHE_CLEANUP ?? "1") === "1";
const TTL_DAYS = numEnv("OC_PREFIX_CACHE_TTL_DAYS", 30);
const KEEP_RECENT = numEnv("OC_PREFIX_CACHE_KEEP_RECENT", 3);
const MAX_GB = numEnv("OC_PREFIX_CACHE_MAX_GB", 0);
const SERVER_ID = process.env.OC_PREFIX_CACHE_SERVER_ID ?? "";
/** Also mirror the log to a file: opencode's TUI owns stdout, so plugin
 *  console output is otherwise invisible during interactive use. */
const LOG_TO_FILE = (process.env.OC_PREFIX_CACHE_LOG ?? "1") === "1";
/** Verify a restored checkpoint only if ≥ this fraction of the warm-up
 *  prompt was served from cache (identical-prompt reuse is unbounded, so
 *  a valid restore must score near 1.0). */
const VERIFY_THRESHOLD = 0.9;

const log = (line: string) => {
  const out = `[prefix-cache] ${line}`;
  try {
    console.log(out);
  } catch {
    /* stdout may be closed in server mode */
  }
  if (LOG_TO_FILE) {
    try {
      appendFileSync(LOG_FILE, out + "\n");
    } catch {
      /* best effort */
    }
  }
};

const norm = (url: string) => url.replace(/\/+$/, "");

/** Physical model identity for the GGUF behind this server. opencode
 *  sends the short model id ("qwen"), which carries no file identity, so
 *  the path comes from the env (OC_PREFIX_CACHE_MODEL_PATH) when set;
 *  otherwise the model id alone is the identity. */
function modelIdentity(model: string): string {
  const path = process.env.OC_PREFIX_CACHE_MODEL_PATH ?? "";
  if (path.startsWith("/")) {
    try {
      const st = statSync(path);
      return `${path}|size=${st.size}|mtime=${st.mtimeMs}`;
    } catch {
      return path;
    }
  }
  return model;
}

let opencodeVersion = "unknown";
function resolveOpencodeVersion(): string {
  try {
    return execFileSync("opencode", ["--version"], {
      encoding: "utf8",
      timeout: 10_000,
      stdio: ["ignore", "pipe", "ignore"],
    }).trim();
  } catch {
    return "unknown";
  }
}

interface WireBody {
  model?: string;
  messages?: Array<{ role?: string; content?: unknown }>;
  tools?: unknown;
  chat_template_kwargs?: unknown;
}

export default async function ocPrefixCache(): Promise<Record<string, unknown>> {
  opencodeVersion = resolveOpencodeVersion();

  if (!BASE_URL) {
    log("disabled (OC_PREFIX_CACHE_BASE_URL unset) — plugin inert, no interception");
    return {};
  }

  log(`plugin loaded ext=${EXTENSION_VERSION} opencode=${opencodeVersion} base=${BASE_URL} slot=${SLOT} persist=${PERSIST}`);

  /** Fingerprint warmed in this process; null = nothing warmed yet. */
  let warmedFingerprint: string | null = null;
  /** Stale-cache cleanup runs at most once per process. */
  let cleaned = false;
  /** Serialize warm/restore work: opencode can run subagents and
   *  background agents concurrently inside one process, and two
   *  cold builds on one slot would just thrash it. */
  let chain: Promise<unknown> = Promise.resolve();
  /** Diagnostics after the real request: slot counters reflect the LAST
   *  prompt, so we only report when they moved past our own warm-up. */
  let warmPromptTokens = -1;

  const ensureWarm = (body: WireBody, fp: string, stable: string, sysRole: string): Promise<void> => {
    const run = async (): Promise<void> => {
      const model = body.model ?? "";
      const tools = body.tools;
      const meta = {
        fileName: `oc-prefix-${fp.slice(0, 16)}.bin`,
        model,
        baseUrl: BASE_URL,
        slot: SLOT,
        opencodeVersion,
        extVersion: EXTENSION_VERSION,
      };
      const warm = (userContent = "") =>
        warmSlot(BASE_URL, model, sysRole, stable, body.chat_template_kwargs ?? {}, SLOT, userContent, tools);

      let restored = false;

      if (PERSIST) {
        const r = await restoreSlotFile(BASE_URL, SLOT, meta.fileName);
        if (r.ok && (r.n_restored ?? 0) > 0) {
          log(
            `RESTORE ${meta.fileName}: n_restored=${r.n_restored} ` +
              `(${r.timings?.restore_ms?.toFixed(1) ?? "?"} ms)`,
          );
          // (b1) identical-prompt gate: raw KV must be present.
          // NOTE: llama-server's timings.prompt_n counts only the tokens
          // PROCESSED by this request; prompt_n + cache_n is the total
          // prompt. Ratios must use the total, otherwise a fully cached
          // verification reads as prompt_n=0 and every gate passes.
          const v = await warm();
          const vt = v.timings ?? {};
          const vTotal = (vt.prompt_n ?? 0) + (vt.cache_n ?? 0);
          const vc = vt.cache_n ?? 0;
          if (vTotal > 0 && vc >= VERIFY_THRESHOLD * vTotal) {
            // (b2) divergent probe: an identical-prompt hit is a plain
            // KV token match and does not prove the restored checkpoint
            // chain can resume before the divergence point.
            const p = await probeSlot(
              BASE_URL,
              model,
              sysRole,
              stable,
              body.chat_template_kwargs ?? {},
              SLOT,
              tools,
            );
            const pt = p.timings ?? {};
            const pn = (pt.prompt_n ?? 0) + (pt.cache_n ?? 0);
            const pc = pt.cache_n ?? 0;
            if (pn > 0 && pc > 0) {
              log(`RESTORE verified: identical ${vc}/${vTotal}, divergent probe ${pc}/${pn}`);
            } else {
              log(
                `RESTORE probe FAILED (divergent cache_n=${pc}/${pn}): no checkpoint resume — ` +
                  `check --checkpoint-min-step/-ub/--ctx-checkpoints; proceeding (full re-prefill expected)`,
              );
            }
            restored = true;
            try {
              recordUse(fp, { ...meta, tokenCount: vTotal });
            } catch (err) {
              log(`metadata update failed (fail-open): ${String(err)}`);
            }
          } else {
            log(`RESTORE INVALID (verify cache_n=${vc}/${vTotal}) → erase + rebuild`);
            const e = await eraseSlotFile(BASE_URL, SLOT);
            if (e.ok) log(`slot erased (n_erased=${e.n_erased})`);
          }
        } else {
          log(`no checkpoint (restore HTTP ${r.status}${r.error ? `: ${r.error}` : ""}) → cold build`);
        }
      }

      if (!restored) {
        log(`cache=MISS → warming slot=${SLOT}${PERSIST ? " (cold)" : ""}`);
        const before = await getSlotStats(BASE_URL, SLOT);
        if (before) log(`slot before: tokens=${before.n_prompt_tokens}`);
        const t0 = Date.now();
        const res = await warm();
        const t = res.timings ?? {};
        const elapsed = Date.now() - t0;
        const total = (t.prompt_n ?? 0) + (t.cache_n ?? 0);
        const rate = total > 0 ? ((1000 * total) / elapsed).toFixed(0) : "?";
        log(
          `warm-up done: HTTP ${res.status} total=${total} new=${t.prompt_n ?? "?"} from_cache=${t.cache_n ?? "?"} ` +
            `(${(elapsed / 1000).toFixed(1)}s, ${rate} tok/s)`,
        );
        if (PERSIST) {
          const s = await saveSlotFile(BASE_URL, SLOT, meta.fileName);
          if (s.ok) {
            log(`SAVED ${meta.fileName}: n_saved=${s.n_saved} (${s.timings?.save_ms?.toFixed(1) ?? "?"} ms)`);
            try {
              recordUse(fp, { ...meta, tokenCount: (t.prompt_n ?? 0) + (t.cache_n ?? 0) });
            } catch (err) {
              log(`metadata update failed (fail-open): ${String(err)}`);
            }
          } else {
            log(`SAVE FAILED: HTTP ${s.status} ${s.error ?? ""}`);
          }
        }
      }

      try {
        const s = await getSlotStats(BASE_URL, SLOT);
        warmPromptTokens = s?.n_prompt_tokens ?? -1;
      } catch {
        /* diagnostics only */
      }
      warmedFingerprint = fp;
    };

    chain = chain.then(run, run);
    return chain as Promise<void>;
  };

  /** One-time stale-checkpoint cleanup (fail-open). */
  const maybeCleanup = (fp: string) => {
    if (cleaned || !PERSIST || !CLEANUP || !SLOT_DIR) return;
    cleaned = true;
    try {
      const c = cleanupStale(SLOT_DIR, [fp], {
        ttlDays: TTL_DAYS,
        keepRecent: KEEP_RECENT,
        maxBytes: MAX_GB > 0 ? MAX_GB * 1e9 : undefined,
      });
      if (c.skipped) log(`cleanup SKIPPED (fail-safe, index untrusted): ${c.skipped}`);
      else if (c.deleted.length > 0) {
        log(
          `cleanup: deleted ${c.deleted.length} stale checkpoint(s), freed ${(c.totalBytesDeleted / 1e9).toFixed(2)} GB`,
        );
      }
    } catch (err) {
      log(`cleanup failed (fail-open): ${String(err)}`);
    }
  };

  /**
   * Install the interception on every configured provider that points at
   * BASE_URL. The wrapper is fail-open: on any internal error it simply
   * forwards the original request.
   */
  const hookProvider = (options: Record<string, unknown>, label: string, original: typeof fetch) => {
    const base = typeof options.baseURL === "string" ? norm(options.baseURL) : "";
    if (!base || base !== BASE_URL) return false;

    const wrapper: typeof fetch = async (input, init) => {
      let body: WireBody | null = null;
      let raw: string | null = null;
      try {
        const url = typeof input === "string" ? input : input instanceof URL ? input.href : input.url;
        if (!/\/chat\/completions(\?|$)/.test(url)) return original(input, init);
        if (typeof init?.body !== "string") return original(input, init);
        raw = init.body;
        body = JSON.parse(raw) as WireBody;
      } catch {
        return original(input, init);
      }

      const messages = body.messages ?? [];
      const tools = body.tools;
      // Agent turns only: the title generator and other sub-requests
      // carry no tools and a tiny system prompt.
      if (!Array.isArray(tools) || tools.length === 0) return original(input, init);
      const sys = findSystemMessage(messages);
      if (!sys) return original(input, init);

      const sp = extractStablePrefix(sys.content);
      if (!sp.ok) {
        log(`skip (fail-open): ${sp.reason}`);
        return original(input, init);
      }

      const fp = computeFingerprint({
        model: body.model ?? "",
        modelIdentity: modelIdentity(body.model ?? ""),
        serverId: SERVER_ID,
        opencodeVersion,
        extensionVersion: EXTENSION_VERSION,
        stablePrefix: sp.text,
        tools,
        chatTemplateKwargs: body.chat_template_kwargs ?? {},
      });

      log(
        `request: provider=${label} model=${body.model} boundary=${sp.boundary}@${sp.boundaryIndex} ` +
          `stable-chars=${sp.text.length} tools=${tools.length} fingerprint=${fp.slice(0, 16)}`,
      );
      maybeCleanup(fp);

      try {
        if (warmedFingerprint === fp) {
          // A server restart clears the slot while we still believe it
          // is warm; a fresh slot has a zeroed prompt counter.
          let live = true;
          try {
            const s = await getSlotStats(BASE_URL, SLOT);
            live = !!s && s.n_prompt_tokens > 0;
          } catch {
            /* keep in-memory state */
          }
          if (live) {
            log("cache=HIT (slot already warm for this fingerprint)");
          } else {
            log("warm state stale (slot empty — server restarted?) → restore/re-warm");
            warmedFingerprint = null;
            await ensureWarm(body, fp, sp.text, sys.role);
          }
        } else {
          await ensureWarm(body, fp, sp.text, sys.role);
        }
      } catch (err) {
        log(`warm/restore failed (fail-open, request proceeds): ${String(err)}`);
      }

      // Forward the request untouched.
      const res = await original(input, init);

      // Diagnostics: the slot's counters now describe the real request.
      if (warmPromptTokens >= 0) {
        try {
          const s = await getSlotStats(BASE_URL, SLOT);
          const pn = s?.n_prompt_tokens ?? -1;
          const cn = s?.n_prompt_tokens_cache ?? 0;
          if (pn > warmPromptTokens) {
            log(
              `observed real request: prompt_n=${pn} cache_n=${cn} ` +
                `(${pn > 0 ? Math.round((cn / pn) * 100) : 0}% from cache)`,
            );
          }
        } catch {
          /* diagnostics only */
        }
      }
      return res;
    };

    options.fetch = wrapper;
    log(`interception armed on provider "${label}" (baseURL ${base})`);
    return true;
  };

  return {
    config: async (cfg: Record<string, any>) => {
      const providers = cfg?.provider ?? {};
      let armed = 0;
      for (const [id, prov] of Object.entries(providers)) {
        if (!prov || typeof prov !== "object") continue;
        const options = (prov as Record<string, any>).options;
        if (!options || typeof options !== "object") continue;
        const original = (options.fetch as typeof fetch) ?? globalThis.fetch;
        if (hookProvider(options as Record<string, unknown>, id, original)) armed++;
      }
      if (armed === 0) {
        log(
          `no provider matched OC_PREFIX_CACHE_BASE_URL=${BASE_URL} — ` +
            `check provider.${Object.keys(providers).join(",") || "<none>"}.options.baseURL`,
        );
      }
    },
  };
}
