/**
 * memgraph-autostart -- Auto-launch Memgraph on session start.
 *
 * On session_start hook:
 * 1. Check if Bolt port 7687 is reachable (fast TCP probe)
 * 2. If yes → done (no-op)
 * 3. If no → check if Docker is available
 * 4. If Docker available → start/create Memgraph container
 * 5. If Docker not available → warn and continue
 *
 * All work is async and non-blocking -- session starts immediately.
 */
import { execFile } from 'child_process';
import { promisify } from 'util';
import * as net from 'net';
import { join } from 'path';
import { homedir } from 'node:os';
import bus from '../lib/event-bus.mts';
import { registerStartupPhase, StartupPhase, getStartupState } from '../lib/startup-lifecycle.ts';
import { resolve as di, updateScopedValue } from '../lib/container.ts';
import { registerService } from '../lib/infrastructure-registry.ts';
import { createRequire } from 'module';
const require = createRequire(import.meta.url);
const QUIET = process.env.PI_QUIET === '1' || process.env.PI_LOG_LEVEL === 'error';

// Graph availability singleton — updated whenever Memgraph comes up or goes down
let _graphAvailability: { setMemgraphAvailable: (v: boolean) => void } | null = null;
function getGraphAvailability() {
  if (!_graphAvailability) {
    try {
      _graphAvailability = require('../lib/graph-availability');
    } catch (_) {}
  }
  return _graphAvailability;
}

// ── Event emission helpers ───────────────────────────────────────────────────

/** Load MAGE query modules after Bolt becomes reachable.
 *  Runs mg.load_all() which loads all compiled C++ and Python MAGE modules.
 *  Non-blocking, fire-and-forget — failure is logged but doesn't block startup.
 *  Uses a direct bolt connection (not the broker) since this runs early in startup.
 */
async function _loadMageModules(log: (msg: string) => void): Promise<void> {
    try {
        const { ensureMageLoaded, resetVerification } = require('../lib/mage-loader.js');
        resetVerification(); // Force re-check after restart/recovery
        const neo4j = require('neo4j-driver');
        const safeMemgraph = require('../lib/safe-memgraph.js');
        // AR-H16: Guard against null driver — broker may not be ready yet at this call site.
        // getDriver() can return null before the connection broker has finished initialization.
        const driver = safeMemgraph.getDriver?.();
        if (!driver) {
            process.stderr.write('[memgraph-autostart] ⚠️ Driver not available yet — skipping MAGE module load (will retry on next tick)\n');
            return;
        }
        const session = driver.session();
        try {
            const ok = await ensureMageLoaded(session);
            if (ok) log('[memgraph-autostart] ✅ MAGE modules loaded and verified');
            else log('[memgraph-autostart] ⚠️ Some MAGE procedures missing — check mage-loader.js');

            // Ensure vector indexes (skill_def_emb, tool_def_emb, etc.) exist.
             // These are created idempotently — no-op if already present.
             // Required for tool-selector.ts Path C and neural-search-gate.
             // VECTOR INDEX GUARD: two independent flags control two independent concerns.
             // HELIOS_SKIP_VECTOR_INDEX=1 → skip HNSW index creation entirely (explicit opt-out).
             // HELIOS_SKIP_INDEX=1 or (PI_AGENT_NAME=worker + HELIOS_SELF_ID) → eval container:
             //   skip codebase re-indexing (20s+ per batch, burns eval budget).
             //   Eval containers connect to production Memgraph which already has indexes; the
             //   backfill of 800+ Episode nodes is not needed for eval correctness.
             // Default (neither set): create indexes AND run backfill.
             const _skipVectorIndex = process.env.HELIOS_SKIP_VECTOR_INDEX === '1';
             const _skipEmbeddingBackfill = process.env.HELIOS_SKIP_INDEX === '1' ||
               !!(process.env.PI_AGENT_NAME === 'worker' && process.env.HELIOS_SELF_ID);
             // Propagate to child modules via the canonical flag ensure-mage checks
             if (_skipEmbeddingBackfill) process.env.HELIOS_SKIP_EMBEDDING_BACKFILL = '1';
             if (!_skipVectorIndex) {
               try {
                 const { ensureMage } = require('../lib/ensure-mage.ts');
                 await ensureMage(); // creates missing HNSW indexes + backfills embeddings (unless skipped)
                 log('[memgraph-autostart] vector search indexes ensured' + (_skipEmbeddingBackfill ? ' (embedding backfill skipped — eval container)' : ''));
               } catch (idxErr: any) {
                 log('[memgraph-autostart] ⚠️ vector index ensure failed (non-fatal): ' + idxErr?.message?.slice(0, 80));
               }
             } else {
               log('[memgraph-autostart] vector search indexes skipped (HELIOS_SKIP_VECTOR_INDEX=1)');
             }

            log('[memgraph-autostart] text search indexes ensured');
        } finally {
            await session.close();
            // B-26: Do NOT call driver.close() — the shared driver from safe-memgraph.getDriver()
            // must not be closed by consumers. Only safe-memgraph.shutdown() should close it.
        }
    } catch (err: any) {
        log('[memgraph-autostart] ⚠️ MAGE loading failed: ' + (err?.message || err));
    }
}

