from __future__ import annotations

import hashlib
import json
import os
import re
import zipfile
from dataclasses import dataclass, field
from decimal import Decimal, InvalidOperation
from pathlib import Path
from sqlite3 import Connection
from typing import Sequence
from urllib import request, error
import xml.etree.ElementTree as ET

from jarvis_finance.audit.log import record_audit_event
from jarvis_finance.equity.manage import add_manual_position_from_catalog
from jarvis_finance.imports.common import stable_id, utc_now
from jarvis_finance.market_data.catalog import CatalogEntryInput, upsert_catalog_entry
from jarvis_finance.quality.alerts import create_alert

VALID_MAPPING_STATUS = {"exact_isin_match", "probable", "needs_manual_review", "rejected", "confirmed"}
VALID_EXTRACTORS = {"deterministic_parser", "qwen_local", "manual"}
ISIN_RE = re.compile(r"\b[A-Z]{2}[A-Z0-9]{9}[0-9]\b")
CURRENCY_RE = re.compile(r"\b(CHF|USD|EUR|GBP|JPY|CAD|AUD)\b", re.I)
TICKER_RE = re.compile(r"\b[A-Z0-9]{1,6}(?:[.:-][A-Z0-9]{1,4})?\b")
DECIMAL_RE = re.compile(r"(?<![A-Z0-9])[-+]?\d{1,3}(?:[ '\u00a0]\d{3})*(?:[.,]\d+)?|(?<![A-Z0-9])[-+]?\d+(?:[.,]\d+)?")


@dataclass(frozen=True)
class CandidateInput:
    source_file: str
    platform: str
    account: str = ""
    source_row_number: str = ""
    raw_name: str = ""
    raw_ticker: str = ""
    raw_isin: str = ""
    raw_currency: str = ""
    raw_quantity: str = ""
    raw_asset_class: str = ""
    extracted_by: str = "deterministic_parser"
    extraction_confidence: str = "low"
    proposed_name: str = ""
    proposed_ticker: str = ""
    proposed_isin: str = ""
    proposed_exchange: str = ""
    proposed_currency: str = ""
    proposed_asset_class: str = ""
    provider: str = ""
    provider_symbol: str = ""
    mapping_status: str = "needs_manual_review"
    review_note: str = ""


@dataclass(frozen=True)
class CandidateSummary:
    candidates_total: int = 0
    with_isin: int = 0
    without_isin: int = 0
    with_ticker: int = 0
    with_currency: int = 0
    exact_isin_match: int = 0
    probable: int = 0
    needs_manual_review: int = 0
    confirmed: int = 0
    rejected: int = 0
    importable_confirmed: int = 0
    blocked: int = 0
    qwen_used: bool = False
    parser_qwen_agree: int = 0
    parser_qwen_disagree: int = 0
    transactions_written: int = 0
    audit_logs_written: int = 0
    warnings: int = 0
    errors: int = 0
    dry_run: bool = False
    qwen_available: bool = False


@dataclass(frozen=True)
class QwenResolveSummary:
    qwen_available: bool = False
    qwen_used: bool = False
    model: str = ""
    candidates_total: int = 0
    exact_before: int = 0
    exact_after: int = 0
    probable_before: int = 0
    probable_after: int = 0
    needs_review_before: int = 0
    needs_review_after: int = 0
    confirmed_after: int = 0
    improved: int = 0
    manual_review_remaining: int = 0
    warnings: int = 0
    errors: int = 0
    dry_run: bool = False

def _norm_space(text: str) -> str:
    return re.sub(r"\s+", " ", str(text or "")).strip()


def _source_label(source_file: str) -> str:
    p = Path(source_file)
    return hashlib.sha256(str(p.name).encode("utf-8")).hexdigest()[:12]


def _candidate_id(inp: CandidateInput) -> str:
    return stable_id(
        "instrument_candidate",
        _source_label(inp.source_file), inp.platform, inp.account, inp.source_row_number,
        inp.raw_name, inp.raw_isin, inp.raw_ticker, inp.raw_currency, inp.raw_quantity,
    )


def _file_label(path: str | Path) -> str:
    return Path(path).name


def extract_docx_text(path: str | Path) -> list[str]:
    """Return sanitized paragraph/table text lines from a DOCX without storing broker content in repo."""
    with zipfile.ZipFile(path) as zf:
        xml = zf.read("word/document.xml")
    root = ET.fromstring(xml)
    ns = {"w": "http://schemas.openxmlformats.org/wordprocessingml/2006/main"}
    lines: list[str] = []
    for para in root.findall(".//w:p", ns):
        text = _norm_space("".join(node.text or "" for node in para.findall(".//w:t", ns)))
        if text:
            lines.append(text)
    return lines


