"""Pull connector: every Slflo order waiting for the ERP becomes one ERP sales order, then an ack.

  python3 pull_orders.py           pull new orders, book them, ack them (run this on a schedule)
  python3 pull_orders.py pull      only pull and stage
  python3 pull_orders.py book      only book what is staged
  python3 pull_orders.py recover   the watermark or staging was lost: work from filter[ack]=pending

State (watermark and staging) is kept in STATE_DIR/staging.db; the ERP stand-in in fake_erp.py.
"""
import json
import os
import sqlite3
import sys

import fake_erp
from slflo import ApiError, Slflo, TokenRefused, report

STATE_DIR = os.environ.get("STATE_DIR", "state")


# region staging
def open_staging():
    os.makedirs(STATE_DIR, exist_ok=True)
    db = sqlite3.connect(os.path.join(STATE_DIR, "staging.db"))
    db.executescript(
        """
        create table if not exists watermark (stream text primary key, version integer not null);
        create table if not exists order_line (
            order_id text not null,
            version integer not null,
            line_number integer not null,
            document text not null,          -- the order as pulled
            booked integer not null default 0,
            primary key (order_id, version, line_number)   -- the unique key that makes re-reads harmless
        );
        """
    )
    return db


def stage(db, order):
    """One row per line under (id, version, line_number); a second copy inserts nothing."""
    lines = [line["line_number"] for line in order["lines"]] or [0]
    db.executemany(
        "insert or ignore into order_line (order_id, version, line_number, document) values (?, ?, ?, ?)",
        [(order["id"], order["version"], n, json.dumps(order)) for n in lines],
    )
# endregion


def watermark(db):
    row = db.execute("select version from watermark where stream = 'orders'").fetchone()
    return row[0] if row else 0


def save_watermark(db, version):
    db.execute(
        "insert into watermark (stream, version) values ('orders', ?) "
        "on conflict (stream) do update set version = max(version, excluded.version)",
        (version,),
    )


# region pull
def pull(api, db, recover=False):
    """Stream style: from the watermark, only orders waiting for the ERP.
    Recover: from zero, only orders not acked at their current version."""
    query = {"filter[ack]": "pending"} if recover else {"filter[status]": "ERP_PENDING"}
    pulled = 0
    for page in api.pages("/orders", query, after=0 if recover else watermark(db)):
        for order in page["data"]:
            stage(db, order)
            pulled += 1
        if not recover:
            save_watermark(db, page["meta"]["next_version_after"])
        db.commit()  # the watermark moves only once its page is staged
    return pulled
# endregion


# region book
def book(api, db, erp):
    result = {"created": 0, "found": 0, "accepted": 0, "rejected": 0, "already_answered": 0}
    staged = db.execute(
        "select distinct order_id, version, document from order_line where booked = 0 order by version"
    ).fetchall()
    for order_id, version, document in staged:
        order = json.loads(document)
        account = order["customer"]["external_ref"] or order["customer"]["code"]
        if fake_erp.on_hold(account):
            answer = {"outcome": "REJECTED", "reason_code": "CUSTOMER_ON_HOLD",
                      "message": f"Customer {account} is on hold in the ERP"}
            result["rejected"] += 1
        else:
            # Find before create: a run that died after creating the sales order but
            # before the ack finds it here instead of creating a twin.
            number = fake_erp.find_sales_order(erp, order["id"])
            if number:
                result["found"] += 1
            else:
                number = fake_erp.create_sales_order(erp, order["id"], account, order["totals"]["gross"])
                result["created"] += 1
            answer = {"outcome": "ACCEPTED", "external_ref": number, "external_number": number}
            result["accepted"] += 1
        try:
            # The key names the order and the version answered, so a retry is a replay
            # and an order sent again (a new version) gets a fresh answer.
            api.post(f"/orders/{order_id}/ack", answer, key=f"{order_id}@{version}")
        except ApiError as refused:
            if refused.code not in ("ALREADY_ACKNOWLEDGED", "ORDER_NOT_AWAITING_ACK"):
                raise
            result["already_answered"] += 1  # answered elsewhere, or no longer waiting
        db.execute("update order_line set booked = 1 where order_id = ? and version = ?", (order_id, version))
        db.commit()
    return result
# endregion


def main(mode):
    api, db, erp = Slflo(), open_staging(), fake_erp.connect()
    summary = {"mode": mode, "pulled": 0}
    try:
        if mode in ("run", "pull", "recover"):
            summary["pulled"] = pull(api, db, recover=mode == "recover")
        if mode in ("run", "book", "recover"):
            summary.update(book(api, db, erp))
    except TokenRefused as stop:
        summary["stopped_by"] = stop.code  # alert someone: retrying won't help
        report(summary)
        return 3
    report(summary)
    return 0


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