function emitMemgraphReady(log: (msg: string) => void): void {
    // Dedup: globalThis guard survives jiti moduleCache:false reloads
    if ((globalThis as any).__memgraph_ready_emitted) return;
    (globalThis as any).__memgraph_ready_emitted = true;

    // Load MAGE query modules (vector_search, text_search, pagerank, etc.)
    // These don't persist across Memgraph restarts, so we load them every time
    // Bolt becomes reachable. Docker's --init-data-file also loads them, but
    // this covers recovery scenarios where the container didn't restart.
    _loadMageModules(log).catch(e => log('[memgraph-autostart] MAGE load failed (non-fatal): ' + (e?.message || e)));

    updateScopedValue('memgraphReadyEmitted', true);
  registerService('memgraph', 'up', { host: BOLT_HOST, port: BOLT_PORT });
  const payload = { timestamp: Date.now(), host: BOLT_HOST, port: BOLT_PORT };
  updateScopedValue('memgraphReady', true);
  updateScopedValue('memgraphReadyAt', payload.timestamp);
  setImmediate(() => { (di('mesh') as any)?.bus?.publish?.('memgraph_ready', payload); });
  try {
    // Also emit on the typed internal bus
    bus.emit('memgraph_ready', payload);
  } catch (e) { process.stderr.write(`[extensions] non-fatal if bus unavailable: ${String(e)}\n`); }
  log(`[memgraph-autostart] Memgraph ready — bolt reachable at ${BOLT_HOST}:${BOLT_PORT}`);
  getGraphAvailability()?.setMemgraphAvailable(true);
}

/**
 * Ensure critical indexes exist after every Memgraph start/recovery.
 * Indexes are in-memory and may not survive unclean restarts.
 * Idempotent: CREATE INDEX IF NOT EXISTS is a no-op if index already exists in Memgraph.
 */
async function ensureIndexes(log: (msg: string) => void): Promise<void> {
  // Memgraph syntax: CREATE INDEX FOR (n:Label) ON (n.prop) — no IF NOT EXISTS clause.
  // Memgraph ignores duplicate CREATE INDEX calls, so this is already idempotent.
  const INDEXES = [
    'CREATE INDEX FOR (n:Person) ON (n.id)',
    'CREATE INDEX FOR (n:KeyFact) ON (n.personId)',
    'CREATE INDEX FOR (n:Organization) ON (n.id)',
    'CREATE INDEX FOR (n:Episode) ON (n.id)',
    'CREATE INDEX FOR (n:Entity) ON (n.id)',
    // Critical production indexes — lost on Memgraph crash-restart.
    // Without these, skill routing, BrainV2, heartbeat, and plan queries
    // do full-table scans until each extension re-initializes.
    'CREATE INDEX FOR (n:SkillDef) ON (n.name)',
    'CREATE INDEX FOR (n:Extension) ON (n.name)',
    'CREATE INDEX FOR (n:Extension) ON (n.health)',
    'CREATE INDEX FOR (n:BusinessAgent) ON (n.id)',
    'CREATE INDEX FOR (n:BusinessAgent) ON (n.status)',
    'CREATE INDEX FOR (n:MissionRun) ON (n.id)',
    'CREATE INDEX FOR (n:Plan) ON (n.id)',
    'CREATE INDEX FOR (n:SessionOutcome) ON (n.session_id)',
    'CREATE INDEX FOR (n:CortexDecisionLog) ON (n.timestamp)',
    'CREATE INDEX FOR (n:CausalLesson) ON (n.timestamp)',
    'CREATE INDEX FOR (n:OrchestratorTip) ON (n.createdAt)',
    'CREATE INDEX FOR (n:SessionEpisode) ON (n.duplicate)',
    'CREATE INDEX FOR (n:GovernanceEvent) ON (n.timestamp)',
    // MissionLoop durability indexes — required before Phase 1 loop driving
    'CREATE INDEX FOR (n:LoopCheckpoint) ON (n.missionRunId)',
    'CREATE INDEX FOR (n:LoopCheckpoint) ON (n.sessionId)',

    // ── Mental Model / Signal indexes ────────────────────────────────────────
    // These are also in schema-v2.ts (applied during backfill) but must be here
    // too so they survive Memgraph crash-restart and are present on every session
    // start — not only after a backfill run. Without these, signal scoring does
    // full-label scans on every triage call.

    // Person graph core
    'CREATE INDEX FOR (n:Person) ON (n.dunbarLayer)',

    // Episode temporal + dedup
    'CREATE INDEX FOR (n:Episode) ON (n.receivedAt)',
    'CREATE INDEX FOR (n:Episode) ON (n.rawId)',
    'CREATE INDEX FOR (n:Episode) ON (n.platform)',

    // Identity resolution
    'CREATE INDEX FOR (n:Identity) ON (n.id)',
    'CREATE INDEX FOR (n:Identity) ON (n.handle)',
    'CREATE INDEX FOR (n:Identity) ON (n.platform)',

    // Topic lifecycle
    'CREATE INDEX FOR (n:Topic) ON (n.status)',
    'CREATE INDEX FOR (n:Topic) ON (n.name)',

    // Question, Commitment, Conversation
    'CREATE INDEX FOR (n:Question) ON (n.status)',
    'CREATE INDEX FOR (n:Commitment) ON (n.status)',
    'CREATE INDEX FOR (n:Commitment) ON (n.dueDate)',
    'CREATE INDEX FOR (n:Conversation) ON (n.platform)',
    'CREATE INDEX FOR (n:Conversation) ON (n.threadId)',

    // GoalNode — goal_relevance signal
    'CREATE INDEX FOR (n:GoalNode) ON (n.status)',

    // FAVEESnapshot — quality scoring history
    'CREATE INDEX FOR (n:FAVEESnapshot) ON (n.personId)',
    'CREATE INDEX FOR (n:FAVEESnapshot) ON (n.week)',

    // KeyFact enrichment
    'CREATE INDEX FOR (n:KeyFact) ON (n.category)',
    'CREATE INDEX FOR (n:KeyFact) ON (n.invalidAt)',

    // Strength/trajectory history
    'CREATE INDEX FOR (n:StrengthSnapshot) ON (n.personId)',
    'CREATE INDEX FOR (n:StrengthSnapshot) ON (n.snapshotAt)',
    'CREATE INDEX FOR (n:LifeEvent) ON (n.personId)',
    'CREATE INDEX FOR (n:LifeEvent) ON (n.date)',
  ];
  try {
    // B-27: Use safe-memgraph.rawWrite() instead of direct neo4j.driver() — the direct
    // driver bypasses the broker and leaves connections open after close() isn't called.
    const safeMemgraph = require('../lib/safe-memgraph.js');
    let created = 0;
    for (const query of INDEXES) {
      try {
        await safeMemgraph.rawWrite(query, {}); // swallow "already exists"
        created++;
      } catch (e: any) {
        if (!e.message?.includes('already exists')) {
          log(`[memgraph-autostart] Index warning: ${query} → ${e.message?.slice(0, 60)}`);
        }
      }
    }
    if (created > 0) {
      log(`[memgraph-autostart] ✅ Ensured ${created} indexes`);
      // Trigger snapshot to persist indexes across restarts
      try {
        await safeMemgraph.rawWrite('CREATE SNAPSHOT', {});
        log('[memgraph-autostart] ✅ Snapshot taken (indexes persisted)');
      } catch (e: any) {
        log(`[memgraph-autostart] ⚠️ Snapshot failed: ${e.message?.slice(0, 60)}`);
      }
    }
  } catch (e: any) {
    log(`[memgraph-autostart] ⚠️ ensureIndexes failed: ${e.message?.slice(0, 60)}`);
  }
}