def _looks_position_line(line: str) -> bool:
    up = line.upper()
    if any(stop in up for stop in ("TOTAL", "SUMME", "CASH", "LIQUIDIT", "KONTO", "MARKTWERT", "PERFORMANCE")):
        return False
    has_identifier = bool(ISIN_RE.search(up) or CURRENCY_RE.search(up))
    has_quantity_like = len(DECIMAL_RE.findall(line)) >= 1
    has_name = len(re.sub(r"[^A-Za-z]", "", line)) >= 4
    return has_identifier and has_quantity_like and has_name


def _asset_class(line: str) -> str:
    up = line.upper()
    if "ETF" in up or "FUND" in up or "FONDS" in up:
        return "etf"
    return "equity"


def _extract_ticker(line: str, isin: str, currency: str) -> str:
    words = [m.group(0).upper() for m in TICKER_RE.finditer(line.upper())]
    ignored = {isin, currency, "CHF", "USD", "EUR", "GBP", "ETF", "FUND", "ISIN"}
    for w in words:
        if w not in ignored and not w.isdigit() and len(w) <= 8:
            return w
    return ""


def _extract_quantity(line: str) -> str:
    # Keep as source text; callers must not print values to chat. Prefer the first decimal after an identifier-ish row.
    vals = DECIMAL_RE.findall(line)
    return _norm_space(vals[0]).replace("'", "").replace("\u00a0", " ") if vals else ""


def deterministic_parse_docx(path: str | Path, *, platform: str, account: str = "") -> list[CandidateInput]:
    lines = extract_docx_text(path)
    out: list[CandidateInput] = []
    seen: set[tuple[str, str, str]] = set()

    def add(inp: CandidateInput) -> None:
        key = (inp.source_row_number, inp.raw_isin or inp.proposed_isin, inp.raw_ticker or inp.proposed_ticker)
        if key not in seen:
            seen.add(key)
            out.append(inp)

    # True Wealth-style blocks: instrument name followed by explicit ISIN and currency/quantity nearby.
    for idx, line in enumerate(lines, start=1):
        raw = _norm_space(line)
        isin_match = ISIN_RE.search(raw.upper())
        if not isin_match:
            continue
        isin = isin_match.group(0)
        prev_candidates = [l for l in lines[max(0, idx - 5):idx - 1] if not any(x in l.upper() for x in ("POSITIONSWERT", "ISIN", "CHF ", "USD ", "EUR "))]
        name = _norm_space(prev_candidates[-1] if prev_candidates else raw.replace(isin, ""))[:240]
        window = " ".join(lines[idx:min(len(lines), idx + 7)])
        curr_match = CURRENCY_RE.search(window.upper())
        currency = curr_match.group(1).upper() if curr_match else ""
        ticker = _extract_ticker(name, isin, currency)
        add(CandidateInput(
            source_file=_file_label(path), platform=platform, account=account, source_row_number=str(idx),
            raw_name=name, raw_ticker=ticker, raw_isin=isin, raw_currency=currency,
            raw_quantity=_extract_quantity(window), raw_asset_class=_asset_class(name), extracted_by="deterministic_parser",
            extraction_confidence="high", proposed_name=name, proposed_ticker=ticker,
            proposed_isin=isin, proposed_currency=currency, proposed_asset_class=_asset_class(name),
            mapping_status="exact_isin_match", review_note="deterministic DOCX staging candidate with source ISIN; no productive booking",
        ))

    # Fallback single-line extraction for synthetic/simple exports.
    for idx, line in enumerate(lines, start=1):
        if not _looks_position_line(line):
            continue
        raw = _norm_space(line)
        isin_match = ISIN_RE.search(raw.upper())
        isin = isin_match.group(0) if isin_match else ""
        curr_match = CURRENCY_RE.search(raw.upper())
        currency = curr_match.group(1).upper() if curr_match else ""
        ticker = _extract_ticker(raw, isin, currency)
        name = DECIMAL_RE.sub(" ", raw)
        if isin:
            name = name.replace(isin, " ")
        if currency:
            name = re.sub(rf"\b{re.escape(currency)}\b", " ", name, flags=re.I)
        name = _norm_space(name)[:240]
        status = "exact_isin_match" if isin else "needs_manual_review"
        confidence = "high" if isin else ("medium" if ticker and currency else "low")
        add(CandidateInput(
            source_file=_file_label(path), platform=platform, account=account, source_row_number=str(idx),
            raw_name=name or raw[:240], raw_ticker=ticker, raw_isin=isin, raw_currency=currency,
            raw_quantity=_extract_quantity(raw), raw_asset_class=_asset_class(raw), extracted_by="deterministic_parser",
            extraction_confidence=confidence, proposed_name=name or raw[:240], proposed_ticker=ticker,
            proposed_isin=isin, proposed_currency=currency, proposed_asset_class=_asset_class(raw),
            mapping_status=status, review_note="deterministic DOCX staging candidate; no productive booking",
        ))

    # PostFinance-style split table: ticker-only rows under Produkt; never auto-confirm ticker-only rows.
    current_asset_class = "equity"
    in_positions = False
    headers = {"PRODUKT", "ANZAHL", "EINSTANDSKURS", "TOTALWERT", "DIFF. VORTAG", "GELDKURS WÄHRUNG", "G&V CHF", "TOTALWERT CHF", "POSITIONEN %"}
    for idx, line in enumerate(lines, start=1):
        up = _norm_space(line).upper()
        if up in {"AKTIEN", "ETFS", "ETF"}:
            current_asset_class = "etf" if "ETF" in up else "equity"
            continue
        if up == "PRODUKT":
            in_positions = True
            continue
        if not in_positions or up in headers or up.startswith("PAGE "):
            continue
        if not re.fullmatch(r"[A-Z][A-Z0-9.]{0,7}", up):
            continue
        window = " ".join(lines[idx:min(len(lines), idx + 9)])
        curr_match = CURRENCY_RE.search(window.upper())
        currency = curr_match.group(1).upper() if curr_match else ""
        qty = _extract_quantity(" ".join(lines[idx:idx + 3]))
        if not qty:
            continue
        add(CandidateInput(
            source_file=_file_label(path), platform=platform, account=account, source_row_number=str(idx),
            raw_name=up, raw_ticker=up, raw_isin="", raw_currency=currency, raw_quantity=qty,
            raw_asset_class=current_asset_class, extracted_by="deterministic_parser", extraction_confidence="medium" if currency else "low",
            proposed_name=up, proposed_ticker=up, proposed_currency=currency, proposed_asset_class=current_asset_class,
            mapping_status="needs_manual_review", review_note="ticker-only split-table candidate; manual ISIN/exchange review required",
        ))
    return out


