#!/usr/bin/env python3
"""Capordia migration-v1 : exercice CPU local, Python >= 3.10, sans réseau.

Les résultats acceptés et l'état terminé sont dans la même transaction SQLite.
Un verrou local et une source clôturée aident la démonstration ; ils ne sont pas
un mécanisme de fencing distribué. Voir README.md et PROCEDURE.md.
"""

import argparse
import csv
import hashlib
import io
import json
import os
import platform
import re
import shutil
import sqlite3
import subprocess
import sys
from contextlib import contextmanager
from datetime import datetime, timezone
from pathlib import Path

VERSION = "ferroscale-migration-1.0.0"
RECETTE = "controle-coupons-entiers-v1"
BASE = Path(__file__).resolve().parent


class Refus(Exception):
    pass


class Interruption(Exception):
    pass


def canon(value):
    return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":"))


def sha(data):
    return hashlib.sha256(data).hexdigest()


def ecrire_json(path, value):
    path.write_text(json.dumps(value, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")


def lire_jeu(path):
    raw = (path / "manifeste.json").read_bytes()
    manifest = json.loads(raw)
    if manifest.get("version") != 1 or manifest.get("recette") != RECETTE:
        raise Refus("MANIFESTE_INCOMPATIBLE")
    lots = manifest.get("lots")
    if not isinstance(lots, list) or not lots:
        raise Refus("MANIFESTE_VIDE")
    ids = set()
    for lot in lots:
        ident = lot.get("id", "")
        if not re.fullmatch(r"LOT-[0-9]{3}", ident):
            raise Refus("IDENTIFIANT_INVALIDE")
        if ident in ids:
            raise Refus("DOUBLON_MANIFESTE: " + ident)
        ids.add(ident)
        if lot.get("fichier") != "donnees/" + ident + ".csv":
            raise Refus("CHEMIN_INVALIDE: " + ident)
        data = (path / lot["fichier"]).read_bytes()
        if sha(data) != lot.get("sha256"):
            raise Refus("ENTREE_CORROMPUE: " + ident)
        calculer(ident, data)
    return manifest, sha(raw)


def calculer(ident, data):
    reader = csv.DictReader(io.StringIO(data.decode("utf-8")))
    if reader.fieldnames != ["coupon", "cible_um", "mesure_um"]:
        raise Refus("COLONNES_INVALIDES: " + ident)
    coupons, ecarts = set(), []
    for row in reader:
        coupon = row["coupon"]
        if not re.fullmatch(r"CP-[0-9]{3}", coupon) or coupon in coupons:
            raise Refus("COUPON_INVALIDE_OU_DOUBLE: " + ident)
        coupons.add(coupon)
        try:
            cible, mesure = int(row["cible_um"]), int(row["mesure_um"])
        except (TypeError, ValueError) as exc:
            raise Refus("MESURE_INVALIDE: " + ident) from exc
        if not 1 <= cible <= 100000 or not 1 <= mesure <= 100000:
            raise Refus("MESURE_HORS_EXERCICE: " + ident)
        ecarts.append(abs(mesure - cible))
    if not ecarts:
        raise Refus("LOT_VIDE: " + ident)
    return {
        "lot": ident, "recette": RECETTE, "entree_sha256": sha(data),
        "coupons": len(ecarts), "ecart_absolu_total_um": sum(ecarts),
        "ecart_max_um": max(ecarts), "dans_tolerance_15_um": sum(x <= 15 for x in ecarts),
    }


@contextmanager
def verrou(dossier):
    path = dossier / "worker.lock"
    try:
        fd = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
    except FileExistsError as exc:
        raise Refus("WORKER_OCCUPE: verrou present ; ne pas le supprimer sans verifier l'arret du processus") from exc
    try:
        os.write(fd, str(os.getpid()).encode("ascii"))
        os.close(fd)
        yield
    finally:
        path.unlink(missing_ok=True)


@contextmanager
def connexion(dossier):
    db = dossier / "registre.sqlite"
    if not db.is_file():
        raise Refus("REGISTRE_ABSENT")
    con = sqlite3.connect(db, timeout=2, isolation_level="DEFERRED")
    con.row_factory = sqlite3.Row
    con.execute("PRAGMA foreign_keys=ON")
    con.execute("PRAGMA synchronous=FULL")
    try:
        yield con
    finally:
        con.close()


def meta(con, key):
    row = con.execute("SELECT valeur FROM meta WHERE cle=?", (key,)).fetchone()
    if row is None:
        raise Refus("META_ABSENTE: " + key)
    return row[0]


def autoriser(con):
    if meta(con, "cloture") != "non":
        raise Refus("SOURCE_CLOTUREE: cette copie ne doit plus produire")


def controler_jeu(con, dossier):
    manifest, digest = lire_jeu(dossier)
    if meta(con, "manifeste_sha256") != digest:
        raise Refus("MANIFESTE_MODIFIE")
    rows = [dict(r) for r in con.execute("SELECT id,fichier,entree_sha256 FROM lots ORDER BY id")]
    expected = [{"id": lot["id"], "fichier": lot["fichier"], "entree_sha256": lot["sha256"]} for lot in sorted(manifest["lots"], key=lambda x: x["id"])]
    if rows != expected:
        raise Refus("REGISTRE_ENTREES_DIVERGENT")
    return manifest


def instantane(con):
    lots = [dict(row) for row in con.execute("SELECT l.id,l.etat,l.tentatives,l.entree_sha256,a.sortie_sha256 FROM lots l LEFT JOIN acceptes a ON a.lot_id=l.id ORDER BY l.id")]
    counts = {state: sum(lot["etat"] == state for lot in lots) for state in ("termine", "en_cours", "a_reprendre")}
    return {"version": VERSION, "proprietaire": meta(con, "proprietaire"), "cloture": meta(con, "cloture"), "etats": counts, "lots": lots}


def exporter(dossier):
    with connexion(dossier) as con:
        snapshot = instantane(con)
        outputs = [json.loads(row[0]) for row in con.execute("SELECT contenu FROM acceptes ORDER BY lot_id")]
    ecrire_json(dossier / "etat.json", snapshot)
    ecrire_json(dossier / "resultats.json", outputs)
    with (dossier / "registre-lots.csv").open("w", encoding="utf-8", newline="") as file:
        writer = csv.DictWriter(file, fieldnames=["id", "etat", "tentatives", "entree_sha256", "sortie_sha256"])
        writer.writeheader()
        writer.writerows(snapshot["lots"])
    return snapshot


def initialiser(dossier, jeu, proprietaire):
    manifest, digest = lire_jeu(jeu)
    if dossier.exists():
        raise Refus("DOSSIER_EXISTANT: choisir un nouveau nom")
    dossier.mkdir(parents=True)
    shutil.copy2(jeu / "manifeste.json", dossier / "manifeste.json")
    (dossier / "donnees").mkdir()
    for lot in manifest["lots"]:
        shutil.copy2(jeu / lot["fichier"], dossier / lot["fichier"])
    con = sqlite3.connect(dossier / "registre.sqlite")
    try:
        con.executescript("""
            PRAGMA journal_mode=DELETE;
            PRAGMA synchronous=FULL;
            PRAGMA foreign_keys=ON;
            CREATE TABLE meta(cle TEXT PRIMARY KEY NOT NULL, valeur TEXT NOT NULL);
            CREATE TABLE lots(
                id TEXT PRIMARY KEY NOT NULL, fichier TEXT NOT NULL,
                entree_sha256 TEXT NOT NULL,
                etat TEXT NOT NULL CHECK(etat IN ('a_reprendre','en_cours','termine')),
                tentatives INTEGER NOT NULL DEFAULT 0);
            CREATE TABLE acceptes(
                lot_id TEXT PRIMARY KEY NOT NULL REFERENCES lots(id),
                contenu TEXT NOT NULL, sortie_sha256 TEXT NOT NULL);
        """)
        with con:
            con.executemany("INSERT INTO meta VALUES(?,?)", [("version", VERSION), ("manifeste_sha256", digest), ("proprietaire", proprietaire), ("cloture", "non")])
            con.executemany("INSERT INTO lots(id,fichier,entree_sha256,etat) VALUES(?,?,?,'a_reprendre')", [(lot["id"], lot["fichier"], lot["sha256"]) for lot in manifest["lots"]])
    finally:
        con.close()
    return exporter(dossier)


def executer(dossier, interrompre=None, apres_insertion=False):
    with verrou(dossier):
        try:
            with connexion(dossier) as con:
                autoriser(con)
                controler_jeu(con, dossier)
                verifier(dossier)
                if interrompre and not con.execute("SELECT 1 FROM lots WHERE id=? AND etat='a_reprendre'", (interrompre,)).fetchone():
                    raise Refus("LOT_DE_COUPURE_NON_DISPONIBLE")
                if apres_insertion and not interrompre:
                    raise Refus("LOT_DE_COUPURE_REQUIS")
                if con.execute("SELECT 1 FROM lots WHERE etat='en_cours'").fetchone():
                    raise Refus("REPRISE_A_PREPARER: examiner puis utiliser preparer-reprise")
                rows = con.execute("SELECT * FROM lots WHERE etat='a_reprendre' ORDER BY id").fetchall()
                for row in rows:
                    with con:
                        con.execute("UPDATE lots SET etat='en_cours',tentatives=tentatives+1 WHERE id=?", (row["id"],))
                    data = (dossier / row["fichier"]).read_bytes()
                    if sha(data) != row["entree_sha256"]:
                        raise Refus("ENTREE_CORROMPUE: " + row["id"])
                    output = canon(calculer(row["id"], data))
                    if row["id"] == interrompre and not apres_insertion:
                        raise Interruption("COUPURE_SIMULEE avant acceptation de " + row["id"])
                    with con:
                        con.execute("INSERT INTO acceptes VALUES(?,?,?)", (row["id"], output, sha(output.encode("utf-8"))))
                        if row["id"] == interrompre:
                            raise Interruption("COUPURE_SIMULEE dans la transaction de " + row["id"])
                        con.execute("UPDATE lots SET etat='termine' WHERE id=?", (row["id"],))
        finally:
            snapshot = exporter(dossier)
    return snapshot


def preparer_reprise(dossier):
    with verrou(dossier), connexion(dossier) as con:
        autoriser(con)
        controler_jeu(con, dossier)
        if con.execute("SELECT 1 FROM acceptes a JOIN lots l ON l.id=a.lot_id WHERE l.etat != 'termine'").fetchone():
            raise Refus("REGISTRE_INCOHERENT: ne pas recommencer automatiquement")
        with con:
            con.execute("UPDATE lots SET etat='a_reprendre' WHERE etat='en_cours'")
    return exporter(dossier)


def verifier(dossier, complet=False):
    with connexion(dossier) as con:
        if con.execute("PRAGMA integrity_check").fetchone()[0] != "ok":
            raise Refus("SQLITE_INTEGRE_NON")
        if con.execute("PRAGMA foreign_key_check").fetchall():
            raise Refus("REFERENCE_INCOHERENTE")
        controler_jeu(con, dossier)
        rows = con.execute("SELECT l.*,a.contenu,a.sortie_sha256 FROM lots l LEFT JOIN acceptes a ON a.lot_id=l.id ORDER BY l.id").fetchall()
        accepted = []
        for row in rows:
            if (row["etat"] == "termine") != (row["contenu"] is not None):
                raise Refus("ETAT_ACCEPTATION_INCOHERENT: " + row["id"])
            if row["contenu"] is not None:
                expected = canon(calculer(row["id"], (dossier / row["fichier"]).read_bytes()))
                if row["contenu"] != expected or sha(expected.encode("utf-8")) != row["sortie_sha256"]:
                    raise Refus("SORTIE_CORROMPUE: " + row["id"])
                accepted.append(json.loads(expected))
        if complet and len(accepted) != len(rows):
            raise Refus("TRAVAIL_INCOMPLET")
        return {"lots": len(rows), "acceptes": len(accepted), "resultats_sha256": sha(canon(accepted).encode("utf-8")), "integre": True}


def basculer(source, destination, proprietaire):
    if destination.exists():
        raise Refus("DESTINATION_EXISTANTE")
    if source == destination or source in destination.parents or destination in source.parents:
        raise Refus("DOSSIERS_IMBRIQUES")
    with verrou(source):
        verifier(source)
        with connexion(source) as con:
            autoriser(con)
            # Clôturer avant de copier : un échec de copie arrête le travail.
            with con:
                con.execute("UPDATE meta SET valeur='oui' WHERE cle='cloture'")
        exporter(source)
        # Le verrou interdit notre worker ; toutes les connexions sont fermées.
        shutil.copytree(source, destination, ignore=shutil.ignore_patterns("worker.lock"))
        verifier(destination)
        with connexion(destination) as con:
            with con:
                con.execute("UPDATE meta SET valeur=? WHERE cle='proprietaire'", (proprietaire,))
                con.execute("UPDATE meta SET valeur='non' WHERE cle='cloture'")
                con.execute("INSERT OR REPLACE INTO meta VALUES('provenance',?)", (source.name,))
    return exporter(destination)


def comparer(reference, candidat):
    left, right = verifier(reference, True), verifier(candidat, True)
    if left != right:
        raise Refus("COMPARAISON_DIFFERENTE")
    return {"identiques": True, **right}


def essayer_doublon(dossier):
    with verrou(dossier), connexion(dossier) as con:
        autoriser(con)
        row = con.execute("SELECT * FROM acceptes ORDER BY lot_id LIMIT 1").fetchone()
        if row is None:
            raise Refus("AUCUNE_SORTIE_A_REJOUER")
        try:
            with con:
                con.execute("INSERT INTO acceptes VALUES(?,?,?)", tuple(row))
        except sqlite3.IntegrityError as exc:
            raise Refus("DOUBLON_REFUSE: " + row["lot_id"]) from exc
        raise RuntimeError("Le doublon a ete accepte : anomalie de schema")


def demonstration(root):
    if root.exists():
        raise Refus("DOSSIER_EXISTANT: choisir un nouveau nom")
    root.mkdir(parents=True)
    checks, commands = [], []

    def check(name, condition, details=None):
        if not condition:
            raise RuntimeError("ECHEC_DE_CONTROLE: " + name)
        checks.append({"controle": name, "reussi": True, "detail": details})

    def call(*args, attendu=0, motif=None):
        command = [sys.executable, str(Path(__file__).resolve()), *map(str, args)]
        done = subprocess.run(command, capture_output=True, text=True, encoding="utf-8", env={**os.environ, "PYTHONIOENCODING": "utf-8"})
        commands.append({"arguments": [str(a).replace(str(root), "ESSAI") for a in args], "code": done.returncode, "sortie": done.stdout.strip(), "erreur": done.stderr.strip()})
        if done.returncode != attendu or (motif and motif not in done.stderr):
            raise RuntimeError("Commande inattendue : " + canon(commands[-1]))
        return json.loads(done.stdout) if done.stdout.strip() else None

    ref, source, cible = (root / name for name in ("reference", "source", "cible"))
    for folder, owner in ((ref, "reference"), (source, "atelier-source")):
        call("initialiser", "--dossier", folder, "--proprietaire", owner)
    call("executer", "--dossier", ref)
    valid = call("verifier", "--dossier", ref, "--complet")
    check("Reference continue", valid["lots"] == valid["acceptes"] == 8, valid)
    call("executer", "--dossier", source, "--interrompre-sur", "LOT-004", attendu=75, motif="COUPURE_SIMULEE")
    interrupted = call("etat", "--dossier", source)
    check("Coupure : 3 termines, 1 en cours, 4 a reprendre", interrupted["etats"] == {"termine": 3, "en_cours": 1, "a_reprendre": 4}, interrupted["etats"])
    call("executer", "--dossier", source, attendu=2, motif="REPRISE_A_PREPARER")
    check("Relance ambigue refusee", True)
    call("basculer", "--source", source, "--destination", cible, "--proprietaire", "atelier-cible")
    call("executer", "--dossier", source, attendu=2, motif="SOURCE_CLOTUREE")
    check("Ancienne source cloturee", True)
    prepared = call("preparer-reprise", "--dossier", cible)
    check("Reconciliation : 3 termines, 0 en cours, 5 a reprendre", prepared["etats"] == {"termine": 3, "en_cours": 0, "a_reprendre": 5})
    call("executer", "--dossier", cible)
    identical = call("comparer", "--reference", ref, "--candidat", cible)
    check("Copie puis reprise : 8 sorties identiques sans lot perdu", identical["identiques"] and identical["acceptes"] == 8, identical)
    state = call("etat", "--dossier", cible)
    check("Lot interrompu retente, lots termines non recalcules", all(lot["tentatives"] == (2 if lot["id"] == "LOT-004" else 1) for lot in state["lots"]))
    before = call("verifier", "--dossier", cible, "--complet")
    call("executer", "--dossier", cible)
    after = call("verifier", "--dossier", cible, "--complet")
    check("Relance terminee idempotente", before == after)
    call("essayer-doublon", "--dossier", cible, attendu=2, motif="DOUBLON_REFUSE")
    check("Deuxieme acceptation du meme lot refusee", call("verifier", "--dossier", cible, "--complet") == before)
    retour = root / "retour"
    call("basculer", "--source", cible, "--destination", retour, "--proprietaire", "atelier-retour")
    call("executer", "--dossier", cible, attendu=2, motif="SOURCE_CLOTUREE")
    check("Retour avec etat recent et cible cloturee", call("comparer", "--reference", ref, "--candidat", retour)["identiques"])
    mauvais = root / "entree-corrompue"
    call("initialiser", "--dossier", mauvais)
    path = mauvais / "donnees" / "LOT-006.csv"
    path.write_bytes(path.read_bytes() + b"CP-999,3000,3001\n")
    call("executer", "--dossier", mauvais, attendu=2, motif="ENTREE_CORROMPUE")
    check("Entree corrompue refusee avant toute acceptation", call("etat", "--dossier", mauvais)["etats"]["termine"] == 0)
    double = root / "jeu-doublon"
    double.mkdir()
    shutil.copytree(BASE / "donnees", double / "donnees")
    manifest = json.loads((BASE / "manifeste.json").read_text(encoding="utf-8"))
    manifest["lots"].append(manifest["lots"][0])
    ecrire_json(double / "manifeste.json", manifest)
    call("initialiser", "--dossier", root / "refuse-doublon", "--jeu", double, attendu=2, motif="DOUBLON_MANIFESTE")
    check("Identifiant duplique dans le manifeste refuse", not (root / "refuse-doublon").exists())
    transaction = root / "transaction-interrompue"
    call("initialiser", "--dossier", transaction)
    call("executer", "--dossier", transaction, "--interrompre-sur", "LOT-004", "--apres-insertion", attendu=75, motif="dans la transaction")
    check("Rollback transaction : aucune sortie partiellement acceptee", call("verifier", "--dossier", transaction)["acceptes"] == 3)
    call("preparer-reprise", "--dossier", transaction)
    call("executer", "--dossier", transaction)
    check("Reprise apres rollback egale a la reference", call("comparer", "--reference", ref, "--candidat", transaction)["identiques"])
    locked = root / "verrouille"
    call("initialiser", "--dossier", locked)
    (locked / "worker.lock").write_text("verrou de test, aucun processus reel\n", encoding="utf-8")
    call("executer", "--dossier", locked, attendu=2, motif="WORKER_OCCUPE")
    check("Verrou local existant refuse", call("etat", "--dossier", locked)["etats"]["termine"] == 0)
    (locked / "worker.lock").unlink()
    output_bad = root / "sortie-corrompue"
    call("initialiser", "--dossier", output_bad)
    call("executer", "--dossier", output_bad)
    with connexion(output_bad) as con:
        with con:
            con.execute("UPDATE acceptes SET contenu='{}' WHERE lot_id='LOT-002'")
    call("verifier", "--dossier", output_bad, "--complet", attendu=2, motif="SORTIE_CORROMPUE")
    check("Alteration d'une sortie detectee", True)
    report = {
        "version": VERSION, "date_utc": datetime.now(timezone.utc).isoformat(),
        "environnement": {"python": platform.python_version(), "sqlite": sqlite3.sqlite_version, "systeme": platform.system(), "architecture": platform.machine()},
        "portee": "CPU et fichiers locaux ; aucune execution GPU, aucun test de panne electrique, aucun fencing distribue",
        "controles_reussis": len(checks), "commandes_executees": len(commands),
        "comparaison_principale": identical, "controles": checks, "commandes": commands,
    }
    ecrire_json(root / "rapport-verification.json", report)
    return {key: value for key, value in report.items() if key != "commandes"}


def main():
    parser = argparse.ArgumentParser(description=__doc__)
    sub = parser.add_subparsers(dest="commande", required=True)
    for name in ("initialiser", "executer", "etat", "preparer-reprise", "verifier", "essayer-doublon", "demonstration"):
        command = sub.add_parser(name)
        command.add_argument("--dossier", type=Path, required=True)
        if name == "initialiser":
            command.add_argument("--jeu", type=Path, default=BASE)
            command.add_argument("--proprietaire", default="atelier-local")
        if name == "executer":
            command.add_argument("--interrompre-sur", metavar="LOT-004")
            command.add_argument("--apres-insertion", action="store_true")
        if name == "verifier":
            command.add_argument("--complet", action="store_true")
    command = sub.add_parser("basculer")
    command.add_argument("--source", type=Path, required=True)
    command.add_argument("--destination", type=Path, required=True)
    command.add_argument("--proprietaire", required=True)
    command = sub.add_parser("comparer")
    command.add_argument("--reference", type=Path, required=True)
    command.add_argument("--candidat", type=Path, required=True)
    args = parser.parse_args()
    for key in ("dossier", "jeu", "source", "destination", "reference", "candidat"):
        if hasattr(args, key):
            setattr(args, key, getattr(args, key).resolve())
    try:
        if args.commande == "initialiser":
            result = initialiser(args.dossier, args.jeu, args.proprietaire)
        elif args.commande == "executer":
            result = executer(args.dossier, args.interrompre_sur, args.apres_insertion)
        elif args.commande == "etat":
            with connexion(args.dossier) as con:
                result = instantane(con)
        elif args.commande == "preparer-reprise":
            result = preparer_reprise(args.dossier)
        elif args.commande == "verifier":
            result = verifier(args.dossier, args.complet)
        elif args.commande == "essayer-doublon":
            result = essayer_doublon(args.dossier)
        elif args.commande == "demonstration":
            result = demonstration(args.dossier)
        elif args.commande == "basculer":
            result = basculer(args.source, args.destination, args.proprietaire)
        else:
            result = comparer(args.reference, args.candidat)
        print(json.dumps(result, ensure_ascii=False, indent=2))
        return 0
    except Refus as exc:
        print("REFUS: " + str(exc), file=sys.stderr)
        return 2
    except Interruption as exc:
        print(str(exc), file=sys.stderr)
        return 75
    except (OSError, sqlite3.Error, ValueError) as exc:
        print("ERREUR: " + str(exc), file=sys.stderr)
        return 1


if __name__ == "__main__":
    raise SystemExit(main())