function emitMemgraphDegraded(log: (msg: string) => void, reason: string): void {
  // R3-H47: Reset ready-emitted flag so emitMemgraphReady() can fire again on recovery
  (globalThis as any).__memgraph_ready_emitted = false;
  const payload = { timestamp: Date.now(), reason };
  updateScopedValue('memgraphReady', false);
  registerService('memgraph', 'degraded', { reason });
  setImmediate(() => { (di('mesh') as any)?.bus?.publish?.('memgraph_degraded', payload); });
  try {
    bus.emit('memgraph_degraded', payload);
  } catch (e) { process.stderr.write(`[extensions] non-fatal: ${String(e)}\n`); }
  getGraphAvailability()?.setMemgraphAvailable(false);
}

/** Returns true if Memgraph bolt was confirmed reachable during this session. */
export function isMemgraphReady(): boolean {
  return !!(di('helios_memgraph_ready') ?? di('memgraphReady'));
}


const execFileAsync = promisify(execFile);
const CONTAINER_NAME = 'helios-memgraph';
const CANONICAL_CONTAINERS = new Set(['helios-memgraph', 'memgraph-autoheal']);
// Dynamic port resolution: read from memgraph.lock (written by start-memgraph.sh).
// Falls back to MEMGRAPH_BOLT_URL env var → default bolt://127.0.0.1:7687.
const _mgPort = (() => { try { return require('../lib/memgraph-port.js'); } catch { return null; } })();
const _envBoltUrl = _mgPort ? _mgPort.getMemgraphBoltUrl() : (process.env.MEMGRAPH_BOLT_URL || 'bolt://127.0.0.1:7687');
const _boltUrlParts = _envBoltUrl.match(/bolt:\/\/([^:]+):(\d+)/);
const BOLT_HOST = _boltUrlParts?.[1] || 'localhost';
const BOLT_PORT = parseInt(_boltUrlParts?.[2] || '7687', 10);
// AGENT_DIR: resolved from HELIOS_ROOT env var (canonical), then from helios-root.js
// (which uses __dirname relative to the repo root), then ~/.helios as last resort.
// Never hardcode a directory name — different users clone to different locations.
const { HELIOS_ROOT: _HELIOS_ROOT_AUTOSTART } = require('../lib/helios-root.js') as { HELIOS_ROOT: string };
const AGENT_DIR = process.env.HELIOS_ROOT ?? _HELIOS_ROOT_AUTOSTART;

// MEMGRAPH_WSL_BINARY: override if Memgraph is installed to a non-standard path in WSL
const MEMGRAPH_WSL_BINARY = process.env.MEMGRAPH_WSL_BINARY ?? '/usr/lib/memgraph/memgraph';

async function isBoltReachable(host = BOLT_HOST, port = BOLT_PORT, timeoutMs = 2000): Promise<boolean> {
  // Query-level probe: run RETURN 1 via the Bolt driver.
  // A plain TCP connect succeeds while Memgraph is still replaying WAL/rebuilding indexes
  // (port opens before the database is queryable). This causes ensureIndexes() to fire
  // CREATE INDEX which requires a unique storage lock, racing with WAL replay and causing
  // SIGABRT. Only declare Memgraph ready when it actually responds to a query.
  // Short timeoutMs values (e.g. 100ms from fast-poll) fall back to TCP-only to avoid
  // spamming neo4j-driver connection setup overhead during the initial poll phase.
  if (timeoutMs <= 150) {
    // Fast-poll phase: TCP only to avoid driver overhead per 10ms interval
    return new Promise((resolve) => {
      const socket = net.createConnection({ host, port }, () => { socket.destroy(); resolve(true); });
      socket.on('error', () => resolve(false));
      socket.setTimeout(timeoutMs, () => { socket.destroy(); resolve(false); });
    });
  }
  // Slow-poll phase: full query probe
  try {
    const neo4j = (() => { try { return require('../lib/safe-memgraph.js') as any; } catch { return null; } })();
    if (neo4j?.rawRead) {
      const result = await Promise.race([
        neo4j.rawRead('RETURN 1 AS ok'),
        new Promise<never>((_, rej) => setTimeout(() => rej(new Error('probe timeout')), timeoutMs)),
      ]) as any[];
      return Array.isArray(result) && result.length > 0;
    }
  } catch { /* fall through to TCP */ }
  // Fallback: TCP if safe-memgraph not loaded yet
  return new Promise((resolve) => {
    const socket = net.createConnection({ host, port }, () => { socket.destroy(); resolve(true); });
    socket.on('error', () => resolve(false));
    socket.setTimeout(timeoutMs, () => { socket.destroy(); resolve(false); });
  });
}

async function isDockerAvailable(): Promise<boolean> {
  try {
    await execFileAsync('docker', ['info'], { timeout: 5000 });
    return true;
  } catch { return false; }
}