def _sanitize_for_qwen(text: str) -> str:
    sanitized = DECIMAL_RE.sub("<number_redacted>", _norm_space(text))
    sanitized = re.sub(r"\b(?:IBAN|KONTO|ACCOUNT|DEPOT|CLIENT|KUNDE|KUNDENNR)[:# ]*[A-Z0-9-]+\b", "<identifier_redacted>", sanitized, flags=re.I)
    return sanitized[:500]


def qwen_extract_candidates(text_lines: list[str], *, source_file: str, platform: str, account: str = "", endpoint: str | None = None, model: str | None = None, timeout: int = 20) -> tuple[list[CandidateInput], bool]:
    """Use local OpenAI-compatible Qwen as assistive parser. It returns staging candidates only."""
    endpoint = endpoint or os.getenv("QWEN_OPENAI_BASE_URL") or os.getenv("LOCAL_QWEN_BASE_URL") or "http://100.101.173.25:11435/v1"
    model = model or os.getenv("QWEN_MODEL", "qwen-3.6-agent")
    payload = {
        "model": model,
        "temperature": 0,
        "messages": [
            {"role": "system", "content": "Extract broker equity/ETF position candidates as JSON only from sanitized text. Do not invent ISIN, history, prices, quantities, account data, or IDs. Output {\"positions\":[{name,ticker,isin,currency,asset_class,row_number,confidence,note}]}"},
            {"role": "user", "content": "\n".join(f"{i+1}: {_sanitize_for_qwen(line)}" for i, line in enumerate(text_lines[:300]))},
        ],
    }
    req = request.Request(f"{endpoint.rstrip('/')}/chat/completions", data=json.dumps(payload).encode("utf-8"), headers={"Content-Type": "application/json"})
    try:
        with request.urlopen(req, timeout=timeout) as resp:
            data = json.loads(resp.read().decode("utf-8"))
    except (OSError, error.URLError, TimeoutError, json.JSONDecodeError):
        return [], False
    content = data.get("choices", [{}])[0].get("message", {}).get("content", "")
    match = re.search(r"\{.*\}", content, re.S)
    if not match:
        return [], True
    try:
        parsed = json.loads(match.group(0))
    except json.JSONDecodeError:
        return [], True
    out: list[CandidateInput] = []
    for pos in parsed.get("positions", []):
        isin = str(pos.get("isin") or "").strip().upper()
        if isin and not ISIN_RE.fullmatch(isin):
            isin = ""
        ticker = str(pos.get("ticker") or "").strip().upper()
        currency = str(pos.get("currency") or "").strip().upper()
        status = "exact_isin_match" if isin else "needs_manual_review"
        out.append(CandidateInput(
            source_file=_file_label(source_file), platform=platform, account=account,
            source_row_number=str(pos.get("row_number") or ""), raw_name=_norm_space(str(pos.get("name") or ""))[:240],
            raw_ticker=ticker, raw_isin=isin, raw_currency=currency, raw_quantity="",
            raw_asset_class=str(pos.get("asset_class") or "").lower() or "equity", extracted_by="qwen_local",
            extraction_confidence=str(pos.get("confidence") or "low"), proposed_name=_norm_space(str(pos.get("name") or ""))[:240],
            proposed_ticker=ticker, proposed_isin=isin, proposed_currency=currency,
            proposed_asset_class=str(pos.get("asset_class") or "").lower() or "equity", mapping_status=status,
            review_note=_norm_space(str(pos.get("note") or "qwen assistive extraction; review required")),
        ))
    return out, True


