import { afterEach, beforeEach, describe, expect, test } from "bun:test";
import { Database } from "bun:sqlite";
import { existsSync, mkdtempSync, rmSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { ClusterScheduler } from "./scheduler";
import { ControllerStatusSchema, buildControllerStatus } from "./status";
import { ControllerStore, type WorkspaceRecord } from "./store";
import { TransitionEngine } from "./transitions";

describe("ClusterScheduler", () => {
  let dir: string;
  let dbPath: string;
  let store: ControllerStore;
  let now: number;

  beforeEach(() => {
    dir = mkdtempSync(join(tmpdir(), "controller-scheduler-"));
    dbPath = join(dir, "state.sqlite");
    now = Date.UTC(2026, 6, 19);
    store = new ControllerStore(dbPath, { now: () => now });
  });

  afterEach(() => {
    store.close();
    rmSync(dir, { recursive: true, force: true });
  });

  test("dispatches durable FIFO to least-loaded eligible builders", () => {
    store.upsertHost({ hostname: "loaded", slotsTotal: 4, slotsUsed: 2 });
    store.upsertHost({ hostname: "idle", slotsTotal: 1, slotsUsed: 0 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("first", 101));
    scheduler.enqueue(ticket("second", 102));

    scheduler.reconcile();

    const tickets = store.listQueueTickets();
    expect(tickets[0]?.key).toBe("first");
    expect(tickets[0]?.placement?.host).toBe("idle");
    expect(tickets[1]?.placement?.host).toBe("loaded");
    expect(store.countHostSlotReservations("idle")).toBe(1);
  });

  test("requires available capability-passed non-quarantined builder", () => {
    store.upsertHost({ hostname: "draining", state: "draining", slotsTotal: 1 });
    store.upsertHost({ hostname: "failed-probe", capabilityOk: false, slotsTotal: 1 });
    store.upsertHost({ hostname: "quarantined", slotsTotal: 1 });
    store.setCapabilityBreaker({
      hostname: "quarantined",
      command: "build",
      state: "open",
      failureCount: 2,
      missingEventEmitted: true,
    });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("blocked", 103));

    scheduler.reconcile();

    expect(store.listQueueTickets()[0]?.state).toBe("queued");
    expect(store.countHostSlotReservations()).toBe(0);
  });

  test("atomically prevents concurrent per-host oversubscription", async () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("race-a", 201));
    scheduler.enqueue(ticket("race-b", 202));
    const [first, second] = store.listQueueTickets();
    const barrier = join(dir, "reservation-race.start");
    const storeUrl = new URL("./store.ts", import.meta.url).href;
    const runPlacement = (position: number, ready: string) => Bun.spawn({
      cmd: [
        process.execPath,
        "-e",
        `import { existsSync, writeFileSync } from "node:fs";
         import { ControllerStore } from ${JSON.stringify(storeUrl)};
         const store = new ControllerStore(${JSON.stringify(dbPath)});
         writeFileSync(${JSON.stringify(ready)}, "ready");
         while (!existsSync(${JSON.stringify(barrier)})) await Bun.sleep(1);
         console.log(store.tryReserveHostSlot(${position}, "debian1", Date.now()));
         store.close();`,
      ],
      stdout: "pipe",
      stderr: "pipe",
    });
    const ready = [join(dir, "race-a.ready"), join(dir, "race-b.ready")];
    const processes = [
      runPlacement(first!.position, ready[0]!),
      runPlacement(second!.position, ready[1]!),
    ];
    while (!ready.every(existsSync)) await Bun.sleep(1);
    writeFileSync(barrier, "go");
    const results = await Promise.all(processes.map(async (process) => ({
      exit: await process.exited,
      stdout: (await new Response(process.stdout).text()).trim(),
      stderr: (await new Response(process.stderr).text()).trim(),
    })));

    expect(results.map((result) => result.exit)).toEqual([0, 0]);
    expect(results.map((result) => result.stderr)).toEqual(["", ""]);
    expect(results.map((result) => result.stdout).sort()).toEqual(["false", "true"]);
    expect(store.countHostSlotReservations("debian1")).toBe(1);
  });

  test("reconciles terminal jobs after restart without leaking host slots", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("restart-job", 301));
    scheduler.reconcile();
    store.upsertJob({
      id: "restart-job",
      repo: "owner/repo",
      host: "debian1",
      snapshot: "abc",
      stage: "building",
      attempt: 1,
      rc: 0,
      infraFailure: false,
    });
    store.close();

    const crashed = new Database(dbPath);
    crashed.prepare("UPDATE jobs SET stage = 'completed' WHERE id = ?").run("restart-job");
    crashed.close();

    store = new ControllerStore(dbPath, { now: () => now });
    createScheduler(store, now).reconcile();

    expect(store.countHostSlotReservations()).toBe(0);
    expect(store.listQueueTickets()[0]?.state).toBe("completed");
  });

  test("restart requeues a placed ticket when no job ledger was created", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("dispatch-crash", 304));
    scheduler.reconcile();
    expect(store.countHostSlotReservations()).toBe(1);
    store.close();

    store = new ControllerStore(dbPath, { now: () => now });
    createScheduler(store, now, true);

    expect(store.countHostSlotReservations()).toBe(0);
    expect(store.listQueueTickets()[0]?.state).toBe("queued");
    expect(store.listQueueTickets()[0]?.placement).toBeUndefined();
  });

  test("ordinary store opens cannot erase an in-flight reservation", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("observer-race", 305));
    scheduler.reconcile();

    const observer = new ControllerStore(dbPath, { now: () => now });

    expect(observer.countHostSlotReservations()).toBe(1);
    expect(observer.listQueueTickets()[0]?.state).toBe("placed");
    observer.close();
  });

  test("releases the host slot atomically when a job becomes terminal", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("terminal-job", 303));
    scheduler.reconcile();

    store.upsertJob({
      id: "terminal-job",
      repo: "owner/repo",
      host: "debian1",
      snapshot: "abc",
      stage: "failed",
      attempt: 1,
      rc: 1,
      infraFailure: true,
    });

    expect(store.countHostSlotReservations()).toBe(0);
    expect(store.listQueueTickets()[0]?.state).toBe("failed");
  });

  test("releases host reservation exactly once without touching global limit", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    expect(store.createWorkspace(workspace("global-job"), 1)).toBe(true);
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("slot-job", 302));
    scheduler.reconcile();

    expect(scheduler.complete("slot-job", "failed")).toBe(true);
    expect(scheduler.complete("slot-job", "failed")).toBe(false);
    expect(store.countHostSlotReservations()).toBe(0);
    expect(store.countRemoteJobReservations()).toBe(1);
  });

  test("spills only when all builders are overloaded and queue exceeds builder count", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1, slotsUsed: 1 });
    store.upsertHost({ hostname: "debian2", slotsTotal: 1, slotsUsed: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("one", 401));
    scheduler.enqueue(ticket("two", 402));
    scheduler.reconcile();
    expect(store.getLease().active).toBe(false);

    scheduler.enqueue(ticket("three", 403));
    scheduler.reconcile();

    const spill = store.listQueueTickets().find((item) => item.placement?.kind === "spill");
    expect(spill?.key).toBe("one");
    expect(store.getLease()).toEqual({
      active: true,
      expiresAt: new Date(now + 5 * 60_000).toISOString(),
      host: "laptop",
      reason: "cluster-overloaded",
    });
  });

  test("terminal spill clears its lease", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1, slotsUsed: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("spill-terminal", 451));
    scheduler.enqueue(ticket("spill-waiting", 452));
    scheduler.reconcile();
    expect(store.getLease().active).toBe(true);

    store.upsertJob({
      id: "spill-terminal",
      repo: "owner/repo",
      host: "laptop",
      snapshot: "abc",
      stage: "completed",
      attempt: 1,
      rc: 0,
      infraFailure: false,
    });

    expect(store.getLease().active).toBe(false);
  });

  test("spill placement atomically rechecks overload predicate", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1, slotsUsed: 0 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("race-one", 461));
    scheduler.enqueue(ticket("race-two", 462));

    const placed = store.tryPlaceSpill(
      store.listQueueTickets()[0]!.position,
      now,
      new Date(now + 60_000).toISOString(),
    );

    expect(placed).toBe(false);
    expect(store.getLease().active).toBe(false);
  });

  test("newly enrolled builder receives new dispatch without a host allowlist", () => {
    store.upsertHost({ hostname: "builder-a", slotsTotal: 1, slotsUsed: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("waiting", 501));
    scheduler.reconcile();
    expect(store.listQueueTickets()[0]?.state).toBe("queued");

    store.upsertHost({ hostname: "builder-new", slotsTotal: 3 });
    scheduler.reconcile();

    expect(store.listQueueTickets()[0]?.placement?.host).toBe("builder-new");
  });

  test("reclaims dead tickets by PID starttime identity", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = new ClusterScheduler(store, {
      now: () => now,
      ownerAlive: (owner) => owner.pid === 600 && owner.starttime === 99,
    });
    scheduler.enqueue(ticket("pid-reused", 600, 98));

    scheduler.reconcile();

    expect(store.listQueueTickets()).toEqual([]);
    expect(store.countHostSlotReservations()).toBe(0);
  });

  test("default owner probe rejects a reused PID with mismatched starttime", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = new ClusterScheduler(store, { now: () => now });
    scheduler.enqueue(ticket("real-pid-reused", process.pid, Number.MAX_SAFE_INTEGER));

    scheduler.reconcile();

    expect(store.listQueueTickets()).toEqual([]);
  });

  test("transition verbs reconcile scheduler placements and recall spill state", async () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("transition-job", 701));
    const engine = new TransitionEngine(store, () => now, undefined, scheduler);

    const reconcile = await engine.handle("admission-reconcile", {
      expectedRevision: 0,
      idempotencyKey: "scheduler-reconcile",
      args: {},
    });
    expect(reconcile.status).toBe(200);
    expect(store.listQueueTickets()[0]?.state).toBe("placed");

    store.upsertHost({ hostname: "debian1", slotsUsed: 1 });
    scheduler.enqueue(ticket("spill-a", 702));
    scheduler.enqueue(ticket("spill-b", 703));
    scheduler.reconcile();
    expect(store.getLease().active).toBe(true);

    const recall = await engine.handle("recall-spill", {
      expectedRevision: 1,
      idempotencyKey: "scheduler-recall",
      args: { host: "laptop" },
    });
    expect(recall.status).toBe(200);
    expect(store.getLease().active).toBe(false);
    expect(store.listQueueTickets().some((item) => item.state === "recall-requested")).toBe(true);
  });

  test("projects scheduler state through the closed ControllerStatus shape", () => {
    store.upsertHost({ hostname: "debian1", slotsTotal: 1 });
    const scheduler = createScheduler(store, now);
    scheduler.enqueue(ticket("status-job", 801));
    scheduler.reconcile();

    const status = buildControllerStatus(store, undefined, () => now);

    expect(ControllerStatusSchema.parse(status)).toEqual(status);
    expect(() => ControllerStatusSchema.parse({ ...status, schedulerInternals: {} })).toThrow();
    expect(status.queue.tickets[0]?.dispatchTarget).toBe("debian1");
  });

  test("migration deterministically preserves the first legacy duplicate ticket", () => {
    store.close();
    const legacyPath = join(dir, "legacy.sqlite");
    const legacy = new Database(legacyPath);
    legacy.exec(`
      CREATE TABLE queue_tickets (
        position INTEGER PRIMARY KEY,
        ticket_key TEXT NOT NULL,
        repo TEXT NOT NULL,
        owner_pid INTEGER NOT NULL,
        owner_starttime INTEGER NOT NULL,
        owner_label TEXT,
        enqueue_age_seconds REAL NOT NULL DEFAULT 0,
        state TEXT NOT NULL DEFAULT 'queued',
        dispatch_target TEXT
      )
    `);
    legacy.prepare(
      `INSERT INTO queue_tickets
       (position, ticket_key, repo, owner_pid, owner_starttime, state)
       VALUES (?, 'duplicate', 'owner/repo', ?, ?, 'queued')`,
    ).run(1, 901, 9_010);
    legacy.prepare(
      `INSERT INTO queue_tickets
       (position, ticket_key, repo, owner_pid, owner_starttime, state)
       VALUES (?, 'duplicate', 'owner/repo', ?, ?, 'queued')`,
    ).run(2, 902, 9_020);
    legacy.close();

    store = new ControllerStore(legacyPath, { now: () => now });

    expect(store.listQueueTickets().map((item) => item.position)).toEqual([1]);
  });
});

function createScheduler(
  store: ControllerStore,
  now: number,
  recoverOrphans = false,
): ClusterScheduler {
  return new ClusterScheduler(store, {
    now: () => now,
    ownerAlive: () => true,
    recoverOrphans,
  });
}

function ticket(key: string, pid: number, starttime = pid * 10) {
  return {
    key,
    repo: "owner/repo",
    command: "build",
    owner: { pid, starttime },
  };
}

function workspace(jobId: string): WorkspaceRecord {
  const root = `/tmp/${jobId}`;
  return {
    jobId,
    repo: "owner/repo",
    host: "debian1",
    snapshot: "abc",
    checkoutGeneration: "gen-1",
    checkoutPath: `${root}/checkout`,
    publicationPath: `${root}/publish`,
    workspacePath: `${root}/workspace`,
    cachePath: `${root}/cache`,
    snapshotPath: `${root}/snapshot`,
    overlayPath: `${root}/overlay`,
    outputPath: `${root}/output`,
    stagingPath: `${root}/staging`,
    backupPath: `${root}/backup`,
    manifest: null,
    publicationState: "none",
    publicationReason: null,
    transportReattachCount: 0,
    completedAt: null,
    stage: "building",
  };
}