/**
 * Detect and remove rogue Memgraph containers.
 * Only helios-memgraph and memgraph-autoheal are canonical.
 * Everything else with "memgraph" in the name gets stopped + removed.
 */
async function cleanupRogueContainers(log: (msg: string) => void): Promise<void> {
  try {
    const { stdout } = await execFileAsync('docker', [
      'ps', '-a', '--filter', 'name=memgraph', '--format', '{{.Names}}'
    ], { timeout: 5000 });
    const containers = stdout.trim().split('\n').filter(Boolean);
    for (const name of containers) {
      if (CANONICAL_CONTAINERS.has(name)) continue;
      // Migration path: if old canonical 'memgraph' container exists, rename it
      if (name === 'memgraph') {
        log(`[memgraph-autostart] 🔄 Migrating old container "memgraph" → "helios-memgraph"`);
        try {
          await execFileAsync('docker', ['rename', 'memgraph', 'helios-memgraph'], { timeout: 10000 });
          log(`[memgraph-autostart] ✅ Renamed "memgraph" → "helios-memgraph"`);
          continue;
        } catch (e) {
          process.stderr.write(`[memgraph-autostart.ts] operation failed: ${String(e)}\n`);
          log(`[memgraph-autostart] ⚠️ Rename failed — falling through to stop+rm`);
        }
      }
      log(`[memgraph-autostart] ⚠️ Rogue container detected: "${name}" — stopping and removing`);
      try {
        await execFileAsync('docker', ['stop', name], { timeout: 10000 });
        await execFileAsync('docker', ['rm', name], { timeout: 10000 });
        log(`[memgraph-autostart] ✅ Removed rogue container: "${name}"`);
      } catch (err) {
        log(`[memgraph-autostart] ⚠️ Failed to remove "${name}": ${err instanceof Error ? err.message : String(err)}`);
      }
    }
  } catch (e) { process.stderr.write(`[extensions] docker not available or no containers: ${String(e)}\n`); }
}

let _ensurePromise: Promise<void> | null = null;

async function _waitForTextSearchReady(log: (msg: string) => void): Promise<void> {
  const MAX_ATTEMPTS = 10;
  const INTERVAL_MS = 3000;
  for (let i = 0; i < MAX_ATTEMPTS; i++) {
    try {
      const sm = require('../lib/safe-memgraph');
      const rows = await sm.rawRead(
        "CALL text_search.search_all('code_chunk_text', 'test', {limit: 1}) YIELD node RETURN count(node) AS c"
      );
      if (rows && rows.length > 0) {
        if (i > 0) log(`[memgraph-autostart] text_search ready after ${i * INTERVAL_MS / 1000}s`);
        return;
      }
    } catch (e: any) {
      const msg = e?.message || String(e);
      if (msg.includes('Unknown exception') || msg.includes('text_search')) {
        if (i === 0) log('[memgraph-autostart] waiting for text indices to rebuild...');
        await new Promise(r => setTimeout(r, INTERVAL_MS));
        continue;
      }
      return;
    }
  }
  log('[memgraph-autostart] text_search not ready after 30s — proceeding (degraded)');
}

/** Maps transaction_id → first-seen timestamp (ms). Used for 2-cycle age tracking. */
const _seenTransactions = new Map<string, number>();

async function ensureMemgraph(log: (msg: string) => void): Promise<void> {
  if (_ensurePromise) { await _ensurePromise; return; }
  _ensurePromise = _ensureMemgraphImpl(log);
  try { await _ensurePromise; } finally { _ensurePromise = null; }
}