def upsert_candidate(conn: Connection, inp: CandidateInput, *, dry_run: bool = False) -> str:
    if inp.extracted_by not in VALID_EXTRACTORS:
        raise ValueError("invalid extractor")
    if inp.mapping_status not in VALID_MAPPING_STATUS:
        raise ValueError("invalid mapping status")
    cid = _candidate_id(inp)
    if dry_run:
        return cid
    now = utc_now()
    conn.execute(
        """
        INSERT INTO instrument_import_candidates(
            candidate_id, source_file, platform, account, source_row_number, raw_name, raw_ticker,
            raw_isin, raw_currency, raw_quantity, raw_asset_class, extracted_by, extraction_confidence,
            proposed_name, proposed_ticker, proposed_isin, proposed_exchange, proposed_currency,
            proposed_asset_class, provider, provider_symbol, mapping_status, review_note, created_at, updated_at
        ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
        ON CONFLICT(candidate_id) DO UPDATE SET
            raw_name=excluded.raw_name, raw_ticker=excluded.raw_ticker, raw_isin=excluded.raw_isin,
            raw_currency=excluded.raw_currency, raw_quantity=excluded.raw_quantity, raw_asset_class=excluded.raw_asset_class,
            extraction_confidence=excluded.extraction_confidence, proposed_name=excluded.proposed_name,
            proposed_ticker=excluded.proposed_ticker, proposed_isin=excluded.proposed_isin,
            proposed_exchange=excluded.proposed_exchange, proposed_currency=excluded.proposed_currency,
            proposed_asset_class=excluded.proposed_asset_class, provider=excluded.provider,
            provider_symbol=excluded.provider_symbol,
            mapping_status=CASE WHEN instrument_import_candidates.mapping_status='confirmed' THEN 'confirmed' ELSE excluded.mapping_status END,
            review_note=excluded.review_note, updated_at=excluded.updated_at
        """,
        (
            cid, inp.source_file, inp.platform, inp.account, inp.source_row_number, inp.raw_name, inp.raw_ticker,
            inp.raw_isin, inp.raw_currency, inp.raw_quantity, inp.raw_asset_class, inp.extracted_by, inp.extraction_confidence,
            inp.proposed_name, inp.proposed_ticker, inp.proposed_isin, inp.proposed_exchange, inp.proposed_currency,
            inp.proposed_asset_class, inp.provider, inp.provider_symbol, inp.mapping_status, inp.review_note, now, now,
        ),
    )
    if inp.mapping_status != "confirmed":
        create_alert(conn, priority="warnung", category="equity", entity_type="instrument_import_candidate", entity_id=cid, rule_id="instrument_candidate_review_required", message="Instrument import candidate requires manual review before booking.", evidence={"platform": inp.platform, "has_isin": bool(inp.raw_isin or inp.proposed_isin)}, fingerprint="instrument_candidate_review_required")
    return cid


def _summary_from_inputs(items: list[CandidateInput], *, qwen_used: bool, qwen_available: bool, agree: int, disagree: int, warnings: int, errors: int, dry_run: bool) -> CandidateSummary:
    return CandidateSummary(
        candidates_total=len(items),
        with_isin=sum(1 for c in items if c.proposed_isin or c.raw_isin),
        without_isin=sum(1 for c in items if not (c.proposed_isin or c.raw_isin)),
        with_ticker=sum(1 for c in items if c.proposed_ticker or c.raw_ticker),
        with_currency=sum(1 for c in items if c.proposed_currency or c.raw_currency),
        exact_isin_match=sum(1 for c in items if c.mapping_status == "exact_isin_match"),
        probable=sum(1 for c in items if c.mapping_status == "probable"),
        needs_manual_review=sum(1 for c in items if c.mapping_status == "needs_manual_review"),
        confirmed=sum(1 for c in items if c.mapping_status == "confirmed"),
        rejected=sum(1 for c in items if c.mapping_status == "rejected"),
        importable_confirmed=sum(1 for c in items if c.mapping_status == "confirmed" and c.raw_quantity and (c.proposed_currency or c.raw_currency)),
        blocked=sum(1 for c in items if c.mapping_status != "confirmed"),
        qwen_used=qwen_used,
        qwen_available=qwen_available,
        parser_qwen_agree=agree,
        parser_qwen_disagree=disagree,
        warnings=warnings,
        errors=errors,
        dry_run=dry_run,
    )


