"""Invoices, receipts and credit notes by their gap-free counter (seq), into ERP journals.

  python3 invoices_by_seq.py          read what was issued since last time, book it, ack it,
                                      then pick up documents that changed after their ack
  python3 invoices_by_seq.py check    look for holes in the ERP's books and fill them

The last seq read per stream is kept in STATE_DIR/seq.json; the ERP stand-in in fake_erp.py.
"""
import hashlib
import json
import os
import sys

import fake_erp
from slflo import Slflo, TokenRefused, report

STATE_DIR = os.environ.get("STATE_DIR", "state")
STATE = os.path.join(STATE_DIR, "seq.json")
# stream -> (the kind in POST /acks, the journal kind, the amount field)
STREAMS = {
    "invoices": ("invoice", "invoice", lambda d: d["totals"]["total"]),
    "payments": ("payment", "payment", lambda d: d["amount"]),
    "returns": ("return", "return", lambda d: d["total"]),
}


def load_state():
    try:
        with open(STATE) as f:
            return json.load(f)
    except FileNotFoundError:
        return {stream: 0 for stream in STREAMS}  # 0 reads all history; GET /{stream}/cursor starts from now


def save_state(state):
    os.makedirs(STATE_DIR, exist_ok=True)
    with open(STATE, "w") as f:
        json.dump(state, f)


# region read
def read_new(api, erp, stream, last_seq):
    """Everything issued after last_seq, in issue order; stops at the first gap."""
    entity, kind, amount = STREAMS[stream]
    read, gap = 0, None
    while True:
        page = api.get(f"/{stream}", {"filter[seq_after]": last_seq, "per_page": 100})
        documents = page["data"]
        for document in documents:
            if document["seq"] != last_seq + 1:
                # Can't happen on our side: the counter has no holes. If it does, your
                # bookkeeping skipped one: stop here, alert, and re-read from last_seq.
                gap = {"expected": last_seq + 1, "got": document["seq"]}
                break
            journal = fake_erp.post_journal(erp, kind, document["id"], document["version"], amount(document), document["seq"])
            document["journal"] = journal
            last_seq = document["seq"]
            read += 1
        booked = [d for d in documents if "journal" in d]
        if booked:
            ack(api, entity, booked)
        if gap or not page["meta"]["has_more"]:
            return read, last_seq, gap
# endregion


def ack(api, entity, documents):
    items = [{"id": d["id"], "outcome": "ACCEPTED", "external_ref": d["journal"], "external_number": d["journal"]}
             for d in documents]
    # One key per set of (id, version): a retry replays, a later change gets a new key.
    answered = ",".join(f"{d['id']}@{d['version']}" for d in documents)
    key = f"{entity}-acks:" + hashlib.sha256(answered.encode()).hexdigest()[:40]
    answer = api.post("/acks", {"entity": entity, "items": items}, key=key)
    failed = [item for item in answer["data"] if item.get("error")]
    if failed:
        raise RuntimeError(f"acks refused: {failed}")


# region changes
def changed_since_ack(api, erp, stream):
    """A document that changes after its ack (a payment reversed, a cheque bounced, a return
    approved) is pending again. seq never sees it twice; filter[ack]=pending does."""
    entity, kind, amount = STREAMS[stream]
    changed = []
    for page in api.pages(f"/{stream}", {"filter[ack]": "pending"}):
        for document in page["data"]:
            journal = fake_erp.journal_of(erp, kind, document["id"])
            if journal is None:
                continue  # never booked: it arrives through its seq (a return once it is credited)
            if document.get("reversal"):
                fake_erp.post_journal(erp, f"{kind}-reversal", document["id"], document["version"], amount(document))
            document["journal"] = journal  # the answer stays the journal that booked it
            changed.append(document)
    if changed:
        ack(api, entity, changed)
    return len(changed)
# endregion


# region check
def check(api, erp, state):
    """Holes in the ERP's own books (a journal deleted by hand, a restore): re-read and re-book them."""
    found = {}
    for stream, (entity, kind, amount) in STREAMS.items():
        booked = set(fake_erp.booked_seqs(erp, kind))
        missing = [n for n in range(1, state.get(stream, 0) + 1) if n not in booked]
        for seq in missing:
            page = api.get(f"/{stream}", {"filter[seq_after]": seq - 1, "per_page": 1})
            document = page["data"][0]
            fake_erp.post_journal(erp, kind, document["id"], document["version"], amount(document), document["seq"])
        found[stream] = missing
    return found
# endregion


def main(mode):
    api, erp, state = Slflo(), fake_erp.connect(), load_state()
    summary = {}
    try:
        if mode == "check":
            summary["missing"] = check(api, erp, state)
        else:
            for stream in STREAMS:
                read, state[stream], gap = read_new(api, erp, stream, state.get(stream, 0))
                save_state(state)
                summary[stream] = {"read": read, "last_seq": state[stream], "gap": gap}
            summary["changed"] = {stream: changed_since_ack(api, erp, stream) for stream in STREAMS}
    except TokenRefused as stop:
        summary["stopped_by"] = stop.code
        report(summary)
        return 3
    report(summary)
    return 1 if any(isinstance(v, dict) and v.get("gap") for v in summary.values()) else 0


if __name__ == "__main__":
    sys.exit(main(sys.argv[1] if len(sys.argv) > 1 else "run"))