async function _ensureMemgraphImpl(log: (msg: string) => void): Promise<void> {
  // 1. Fast check -- already running?
  if (await isBoltReachable()) {
    await _waitForTextSearchReady(log);
    emitMemgraphReady(log);
    ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
    return;
  }

  log('[memgraph-autostart] Bolt port 7687 not reachable -- checking Docker...');

  // 2a. On Windows: Memgraph is a WSL native service — Docker is not available.
  //     Use wsl.exe to start the systemd service instead.
  //     This is the primary Windows path — no WSL = no Memgraph = no graph features.
  if (typeof process !== 'undefined' && process.platform === 'win32') {
    log('[memgraph-autostart] Windows detected — starting Memgraph via WSL systemctl');
    const { spawnSync } = require('child_process');

    // Check if WSL is available at all before trying to start Memgraph
    const wslCheck = spawnSync('wsl.exe', ['--status'], { timeout: 5000, encoding: 'utf8' });
    if (wslCheck.error || wslCheck.status !== 0) {
      const wslErrMsg = wslCheck.error ? wslCheck.error.message : `exit ${wslCheck.status}`;
      log(`[memgraph-autostart] ❌ WSL not available (${wslErrMsg}) — Memgraph requires WSL2 with Ubuntu on Windows`);
      log('[memgraph-autostart] Install WSL2: wsl --install -d Ubuntu  (run in PowerShell as admin)');
      log('[memgraph-autostart] Then install Memgraph inside Ubuntu: https://memgraph.com/docs/getting-started/install-memgraph/wsl');
      emitMemgraphDegraded(log, 'WSL2 not available — install WSL2 with Ubuntu to enable Memgraph');
      return;
    }

    // Try systemctl with sudo -n (non-interactive, requires NOPASSWD sudoers rule)
    // CRITICAL: Check is-active first — calling 'start' on an already-starting service
    // triggers a full restart cycle (stop → wait-port-7687.sh → start) which causes
    // the port-bind assertion failure. Only start if genuinely not active.
    const result = spawnSync('wsl.exe', ['--', 'bash', '-lc',
      'systemctl is-active memgraph >/dev/null 2>&1 && echo "exit:0 (already active)" || { sudo -n systemctl start memgraph 2>&1; echo "exit:$?"; }'], {
      timeout: 15000, encoding: 'utf8'
    });
    const systemctlOutput = result.stdout || '';
    const systemctlSucceeded = systemctlOutput.includes('exit:0');

    if (systemctlSucceeded) {
      log('[memgraph-autostart] started Memgraph via WSL systemctl');
    } else {
      // systemctl failed (systemd not the init, or memgraph service not installed)
      log('[memgraph-autostart] systemctl failed — trying direct memgraph binary via WSL');
      // Use start-memgraph.sh so dynamic port selection and lock file write happen correctly.
      // Fall back to direct binary with default port only if the script is absent.
      const startScript = '/usr/local/bin/start-memgraph.sh';
      const directResult = spawnSync(
        'wsl.exe',
        ['--', 'bash', '-lc', `[ -x "${startScript}" ] && nohup ${startScript} >/tmp/mg-autostart.log 2>&1 & || sudo -n ${MEMGRAPH_WSL_BINARY} --bolt-port=7687 --data-directory=/var/lib/memgraph --log-level=WARNING &`],
        { timeout: 5000, encoding: 'utf8' }
      );
      if (directResult.status !== 0 && !directResult.stdout?.includes('Memgraph')) {
        // Binary not found in WSL
        log(`[memgraph-autostart] ❌ Memgraph binary not found in WSL at ${MEMGRAPH_WSL_BINARY}`);
        log('[memgraph-autostart] Set MEMGRAPH_WSL_BINARY env var to override the path, or install Memgraph: https://memgraph.com/docs/getting-started/install-memgraph/wsl');
        log('[memgraph-autostart] Or enable systemd in WSL2: add [boot] systemd=true to /etc/wsl.conf');
        emitMemgraphDegraded(log, 'Memgraph binary not found in WSL — install Memgraph in WSL Ubuntu');
        return;
      }
    }

    // P3-3 NX fix: fast 10ms poll for first 2s (NX WAIT_FOR_SERVER_CONFIG pattern),
    // then fall back to 2000ms for the remaining 12s.
    // BUG-3 fix: pass 100ms timeout to isBoltReachable() in the fast-poll phase.
    // Without this, each probe can block for up to 2000ms (default), making 200 × 2000ms
    // = 400s worst case. With 100ms: 200 × (10ms sleep + 100ms probe) = ~22s max, and
    // the common case (ECONNREFUSED = immediate) is still ~2s.
    let bolted = false;
    for (let i = 0; i < 200 && !bolted; i++) {
      await new Promise(r => setTimeout(r, 10)); // 10ms sleep between probes (NX pace)
      if (await isBoltReachable(BOLT_HOST, BOLT_PORT, 100)) { bolted = true; } // 100ms timeout per probe
    }
    if (!bolted) {
      // Fall back to slow poll for remaining time
      for (let i = 0; i < 6 && !bolted; i++) {
        await new Promise(r => setTimeout(r, 2000));
        if (await isBoltReachable()) { bolted = true; }
      }
    }
    if (bolted) {
        emitMemgraphReady(log);
        ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
        return;
    }
    log('[memgraph-autostart] ⚠️ WSL Memgraph started but Bolt not reachable after 14s');
    log('[memgraph-autostart] Check WSL Memgraph logs: wsl -- journalctl -u memgraph -n 30');
    log('[memgraph-autostart] Or check the process: wsl -- ps aux | grep memgraph');
    emitMemgraphDegraded(log, 'Memgraph started in WSL but Bolt port 7687 unreachable — check WSL Memgraph logs');
    return;
  }

  // 2b. Docker available?
  if (!(await isDockerAvailable())) {
    log('[memgraph-autostart] Docker not available -- skipping Memgraph auto-start');
    emitMemgraphDegraded(log, 'unreachable after retries');
    return;
  }

  await cleanupRogueContainers(log);

  // 3. Check if container exists but is stopped
  let foundContainerName: string | null = null;
  let foundState: string | null = null;

  for (const tryName of [CONTAINER_NAME]) {
    try {
      const { stdout: status } = await execFileAsync('docker', [
        'inspect', '--format', '{{.State.Status}}', tryName
      ], { timeout: 5000 });
      const state = status.trim();
      if (state) {
        foundContainerName = tryName;
        foundState = state;
        break;
      }
    } catch (e) { process.stderr.write(`[extensions] not found, try next: ${String(e)}\n`); }
  }

  if (foundContainerName && foundState) {
    const containerName = foundContainerName;
    const state = foundState;

    if (state === 'running') {
      // Container says running but Bolt not reachable yet -- wait a bit
      log('[memgraph-autostart] Container running, waiting for Bolt...');
      for (let i = 0; i < 10; i++) {
        await new Promise(r => setTimeout(r, 2000));
        if (await isBoltReachable()) {
          emitMemgraphReady(log);
          ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
          return;
        }
      }
      log('[memgraph-autostart] ⚠️ Container running but Bolt not responding after 20s');
      emitMemgraphDegraded(log, 'unreachable after retries');
      return;
    }

    if (state === 'exited' || state === 'created') {
      // Use docker compose up -d instead of docker start — compose detects config
      // changes in docker-compose.yml and recreates the container when needed.
      // This ensures new flags (--log-async, --parallel-schema-recovery, etc.)
      // are picked up without manual intervention.
      const composeDir = join(AGENT_DIR, 'proxies', 'memgraph');
      const composeFile = join(composeDir, 'docker-compose.yml');
      log(`[memgraph-autostart] Starting Memgraph via compose (picks up config changes)...`);
      try {
        await execFileAsync('docker', ['compose', '-f', composeFile, 'up', '-d'], {
          timeout: 30000,
          cwd: composeDir,
        });
      } catch (startErr) {
        const msg = startErr instanceof Error ? startErr.message : String(startErr);
        if (msg.includes('port is already allocated') || msg.includes('address already in use')) {
          log('[memgraph-autostart] Port conflict on start -- probing Bolt directly...');
          for (let i = 0; i < 10; i++) {
            await new Promise(r => setTimeout(r, 2000));
            if (await isBoltReachable()) {
              emitMemgraphReady(log);
              ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
              return;
            }
          }
          log('[memgraph-autostart] ⚠️ Port conflict and Bolt not responding');
          emitMemgraphDegraded(log, 'unreachable after retries');
          return;
        }
        // Fallback: try plain docker start if compose fails (e.g. compose file missing)
        log(`[memgraph-autostart] Compose failed, falling back to docker start: ${msg}`);
        try {
          await execFileAsync('docker', ['start', containerName], { timeout: 10000 });
        } catch (e) {
          process.stderr.write(`[memgraph-autostart.ts] operation failed: ${String(e)}\n`);
          throw startErr; // throw original compose error
        }
      }
      // Wait for Bolt
      for (let i = 0; i < 15; i++) {
        await new Promise(r => setTimeout(r, 2000));
        if (await isBoltReachable()) {
          emitMemgraphReady(log);
          ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
          return;
        }
      }
      log('[memgraph-autostart] ⚠️ Container started but Bolt not responding after 30s');
      emitMemgraphDegraded(log, 'unreachable after retries');
      return;
    }
  }

  // 4. Container doesn't exist -- create via docker-compose (health checks + autoheal)
  log('[memgraph-autostart] Creating Memgraph stack via docker-compose...');
  try {
    const composeDir = join(AGENT_DIR, 'proxies', 'memgraph');
    await execFileAsync(
      'docker',
      ['compose', '-f', join(composeDir, 'docker-compose.yml'), 'up', '-d', '--force-recreate'],
      { timeout: 60000, cwd: composeDir }
    );

    // Wait for Bolt
    for (let i = 0; i < 30; i++) {
      await new Promise(r => setTimeout(r, 2000));
      if (await isBoltReachable()) {
        emitMemgraphReady(log);
        ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
        return;
      }
    }
    log('[memgraph-autostart] ⚠️ Stack started but Bolt not responding after 60s');
    emitMemgraphDegraded(log, 'unreachable after retries');
  } catch (err) {
    log(`[memgraph-autostart] ⚠️ Failed to start compose stack: ${err instanceof Error ? err.message : String(err)}`);
    emitMemgraphDegraded(log, 'unreachable after retries');
  }
}