def stage_docx_candidates(conn: Connection, files: Sequence[str | Path], *, use_qwen: bool = False, dry_run: bool = False) -> CandidateSummary:
    qwen_used = False
    qwen_available = False
    parser_signatures: set[tuple[str, str, str]] = set()
    qwen_signatures: set[tuple[str, str, str]] = set()
    staged_items: list[CandidateInput] = []
    warnings = 0
    errors = 0
    for path in files:
        p = Path(path)
        platform = "true_wealth" if "true" in p.name.lower() else ("postfinance" if "post" in p.name.lower() else "unknown")
        try:
            parser_candidates = deterministic_parse_docx(p, platform=platform)
            for c in parser_candidates:
                staged_items.append(c)
                parser_signatures.add(((c.proposed_isin or c.raw_isin).upper(), (c.proposed_ticker or c.raw_ticker).upper(), (c.proposed_currency or c.raw_currency).upper()))
                upsert_candidate(conn, c, dry_run=dry_run)
            if use_qwen:
                lines = extract_docx_text(p)
                q_candidates, available = qwen_extract_candidates(lines, source_file=p.name, platform=platform)
                qwen_available = qwen_available or available
                qwen_used = qwen_used or bool(q_candidates)
                for c in q_candidates:
                    staged_items.append(c)
                    qwen_signatures.add(((c.proposed_isin or c.raw_isin).upper(), (c.proposed_ticker or c.raw_ticker).upper(), (c.proposed_currency or c.raw_currency).upper()))
                    upsert_candidate(conn, c, dry_run=dry_run)
        except Exception:
            errors += 1
    if not dry_run:
        conn.commit()
    summary = summarize_candidates(conn, dry_run=dry_run)
    agree = len((parser_signatures - {('', '', '')}) & (qwen_signatures - {('', '', '')}))
    disagree = len((parser_signatures ^ qwen_signatures) - {('', '', '')}) if qwen_signatures else 0
    if dry_run:
        return _summary_from_inputs(staged_items, qwen_used=qwen_used, qwen_available=qwen_available, agree=agree, disagree=disagree, warnings=warnings, errors=errors, dry_run=True)
    return CandidateSummary(**{**summary.__dict__, "qwen_used": qwen_used, "qwen_available": qwen_available, "parser_qwen_agree": agree, "parser_qwen_disagree": disagree, "warnings": warnings, "errors": errors, "dry_run": dry_run})


def summarize_candidates(conn: Connection, *, dry_run: bool = False) -> CandidateSummary:
    def count(where: str = "1=1") -> int:
        return int(conn.execute(f"SELECT COUNT(*) AS c FROM instrument_import_candidates WHERE {where}").fetchone()["c"] or 0)
    return CandidateSummary(
        candidates_total=count(),
        with_isin=count("COALESCE(proposed_isin, raw_isin, '')<>''"),
        without_isin=count("COALESCE(proposed_isin, raw_isin, '')=''"),
        with_ticker=count("COALESCE(proposed_ticker, raw_ticker, '')<>''"),
        with_currency=count("COALESCE(proposed_currency, raw_currency, '')<>''"),
        exact_isin_match=count("mapping_status='exact_isin_match'"),
        probable=count("mapping_status='probable'"),
        needs_manual_review=count("mapping_status='needs_manual_review'"),
        confirmed=count("mapping_status='confirmed'"),
        rejected=count("mapping_status='rejected'"),
        importable_confirmed=count("mapping_status='confirmed' AND COALESCE(raw_quantity,'')<>'' AND COALESCE(proposed_currency,raw_currency,'')<>''"),
        blocked=count("mapping_status!='confirmed'"),
        dry_run=dry_run,
    )


def resolve_candidate(conn: Connection, candidate_id: str) -> dict[str, str]:
    row = conn.execute("SELECT * FROM instrument_import_candidates WHERE candidate_id=?", (candidate_id,)).fetchone()
    if row is None:
        raise ValueError("candidate not found")
    item = dict(row)
    isin = (item.get("proposed_isin") or item.get("raw_isin") or "").upper()
    ticker = (item.get("proposed_ticker") or item.get("raw_ticker") or "").upper()
    exchange = (item.get("proposed_exchange") or "").upper()
    currency = (item.get("proposed_currency") or item.get("raw_currency") or "").upper()
    status = "needs_manual_review"
    if isin:
        matches = conn.execute("SELECT COUNT(*) AS c FROM instrument_catalog_entries WHERE isin=?", (isin,)).fetchone()["c"]
        status = "exact_isin_match" if int(matches or 0) == 1 else "needs_manual_review"
    elif ticker and exchange and currency:
        matches = conn.execute("SELECT COUNT(*) AS c FROM instrument_catalog_entries WHERE upper(ticker)=? AND upper(exchange)=? AND upper(COALESCE(trading_currency,instrument_currency,''))=?", (ticker, exchange, currency)).fetchone()["c"]
        status = "probable" if int(matches or 0) == 1 else "needs_manual_review"
    conn.execute("UPDATE instrument_import_candidates SET mapping_status=?, updated_at=? WHERE candidate_id=? AND mapping_status!='confirmed'", (status, utc_now(), candidate_id))
    conn.commit()
    return {"candidate_id": candidate_id, "mapping_status": status}


def confirm_candidate(conn: Connection, candidate_id: str, *, note: str, proposed_isin: str | None = None, proposed_ticker: str | None = None, proposed_exchange: str | None = None, proposed_currency: str | None = None) -> None:
    note = _norm_space(note)
    if not note:
        raise ValueError("review note is required")
    row = conn.execute("SELECT * FROM instrument_import_candidates WHERE candidate_id=?", (candidate_id,)).fetchone()
    if row is None:
        raise ValueError("candidate not found")
    isin = (proposed_isin if proposed_isin is not None else row["proposed_isin"] or row["raw_isin"] or "").strip().upper()
    ticker = (proposed_ticker if proposed_ticker is not None else row["proposed_ticker"] or row["raw_ticker"] or "").strip().upper()
    exchange = (proposed_exchange if proposed_exchange is not None else row["proposed_exchange"] or "").strip().upper()
    currency = (proposed_currency if proposed_currency is not None else row["proposed_currency"] or row["raw_currency"] or "").strip().upper()
    if not isin and not (ticker and exchange and currency):
        raise ValueError("confirmed candidate requires ISIN or ticker+exchange+currency")
    conn.execute(
        "UPDATE instrument_import_candidates SET proposed_isin=?, proposed_ticker=?, proposed_exchange=?, proposed_currency=?, mapping_status='confirmed', review_note=?, updated_at=? WHERE candidate_id=?",
        (isin, ticker, exchange, currency, note, utc_now(), candidate_id),
    )
    record_audit_event(conn, source="instrument_import_review", action="confirm_candidate", entity_type="instrument_import_candidate", entity_id=candidate_id, new_values={"mapping_status": "confirmed", "has_isin": bool(isin)}, user_text_note=note, confirmed=True, created_by="dashboard")
    conn.commit()


def reject_candidate(conn: Connection, candidate_id: str, *, note: str) -> None:
    note = _norm_space(note)
    if not note:
        raise ValueError("review note is required")
    conn.execute("UPDATE instrument_import_candidates SET mapping_status='rejected', review_note=?, updated_at=? WHERE candidate_id=?", (note, utc_now(), candidate_id))
    record_audit_event(conn, source="instrument_import_review", action="reject_candidate", entity_type="instrument_import_candidate", entity_id=candidate_id, new_values={"mapping_status": "rejected"}, user_text_note=note, confirmed=True, created_by="dashboard")
    conn.commit()



def defer_candidate(conn: Connection, candidate_id: str, *, note: str = "later review") -> None:
    note = _norm_space(note) or "later review"
    row = conn.execute("SELECT candidate_id FROM instrument_import_candidates WHERE candidate_id=?", (candidate_id,)).fetchone()
    if row is None:
        raise ValueError("candidate not found")
    conn.execute("UPDATE instrument_import_candidates SET mapping_status='needs_manual_review', review_note=?, updated_at=? WHERE candidate_id=? AND mapping_status!='confirmed'", (note, utc_now(), candidate_id))
    record_audit_event(conn, source="instrument_import_review", action="defer_candidate", entity_type="instrument_import_candidate", entity_id=candidate_id, new_values={"mapping_status": "needs_manual_review"}, user_text_note=note, confirmed=True, created_by="dashboard")
    conn.commit()


def import_confirmed_candidate_snapshot(conn: Connection, candidate_id: str, *, account_id: str, snapshot_date: str, note: str, confirm: bool) -> str:
    if not confirm:
        raise ValueError("explicit confirm is required")
    row = conn.execute("SELECT * FROM instrument_import_candidates WHERE candidate_id=?", (candidate_id,)).fetchone()
    if row is None:
        raise ValueError("candidate not found")
    if row["mapping_status"] != "confirmed":
        raise ValueError("only confirmed candidates can be imported")
    isin = (row["proposed_isin"] or row["raw_isin"] or "").strip().upper()
    ticker = (row["proposed_ticker"] or row["raw_ticker"] or "").strip().upper()
    exchange = (row["proposed_exchange"] or "").strip().upper()
    currency = (row["proposed_currency"] or row["raw_currency"] or "CHF").strip().upper()
    name = _norm_space(row["proposed_name"] or row["raw_name"] or "Reviewed instrument")
    if not isin and not (ticker and exchange and currency):
        raise ValueError("confirmed import requires ISIN or ticker+exchange+currency")
    try:
        Decimal(str(row["raw_quantity"]).replace("'", "").replace(" ", "").replace(",", "."))
    except (InvalidOperation, ValueError) as exc:
        raise ValueError("candidate quantity must be Decimal text") from exc
    cid = upsert_catalog_entry(conn, CatalogEntryInput(asset_class="etf" if row["proposed_asset_class"] == "etf" else "equity", name=name, isin=isin, ticker=ticker, exchange=exchange, trading_currency=currency, instrument_currency=currency, provider=row["provider"] or "manual", provider_symbol=row["provider_symbol"] or None, source="broker_review", source_confidence="high" if isin else "medium"), note=note)
    result = add_manual_position_from_catalog(conn, catalog_entry_id=cid, account_id=account_id, position_type="initial_snapshot", quantity_text=str(row["raw_quantity"]).replace("'", "").replace(" ", "").replace(",", "."), trade_date=snapshot_date, currency=currency, cost_basis_original_text=None, fx_status="not_needed" if currency == "CHF" else "missing", note=note, confirm=True)
    conn.execute("UPDATE transactions SET source_type='broker_import_reviewed_snapshot', source_id=? WHERE transaction_id=?", (candidate_id, result.transaction_id))
    record_audit_event(conn, source="instrument_import_review", action="import_confirmed_candidate_snapshot", entity_type="transaction", entity_id=result.transaction_id, new_values={"candidate_id": candidate_id, "cost_basis_uncertain": True}, user_text_note=note, confirmed=True, created_by="dashboard")
    conn.commit()
    return result.transaction_id