// ── Background health monitor ───────────────────────────────────────────────
function getHealthTimer(): ReturnType<typeof setInterval> | null {
  return (globalThis as any).__helios_health_timer ?? null;
}
function setHealthTimer(t: ReturnType<typeof setInterval> | null): void {
  (globalThis as any).__helios_health_timer = t;
}

function startHealthMonitor(log: (msg: string) => void): void {
  if (getHealthTimer()) return; // already running
  setHealthTimer(setInterval(async () => {
    if (!(await isDockerAvailable())) return;
    await cleanupRogueContainers(log);
    if (await isBoltReachable()) {
      // Bolt healthy — reap stuck transactions using 2-cycle grace period (2026-04-16 fix)
      // Previously killed ALL running transactions immediately; now uses cleanupStuckTransactions
      // which tracks transactions across cycles and only kills those seen in 2+ consecutive checks.
      try {
        await cleanupStuckTransactions(log);
      } catch (e) { process.stderr.write(`[extensions] non-fatal: ${String(e)}\n`); }
      return;
    }
    if (!(await isDockerAvailable())) return;
    log('[memgraph-autostart] ⚠️ Bolt unreachable — attempting recovery...');
    try {
      const composeDir = join(AGENT_DIR, 'proxies', 'memgraph');
      await execFileAsync(
        'docker',
        ['compose', '-f', join(composeDir, 'docker-compose.yml'), 'up', '-d', '--force-recreate'],
        { timeout: 30000, cwd: composeDir }
      );
      // Wait for recovery (up to 20s)
      for (let i = 0; i < 10; i++) {
        await new Promise(r => setTimeout(r, 2000));
        if (await isBoltReachable()) {
          log('[memgraph-autostart] ✅ Memgraph recovered');
          // Reload MAGE modules after recovery (they don't persist across restarts)
          _loadMageModules(log).catch(e => process.stderr.write(`[extensions] operation failed: ${e?.message}\n`));
          ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
          return;
        }
      }
      log('[memgraph-autostart] ⚠️ Recovery failed — circuit breaker will protect queries');
    } catch (e) { process.stderr.write(`[extensions] non-fatal: ${String(e)}\n`); }
  }, 30_000)); // 30s interval — 2-cycle grace = 60s, matches query_execution_timeout_sec=60
  // Don't keep Node alive just for the health monitor
  if ((getHealthTimer() as any)?.unref) (getHealthTimer() as any).unref();
}

function stopHealthMonitor(): void {
  const timer = getHealthTimer();
  if (timer) {
    clearInterval(timer);
    setHealthTimer(null);
  }
}

// ── Extension export ────────────────────────────────────────────────