def _status_counts(conn: Connection) -> dict[str, int]:
    row = conn.execute(
        """
        SELECT
            COUNT(*) AS total,
            SUM(CASE WHEN mapping_status='exact_isin_match' THEN 1 ELSE 0 END) AS exact,
            SUM(CASE WHEN mapping_status='probable' THEN 1 ELSE 0 END) AS probable,
            SUM(CASE WHEN mapping_status='needs_manual_review' THEN 1 ELSE 0 END) AS needs_review,
            SUM(CASE WHEN mapping_status='confirmed' THEN 1 ELSE 0 END) AS confirmed
        FROM instrument_import_candidates
        """
    ).fetchone()
    return {"total": int(row["total"] or 0), "exact": int(row["exact"] or 0), "probable": int(row["probable"] or 0), "needs_review": int(row["needs_review"] or 0), "confirmed": int(row["confirmed"] or 0)}


def _qwen_chat_json(payload: dict[str, object], *, endpoint: str, timeout: int) -> tuple[dict[str, object] | None, bool]:
    req = request.Request(f"{endpoint.rstrip('/')}/chat/completions", data=json.dumps(payload).encode("utf-8"), headers={"Content-Type": "application/json"})
    try:
        with request.urlopen(req, timeout=timeout) as resp:
            data = json.loads(resp.read().decode("utf-8"))
    except (OSError, error.URLError, TimeoutError, json.JSONDecodeError):
        return None, False
    content = data.get("choices", [{}])[0].get("message", {}).get("content", "")
    match = re.search(r"\{.*\}", str(content), re.S)
    if not match:
        return {}, True
    try:
        parsed = json.loads(match.group(0))
    except json.JSONDecodeError:
        return {}, True
    return parsed, True


def qwen_resolve_staged_candidates(conn: Connection, *, endpoint: str | None = None, model: str | None = None, limit: int | None = None, dry_run: bool = False, timeout: int = 45, batch_size: int = 5) -> QwenResolveSummary:
    """Run local Qwen over sanitized staged candidates only. No quantities, values, accounts or ledger writes."""
    endpoint = endpoint or os.getenv("QWEN_OPENAI_BASE_URL") or os.getenv("LOCAL_QWEN_BASE_URL") or "http://100.101.173.25:11435/v1"
    model = model or os.getenv("QWEN_MODEL", "qwen-3.6-agent")
    before = _status_counts(conn)
    rows = conn.execute(
        """
        SELECT candidate_id, platform, source_row_number, raw_name, raw_ticker, raw_isin, raw_currency,
               raw_asset_class, proposed_name, proposed_ticker, proposed_isin, proposed_exchange,
               proposed_currency, proposed_asset_class, mapping_status, review_note
        FROM instrument_import_candidates
        WHERE mapping_status NOT IN ('confirmed','rejected')
        ORDER BY platform, source_row_number
        """
    ).fetchall()
    if limit is not None:
        rows = rows[:limit]
    sanitized = []
    for r in rows:
        d = dict(r)
        sanitized.append({
            "candidate_id": d["candidate_id"],
            "platform": d.get("platform") or "",
            "source_row_id": d.get("source_row_number") or "",
            "product_name": d.get("proposed_name") or d.get("raw_name") or "",
            "ticker": d.get("proposed_ticker") or d.get("raw_ticker") or "",
            "isin": d.get("proposed_isin") or d.get("raw_isin") or "",
            "currency": d.get("proposed_currency") or d.get("raw_currency") or "",
            "asset_class": d.get("proposed_asset_class") or d.get("raw_asset_class") or "",
            "exchange": d.get("proposed_exchange") or "",
            "current_status": d.get("mapping_status") or "needs_manual_review",
        })
    all_updates: list[dict[str, object]] = []
    available = True
    batch_size = max(1, batch_size)
    for start_idx in range(0, len(sanitized), batch_size):
        batch = sanitized[start_idx:start_idx + batch_size]
        payload = {
            "model": model,
            "temperature": 0,
            "max_tokens": 4096,
            "stream": False,
            "messages": [
                {"role": "system", "content": "You are a conservative equity/ETF import review assistant. Return JSON only. Never invent ISINs, quantities, cost basis, transaction history, account numbers or valuations. If ISIN is empty, do not add one. Ticker alone is not enough for confirmation. Use statuses exact_isin_match, probable, needs_manual_review only."},
                {"role": "user", "content": json.dumps({"task": "Classify sanitized instrument import candidates for human review. Do not confirm anything. Return {updates:[{candidate_id, proposed_name, proposed_ticker, proposed_isin, proposed_exchange, proposed_currency, proposed_asset_class, mapping_status, review_note}]}. If an ISIN is not present in input, proposed_isin must be empty.", "candidates": batch}, ensure_ascii=False)},
            ],
        }
        parsed, batch_available = _qwen_chat_json(payload, endpoint=endpoint, timeout=timeout)
        if not batch_available:
            available = False
            break
        if isinstance(parsed, dict):
            updates = parsed.get("updates", [])
            if isinstance(updates, list):
                all_updates.extend([u for u in updates if isinstance(u, dict)])
    if not available:
        after = _status_counts(conn)
        return QwenResolveSummary(qwen_available=False, qwen_used=False, model=model, candidates_total=before["total"], exact_before=before["exact"], exact_after=after["exact"], probable_before=before["probable"], probable_after=after["probable"], needs_review_before=before["needs_review"], needs_review_after=after["needs_review"], confirmed_after=after["confirmed"], manual_review_remaining=after["needs_review"], errors=1, dry_run=dry_run)
    parsed = {"updates": all_updates}
    row_by_id = {str(r["candidate_id"]): dict(r) for r in rows}
    improved = 0
    warnings = 0
    for upd in (parsed or {}).get("updates", []) if isinstance(parsed, dict) else []:
        cid = str(upd.get("candidate_id") or "")
        if cid not in row_by_id:
            warnings += 1
            continue
        current = row_by_id[cid]
        existing_isin = (current.get("proposed_isin") or current.get("raw_isin") or "").strip().upper()
        q_isin = str(upd.get("proposed_isin") or "").strip().upper()
        if q_isin and q_isin != existing_isin:
            # Qwen may not invent or alter ISIN. Keep it out and leave manual review note.
            q_isin = existing_isin
            warnings += 1
        q_ticker = str(upd.get("proposed_ticker") or current.get("proposed_ticker") or current.get("raw_ticker") or "").strip().upper()
        q_exchange = str(upd.get("proposed_exchange") or current.get("proposed_exchange") or "").strip().upper()
        q_currency = str(upd.get("proposed_currency") or current.get("proposed_currency") or current.get("raw_currency") or "").strip().upper()
        q_name = _norm_space(str(upd.get("proposed_name") or current.get("proposed_name") or current.get("raw_name") or ""))[:240]
        q_asset_class = _norm_space(str(upd.get("proposed_asset_class") or current.get("proposed_asset_class") or current.get("raw_asset_class") or "equity")).lower()
        requested_status = str(upd.get("mapping_status") or "needs_manual_review").strip()
        if existing_isin:
            status = "exact_isin_match"
        elif q_ticker and q_exchange and q_currency and requested_status == "probable":
            status = "probable"
        else:
            status = "needs_manual_review"
        note = _norm_space(str(upd.get("review_note") or "Qwen review helper: manual review required; no booking."))[:500]
        if status != current.get("mapping_status"):
            improved += 1
        if not dry_run:
            conn.execute(
                """
                UPDATE instrument_import_candidates
                SET proposed_name=?, proposed_ticker=?, proposed_isin=?, proposed_exchange=?, proposed_currency=?,
                    proposed_asset_class=?, mapping_status=?, review_note=?, updated_at=?
                WHERE candidate_id=? AND mapping_status NOT IN ('confirmed','rejected')
                """,
                (q_name, q_ticker, q_isin, q_exchange, q_currency, q_asset_class, status, note, utc_now(), cid),
            )
    if not dry_run:
        conn.commit()
    after = _status_counts(conn)
    return QwenResolveSummary(qwen_available=True, qwen_used=True, model=model, candidates_total=after["total"], exact_before=before["exact"], exact_after=after["exact"], probable_before=before["probable"], probable_after=after["probable"], needs_review_before=before["needs_review"], needs_review_after=after["needs_review"], confirmed_after=after["confirmed"], improved=improved, manual_review_remaining=after["needs_review"], warnings=warnings, errors=0, dry_run=dry_run)