export function activate(pi: any): void {
  // ── HELIOS_SKIP_INFRA / HELIOS_SKIP_MEMGRAPH: skip Memgraph startup ──
  // HELIOS_SKIP_INFRA=1: full headless mode (no infra at all)
  // HELIOS_SKIP_MEMGRAPH=1: bolt unavailable (e.g. E2B sandbox) but other infra ok
  if (process.env.HELIOS_SKIP_INFRA === '1' || process.env.HELIOS_SKIP_MEMGRAPH === '1') {
    console.warn('[WARN] [memgraph-autostart] HELIOS_SKIP_INFRA/MEMGRAPH=1 — skipping Memgraph startup (headless mode)');
    // Register DATA_LAYER phase with immediate resolution so downstream extensions
    // that await waitForPhase(DATA_LAYER) don't wait for the 8s timeout.
    registerStartupPhase('memgraph-data-layer', StartupPhase.DATA_LAYER, async () => {
      // Headless mode: Memgraph unavailable by design — resolve immediately (degraded)
    });
    return;
  }

  const log = (msg: string) => {
    if (pi?.logger?.info) pi.logger.info(msg);
    else if (!QUIET) process.stderr.write(msg + '\n');
  };

  // Register as INFRASTRUCTURE phase handler in the startup lifecycle
  // MAGMA dual-stream pattern: fast path (probe only) is blocking,
  // slow path (Docker recovery) is fire-and-forget background work.
  registerStartupPhase('memgraph-autostart', StartupPhase.INFRASTRUCTURE, async () => {
    // FAST PATH: 2s TCP probe — the only blocking operation
    if (await isBoltReachable()) {
      emitMemgraphReady(log);
      ensureIndexes(log).catch(e => process.stderr.write('[memgraph-autostart] ensureIndexes error: ' + e?.message + '\n'));
      cleanupStuckTransactions(log).catch(e => process.stderr.write(`[extensions] operation failed: ${e?.message}\n`));
      startHealthMonitor(log);
      return;
    }
    // Bolt not reachable — emit degraded immediately so downstream
    // extensions don't wait. Background recovery will emit memgraph_ready
    // when Docker brings Bolt up.
    log('[memgraph-autostart] Bolt not reachable — starting background recovery');
    emitMemgraphDegraded(log, 'bolt not reachable — background recovery starting');
    // SLOW PATH: fire-and-forget Docker recovery
    ensureMemgraph(log)
      .then(() => cleanupStuckTransactions(log))
      .then(() => startHealthMonitor(log))
      .catch((err) => {
        log(`[memgraph-autostart] Background recovery failed: ${err instanceof Error ? err.message : String(err)}`);
      });
  });

  // DATA_LAYER gate: only resolves when Memgraph bolt is confirmed reachable.
  // Extensions using waitForPhase(DATA_LAYER) depend on this to know graph queries will work.
  // Budget: must complete within phase timeout (3s). If bolt isn't up by then, proceed degraded.
  registerStartupPhase('memgraph-data-layer', StartupPhase.DATA_LAYER, async () => {
    if (await isBoltReachable()) return;
    // Quick retry (2 attempts × 500ms) — fits within 3s phase timeout
    for (let i = 0; i < 2; i++) {
      await new Promise(r => setTimeout(r, 500));
      if (await isBoltReachable()) {
        log('[memgraph-autostart] DATA_LAYER: bolt became reachable');
        return;
      }
    }
    log('[memgraph-autostart] DATA_LAYER: bolt not reachable within budget — degraded mode');
  });

  // Register session_start hook -- check if startup lifecycle already ran this phase
  // R3-L22: Guard with globalThis to prevent handler accumulation on extension reload.
  if (typeof pi?.on === 'function' && !(globalThis as any).__memgraph_autostart_wired) {
    (globalThis as any).__memgraph_autostart_wired = true;
    pi.on('session_start', async () => {
      // P3-1: NX daemon pattern — try to reuse the persistent driver from the previous session.
      // If healthy, this opens the warm gate immediately, skipping the full 2-8s startup probe.
      try {
        const sm = require('../lib/safe-memgraph.js');
        if (sm?.tryPersistentDriverFastPath) {
          const fastPath = await sm.tryPersistentDriverFastPath();
          if (fastPath) {
            startHealthMonitor(log);
            return; // warm start — skip ensureMemgraph entirely
          }
        }
      } catch { /* safe-memgraph not available — fall through to cold start */ }

      const state = getStartupState();
      // If startup lifecycle already ran INFRASTRUCTURE phase, skip ensureMemgraph
      if (state.completedPhases.has(StartupPhase.INFRASTRUCTURE)) {
        log('[memgraph-autostart] INFRASTRUCTURE phase already completed by startup lifecycle — starting health monitor only');
        startHealthMonitor(log);
        return;
      }
      // Fallback: startup lifecycle not active, run as before
      ensureMemgraph(log)
        .then(() => cleanupStuckTransactions(log))
        .then(() => startHealthMonitor(log))
        .catch((err) => {
          log(`[memgraph-autostart] Error: ${err instanceof Error ? err.message : String(err)}`);
        });
    });
    pi.on('session_shutdown', async () => {
      stopHealthMonitor();
      // Close Neo4j drivers and broker client to release all ref'd handles (TCPSocketWrap,
      // PipeWrap) and allow the process to exit naturally after pi -p completes.
      try {
        const mg = require('../lib/safe-memgraph.js');
        // Bug-4 fix: mark shutdown so broker errors are suppressed before closing
        if (mg.markShutdownInitiated) mg.markShutdownInitiated();
        // Also close triage-core mental-model singletons that hold their own drivers
        let cleanupMM: (() => Promise<void>) | undefined;
        try { cleanupMM = require('../lib/triage-core/mental-model/cleanup').cleanupMentalModelDrivers; } catch {}
        // mg.shutdown() now closes both the bolt driver AND the broker client socket.
        // 6s timeout: inner shutdown() takes up to 5s to drain the driver pool.
        await Promise.race([
          Promise.allSettled([
            mg.shutdown(),
            cleanupMM ? cleanupMM() : Promise.resolve(),
          ]),
          new Promise(r => setTimeout(r, 6000)),
        ]);
      } catch (e) {
        // swallow — shutting down
      }
    });
  }

  // ── Performance tier gate ─────────────────────────────────────────────────
  // Progressive enhancement: always emit the expected event so downstream
  // listeners (HEMA, cortex, codebase-index) don't hang waiting.
  const features = di('helios_features') as any;
  if (features && !features.memgraph) {
    log('[memgraph-autostart] Memgraph DISABLED by performance tier (lite) — emitting degraded');
    emitMemgraphDegraded(log, 'disabled by performance tier (lite)');
    return; // No Docker, no Bolt probe, no health monitor. Saves ~1GB RAM.
  }
  if (features && features.memgraphLazy) {
    log('[memgraph-autostart] Memgraph set to LAZY start (standard tier) — deferring to first query');
    // Register lazy activator — unified-graph or first direct query triggers it
    let _lazyStarted = false;
    updateScopedValue('memgraphLazyStart', async () => {
      if (_lazyStarted) return;
      _lazyStarted = true;
      log('[memgraph-autostart] Lazy-starting Memgraph on first graph query...');
      await ensureMemgraph(log);
      await cleanupStuckTransactions(log);
      startHealthMonitor(log);
    });
    emitMemgraphDegraded(log, 'deferred by performance tier (standard/lazy)');
    return; // Don't start eagerly — lazy activator handles it
  }

  log('[memgraph-autostart] Extension registered -- will auto-start Memgraph via INFRASTRUCTURE phase');

  // Graph bootstrap — seeds SkillDef, ToolDef, AgentRole, OrchestratorProfile,
  // and precomputed Extension+SkillDef embeddings on fresh installs.
  // C1 fix: call runGraphBootstrap directly (2s delay) rather than relying on
  // the event bus cross-module singleton which may not be shared.
  try {
    const { runGraphBootstrap } = require('../lib/graph-bootstrap.ts');
    if (typeof runGraphBootstrap === 'function') {
      setTimeout(() => runGraphBootstrap().catch(() => {}), 2000);
    }
  } catch (_e) { /* fail-open: bootstrap is non-critical */ }

  // Removed redundant direct ensureMemgraph() call.
  // The INFRASTRUCTURE phase handler and session_start hook cover all cases.
}


/**
 * Kill ALL stuck Memgraph transactions — age-based, not pattern-based.
 *
 * HISTORY (2026-04-09 incident): The previous implementation used a hardcoded
 * pattern whitelist that only matched 4 specific query shapes. In production,
 * plan-tracker, codebase-index, and session-memory all deadlocked under
 * SNAPSHOT_ISOLATION — but none matched the whitelist, causing cascading hangs.
 *
 * FIX: Kill ANY running transaction that isn't SHOW TRANSACTIONS itself.
 * Safe: all Helios writes use idempotent MERGE (safe to re-submit).
 * Called on: session_start, and periodically by the health monitor.
 */
/**
 * Pure function — processes one cycle of SHOW TRANSACTIONS output.
 * Exported for unit testing without live Memgraph.
 *
 * Strategy: 2-cycle age tracking
 * - First sighting of a TID: record in seenMap, do NOT kill
 * - Second sighting: transaction is stuck → add to toKill list
 * - TIDs absent from current snapshot → obsoleteKeys (caller should delete from map)
 *
 * This guarantees a minimum one full cycle (≥60 s by default) grace period.
 */
export function _processTransactionCycle(
  txns: Array<Record<string, unknown>>,
  seenMap: Map<string, number>,
  now: number,
): { toKill: Array<{ tid: string; query: string }>; obsoleteKeys: string[] } {
  const currentTids = new Set<string>();
  const toKill: Array<{ tid: string; query: string }> = [];

  for (const tx of txns) {
    const tid = String((tx as any).transaction_id ?? (tx as any)[1] ?? '');
    const queryArr = (tx as any).query ?? (tx as any)[2];
    const query = Array.isArray(queryArr) ? (queryArr[0] ?? '') : String(queryArr ?? '');
    const status = String((tx as any).status ?? (tx as any)[3] ?? 'running');

    if (!tid || query.includes('SHOW TRANSACTIONS') || query.includes('mg.load_all') || query.includes('mg.') && query.includes('CALL') || status !== 'running') continue;
    currentTids.add(tid);

    if (seenMap.has(tid)) {
      // Second sighting — transaction persisted across cycles → kill it
      toKill.push({ tid, query: query.slice(0, 80) });
    } else {
      // First sighting — record timestamp, give it one full cycle grace period
      seenMap.set(tid, now);
    }
  }

  // TIDs no longer in SHOW TRANSACTIONS have completed — remove from tracking map
  const obsoleteKeys = [...seenMap.keys()].filter(k => !currentTids.has(k));
  return { toKill, obsoleteKeys };
}

async function cleanupStuckTransactions(log: (msg: string) => void): Promise<void> {
  if ((di('helios_memgraph_ready') ?? di('memgraphReady')) !== true) return;

  try {
    const { autoCommit } = require('../lib/safe-memgraph.js');
    const txns = await autoCommit('SHOW TRANSACTIONS') as Array<Record<string, unknown>>;

    const { toKill, obsoleteKeys } = _processTransactionCycle(txns, _seenTransactions, Date.now());

    // Prune completed transactions from tracking map
    for (const k of obsoleteKeys) _seenTransactions.delete(k);
    // Size cap: if map > 500, clear entries older than 1 hour
    if (_seenTransactions.size > 500) {
      const now = Date.now(); // FIX: was `now` (undefined — `now` is a param of _processTransactionCycle, not in this scope)
      const oneHourAgo = now - 3_600_000;
      for (const [k, ts] of _seenTransactions) {
        if (ts < oneHourAgo) _seenTransactions.delete(k);
      }
    }

    if (!toKill.length) {
      log('[memgraph-autostart] No stuck transactions found');
      return;
    }

    log(`[memgraph-autostart] Terminating ${toKill.length} stuck transaction(s) (seen in 2+ consecutive cycles)...`);
    for (const { tid, query } of toKill) {
      try {
        await autoCommit(`TERMINATE TRANSACTIONS "${tid}"`);
        log(`[memgraph-autostart] Terminated: ${tid.slice(0, 20)}... (${query})`);
      } catch (_e) { /* ignore -- may have already finished */ }
    }
  } catch (err) {
    const _msg = err instanceof Error ? err.message : String(err);
    if (_msg.includes('Pool is closed') || _msg.includes('Connection is closed')) return;
    log(`[memgraph-autostart] Transaction cleanup failed (non-fatal): ${_msg}`);
  }
}

export default activate;
